diff --git a/src/processPaymentCsv.js b/src/processPaymentCsv.js index fc080ea..bc282f4 100644 --- a/src/processPaymentCsv.js +++ b/src/processPaymentCsv.js @@ -1,6 +1,6 @@ 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 { DynamoDBDocumentClient, GetCommand, PutCommand, UpdateCommand, ScanCommand } from "@aws-sdk/lib-dynamodb"; import { SSMClient, GetParameterCommand } from "@aws-sdk/client-ssm"; import { parse } from "csv-parse/sync"; @@ -44,7 +44,194 @@ function toISODate(mdyDate) { 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, 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 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, " ")); @@ -166,8 +353,9 @@ export const handler = async (event) => { ); count++; - // Only process checks for CashPro + // 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 @@ -275,12 +463,22 @@ export const handler = async (event) => { 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 submitToBoA(newChecks, "add_Issue"); + await submitInBatches(newChecks, "add_Issue"); } if (cancelChecks.length) { - await submitToBoA(cancelChecks, "cancel_Issue"); + await submitInBatches(cancelChecks, "cancel_Issue"); } }