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 <noreply@anthropic.com>
This commit is contained in:
Adam Moussa 2026-04-03 14:29:32 -04:00
parent 3854750597
commit 102999165f
5 changed files with 185 additions and 17 deletions

View file

@ -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 `</body>`:
@ -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

View file

@ -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",

99
lambdas/po-sync/index.js Normal file
View file

@ -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);
}

View file

@ -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 = {

View file

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