From 38547505978d35e064aa396c6397ad435c299481 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Fri, 3 Apr 2026 13:58:15 -0400 Subject: [PATCH 1/2] Add PO auto-sync from external purchase-orders table POs now auto-sync from the external DynamoDB purchase-orders table on every GET /pos request, replacing manual import. Adds cursor-based pagination support and SAM template for local development. Co-Authored-By: Claude Opus 4.6 --- .gitignore | 3 + infra/lib/ledgerflow-stack.js | 1 + lambdas/pos/index.js | 136 ++++++++++------- package.json | 5 +- shared/index.js | 11 +- template.yaml | 277 ++++++++++++++++++++++++++++++++++ 6 files changed, 376 insertions(+), 57 deletions(-) create mode 100644 template.yaml diff --git a/.gitignore b/.gitignore index b23b8eb..e73e16e 100644 --- a/.gitignore +++ b/.gitignore @@ -1,5 +1,8 @@ node_modules/ cdk.out/ +.aws-sam/ .env +env.json +env.local.json *.js.map .DS_Store diff --git a/infra/lib/ledgerflow-stack.js b/infra/lib/ledgerflow-stack.js index 7e0c839..1c027f6 100644 --- a/infra/lib/ledgerflow-stack.js +++ b/infra/lib/ledgerflow-stack.js @@ -106,6 +106,7 @@ class LedgerFlowStack extends cdk.Stack { EDI_TX_TABLE: ediTxTable.tableName, SESSIONS_TABLE: sessionsTable.tableName, SETTINGS_TABLE: settingsTable.tableName, + PURCHASE_ORDERS_TABLE: "purchase-orders", EDI_INPUT_BUCKET: ediInputBucket.bucketName, EDI_OUTPUT_BUCKET: ediOutputBucket.bucketName, // Set these via SSM Parameter Store or Secrets Manager in production: diff --git a/lambdas/pos/index.js b/lambdas/pos/index.js index 2c2f9d0..72a88c6 100644 --- a/lambdas/pos/index.js +++ b/lambdas/pos/index.js @@ -1,19 +1,18 @@ // lambdas/pos/index.js // Purchase Orders API -// GET /pos — list all POs (paginated, filterable) +// GET /pos — list all POs (auto-syncs from purchase-orders table) // POST /pos — create a single PO manually // GET /pos/:id — get one PO // PUT /pos/:id — full update // PATCH /pos/:id — partial update (e.g. status change) // DELETE /pos/:id — delete -// POST /pos/import — bulk import (from DynamoDB scan proxy or pasted JSON) -// GET /pos/dynamo-scan — live scan of the customer's own DynamoDB PO table +// POST /pos/import — bulk import (from pasted JSON) const { ScanCommand, GetCommand, PutCommand, UpdateCommand, DeleteCommand, } = require("@aws-sdk/lib-dynamodb"); const { - ok, created, noContent, badRequest, notFound, conflict, serverError, + ok, created, noContent, badRequest, notFound, conflict, getDocClient, TABLES, genId, parseBody, require_fields, handler, parsePagination, paginatedResponse, decodeCursor, } = require("@ledgerflow/shared"); @@ -32,8 +31,6 @@ exports.handler = handler(async (event, _ctx, user) => { // POST /pos/import if (method === "POST" && id === "import") return bulkImport(event, user); - // GET /pos/dynamo-scan - if (method === "GET" && id === "dynamo-scan") return dynamoScan(event, user); if (!id) { if (method === "GET") return listPOs(qs, user); @@ -48,12 +45,94 @@ exports.handler = handler(async (event, _ctx, user) => { return { statusCode: 405, body: JSON.stringify({ error: "Method Not Allowed" }) }; }); +// ─── Auto-Sync from External purchase-orders Table ────────────────────────── +// Scans the external purchase-orders table and upserts any new POs into +// the internal ledgerflow-pos table. Runs before every list to keep in sync. + +async function syncFromPurchaseOrders(db, user) { + const EXTERNAL = TABLES.PURCHASE_ORDERS; + const now = new Date().toISOString(); + + // Scan all items from the external table (paginate through all pages) + let externalItems = []; + let lastKey; + do { + const params = { TableName: EXTERNAL, ...(lastKey && { ExclusiveStartKey: lastKey }) }; + const result = await db.send(new ScanCommand(params)); + externalItems = externalItems.concat(result.Items || []); + lastKey = result.LastEvaluatedKey; + } while (lastKey); + + if (externalItems.length === 0) return; + + // Get all existing PO numbers to avoid duplicates + let existingNumbers = new Set(); + let lek; + do { + const scan = await db.send(new ScanCommand({ + TableName: TABLE, + ProjectionExpression: "poNumber", + ...(lek && { ExclusiveStartKey: lek }), + })); + (scan.Items || []).forEach(i => existingNumbers.add(i.poNumber)); + lek = scan.LastEvaluatedKey; + } while (lek); + + // Insert any POs that don't already exist + const newItems = externalItems.filter(item => { + const poNumber = (item.po_number || item.poNumber || "").toString().trim(); + return poNumber && !existingNumbers.has(poNumber); + }); + + const batchSize = 25; + for (let i = 0; i < newItems.length; i += batchSize) { + const batch = newItems.slice(i, i + batchSize); + await Promise.all(batch.map(async (raw) => { + try { + 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) return; + + const firstLine = Array.isArray(raw.line_items) ? raw.line_items[0] : null; + + const 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: user?.email || "system", + createdAt: now, + updatedAt: now, + }; + + await db.send(new PutCommand({ TableName: TABLE, Item: item })); + } catch (err) { + console.warn("[auto-sync] skipped item:", err.message); + } + })); + } +} + // ─── List ───────────────────────────────────────────────────────────────────── 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); + } + const params = { TableName: TABLE, Limit: limit, @@ -238,48 +317,3 @@ async function bulkImport(event, user) { return ok({ message: `Import complete`, ...results }); } -// ─── Live DynamoDB Scan (Customer's Own Table) ─────────────────────────────── -// Proxies a scan against the customer's configured DynamoDB table. -// The customer's table name/field config is stored in their settings. - -async function dynamoScan(event, user) { - const qs = event.queryStringParameters || {}; - const filterStatus = qs.status; - const limit = Math.min(parseInt(qs.limit || "50"), 200); - - // Retrieve customer's DynamoDB config from the settings table - const db = getDocClient(); - const settingsResult = await db.send(new GetCommand({ TableName: TABLES.SETTINGS, Key: { userId: user.userId } })); - const dynamoCfg = settingsResult.Item?.config?.dynamo; - - if (!dynamoCfg?.table) return badRequest("No DynamoDB table configured — update your settings first"); - - const params = { TableName: dynamoCfg.table, Limit: limit }; - - if (filterStatus) { - params.FilterExpression = "#s = :s"; - params.ExpressionAttributeNames = { "#s": "status" }; - params.ExpressionAttributeValues = { ":s": filterStatus }; - } - - try { - const result = await db.send(new ScanCommand(params)); - - // Map customer's field names to LedgerFlow's expected shape - const fields = dynamoCfg.fields || {}; - const items = (result.Items || []).map(item => ({ - poNumber: item[fields.poNumber || "po_number"] || "", - vendor: item[fields.vendor || "vendor_name"] || "", - amount: parseFloat(item[fields.amount || "total_amount"]) || 0, - status: item[fields.status || "status"] || "open", - issueDate: item[fields.issueDate || "issue_date"] || "", - dueDate: item[fields.dueDate || "due_date"] || "", - _raw: item, - })); - - return ok({ items, count: result.Count }); - } catch (err) { - if (err.name === "ResourceNotFoundException") return notFound(`Table ${dynamoCfg.table}`); - return serverError(`DynamoDB scan failed: ${err.message}`, err); - } -} diff --git a/package.json b/package.json index ef04aae..edc473d 100644 --- a/package.json +++ b/package.json @@ -12,7 +12,10 @@ "deploy": "cd infra && npx aws-cdk deploy --all", "deploy:dev": "cd infra && npx aws-cdk deploy --all --context env=dev", "destroy": "cd infra && npx aws-cdk destroy --all", - "test": "jest --passWithNoTests" + "test": "jest --passWithNoTests", + "local": "sam local start-api --env-vars env.local.json --warm-containers EAGER --port 3001", + "local:build": "sam build", + "local:invoke": "sam local invoke --env-vars env.local.json" }, "devDependencies": { "@types/aws-lambda": "^8.10.136", diff --git a/shared/index.js b/shared/index.js index 4644529..4b1b7d3 100644 --- a/shared/index.js +++ b/shared/index.js @@ -112,11 +112,12 @@ function getDocClient() { // ─── Table Names ───────────────────────────────────────────────────────────── const TABLES = { - POS: process.env.POS_TABLE || "ledgerflow-pos", - INVOICES: process.env.INVOICES_TABLE || "ledgerflow-invoices", - EDI_TX: process.env.EDI_TX_TABLE || "ledgerflow-edi-transactions", - SESSIONS: process.env.SESSIONS_TABLE || "ledgerflow-sessions", - SETTINGS: process.env.SETTINGS_TABLE || "ledgerflow-settings", + POS: process.env.POS_TABLE || "ledgerflow-pos", + PURCHASE_ORDERS: process.env.PURCHASE_ORDERS_TABLE || "purchase-orders", + INVOICES: process.env.INVOICES_TABLE || "ledgerflow-invoices", + EDI_TX: process.env.EDI_TX_TABLE || "ledgerflow-edi-transactions", + SESSIONS: process.env.SESSIONS_TABLE || "ledgerflow-sessions", + SETTINGS: process.env.SETTINGS_TABLE || "ledgerflow-settings", }; // ─── Pagination Helper ─────────────────────────────────────────────────────── diff --git a/template.yaml b/template.yaml new file mode 100644 index 0000000..d717d0f --- /dev/null +++ b/template.yaml @@ -0,0 +1,277 @@ +AWSTemplateFormatVersion: "2010-09-09" +Transform: AWS::Serverless-2016-10-31 +Description: LedgerFlow B2B Accounting — local SAM development + +Globals: + Function: + Runtime: nodejs20.x + Architectures: [arm64] + MemorySize: 512 + Timeout: 30 + Environment: + Variables: + NODE_ENV: dev + AWS_NODEJS_CONNECTION_REUSE_ENABLED: "1" + POS_TABLE: ledgerflow-pos + INVOICES_TABLE: ledgerflow-invoices + EDI_TX_TABLE: ledgerflow-edi-transactions + SESSIONS_TABLE: ledgerflow-sessions + SETTINGS_TABLE: ledgerflow-settings + PURCHASE_ORDERS_TABLE: purchase-orders + EDI_INPUT_BUCKET: !Sub "ledgerflow-edi-input-${AWS::AccountId}" + EDI_OUTPUT_BUCKET: !Sub "ledgerflow-edi-output-${AWS::AccountId}" + GOOGLE_CLIENT_ID: "" + ALLOWED_ORIGIN: "*" + ALLOWED_DOMAINS: "" + EDI_PARTNERSHIP_ID: "" + EDI_TRANSFORMER_ID: "" + EDI_SENDER_ID: "" + EDI_RECEIVER_ID: "" + +Resources: + # ── API Gateway ────────────────────────────────────────────────────────────── + + LedgerFlowAPI: + Type: AWS::Serverless::HttpApi + Properties: + StageName: $default + CorsConfiguration: + AllowOrigins: + - "*" + - "http://localhost:3000" + AllowMethods: + - GET + - POST + - PUT + - PATCH + - DELETE + - OPTIONS + AllowHeaders: + - Content-Type + - Authorization + MaxAge: 86400 + Auth: + DefaultAuthorizer: GoogleJWT + Authorizers: + GoogleJWT: + AuthorizationScopes: [] + FunctionArn: !GetAtt AuthorizerFn.Arn + FunctionInvokeRole: !GetAtt AuthorizerInvokeRole.Arn + Identity: + Headers: + - Authorization + AuthorizerPayloadFormatVersion: "2.0" + EnableSimpleResponses: true + + # IAM role allowing API Gateway to invoke the authorizer Lambda + AuthorizerInvokeRole: + Type: AWS::IAM::Role + Properties: + AssumeRolePolicyDocument: + Version: "2012-10-17" + Statement: + - Effect: Allow + Principal: + Service: apigateway.amazonaws.com + Action: sts:AssumeRole + Policies: + - PolicyName: InvokeAuthorizerFn + PolicyDocument: + Version: "2012-10-17" + Statement: + - Effect: Allow + Action: lambda:InvokeFunction + Resource: !GetAtt AuthorizerFn.Arn + + # ── Lambda Functions ───────────────────────────────────────────────────────── + + AuthorizerFn: + Type: AWS::Serverless::Function + Properties: + FunctionName: ledgerflow-authorizer + Handler: index.handler + CodeUri: lambdas/authorizer/ + Metadata: + BuildMethod: esbuild + BuildProperties: + Minify: false + Sourcemap: true + EntryPoints: [index.js] + External: [] + + AuthFn: + Type: AWS::Serverless::Function + Properties: + FunctionName: ledgerflow-auth + Handler: index.handler + CodeUri: lambdas/auth/ + Events: + AuthConfig: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /auth/config + Method: ANY + Auth: + Authorizer: NONE + AuthConfigProxy: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /auth/config/{proxy+} + Method: ANY + Auth: + Authorizer: NONE + AuthMe: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /auth/me + Method: ANY + Auth: + Authorizer: NONE + AuthMeProxy: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /auth/me/{proxy+} + Method: ANY + Auth: + Authorizer: NONE + AuthLogout: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /auth/logout + Method: ANY + Auth: + Authorizer: NONE + AuthLogoutProxy: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /auth/logout/{proxy+} + Method: ANY + Auth: + Authorizer: NONE + Metadata: + BuildMethod: esbuild + BuildProperties: + Minify: false + Sourcemap: true + EntryPoints: [index.js] + External: [] + + POsFn: + Type: AWS::Serverless::Function + Properties: + FunctionName: ledgerflow-pos + Handler: index.handler + CodeUri: lambdas/pos/ + Events: + POs: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /pos + Method: ANY + POsProxy: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /pos/{proxy+} + Method: ANY + Metadata: + BuildMethod: esbuild + BuildProperties: + Minify: false + Sourcemap: true + EntryPoints: [index.js] + External: [] + + InvoicesFn: + Type: AWS::Serverless::Function + Properties: + FunctionName: ledgerflow-invoices + Handler: index.handler + CodeUri: lambdas/invoices/ + Events: + Invoices: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /invoices + Method: ANY + InvoicesProxy: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /invoices/{proxy+} + Method: ANY + Metadata: + BuildMethod: esbuild + BuildProperties: + Minify: false + Sourcemap: true + EntryPoints: [index.js] + External: [] + + EDIFn: + Type: AWS::Serverless::Function + Properties: + FunctionName: ledgerflow-edi + Handler: index.handler + CodeUri: lambdas/edi/ + Timeout: 60 + Events: + EDI: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /edi + Method: ANY + EDIProxy: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /edi/{proxy+} + Method: ANY + Metadata: + BuildMethod: esbuild + BuildProperties: + Minify: false + Sourcemap: true + EntryPoints: [index.js] + External: [] + + SettingsFn: + Type: AWS::Serverless::Function + Properties: + FunctionName: ledgerflow-settings + Handler: index.handler + CodeUri: lambdas/settings/ + Events: + Settings: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /settings + Method: ANY + SettingsProxy: + Type: HttpApi + Properties: + ApiId: !Ref LedgerFlowAPI + Path: /settings/{proxy+} + Method: ANY + Metadata: + BuildMethod: esbuild + BuildProperties: + Minify: false + Sourcemap: true + EntryPoints: [index.js] + External: [] + +Outputs: + ApiUrl: + Description: Local API Gateway endpoint + Value: !Sub "https://${LedgerFlowAPI}.execute-api.${AWS::Region}.amazonaws.com" From 102999165f56866749e504ee48e948a348acbfcd Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Fri, 3 Apr 2026 14:29:32 -0400 Subject: [PATCH 2/2] Add DynamoDB Streams real-time PO sync and hardcode Google Client ID - New po-sync Lambda triggered by DynamoDB Streams on the external purchase-orders table for near-real-time upsert into ledgerflow-pos - GET /pos skips blocking sync when streams are configured - Hardcode GOOGLE_CLIENT_ID fallback in CDK stack so /auth/config and JWT verification work without env var Co-Authored-By: Claude Opus 4.6 --- README.md | 36 +++++++++---- infra/lib/ledgerflow-stack.js | 40 +++++++++++++- lambdas/po-sync/index.js | 99 +++++++++++++++++++++++++++++++++++ lambdas/pos/index.js | 13 +++-- template.yaml | 14 +++++ 5 files changed, 185 insertions(+), 17 deletions(-) create mode 100644 lambdas/po-sync/index.js 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: