Coupa PO email ingestion pipeline
Find a file
Adam Moussa 47fa35688d
Some checks are pending
Deploy / deploy (push) Waiting to run
fix(po): grant KMS on seahaven-dynamodb CMK to purchase-orders consumers (INFRA-104) (#56)
The purchase-orders table was migrated to SSE-KMS (alias/seahaven-dynamodb,
INFRA-95/M-3) out-of-band, but po_stack never declared the key, so
grant_read_write_data did not propagate kms perms. po-email-processor failed
~99.6% of invocations with kms:Decrypt AccessDeniedException, a data-loss
outage on the PO ingestion write path.

- po_stack: declare encryption_key on purchase-orders (reconciles SSE drift;
  no-op against the already-encrypted live table) so the existing grants add
  kms:Decrypt/GenerateDataKey/DescribeKey to EmailProcessor, WebUI, SiteExtractor.
- wo_stack: pre-emptive grant_encrypt_decrypt on the WO processor role ahead of
  the WorkOrders CMK migration (INFRA-6); tables left unencrypted, no table change.

GPT-4.1 cross-review: no blockers.
2026-06-10 19:31:55 -04:00
.github chore(deps): remove blanket aws-cdk-lib dependabot ignore (#47) 2026-06-05 13:57:53 -04:00
cdk fix(po): grant KMS on seahaven-dynamodb CMK to purchase-orders consumers (INFRA-104) (#56) 2026-06-10 19:31:55 -04:00
lambdas Update anthropic requirement in /lambdas/wo/email_processor (#53) 2026-06-10 22:21:45 +00:00
scripts Merge workorder-ingest into unified procurement repo (#22) 2026-05-12 15:21:06 -04:00
.gitignore Merge workorder-ingest into unified procurement repo (#22) 2026-05-12 15:21:06 -04:00
README.md Reconcile IaC with out-of-band DLQ + Function URL changes (INFRA-74, INFRA-41) (#50) 2026-06-08 16:02:29 -04:00
test_local.py Merge workorder-ingest into unified procurement repo (#22) 2026-05-12 15:21:06 -04:00

Procurement Ingest

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. Claude extracts structured JSON (PO number, status, supplier, site code, trade classification, line items, fiscal year).
  5. Conditional write to DynamoDB:
    • new_po — idempotent insert (no-op if PO exists)
    • revision — unconditional overwrite
    • cancellation — marks existing row Cancelled
  6. 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. Claude extracts structured JSON (work order ID, site code, severity, priority, dates, assigned technician).
  6. 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.

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, plus an ALARM-only CloudWatch Errors alarm (Sum, threshold > 0) wired to the shared site-alerts SNS topic.

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.

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