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