Coupa PO email ingestion pipeline
Find a file
2026-07-14 09:04:31 -04:00
.github chore(ci): SHA-pin org reusable-workflow caller refs (INFRA-50) (#85) 2026-07-06 18:27:18 -04:00
cdk Bump aws-cdk-lib in /cdk in the minor-and-patch group across 1 directory (#86) 2026-07-07 18:20:00 +00:00
lambdas Update boto3 requirement in /lambdas/wo/web_ui (#92) 2026-07-14 09:04:31 -04: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 docs: document cross-stack DynamoDB data contracts (INFRA-138) (#91) 2026-07-08 16:22:38 -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. 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.

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