import { S3Client, GetObjectCommand } from "@aws-sdk/client-s3"; import { DynamoDBClient } from "@aws-sdk/client-dynamodb"; import { DynamoDBDocumentClient, GetCommand, PutCommand, ScanCommand, BatchWriteCommand, } from "@aws-sdk/lib-dynamodb"; import { SQSClient, SendMessageCommand } from "@aws-sdk/client-sqs"; import { SSMClient, GetParameterCommand } from "@aws-sdk/client-ssm"; import { simpleParser } from "mailparser"; const s3 = new S3Client(); const ddb = DynamoDBDocumentClient.from(new DynamoDBClient()); const sqs = new SQSClient(); const ssm = new SSMClient(); const TABLE_NAME = process.env.TABLE_NAME; const CHANNEL_ID = process.env.PAYROLL_CHANNEL_ID; const QUEUE_URL = process.env.CONTRACTOR_BATCH_QUEUE_URL; let cachedToken; async function getSlackToken() { if (cachedToken) return cachedToken; const { Parameter } = await ssm.send( new GetParameterCommand({ Name: process.env.SLACK_BOT_TOKEN_PARAM, WithDecryption: true, }) ); cachedToken = Parameter.Value; return cachedToken; } const formatCurrency = (v) => new Intl.NumberFormat("en-US", { style: "currency", currency: "USD", }).format(Number(v || 0)); function classifyEmail(from, subject, text) { const fromAddr = (from?.text || from || "").toLowerCase(); const subj = (subject || "").toLowerCase(); if ( fromAddr.includes("automated@gusto.com") && subj.includes("payroll confirmation") ) { return "employee"; } if ( fromAddr.includes("gustonoreply@gusto.com") && subj.includes("payment confirmation") ) { return "contractor"; } const strippedSubj = subj.replace(/^fwd?:\s*/i, ""); const body = (text || "").toLowerCase(); const bodyMentionsGusto = body.includes("gusto.com"); if (bodyMentionsGusto && strippedSubj.includes("payroll confirmation")) { return "employee"; } if (bodyMentionsGusto && strippedSubj.includes("payment confirmation")) { return "contractor"; } return null; } function parseDateFromSubject(subject) { const m = subject.match(/:\s*(\w+,\s*\w+\s+\d+)\s+(?:payroll|payment)\s+confirmation/i); if (!m) return null; const raw = m[1].trim(); const year = new Date().getFullYear(); const d = new Date(`${raw}, ${year}`); if (isNaN(d)) return null; return { display: raw, iso: d.toISOString().slice(0, 10), }; } function parseEmployeePayroll(subject, text) { const checkDate = parseDateFromSubject(subject); if (!checkDate) return null; const norm = text.replace(/\n/g, " ").replace(/\s+/g, " "); const debitMatch = norm.match( /we.ll debit \$([\d,]+\.\d{2}) from the bank account ending in (\d+)/i ); if (!debitMatch) return null; const totalDebit = parseFloat(debitMatch[1].replace(/,/g, "")); const bankSuffix = debitMatch[2]; const payPeriodMatch = norm.match(/payroll for the (.+?)\s+pay period/i); const netPayMatch = norm.match( /\$([\d,]+\.\d{2}) will be for your employee net pay/i ); const taxesMatch = norm.match(/\$([\d,]+\.\d{2}) will be for taxes/i); const reimbursementsMatch = norm.match( /\$([\d,]+\.\d{2}) will be for reimbursements/i ); return { checkDate: checkDate.display, dateKey: checkDate.iso, totalDebit, bankSuffix, payPeriod: payPeriodMatch ? payPeriodMatch[1].trim() : null, netPay: netPayMatch ? parseFloat(netPayMatch[1].replace(/,/g, "")) : null, taxes: taxesMatch ? parseFloat(taxesMatch[1].replace(/,/g, "")) : null, reimbursements: reimbursementsMatch ? parseFloat(reimbursementsMatch[1].replace(/,/g, "")) : null, }; } function parseContractorPayment(subject, text) { const checkDate = parseDateFromSubject(subject); if (!checkDate) return null; const norm = text.replace(/\n/g, " ").replace(/\s+/g, " "); const debitMatch = norm.match( /we.ll debit \$([\d,]+\.\d{2}) from the bank account ending in (\d+)/i ); if (!debitMatch) return null; const totalDebit = parseFloat(debitMatch[1].replace(/,/g, "")); const bankSuffix = debitMatch[2]; const nameMatch = norm.match(/at that time for (.+?) to be paid on/i); const contractorName = nameMatch ? nameMatch[1].trim() : "Unknown contractor"; return { checkDate: checkDate.display, dateKey: checkDate.iso, totalDebit, bankSuffix, contractorName, }; } function buildCombinedBlocks(dateDisplay, employee, contractors) { const blocks = [ { type: "header", text: { type: "plain_text", text: `Payroll Processed — ${dateDisplay}`, }, }, ]; if (employee) { const topFields = [ { type: "mrkdwn", text: `*Check Date*\n${employee.checkDate}` }, ]; if (employee.payPeriod) { topFields.unshift({ type: "mrkdwn", text: `*Pay Period*\n${employee.payPeriod}`, }); } blocks.push({ type: "section", fields: topFields }); blocks.push({ type: "divider" }); const detailFields = []; if (employee.netPay != null) { detailFields.push({ type: "mrkdwn", text: `*Employee Net Pay*\n${formatCurrency(employee.netPay)}`, }); } if (employee.taxes != null) { detailFields.push({ type: "mrkdwn", text: `*Taxes*\n${formatCurrency(employee.taxes)}`, }); } if (detailFields.length) blocks.push({ type: "section", fields: detailFields }); const extraFields = []; if (employee.reimbursements != null) { extraFields.push({ type: "mrkdwn", text: `*Reimbursements*\n${formatCurrency(employee.reimbursements)}`, }); } extraFields.push({ type: "mrkdwn", text: `*Bank Account*\n...${employee.bankSuffix}`, }); blocks.push({ type: "section", fields: extraFields }); } if (contractors.length > 0) { blocks.push({ type: "divider" }); blocks.push({ type: "section", text: { type: "mrkdwn", text: `*Contractor Payments*`, }, }); for (const c of contractors) { blocks.push({ type: "section", fields: [ { type: "mrkdwn", text: `*Contractor*\n${c.contractorName}` }, { type: "mrkdwn", text: `*Amount*\n${formatCurrency(c.totalDebit)}` }, ], }); } } const grandTotal = (employee ? employee.totalDebit : 0) + contractors.reduce((sum, c) => sum + c.totalDebit, 0); blocks.push({ type: "divider" }); blocks.push({ type: "section", text: { type: "mrkdwn", text: `:moneybag: *Total Bank Debit: ${formatCurrency(grandTotal)}*`, }, }); blocks.push({ type: "context", elements: [ { type: "mrkdwn", text: "Source: Gusto payroll confirmation emails" }, ], }); return { blocks, grandTotal }; } async function postSlack(token, blocks, text) { const res = await fetch("https://slack.com/api/chat.postMessage", { method: "POST", headers: { Authorization: `Bearer ${token}`, "Content-Type": "application/json", }, body: JSON.stringify({ channel: CHANNEL_ID, blocks, text }), }); const data = await res.json(); if (!data.ok) { console.error("Slack post failed:", data.error); throw new Error(`Slack API error: ${data.error}`); } return data.ts; } const ttl90Days = () => Math.floor(Date.now() / 1000) + 90 * 24 * 60 * 60; async function handleS3Event(record) { const bucket = record.s3.bucket.name; const key = decodeURIComponent(record.s3.object.key.replace(/\+/g, " ")); const { Body } = await s3.send( new GetObjectCommand({ Bucket: bucket, Key: key }) ); const rawEmail = await Body.transformToByteArray(); const parsed = await simpleParser(Buffer.from(rawEmail)); const type = classifyEmail(parsed.from, parsed.subject, parsed.text); if (!type) { console.log("Unrecognized email, skipping:", parsed.subject); return; } if (type === "employee") { const data = parseEmployeePayroll(parsed.subject, parsed.text); if (!data) { console.error("Failed to parse employee payroll:", parsed.subject); return; } const pk = `PAYROLL_EMAIL#employee#${data.dateKey}`; const existing = await ddb.send( new GetCommand({ TableName: TABLE_NAME, Key: { pk } }) ); if (existing.Item) { console.log(`Already recorded: employee ${data.dateKey}`); return; } await ddb.send( new PutCommand({ TableName: TABLE_NAME, Item: { pk, type: "employee", status: "pending", checkDate: data.checkDate, dateKey: data.dateKey, totalDebit: data.totalDebit, bankSuffix: data.bankSuffix, payPeriod: data.payPeriod, netPay: data.netPay, taxes: data.taxes, reimbursements: data.reimbursements, processedAt: new Date().toISOString(), ttl: ttl90Days(), }, }) ); await sqs.send( new SendMessageCommand({ QueueUrl: QUEUE_URL, MessageBody: JSON.stringify({ dateKey: data.dateKey }), }) ); console.log(`Queued employee payroll for ${data.dateKey}`); } if (type === "contractor") { const data = parseContractorPayment(parsed.subject, parsed.text); if (!data) { console.error("Failed to parse contractor payment:", parsed.subject); return; } const pk = `PAYROLL_EMAIL#contractor#${data.dateKey}#${data.bankSuffix}`; const existing = await ddb.send( new GetCommand({ TableName: TABLE_NAME, Key: { pk } }) ); if (existing.Item) { console.log( `Already recorded: contractor ${data.contractorName} ${data.dateKey}` ); return; } await ddb.send( new PutCommand({ TableName: TABLE_NAME, Item: { pk, type: "contractor", status: "pending", checkDate: data.checkDate, dateKey: data.dateKey, totalDebit: data.totalDebit, bankSuffix: data.bankSuffix, contractorName: data.contractorName, processedAt: new Date().toISOString(), ttl: ttl90Days(), }, }) ); await sqs.send( new SendMessageCommand({ QueueUrl: QUEUE_URL, MessageBody: JSON.stringify({ dateKey: data.dateKey }), }) ); console.log( `Queued contractor: ${data.contractorName} for ${data.dateKey}` ); } } async function handleSqsEvent(sqsRecord) { const { dateKey } = JSON.parse(sqsRecord.body); const batchPk = `PAYROLL_BATCH#${dateKey}`; const batchCheck = await ddb.send( new GetCommand({ TableName: TABLE_NAME, Key: { pk: batchPk } }) ); if (batchCheck.Item) { console.log(`Batch already posted for ${dateKey}`); return; } const { Items } = await ddb.send( new ScanCommand({ TableName: TABLE_NAME, FilterExpression: "begins_with(pk, :prefix) AND #s = :pending", ExpressionAttributeNames: { "#s": "status" }, ExpressionAttributeValues: { ":prefix": `PAYROLL_EMAIL#`, ":pending": "pending", }, }) ); const pending = (Items || []).filter((i) => i.dateKey === dateKey); if (pending.length === 0) { console.log(`No pending items for ${dateKey}`); return; } const employee = pending.find((i) => i.type === "employee") || null; const contractors = pending.filter((i) => i.type === "contractor"); const dateDisplay = pending[0].checkDate; const token = await getSlackToken(); const { blocks, grandTotal } = buildCombinedBlocks(dateDisplay, employee, contractors); const slackTs = await postSlack( token, blocks, `Payroll processed for ${dateDisplay}: ${formatCurrency(grandTotal)}` ); await ddb.send( new PutCommand({ TableName: TABLE_NAME, Item: { pk: batchPk, dateKey, checkDate: dateDisplay, employeeCount: employee ? 1 : 0, contractorCount: contractors.length, totalDebit: grandTotal, slackTs, processedAt: new Date().toISOString(), ttl: ttl90Days(), }, }) ); const updateBatch = pending.map((item) => ({ PutRequest: { Item: { ...item, status: "notified" }, }, })); for (let i = 0; i < updateBatch.length; i += 25) { await ddb.send( new BatchWriteCommand({ RequestItems: { [TABLE_NAME]: updateBatch.slice(i, i + 25) }, }) ); } console.log( `Posted batch: ${employee ? 1 : 0} employee + ${contractors.length} contractor(s) for ${dateKey}` ); } export const handler = async (event) => { if (event.Records?.[0]?.eventSource === "aws:s3") { for (const record of event.Records) { await handleS3Event(record); } } else if (event.Records?.[0]?.eventSource === "aws:sqs") { for (const record of event.Records) { await handleSqsEvent(record); } } else { console.log("Unknown event source:", JSON.stringify(event).slice(0, 200)); } return { statusCode: 200 }; };