From f67d8b9907ad3089b570b0ad718fedd7e2008abf Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Thu, 23 Jul 2026 19:32:20 -0400 Subject: [PATCH] feat(api): procurement-api read stack + OpenAPI docs (SHOC reconciliation path) (#127) * feat(api): add procurement-api stack - read API + OpenAPI docs page Third CDK stack: API Gateway REST API (IAM SigV4) over both pipelines' tables, replacing SHOC's retired SyncController cross-account DynamoDB scan as the reconciliation/backfill path. - lambdas/api/: handler (healthcheck + docs-token gate + router dispatch), router (single route table), pagination (opaque cursor, hostile -> 400), Decimal-safe serialization, wo_repo/po_repo reads. No VendorReplies. - OpenAPI 3.1 spec as source of truth incl. top-level webhooks section documenting the outbound SHOC feed; phase-2 write endpoints x-planned (router answers 501). Self-contained /docs page, no CDN. - Auth: AWS_IAM on data routes + resource policy scoped to exactly arn:aws:iam::396287094661:role/shoc-backend-dev on GET/*; /docs and /openapi.json carve-out is token-gated in the Lambda via shared web_ui_auth (fail-closed, INFRA-74 posture). - KMS: explicit Decrypt/DescribeKey on the DynamoDB CMK from SSM (name-imported table drops the key association - INFRA-104 class). - Alarms: errors/throttles/duration(p99>=22.5s) + gateway 5xx, ALARM-only to site-alerts. No access logging in v1 (docs ?token= shim stays out of logs); cloud_watch_role=False. - Tests: handler auth-seam + routing + Decimal round-trip; moto cursor pagination incl. hostile cursors; spec<->router drift gate; bundle AST pins for the api command; pytest.ini --cov + loader siblings. - Deploy role: third stack DescribeStacks ARN + procurement-api smoke invoke ARN (re-run create-deploy-role.sh before merge). * harden(api): apply sh-security-review findings to procurement-api Fan-out (6 detectors) + review findings resolved: Correctness / DoS: - pagination: require EXACT key-set match (was subset) so a partial/foreign composite cursor can't reach DynamoDB as an inconsistent ExclusiveStartKey -> ValidationException -> 500; comments Query now pins the cursor's work_order_id to the path entity. - handler: map botocore ValidationException to 400 (defense in depth) so a crafted cursor can't drive the zero-threshold 5xx alarm. - web_ui_auth: compare tokens as bytes; a non-ASCII presented token now fails closed (401) instead of crashing hmac.compare_digest into a 500. Resolves the pre-existing xfail(strict) follow-up test; hardens the web UIs too. Docs page: - typeStr() now escapes the one spec-derived string that reached innerHTML. - spec inlined into the docs breakout guard); /openapi.json still served byte-faithful. - Cache-Control: no-store + Referrer-Policy: no-referrer on docs responses so the ?token= URL stays out of caches/Referer. - spec-drift test asserts the committed spec carries no "-sender-auth-rejected` alarm (see [CloudWatch alarms](#cloudwatch-alarms)) — that alarm pages at ≥1 match in its window, so repeated healthcheck invokes across deploys (two deploys in ~30 min is routine) must never contribute to it. Unit coverage: `lambdas/po/email_processor/tests/test_po_healthcheck.py` and `lambdas/wo/email_processor/tests/test_healthcheck.py`. -**Post-deploy smoke gate.** `scripts/post-deploy-smoke.sh` is wired into the CD workflow as `cd-cdk.yaml`'s `post-deploy-script` input (see [CI/CD](#cicd)) and runs synchronously after every deploy, before the workflow is considered green. It invokes both `po-email-processor` and `workorder-email-processor` with `aws lambda invoke --invocation-type RequestResponse --payload '{"healthcheck": true}'` (region `us-east-1`) and asserts, per function: +**Post-deploy smoke gate.** `scripts/post-deploy-smoke.sh` is wired into the CD workflow as `cd-cdk.yaml`'s `post-deploy-script` input (see [CI/CD](#cicd)) and runs synchronously after every deploy, before the workflow is considered green. It invokes `po-email-processor`, `workorder-email-processor`, and `procurement-api` with `aws lambda invoke --invocation-type RequestResponse --payload '{"healthcheck": true}'` (region `us-east-1`) and asserts, per function: 1. The invoke response's **`FunctionError` field is absent** — this is the load-bearing check. A broken bundle (e.g. an `ImportError` at module init from a missing sibling module) still returns HTTP 200 from the Lambda Invoke API with `FunctionError=Unhandled`; a bare exit-code check on `aws lambda invoke` would false-pass on exactly the failure this gate exists to catch. 2. The returned payload is **exactly** `{"healthcheck": "ok"}`. @@ -248,13 +266,13 @@ The `purchase-orders` DynamoDB table is **owned by this repo's `po-ingest` stack **Consumers (read-only):** -| Repo | How it reads | Purpose | +| Consumer | How it reads | Purpose | |---|---|---| -| `seahaven-slack-bot` | `po-sync` (DynamoDB Streams + daily scan) and `wo-po-lookup` | Daily KB sync + Bedrock agent PO lookups | +| `procurement-api` (this repo) | `Table.from_table_name` + `grant_read_data` | REST reads (`GET /purchase-orders*`) | -The consumer imports the table via `Table.fromTableName(...)` and is granted read-only access (`grantReadData`); it does not own or define it. +`seahaven-slack-bot`, the former external consumer, was decommissioned 2026-07-23 and its cross-account grants removed. External consumers now read through the `procurement-api` REST contract (`lambdas/api/openapi.json`), never by importing the table directly. -**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. +**Schema-coordination rule:** Any change to the `purchase-orders` schema (partition key, item shape, attribute names, streams view type) must be reflected in the API's OpenAPI schemas in the same PR (the spec-drift test pins paths, not item shapes — the schema fields are maintained by hand). Treat schema changes as a contract 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. @@ -284,7 +302,7 @@ Both tables are **owned by this repo's `WorkorderIngestStack`** (`cdk/wo_stack.p Both currently use default DynamoDB encryption — they are **not** yet on the shared customer-managed CMK (`alias/seahaven-dynamodb`, INFRA-95 / M-3); that migration is tracked in INFRA-6. The `workorder-email-processor` role no longer holds a pre-emptive encrypt/decrypt grant on that CMK (removed in the 2026-06-17 security sweep — it was unused while the tables are unencrypted and extended the role's decrypt reach to the CMK protecting `purchase-orders`). Re-add the grant as part of the INFRA-6 migration, at which point `grant_read_write_data` on the then-encrypted tables propagates the needed key permissions automatically. -**Consumer (read-only) — data contract:** `seahaven-slack-bot` imports both tables via `Table.fromTableName(...)` (`grantReadData`) 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 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 configuration (e.g. the INFRA-6 CMK migration) 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). +**Consumers (read-only) — data contract:** `seahaven-slack-bot` (the former external reader) was decommissioned 2026-07-23; its grants are gone. Current consumers: the `procurement-api` Lambda (this repo, `Table.from_table_name` + `grant_read_data`, serving `GET /work-orders*`), and — once the SHOC webhook ships — the `workorder-shoc-emitter` stream consumer. External readers (SHOC) consume through the REST contract (`lambdas/api/openapi.json`) and the webhook contract (`docs/shoc-webhook-contract.md`); both documents' field lists mirror `lambdas/wo/email_processor/persistence.py`. Any change to table name, key schema, attribute names, or encryption configuration (e.g. the INFRA-6 CMK migration) must update those two contracts in the same PR — the tables are imported by name, so there is no compile-time link and breakage surfaces at runtime. **Stream field contract (`site_code`).** Three definitions of "is this a valid site code" have existed in this repo at once. The canonical shape is `derived_fields._STRICT_CODE_RE` (`[A-Z][A-Z0-9]{2,4}`, `fullmatch`) plus its skip-list semantics (a code-shaped token is only a real site code if it is *not* skip-listed, e.g. `LLC`/`INC`/`CORP`/`LTD`/`ATTN`), used by `derive_site_code()` (the `enrich_parsed()` classifier, see the **Derived fields** note in the Purchase Orders flow above). Honestly noted, not papered over: `lambdas/po/site_extractor/handler.py`'s own `SITE_CODE_PATTERN` (`[A-Z]{2,4}\d{1,2}`, prefix-anchored `.match`, digit-requiring) still diverges from the canonical shape as of this phase — it rejects valid all-letter codes like `KLAL` (which surface as permanent `pending-site-review` rows) and accepts overlong junk like `DLI6X`/`SNY55` that the canonical `fullmatch` would not. Reconciling `po-ingest-site-extractor`'s direct-field validation onto the canonical `derived_fields` shape is Phase 6 scope, tracked separately — this paragraph will be updated when it lands. @@ -292,9 +310,7 @@ Both currently use default DynamoDB encryption — they are **not** yet on the s Owned by this repo's `po-ingest` stack (`cdk/po_stack.py`). PK `siteCode` (S); default DynamoDB 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. +**Consumer (read-only) — data contract:** `seahaven-slack-bot` (former reader, incl. the INFRA-138/INFRA-180 `by-state` GSI drift saga) was decommissioned 2026-07-23 — the GSI-drift issue died with it. Current consumer: the `procurement-api` Lambda (`GET /verified-sites*`). The API's `VerifiedSite` schema deliberately promises only pipeline-written fields (`siteCode`, `address`, `city`, `state`, `zip`, `fullAddress`, `locationCode`, `poCount`, `sourcePOs`) — `latitude`/`longitude`/`notes` were manually curated legacy attributes the extractor never writes, and are not part of the contract. ## Documentation @@ -306,7 +322,7 @@ The canonical map of Sea Haven's AWS infrastructure lives in Confluence. This pr GitHub Actions with reusable workflows from `Sea-Haven-Industries/.github` (all pinned to a commit SHA of `main`): - **CI** (`ci.yaml`, PR to `main`): linting + `cdk synth` via `ci-python-sam.yaml`. `cdk synth`'s Docker-bundled asset build for `po-email-processor` and `workorder-email-processor` mounts the widened `../lambdas` asset root (Phase 2, see [Deploy-Pipeline Guards](#deploy-pipeline-guards-phase-0)) as build context, not just each function's own subdirectory — the `exclude` list on both `from_asset` calls strips local-only `__pycache__`/`package/` (and `tests/`) cruft from that wider mount's source fingerprint, so CI's asset hash matches a clean local checkout, and each function's scoped `cp` glob copies only its own pipeline's `*.py` into the zip. (The exclude does not, and under `SOURCE` hashing cannot, keep the *other* pipeline's tracked source out of the fingerprint (see the PO bundling note above on `SOURCE` hashing) — but that source is identical in CI and local, so it does not cause hash divergence.) -- **CD** (`deploy.yaml`, push to `main`): CDK deploy via `cd-cdk.yaml` (OIDC auth), followed by the synchronous `post-deploy-script: scripts/post-deploy-smoke.sh` healthcheck gate (see [Deploy-Pipeline Guards](#deploy-pipeline-guards-phase-0)) — `cd-cdk.yaml`'s `stack-name` input only accepts one stack, so the smoke script itself enumerates both `po-email-processor` and `workorder-email-processor` +- **CD** (`deploy.yaml`, push to `main`): CDK deploy via `cd-cdk.yaml` (OIDC auth), followed by the synchronous `post-deploy-script: scripts/post-deploy-smoke.sh` healthcheck gate (see [Deploy-Pipeline Guards](#deploy-pipeline-guards-phase-0)) — `cd-cdk.yaml`'s `stack-name` input only accepts one stack, so the smoke script itself enumerates `po-email-processor`, `workorder-email-processor`, and `procurement-api` - Plus dependency review and PR labeler workflows on every PR Branch protection on `main` — all changes through PR. @@ -392,12 +408,13 @@ python scripts/backfill_sites.py ``` cdk/ - app.py # Two stacks: po-ingest + WorkorderIngestStack (region-only env) + app.py # Three stacks: po-ingest + WorkorderIngestStack + procurement-api (region-only env) common.py # Phase 4: shared plain-function CDK helpers (alarms, Bedrock grant, # email bucket, processor DLQ, fallback-rate alarm) -- called with # each stack's own scope + literal construct ids, logical-ID-safe po_stack.py # Purchase order pipeline resources wo_stack.py # Work order pipeline resources + procurement_api_stack.py # REST API (IAM SigV4 + resource policy) over both pipelines' tables + token-gated /docs lambdas/ # Phase 2: shared Code.from_asset("../lambdas") bundling root for # BOTH po-email-processor and workorder-email-processor (and, since # Phase 3, both web_ui functions) -- each email-processor command @@ -408,9 +425,18 @@ lambdas/ # Phase 2: shared Code.from_asset("../lambdas") bundling shared/ # Phase 3: single-sourced first-party modules, landed FLAT (no # __init__.py -- bare-name imports) into each bundle by the cp above ses_auth.py # fail-closed SES sender-auth (INFRA-107) -- one copy, one fix - web_ui_auth.py # fail-closed X-Auth-Token gate + token cache (INFRA-74) + web_ui_auth.py # fail-closed X-Auth-Token gate + token cache (INFRA-74; also gates the API /docs routes) email_parsing.py # parse_raw_email (WO superset; returns cc unconditionally) emf.py # generic CloudWatch EMF emitter (dimension-sets pinned once) + api/ # procurement-api Lambda (REST reads over both pipelines' tables) + handler.py # healthcheck early-return + docs-token gate + route dispatch + error mapping + router.py # single route table (spec-drift test pins it to openapi.json) + pagination.py # opaque base64(LastEvaluatedKey) cursor encode/validate (hostile cursor -> 400) + serialization.py # Decimal-safe JSON responses + wo_repo.py # WorkOrders/WorkOrderComments reads (paginated Scan / PK Query) + po_repo.py # purchase-orders/verified-sites reads (no VendorReplies -- dead table) + openapi.json # OpenAPI 3.1 source of truth (paths + outbound `webhooks` section) + docs.html # self-contained reference page (spec inlined at request time; no CDN) po/ # PO pipeline Lambdas email_processor/ # Phase 5: God-handler decomposed into flat siblings (bare-name # imports; the Phase 0/2/3 `cp /email_processor/*.py` @@ -460,7 +486,7 @@ lambdas/ # Phase 2: shared Code.from_asset("../lambdas") bundling scripts/ reprocess.py backfill_sites.py - post-deploy-smoke.sh # CD gate: synchronous healthcheck invoke of both processors, checks FunctionError + post-deploy-smoke.sh # CD gate: synchronous healthcheck invoke of the three smoke-gated functions, checks FunctionError conftest.py # Phase 8: THE repo-root session-invariant conftest. rootdir is pinned by # pytest.ini at repo root, so this loads for every pytest invocation shape -- # including a standalone `pytest lambdas/po/email_processor/tests` -- before diff --git a/cdk/app.py b/cdk/app.py index 6f237be..f1289ab 100644 --- a/cdk/app.py +++ b/cdk/app.py @@ -1,6 +1,7 @@ #!/usr/bin/env python3 import aws_cdk as cdk from po_stack import PoIngestStack +from procurement_api_stack import ProcurementApiStack from wo_stack import WorkorderIngestStack app = cdk.App() @@ -19,4 +20,11 @@ WorkorderIngestStack( env=cdk.Environment(region="us-east-1"), ) +ProcurementApiStack( + app, + "procurement-api", + stack_name="procurement-api", + env=cdk.Environment(region="us-east-1"), +) + app.synth() diff --git a/cdk/procurement_api_stack.py b/cdk/procurement_api_stack.py new file mode 100644 index 0000000..fccb655 --- /dev/null +++ b/cdk/procurement_api_stack.py @@ -0,0 +1,361 @@ +"""procurement-api stack: read-only REST API over both pipelines' tables. + +Third stack in the app. Serves work orders + comments (WO stack tables) and +purchase orders + verified sites (PO stack tables) to SigV4 callers -- the +primary consumer is the SHOC backend, for which this API replaces the retired +SyncController cross-account DynamoDB scan as the reconciliation/backfill +path. Also hosts the token-gated OpenAPI docs page (/docs, /openapi.json). + +Tables are imported by fixed physical name (Table.from_table_name), NOT +passed as cross-stack objects: object passing would synthesize CFN Exports +from the owning stacks and lock them against future changes to the tables. +The one thing name-import does NOT carry is the purchase-orders CMK +association -- see the explicit KMS grant below. +""" + +import aws_cdk as cdk +from aws_cdk import ( + Duration, + Stack, + aws_apigateway as apigateway, + aws_cloudwatch as cloudwatch, + aws_cloudwatch_actions as cw_actions, + aws_dynamodb as dynamodb, + aws_iam as iam, + aws_kms as kms, + aws_lambda as lambda_, + aws_secretsmanager as secretsmanager, + aws_sns as sns, + aws_ssm as ssm, +) +from constructs import Construct + +import common + +# The only cross-account caller. Grants are to this exact role ARN -- future +# shoc-backend-staging/-prod roles are each a deliberate, individually +# cross-reviewed addition (no wildcard/prefix trust). +SHOC_BACKEND_DEV_ROLE_ARN = "arn:aws:iam::396287094661:role/shoc-backend-dev" + +STAGE_NAME = "prod" + + +class ProcurementApiStack(Stack): + def __init__(self, scope: Construct, construct_id: str, **kwargs) -> None: + super().__init__(scope, construct_id, **kwargs) + + alarm_topic = sns.Topic.from_topic_arn( + self, + "SiteAlertsTopic", + f"arn:aws:sns:{self.region}:{self.account}:site-alerts", + ) + + work_orders_table = dynamodb.Table.from_table_name( + self, "WorkOrdersTable", "WorkOrders" + ) + comments_table = dynamodb.Table.from_table_name( + self, "CommentsTable", "WorkOrderComments" + ) + po_table = dynamodb.Table.from_table_name( + self, "PurchaseOrdersTable", "purchase-orders" + ) + verified_sites_table = dynamodb.Table.from_table_name( + self, "VerifiedSitesTable", "verified-sites" + ) + + web_ui_auth_secret = secretsmanager.Secret.from_secret_name_v2( + self, + "WebUiAuthToken", + "procurement-ingest/web-ui-auth-token", + ) + + # purchase-orders is SSE-KMS encrypted with the org DynamoDB CMK. A + # name-imported Table has no encryption-key association, so + # grant_read_data alone leaves the reader without kms:Decrypt and every + # purchase-orders read AccessDenies at runtime (the INFRA-104 failure + # class). Import the key from the same SSM parameter po_stack uses and + # grant it explicitly. + dynamodb_cmk = kms.Key.from_key_arn( + self, + "DynamoDbCmk", + ssm.StringParameter.value_for_string_parameter( + self, "/seahaven/dynamodb/cmk-arn" + ), + ) + + # RETAIN log group (INFRA-114). This is the only stateful resource in + # this otherwise-stateless stack, so it carries the fresh-deploy + # rollback trap: if the FIRST create fails after the group exists, + # rollback deletes everything else but retains the group, and the retry + # CREATE then collides on "/aws/lambda/procurement-api already exists". + # Recovery: delete that log group before re-running a failed first + # deploy (same class as the RETAIN-orphan recovery in the deploy-role + # runbook / reference_cfn_deploy_role_gotchas). + api_log_group = common.make_function_log_group( + self, "ProcurementApi", "procurement-api" + ) + api_fn = lambda_.Function( + self, + "ProcurementApi", + function_name="procurement-api", + runtime=lambda_.Runtime.PYTHON_3_12, + architecture=lambda_.Architecture.ARM_64, + handler="handler.handler", + code=lambda_.Code.from_asset( + "../lambdas", + exclude=["**/__pycache__/**"], + bundling=cdk.BundlingOptions( + image=lambda_.Runtime.PYTHON_3_12.bundling_image, + command=[ + "bash", + "-c", + # Non-recursive glob ships every api/ sibling (the + # allowlist-omission trap from PR #105/PR #2); the spec + # and docs page ride along because the handler serves + # them from its own package dir. web_ui_auth.py must + # land FLAT beside handler.py for the bare import. + "cp api/*.py /asset-output/ && " + "cp api/openapi.json /asset-output/ && " + "cp api/docs.html /asset-output/ && " + "cp shared/web_ui_auth.py /asset-output/ && " + "rm -rf /asset-output/__pycache__", + ], + ), + ), + timeout=Duration.seconds(30), + memory_size=256, + log_group=api_log_group, + environment={ + "WORK_ORDERS_TABLE": work_orders_table.table_name, + "COMMENTS_TABLE": comments_table.table_name, + "PO_TABLE": po_table.table_name, + "VERIFIED_SITES_TABLE": verified_sites_table.table_name, + "WEB_UI_AUTH_TOKEN_SECRET_ARN": web_ui_auth_secret.secret_arn, + }, + ) + + work_orders_table.grant_read_data(api_fn) + comments_table.grant_read_data(api_fn) + po_table.grant_read_data(api_fn) + verified_sites_table.grant_read_data(api_fn) + web_ui_auth_secret.grant_read(api_fn) + api_fn.add_to_role_policy( + iam.PolicyStatement( + actions=["kms:Decrypt", "kms:DescribeKey"], + resources=[dynamodb_cmk.key_arn], + # Scope the grant to the DynamoDB data path only: the role can + # decrypt purchase-orders items via DynamoDB, never call + # kms:Decrypt directly on arbitrary ciphertext under the shared + # org CMK. + conditions={ + "StringEquals": { + "kms:ViaService": f"dynamodb.{self.region}.amazonaws.com" + } + }, + ) + ) + + # Resource policy: once a REST API has one, anything not explicitly + # allowed is denied -- so the docs routes need their own Allow or the + # NONE-auth methods would be black-holed. The docs routes are still + # token-gated inside the Lambda (fail-closed), so this is not an + # unauthenticated data path (INFRA-74 posture). + # + # The SHOC data grant enumerates the exact GET resources rather than + # GET/*: adding a RESOURCE must be as reviewable as adding a PRINCIPAL, + # so a future GET route can't silently inherit cross-account reach + # without a policy diff. (Note: this resource policy binds CROSS-account + # callers only; a same-account principal holding execute-api:Invoke is + # authorized by its own identity policy under AWS union semantics -- it + # is NOT constrained here, including on the planned PATCH/POST methods, + # which is why those also rely on the 501 handler + absent write grant, + # not on this policy, until phase 2.) + shoc_data_resources = [ + f"execute-api:/{STAGE_NAME}/GET/work-orders", + f"execute-api:/{STAGE_NAME}/GET/work-orders/*", + f"execute-api:/{STAGE_NAME}/GET/purchase-orders", + f"execute-api:/{STAGE_NAME}/GET/purchase-orders/*", + f"execute-api:/{STAGE_NAME}/GET/verified-sites", + f"execute-api:/{STAGE_NAME}/GET/verified-sites/*", + ] + api_policy = iam.PolicyDocument( + statements=[ + iam.PolicyStatement( + sid="ShocBackendDevDataRead", + principals=[iam.ArnPrincipal(SHOC_BACKEND_DEV_ROLE_ARN)], + actions=["execute-api:Invoke"], + resources=shoc_data_resources, + ), + # SECURITY INVARIANT: this AnyPrincipal allow is safe only + # while /docs//openapi.json serve static docs (handler + # enforces the shared token, fail-closed). Widening these + # routes to dynamic data, or enabling access logging (which + # would record the docs ?token= shim), requires a security + # re-review + docs-token rotation. + iam.PolicyStatement( + sid="DocsTokenGatedRoutes", + principals=[iam.AnyPrincipal()], + actions=["execute-api:Invoke"], + resources=[ + f"execute-api:/{STAGE_NAME}/GET/docs", + f"execute-api:/{STAGE_NAME}/GET/openapi.json", + ], + ), + ] + ) + + # OPERATIONAL NOTE: API Gateway serves the resource policy from the + # deployed stage snapshot, and CDK's Deployment hash is computed from + # resources/methods, not the RestApi Policy. A later policy-ONLY change + # (e.g. revoking the SHOC role) will UPDATE the RestApi but keep serving + # the old policy until a new Deployment is forced (any method/resource + # change, or a salted deployment). When tightening this policy, force a + # redeploy and verify the effective policy post-deploy. + api = apigateway.RestApi( + self, + "ProcurementRestApi", + rest_api_name="procurement-api", + description=( + "Read API over procurement-ingest work orders + purchase " + "orders; token-gated OpenAPI docs at /docs" + ), + endpoint_types=[apigateway.EndpointType.REGIONAL], + policy=api_policy, + deploy_options=apigateway.StageOptions( + stage_name=STAGE_NAME, + # Bound the blast radius of the unauthenticated /docs routes + # (and the whole API) below the 10k rps account default -- this + # is a low-volume reconciliation/backfill API, not a hot path. + throttling_rate_limit=50, + throttling_burst_limit=100, + ), + # No access/execution logging in v1: avoids the account-level API + # Gateway CloudWatch role prerequisite AND keeps the docs ?token= + # query shim out of any log. Rotate the docs token before ever + # enabling access logging here. + cloud_watch_role=False, + ) + + integration = apigateway.LambdaIntegration(api_fn) + iam_auth = {"authorization_type": apigateway.AuthorizationType.IAM} + + work_orders = api.root.add_resource("work-orders") + work_orders.add_method("GET", integration, **iam_auth) + wo_by_id = work_orders.add_resource("{workOrderId}") + wo_by_id.add_method("GET", integration, **iam_auth) + # Phase-2 planned write endpoints: deployed but the handler answers 501 + # and the role holds no DynamoDB write grant. The resource policy denies + # these to the cross-account SHOC role (GET-only enumeration above); a + # same-account caller is NOT blocked by the resource policy, so the 501 + # handler + absent write grant are the real gate until phase 2 lands the + # deliberate policy + handler + write-grant change with its own review. + wo_by_id.add_method("PATCH", integration, **iam_auth) + wo_comments = wo_by_id.add_resource("comments") + wo_comments.add_method("GET", integration, **iam_auth) + wo_comments.add_method("POST", integration, **iam_auth) + + purchase_orders = api.root.add_resource("purchase-orders") + purchase_orders.add_method("GET", integration, **iam_auth) + purchase_orders.add_resource("{poNumber}").add_method( + "GET", integration, **iam_auth + ) + + verified_sites = api.root.add_resource("verified-sites") + verified_sites.add_method("GET", integration, **iam_auth) + verified_sites.add_resource("{siteCode}").add_method( + "GET", integration, **iam_auth + ) + + api.root.add_resource("docs").add_method( + "GET", + integration, + authorization_type=apigateway.AuthorizationType.NONE, + ) + api.root.add_resource("openapi.json").add_method( + "GET", + integration, + authorization_type=apigateway.AuthorizationType.NONE, + ) + + # --- Alarms (ALARM-only -> site-alerts, NOT_BREACHING) --- + # common.add_standard_lambda_alarms is NOT used here: its duration + # threshold is a fixed 45000 ms (75% of the processors' 60 s timeout), + # which this function's 30 s timeout can never reach. Same idiom, + # right-sized thresholds. + api_fn.metric_errors(period=Duration.minutes(5), statistic="Sum").create_alarm( + self, + "ProcurementApiErrorsAlarm", + alarm_name="procurement-api-errors", + alarm_description=( + "procurement-api Lambda raised (bundle/init failures; the " + "handler catches request errors, so any signal here is " + "structural)" + ), + threshold=0, + evaluation_periods=1, + comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD, + treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING, + ).add_alarm_action(cw_actions.SnsAction(alarm_topic)) + + api_fn.metric_throttles( + period=Duration.minutes(5), statistic="Sum" + ).create_alarm( + self, + "ProcurementApiThrottlesAlarm", + alarm_name="procurement-api-throttles", + alarm_description="procurement-api Lambda throttled", + threshold=0, + evaluation_periods=1, + comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD, + treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING, + ).add_alarm_action(cw_actions.SnsAction(alarm_topic)) + + api_fn.metric_duration( + period=Duration.minutes(5), statistic="p99" + ).create_alarm( + self, + "ProcurementApiDurationAlarm", + alarm_name="procurement-api-duration", + alarm_description=( + "procurement-api p99 duration >= 22.5s (75% of the 30s " + "timeout; scans degrading toward timeout)" + ), + threshold=22500, + evaluation_periods=3, + datapoints_to_alarm=2, + comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_OR_EQUAL_TO_THRESHOLD, + treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING, + ).add_alarm_action(cw_actions.SnsAction(alarm_topic)) + + # Gateway-side 5xx: catches what the Lambda's own Errors metric can't + # (the handler returns clean 500s; integration faults surface here). + # No 4XX alarm -- 401/403/404 are expected traffic. + cloudwatch.Metric( + namespace="AWS/ApiGateway", + metric_name="5XXError", + dimensions_map={"ApiName": "procurement-api", "Stage": STAGE_NAME}, + period=Duration.minutes(5), + statistic="Sum", + ).create_alarm( + self, + "ProcurementApi5xxAlarm", + alarm_name="procurement-api-5xx", + alarm_description="procurement-api gateway 5XX responses", + threshold=0, + evaluation_periods=1, + comparison_operator=cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD, + treat_missing_data=cloudwatch.TreatMissingData.NOT_BREACHING, + ).add_alarm_action(cw_actions.SnsAction(alarm_topic)) + + cdk.CfnOutput( + self, + "ApiEndpointUrl", + value=api.url, + description="procurement-api invoke URL (stage prod)", + ) + cdk.CfnOutput( + self, + "ProcurementApiFunctionArn", + value=api_fn.function_arn, + description="ARN of the procurement-api Lambda", + ) diff --git a/infra/deploy-role/README.md b/infra/deploy-role/README.md index 5dbfec5..94ea7a1 100644 --- a/infra/deploy-role/README.md +++ b/infra/deploy-role/README.md @@ -10,19 +10,20 @@ stays untouched until decommission as the emergency mgmt deploy path. | File | Purpose | |---|---| | `trust-policy.json` | OIDC trust: `repo:Sea-Haven-Industries/procurement-ingest:ref:refs/heads/main` only | -| `permissions-policy.json` | `sts:AssumeRole` on the four `cdk-hnb659fds-*` bootstrap roles, `cloudformation:DescribeStacks` scoped to this repo's stacks + `CDKToolkit` (cd-cdk health check), and `lambda:InvokeFunction` on exactly the two email-processor function ARNs (post-deploy smoke gate) | +| `permissions-policy.json` | `sts:AssumeRole` on the four `cdk-hnb659fds-*` bootstrap roles, `cloudformation:DescribeStacks` scoped to this repo's stacks + `CDKToolkit` (cd-cdk health check), and `lambda:InvokeFunction` on exactly the three smoke-gated function ARNs (post-deploy smoke gate) | | `create-deploy-role.sh` | Idempotent create-or-update from the two JSON files, profile `seahaven-prod` | -> **Maintenance note:** `DescribeStacks` is scoped to `stack/po-ingest/*`, `stack/WorkorderIngestStack/*`, and `stack/CDKToolkit/*`. If a third stack is ever added to this CDK app, add its ARN pattern here and re-run the review-then-apply flow — otherwise the cd-cdk health check on the new stack will `AccessDenied`. +> **Maintenance note:** `DescribeStacks` is scoped to `stack/po-ingest/*`, `stack/WorkorderIngestStack/*`, `stack/procurement-api/*` (added with the procurement-api stack), and `stack/CDKToolkit/*`. If another stack is ever added to this CDK app, add its ARN pattern here and re-run the review-then-apply flow — otherwise the cd-cdk health check on the new stack will `AccessDenied`. ## Why the SmokeInvokeLambda statement exists `deploy.yaml` runs `scripts/post-deploy-smoke.sh` under the deploy role's own session, not the assumed `cdk-*` roles. Without `lambda:InvokeFunction` on the -two function ARNs the smoke gate hits AccessDenied and every deploy fails -closed. The mgmt-era grant was applied out-of-band and undocumented; keeping -it in these reviewed artifacts closes that gap. Scope it to exactly the two -ARNs, never `Resource: "*"`. +smoke-gated function ARNs (the two email processors + `procurement-api`) the +smoke gate hits AccessDenied and every deploy fails closed. The mgmt-era grant +was applied out-of-band and undocumented; keeping it in these reviewed +artifacts closes that gap. Scope it to exactly the named ARNs, never +`Resource: "*"`. ## Change process diff --git a/infra/deploy-role/permissions-policy.json b/infra/deploy-role/permissions-policy.json index e29ba7e..fc7261b 100644 --- a/infra/deploy-role/permissions-policy.json +++ b/infra/deploy-role/permissions-policy.json @@ -19,6 +19,7 @@ "Resource": [ "arn:aws:cloudformation:us-east-1:011934824531:stack/po-ingest/*", "arn:aws:cloudformation:us-east-1:011934824531:stack/WorkorderIngestStack/*", + "arn:aws:cloudformation:us-east-1:011934824531:stack/procurement-api/*", "arn:aws:cloudformation:us-east-1:011934824531:stack/CDKToolkit/*" ], "Condition": { @@ -33,7 +34,8 @@ "Action": "lambda:InvokeFunction", "Resource": [ "arn:aws:lambda:us-east-1:011934824531:function:po-email-processor", - "arn:aws:lambda:us-east-1:011934824531:function:workorder-email-processor" + "arn:aws:lambda:us-east-1:011934824531:function:workorder-email-processor", + "arn:aws:lambda:us-east-1:011934824531:function:procurement-api" ] } ] diff --git a/lambdas/api/docs.html b/lambdas/api/docs.html new file mode 100644 index 0000000..06f5aea --- /dev/null +++ b/lambdas/api/docs.html @@ -0,0 +1,229 @@ + + + + + +Procurement Ingest API + + + +

Loading spec…

+ + + + diff --git a/lambdas/api/handler.py b/lambdas/api/handler.py new file mode 100644 index 0000000..283bcb8 --- /dev/null +++ b/lambdas/api/handler.py @@ -0,0 +1,193 @@ +"""Procurement API Lambda (API Gateway REST proxy integration). + +Auth is split by route class and enforced at two layers: +- Data routes: AWS_IAM at the gateway (SigV4; cross-account callers allowed by + the API resource policy). The handler does NOT re-check a token there -- + authorization is API Gateway's job on those routes. +- Docs routes (/docs, /openapi.json): reachable at the gateway (auth NONE + + resource-policy carve-out) but the handler fails closed on the shared + header token via web_ui_auth (same secret + constant-time compare as the + web UIs), so they are never an unauthenticated data path (INFRA-74). +""" + +import logging +from pathlib import Path + +import po_repo +import wo_repo +from botocore.exceptions import ClientError +from pagination import BadCursor, clamp_limit +from router import DATA_ROUTES, DOCS_ROUTES, PLANNED_ROUTES +from serialization import error_response, json_response +from web_ui_auth import is_authenticated + +logger = logging.getLogger() +logger.setLevel(logging.INFO) + +_MODULE_DIR = Path(__file__).resolve().parent +_SPEC_PATH = _MODULE_DIR / "openapi.json" +_DOCS_PATH = _MODULE_DIR / "docs.html" +# The docs page template carries this placeholder where the spec JSON is +# inlined, so /docs is a single token-gated request (a browser can't attach +# the auth header to a follow-up asset fetch). +_SPEC_PLACEHOLDER = "__OPENAPI_SPEC_JSON__" + +_spec_cache = None +_docs_cache = None + +_REPO_FUNCS = { + "list_work_orders": wo_repo.list_work_orders, + "get_work_order": wo_repo.get_work_order, + "list_comments": wo_repo.list_comments, + "list_purchase_orders": po_repo.list_purchase_orders, + "get_purchase_order": po_repo.get_purchase_order, + "list_verified_sites": po_repo.list_verified_sites, + "get_verified_site": po_repo.get_verified_site, +} + +_PATH_PARAM_BY_RESOURCE = { + "/work-orders/{workOrderId}": "workOrderId", + "/work-orders/{workOrderId}/comments": "workOrderId", + "/purchase-orders/{poNumber}": "poNumber", + "/verified-sites/{siteCode}": "siteCode", +} + + +# Headers on the docs responses: the token rides in the ?token= query shim, +# so keep the token-keyed URL and page out of shared/browser caches and out of +# any Referer sent to a followed link. +_DOCS_SECURITY_HEADERS = { + "Cache-Control": "no-store", + "Referrer-Policy": "no-referrer", +} + + +def _load_spec() -> str: + """Raw spec bytes, served verbatim at /openapi.json (byte-faithful JSON).""" + global _spec_cache + if _spec_cache is None: + _spec_cache = _SPEC_PATH.read_text(encoding="utf-8") + return _spec_cache + + +def _load_docs_html() -> str: + global _docs_cache + if _docs_cache is None: + # Neutralize any " block + # where the HTML parser ends the element at the first literal "<" + # sequence regardless of JSON quoting. "<" is the same JSON + # string value ("<") to JSON.parse but can never close the script tag. + inlined_spec = _load_spec().replace("<", "\\u003c") + _docs_cache = _DOCS_PATH.read_text(encoding="utf-8").replace( + _SPEC_PLACEHOLDER, inlined_spec + ) + return _docs_cache + + +def _with_token_shim(event: dict) -> dict: + """Copy a ?token= query parameter into an x-auth-token header. + + Browsers can't set headers on plain navigation, so /docs accepts the + shared token as a query parameter too. The constant-time compare still + happens inside web_ui_auth -- this only synthesizes the header on a + shallow copy. Acceptable only while API Gateway access logging stays off + (nothing at the gateway records the query string); rotate the token if + access logging is ever enabled. + """ + qs = event.get("queryStringParameters") or {} + token = qs.get("token") + if not token: + return event + shimmed = dict(event) + headers = dict(event.get("headers") or {}) + headers["x-auth-token"] = token + shimmed["headers"] = headers + return shimmed + + +def _docs_response(resource: str, event: dict) -> dict: + # SECURITY INVARIANT: these routes serve ONLY the committed spec and the + # static docs page -- never table data. The gateway resource policy allows + # Principal "*" on exactly these two GETs on the strength of that; serving + # anything dynamic here requires a resource-policy + security re-review. + if not is_authenticated(_with_token_shim(event)): + return error_response(401, "unauthorized") + if resource == "/openapi.json": + return { + "statusCode": 200, + "headers": {"Content-Type": "application/json", **_DOCS_SECURITY_HEADERS}, + "body": _load_spec(), + } + return { + "statusCode": 200, + "headers": { + "Content-Type": "text/html; charset=utf-8", + **_DOCS_SECURITY_HEADERS, + }, + "body": _load_docs_html(), + } + + +def _data_response(route_key: tuple, event: dict) -> dict: + method, resource = route_key + func = _REPO_FUNCS[DATA_ROUTES[route_key]] + qs = event.get("queryStringParameters") or {} + path_params = event.get("pathParameters") or {} + + if DATA_ROUTES[route_key].startswith("list_"): + limit = clamp_limit(qs.get("limit")) + cursor = qs.get("cursor") + if resource in _PATH_PARAM_BY_RESOURCE: + entity_id = path_params.get(_PATH_PARAM_BY_RESOURCE[resource], "") + items, next_cursor = func(entity_id, limit, cursor) + else: + items, next_cursor = func(limit, cursor) + return json_response(200, {"items": items, "next_cursor": next_cursor}) + + entity_id = path_params.get(_PATH_PARAM_BY_RESOURCE[resource], "") + item = func(entity_id) + if item is None: + return error_response(404, "not found") + return json_response(200, item) + + +def _dispatch(route_key: tuple, event: dict) -> dict: + if route_key in PLANNED_ROUTES: + return error_response(501, "planned endpoint - not implemented (phase 2)") + if route_key in DOCS_ROUTES: + return _docs_response(route_key[1], event) + if route_key in DATA_ROUTES: + return _data_response(route_key, event) + return error_response(404, "not found") + + +def handler(event, context): + # Deploy-guard healthcheck: a direct-invoke {"healthcheck": true} probe + # returns before any routing/auth so the post-deploy smoke gate can verify + # the bundle imports and the runtime boots. + if isinstance(event, dict) and event.get("healthcheck") is True: + return {"healthcheck": "ok"} + + method = (event.get("httpMethod") or "").upper() + resource = event.get("resource") or "" + + try: + return _dispatch((method, resource), event) + except BadCursor as exc: + return error_response(400, str(exc)) + except ClientError as exc: + # A client-supplied cursor that survives validation but is still + # inconsistent at the data layer makes DynamoDB raise + # ValidationException; map it to 400, not 500, so a crafted cursor + # can't drive the 5xx alarm. Any other AWS error is a real 500. + if exc.response.get("Error", {}).get("Code") == "ValidationException": + return error_response(400, "cursor is not valid") + logger.exception("AWS error serving %s %s", method, resource) + return error_response(500, "internal error") + except Exception: + # A raised exception would surface as an opaque 502 from the proxy + # integration; return a clean 500 instead. The API Gateway 5XX alarm + # pages on these; the exception (never the request token) is logged. + logger.exception("Unhandled error serving %s %s", method, resource) + return error_response(500, "internal error") diff --git a/lambdas/api/openapi.json b/lambdas/api/openapi.json new file mode 100644 index 0000000..2edde02 --- /dev/null +++ b/lambdas/api/openapi.json @@ -0,0 +1,1022 @@ +{ + "openapi": "3.1.0", + "info": { + "title": "Procurement Ingest API", + "version": "1.0.0", + "description": "Read API over the procurement-ingest pipelines (work orders + purchase orders), plus the outbound SHOC work-order webhook feed (see `webhooks`).\n\n**Purpose:** reconciliation and backfill for downstream consumers (primarily SHOC) - this API replaces SHOC's retired SyncController DynamoDB scan. Listings are **unordered** paginated scans: follow `next_cursor` until it is `null`. `cursor` is opaque; a malformed cursor returns `400`. Field names mirror the DynamoDB attributes written by the pipelines (source of truth: `lambdas/wo/email_processor/persistence.py` and `lambdas/po/email_processor/persistence.py`).\n\n**Auth:** data endpoints require AWS IAM SigV4 (service `execute-api`, region `us-east-1`); cross-account callers must also be allowed by the API resource policy. `/docs` and `/openapi.json` use the shared docs token instead (header `X-Auth-Token`, or `?token=` in a browser). No CORS is configured (server-to-server and Postman callers only).\n\n**Write endpoints** marked `x-planned` are phase 2: documented here for contract visibility, the API answers `501` until they ship." + }, + "servers": [ + { + "url": "https://{apiId}.execute-api.us-east-1.amazonaws.com/prod", + "description": "seahaven-prod (011934824531). The concrete apiId is in the procurement-api stack output ApiEndpointUrl.", + "variables": { + "apiId": { + "default": "SEE-STACK-OUTPUT" + } + } + } + ], + "security": [ + { + "sigv4": [] + } + ], + "paths": { + "/work-orders": { + "get": { + "operationId": "listWorkOrders", + "summary": "List work orders (unordered, paginated)", + "parameters": [ + { + "$ref": "#/components/parameters/Limit" + }, + { + "$ref": "#/components/parameters/Cursor" + } + ], + "responses": { + "200": { + "description": "One page of work orders.", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["items", "next_cursor"], + "properties": { + "items": { + "type": "array", + "items": { + "$ref": "#/components/schemas/WorkOrder" + } + }, + "next_cursor": { + "$ref": "#/components/schemas/NextCursor" + } + } + } + } + } + }, + "400": { + "$ref": "#/components/responses/BadRequest" + }, + "403": { + "$ref": "#/components/responses/Forbidden" + } + } + } + }, + "/work-orders/{workOrderId}": { + "get": { + "operationId": "getWorkOrder", + "summary": "Get one work order", + "parameters": [ + { + "$ref": "#/components/parameters/WorkOrderId" + } + ], + "responses": { + "200": { + "description": "The work order.", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/WorkOrder" + } + } + } + }, + "403": { + "$ref": "#/components/responses/Forbidden" + }, + "404": { + "$ref": "#/components/responses/NotFound" + } + } + }, + "patch": { + "x-planned": true, + "operationId": "patchWorkOrder", + "summary": "PLANNED (phase 2): update dispatch fields on a work order", + "description": "Not implemented - returns 501. Phase-2 write-back for SHOC dispatch workflow (status/assignment). Writes will stamp `write_origin: shoc-write-api` so the outbound webhook never echoes SHOC's own writes back at it. Ships with its own IAM diff and cross-family review.", + "responses": { + "501": { + "$ref": "#/components/responses/NotImplemented" + } + } + } + }, + "/work-orders/{workOrderId}/comments": { + "get": { + "operationId": "listWorkOrderComments", + "summary": "List comments/events for a work order (paginated)", + "parameters": [ + { + "$ref": "#/components/parameters/WorkOrderId" + }, + { + "$ref": "#/components/parameters/Limit" + }, + { + "$ref": "#/components/parameters/Cursor" + } + ], + "responses": { + "200": { + "description": "One page of comments (all source-email event records: comments, updates, cancellations).", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["items", "next_cursor"], + "properties": { + "items": { + "type": "array", + "items": { + "$ref": "#/components/schemas/WorkOrderComment" + } + }, + "next_cursor": { + "$ref": "#/components/schemas/NextCursor" + } + } + } + } + } + }, + "400": { + "$ref": "#/components/responses/BadRequest" + }, + "403": { + "$ref": "#/components/responses/Forbidden" + } + } + }, + "post": { + "x-planned": true, + "operationId": "createWorkOrderComment", + "summary": "PLANNED (phase 2): append a SHOC-authored comment", + "description": "Not implemented - returns 501. Phase-2 write-back: SHOC dispatch notes land in WorkOrderComments with `write_origin: shoc-write-api` (append-only; no field conflicts with the email pipeline). Ships with its own IAM diff and cross-family review.", + "responses": { + "501": { + "$ref": "#/components/responses/NotImplemented" + } + } + } + }, + "/purchase-orders": { + "get": { + "operationId": "listPurchaseOrders", + "summary": "List purchase orders (unordered, paginated)", + "parameters": [ + { + "$ref": "#/components/parameters/Limit" + }, + { + "$ref": "#/components/parameters/Cursor" + } + ], + "responses": { + "200": { + "description": "One page of purchase orders.", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["items", "next_cursor"], + "properties": { + "items": { + "type": "array", + "items": { + "$ref": "#/components/schemas/PurchaseOrder" + } + }, + "next_cursor": { + "$ref": "#/components/schemas/NextCursor" + } + } + } + } + } + }, + "400": { + "$ref": "#/components/responses/BadRequest" + }, + "403": { + "$ref": "#/components/responses/Forbidden" + } + } + } + }, + "/purchase-orders/{poNumber}": { + "get": { + "operationId": "getPurchaseOrder", + "summary": "Get one purchase order", + "parameters": [ + { + "name": "poNumber", + "in": "path", + "required": true, + "schema": { + "type": "string" + }, + "description": "Coupa PO number, e.g. `2D-22030794`." + } + ], + "responses": { + "200": { + "description": "The purchase order.", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/PurchaseOrder" + } + } + } + }, + "403": { + "$ref": "#/components/responses/Forbidden" + }, + "404": { + "$ref": "#/components/responses/NotFound" + } + } + } + }, + "/verified-sites": { + "get": { + "operationId": "listVerifiedSites", + "summary": "List verified Amazon sites (unordered, paginated)", + "parameters": [ + { + "$ref": "#/components/parameters/Limit" + }, + { + "$ref": "#/components/parameters/Cursor" + } + ], + "responses": { + "200": { + "description": "One page of verified sites.", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["items", "next_cursor"], + "properties": { + "items": { + "type": "array", + "items": { + "$ref": "#/components/schemas/VerifiedSite" + } + }, + "next_cursor": { + "$ref": "#/components/schemas/NextCursor" + } + } + } + } + } + }, + "400": { + "$ref": "#/components/responses/BadRequest" + }, + "403": { + "$ref": "#/components/responses/Forbidden" + } + } + } + }, + "/verified-sites/{siteCode}": { + "get": { + "operationId": "getVerifiedSite", + "summary": "Get one verified site", + "parameters": [ + { + "name": "siteCode", + "in": "path", + "required": true, + "schema": { + "type": "string" + }, + "description": "Amazon site code, e.g. `JFK8`." + } + ], + "responses": { + "200": { + "description": "The verified site.", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/VerifiedSite" + } + } + } + }, + "403": { + "$ref": "#/components/responses/Forbidden" + }, + "404": { + "$ref": "#/components/responses/NotFound" + } + } + } + }, + "/docs": { + "get": { + "operationId": "getDocs", + "summary": "This documentation page (token-gated)", + "security": [ + { + "docsToken": [] + } + ], + "parameters": [ + { + "name": "token", + "in": "query", + "required": false, + "schema": { + "type": "string" + }, + "description": "Browser convenience: the docs token as a query parameter (browsers can't set headers on navigation). Prefer the X-Auth-Token header from tooling." + } + ], + "responses": { + "200": { + "description": "Self-contained HTML reference page.", + "content": { + "text/html": {} + } + }, + "401": { + "$ref": "#/components/responses/Unauthorized" + } + } + } + }, + "/openapi.json": { + "get": { + "operationId": "getOpenApiSpec", + "summary": "This spec (token-gated)", + "security": [ + { + "docsToken": [] + } + ], + "responses": { + "200": { + "description": "The OpenAPI 3.1 document.", + "content": { + "application/json": {} + } + }, + "401": { + "$ref": "#/components/responses/Unauthorized" + } + } + } + } + }, + "webhooks": { + "work_order.created": { + "post": { + "summary": "Outbound: a work order was created", + "description": "Sent by `workorder-shoc-emitter` (seahaven-prod) to the configured SHOC endpoint. Full contract incl. HMAC verification, ordering, and retry semantics: `docs/shoc-webhook-contract.md` (Rev 2026-07-23). Requests carry `X-SH-Timestamp`, `X-SH-Key-Id`, and `X-SH-Signature: v1=hex(HMAC_SHA256(secret, \"{timestamp}.{raw_body}\"))`; verify over the raw body, constant-time, +/-300s window, fail closed.", + "requestBody": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/WorkOrderEventEnvelope" + } + } + } + }, + "responses": { + "2XX": { + "description": "Delivered. Respond fast (<10s); accept-and-enqueue if processing is slow. 429/5xx/timeouts are retried in order for 24h; other 4xx park immediately." + } + } + } + }, + "work_order.updated": { + "post": { + "summary": "Outbound: a work order changed", + "description": "Same envelope and semantics as work_order.created; `data` is the full current state, not a diff.", + "requestBody": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/WorkOrderEventEnvelope" + } + } + } + }, + "responses": { + "2XX": { + "description": "Delivered." + } + } + } + }, + "work_order.cancelled": { + "post": { + "summary": "Outbound: a work order transitioned to cancelled", + "description": "A specialization of work_order.updated (same body) emitted when `wo_status` transitions to `cancelled`.", + "requestBody": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/WorkOrderEventEnvelope" + } + } + } + }, + "responses": { + "2XX": { + "description": "Delivered." + } + } + } + }, + "work_order.comment_added": { + "post": { + "summary": "Outbound: a comment/event record was ingested", + "description": "One per source email (comments, updates, and cancellation event records). Dedupe on `delivery_id` or `data.comment_id`. May occasionally arrive before the work order's created event - upsert a skeleton work order and let the state event backfill it.", + "requestBody": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/CommentEventEnvelope" + } + } + } + }, + "responses": { + "2XX": { + "description": "Delivered." + } + } + } + } + }, + "components": { + "securitySchemes": { + "sigv4": { + "type": "apiKey", + "name": "Authorization", + "in": "header", + "x-amazon-apigateway-authtype": "awsSigv4", + "description": "AWS IAM SigV4 (service execute-api, region us-east-1). Cross-account callers must be allowed by the API resource policy; Postman signs natively via Authorization type 'AWS Signature'." + }, + "docsToken": { + "type": "apiKey", + "name": "X-Auth-Token", + "in": "header", + "description": "Shared docs token (secret procurement-ingest/web-ui-auth-token). Docs routes only." + } + }, + "parameters": { + "Limit": { + "name": "limit", + "in": "query", + "required": false, + "schema": { + "type": "integer", + "minimum": 1, + "maximum": 500, + "default": 100 + }, + "description": "Page size; values outside 1-500 are clamped." + }, + "Cursor": { + "name": "cursor", + "in": "query", + "required": false, + "schema": { + "type": "string" + }, + "description": "Opaque pagination cursor from the previous page's `next_cursor`. Malformed cursors return 400." + }, + "WorkOrderId": { + "name": "workOrderId", + "in": "path", + "required": true, + "schema": { + "type": "string", + "pattern": "^[0-9]+$" + }, + "description": "Numeric APM work-order id, e.g. `11144580730`." + } + }, + "responses": { + "BadRequest": { + "description": "Malformed cursor or limit.", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Error" + } + } + } + }, + "Unauthorized": { + "description": "Missing or wrong docs token.", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Error" + } + } + } + }, + "Forbidden": { + "description": "SigV4 auth failed or the caller is not allowed by the API resource policy (returned by API Gateway, not the Lambda)." + }, + "NotFound": { + "description": "No record with that id.", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Error" + } + } + } + }, + "NotImplemented": { + "description": "Planned phase-2 endpoint; not implemented yet.", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Error" + } + } + } + } + }, + "schemas": { + "Error": { + "type": "object", + "required": ["error"], + "properties": { + "error": { + "type": "string" + } + } + }, + "NextCursor": { + "type": ["string", "null"], + "description": "Pass as `cursor` to fetch the next page; `null` means this is the last page." + }, + "WorkOrder": { + "type": "object", + "description": "Mirrors the WorkOrders DynamoDB item (PK work_order_id). Fields absent from the source email are absent or null.", + "required": ["work_order_id"], + "properties": { + "work_order_id": { + "type": "string", + "pattern": "^[0-9]+$" + }, + "wo_status": { + "type": ["string", "null"], + "enum": ["new", "assigned", "in_progress", "on_hold", "completed", "cancelled", "unknown", null], + "description": "`unknown` is a real emitted value - map it explicitly." + }, + "description": { + "type": ["string", "null"] + }, + "customer": { + "type": ["string", "null"] + }, + "site_code": { + "type": ["string", "null"] + }, + "building": { + "type": ["string", "null"] + }, + "address": { + "type": ["string", "null"] + }, + "severity": { + "type": ["string", "null"] + }, + "priority": { + "type": ["string", "null"] + }, + "assigned_to": { + "type": ["string", "null"] + }, + "date_reported": { + "type": ["string", "null"] + }, + "scheduled_start": { + "type": ["string", "null"] + }, + "due_date": { + "type": ["string", "null"] + }, + "record_type": { + "type": ["string", "null"], + "enum": ["new_work_order", "update", "comment", "cancellation", null], + "description": "Type of the most recent source email." + }, + "created_at": { + "type": ["string", "null"], + "description": "ISO 8601 with UTC offset." + }, + "updated_at": { + "type": ["string", "null"], + "description": "ISO 8601 with UTC offset." + }, + "source_email_s3_key": { + "type": ["string", "null"] + } + }, + "additionalProperties": true + }, + "WorkOrderComment": { + "type": "object", + "description": "Mirrors the WorkOrderComments DynamoDB item (PK work_order_id, SK comment_id). One row per source email event.", + "required": ["work_order_id", "comment_id"], + "properties": { + "work_order_id": { + "type": "string" + }, + "comment_id": { + "type": "string", + "description": "`{work_order_id}#{time|nocomment}#{sha256(s3_key)[:12]}` - unique per source email, stable across retries; a dedupe key." + }, + "record_type": { + "type": ["string", "null"], + "enum": ["new_work_order", "update", "comment", "cancellation", null] + }, + "commenter": { + "type": ["string", "null"] + }, + "text": { + "type": ["string", "null"] + }, + "created_at": { + "type": ["string", "null"], + "description": "Display timestamp (comment time when the source email carried one)." + }, + "ingested_at": { + "type": ["string", "null"], + "description": "When the pipeline wrote the row; ISO 8601 with UTC offset." + }, + "source_email_s3_key": { + "type": ["string", "null"] + } + }, + "additionalProperties": true + }, + "PurchaseOrder": { + "type": "object", + "description": "Mirrors the purchase-orders DynamoDB item (PK po_number). Shape follows the Coupa extraction schema; older records may lack newer fields.", + "required": ["po_number"], + "properties": { + "po_number": { + "type": "string" + }, + "po_status": { + "type": ["string", "null"] + }, + "email_type": { + "type": ["string", "null"] + }, + "source_system": { + "type": ["string", "null"] + }, + "submitted_by": { + "type": ["string", "null"] + }, + "on_behalf_of": { + "type": ["string", "null"] + }, + "order_date": { + "type": ["string", "null"] + }, + "revision_date": { + "type": ["string", "null"] + }, + "payment_terms": { + "type": ["string", "null"] + }, + "requisition_number": { + "type": ["string", "null"] + }, + "department": { + "type": ["string", "null"] + }, + "view_order_url": { + "type": ["string", "null"] + }, + "supplier": { + "type": ["object", "null"], + "properties": { + "name": { + "type": ["string", "null"] + } + }, + "additionalProperties": true + }, + "site_code": { + "type": ["string", "null"] + }, + "ship_to": { + "type": ["object", "null"], + "properties": { + "name": { + "type": ["string", "null"] + }, + "address": { + "type": ["string", "null"] + }, + "street": { + "type": ["string", "null"] + }, + "city": { + "type": ["string", "null"] + }, + "state": { + "type": ["string", "null"] + }, + "zip": { + "type": ["string", "null"] + }, + "location_code": { + "type": ["string", "null"] + }, + "attn": { + "type": ["string", "null"] + } + }, + "additionalProperties": true + }, + "total_amount": { + "type": ["number", "null"], + "description": "Stored as Decimal; serialized as a JSON number." + }, + "currency": { + "type": ["string", "null"] + }, + "fiscal_year": { + "type": ["string", "null"] + }, + "trade": { + "type": ["string", "null"] + }, + "coupa_category": { + "type": ["string", "null"] + }, + "line_items": { + "type": ["array", "null"], + "items": { + "type": "object", + "properties": { + "description": { + "type": ["string", "null"] + }, + "quantity": { + "type": ["number", "null"] + }, + "unit": { + "type": ["string", "null"] + }, + "price": { + "type": ["number", "null"] + }, + "amount": { + "type": ["number", "null"] + }, + "currency": { + "type": ["string", "null"] + }, + "need_by": { + "type": ["string", "null"] + }, + "category": { + "type": ["string", "null"] + }, + "account_code": { + "type": ["string", "null"] + }, + "period": { + "type": ["string", "null"] + } + }, + "additionalProperties": true + } + }, + "cancelled_at": { + "type": ["string", "null"] + }, + "processed_at": { + "type": ["string", "null"] + }, + "raw_s3_key": { + "type": ["string", "null"] + }, + "data_source": { + "type": ["string", "null"] + } + }, + "additionalProperties": true + }, + "VerifiedSite": { + "type": "object", + "description": "Mirrors the verified-sites DynamoDB item (PK siteCode), maintained by the PO site extractor.", + "required": ["siteCode"], + "properties": { + "siteCode": { + "type": "string" + }, + "address": { + "type": ["string", "null"] + }, + "city": { + "type": ["string", "null"] + }, + "state": { + "type": ["string", "null"] + }, + "zip": { + "type": ["string", "null"] + }, + "fullAddress": { + "type": ["string", "null"] + }, + "locationCode": { + "type": ["string", "null"] + }, + "poCount": { + "type": ["integer", "null"], + "description": "Count of POs that referenced this site." + }, + "sourcePOs": { + "type": ["array", "null"], + "items": { + "type": "string" + }, + "description": "PO numbers that referenced this site (string set; serialized sorted)." + } + }, + "additionalProperties": true + }, + "WebhookEnvelope": { + "type": "object", + "description": "Common webhook envelope (contract section 3). Headers: X-SH-Timestamp (unix seconds), X-SH-Key-Id, X-SH-Signature (v1=hex HMAC-SHA256 over '{timestamp}.{raw_body}').", + "required": ["schema_version", "delivery_id", "event_type", "occurred_at", "source", "replay", "data"], + "properties": { + "schema_version": { + "type": "integer", + "description": "Reject versions you don't know." + }, + "delivery_id": { + "type": "string", + "description": "Unique per source event, stable across producer retries - the idempotency key." + }, + "event_type": { + "type": "string", + "enum": ["work_order.created", "work_order.updated", "work_order.cancelled", "work_order.comment_added"] + }, + "occurred_at": { + "type": "string", + "description": "ISO 8601 with UTC offset." + }, + "source": { + "type": "string", + "const": "procurement-ingest/workorder-shoc-emitter" + }, + "replay": { + "type": "boolean", + "description": "true when re-sent by the operator replay tool; same idempotency rules." + }, + "data": { + "type": "object" + } + } + }, + "WorkOrderEventEnvelope": { + "allOf": [ + { + "$ref": "#/components/schemas/WebhookEnvelope" + }, + { + "type": "object", + "properties": { + "data": { + "$ref": "#/components/schemas/WorkOrderEventData" + } + } + } + ] + }, + "CommentEventEnvelope": { + "allOf": [ + { + "$ref": "#/components/schemas/WebhookEnvelope" + }, + { + "type": "object", + "properties": { + "data": { + "$ref": "#/components/schemas/CommentEventData" + } + } + } + ] + }, + "WorkOrderEventData": { + "type": "object", + "description": "Full current work-order state (not a diff); contract section 4.1. Same fields as WorkOrder minus source_email_s3_key.", + "required": ["work_order_id"], + "properties": { + "work_order_id": { + "type": "string", + "pattern": "^[0-9]+$" + }, + "wo_status": { + "type": ["string", "null"], + "enum": ["new", "assigned", "in_progress", "on_hold", "completed", "cancelled", "unknown", null] + }, + "description": { + "type": ["string", "null"] + }, + "customer": { + "type": ["string", "null"] + }, + "site_code": { + "type": ["string", "null"] + }, + "building": { + "type": ["string", "null"] + }, + "address": { + "type": ["string", "null"] + }, + "severity": { + "type": ["string", "null"] + }, + "priority": { + "type": ["string", "null"] + }, + "assigned_to": { + "type": ["string", "null"] + }, + "date_reported": { + "type": ["string", "null"] + }, + "scheduled_start": { + "type": ["string", "null"] + }, + "due_date": { + "type": ["string", "null"] + }, + "record_type": { + "type": ["string", "null"], + "enum": ["new_work_order", "update", "comment", "cancellation", null] + }, + "created_at": { + "type": ["string", "null"] + }, + "updated_at": { + "type": ["string", "null"] + }, + "write_origin": { + "type": ["string", "null"], + "description": "Forward-compat (phase 2): present on records written through the write-back API; the emitter skips those, so receivers should tolerate but never see it." + } + } + }, + "CommentEventData": { + "type": "object", + "description": "Contract section 4.2. Same fields as WorkOrderComment minus source_email_s3_key.", + "required": ["work_order_id", "comment_id"], + "properties": { + "work_order_id": { + "type": "string" + }, + "comment_id": { + "type": "string" + }, + "record_type": { + "type": ["string", "null"] + }, + "commenter": { + "type": ["string", "null"] + }, + "text": { + "type": ["string", "null"] + }, + "created_at": { + "type": ["string", "null"] + }, + "ingested_at": { + "type": ["string", "null"] + } + } + } + } + } +} diff --git a/lambdas/api/pagination.py b/lambdas/api/pagination.py new file mode 100644 index 0000000..affb200 --- /dev/null +++ b/lambdas/api/pagination.py @@ -0,0 +1,62 @@ +"""Opaque cursor pagination over DynamoDB ``LastEvaluatedKey``. + +The cursor is base64url(JSON(LastEvaluatedKey)). It is untrusted client input: +``decode_cursor`` validates shape strictly (a flat dict whose keys are exactly +a subset of the table's key attributes and whose values are non-empty strings) +and raises ``BadCursor`` -- mapped to HTTP 400 by the handler -- on anything +else, so a malformed or tampered cursor can never reach DynamoDB as an +arbitrary ``ExclusiveStartKey`` or surface as a 500. +""" + +import base64 +import binascii +import json + +DEFAULT_LIMIT = 100 +MAX_LIMIT = 500 +_MAX_CURSOR_CHARS = 2048 + + +class BadCursor(ValueError): + pass + + +def clamp_limit(raw) -> int: + if raw is None or raw == "": + return DEFAULT_LIMIT + try: + value = int(raw) + except (TypeError, ValueError): + raise BadCursor("limit must be an integer") from None + return max(1, min(MAX_LIMIT, value)) + + +def encode_cursor(last_evaluated_key: dict) -> str: + raw = json.dumps(last_evaluated_key, default=str, sort_keys=True) + return base64.urlsafe_b64encode(raw.encode()).decode() + + +def decode_cursor(cursor: str, key_attrs: frozenset) -> dict: + """Decode a cursor and require it to be EXACTLY the table's key attributes. + + ``key_attrs`` is the full key schema of the operation the cursor feeds + (e.g. ``{work_order_id, comment_id}`` for the comments Query). An exact + match -- not a subset -- is required: a partial composite key or a cursor + minted for a different endpoint would otherwise pass a subset check, reach + DynamoDB as an incomplete/inconsistent ``ExclusiveStartKey``, and raise a + ValidationException that surfaces as a 500 (and pages the 5xx alarm). Here + it fails closed as a 400 instead. Callers additionally pin the + partition-key value to the request path (see ``wo_repo.list_comments``). + """ + if len(cursor) > _MAX_CURSOR_CHARS: + raise BadCursor("cursor too long") + try: + decoded = json.loads(base64.urlsafe_b64decode(cursor.encode())) + except (binascii.Error, UnicodeDecodeError, json.JSONDecodeError, ValueError): + raise BadCursor("cursor is not valid") from None + if not isinstance(decoded, dict) or set(decoded) != key_attrs: + raise BadCursor("cursor is not valid") + for value in decoded.values(): + if not isinstance(value, str) or not value: + raise BadCursor("cursor is not valid") + return decoded diff --git a/lambdas/api/po_repo.py b/lambdas/api/po_repo.py new file mode 100644 index 0000000..9d3f486 --- /dev/null +++ b/lambdas/api/po_repo.py @@ -0,0 +1,53 @@ +"""Read access to the purchase-order tables for the procurement API. + +Same unordered cursor-paginated Scan pattern as wo_repo (no GSIs; listing +serves reconciliation, not ranked queries). VendorReplies is deliberately +absent: the table is dead (its writer was deleted) and must not gain new +consumers. +""" + +import os + +import boto3 + +from pagination import decode_cursor, encode_cursor + +dynamodb = boto3.resource("dynamodb") + +PO_TABLE = os.environ.get("PO_TABLE", "purchase-orders") +VERIFIED_SITES_TABLE = os.environ.get("VERIFIED_SITES_TABLE", "verified-sites") + +PO_KEY_ATTRS = frozenset({"po_number"}) +SITE_KEY_ATTRS = frozenset({"siteCode"}) + + +def _page(response) -> tuple[list, str | None]: + items = response.get("Items", []) + lek = response.get("LastEvaluatedKey") + return items, (encode_cursor(lek) if lek else None) + + +def list_purchase_orders(limit: int, cursor: str | None) -> tuple[list, str | None]: + table = dynamodb.Table(PO_TABLE) + kwargs = {"Limit": limit} + if cursor: + kwargs["ExclusiveStartKey"] = decode_cursor(cursor, PO_KEY_ATTRS) + return _page(table.scan(**kwargs)) + + +def get_purchase_order(po_number: str) -> dict | None: + table = dynamodb.Table(PO_TABLE) + return table.get_item(Key={"po_number": po_number}).get("Item") + + +def list_verified_sites(limit: int, cursor: str | None) -> tuple[list, str | None]: + table = dynamodb.Table(VERIFIED_SITES_TABLE) + kwargs = {"Limit": limit} + if cursor: + kwargs["ExclusiveStartKey"] = decode_cursor(cursor, SITE_KEY_ATTRS) + return _page(table.scan(**kwargs)) + + +def get_verified_site(site_code: str) -> dict | None: + table = dynamodb.Table(VERIFIED_SITES_TABLE) + return table.get_item(Key={"siteCode": site_code}).get("Item") diff --git a/lambdas/api/requirements.txt b/lambdas/api/requirements.txt new file mode 100644 index 0000000..10b0d4b --- /dev/null +++ b/lambdas/api/requirements.txt @@ -0,0 +1 @@ +# Dependabot anchor only. The procurement-api Lambda is stdlib+boto3 (provided by the runtime); nothing is pip-installed into the bundle. diff --git a/lambdas/api/router.py b/lambdas/api/router.py new file mode 100644 index 0000000..8bf0539 --- /dev/null +++ b/lambdas/api/router.py @@ -0,0 +1,33 @@ +"""Route tables for the procurement API. + +Single source of truth for what the API implements: the spec-drift test +(tests/test_api_spec_drift.py) asserts these tables match lambdas/api/ +openapi.json exactly, so an endpoint can't be added, removed, or renamed on +one side without failing CI. Keys are (httpMethod, resource) as API Gateway's +Lambda-proxy event presents them (resource = the templated path). + +DATA_ROUTES values are handler-module attribute names resolved by handler.py. +PLANNED_ROUTES are phase-2 write endpoints: present in the spec (x-planned) +and at the gateway, but the handler answers 501 until they ship with their +own IAM diff + cross-family review. +""" + +DATA_ROUTES = { + ("GET", "/work-orders"): "list_work_orders", + ("GET", "/work-orders/{workOrderId}"): "get_work_order", + ("GET", "/work-orders/{workOrderId}/comments"): "list_comments", + ("GET", "/purchase-orders"): "list_purchase_orders", + ("GET", "/purchase-orders/{poNumber}"): "get_purchase_order", + ("GET", "/verified-sites"): "list_verified_sites", + ("GET", "/verified-sites/{siteCode}"): "get_verified_site", +} + +DOCS_ROUTES = { + ("GET", "/docs"), + ("GET", "/openapi.json"), +} + +PLANNED_ROUTES = { + ("POST", "/work-orders/{workOrderId}/comments"), + ("PATCH", "/work-orders/{workOrderId}"), +} diff --git a/lambdas/api/serialization.py b/lambdas/api/serialization.py new file mode 100644 index 0000000..a9271c2 --- /dev/null +++ b/lambdas/api/serialization.py @@ -0,0 +1,43 @@ +"""JSON response helpers for the procurement API. + +DynamoDB items come back through boto3 with numeric attributes as +``decimal.Decimal`` (e.g. purchase-orders ``total_amount``/``line_items`` +prices), which ``json.dumps`` rejects. The encoder below renders integral +Decimals as JSON integers (exact) and non-integral Decimals as floats. For the +values this API serves -- currency amounts at PO magnitudes, well under ~15 +significant digits -- the shortest-float repr round-trips the original digits +exactly. A Decimal carrying more precision than a float can hold (DynamoDB +Number supports 38 digits) would be truncated silently; the pipeline never +writes such values, but if that ever changes, switch non-integral Decimals to +``str(o)`` and update the OpenAPI number types. +""" + +import decimal +import json + + +class _DecimalSafeEncoder(json.JSONEncoder): + def default(self, o): + if isinstance(o, decimal.Decimal): + if o == o.to_integral_value(): + return int(o) + return float(o) + if isinstance(o, set): + return sorted(o) + return super().default(o) + + +def dumps(payload) -> str: + return json.dumps(payload, cls=_DecimalSafeEncoder) + + +def json_response(status_code: int, payload) -> dict: + return { + "statusCode": status_code, + "headers": {"Content-Type": "application/json"}, + "body": dumps(payload), + } + + +def error_response(status_code: int, message: str) -> dict: + return json_response(status_code, {"error": message}) diff --git a/lambdas/api/wo_repo.py b/lambdas/api/wo_repo.py new file mode 100644 index 0000000..3521468 --- /dev/null +++ b/lambdas/api/wo_repo.py @@ -0,0 +1,62 @@ +"""Read access to the work-order tables for the procurement API. + +Listing is a cursor-paginated Scan -- deliberately: the WO tables carry no +GSIs (removed 2026-06-03 for zero reads), and the API's list consumers +(SHOC reconciliation/backfill) walk the full table anyway, which is exactly +the access pattern SHOC's retired SyncController used. Listings are therefore +UNORDERED across pages; per-work-order comment listing is a cheap Query on +the partition key. +""" + +import os + +import boto3 +from boto3.dynamodb.conditions import Key + +from pagination import BadCursor, decode_cursor, encode_cursor + +dynamodb = boto3.resource("dynamodb") + +WORK_ORDERS_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders") +COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments") + +WO_KEY_ATTRS = frozenset({"work_order_id"}) +COMMENT_KEY_ATTRS = frozenset({"work_order_id", "comment_id"}) + + +def _page(response) -> tuple[list, str | None]: + items = response.get("Items", []) + lek = response.get("LastEvaluatedKey") + return items, (encode_cursor(lek) if lek else None) + + +def list_work_orders(limit: int, cursor: str | None) -> tuple[list, str | None]: + table = dynamodb.Table(WORK_ORDERS_TABLE) + kwargs = {"Limit": limit} + if cursor: + kwargs["ExclusiveStartKey"] = decode_cursor(cursor, WO_KEY_ATTRS) + return _page(table.scan(**kwargs)) + + +def get_work_order(work_order_id: str) -> dict | None: + table = dynamodb.Table(WORK_ORDERS_TABLE) + return table.get_item(Key={"work_order_id": work_order_id}).get("Item") + + +def list_comments( + work_order_id: str, limit: int, cursor: str | None +) -> tuple[list, str | None]: + table = dynamodb.Table(COMMENTS_TABLE) + kwargs = { + "KeyConditionExpression": Key("work_order_id").eq(work_order_id), + "Limit": limit, + } + if cursor: + start_key = decode_cursor(cursor, COMMENT_KEY_ATTRS) + # Pin the cursor's partition to the path entity: a cursor minted for + # WO A must never resume WO B's Query (DynamoDB would reject the + # partition mismatch as a 500; reject it as a 400 here instead). + if start_key["work_order_id"] != work_order_id: + raise BadCursor("cursor is not valid") + kwargs["ExclusiveStartKey"] = start_key + return _page(table.query(**kwargs)) diff --git a/lambdas/shared/web_ui_auth.py b/lambdas/shared/web_ui_auth.py index 4b2d245..0359298 100644 --- a/lambdas/shared/web_ui_auth.py +++ b/lambdas/shared/web_ui_auth.py @@ -74,4 +74,8 @@ def is_authenticated(event: dict) -> bool: presented = auth[7:].strip() if not presented: return False - return hmac.compare_digest(presented, token) + # Compare as bytes: hmac.compare_digest raises TypeError on non-ASCII str + # operands, which a crafted token (?token=%C3%A9 or a non-ASCII header) + # would otherwise turn into an uncaught 500. Bytes always compare in + # constant time, so a non-matching token fails closed (401) instead. + return hmac.compare_digest(presented.encode("utf-8"), token.encode("utf-8")) diff --git a/pytest.ini b/pytest.ini index 9511b4a..1bacae8 100644 --- a/pytest.ini +++ b/pytest.ini @@ -12,6 +12,7 @@ addopts = --cov=lambdas/po/web_ui --cov=lambdas/wo/web_ui --cov=lambdas/po/site_extractor + --cov=lambdas/api --cov=lambdas/shared --cov-report=term-missing --cov-fail-under=80 diff --git a/scripts/post-deploy-smoke.sh b/scripts/post-deploy-smoke.sh index 9daffc2..5988a19 100755 --- a/scripts/post-deploy-smoke.sh +++ b/scripts/post-deploy-smoke.sh @@ -14,7 +14,7 @@ set -euo pipefail REGION="us-east-1" -FUNCTION_NAMES=("po-email-processor" "workorder-email-processor") +FUNCTION_NAMES=("po-email-processor" "workorder-email-processor" "procurement-api") HEALTHCHECK_PAYLOAD='{"healthcheck": true}' work_dir="$(mktemp -d)" diff --git a/tests/support/loader.py b/tests/support/loader.py index a826b78..ad959aa 100644 --- a/tests/support/loader.py +++ b/tests/support/loader.py @@ -50,6 +50,15 @@ _SIBLING_MODULES = ( "email_parsing", "emf", "web_ui_auth", + # procurement-api siblings (lambdas/api/): pagination/serialization/router + # have no sibling deps; wo_repo/po_repo import pagination, so they follow + # it. These names exist only under lambdas/api/, so the po/wo handler + # loads skip them via the exists() check. + "pagination", + "serialization", + "router", + "wo_repo", + "po_repo", "prompts", "template_parser", "derived_fields", diff --git a/tests/test_api_handlers.py b/tests/test_api_handlers.py new file mode 100644 index 0000000..ceb2280 --- /dev/null +++ b/tests/test_api_handlers.py @@ -0,0 +1,271 @@ +"""procurement-api handler tests (offline, fake Dynamo). + +Pins the auth seam both ways: docs routes fail closed on the shared token +(including the browser ?token= shim), while data routes must NOT consult the +token at all -- SigV4 authorization is API Gateway's job upstream, and a +handler-side token check there would break every legitimate SigV4 caller. +Also pins routing, 404/501/400 mapping, Decimal-safe serialization, and limit +clamping. +""" + +import json +from decimal import Decimal + +import pytest + +from tests.support import load_lambda_module + + +class _FakeTable: + def __init__(self, items): + self._items = items + self.last_kwargs = None + + def scan(self, **kwargs): + self.last_kwargs = kwargs + return {"Items": list(self._items)} + + def query(self, **kwargs): + self.last_kwargs = kwargs + return {"Items": list(self._items)} + + def get_item(self, Key): # noqa: N803 (boto3 kwarg name) + pk_attr, pk_val = next(iter(Key.items())) + for item in self._items: + if item.get(pk_attr) == pk_val: + return {"Item": item} + return {} + + +class _FakeDynamo: + def __init__(self, tables): + self.tables = tables + + def Table(self, name): # noqa: N802 (boto3 method name) + return self.tables[name] + + +def _event(method, resource, path_params=None, qs=None, headers=None): + return { + "httpMethod": method, + "resource": resource, + "pathParameters": path_params or {}, + "queryStringParameters": qs or {}, + "headers": headers or {}, + } + + +@pytest.fixture() +def api(monkeypatch): + mod = load_lambda_module("api", "handler") + wo_items = [ + { + "work_order_id": "11144580730", + "wo_status": "assigned", + "description": "Dock door 14", + } + ] + comment_items = [ + { + "work_order_id": "11144580730", + "comment_id": "11144580730#2026-01-01T00:00:00#abc123def456", + "text": "Vendor dispatched", + } + ] + po_items = [ + { + "po_number": "2D-22030794", + "total_amount": Decimal("123.45"), + "line_items": [{"quantity": Decimal("5"), "price": Decimal("24.69")}], + } + ] + site_items = [{"siteCode": "JFK8", "state": "NY", "poCount": Decimal("7")}] + + tables = { + mod.wo_repo.WORK_ORDERS_TABLE: _FakeTable(wo_items), + mod.wo_repo.COMMENTS_TABLE: _FakeTable(comment_items), + mod.po_repo.PO_TABLE: _FakeTable(po_items), + mod.po_repo.VERIFIED_SITES_TABLE: _FakeTable(site_items), + } + fake = _FakeDynamo(tables) + monkeypatch.setattr(mod.wo_repo, "dynamodb", fake) + monkeypatch.setattr(mod.po_repo, "dynamodb", fake) + return mod, fake + + +def test_healthcheck_short_circuits(api): + mod, _ = api + assert mod.handler({"healthcheck": True}, None) == {"healthcheck": "ok"} + + +def test_docs_fail_closed_without_token(api): + # No WEB_UI_AUTH_TOKEN_SECRET_ARN configured -> web_ui_auth returns None + # -> the docs routes must 401, not serve the spec. + mod, _ = api + for resource in ("/docs", "/openapi.json"): + response = mod.handler(_event("GET", resource), None) + assert response["statusCode"] == 401 + + +def test_docs_served_when_authenticated(api, monkeypatch): + mod, _ = api + monkeypatch.setattr(mod, "is_authenticated", lambda event: True) + docs = mod.handler(_event("GET", "/docs"), None) + assert docs["statusCode"] == 200 + assert docs["headers"]["Content-Type"].startswith("text/html") + assert "Procurement Ingest API" in docs["body"] + assert mod._SPEC_PLACEHOLDER not in docs["body"] + + spec = mod.handler(_event("GET", "/openapi.json"), None) + assert spec["statusCode"] == 200 + assert json.loads(spec["body"])["openapi"] == "3.1.0" + + +def test_token_query_shim_synthesizes_header(api, monkeypatch): + mod, _ = api + seen = {} + + def capture(event): + seen.update(event.get("headers") or {}) + return False + + monkeypatch.setattr(mod, "is_authenticated", capture) + original = _event("GET", "/docs", qs={"token": "sekrit"}) + mod.handler(original, None) + assert seen.get("x-auth-token") == "sekrit" + # The shim must not mutate the caller's event in place. + assert "x-auth-token" not in original["headers"] + + +def test_non_ascii_token_fails_closed_not_500(api): + # A non-ASCII presented token must 401 (fail closed), not crash + # hmac.compare_digest into a 500 that pages the 5xx alarm. Reach the + # web_ui_auth module namespace through the imported is_authenticated's + # globals and seed a valid cached token. + mod, _ = api + auth_globals = mod.is_authenticated.__globals__ + auth_globals["_auth_token_cache"] = "correct-token" + auth_globals["_auth_token_cached_at"] = float("inf") + try: + response = mod.handler( + _event("GET", "/docs", headers={"x-auth-token": "é"}), None + ) + assert response["statusCode"] == 401 + finally: + auth_globals["_auth_token_cache"] = None + auth_globals["_auth_token_cached_at"] = 0.0 + + +def test_data_routes_never_consult_docs_token(api, monkeypatch): + mod, _ = api + + def explode(event): + raise AssertionError("data routes must not consult the docs token") + + monkeypatch.setattr(mod, "is_authenticated", explode) + response = mod.handler(_event("GET", "/work-orders"), None) + assert response["statusCode"] == 200 + + +def test_list_and_get_routing(api): + mod, _ = api + listing = mod.handler(_event("GET", "/work-orders"), None) + body = json.loads(listing["body"]) + assert body["items"][0]["work_order_id"] == "11144580730" + assert body["next_cursor"] is None + + hit = mod.handler( + _event( + "GET", + "/work-orders/{workOrderId}", + path_params={"workOrderId": "11144580730"}, + ), + None, + ) + assert json.loads(hit["body"])["wo_status"] == "assigned" + + miss = mod.handler( + _event( + "GET", + "/work-orders/{workOrderId}", + path_params={"workOrderId": "999"}, + ), + None, + ) + assert miss["statusCode"] == 404 + + comments = mod.handler( + _event( + "GET", + "/work-orders/{workOrderId}/comments", + path_params={"workOrderId": "11144580730"}, + ), + None, + ) + assert json.loads(comments["body"])["items"][0]["text"] == "Vendor dispatched" + + sites = mod.handler( + _event( + "GET", + "/verified-sites/{siteCode}", + path_params={"siteCode": "JFK8"}, + ), + None, + ) + assert json.loads(sites["body"])["state"] == "NY" + + +def test_decimal_serialization_round_trips(api): + mod, _ = api + response = mod.handler( + _event( + "GET", + "/purchase-orders/{poNumber}", + path_params={"poNumber": "2D-22030794"}, + ), + None, + ) + body = json.loads(response["body"]) + assert body["total_amount"] == 123.45 + assert body["line_items"][0]["quantity"] == 5 + assert isinstance(body["line_items"][0]["quantity"], int) + assert "Decimal" not in response["body"] + + +def test_limit_clamped_to_bounds(api): + mod, fake = api + mod.handler(_event("GET", "/work-orders", qs={"limit": "9999"}), None) + assert fake.tables[mod.wo_repo.WORK_ORDERS_TABLE].last_kwargs["Limit"] == 500 + mod.handler(_event("GET", "/work-orders", qs={"limit": "0"}), None) + assert fake.tables[mod.wo_repo.WORK_ORDERS_TABLE].last_kwargs["Limit"] == 1 + + +def test_non_integer_limit_is_400(api): + mod, _ = api + response = mod.handler(_event("GET", "/work-orders", qs={"limit": "abc"}), None) + assert response["statusCode"] == 400 + + +def test_planned_routes_return_501(api): + mod, _ = api + for method, resource in mod.PLANNED_ROUTES: + response = mod.handler(_event(method, resource), None) + assert response["statusCode"] == 501, (method, resource) + + +def test_unknown_route_is_404(api): + mod, _ = api + response = mod.handler(_event("GET", "/nope"), None) + assert response["statusCode"] == 404 + + +def test_unexpected_error_maps_to_clean_500(api, monkeypatch): + mod, fake = api + + def boom(**kwargs): + raise RuntimeError("dynamo fell over") + + monkeypatch.setattr(fake.tables[mod.wo_repo.WORK_ORDERS_TABLE], "scan", boom) + response = mod.handler(_event("GET", "/work-orders"), None) + assert response["statusCode"] == 500 + assert json.loads(response["body"]) == {"error": "internal error"} diff --git a/tests/test_api_pagination_moto.py b/tests/test_api_pagination_moto.py new file mode 100644 index 0000000..10584ec --- /dev/null +++ b/tests/test_api_pagination_moto.py @@ -0,0 +1,232 @@ +"""Cursor pagination against real DynamoDB semantics (moto). + +The offline handler tests fake Dynamo, which can't exercise real +LastEvaluatedKey behavior. Here moto-backed tables with the production key +schemas prove: pages are disjoint and complete, the final page's next_cursor +is null, comment Query pagination respects the composite key, and hostile +cursors (garbage, wrong key attrs, huge) map to 400 -- never 500, never a raw +ExclusiveStartKey error. + +The repo modules build their boto3 resource at import time; a resource created +before mock_aws starts is NOT intercepted, so each test rebinds the modules' +``dynamodb`` to a fresh resource inside the mock (same trap as test_po_merge's +moto-before-handler ordering). +""" + +import base64 +import json + +import boto3 +import pytest +from moto import mock_aws + +from tests.support import load_lambda_module + + +def _event(method, resource, path_params=None, qs=None): + return { + "httpMethod": method, + "resource": resource, + "pathParameters": path_params or {}, + "queryStringParameters": qs or {}, + "headers": {}, + } + + +@pytest.fixture() +def api(monkeypatch): + with mock_aws(): + mod = load_lambda_module("api", "handler") + resource = boto3.resource("dynamodb", region_name="us-east-1") + monkeypatch.setattr(mod.wo_repo, "dynamodb", resource) + monkeypatch.setattr(mod.po_repo, "dynamodb", resource) + + resource.create_table( + TableName=mod.wo_repo.WORK_ORDERS_TABLE, + KeySchema=[{"AttributeName": "work_order_id", "KeyType": "HASH"}], + AttributeDefinitions=[ + {"AttributeName": "work_order_id", "AttributeType": "S"} + ], + BillingMode="PAY_PER_REQUEST", + ) + resource.create_table( + TableName=mod.wo_repo.COMMENTS_TABLE, + KeySchema=[ + {"AttributeName": "work_order_id", "KeyType": "HASH"}, + {"AttributeName": "comment_id", "KeyType": "RANGE"}, + ], + AttributeDefinitions=[ + {"AttributeName": "work_order_id", "AttributeType": "S"}, + {"AttributeName": "comment_id", "AttributeType": "S"}, + ], + BillingMode="PAY_PER_REQUEST", + ) + resource.create_table( + TableName=mod.po_repo.PO_TABLE, + KeySchema=[{"AttributeName": "po_number", "KeyType": "HASH"}], + AttributeDefinitions=[{"AttributeName": "po_number", "AttributeType": "S"}], + BillingMode="PAY_PER_REQUEST", + ) + resource.create_table( + TableName=mod.po_repo.VERIFIED_SITES_TABLE, + KeySchema=[{"AttributeName": "siteCode", "KeyType": "HASH"}], + AttributeDefinitions=[{"AttributeName": "siteCode", "AttributeType": "S"}], + BillingMode="PAY_PER_REQUEST", + ) + yield mod, resource + + +def _walk_pages(mod, resource_path, qs_extra=None, path_params=None, limit=5): + seen, pages = [], 0 + cursor = None + while True: + qs = {"limit": str(limit), **(qs_extra or {})} + if cursor: + qs["cursor"] = cursor + response = mod.handler( + _event("GET", resource_path, path_params=path_params, qs=qs), None + ) + assert response["statusCode"] == 200 + body = json.loads(response["body"]) + seen.extend(body["items"]) + pages += 1 + cursor = body["next_cursor"] + assert pages < 50, "pagination failed to terminate" + if cursor is None: + return seen, pages + + +def test_work_order_scan_pages_are_disjoint_and_complete(api): + mod, resource = api + table = resource.Table(mod.wo_repo.WORK_ORDERS_TABLE) + for i in range(12): + table.put_item(Item={"work_order_id": str(10000 + i), "wo_status": "new"}) + + items, pages = _walk_pages(mod, "/work-orders", limit=5) + ids = [item["work_order_id"] for item in items] + assert len(ids) == 12 + assert len(set(ids)) == 12 + assert pages >= 3 + + +def test_comment_query_pagination_composite_key(api): + mod, resource = api + table = resource.Table(mod.wo_repo.COMMENTS_TABLE) + for i in range(7): + table.put_item( + Item={ + "work_order_id": "555", + "comment_id": f"555#2026-01-0{i + 1}T00:00:00#{i:012d}", + "text": f"comment {i}", + } + ) + # A second work order's comments must never bleed into the page. + table.put_item( + Item={"work_order_id": "666", "comment_id": "666#x#0", "text": "other"} + ) + + items, _ = _walk_pages( + mod, + "/work-orders/{workOrderId}/comments", + path_params={"workOrderId": "555"}, + limit=3, + ) + assert len(items) == 7 + assert {item["work_order_id"] for item in items} == {"555"} + + +def test_purchase_order_and_site_pagination(api): + mod, resource = api + po_table = resource.Table(mod.po_repo.PO_TABLE) + for i in range(6): + po_table.put_item(Item={"po_number": f"2D-{i:08d}"}) + sites_table = resource.Table(mod.po_repo.VERIFIED_SITES_TABLE) + for code in ("JFK8", "WNY2", "UVA5"): + sites_table.put_item(Item={"siteCode": code}) + + po_items, _ = _walk_pages(mod, "/purchase-orders", limit=4) + assert len(po_items) == 6 + site_items, _ = _walk_pages(mod, "/verified-sites", limit=2) + assert {item["siteCode"] for item in site_items} == {"JFK8", "WNY2", "UVA5"} + + +@pytest.mark.parametrize( + "cursor", + [ + "not-base64!!!", + "aGVsbG8=", # base64("hello") -- not JSON-dict + "e30=", # base64("{}") -- empty dict + # base64 of {"evil_attr": "x"} -- key attr not allowed for this table + "eyJldmlsX2F0dHIiOiAieCJ9", + "A" * 3000, # oversized + ], +) +def test_hostile_cursor_is_400_not_500(api, cursor): + mod, resource = api + resource.Table(mod.wo_repo.WORK_ORDERS_TABLE).put_item(Item={"work_order_id": "1"}) + response = mod.handler(_event("GET", "/work-orders", qs={"cursor": cursor}), None) + assert response["statusCode"] == 400 + assert "error" in json.loads(response["body"]) + + +def _cursor(payload): + return base64.urlsafe_b64encode(json.dumps(payload).encode()).decode() + + +def test_partial_composite_cursor_is_400_not_500(api): + # A work-orders-shaped cursor {work_order_id} is a partial key for the + # comments Query (needs work_order_id+comment_id). Exact-key match rejects + # it as 400 rather than letting DynamoDB ValidationException -> 500. + mod, resource = api + resource.Table(mod.wo_repo.COMMENTS_TABLE).put_item( + Item={"work_order_id": "555", "comment_id": "555#x#0"} + ) + response = mod.handler( + _event( + "GET", + "/work-orders/{workOrderId}/comments", + path_params={"workOrderId": "555"}, + qs={"cursor": _cursor({"work_order_id": "555"})}, + ), + None, + ) + assert response["statusCode"] == 400 + + +def test_cross_work_order_comment_cursor_is_400(api): + # A full comments cursor minted for WO A must not resume WO B's Query. + mod, resource = api + response = mod.handler( + _event( + "GET", + "/work-orders/{workOrderId}/comments", + path_params={"workOrderId": "B"}, + qs={"cursor": _cursor({"work_order_id": "A", "comment_id": "A#x#0"})}, + ), + None, + ) + assert response["statusCode"] == 400 + + +def test_dynamodb_validation_error_maps_to_400(api, monkeypatch): + # Defense in depth: even if some crafted-but-valid cursor reached DynamoDB + # and raised ValidationException, the handler maps it to 400, not a 500 + # that would page the zero-threshold 5xx alarm. + from botocore.exceptions import ClientError + + mod, _ = api + + class _RaisingTable: + def scan(self, **kwargs): + raise ClientError( + {"Error": {"Code": "ValidationException", "Message": "bad key"}}, + "Scan", + ) + + monkeypatch.setattr( + mod.wo_repo, + "dynamodb", + type("D", (), {"Table": lambda self, n: _RaisingTable()})(), + ) + response = mod.handler(_event("GET", "/work-orders"), None) + assert response["statusCode"] == 400 diff --git a/tests/test_api_spec_drift.py b/tests/test_api_spec_drift.py new file mode 100644 index 0000000..531b32f --- /dev/null +++ b/tests/test_api_spec_drift.py @@ -0,0 +1,114 @@ +"""Spec <-> implementation drift gate. + +lambdas/api/openapi.json is the published contract; lambdas/api/router.py is +what the Lambda actually serves. This test makes them the SAME set: an +endpoint added/removed/renamed on one side without the other fails CI here, +so the docs page can never silently lie. Also pins the phase-2 x-planned +markers, the four outbound webhook events, and the enum values the spec +promises against the WO extraction contract (prompts.py). +""" + +import json +from pathlib import Path + +from tests.support import REPO_ROOT, load_lambda_module + +_SPEC_PATH = Path(REPO_ROOT) / "lambdas" / "api" / "openapi.json" + +_HTTP_METHODS = {"get", "put", "post", "patch", "delete", "head", "options"} + + +def _spec(): + return json.loads(_SPEC_PATH.read_text(encoding="utf-8")) + + +def _spec_routes(spec): + implemented, planned = set(), set() + for path, ops in spec["paths"].items(): + for method, op in ops.items(): + if method not in _HTTP_METHODS: + continue + key = (method.upper(), path) + if op.get("x-planned"): + planned.add(key) + else: + implemented.add(key) + return implemented, planned + + +def test_spec_is_openapi_31(): + assert _spec()["openapi"] == "3.1.0" + + +def test_implemented_routes_match_router_exactly(): + router = load_lambda_module("api", "router") + implemented, planned = _spec_routes(_spec()) + assert implemented == set(router.DATA_ROUTES) | set(router.DOCS_ROUTES) + assert planned == set(router.PLANNED_ROUTES) + + +def test_every_data_route_has_a_repo_function(): + router = load_lambda_module("api", "router") + handler = load_lambda_module("api", "handler") + assert set(router.DATA_ROUTES.values()) == set(handler._REPO_FUNCS) + + +def test_webhooks_section_documents_all_four_events(): + spec = _spec() + assert set(spec["webhooks"]) == { + "work_order.created", + "work_order.updated", + "work_order.cancelled", + "work_order.comment_added", + } + for ops in spec["webhooks"].values(): + assert "post" in ops + + +def test_spec_enums_match_wo_extraction_contract(): + # The WO pipeline's prompts.py is the enum source of truth (the webhook + # contract pins it too). If the pipeline ever widens wo_status or + # record_type, the published spec must move in the same PR. + prompts = load_lambda_module("wo", "email_processor/prompts") + prompt_text = prompts.EXTRACTION_PROMPT + spec = _spec() + schemas = spec["components"]["schemas"] + + wo_status_enum = { + value + for value in schemas["WorkOrder"]["properties"]["wo_status"]["enum"] + if value is not None + } + record_type_enum = { + value + for value in schemas["WorkOrder"]["properties"]["record_type"]["enum"] + if value is not None + } + for value in wo_status_enum - {"unknown"}: + assert value in prompt_text, f"wo_status {value!r} not in extraction prompt" + for value in record_type_enum: + assert value in prompt_text, f"record_type {value!r} not in extraction prompt" + + event_data = schemas["WorkOrderEventData"]["properties"] + assert set(event_data["wo_status"]["enum"]) == set( + schemas["WorkOrder"]["properties"]["wo_status"]["enum"] + ) + + +def test_spec_has_no_script_breakout_sequence(): + # The spec is inlined into a