import { S3Client, GetObjectCommand } from "@aws-sdk/client-s3"; import { DynamoDBClient } from "@aws-sdk/client-dynamodb"; import { DynamoDBDocumentClient, PutCommand, UpdateCommand } from "@aws-sdk/lib-dynamodb"; import { parse } from "csv-parse/sync"; const s3 = new S3Client(); const ddb = DynamoDBDocumentClient.from(new DynamoDBClient()); const TABLE_NAME = process.env.TABLE_NAME; 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; }; // Upsert each payment, preserving clear_status/bank_reference/cleared_date if they exist let count = 0; for (const row of normalizedRows) { const checkNumber = (row["Check Number"] || "").trim(); const pk = `payment#${checkNumber}`; 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": (row["Method"] || "").trim(), ":payee": (row["Payee"] || "").trim(), ":check_number": checkNumber, ":invoice_numbers": (row["Invoice Numbers"] || "").trim(), ":send_payment_on": (row["Send Payment On"] || "").trim(), ":amount_usd": parseAmount(row["Amount in USD"]), ":status": (row["Status"] || "").trim(), ":company_subsidiary": (row["Company/Subsidiary"] || "").trim(), }, }) ); count++; } // 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 from ${key}`); return { statusCode: 200, body: `Upserted ${count} payments` }; };