Skip to content

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


  1. RBAC Architecture
  2. Trigger.dev Python Extension Integration
  3. Apps/Trigger Separation
  4. Monorepo Organization
  5. Code Reuse Strategy
  6. OpenTelemetry Integration
  7. Updated Architecture Diagrams (Mermaid)

// Current: Only org-level isolation
interface ClerkJWT {
sub: string; // User ID
org_id: string; // Organization ID
// ❌ No role information
}
// All users in org have same permissions
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) | βœ“ | βœ“ | βœ— | βœ— | βœ—
packages/auth/src/rbac.ts
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 };
}
packages/data-backend/cube/models/erp_products.yml
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')"
apps/webapp/src/middleware/rbac.ts
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();
});
packages/data-backend/cube/cube.js
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;
}
}

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.

apps/webapp/trigger.config.ts
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',
}),
],
},
});
apps/webapp/src/trigger/python/requirements.txt
dlt[parquet,s3]==0.4.12
pyarrow>=14.0.0
s3fs>=2023.12.0
pandas>=2.0.0
requests>=2.31.0
pydantic>=2.5.0
apps/webapp/src/trigger/python/etl/prophet21_products.py
"""
Prophet21 Products ETL Pipeline
Extracts products from Prophet21 OData API and loads to Hudi lakehouse
"""
import dlt
import requests
import os
from typing import Iterator, Dict, Any
from 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))
apps/trigger/src/tasks/erp/sync-products-v2.ts
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;
}
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)

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:#fff
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.md
{
"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"
}
}
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 production
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 isolated

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 layers
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.md
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)

packages/data-backend/client/src/cube-client.ts
import { logger } from '@repo/observability'; // ← REUSE observability
import 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();
}
}
packages/db/src/schema/hudi-mapping.ts
/**
* 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];
}
packages/observability/src/cube-instrumentation.ts
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();
}
});
}
// packages/db/src/schema/erp-types.ts (EXISTING - extend it)
import { z } from 'zod';
// Add Hudi-specific schemas alongside existing Drizzle schemas
export 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>;
packages/auth/src/middleware/rbac.ts
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.ts
export const prerender = false;
export const middleware = sequence(createRBACMiddleware('erp_data:read'));

From packages/observability/src/otel-config.ts, you already have:

  • βœ… OpenTelemetry SDK setup
  • βœ… Better Stack exporter
  • βœ… Trace context propagation
  • βœ… Structured logging

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:

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:

Terminal window
# .env (root)
OTEL_EXPORTER_OTLP_ENDPOINT=https://in-otel.betterstack.com
BETTER_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 header

How 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:

Terminal window
# 1. Deploy trigger tasks with new config
cd apps/trigger
pnpm trigger:deploy
# 2. Trigger a test task
pnpm trigger:test
# 3. Check Better Stack dashboard
# Traces should appear under service: "erp-unlocked-trigger-tasks"
packages/observability/src/instrumentation/trino.ts
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
}
);
}
packages/data-backend/client/src/instrumentation.ts
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();
}
});
}
}
apps/webapp/src/trigger/python/etl/instrumentation.py
"""
OpenTelemetry instrumentation for Python ETL pipelines
"""
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
import os
# Initialize tracer (reads OTEL_EXPORTER_OTLP_ENDPOINT from env)
trace.set_tracer_provider(TracerProvider())
tracer = trace.get_tracer("dlt-pipeline")
# Export to Better Stack
otlp_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
pass
apps/trigger/src/tasks/erp/sync-products-v2.ts
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;
},
});
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:#fff

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]
end
subgraph "Observability"
N[Better Stack<br/>OpenTelemetry]
end
A -->|JWT auth| F
A -->|orders| C
A -->|PDFs| D
D --> E
E --> M
F -->|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 observability
sequenceDiagram
actor User
participant Clerk
participant Webapp
participant Middleware
participant CubeAPI
participant Trino
participant Hudi
User->>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 context
CubeAPI->>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
end
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/]
end
subgraph "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| H
B -->|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 observability
sequenceDiagram
actor User
participant Frontend
participant CubeClient
participant CubeAPI
participant Trino
participant Hudi
participant BetterStack
User->>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 results
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 --> M
B --> 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 external
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| D
B --> 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 external

  1. 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
  2. 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
  3. 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
  4. 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
  5. 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
  6. 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
  7. 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
  8. 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
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
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,300
3-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/month

Next Steps:

  1. Review RBAC role definitions with product team
  2. Prototype Python extension integration in Trigger.dev
  3. Define package boundaries and dependencies
  4. Extend OpenTelemetry instrumentation
  5. Begin Phase 1 pilot with products entity

References: