procurement-ingest/README.md
Adam Moussa 3f42078f81
Some checks are pending
Deploy / deploy (push) Waiting to run
docs: link Confluence AWS Architecture Map (INFRA-53) (#83)
2026-07-06 17:43:51 -04:00

173 lines
8.8 KiB
Markdown

# Procurement Ingest
![Python](https://img.shields.io/badge/Python-3776AB?logo=python&logoColor=white)
![AWS CDK](https://img.shields.io/badge/AWS-CDK-FF9900?logo=amazonaws&logoColor=white)
![CI](https://github.com/Sea-Haven-Industries/procurement-ingest/actions/workflows/ci.yaml/badge.svg)
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.
## 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.
- **[AWS Architecture Map](https://seahaven.atlassian.net/wiki/spaces/IT/pages/1540098)** (Confluence, IT space, page 1540098)
## 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:
```bash
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:
```bash
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):
```bash
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):
```bash
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
```