diff --git a/README.md b/README.md index a075c36..f539be8 100644 --- a/README.md +++ b/README.md @@ -20,13 +20,15 @@ API Gateway (HTTP API) ├── GET /auth/config ← public, returns Google Client ID │ ├── /pos/** ← Purchase Orders Lambda - │ ├── GET /pos list (paginated) + │ ├── GET /pos list (paginated, auto-syncs from external table) │ ├── POST /pos create │ ├── GET /pos/:id get one │ ├── PUT /pos/:id update │ ├── DELETE /pos/:id delete - │ ├── POST /pos/import bulk import - │ └── GET /pos/dynamo-scan proxy scan of customer's DynamoDB table + │ └── POST /pos/import bulk import + │ + ├── DynamoDB Streams ← PO Sync Lambda (ledgerflow-po-sync) + │ purchase-orders → ledgerflow-pos (near-real-time upsert) │ ├── /invoices/** ← Invoices Lambda │ ├── GET /invoices list @@ -55,7 +57,7 @@ DynamoDB Tables: ledgerflow-settings User/company settings External Tables: - purchase-orders Coupa POs (read-only, ship_to address for EDI) + purchase-orders Coupa POs (DynamoDB Streams → auto-sync to ledgerflow-pos) S3 Buckets (EDI file exchange): ledgerflow-edi-input-{account} @@ -122,21 +124,33 @@ Key variables: 5. Create a **Transformer** for inbound/outbound mapping 6. Copy the Partnership ID and Transformer ID to your `.env` -### 5. Deploy +### 5. Enable DynamoDB Streams (for real-time PO sync) ```bash -# Load env vars and deploy -export $(cat .env | xargs) -npm run deploy +# Enable streams on the external purchase-orders table +aws dynamodb update-table --table-name purchase-orders \ + --stream-specification StreamEnabled=true,StreamViewType=NEW_IMAGE ``` +Copy the stream ARN from the output — you'll pass it during deploy. + +### 6. Deploy + +```bash +# Load env vars and deploy (include stream ARN for real-time sync) +export $(cat .env | xargs) +npx cdk deploy -c purchaseOrdersStreamArn="arn:aws:dynamodb:us-east-1:ACCOUNT:table/purchase-orders/stream/..." +``` + +If you omit the stream ARN, the PO Lambda falls back to on-demand sync (scans the external table on every GET /pos request). + CDK will output your **API Gateway URL**: ``` Outputs: LedgerFlow.ApiUrl = https://abc123.execute-api.us-east-1.amazonaws.com ``` -### 6. Connect the frontend +### 7. Connect the frontend In your `accounting-app.html`, add before ``: @@ -221,7 +235,9 @@ ledgerflow-backend/ │ │ └── index.js │ ├── auth/ Auth endpoints (/auth/me, /auth/config) │ │ └── index.js -│ ├── pos/ Purchase Orders CRUD + DynamoDB import +│ ├── pos/ Purchase Orders CRUD + on-demand sync fallback +│ │ └── index.js +│ ├── po-sync/ DynamoDB Streams handler (purchase-orders → ledgerflow-pos) │ │ └── index.js │ ├── invoices/ Invoices CRUD + PO balance tracking │ │ └── index.js diff --git a/infra/lib/ledgerflow-stack.js b/infra/lib/ledgerflow-stack.js index 1c027f6..4699a14 100644 --- a/infra/lib/ledgerflow-stack.js +++ b/infra/lib/ledgerflow-stack.js @@ -18,6 +18,7 @@ const acm = require("aws-cdk-lib/aws-certificatemanager"); const route53 = require("aws-cdk-lib/aws-route53"); const targets = require("aws-cdk-lib/aws-route53-targets"); const logs = require("aws-cdk-lib/aws-logs"); +const evtSrc = require("aws-cdk-lib/aws-lambda-event-sources"); const path = require("path"); class LedgerFlowStack extends cdk.Stack { @@ -96,6 +97,9 @@ class LedgerFlowStack extends cdk.Stack { lifecycleRules: [{ expiration: cdk.Duration.days(365), id: "expire-old-output" }], }); + // ── External Table Stream ARN (resolved early for use in env + event source) ─ + const poStreamArn = this.node.tryGetContext("purchaseOrdersStreamArn") || process.env.PURCHASE_ORDERS_STREAM_ARN || ""; + // ── Shared Lambda Environment ───────────────────────────────────────────── const commonEnv = { @@ -107,12 +111,13 @@ class LedgerFlowStack extends cdk.Stack { SESSIONS_TABLE: sessionsTable.tableName, SETTINGS_TABLE: settingsTable.tableName, PURCHASE_ORDERS_TABLE: "purchase-orders", + PURCHASE_ORDERS_STREAM_ARN: poStreamArn, EDI_INPUT_BUCKET: ediInputBucket.bucketName, EDI_OUTPUT_BUCKET: ediOutputBucket.bucketName, // Set these via SSM Parameter Store or Secrets Manager in production: // GOOGLE_CLIENT_ID, EDI_PARTNERSHIP_ID, EDI_TRANSFORMER_ID, // EDI_SENDER_ID, EDI_RECEIVER_ID, ALLOWED_DOMAINS, ALLOWED_ORIGIN - GOOGLE_CLIENT_ID: process.env.GOOGLE_CLIENT_ID || "", + GOOGLE_CLIENT_ID: process.env.GOOGLE_CLIENT_ID || "510349952236-eqo5crd0eiae63hifdmu58b0j8qml3ta.apps.googleusercontent.com", ALLOWED_ORIGIN: process.env.ALLOWED_ORIGIN || "*", ALLOWED_DOMAINS: process.env.ALLOWED_DOMAINS || "", EDI_PARTNERSHIP_ID: process.env.EDI_PARTNERSHIP_ID || "", @@ -164,9 +169,40 @@ class LedgerFlowStack extends cdk.Stack { settingsTable.grantReadData(posFn); // External purchase-orders table (Coupa POs) — grant read to POs and EDI lambdas - const purchaseOrdersTable = dynamo.Table.fromTableName(this, "ExternalPOTable", "purchase-orders"); + // Stream ARN is required for the real-time sync Lambda. + // Enable streams on the table: aws dynamodb update-table --table-name purchase-orders --stream-specification StreamEnabled=true,StreamViewType=NEW_IMAGE + const purchaseOrdersTable = poStreamArn + ? dynamo.Table.fromTableAttributes(this, "ExternalPOTable", { + tableName: "purchase-orders", + tableStreamArn: poStreamArn, + }) + : dynamo.Table.fromTableName(this, "ExternalPOTable", "purchase-orders"); purchaseOrdersTable.grantReadData(posFn); + // ── PO Stream Sync Lambda ──────────────────────────────────────────────── + // Triggered by DynamoDB Streams on the external purchase-orders table. + // Upserts POs into ledgerflow-pos in near-real-time. + const poSyncFn = new nodejsFn.NodejsFunction(this, "POSyncFn", { + ...lambdaDefaults, + functionName: "ledgerflow-po-sync", + entry: path.join(__dirname, "../../lambdas/po-sync/index.js"), + handler: "handler", + environment: { ...commonEnv }, + description: "DynamoDB Streams sync — purchase-orders → ledgerflow-pos", + }); + posTable.grantReadWriteData(poSyncFn); + purchaseOrdersTable.grantReadData(poSyncFn); + + if (poStreamArn) { + poSyncFn.addEventSource(new evtSrc.DynamoEventSource(purchaseOrdersTable, { + startingPosition: lambda.StartingPosition.TRIM_HORIZON, + batchSize: 25, + maxBatchingWindow: cdk.Duration.seconds(5), + retryAttempts: 3, + bisectBatchOnError: true, + })); + } + const invoicesFn = new nodejsFn.NodejsFunction(this, "InvoicesFn", { ...lambdaDefaults, functionName: "ledgerflow-invoices", diff --git a/lambdas/po-sync/index.js b/lambdas/po-sync/index.js new file mode 100644 index 0000000..c2d9805 --- /dev/null +++ b/lambdas/po-sync/index.js @@ -0,0 +1,99 @@ +// lambdas/po-sync/index.js +// DynamoDB Streams handler — syncs new/updated POs from the external +// purchase-orders table into ledgerflow-pos in near-real-time. + +const { + GetCommand, PutCommand, UpdateCommand, ScanCommand, +} = require("@aws-sdk/lib-dynamodb"); +const { getDocClient, TABLES, genId } = require("@ledgerflow/shared"); + +const TABLE = TABLES.POS; + +exports.handler = async (event) => { + const db = getDocClient(); + let synced = 0; + let skipped = 0; + + for (const record of event.Records) { + try { + // Only process inserts and modifications + if (record.eventName !== "INSERT" && record.eventName !== "MODIFY") continue; + + const raw = unmarshallImage(record.dynamodb.NewImage); + const poNumber = (raw.po_number || raw.poNumber || "").toString().trim(); + const vendor = (raw.supplier?.name || raw.vendor_name || raw.vendor || "").toString().trim(); + const amount = parseFloat(raw.total_amount || raw.amount || 0); + + if (!poNumber || !vendor) { skipped++; continue; } + + // Check if this PO already exists in our table + const existing = await db.send(new ScanCommand({ + TableName: TABLE, + FilterExpression: "poNumber = :p", + ExpressionAttributeValues: { ":p": poNumber }, + ProjectionExpression: "id, poNumber", + Limit: 1, + })); + + const now = new Date().toISOString(); + const firstLine = Array.isArray(raw.line_items) ? raw.line_items[0] : null; + + if (existing.Count > 0 && record.eventName === "MODIFY") { + // Update existing PO with latest external data + const existingId = existing.Items[0].id; + await db.send(new UpdateCommand({ + TableName: TABLE, + Key: { id: existingId }, + UpdateExpression: "SET vendor = :v, amount = :a, #s = :s, issueDate = :isd, dueDate = :dd, notes = :n, updatedAt = :u", + ExpressionAttributeNames: { "#s": "status" }, + ExpressionAttributeValues: { + ":v": vendor, + ":a": amount, + ":s": (raw.po_status || raw.status || "open").toLowerCase(), + ":isd": raw.order_date || raw.issue_date || raw.issueDate || now.split("T")[0], + ":dd": firstLine?.need_by || raw.due_date || raw.dueDate || null, + ":n": firstLine?.description || raw.notes || raw.description || "", + ":u": now, + }, + })); + synced++; + } else if (existing.Count === 0) { + // Insert new PO + await db.send(new PutCommand({ + TableName: TABLE, + Item: { + id: genId("PO"), + poNumber, + vendor, + amount, + billed: 0, + status: (raw.po_status || raw.status || "open").toLowerCase(), + issueDate: raw.order_date || raw.issue_date || raw.issueDate || now.split("T")[0], + dueDate: firstLine?.need_by || raw.due_date || raw.dueDate || null, + notes: firstLine?.description || raw.notes || raw.description || "", + source: "auto-sync", + createdBy: "stream", + createdAt: now, + updatedAt: now, + }, + })); + synced++; + } else { + skipped++; + } + } catch (err) { + console.error("[po-sync] Error processing record:", err.message, JSON.stringify(record)); + } + } + + console.log(`[po-sync] Processed ${event.Records.length} records: ${synced} synced, ${skipped} skipped`); + return { synced, skipped }; +}; + +// Convert DynamoDB stream image (marshalled) to plain object. +// Stream records use raw DynamoDB attribute format, not DocumentClient format. +function unmarshallImage(image) { + if (!image) return {}; + const { unmarshall } = require("@aws-sdk/util-dynamodb"); + return unmarshall(image); +} diff --git a/lambdas/pos/index.js b/lambdas/pos/index.js index 72a88c6..db194ff 100644 --- a/lambdas/pos/index.js +++ b/lambdas/pos/index.js @@ -126,11 +126,14 @@ async function listPOs(qs, user) { const { limit, cursor } = parsePagination(qs); const db = getDocClient(); - // Auto-sync new POs from external purchase-orders table - try { - await syncFromPurchaseOrders(db, user); - } catch (err) { - console.warn("[auto-sync] failed, returning cached POs:", err.message); + // Sync is now handled by DynamoDB Streams (po-sync Lambda). + // Fall back to on-demand sync only if streams are not configured. + if (!process.env.PURCHASE_ORDERS_STREAM_ARN) { + try { + await syncFromPurchaseOrders(db, user); + } catch (err) { + console.warn("[auto-sync] failed, returning cached POs:", err.message); + } } const params = { diff --git a/template.yaml b/template.yaml index d717d0f..46f8cac 100644 --- a/template.yaml +++ b/template.yaml @@ -244,6 +244,20 @@ Resources: EntryPoints: [index.js] External: [] + POSyncFn: + Type: AWS::Serverless::Function + Properties: + FunctionName: ledgerflow-po-sync + Handler: index.handler + CodeUri: lambdas/po-sync/ + Metadata: + BuildMethod: esbuild + BuildProperties: + Minify: false + Sourcemap: true + EntryPoints: [index.js] + External: [] + SettingsFn: Type: AWS::Serverless::Function Properties: