import { S3Client, GetObjectCommand } from "@aws-sdk/client-s3"; import { DynamoDBClient } from "@aws-sdk/client-dynamodb"; import { DynamoDBDocumentClient, GetCommand, PutCommand, UpdateCommand } from "@aws-sdk/lib-dynamodb"; import { SSMClient, GetParameterCommand } from "@aws-sdk/client-ssm"; import { parse } from "csv-parse/sync"; const s3 = new S3Client(); const ddb = DynamoDBDocumentClient.from(new DynamoDBClient()); const ssm = new SSMClient(); const TABLE_NAME = process.env.TABLE_NAME; const BOA_BASE_URL = process.env.BOA_BASE_URL; async function getSSMParam(name) { const { Parameter } = await ssm.send( new GetParameterCommand({ Name: name, WithDecryption: true }) ); return Parameter.Value; } 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")}`; } export const handler = async (event) => { 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 if (method !== "Check") 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, clientSecret, accountNumber, companyId] = await Promise.all([ getSSMParam(process.env.BOA_CHECK_MGMT_APP_ID_PARAM), getSSMParam(process.env.BOA_CHECK_MGMT_CLIENT_ID_PARAM), getSSMParam(process.env.BOA_CHECK_MGMT_SECRET_PARAM), getSSMParam(process.env.BOA_ACCOUNT_NUMBER_PARAM), getSSMParam(process.env.BOA_COMPANY_ID_PARAM), ]); 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, })); 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 data = await res.json(); if (!res.ok) { console.error(`BoA ${action} error ${res.status}:`, JSON.stringify(data)); throw new Error(`BoA ${action} failed: ${res.status}`); } console.log( `BoA ${action}: ${data.processedItems}/${data.totalItems} processed, ${data.unprocessedItems} failed` ); return data; }; if (newChecks.length) { await submitToBoA(newChecks, "add_Issue"); } if (cancelChecks.length) { await submitToBoA(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}`, }; };