procurement-ingest/lambdas/wo/shoc_emitter/handler.py
Adam Moussa 7569bd8250
harden(webhook): resolve /sh-security-review findings (1 confirmed medium + cheap fixes)
High-recall detector fan-out (injection/authz/secrets-crypto/iac-iam/logic)
+ proof-or-kill verifier. Gate PASSES: 1 confirmed medium, 0 confirmed
critical/high. Confirmed finding fixed; several unverified-but-cheap
hardenings applied since the emitter ships dark and activation is weeks out.

- CONFIRMED medium (confused deputy): the rotation Lambda's generated
  invoke permission for secretsmanager.amazonaws.com carried no
  SourceAccount/SourceArn, so any account's Secrets Manager could invoke
  the rotator. Patched the generated CfnPermission in place (a second
  permission would be additive, not restrictive) to pin account + this
  secret ARN.
- delivery + replay: refuse to follow receiver 3xx redirects (no-redirect
  opener) so live X-SH-* auth headers can't be forwarded to a
  receiver-chosen Location and an http:// Location can't slip past the
  https guard. Fixed the "unfollowed 3xx" comment that was factually wrong.
- delivery: classify 401/403 as retryable (invalidate key cache + retry in
  order) instead of parking -- transient auth failures (rotation outran the
  TTL cache, clock skew) are availability events, not contract bugs.
- envelope: build_event now genuinely total (guarded eventID /
  ApproximateCreationDateTime subscripts) per its own never-raise contract.
- handler: catch-all so an unexpected per-record error (e.g. SQS park
  failure) reports only that record instead of failing the whole batch
  (which would re-deliver every earlier success for 24h); per-invocation
  emit/skip batch summary so a systemic silent drop is queryable/alarmable.
- rotator: narrow the AWSCURRENT-read except to ResourceNotFound/JSONDecode
  (transient SM/KMS errors re-raise so the overlap key isn't silently
  dropped); kid uniqueness checked against ALL retained kids with a random
  suffix on collision (never reissue a kid for a different secret).
- contract: skeleton-upsert required on ANY unknown work_order_id (not just
  comment-before-create) + monotonicity guard (ignore older updated_at), so
  a parked created or an out-of-order replay can't corrupt receiver state.

Unverified/refuted findings left as-is with rationale: the two "high" logic
claims (whole-batch crash triggers, ordering violation) were refuted on
reachability (real stream records carry required fields; persistence writes
strings only; full-state idempotent upsert absorbs the ordering gap). Signed
kid/version binding (AUTHZ-002) declined: coordinated contract change, not
cheap, no exploit with one algorithm/key.
2026-07-24 15:23:24 -04:00

162 lines
6.1 KiB
Python

"""
SHOC work-order webhook emitter Lambda.
Consumes the WorkOrders/WorkOrderComments DynamoDB streams and pushes each
mutation to SHOC as an HMAC-signed HTTPS POST (docs/shoc-webhook-contract.md).
This handler is the thin event loop; the work lives in flat sibling modules
(bare-name imports resolve via the same flat-landing bundling as the email
processor's siblings):
envelope.py -- stream record -> envelope mapping + event classification
delivery.py -- secret cache, signing, POST, response classification
Ships DARK: both event-source mappings deploy with enabled=False, so this code
runs zero deliveries until the activation PR flips them on after SHOC's
receiver passes the shared HMAC test vectors.
Ordering semantics: the ESMs run parallelization_factor=1 with bisect-on-error
off, and this loop processes records strictly in order. A retryable failure
(429/5xx/timeout/connection error) stops the batch immediately and reports
that record via report_batch_item_failures -- earlier successes are not
re-delivered, and the ESM blocks the shard and retries from the failed record,
preserving per-work-order commit order. Non-retryable 4xx responses are a
contract bug, not an availability blip: the full envelope is parked on the
rejected queue (alarmed) and the loop continues, so a bad payload can never
block the shard for 24 hours.
No healthcheck branch on purpose: stream consumers are not smoke-gated
(po-ingest-site-extractor precedent) -- there is no direct-invoke path to
probe, and a synthetic stream record would be a real delivery.
"""
import json
import logging
import os
import time
import boto3
import delivery
import envelope
logger = logging.getLogger()
logger.setLevel(logging.INFO)
REJECTED_QUEUE_URL = os.environ.get("REJECTED_QUEUE_URL")
# Lazy cached SQS client. Keeps the public attribute name ``sqs`` so the test
# monkeypatch target changes module only, not attribute name.
sqs = None
def _get_sqs():
global sqs
if sqs is None:
sqs = boto3.client("sqs")
return sqs
def _delivery_log(event: dict, status_code, latency_ms: int, outcome: str) -> str:
"""One structured record per delivery attempt (Logs-Insights-queryable)."""
return json.dumps(
{
"event": "shoc_delivery",
"delivery_id": event["delivery_id"],
"event_type": event["event_type"],
"work_order_id": event["data"].get("work_order_id"),
"status_code": status_code,
"latency_ms": latency_ms,
"outcome": outcome,
}
)
def _park_rejected(event: dict, status_code: int):
"""Send the full envelope to the rejected queue for operator replay."""
_get_sqs().send_message(
QueueUrl=REJECTED_QUEUE_URL,
MessageBody=json.dumps({"envelope": event, "response_status": status_code}),
)
logger.warning(
json.dumps(
{
"event": "shoc_delivery_rejected",
"delivery_id": event["delivery_id"],
"event_type": event["event_type"],
"work_order_id": event["data"].get("work_order_id"),
"response_status": status_code,
}
)
)
def _batch_summary(records: int, emitted: int, skipped: int) -> str:
"""Per-invocation emit/skip tally.
Skips (REMOVE, echo guard, comment MODIFYs, unmappable records) return no
delivery and would otherwise be invisible: a systemic classification break
-- e.g. a table rename desyncing the eventSourceARN parse -- would drop
100%% of events while every batch still reports success. Logging the ratio
makes that queryable in Logs Insights and alarmable.
"""
return json.dumps(
{
"event": "shoc_batch_summary",
"records": records,
"emitted": emitted,
"skipped": skipped,
}
)
def handler(event, context):
"""Lambda entry point. Triggered by the two WO-table stream ESMs."""
records = event.get("Records", [])
emitted = 0
skipped = 0
for record in records:
webhook_event = envelope.build_event(record)
if webhook_event is None:
skipped += 1
continue
emitted += 1
start = time.monotonic()
try:
outcome, status_code = delivery.deliver(webhook_event)
latency_ms = int((time.monotonic() - start) * 1000)
logger.info(_delivery_log(webhook_event, status_code, latency_ms, outcome))
if outcome == "rejected":
_park_rejected(webhook_event, status_code)
except delivery.RetryableDeliveryError as exc:
latency_ms = int((time.monotonic() - start) * 1000)
logger.warning(
_delivery_log(webhook_event, exc.status_code, latency_ms, "retryable")
)
logger.info(_batch_summary(len(records), emitted, skipped))
# Stop here: reporting this record's sequence number makes the ESM
# retry from it in order; earlier successes are not re-delivered.
return {
"batchItemFailures": [
{"itemIdentifier": record["dynamodb"]["SequenceNumber"]}
]
}
except Exception:
# Catch-all so an unexpected error (e.g. an SQS park failure)
# cannot escape and fail the WHOLE invocation -- that would make
# the ESM re-deliver every earlier success in the batch for up to
# 24h. Report only THIS record so the ESM retries from it in order.
logger.exception(
json.dumps(
{
"event": "shoc_delivery_error",
"delivery_id": webhook_event["delivery_id"],
"event_type": webhook_event["event_type"],
}
)
)
logger.info(_batch_summary(len(records), emitted, skipped))
return {
"batchItemFailures": [
{"itemIdentifier": record["dynamodb"]["SequenceNumber"]}
]
}
logger.info(_batch_summary(len(records), emitted, skipped))
return {"batchItemFailures": []}