Coupa PO email ingestion pipeline
Find a file
Adam Moussa fa5f60208f
Harden AR parser per cross-family review
Cross-family (GPT-4.1) review findings: terminate the dkim result
token at end-of-clause, whitespace, or a comment so a value like
"dkim=pass-fake" can never be read as a pass; normalize trailing
dots off allowlist entries so "seahaven.com." matches; make the
compat32 parser policy explicit. Adds tests for result-token
boundaries, comments after the result, quoted domain values, and
folding inside a dkim clause.

Refs: INFRA-107
2026-07-15 18:55:32 -04:00
.github chore(ci): SHA-pin org reusable-workflow caller refs (INFRA-50) (#85) 2026-07-06 18:27:18 -04:00
cdk Add fail-closed SES sender authentication 2026-07-15 18:53:46 -04:00
lambdas Harden AR parser per cross-family review 2026-07-15 18:55:32 -04:00
scripts Merge workorder-ingest into unified procurement repo (#22) 2026-05-12 15:21:06 -04:00
tests Harden AR parser per cross-family review 2026-07-15 18:55:32 -04:00
.gitignore Merge workorder-ingest into unified procurement repo (#22) 2026-05-12 15:21:06 -04:00
README.md Add fail-closed SES sender authentication 2026-07-15 18:53:46 -04:00
test_local.py Merge workorder-ingest into unified procurement repo (#22) 2026-05-12 15:21:06 -04:00

Procurement Ingest

Python AWS CDK CI

Unified email ingestion pipelines for Amazon procurement data. Two independent pipelines — purchase orders (Coupa) and work orders (APM/Hexagon EAM) — share a single repo and CDK app but deploy as separate CloudFormation stacks.

Pipelines

Purchase Orders (po-ingest stack)

Coupa PO emails are received at amazon_po@int.seahaven.com, parsed by Claude Haiku 4.5, and written to the purchase-orders DynamoDB table.

Flow:

  1. Coupa sends a PO email (new, revision, or cancellation).
  2. SES (INBOUND_MAIL rule set) drops the raw MIME into s3://po-ingest-emails-{AccountId}/inbound/.
  3. S3 ObjectCreated triggers the po-email-processor Lambda.
  4. Fail-closed sender authentication (INFRA-107): the SES-stamped Authentication-Results header must show dkim=pass for amazon.coupahost.com (see Sender authentication); otherwise the email is logged and dropped.
  5. Claude extracts structured JSON (PO number, status, supplier, site code, trade classification, line items, fiscal year).
  6. Conditional write to DynamoDB:
    • new_po — idempotent insert (no-op if PO exists)
    • revision — unconditional overwrite
    • cancellation — marks existing row Cancelled
  7. DynamoDB Streams feeds downstream consumers:
    • LedgerFlow (seahaven-slack-bot/po-sync) — daily KB sync
    • Site extractor (po-ingest-site-extractor) — real-time site address extraction into verified-sites table

Lambdas (lambdas/po/):

Function Trigger Purpose
po-email-processor S3 ObjectCreated Claude extraction + DynamoDB write
po-ingest-site-extractor DynamoDB Streams Site code/address extraction -> verified-sites
po-web-ui Manual invoke HTML dashboard (public Function URL removed 2026-06-08, INFRA-74)

Tables:

  • purchase-orders (PK: po_number, Streams: NEW_IMAGE) — shared with seahaven-slack-bot (read-only; see Shared Resources)
  • verified-sites (PK: siteCode, GSI: by-state) — ~1,100 unique Amazon facility sites
  • pending-site-review (PK: po_number) — unresolvable POs for manual Payee Central verification

Work Orders (WorkorderIngestStack stack)

Amazon APM work order emails (from Hexagon EAM / HxGN SmartCloud) are received at apm@int.seahaven.com, parsed by Claude Haiku 4.5, and written to the WorkOrders DynamoDB table.

Flow:

  1. Hexagon EAM sends email notifications (new assignments, comments, updates, cancellations) to amazon@seahavenind.com.
  2. Gmail filter forwards APM emails to apm@int.seahaven.com (SES).
  3. SES drops the raw MIME into s3://workorder-ingest-emails-{AccountId}/inbound/.
  4. S3 triggers the workorder-email-processor Lambda.
  5. Fail-closed sender authentication (INFRA-107): the SES-stamped Authentication-Results header must show dkim=pass for seahaven.com (the Google Groups forward re-signs the mail; see Sender authentication); otherwise the email is logged and dropped.
  6. Claude extracts structured JSON (work order ID, site code, severity, priority, dates, assigned technician).
  7. Work order upserted to WorkOrders, event/comment appended to WorkOrderComments.

Lambdas (lambdas/wo/):

Function Trigger Purpose
workorder-email-processor S3 ObjectCreated Claude extraction + DynamoDB write
workorder-web-ui Manual invoke HTML dashboard (public Function URL removed 2026-06-08, INFRA-74)

Tables:

  • WorkOrders (PK: work_order_id, GSIs: site-code-index, status-index)
  • WorkOrderComments (PK: work_order_id, SK: comment_id)

Architecture

IaC: AWS CDK (Python), two stacks in one app, region us-east-1.

All Lambdas: Python 3.12, ARM64, 60-day log retention.

Secrets:

  • po-ingest/anthropic-api-key — Anthropic API key for PO parsing
  • workorder-ingest/anthropic-api-key — Anthropic API key for WO parsing

SES: Both stacks add rules to the shared INBOUND_MAIL receipt rule set on int.seahaven.com.

Sender authentication (INFRA-107)

The From header and any Authentication-Results header inside the raw MIME are attacker-forgeable, so neither is trusted. Instead, both email processors (lambdas/*/email_processor/ses_auth.py) authenticate the sender against the verdicts SES itself stamps at delivery time, failing closed:

  1. Take only the topmost Authentication-Results header (SES prepends its trace headers; any lower copies arrived inside the message and are ignored).
  2. Require its authserv-id to be amazonses.com.
  3. Require a dkim=pass clause whose header.d=/header.i= domain is in the pipeline's allowlist.

The allowlist is the ALLOWED_DKIM_DOMAINS Lambda environment variable (comma-separated, set per stack in CDK — no code change needed to adjust):

Pipeline ALLOWED_DKIM_DOMAINS Why
Work orders seahaven.com APM mail arrives via the apm@ Google Groups forward, which re-signs as seahaven.com (the original hxgnsmartcloud.com signature does not survive the forward)
Purchase orders amazon.coupahost.com Coupa signs as amazon.coupahost.com. amazonses.com also passes but is deliberately not allowlisted — every SES customer's mail passes for it

On any failure (env var unset, header missing/unparseable, verdict fail, unaligned domain) the processor logs a structured sender_auth_rejected warning with the reason and S3 key, skips the email, and returns normally — rejected mail never triggers Lambda retries or DLQ messages. Unit tests live in tests/test_ses_auth.py.

Failure handling (INFRA-41): Each email-processor is async-invoked (S3 → Lambda). Both have a CDK-managed SQS dead-letter queue (dead_letter_queue=, 14-day retention, SSL-enforced) so a failed parse is captured rather than silently dropped after Lambda's retries.

CloudWatch alarms

Every alarm is ALARM-only (no OK action), sends to the shared site-alerts SNS topic (imported once per stack via Topic.from_topic_arn), and uses TreatMissingData.NOT_BREACHING.

Lambda alarms (AWS/Lambda, FunctionName dimension):

Alarm Functions Metric / config
<fn>-errors po-email-processor, po-ingest-site-extractor, workorder-email-processor Errors Sum, 5 min, > 0, eval 1
<fn>-throttles po-email-processor, po-ingest-site-extractor, po-web-ui, workorder-email-processor Throttles Sum, 5 min, > 0, eval 1
<fn>-duration po-email-processor, po-ingest-site-extractor, po-web-ui (p99); workorder-email-processor (p95) Duration percentile, 5 min, >= 45000 ms (75% of the 60s timeout), eval 3 / datapoints 2

The <fn>-duration and <fn>-throttles alarms for po-email-processor and workorder-email-processor supersede the orphaned, CLI-created Lambda-Duration-* / Lambda-Throttles-* alarms (deleted post-deploy).

DynamoDB alarms (AWS/DynamoDB): each owned table gets <table>-throttles (ThrottledRequests) and <table>-system-errors (SystemErrors). These metrics emit only at the TableName + Operation dimension set, so each alarm is a Sum math expression across the operations the table uses (Get/BatchGet/Query/Scan/Put/Update/Delete/BatchWrite). Tables covered: purchase-orders, verified-sites, pending-site-review (po-ingest); WorkOrders, WorkOrderComments (workorder-ingest).

Shared Resources

purchase-orders table (owned here)

The purchase-orders DynamoDB table is owned by this repo's po-ingest stack (defined in cdk/po_stack.py with RemovalPolicy.RETAIN and StreamViewType.NEW_IMAGE). The po-email-processor Lambda is the authoritative writer — it performs the conditional inserts, revision overwrites, and cancellation updates described above.

Consumers (read-only):

Repo How it reads Purpose
seahaven-slack-bot po-sync (DynamoDB Streams + daily scan) and wo-po-lookup Daily KB sync + Bedrock agent PO lookups

The consumer imports the table via Table.fromTableName(...) and is granted read-only access (grantReadData); it does not own or define it.

Schema-coordination rule: Any change to the purchase-orders schema (partition key, item shape, attribute names, streams view type) must be coordinated with seahaven-slack-bot. The owner here ships the change; the consumer must be updated in lockstep so its readers do not break. Treat schema changes as a cross-repo migration, not a local edit.

Known exception (INFRA-51): amazon-po-parser currently writes directly to purchase-orders outside this stack (backfill/enrichment scripts). This second writer is being folded into the po-ingest pipeline so this stack is the sole writer; until INFRA-51 closes, coordinate any schema change with amazon-po-parser as well.

WorkOrders and WorkOrderComments tables (owned here)

Both tables are owned by this repo's WorkorderIngestStack (cdk/wo_stack.py, RemovalPolicy.RETAIN, shared customer-managed CMK per INFRA-95 / M-3):

  • WorkOrders — PK work_order_id (S).
  • WorkOrderComments — PK work_order_id (S), SK comment_id (S).

Consumer (read-only) — data contract: seahaven-slack-bot imports both tables via Table.fromTableName(...) (grantReadData plus an explicit kms:Decrypt grant on the shared CMK) and reads them from two Lambdas: workorder-sync (daily full-table scan into the Bedrock knowledge base) and wo-po-lookup (the Bedrock agent's direct WO lookup action group). The bot depends on the PK/SK schema above, the CMK encryption, and these attributes: on WorkOrders — description, wo_status, customer, site_code, building, address, severity, priority, assigned_to, date_reported, scheduled_start, due_date, updated_at; on WorkOrderComments — created_at (used to sort comments), commenter, text. Any change to table name, key schema, these attribute names, or the encryption key must be coordinated with seahaven-slack-bot before it ships, or the Bedrock agent breaks at runtime (not at deploy — the tables are imported by name, so there is no compile-time link).

verified-sites table (owned here)

Owned by this repo's po-ingest stack (cdk/po_stack.py). PK siteCode (S); AWS-managed encryption (NOT the shared CMK).

Consumer (read-only) — data contract: seahaven-slack-bot's wo-po-lookup Lambda imports this table via Table.fromTableName(...) for the Bedrock agent's lookup_site action. It does point lookups by siteCode and reads address, fullAddress, city, state, zip, latitude, longitude, notes. Coordinate any change to the table name, key schema, or these attribute names with seahaven-slack-bot.

GSI drift (INFRA-138): the by-state GSI was removed here on 2026-06-03 (audit M-20, "0 reads in 30d"), but seahaven-slack-bot still queries IndexName: 'by-state' for its state-listing path, so that path fails at runtime today. Restoring the GSI or removing the consumer's state path needs to be reconciled cross-repo. This is the kind of silent owner-side lifecycle change this data-contract note exists to prevent.

Documentation

The canonical map of Sea Haven's AWS infrastructure lives in Confluence. This project's po-ingest and workorder-ingest stacks are represented there as Mermaid subgraphs.

CI/CD

GitHub Actions with reusable workflows from Sea-Haven-Industries/.github:

  • CI (PR to main): linting + cdk synth via ci-python-sam.yaml@main
  • CD (push to main): cdk deploy --all via cd-cdk.yaml@main (OIDC auth)

Branch protection on main — all changes through PR.

Setup

  1. Bootstrap CDK: cdk bootstrap aws://{AccountId}/us-east-1
  2. Store Anthropic API keys:
    aws secretsmanager create-secret --name po-ingest/anthropic-api-key --secret-string "sk-ant-..."
    aws secretsmanager create-secret --name workorder-ingest/anthropic-api-key --secret-string "sk-ant-..."
    
  3. Deploy both stacks:
    cd cdk
    pip install -r requirements.txt
    cdk deploy --all
    
  4. CloudFormation outputs include WebUIUrl for each stack's dashboard.

Scripts

Reprocess PO emails (re-run parser against all emails still in S3):

python scripts/reprocess.py            # dry-run
python scripts/reprocess.py --execute  # invoke po-email-processor for each

Backfill verified sites (one-time scan of historical POs):

python scripts/backfill_sites.py

Directory Structure

cdk/
  app.py              # Two stacks: po-ingest + WorkorderIngestStack
  po_stack.py          # Purchase order pipeline resources
  wo_stack.py          # Work order pipeline resources
lambdas/
  po/                  # PO pipeline Lambdas
    email_processor/
    site_extractor/
    web_ui/
  wo/                  # WO pipeline Lambdas
    email_processor/
    web_ui/
shared/
  models.py            # Work order dataclasses/enums
scripts/
  reprocess.py
  backfill_sites.py
tests/
  conftest.py          # AWS env stubs + per-pipeline module loader
  test_ses_auth.py     # Sender-authentication parser tests (INFRA-107)