ledgerflow-backend/lambdas/pos/index.js
Adam Moussa 3854750597 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 <noreply@anthropic.com>
2026-04-03 13:58:15 -04:00

319 lines
12 KiB
JavaScript

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