ledgerflow-backend/lambdas/pos/index.js
Adam Moussa 59127d5ab8 Initial commit — LedgerFlow backend
Lambda-based serverless backend with Google SSO, purchase orders,
invoices, and X12 810 EDI generation for Amazon Payee Central.
Includes bill-to/ship-to address support from Coupa purchase-orders table.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-02 18:19:38 -04:00

285 lines
11 KiB
JavaScript

// lambdas/pos/index.js
// Purchase Orders API
// GET /pos — list all POs (paginated, filterable)
// 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
const {
ScanCommand, GetCommand, PutCommand, UpdateCommand, DeleteCommand,
} = require("@aws-sdk/lib-dynamodb");
const {
ok, created, noContent, badRequest, notFound, conflict, serverError,
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);
// GET /pos/dynamo-scan
if (method === "GET" && id === "dynamo-scan") return dynamoScan(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" }) };
});
// ─── List ─────────────────────────────────────────────────────────────────────
async function listPOs(qs, user) {
const { limit, cursor } = parsePagination(qs);
const db = getDocClient();
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 });
}
// ─── 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);
}
}