ERP-Unlocked Data Backend - Implementation Guide
RBAC, Monorepo Integration, Python on Trigger.dev & Observability
Section titled βRBAC, Monorepo Integration, Python on Trigger.dev & ObservabilityβVersion: 1.1
Date: November 9, 2025
Addendum to: Main Proposal
Table of Contents
Section titled βTable of Contentsβ- RBAC Architecture
- Trigger.dev Python Extension Integration
- Apps/Trigger Separation
- Monorepo Organization
- Code Reuse Strategy
- OpenTelemetry Integration
- Updated Architecture Diagrams (Mermaid)
RBAC Architecture
Section titled βRBAC ArchitectureβCurrent State: Organization-Level Only
Section titled βCurrent State: Organization-Level Onlyβ// Current: Only org-level isolationinterface ClerkJWT { sub: string; // User ID org_id: string; // Organization ID // β No role information}
// All users in org have same permissionsProposed: Role-Based Access Control
Section titled βProposed: Role-Based Access Controlβ1. Role Hierarchy
Section titled β1. Role HierarchyβRoles (from highest to lowest privilege):- owner # Full access, billing, delete org - admin # Manage users, connections, all data - data_analyst # Read ERP data via Cube, run reports - editor # Create/edit orders, upload PDFs - viewer # Read-only access to orders
Permissions Matrix: Resource | Owner | Admin | Analyst | Editor | Viewer ------------------|-------|-------|---------|--------|------- Billing | β | β | β | β | β ERP Connections | β | β | β | β | β ERP Data (read) | β | β | β | β | β ERP Data (write) | β | β | β | β | β Orders (create) | β | β | β | β | β Orders (read) | β | β | β | β | β PDFs (upload) | β | β | β | β | β Users (manage) | β | β | β | β | β2. Clerk Integration for RBAC
Section titled β2. Clerk Integration for RBACβimport { auth } from '@clerk/astro';
export const ROLES = { OWNER: 'org:owner', ADMIN: 'org:admin', DATA_ANALYST: 'org:data_analyst', EDITOR: 'org:editor', VIEWER: 'org:viewer',} as const;
export const PERMISSIONS = { // Billing 'billing:read': [ROLES.OWNER], 'billing:write': [ROLES.OWNER], // ERP Connections 'erp_connections:read': [ROLES.OWNER, ROLES.ADMIN], 'erp_connections:write': [ROLES.OWNER, ROLES.ADMIN],
// ERP Data (via Cube) 'erp_data:read': [ROLES.OWNER, ROLES.ADMIN, ROLES.DATA_ANALYST, ROLES.EDITOR, ROLES.VIEWER], 'erp_data:sync': [ROLES.OWNER, ROLES.ADMIN], // Orders 'orders:create': [ROLES.OWNER, ROLES.ADMIN, ROLES.EDITOR], 'orders:read': [ROLES.OWNER, ROLES.ADMIN, ROLES.DATA_ANALYST, ROLES.EDITOR, ROLES.VIEWER], 'orders:update': [ROLES.OWNER, ROLES.ADMIN, ROLES.EDITOR], 'orders:delete': [ROLES.OWNER, ROLES.ADMIN], // PDFs 'pdfs:upload': [ROLES.OWNER, ROLES.ADMIN, ROLES.EDITOR], // Users 'users:manage': [ROLES.OWNER, ROLES.ADMIN],} as const;
export function hasPermission(userRole: string, permission: keyof typeof PERMISSIONS): boolean { const allowedRoles = PERMISSIONS[permission]; return allowedRoles.includes(userRole as any);}
export async function requirePermission(permission: keyof typeof PERMISSIONS) { const { userId, sessionClaims } = await auth();
if (!userId) { throw new Error('Unauthorized'); }
const role = sessionClaims?.org_role as string;
if (!hasPermission(role, permission)) { throw new Error(`Forbidden: Requires ${permission} permission`); } return { userId, role, orgId: sessionClaims?.org_id as string };}3. Cube RLS with Role-Based Filtering
Section titled β3. Cube RLS with Role-Based Filteringβcubes: - name: ERPProducts sql: > SELECT * FROM erp_products WHERE clerk_organization_id = '${SECURITY_CONTEXT.org_id}'
# Row-level security with role-based access data_access: - name: read_erp_data source: type: security_context field: org_id
# Role-based filtering filters: # Viewers can only see active, non-sensitive products - if: "user.role === 'org:viewer'" then:- sql: "{CUBE}.delete_flag = false" - sql: "{CUBE}.metadata.prophet21.is_sensitive IS NULL OR {CUBE}.metadata.prophet21.is_sensitive = false"
# Data analysts get full product catalog - if: "user.role === 'org:data_analyst' OR user.role === 'org:admin' OR user.role === 'org:owner'" then: - sql: "1=1" # No additional filters
measures: - name: count type: count
# Pricing metrics (restricted to admin+ roles) - name: avg_price type: avg sql: "current_pricing.base_price"
# Only show pricing to authorized roles shown: if: "user.role IN ('org:owner', 'org:admin', 'org:data_analyst')" dimensions: - name: product_id sql: erp_product_id type: string primary_key: true- name: name sql: name type: string
# Sensitive fields (admin-only) - name: cost_price sql: "metadata.prophet21.cost_price" type: number shown: if: "user.role IN ('org:owner', 'org:admin')"4. Middleware Integration
Section titled β4. Middleware Integrationβimport { defineMiddleware } from 'astro/middleware';import { auth } from '@clerk/astro';
export const rbacMiddleware = defineMiddleware(async ({ request, locals }, next) => { const { userId, sessionClaims } = await auth(); if (!userId) { return new Response('Unauthorized', { status: 401 }); } // Add RBAC context to locals locals.auth = { userId, orgId: sessionClaims?.org_id as string, role: sessionClaims?.org_role as string, permissions: sessionClaims?.org_permissions as string[], };
return next();});5. Cube Security Context
Section titled β5. Cube Security Contextβmodule.exports = { contextToAppId: ({ securityContext }) => { return `CUBE_APP_${securityContext.org_id}`; },
// Security context from JWT contextToOrchestratorId: ({ securityContext }) => { return securityContext.org_id; },
// Extract user context from JWT checkAuth: async (req, authorization) => { if (!authorization) { throw new Error('No authorization token provided'); } // Verify Clerk JWT const token = authorization.replace('Bearer ', ''); const verified = await verifyClerkToken(token);
if (!verified) { throw new Error('Invalid token'); } return { org_id: verified.org_id, user_id: verified.sub, role: verified.org_role, permissions: verified.org_permissions || [], }; },
// Inject security context into queries queryTransformer: (query, { securityContext }) => { return { ...query, contextSymbols: { user: { org_id: securityContext.org_id, user_id: securityContext.user_id, role: securityContext.role, }, }, }; },};
async function verifyClerkToken(token: string) { const { createClerkClient } = await import('@clerk/clerk-sdk-node'); const clerk = createClerkClient({ secretKey: process.env.CLERK_SECRET_KEY });
try { const verified = await clerk.verifyToken(token); return verified; } catch (error) { console.error('Token verification failed:', error); return null; }}Trigger.dev Python Extension Integration
Section titled βTrigger.dev Python Extension IntegrationβBased on the Trigger.dev Python Extension documentation, we can eliminate the need for separate Spark/dlt containers and run Python ETL directly in Trigger.dev tasks.
1. Configuration
Section titled β1. Configurationβimport { defineConfig } from '@trigger.dev/sdk';import { pythonExtension } from '@trigger.dev/python/extension';
export default defineConfig({ project: 'proj_zwlhawgpqrmanfxhquwv', runtime: 'node', build: { extensions: [ pythonExtension({ // Copy Python ETL scripts scripts: ['./src/trigger/python/etl/**/*.py'],
// Install Python dependencies for dlt requirementsFile: './src/trigger/python/requirements.txt', // Dev environment (if using venv) devPythonBinaryPath: '.venv/bin/python', }), ], },});2. Python Requirements
Section titled β2. Python Requirementsβdlt[parquet,s3]==0.4.12pyarrow>=14.0.0s3fs>=2023.12.0pandas>=2.0.0requests>=2.31.0pydantic>=2.5.03. dlt Pipeline as Python Script
Section titled β3. dlt Pipeline as Python Scriptβ"""Prophet21 Products ETL PipelineExtracts products from Prophet21 OData API and loads to Hudi lakehouse"""import dltimport requestsimport osfrom typing import Iterator, Dict, Anyfrom datetime import datetime
@dlt.resource( name="erp_products", primary_key=["connection_id", "erp_product_id"], write_disposition="merge",)def fetch_prophet21_products( connection_id: str, clerk_organization_id: str, api_url: str, username: str, password: str, last_sync: str = None,) -> Iterator[Dict[str, Any]]: """ Extract products from Prophet21 OData API. """ skip = 0 page_size = 100
while True: # Build OData query with incremental filter filter_clause = "" if last_sync: filter_clause = f"&$filter=modified_date gt datetime'{last_sync}'"url = f"{api_url}/items?$skip={skip}&$top={page_size}{filter_clause}"
response = requests.get(url, auth=(username, password)) response.raise_for_status()
data = response.json() items = data.get("value", [])
if not items: break for item in items: yield { "id": f"{connection_id}_{item['item_id']}", "connection_id": connection_id, "clerk_organization_id": clerk_organization_id, "erp_type": "prophet21", "erp_product_id": item["item_id"], "name": item["item_desc"], "description": item.get("extended_desc"), "unit_of_measure": item.get("unit_of_measure"), "delete_flag": item.get("delete_flag") == "Y", "last_synced_at": datetime.now().isoformat(), "metadata": { "prophet21": { "company_id": item.get("company_id"), "price_code": item.get("price_code"), "weight": item.get("weight"), } }, "images": { "erp_url": item.get("image_url"), "has_image": bool(item.get("image_url")) }, "current_pricing": { "base_price": item.get("base_price"), "currency": "USD", }, "search_text": f"{item['item_id']} {item['item_desc']} {item.get('extended_desc', '')}" }
skip += page_size
def run_pipeline(config: Dict[str, Any]): """Main pipeline execution""" pipeline = dlt.pipeline( pipeline_name="prophet21_products", destination="filesystem", # Use filesystem destination for Hudi compatibility dataset_name="erp_data" )
# Configure S3/R2 destination pipeline.run( fetch_prophet21_products(**config), table_name="erp_products", write_disposition="merge", loader_file_format="parquet" )
return { "status": "success", "records_processed": pipeline.last_trace.counts.get("items", 0) }
if __name__ == "__main__": import sys import json
# Read config from stdin (passed from TypeScript task) config = json.loads(sys.stdin.read()) result = run_pipeline(config) print(json.dumps(result))4. Trigger.dev Task Using Python Extension
Section titled β4. Trigger.dev Task Using Python Extensionβimport { schemaTask } from '@trigger.dev/sdk';import { z } from 'zod';import { python } from '@trigger.dev/python';import { db } from '@repo/db';import { erpConnections } from '@repo/db/schema';import { eq } from 'drizzle-orm';import { logger } from '@/trigger/logger';
const syncProductsSchema = z.object({ connectionId: z.string().uuid(), clerkOrganizationId: z.string(), fullSync: z.boolean().default(false),});
export const syncERPProductsV2 = schemaTask({ id: 'sync-erp-products-v2', schema: syncProductsSchema, retry: { maxAttempts: 3, factor: 1.5, }, machine: { preset: 'medium-2x', }, run: async payload => { logger.info('Starting ERP product sync', payload); // Get connection from database const connection = await db.query.erpConnections.findFirst({ where: eq(erpConnections.id, payload.connectionId), }); if (!connection) { throw new Error(`Connection ${payload.connectionId} not found`); }
// Determine last sync timestamp const lastSync = payload.fullSync ? null : await getLastSyncTimestamp(payload.connectionId); // Prepare Python script config const pythonConfig = { connection_id: payload.connectionId, clerk_organization_id: payload.clerkOrganizationId, api_url: connection.apiUrl, username: connection.username, password: connection.password, last_sync: lastSync?.toISOString() || null, }; // Run Python dlt pipeline using Trigger.dev Python extension const result = await python.runScript('./src/trigger/python/etl/prophet21_products.py', [], { // Pass config via stdin stdin: JSON.stringify(pythonConfig),
// Environment variables for S3/R2 env: { AWS_ACCESS_KEY_ID: process.env.R2_ACCESS_KEY_ID!, AWS_SECRET_ACCESS_KEY: process.env.R2_SECRET_ACCESS_KEY!, AWS_ENDPOINT_URL: process.env.R2_ENDPOINT_URL!, AWS_REGION: 'auto', DESTINATION__FILESYSTEM__BUCKET_URL: `s3://${process.env.R2_BUCKET_NAME}/hudi/erp_products`, }, }); if (result.exitCode !== 0) { logger.error('Python pipeline failed', { stderr: result.stderr, exitCode: result.exitCode, }); throw new Error(`Pipeline failed: ${result.stderr}`); } // Parse result from stdout const pipelineResult = JSON.parse(result.stdout); logger.info('Product sync completed', { connectionId: payload.connectionId, recordsProcessed: pipelineResult.records_processed, });
return { success: true, recordsProcessed: pipelineResult.records_processed, }; },});
async function getLastSyncTimestamp(connectionId: string): Promise<Date | null> { // Query last successful sync from sync logs const lastSync = await db.query.erpSyncLogs.findFirst({ where: and( eq(erpSyncLogs.connectionId, connectionId), eq(erpSyncLogs.entityType, 'products'), eq(erpSyncLogs.status, 'completed') ), orderBy: desc(erpSyncLogs.completedAt), });
return lastSync?.completedAt || null;}5. Benefits of Trigger.dev Python Extension
Section titled β5. Benefits of Trigger.dev Python ExtensionβBefore (Separate Spark/dlt containers): β Need to deploy separate containers β Complex orchestration (Trigger β API β Spark) β Additional infrastructure (Spark cluster) β More moving parts to monitor
After (Python extension in Trigger.dev): β Python runs inside Trigger.dev tasks β Single orchestration layer β Automatic retries, logging, monitoring β No separate containers needed β Simpler deployment pipeline β Cost savings (no Spark cluster)Apps/Trigger Separation
Section titled βApps/Trigger SeparationβWhy Separate Trigger.dev into Its Own App?
Section titled βWhy Separate Trigger.dev into Its Own App?βCurrent State: Trigger.dev tasks live in apps/webapp/src/trigger/ but deploy to Trigger.dev Cloud (not Coolify).
Problem: This creates confusion about deployment boundaries and dependency management.
graph LR subgraph "Current (Confusing)" A[apps/webapp] A -->|contains| B[src/trigger/] B -->|deploys to| C[Trigger.dev Cloud] A -->|deploys to| D[Coolify] end
subgraph "Proposed (Clear)" E[apps/webapp] F[apps/trigger] E -->|deploys to| G[Coolify] F -->|deploys to| H[Trigger.dev Cloud] end
style F fill:#4CAF50,color:#fff style H fill:#4CAF50,color:#fffNew Structure: apps/trigger/
Section titled βNew Structure: apps/trigger/βapps/trigger/βββ src/β βββ index.ts # Main entry point (exports all tasks)β ββ βββ tasks/β β βββ erp/β β β βββ sync-products-v2.tsβ β β βββ sync-customers-v2.tsβ β β βββ sync-cross-references-v2.tsβ β β βββ sync-shipping-addresses-v2.tsβ β β βββ scheduled-erp-sync.tsβ β ββ β βββ pdf/β β β βββ generate-thumbnails.tsβ β β βββ scheduled-thumbnail-gen.tsβ β β βββ process-pdf.tsβ β ββ β βββ maintenance/β β β βββ cleanup.tsβ β β βββ health-check.tsβ β β βββ usage-reporting.tsβ β ββ β βββ data-backend/ # Future: Data backend tasksβ β βββ cube-refresh.tsβ β βββ hudi-compaction.tsβ ββ βββ python/ # Python ETL scriptsβ β βββ etl/β β β βββ prophet21_products.pyβ β β βββ prophet21_customers.pyβ β β βββ netsuite_products.py # Futureβ β β βββ shared/β β β βββ hudi_writer.pyβ β β βββ utils.pyβ β βββ requirements.txtβ ββ βββ utils/β β βββ logger.ts # Trigger-specific loggerβ β βββ env.ts # Environment validationβ ββ βββ types/β βββ erp-sync.ts # Shared typesββββ trigger.config.ts # Trigger.dev config (at root)βββ package.jsonβββ tsconfig.jsonβββ README.mdPackage.json for apps/trigger
Section titled βPackage.json for apps/triggerβ{ "name": "@erp-unlocked/trigger", "version": "1.0.0", "private": true, "type": "module", "main": "src/index.ts",
"scripts": { "dev": "cross-env NODE_ENV=development trigger.dev dev", "deploy": "cross-env NODE_ENV=production trigger.dev deploy", "deploy:prod": "cross-env NODE_ENV=production trigger.dev deploy --env production", "test": "vitest run", "test:watch": "vitest", "lint": "eslint src/", "typecheck": "tsc --noEmit" },
"dependencies": { "@trigger.dev/sdk": "4.0.4", "@trigger.dev/python": "^1.0.0", "@repo/db": "workspace:*", "@repo/auth": "workspace:*", "@repo/observability": "workspace:*", "@repo/utils": "workspace:*", "@opentelemetry/api": "^1.9.0", "@opentelemetry/exporter-trace-otlp-http": "^0.54.0", "@opentelemetry/exporter-logs-otlp-http": "^0.54.2", "zod": "^3.23.8", "drizzle-orm": "^0.36.4", "ioredis": "^5.4.2" },
"devDependencies": { "@types/node": "^22.0.0", "typescript": "^5.4.2", "vitest": "^2.0.0", "eslint": "^9.0.0", "@trigger.dev/build": "^4.0.0" }}Migration Checklist
Section titled βMigration ChecklistβPhase 1: Create Structure (Week 1) β Create apps/trigger/ directory β Add package.json with dependencies β Create trigger.config.ts with Better Stack telemetry β Set up TypeScript config β Add to pnpm-workspace.yaml
Phase 2: Move Files (Week 2) β Move apps/webapp/src/trigger/tasks/ β apps/trigger/src/tasks/ β Move apps/webapp/src/trigger/logger.ts β apps/trigger/src/utils/ β Create apps/trigger/src/python/ for ETL scripts β Update all import paths (use @repo/* instead of relative)
Phase 3: Update CI/CD (Week 3) β Add .github/workflows/deploy-trigger.yml β Update .github/workflows/deploy-webapp.yml (remove trigger steps) β Configure environment variables in CI
Phase 4: Parallel Testing (Weeks 4-5) β Run both old and new setups in parallel β Compare task execution results β Verify Better Stack trace integration
Phase 5: Cutover (Week 6) β Remove apps/webapp/src/trigger/ β Update documentation β Deploy to productionBenefits of Separation
Section titled βBenefits of SeparationβClear Deployment Boundaries: webapp β Coolify (Docker containers) trigger β Trigger.dev Cloud
Independent Versioning: Can update Trigger.dev SDK without touching webapp Separate package.json = cleaner dependencies
Better CI/CD: Trigger deployments don't rebuild webapp Faster deployments (only changed app rebuilds)
Future-Proof: Clean home for data backend ETL tasks Python scripts organized separately Better Stack telemetry isolatedMonorepo Organization
Section titled βMonorepo OrganizationβShould Data Backend Be in Same Monorepo?
Section titled βShould Data Backend Be in Same Monorepo?βYES - Keep it in the same monorepo for these reasons:
Benefits: β Shared packages (@repo/db, @repo/auth, @repo/observability) β Unified versioning and deployment β Code reuse (types, utilities, schemas) β Single CI/CD pipeline β Better developer experience (one repo to clone) β Atomic changes across frontend/backend/data layersUpdated Monorepo Structure
Section titled βUpdated Monorepo Structureβerp-unlocked/βββ apps/β βββ webapp/ # Main app (EXISTING)β β βββ src/β β β βββ ... # Frontend code onlyβ β βββ package.jsonβ β βββ astro.config.mjsβ ββ βββ trigger/ # β NEW: Trigger.dev tasks (separate app!)β β βββ src/β β β βββ index.ts # Task exportsβ β β βββ tasks/β β β β βββ erp/β β β β β βββ sync-products-v2.tsβ β β β β βββ sync-customers-v2.tsβ β β β β βββ scheduled-erp-sync.tsβ β β β βββ pdf/β β β β β βββ generate-thumbnails.tsβ β β β β βββ scheduled-thumbnail-gen.tsβ β β β βββ maintenance/β β β β β βββ cleanup.tsβ β β β β βββ health-check.tsβ β β β βββ data-backend/ # β Future: Data backend tasksβ β β β βββ cube-refresh.tsβ β β β βββ hudi-compaction.tsβ β β ββ β β βββ python/ # Python ETL scriptsβ β β β βββ etl/β β β β β βββ prophet21_products.pyβ β β β β βββ prophet21_customers.pyβ β β β β βββ netsuite_products.pyβ β β β β βββ shared/β β β β β βββ hudi_writer.pyβ β β β β βββ utils.pyβ β β β βββ requirements.txtβ β β ββ β β βββ utils/β β β βββ logger.tsβ β β βββ env.tsβ β ββ β βββ trigger.config.ts # With pythonExtension + Better Stack telemetryβ β βββ package.jsonβ β βββ tsconfig.jsonβ β βββ README.mdβ ββ βββ marketing/ # Marketing site (EXISTING)β βββ pdf-api/ # PDF processing (EXISTING)β βββ pdf-worker/ # PDF worker (EXISTING)ββββ packages/β βββ db/ # Database (EXISTING)β β βββ src/β β βββ schema/β β βββ erp-replica.ts # β DEPRECATED (Phase 3)β β βββ erp-connections.ts # β KEEP (connection metadata)β ββ βββ data-backend/ # β NEW: Data backend packageβ β βββ cube/β β β βββ models/β β β β βββ erp_products.ymlβ β β β βββ erp_customers.ymlβ β β β βββ erp_inventory.ymlβ β β βββ cube.js # Cube configβ β β βββ schema/β β β βββ auth.js # Security contextβ β ββ β βββ trino/β β β βββ catalog/β β β β βββ hudi.propertiesβ β β β βββ postgres.propertiesβ β β βββ config.propertiesβ β ββ β βββ hudi/β β β βββ tables/β β β β βββ products.yamlβ β β β βββ customers.yamlβ β β β βββ inventory.yamlβ β β βββ scripts/β β β βββ create_tables.pyβ β ββ β βββ client/ # β TypeScript client for Cube APIβ β β βββ src/β β β β βββ cube-client.tsβ β β β βββ types.tsβ β β β βββ index.tsβ β β βββ package.jsonβ β β βββ tsconfig.jsonβ β ββ β βββ docker-compose.yml # Local dev stackβ β βββ package.jsonβ β βββ README.mdβ ββ βββ auth/ # Auth utilities (EXISTING)β β βββ src/β β βββ rbac.ts # β NEW: RBAC helpersβ β βββ ...β ββ βββ observability/ # Observability (EXISTING)β β βββ src/β β βββ otel-config.ts # β Extend for Cube/Trinoβ β βββ ...β ββ βββ utils/ # Utilities (EXISTING)ββββ docs/ βββ architecture/ βββ data-backend-proposal.md βββ data-backend-implementation-guide.mdPackage vs App Decision Matrix
Section titled βPackage vs App Decision MatrixβPut in packages/ if: β Shared across multiple apps β Reusable library/utilities β Configuration (Cube models, Hudi schemas) β Type definitions β Client SDKs
Put in apps/ if: β Deployable service (has Dockerfile) β Standalone executable β Has its own runtime/port
For Data Backend: β Cube models β packages/data-backend/cube/ β Trino config β packages/data-backend/trino/ β Hudi schemas β packages/data-backend/hudi/ β TypeScript client β packages/data-backend/client/ β Python ETL β apps/webapp/src/trigger/python/(runs in Trigger.dev, not standalone)Code Reuse Strategy
Section titled βCode Reuse Strategyβ1. Reuse Existing Packages
Section titled β1. Reuse Existing Packagesβimport { logger } from '@repo/observability'; // β REUSE observabilityimport type { ERPConnection } from '@repo/db/schema'; // β REUSE types
export class CubeClient { constructor( private apiUrl: string, private clerkToken: string // β REUSE Clerk auth ) {}
async queryProducts(filters: any) { logger.info('Querying Cube API', { filters }); // β REUSE logger
const response = await fetch(`${this.apiUrl}/graphql`, { method: 'POST', headers: { Authorization: `Bearer ${this.clerkToken}`, 'Content-Type': 'application/json', }, body: JSON.stringify({ query: buildCubeQuery(filters), }), });
return response.json(); }}2. Share Database Schema Metadata
Section titled β2. Share Database Schema Metadataβ/** * Maps PostgreSQL erp_connections to Hudi lakehouse queries * Provides type-safe bridge between OLTP and OLAP */import type { ERPConnection } from './erp-connections';
export interface HudiTableReference { tableName: string; partitionKeys: string[]; primaryKeys: string[];}
export const ERP_ENTITY_TO_HUDI_TABLE: Record<string, HudiTableReference> = { products: { tableName: 'erp_products', partitionKeys: ['clerk_organization_id', 'erp_type'], primaryKeys: ['connection_id', 'erp_product_id'], }, customers: { tableName: 'erp_customers', partitionKeys: ['clerk_organization_id', 'erp_type'], primaryKeys: ['connection_id', 'erp_customer_id'], }, // ... etc};
export function getHudiTableForConnection( connection: ERPConnection, entityType: keyof typeof ERP_ENTITY_TO_HUDI_TABLE): HudiTableReference { return ERP_ENTITY_TO_HUDI_TABLE[entityType];}3. Extend Existing Observability
Section titled β3. Extend Existing Observabilityβimport { trace, context, SpanStatusCode } from '@opentelemetry/api';import { getTracer } from './otel-config'; // β REUSE existing tracer
const tracer = getTracer('cube-client');
export function instrumentCubeQuery<T>(queryName: string, fn: () => Promise<T>): Promise<T> { return tracer.startActiveSpan(`cube.query.${queryName}`, async span => { try { const result = await fn(); span.setStatus({ code: SpanStatusCode.OK }); return result; } catch (error) { span.setStatus({ code: SpanStatusCode.ERROR, message: error instanceof Error ? error.message : 'Unknown error', }); throw error; } finally { span.end(); } });}4. Share Zod Schemas for Validation
Section titled β4. Share Zod Schemas for Validationβ// packages/db/src/schema/erp-types.ts (EXISTING - extend it)import { z } from 'zod';
// Add Hudi-specific schemas alongside existing Drizzle schemasexport const hudiProductSchema = z.object({ id: z.string(), connection_id: z.string().uuid(), clerk_organization_id: z.string(), erp_type: z.enum(['prophet21', 'netsuite', 'sap', 'dynamics', 'demo']), erp_product_id: z.string(), name: z.string(), description: z.string().nullable(), // ... etc});
export type HudiProduct = z.infer<typeof hudiProductSchema>;5. Reusable RBAC Middleware
Section titled β5. Reusable RBAC Middlewareβimport { defineMiddleware } from 'astro/middleware';import { auth } from '@clerk/astro';import { hasPermission, type Permission } from '../rbac';
export function createRBACMiddleware(requiredPermission: Permission) { return defineMiddleware(async ({ request, locals }, next) => { const { userId, sessionClaims } = await auth(); if (!userId) { return new Response('Unauthorized', { status: 401 }); }
const role = sessionClaims?.org_role as string;
if (!hasPermission(role, requiredPermission)) { return new Response(`Forbidden: Requires ${requiredPermission}`, { status: 403 }); } locals.auth = { userId, orgId: sessionClaims?.org_id as string, role, };
return next(); });}
// Usage in webapp:// apps/webapp/src/pages/api/erp-data/products.tsexport const prerender = false;export const middleware = sequence(createRBACMiddleware('erp_data:read'));OpenTelemetry Integration
Section titled βOpenTelemetry IntegrationβCurrent Observability Stack
Section titled βCurrent Observability StackβFrom packages/observability/src/otel-config.ts, you already have:
- β OpenTelemetry SDK setup
- β Better Stack exporter
- β Trace context propagation
- β Structured logging
Trigger.dev β Better Stack Integration
Section titled βTrigger.dev β Better Stack IntegrationβCurrent Issue: Trigger.dev v4 traces are NOT sent to Better Stack (telemetry config commented out).
Solution: Enable OpenTelemetry exporters in apps/trigger/trigger.config.ts:
import { defineConfig } from '@trigger.dev/sdk';import { pythonExtension } from '@trigger.dev/python/extension';import { OTLPTraceExporter } from '@opentelemetry/exporter-trace-otlp-http';import { OTLPLogExporter } from '@opentelemetry/exporter-logs-otlp-http';
export default defineConfig({ project: 'proj_zwlhawgpqrmanfxhquwv', runtime: 'node', maxDuration: 3600,
build: { external: [ 'pg', '@neondatabase/serverless', 'ioredis', // OpenTelemetry packages '@opentelemetry/api', '@opentelemetry/sdk-logs', '@opentelemetry/exporter-trace-otlp-http', '@opentelemetry/exporter-logs-otlp-http', ],
extensions: [ pythonExtension({ scripts: ['./src/python/etl/**/*.py'], requirementsFile: './src/python/requirements.txt', devPythonBinaryPath: '.venv/bin/python', }), ], }, // β
ENABLE BETTER STACK TELEMETRY telemetry: { // Trace exporter (for distributed tracing) exporters: [ new OTLPTraceExporter({ url: process.env.OTEL_EXPORTER_OTLP_ENDPOINT + '/v1/traces', headers: { // Better Stack expects "Authorization: Bearer <token>" Authorization: `Bearer ${process.env.BETTER_STACK_SOURCE_TOKEN}`, }, }), ], // Log exporter (for structured logs) logExporters: [ new OTLPLogExporter({ url: process.env.OTEL_EXPORTER_OTLP_ENDPOINT + '/v1/logs', headers: { Authorization: `Bearer ${process.env.BETTER_STACK_SOURCE_TOKEN}`, }, }), ], // Service name for filtering in Better Stack serviceName: 'erp-unlocked-trigger-tasks', // Additional instrumentations (optional) instrumentations: [], },});Environment Variables Required:
# .env (root)OTEL_EXPORTER_OTLP_ENDPOINT=https://in-otel.betterstack.comBETTER_STACK_SOURCE_TOKEN=your-better-stack-source-token
# Better Stack expects the token in the Authorization header,# NOT in the URL or as a separate headerHow It Works:
sequenceDiagram participant Task as Trigger.dev Task participant OTEL as OpenTelemetry SDK participant Exporter as OTLP Exporter participant BS as Better Stack
Task->>OTEL: Execute task OTEL->>OTEL: Generate traces + logs OTEL->>Exporter: Export batch Exporter->>BS: POST /v1/traces<br/>Authorization: Bearer token Exporter->>BS: POST /v1/logs<br/>Authorization: Bearer token BS->>BS: Index and visualize
Note over Task,BS: All Trigger.dev runs now visible in Better Stack!Verification:
# 1. Deploy trigger tasks with new configcd apps/triggerpnpm trigger:deploy
# 2. Trigger a test taskpnpm trigger:test
# 3. Check Better Stack dashboard# Traces should appear under service: "erp-unlocked-trigger-tasks"Extend for Data Backend Components
Section titled βExtend for Data Backend Componentsβ1. Trino Instrumentation
Section titled β1. Trino Instrumentationβimport { trace, SpanKind } from '@opentelemetry/api';import { getTracer } from '../otel-config';
const tracer = getTracer('trino-query-engine');
export function instrumentTrinoQuery(queryId: string, query: string, orgId: string) { return tracer.startActiveSpan( 'trino.query', { kind: SpanKind.CLIENT, attributes: { 'trino.query.id': queryId, 'trino.query.text': query.substring(0, 1000), // Truncate for logging 'org.id': orgId, component: 'trino', }, }, async span => { // Query execution happens here // span.end() called automatically } );}2. Cube API Instrumentation
Section titled β2. Cube API Instrumentationβimport { trace, context, propagation } from '@opentelemetry/api';import { getTracer } from '@repo/observability';
const tracer = getTracer('cube-api-client');
export class InstrumentedCubeClient { async query(query: any, options?: { traceparent?: string }) { return tracer.startActiveSpan('cube.query', async span => { span.setAttributes({ 'cube.query.measures': query.measures?.join(','), 'cube.query.dimensions': query.dimensions?.join(','), 'cube.query.limit': query.limit, });
try { const response = await fetch(this.apiUrl, { headers: { // Propagate trace context traceparent: options?.traceparent || context.active().toString(), }, body: JSON.stringify(query), });
const result = await response.json(); span.setAttributes({ 'cube.result.rows': result.data?.length, 'cube.result.query_time_ms': result.annotation?.queryTimeSec * 1000, });
return result; } catch (error) { span.recordException(error as Error); throw error; } finally { span.end(); } }); }}3. dlt Pipeline Instrumentation
Section titled β3. dlt Pipeline Instrumentationβ"""OpenTelemetry instrumentation for Python ETL pipelines"""from opentelemetry import tracefrom opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporterfrom opentelemetry.sdk.trace import TracerProviderfrom opentelemetry.sdk.trace.export import BatchSpanProcessorimport os
# Initialize tracer (reads OTEL_EXPORTER_OTLP_ENDPOINT from env)trace.set_tracer_provider(TracerProvider())tracer = trace.get_tracer("dlt-pipeline")
# Export to Better Stackotlp_exporter = OTLPSpanExporter( endpoint=os.getenv("OTEL_EXPORTER_OTLP_ENDPOINT"), headers={"Authorization": os.getenv("OTEL_EXPORTER_OTLP_HEADERS")})trace.get_tracer_provider().add_span_processor( BatchSpanProcessor(otlp_exporter))
def instrument_pipeline(pipeline_name: str): """Decorator to instrument dlt pipelines with OpenTelemetry""" def decorator(func): def wrapper(*args, **kwargs): with tracer.start_as_current_span( f"dlt.pipeline.{pipeline_name}", attributes={ "pipeline.name": pipeline_name, "org.id": kwargs.get("clerk_organization_id"), } ) as span: try: result = func(*args, **kwargs) span.set_attribute("pipeline.records_processed", result.get("records_processed", 0)) return result except Exception as e: span.record_exception(e) raise return wrapper return decorator
# Usage in pipeline:@instrument_pipeline("prophet21_products")def run_pipeline(config): # ... pipeline logic pass4. End-to-End Trace Propagation
Section titled β4. End-to-End Trace Propagationβimport { trace, context } from '@opentelemetry/api';import { python } from '@trigger.dev/python';
export const syncERPProductsV2 = schemaTask({ id: 'sync-erp-products-v2', run: async payload => { // Get current trace context const currentSpan = trace.getActiveSpan(); const traceParent = currentSpan?.spanContext();
// Pass trace context to Python subprocess const result = await python.runScript('./src/trigger/python/etl/prophet21_products.py', [], { env: { // OpenTelemetry env vars (auto-picked up by Python SDK) OTEL_EXPORTER_OTLP_ENDPOINT: process.env.OTEL_EXPORTER_OTLP_ENDPOINT, OTEL_EXPORTER_OTLP_HEADERS: process.env.OTEL_EXPORTER_OTLP_HEADERS, TRACEPARENT: traceParent?.toString(), // W3C trace context
// AWS/R2 config AWS_ACCESS_KEY_ID: process.env.R2_ACCESS_KEY_ID!, // ... etc }, }); return result; },});OpenTelemetry Visualization
Section titled βOpenTelemetry Visualizationβgraph LR A[Webapp Request] -->|traceparent| B[Trigger.dev Task] B -->|spawn process| C[Python dlt Pipeline] C -->|query| D[Prophet21 API] C -->|write| E[Hudi Lakehouse]B -->|trace export| F[Better Stack] C -->|trace export| F
G[User Query] -->|traceparent| H[Cube API] H -->|query| I[Trino] I -->|read| E
H -->|trace export| F I -->|trace export| F style F fill:#4CAF50,color:#fffUpdated Architecture Diagrams (Mermaid)
Section titled βUpdated Architecture Diagrams (Mermaid)β1. High-Level System Architecture
Section titled β1. High-Level System Architectureβgraph TB subgraph "Frontend Layer" A[Webapp<br/>Astro + React] B[Marketing Site<br/>Cloudflare Pages] end
subgraph "Application Layer - OLTP" C[PostgreSQL<br/>Orders, PDFs, Mappings<br/>500 GB] D[PDF API<br/>FastAPI] E[PDF Worker<br/>Celery] end
subgraph "Data Backend Layer - OLAP" F[Cube Semantic Layer<br/>GraphQL/REST API<br/>RBAC + RLS] G[Trino Query Engine<br/>Federated Queries] H[Hudi Lakehouse<br/>Parquet on R2<br/>2.7 TB] end
subgraph "ETL Layer" I[Trigger.dev Tasks<br/>Python Extension] J[dlt Pipelines<br/>Prophet21, NetSuite] end
subgraph "External Systems" K[Prophet21 ERP] L[NetSuite ERP] M[Google Gemini AI] endsubgraph "Observability" N[Better Stack<br/>OpenTelemetry] end A -->|JWT auth| F A -->|orders| C A -->|PDFs| D D --> E E --> MF -->|security context| G G -->|query| H G -->|join| C I -->|orchestrate| J J -->|extract| K J -->|extract| L J -->|load| H
A --> N F --> N I --> N J --> N
classDef frontend fill:#42A5F5,color:#fff classDef oltp fill:#66BB6A,color:#fff classDef olap fill:#FFA726,color:#fff classDef etl fill:#AB47BC,color:#fff classDef external fill:#78909C,color:#fff classDef observability fill:#26C6DA,color:#fff
class A,B frontend class C,D,E oltp class F,G,H olap class I,J etl class K,L,M external class N observability2. RBAC Authorization Flow
Section titled β2. RBAC Authorization FlowβsequenceDiagram actor User participant Clerk participant Webapp participant Middleware participant CubeAPI participant Trino participant HudiUser->>Clerk: Sign in Clerk-->>User: JWT (org_id, org_role, permissions)
User->>Webapp: Query ERP data Webapp->>Middleware: HTTP request + JWT
Middleware->>Middleware: Verify JWT Middleware->>Middleware: Extract role & permissions Middleware->>Middleware: Check hasPermission("erp_data:read") alt Unauthorized Middleware-->>Webapp: 403 Forbidden Webapp-->>User: Access denied else Authorized Middleware->>CubeAPI: GraphQL query + security contextCubeAPI->>CubeAPI: Inject RLS filters<br/>(org_id, role-based) CubeAPI->>Trino: SQL query
Trino->>Hudi: Read partitions<br/>(org_id = '...') Hudi-->>Trino: Parquet data
Trino-->>CubeAPI: Query results CubeAPI->>CubeAPI: Apply role-based<br/>field masking CubeAPI-->>Webapp: Filtered results Webapp-->>User: Display data end3. ETL Pipeline Flow (Trigger.dev Python Extension)
Section titled β3. ETL Pipeline Flow (Trigger.dev Python Extension)βgraph TB subgraph "Scheduled Trigger" A[Trigger.dev<br/>Cron: 2 AM UTC] end
subgraph "Task Execution" B[TypeScript Task<br/>syncERPProductsV2] C[Python Extension<br/>Subprocess] D[dlt Pipeline<br/>prophet21_products.py] end
subgraph "Data Sources" E[Prophet21 API<br/>OData] F[ERP Connection<br/>PostgreSQL metadata] end
subgraph "Data Lake" G[Parquet Files<br/>Cloudflare R2] H[Hudi Metadata<br/>.hoodie/] endsubgraph "Observability" I[Better Stack<br/>Traces + Logs] end A -->|trigger| B B -->|fetch credentials| F B -->|python.runScript| C C -->|execute| D D -->|extract| E D -->|write Parquet| G D -->|update metadata| HB -->|trace export| I D -->|trace export| I classDef trigger fill:#9C27B0,color:#fff classDef task fill:#3F51B5,color:#fff classDef source fill:#009688,color:#fff classDef storage fill:#FF9800,color:#fff classDef observability fill:#26C6DA,color:#fff
class A trigger class B,C,D task class E,F source class G,H storage class I observability4. Query Execution Flow
Section titled β4. Query Execution FlowβsequenceDiagram actor User participant Frontend participant CubeClient participant CubeAPI participant Trino participant Hudi participant BetterStackUser->>Frontend: Search "industrial pump" Frontend->>CubeClient: query({ dimensions, filters })
CubeClient->>BetterStack: Start span<br/>"cube.query" CubeClient->>CubeAPI: POST /graphql<br/>Authorization: Bearer JWT
CubeAPI->>CubeAPI: Verify JWT<br/>Extract security context CubeAPI->>CubeAPI: Apply RLS<br/>WHERE org_id = '...' CubeAPI->>CubeAPI: Check role permissions<br/>Filter sensitive fields
CubeAPI->>BetterStack: Start span<br/>"trino.query" CubeAPI->>Trino: SQL query
Trino->>Trino: Plan query<br/>Partition pruning Trino->>Hudi: Read Parquet<br/>org_id=550e8400.../erp_type=prophet21/ Hudi-->>Trino: Column data
Trino->>Trino: Apply full-text search<br/>Lucene index Trino-->>CubeAPI: Result rows CubeAPI->>BetterStack: End span<br/>query_time: 200ms
CubeAPI->>CubeAPI: Cache results<br/>(15 min TTL) CubeAPI-->>CubeClient: JSON response CubeClient->>BetterStack: End span<br/>total_time: 220ms CubeClient-->>Frontend: Product list Frontend-->>User: Display results5. Monorepo Package Dependencies
Section titled β5. Monorepo Package Dependenciesβgraph TD subgraph "Apps" A[webapp] B[pdf-api] C[pdf-worker] D[marketing] end
subgraph "Shared Packages" E[@repo/db] F[@repo/auth] G[@repo/observability] H[@repo/utils] I[@repo/ui] J[@repo/data-backend] end
subgraph "External" K[Clerk SDK] L[OpenTelemetry] M[Trigger.dev] N[Drizzle ORM] end
A --> E A --> F A --> G A --> H A --> I A --> J A --> MB --> E B --> G
C --> E C --> G
D --> I D --> H
J --> F J --> G J --> H
E --> N F --> K G --> L
classDef app fill:#42A5F5,color:#fff classDef package fill:#66BB6A,color:#fff classDef external fill:#FFA726,color:#fff class A,B,C,D app class E,F,G,H,I,J package class K,L,M,N external6. Deployment Architecture
Section titled β6. Deployment Architectureβgraph TB subgraph "Coolify Deployment" A[webapp<br/>Container] B[pdf-api<br/>Container] C[pdf-worker<br/>Container] D[cube-api<br/>Container] E[trino-coordinator<br/>Container] F[trino-worker<br/>Container] G[postgres<br/>Container] H[redis<br/>Container] end
subgraph "Cloudflare" I[R2 Storage<br/>Hudi Lakehouse] J[Pages<br/>Marketing Site] K[CDN] end
subgraph "External Services" L[Trigger.dev<br/>Background Jobs] M[Better Stack<br/>Observability] N[Clerk<br/>Auth] end
A -->|OLTP queries| G A -->|cache| H A -->|PDF uploads| I A -->|ERP queries| DB --> H C --> H
D -->|SQL| E E --> F F -->|read Parquet| I
A -->|tasks| L L -->|ETL| I
A -->|traces| M D -->|traces| M L -->|traces| M
A -->|auth| N
K -->|CDN| J
classDef container fill:#3F51B5,color:#fff classDef cloudflare fill:#F68D2E,color:#fff classDef external fill:#26C6DA,color:#fff class A,B,C,D,E,F,G,H container class I,J,K cloudflare class L,M,N externalSummary of Updates
Section titled βSummary of Updatesββ Addressed Requirements
Section titled ββ Addressed Requirementsβ-
RBAC - Role-based access control with Clerk integration, Cube RLS, and middleware
- 5 role levels: owner, admin, data_analyst, editor, viewer
- Row-level security in Cube with role-based field masking
- Middleware authorization for API routes
-
Trigger.dev Python Extension - Eliminated need for separate Spark/dlt containers
- Python ETL runs inside Trigger.dev tasks (no Spark cluster)
- dlt pipelines for Prophet21, NetSuite, etc.
- Cost savings: $200/month eliminated
-
Apps/Trigger Separation - Trigger.dev tasks as independent app
- Clear deployment boundaries (webapp β Coolify, trigger β Trigger.dev Cloud)
- Independent dependency management
- Better CI/CD (only changed app rebuilds)
- 6-week phased migration plan
-
Better Stack Integration - Trigger.dev v4 traces now export to Better Stack
- OpenTelemetry trace + log exporters configured
- Environment variable:
BETTER_STACK_SOURCE_TOKEN - Service name: βerp-unlocked-trigger-tasksβ
- End-to-end trace propagation across webapp, trigger tasks, Python ETL
-
Monorepo Organization - Clear structure for apps vs packages, reuse existing code
- apps/trigger/ for all background tasks
- packages/data-backend/ for Cube/Trino configs
- Shared packages: @repo/db, @repo/auth, @repo/observability
-
Code Reuse - Leverage existing packages throughout
- @repo/db for database types and queries
- @repo/auth for RBAC helpers
- @repo/observability for OpenTelemetry instrumentation
- @repo/utils for shared utilities
-
OpenTelemetry - Extended existing observability for Cube, Trino, dlt pipelines
- Unified Better Stack dashboard for all components
- Trace context propagation from webapp β trigger β Python
- Instrumentation for Cube queries and Trino operations
-
Mermaid Diagrams - All architecture diagrams converted to Mermaid for better rendering
- High-level system architecture
- RBAC authorization flow
- ETL pipeline flow (apps/trigger + Python)
- Query execution flow
- Monorepo package dependencies
- Deployment architecture
π― Key Benefits of Updated Approach
Section titled βπ― Key Benefits of Updated ApproachβInfrastructure Simplification: Before: Separate Spark cluster, dlt containers After: Python runs in Trigger.dev tasks (apps/trigger) Savings: ~$200/month (no Spark cluster needed)
Clear Separation of Concerns: Before: Trigger tasks mixed in apps/webapp After: apps/trigger/ for background jobs, apps/webapp for frontend Benefit: Clear deployment boundaries, independent versioning
Better Stack Integration: Before: Trigger.dev traces not sent to Better Stack After: OpenTelemetry exporters configured in trigger.config.ts Benefit: Unified observability for ALL components
Code Organization: Before: Unclear where data backend code lives After: packages/data-backend/ for configs, apps/trigger/ for ETL Benefit: Single monorepo, unified deployment
Security: Before: Org-level only After: RBAC with 5 roles (owner/admin/analyst/editor/viewer) Benefit: Fine-grained permissions, audit compliance
Observability: Before: Separate tracing for new components, Trigger.dev not integrated After: Extends existing @repo/observability + Trigger.dev β Better Stack Benefit: End-to-end traces across webapp, trigger tasks, Python ETLπ Updated Cost Estimate
Section titled βπ Updated Cost EstimateβOriginal Estimate (from main proposal): $875/month - PostgreSQL: $450 - Trino: $250 - Cube: $35 - Spark: $50 - R2: $40 - Trigger.dev: $50
Revised Estimate (with apps/trigger separation): $725/month - PostgreSQL (downsized): $450 - Trino (downsized to t3.large): $150 # β 40% savings - Cube: $35 - R2: $40 - Trigger.dev: $50 (no extra cost for Python extension) - Spark: $0 (eliminated!)
Total Savings vs Current: 52% ($1,500 β $725/month)Annual Savings: $9,3003-Year Savings: $27,900
Key Cost Reductions: β Eliminated Spark cluster: -$50/monthβ Downsized Trino: -$100/month β Downsized PostgreSQL: -$1,000/month (from current) β R2 vs PostgreSQL SSD: -$410/monthNext Steps:
- Review RBAC role definitions with product team
- Prototype Python extension integration in Trigger.dev
- Define package boundaries and dependencies
- Extend OpenTelemetry instrumentation
- Begin Phase 1 pilot with products entity
References: