import { S3Client, GetObjectCommand } from "@aws-sdk/client-s3"; import { DynamoDBClient } from "@aws-sdk/client-dynamodb"; import { DynamoDBDocumentClient, BatchWriteCommand, PutCommand } 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; }; const payments = normalizedRows.map((row) => ({ pk: `payment#${(row["Check Number"] || "").trim()}`, method: (row["Method"] || "").trim(), payee: (row["Payee"] || "").trim(), check_number: (row["Check Number"] || "").trim(), 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(), })); // Write payments in batches of 25 (DynamoDB BatchWrite limit) const batches = []; for (let i = 0; i < payments.length; i += 25) { const batch = payments.slice(i, i + 25).map((item) => ({ PutRequest: { Item: item }, })); batches.push(batch); } for (const batch of batches) { await ddb.send( new BatchWriteCommand({ RequestItems: { [TABLE_NAME]: batch }, }) ); } // 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: payments.length, }, }) ); console.log(`Upserted ${payments.length} payments from ${key}`); return { statusCode: 200, body: `Upserted ${payments.length} payments` }; };