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"; 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) { const text = await res.text(); throw new Error(`OAuth token exchange failed: ${res.status} - ${text}`); } const data = await res.json(); return data.access_token; } // Convert MM/DD/YYYY to YYYY-MM-DD function toISODate(mdyDate) { const parts = String(mdyDate).split("/"); if (parts.length !== 3) return null; const [mm, dd, yyyy] = parts; return `${yyyy}-${mm.padStart(2, "0")}-${dd.padStart(2, "0")}`; } 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; toSubmit.push({ checkNumber: item.check_number, amount: Number(item.amount_usd).toFixed(2), issueDate: toISODate(item.send_payment_on), }); } 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), response_headers: JSON.stringify(headers), response_body: text, 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 }; }; 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) { console.error(`Backfill batch ${batchNum} retry failed:`, JSON.stringify(retryResult.data)); const detail = retryResult.parseError ? "BoA returned a non-JSON response" : `BoA returned ${JSON.stringify(retryResult.data)}`; 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 parseAmount = (value) => { const num = parseFloat(String(value || "0").replace(/,/g, "").trim()); return isNaN(num) ? 0 : num; }; 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 = []; // 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 sendOn = (row["Send Payment On"] || "").trim(); // ACH payments clear automatically on their send date if (method === "ACH" && !cancelStatuses.includes(status.toLowerCase())) { const sendDate = toISODate(sendOn); const today = new Date().toISOString().slice(0, 10); if (sendDate && sendDate <= today) { status = "Cleared"; } } 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; } } await ddb.send( new UpdateCommand({ TableName: TABLE_NAME, Key: { pk }, UpdateExpression: ` SET #method = :method, payee = :payee, check_number = :check_number, invoice_numbers = :invoice_numbers, send_payment_on = :send_payment_on, amount_usd = :amount_usd, #status = :status, company_subsidiary = :company_subsidiary `, ExpressionAttributeNames: { "#method": "method", "#status": "status", }, ExpressionAttributeValues: { ":method": method, ":payee": (row["Payee"] || "").trim(), ":check_number": checkNumber, ":invoice_numbers": (row["Invoice Numbers"] || "").trim(), ":send_payment_on": sendOn, ":amount_usd": parseAmount(row["Amount in USD"]), ":status": status, ":company_subsidiary": (row["Company/Subsidiary"] || "").trim(), }, }) ); count++; // Only process checks for CashPro (BoA rejects check numbers > 10 digits) if (method !== "Check") continue; if (checkNumber.length > 10) continue; if (!existing) { // New check → issue newChecks.push({ checkNumber, amount: parseAmount(row["Amount in USD"]).toFixed(2), issueDate: toISODate(sendOn), }); } else if ( cancelStatuses.includes(status.toLowerCase()) && !cancelStatuses.includes((existing.status || "").toLowerCase()) ) { // Existing check now voided/cancelled → cancel cancelChecks.push({ checkNumber, amount: parseAmount(row["Amount in USD"]).toFixed(2), issueDate: toISODate(sendOn), }); } } // 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), response_headers: JSON.stringify(headers), response_body: text, ttl: Math.floor(Date.now() / 1000) + 90 * 24 * 60 * 60, }; if (!res.ok) { console.error(`BoA ${action} error ${res.status}:`, text); console.error(`BoA ${action} response headers:`, JSON.stringify(headers)); 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}` ); console.log(`BoA ${action} response headers:`, JSON.stringify(headers)); console.log(`BoA ${action} response body:`, text); 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, }, }) ); console.log( `Upserted ${count} payments, ${newChecks.length} issued, ${cancelChecks.length} cancelled` ); return { statusCode: 200, body: `Upserted ${count}, issued ${newChecks.length}, cancelled ${cancelChecks.length}`, }; };