import { S3Client, GetObjectCommand } from "@aws-sdk/client-s3"; import { DynamoDBClient } from "@aws-sdk/client-dynamodb"; import { DynamoDBDocumentClient, GetCommand, PutCommand, UpdateCommand, ScanCommand, } from "@aws-sdk/lib-dynamodb"; import { SecretsManagerClient, GetSecretValueCommand } from "@aws-sdk/client-secrets-manager"; import { parse } from "csv-parse/sync"; import { toISODate, toCanonicalMDY, isPlausibleSendYear } from "./dates.js"; const s3 = new S3Client(); const ddb = DynamoDBDocumentClient.from(new DynamoDBClient()); const secrets = new SecretsManagerClient(); const TABLE_NAME = process.env.TABLE_NAME; const BOA_BASE_URL = process.env.BOA_BASE_URL; let cachedCreds; async function getCheckMgmtCreds() { if (cachedCreds) return cachedCreds; const { SecretString } = await secrets.send( new GetSecretValueCommand({ SecretId: process.env.BOA_CHECK_MGMT_SECRET_NAME }), ); cachedCreds = JSON.parse(SecretString); return cachedCreds; } async function getAccessToken(applicationID, clientId, clientSecret) { const res = await fetch(`${BOA_BASE_URL}/authn/v1/client-authentication`, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ applicationID, authn: { client_id: clientId, client_secret: clientSecret }, }), }); if (!res.ok) { throw new Error(`OAuth token exchange failed: HTTP ${res.status}`); } const data = await res.json(); return data.access_token; } // Untrusted values (CSV cells, record fields) must be encoded before log // interpolation — quoted CSV cells can carry newlines, which would forge // CloudWatch log lines (CWE-117). const logSafe = (v) => JSON.stringify(String(v ?? "").slice(0, 64)); // Comma-tolerant amount parsing, shared by the CSV path and the backfill so // a hand-inserted "1,234.56" string record can't NaN out of registration. const parseAmount = (value) => { const num = parseFloat( String(value || "0") .replace(/,/g, "") .trim(), ); return isNaN(num) ? 0 : num; }; async function backfillBoA() { const cancelStatuses = ["voided", "cancelled", "canceled", "marked as void"]; // Find all successfully submitted check numbers const submitted = new Set(); let txnKey; do { const txnScan = await ddb.send( new ScanCommand({ TableName: TABLE_NAME, FilterExpression: "begins_with(pk, :prefix) AND success = :t AND #action = :add", ExpressionAttributeNames: { "#action": "action" }, ExpressionAttributeValues: { ":prefix": "boa_txn#", ":t": true, ":add": "add_Issue" }, ProjectionExpression: "check_numbers", ...(txnKey && { ExclusiveStartKey: txnKey }), }), ); for (const item of txnScan.Items || []) { for (const cn of item.check_numbers || []) submitted.add(cn); } txnKey = txnScan.LastEvaluatedKey; } while (txnKey); // Find all Check-method records not yet submitted const toSubmit = []; let payKey; do { const payScan = await ddb.send( new ScanCommand({ TableName: TABLE_NAME, FilterExpression: "begins_with(pk, :prefix) AND #method = :check", ExpressionAttributeNames: { "#method": "method", "#status": "status" }, ExpressionAttributeValues: { ":prefix": "payment#", ":check": "Check" }, ProjectionExpression: "check_number, amount_usd, send_payment_on, #status", ...(payKey && { ExclusiveStartKey: payKey }), }), ); for (const item of payScan.Items || []) { if (submitted.has(item.check_number)) continue; if (cancelStatuses.includes((item.status || "").toLowerCase())) continue; if (item.check_number.length > 10) continue; const issueDate = toISODate(item.send_payment_on); const amount = parseAmount(item.amount_usd); if (!issueDate || !(amount > 0)) { console.error( `Backfill skipping check ${logSafe(item.check_number)}: missing issue date or amount`, ); continue; } toSubmit.push({ checkNumber: item.check_number, amount: amount.toFixed(2), issueDate, }); } payKey = payScan.LastEvaluatedKey; } while (payKey); if (!toSubmit.length) { console.log("Backfill: no unsubmitted checks found"); return { statusCode: 200, body: "No checks to backfill" }; } console.log(`Backfill: ${toSubmit.length} checks to submit`); const { appId, clientId, token: clientSecret, accountNumber, companyId, } = await getCheckMgmtCreds(); const bearerToken = await getAccessToken(appId, clientId, clientSecret); const BOA_BATCH_SIZE = 100; const ALREADY_SUBMITTED = new Set(["10096", "10097", "10098"]); let totalProcessed = 0; let totalSkippedDuplicates = 0; const submitBatch = async (items, label) => { const issueList = items.map((item) => ({ accountNumber, issueAction: "add_Issue", checkNumber: item.checkNumber, amount: item.amount, issueDate: item.issueDate, payee: "", })); const timestamp = new Date().toISOString(); const res = await fetch(`${BOA_BASE_URL}/cashpro/checkmanagement/v1/check-issues`, { method: "POST", headers: { "Content-Type": "application/json", Authorization: `Bearer ${bearerToken}`, companyId, }, body: JSON.stringify({ issueList }), }); const text = await res.text(); const headers = Object.fromEntries(res.headers.entries()); const transactionId = headers["transactionid"] || headers["x-transactionid"] || headers["x-correlation-id"] || headers["x-cashpro-transaction-id"] || null; let data = {}; let parseError = null; if (text.trim()) { try { data = JSON.parse(text); } catch (err) { parseError = err; } } else { parseError = new Error("BoA returned an empty response body"); } await ddb.send( new PutCommand({ TableName: TABLE_NAME, Item: { pk: `boa_txn#${timestamp}#add_Issue`, timestamp, action: "add_Issue", backfill: true, http_status: res.status, transaction_id: transactionId, check_numbers: items.map((i) => i.checkNumber), total_amount: items.reduce((sum, i) => sum + parseFloat(i.amount), 0).toFixed(2), success: res.ok, processed_items: data.processedItems || 0, total_items: data.totalItems || 0, ttl: Math.floor(Date.now() / 1000) + 90 * 24 * 60 * 60, }, }), ); if (res.ok && parseError) { throw new Error(`Backfill ${label} failed: BoA returned invalid JSON (HTTP ${res.status})`); } return { ok: res.ok, status: res.status, data, text, parseError, transactionId }; }; for (let i = 0; i < toSubmit.length; i += BOA_BATCH_SIZE) { const batch = toSubmit.slice(i, i + BOA_BATCH_SIZE); const batchNum = Math.floor(i / BOA_BATCH_SIZE) + 1; console.log(`Backfill batch ${batchNum}: ${batch.length} items`); const result = await submitBatch(batch, `batch ${batchNum}`); if (result.ok) { totalProcessed += result.data.processedItems; console.log(`Backfill batch ${batchNum}: ${result.data.processedItems} processed`); continue; } // Parse the error to separate duplicates from genuinely new items const dupeCheckNumbers = new Set(); for (const item of result.data?.issueList || []) { if (ALREADY_SUBMITTED.has(String(item.status))) { dupeCheckNumbers.add(item.checkNumber); } } if (dupeCheckNumbers.size === 0) { const detail = result.parseError ? "BoA returned a non-JSON response" : "BoA returned a non-duplicate error"; throw new Error(`Backfill batch ${batchNum} failed: ${detail} (HTTP ${result.status})`); } totalSkippedDuplicates += dupeCheckNumbers.size; const retryItems = batch.filter((item) => !dupeCheckNumbers.has(item.checkNumber)); console.log( `Backfill batch ${batchNum}: ${dupeCheckNumbers.size} already in BoA, ${retryItems.length} to retry`, ); if (retryItems.length === 0) continue; const retryResult = await submitBatch(retryItems, `batch ${batchNum} retry`); if (!retryResult.ok) { // Do not log/throw the raw BoA response body — it can carry account // numbers/PII. Status + txnId are enough to triage; the full body is // retrievable from BoA support via txnId. const detail = retryResult.parseError ? "BoA returned a non-JSON response" : "BoA returned a non-duplicate error"; console.error( `Backfill batch ${batchNum} retry failed: HTTP ${retryResult.status}, txnId=${retryResult.transactionId}`, ); throw new Error( `Backfill batch ${batchNum} retry failed: ${detail} (HTTP ${retryResult.status})`, ); } totalProcessed += retryResult.data.processedItems; console.log(`Backfill batch ${batchNum} retry: ${retryResult.data.processedItems} processed`); } console.log( `Backfill complete: ${totalProcessed} submitted, ${totalSkippedDuplicates} already in BoA`, ); return { statusCode: 200, body: `Backfilled ${totalProcessed} checks, ${totalSkippedDuplicates} already in BoA`, }; } export const handler = async (event) => { if (event.backfill) return backfillBoA(); const record = event.Records[0]; const bucket = record.s3.bucket.name; const key = decodeURIComponent(record.s3.object.key.replace(/\+/g, " ")); // Download CSV from S3 const { Body } = await s3.send(new GetObjectCommand({ Bucket: bucket, Key: key })); const csvText = await Body.transformToString("utf-8"); // Parse CSV const rows = parse(csvText, { columns: true, skip_empty_lines: true, trim: true, }); const normalizedRows = rows.map((row) => { const clean = {}; for (const [k, v] of Object.entries(row)) { clean[String(k).trim()] = typeof v === "string" ? v.trim() : v; } return clean; }); const cancelStatuses = ["voided", "cancelled", "canceled", "marked as void"]; // Status progression ranks — higher number = further along in lifecycle // Once a payment reaches a higher rank, CSV cannot move it backward const statusRank = { scheduled: 1, "payment submitted": 2, issued: 3, outstanding: 4, cleared: 5, }; const newChecks = []; const cancelChecks = []; const rejectedRows = []; // Upsert each payment, tracking new and cancelled checks let count = 0; for (const row of normalizedRows) { const checkNumber = (row["Check Number"] || "").trim(); if (!checkNumber) continue; const method = (row["Method"] || "").trim(); let status = (row["Status"] || "").trim(); const rawSendOn = (row["Send Payment On"] || "").trim(); // Reject rows with unparseable/implausible dates rather than storing // garbage (Stampli switched MM/DD/YYYY -> M/D/YY once already; a future // format change must fail loudly, not corrupt issueDate/aging). Blank is // allowed — canceled payments legitimately have no send date. const sendOn = rawSendOn ? toCanonicalMDY(rawSendOn) : ""; if (rawSendOn && (sendOn === null || !isPlausibleSendYear(rawSendOn))) { rejectedRows.push({ checkNumber, field: "Send Payment On", value: rawSendOn }); continue; } // ACH rows keep their Stampli status — the bank confirms settlement via // fetchBoaTransactions (#69). The old send-date auto-clear marked ACH // Cleared days before the debit settled (and past bounced settlements // that re-debited later were invisible to it). const pk = `payment#${checkNumber}`; // Check if record already exists (for detecting new vs updated) const { Item: existing } = await ddb.send( new GetCommand({ TableName: TABLE_NAME, Key: { pk } }), ); // Don't overwrite status once the bank has confirmed it as Cleared const bankConfirmed = existing?.clear_status === "Cleared"; if (bankConfirmed) { status = "Cleared"; } // Status progression protection — never allow status to move backward if (existing) { const oldRank = statusRank[(existing.status || "").toLowerCase()] || 0; const newRank = statusRank[status.toLowerCase()] || 0; if (oldRank === 5) { // Cleared is permanent — cannot be voided, cancelled, or anything else status = existing.status; } else if (newRank < oldRank && !cancelStatuses.includes(status.toLowerCase())) { // Non-cancel status regression — keep the existing (higher) status status = existing.status; } } // Stampli blanks "Amount in USD" (and can blank dates) on canceled rows. // Never let a blank overwrite a previously stored real value; fall back // to the "Amount" column, which keeps its value on cancels. const usdAmount = parseAmount(row["Amount in USD"]); const amount = usdAmount > 0 ? usdAmount : parseAmount(row["Amount"]); // Decide the BoA action BEFORE writing DDB. If a row demands a BoA call // we cannot make (no resolvable amount/issue date), reject it with the // record untouched — writing first would flip the status and permanently // close the transition gate, silently losing the cancel. First-seen rows // that are already canceled are stored but never registered: add_Issue // for a voided check would create an active issue with no cancel to follow. const isCancelRow = cancelStatuses.includes(status.toLowerCase()); const boaEligible = method === "Check" && checkNumber.length <= 10; const boaAmount = amount > 0 ? amount : parseAmount(existing?.amount_usd); const boaIssueDate = toISODate(sendOn) || toISODate(existing?.send_payment_on); const needsAdd = boaEligible && !existing && !isCancelRow; const needsCancel = boaEligible && existing && isCancelRow && !cancelStatuses.includes((existing.status || "").toLowerCase()); if ((needsAdd || needsCancel) && (!boaIssueDate || !(boaAmount > 0))) { rejectedRows.push({ checkNumber, field: needsAdd ? "add_Issue data" : "cancel_Issue data", value: `issueDate=${boaIssueDate}, amount=${boaAmount}`, }); continue; } const sets = [ "#method = :method", "payee = :payee", "check_number = :check_number", "invoice_numbers = :invoice_numbers", "#status = :status", "company_subsidiary = :company_subsidiary", ]; const values = { ":method": method, ":payee": (row["Payee"] || "").trim(), ":check_number": checkNumber, ":invoice_numbers": (row["Invoice Numbers"] || "").trim(), ":status": status, ":company_subsidiary": (row["Company/Subsidiary"] || "").trim(), }; if (amount > 0 || !existing) { sets.push("amount_usd = :amount_usd"); values[":amount_usd"] = amount; } if (sendOn || !existing) { sets.push("send_payment_on = :send_payment_on"); values[":send_payment_on"] = sendOn; } await ddb.send( new UpdateCommand({ TableName: TABLE_NAME, Key: { pk }, UpdateExpression: `SET ${sets.join(", ")}`, ExpressionAttributeNames: { "#method": "method", "#status": "status", }, ExpressionAttributeValues: values, }), ); count++; if (needsAdd) { newChecks.push({ checkNumber, amount: boaAmount.toFixed(2), issueDate: boaIssueDate, }); } else if (needsCancel) { cancelChecks.push({ checkNumber, amount: boaAmount.toFixed(2), issueDate: boaIssueDate, }); } } // Submit to CashPro if there are any new issues or cancels if (newChecks.length || cancelChecks.length) { const { appId, clientId, token: clientSecret, accountNumber, companyId, } = await getCheckMgmtCreds(); const bearerToken = await getAccessToken(appId, clientId, clientSecret); const submitToBoA = async (items, action) => { const issueList = items.map((item) => ({ accountNumber, issueAction: action, checkNumber: item.checkNumber, amount: item.amount, issueDate: item.issueDate, payee: "", })); const timestamp = new Date().toISOString(); const checkNumbers = items.map((i) => i.checkNumber); const totalAmount = items.reduce((sum, i) => sum + parseFloat(i.amount), 0); const res = await fetch(`${BOA_BASE_URL}/cashpro/checkmanagement/v1/check-issues`, { method: "POST", headers: { "Content-Type": "application/json", Authorization: `Bearer ${bearerToken}`, companyId, }, body: JSON.stringify({ issueList }), }); const text = await res.text(); const headers = Object.fromEntries(res.headers.entries()); const transactionId = headers["transactionid"] || headers["x-transactionid"] || headers["x-correlation-id"] || headers["x-cashpro-transaction-id"] || null; const record = { pk: `boa_txn#${timestamp}#${action}`, timestamp, action, http_status: res.status, transaction_id: transactionId, check_numbers: checkNumbers, total_amount: totalAmount.toFixed(2), ttl: Math.floor(Date.now() / 1000) + 90 * 24 * 60 * 60, }; if (!res.ok) { // Do not log the raw response body/headers — the BoA CashPro response // can carry account numbers/PII. Status + correlation id are enough to // triage; the full body is retrievable from BoA support via txnId. console.error(`BoA ${action} failed: HTTP ${res.status}, txnId=${transactionId}`); record.success = false; await ddb.send(new PutCommand({ TableName: TABLE_NAME, Item: record })); throw new Error(`BoA ${action} failed: ${res.status}`); } const data = JSON.parse(text); record.success = true; record.processed_items = data.processedItems; record.total_items = data.totalItems; record.unprocessed_items = data.unprocessedItems; record.message = data.issueList?.[0]?.message || null; await ddb.send(new PutCommand({ TableName: TABLE_NAME, Item: record })); console.log( `BoA ${action}: ${data.processedItems}/${data.totalItems} processed, ${data.unprocessedItems} failed, txnId=${transactionId}`, ); return data; }; const BOA_BATCH_SIZE = 100; const submitInBatches = async (items, action) => { for (let i = 0; i < items.length; i += BOA_BATCH_SIZE) { const batch = items.slice(i, i + BOA_BATCH_SIZE); console.log( `BoA ${action} batch ${Math.floor(i / BOA_BATCH_SIZE) + 1}: ${batch.length} items`, ); await submitToBoA(batch, action); } }; if (newChecks.length) { await submitInBatches(newChecks, "add_Issue"); } if (cancelChecks.length) { await submitInBatches(cancelChecks, "cancel_Issue"); } } // Update metadata record await ddb.send( new PutCommand({ TableName: TABLE_NAME, Item: { pk: "metadata", file_name: key.split("/").pop(), last_updated: new Date().toISOString(), last_file_count: count, last_rejected_count: rejectedRows.length, }, }), ); console.log( `Upserted ${count} payments, ${newChecks.length} issued, ${cancelChecks.length} cancelled, ${rejectedRows.length} rejected`, ); // All valid rows are processed and BoA submissions are done; now fail the // invocation so rejects surface via the errors alarm + async DLQ. Retries // are safe: upserts are idempotent and already-existing checks are not // re-offered to BoA. if (rejectedRows.length) { for (const r of rejectedRows) { console.error( `Rejected row: check ${logSafe(r.checkNumber)}, ${r.field}=${logSafe(r.value)}`, ); } throw new Error( `${rejectedRows.length} of ${normalizedRows.length} rows rejected (bad dates or missing BoA data; ${count} valid rows processed)`, ); } return { statusCode: 200, body: `Upserted ${count}, issued ${newChecks.length}, cancelled ${cancelChecks.length}`, }; };