// lambdas/pos/index.js // Purchase Orders API // 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 pasted JSON) const { ScanCommand, GetCommand, PutCommand, UpdateCommand, DeleteCommand, } = require("@aws-sdk/lib-dynamodb"); const { ok, created, noContent, badRequest, notFound, conflict, getDocClient, TABLES, genId, parseBody, require_fields, handler, parsePagination, paginatedResponse, decodeCursor, } = require("@ledgerflow/shared"); const TABLE = TABLES.POS; // ─── Router ────────────────────────────────────────────────────────────────── exports.handler = handler(async (event, _ctx, user) => { const method = event.requestContext?.http?.method || event.httpMethod; const rawPath = event.rawPath || event.path || ""; const segments = rawPath.replace(/^\/pos\/?/, "").split("/").filter(Boolean); const id = segments[0] || null; const sub = segments[1] || null; const qs = event.queryStringParameters || {}; // POST /pos/import if (method === "POST" && id === "import") return bulkImport(event, user); if (!id) { if (method === "GET") return listPOs(qs, user); if (method === "POST") return createPO(event, user); } else { if (method === "GET") return getPO(id); if (method === "PUT") return updatePO(id, event, user, false); if (method === "PATCH") return updatePO(id, event, user, true); if (method === "DELETE") return deletePO(id, 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(); // 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 = { TableName: TABLE, Limit: limit, ExclusiveStartKey: decodeCursor(cursor), }; // Filter by status if provided if (qs.status) { params.FilterExpression = "#s = :s"; params.ExpressionAttributeNames = { "#s": "status" }; params.ExpressionAttributeValues = { ":s": qs.status.toLowerCase() }; } const result = await db.send(new ScanCommand(params)); return ok(paginatedResponse(result.Items || [], result.LastEvaluatedKey, result.Count)); } // ─── Get One ────────────────────────────────────────────────────────────────── async function getPO(id) { const db = getDocClient(); const result = await db.send(new GetCommand({ TableName: TABLE, Key: { id } })); if (!result.Item) return notFound("Purchase Order"); return ok(result.Item); } // ─── Create ─────────────────────────────────────────────────────────────────── async function createPO(event, user) { const body = parseBody(event); require_fields(body, ["poNumber", "vendor", "amount"]); const db = getDocClient(); const id = genId("PO"); const now = new Date().toISOString(); const item = { id, poNumber: body.poNumber.trim(), vendor: body.vendor.trim(), amount: parseFloat(body.amount), billed: 0, status: (body.status || "open").toLowerCase(), issueDate: body.issueDate || now.split("T")[0], dueDate: body.dueDate || null, notes: body.notes || "", source: "manual", createdBy: user?.email || "system", createdAt: now, updatedAt: now, }; // Prevent duplicate PO numbers const existing = await db.send(new ScanCommand({ TableName: TABLE, FilterExpression: "poNumber = :p", ExpressionAttributeValues: { ":p": item.poNumber }, Limit: 1, })); if (existing.Count > 0) return conflict(`PO number ${item.poNumber} already exists`); await db.send(new PutCommand({ TableName: TABLE, Item: item })); return created(item); } // ─── Update ─────────────────────────────────────────────────────────────────── async function updatePO(id, event, user, partial) { const body = parseBody(event); const db = getDocClient(); const existing = await db.send(new GetCommand({ TableName: TABLE, Key: { id } })); if (!existing.Item) return notFound("Purchase Order"); const UPDATABLE = ["vendor", "amount", "status", "issueDate", "dueDate", "notes", "poNumber"]; const updates = partial ? body : Object.fromEntries(UPDATABLE.map(k => [k, body[k] ?? existing.Item[k]])); const expressions = []; const names = {}; const values = { ":updatedAt": new Date().toISOString(), ":updatedBy": user?.email || "system" }; UPDATABLE.forEach(key => { if (updates[key] !== undefined) { expressions.push(`#${key} = :${key}`); names[`#${key}`] = key; values[`:${key}`] = key === "amount" ? parseFloat(updates[key]) : updates[key]; } }); expressions.push("#updatedAt = :updatedAt", "#updatedBy = :updatedBy"); names["#updatedAt"] = "updatedAt"; names["#updatedBy"] = "updatedBy"; const result = await db.send(new UpdateCommand({ TableName: TABLE, Key: { id }, UpdateExpression: "SET " + expressions.join(", "), ExpressionAttributeNames: names, ExpressionAttributeValues: values, ReturnValues: "ALL_NEW", })); return ok(result.Attributes); } // ─── Delete ─────────────────────────────────────────────────────────────────── async function deletePO(id, user) { const db = getDocClient(); const existing = await db.send(new GetCommand({ TableName: TABLE, Key: { id } })); if (!existing.Item) return notFound("Purchase Order"); await db.send(new DeleteCommand({ TableName: TABLE, Key: { id } })); return noContent(); } // ─── Bulk Import ────────────────────────────────────────────────────────────── async function bulkImport(event, user) { const body = parseBody(event); if (!Array.isArray(body.items) || body.items.length === 0) { return badRequest("items must be a non-empty array"); } if (body.items.length > 500) return badRequest("Maximum 500 items per import"); const db = getDocClient(); const now = new Date().toISOString(); const fieldMap = body.fieldMap || {}; // Custom field mapping from client config const mapField = (item, key, fallback) => item[fieldMap[key] || key] ?? item[fallback] ?? null; const results = { imported: 0, skipped: 0, errors: [] }; // Process in batches of 25 (DynamoDB TransactWrite limit) const batchSize = 25; for (let i = 0; i < body.items.length; i += batchSize) { const batch = body.items.slice(i, i + batchSize); await Promise.all(batch.map(async (raw, idx) => { try { const poNumber = mapField(raw, "poNumber", "po_number")?.toString()?.trim(); const vendor = mapField(raw, "vendor", "vendor_name")?.toString()?.trim(); const amount = parseFloat(mapField(raw, "amount", "total_amount") || 0); if (!poNumber || !vendor) { results.errors.push({ index: i + idx, error: "Missing poNumber or vendor" }); results.skipped++; return; } // Skip if already exists const existing = await db.send(new ScanCommand({ TableName: TABLE, FilterExpression: "poNumber = :p", ExpressionAttributeValues: { ":p": poNumber }, Limit: 1, })); if (existing.Count > 0) { results.skipped++; return; } const item = { id: genId("PO"), poNumber, vendor, amount, billed: 0, status: (mapField(raw, "status", "status") || "open").toLowerCase(), issueDate: mapField(raw, "issueDate", "issue_date") || now.split("T")[0], dueDate: mapField(raw, "dueDate", "due_date") || null, notes: raw.notes || raw.description || "", source: body.source || "import", createdBy: user?.email || "system", createdAt: now, updatedAt: now, }; await db.send(new PutCommand({ TableName: TABLE, Item: item })); results.imported++; } catch (err) { results.errors.push({ index: i + idx, error: err.message }); results.skipped++; } })); } return ok({ message: `Import complete`, ...results }); }