procurement-ingest/lambdas/wo/shoc_emitter/delivery.py
Adam Moussa eb56d39b93
feat(webhook): SHOC WO webhook emitter — dark-ship streams, HMAC secret + rotation
Implements docs/shoc-webhook-plan.md Phases 1-5 (PR-2 of the SHOC
call-and-be-called effort). Everything ships DARK: both DynamoDB event
source mappings deploy enabled=False; activation is a deliberate
one-line follow-up PR gated on the SHOC receiver passing the shared
HMAC test vectors.

- Streams: NEW_AND_OLD_IMAGES on WorkOrders + WorkOrderComments
  (in-place update, RETAIN + logical IDs untouched; no existing
  consumers — verified live, neither table had a stream).
- workorder-shoc-emitter (Py3.12/ARM64): stream -> envelope ->
  HMAC-signed POST per docs/shoc-webhook-contract.md; strict per-shard
  ordering (parallelization 1, bisect off, retry until 24h age,
  ReportBatchItemFailures); 429/5xx/timeout block the shard in order,
  other 4xx park to workorder-shoc-emitter-rejected; ESM failures ->
  workorder-shoc-emitter-failures (metadata; replay rebuilds from
  DynamoDB). Echo guard skips write_origin=shoc-write-api.
- Secret workorder-ingest/shoc-webhook-hmac on a dedicated CMK
  (alias workorder-ingest-shoc-webhook-kms); cross-account
  GetSecretValue/DescribeSecret + kms:Decrypt granted to exactly
  arn:aws:iam::396287094661:role/shoc-backend-dev. RemovalPolicy
  DESTROY deliberately (machine-generated material; avoids the
  fixed-name RETAIN-orphan deadlock).
- workorder-shoc-hmac-rotator: 30-day rotation, dual-key overlap,
  64-hex keys, kid = UTC %Y-%m-%dT%H.
- Alarms (ALARM-only -> site-alerts): emitter errors/throttles/
  duration + iterator-age (>=10 min) + failures/rejected queue
  depth; rotator standard trio.
- scripts/replay_shoc_webhooks.py: dry-run-default operator replay
  (rebuilds from tables, replay:true envelopes).
- Tests: 742 passing, 85.56% aggregate; golden HMAC vectors shared
  with SHOC in docs/shoc-webhook-test-vectors.json (emitter + replay
  signing pinned to identical vectors); bundle-consistency AST pins
  for both new bundles.
- README: WO stack + webhook feed section, alarm table, runbooks;
  removed stale seahaven-slack-bot consumer references.
2026-07-24 14:54:07 -04:00

152 lines
5.9 KiB
Python

"""HMAC-signed webhook delivery to the SHOC receiver (contract sections 6-7).
Owns the signed POST and its response classification. The signing key comes
from Secrets Manager (``workorder-ingest/shoc-webhook-hmac``, dual-key shape
``{"keys": [{"kid", "secret"}, ...]}``), cached in the warm container for a
short TTL so 30-day rotation propagates without waiting for the execution
environment to recycle -- the producer always signs with ``keys[0]``.
Response contract (section 7): 2xx -> delivered; 429/5xx/timeout/connection
error -> RetryableDeliveryError (the handler surfaces it as a batch item
failure so the ESM blocks the shard and retries in order); any other status ->
("rejected", code) for the handler to park (a contract bug must not block the
shard for 24 hours). Secret material and signatures are never logged.
``sign_body`` is a module-level pure function on purpose: the golden-vector
test suite and the replay script both pin against it, and it is the shared
definition Luby's receiver verifies with.
"""
import hashlib
import hmac
import json
import logging
import os
import time
import urllib.error
import urllib.request
from http import HTTPStatus
import boto3
logger = logging.getLogger()
logger.setLevel(logging.INFO)
# Endpoint path is TBD by SHOC; the activation PR confirms the final URL.
SHOC_WEBHOOK_URL = os.environ.get(
"SHOC_WEBHOOK_URL", "https://api.dev.seahaven.com/api/webhooks/work-orders"
)
HMAC_SECRET_ARN = os.environ.get("HMAC_SECRET_ARN")
POST_TIMEOUT_SECONDS = 10
USER_AGENT = "workorder-shoc-emitter/1"
# 2xx = delivered; 429 and 5xx retry; everything else parks (contract sec. 7).
_HTTP_SUCCESS_RANGE = range(HTTPStatus.OK, HTTPStatus.MULTIPLE_CHOICES)
_HTTP_SERVER_ERROR_MIN = HTTPStatus.INTERNAL_SERVER_ERROR
# Refresh the cached key material this often so a rotated secret propagates
# without waiting for the execution environment to recycle (contract requires
# receiver-side TTL <= 300s; the producer matches it).
_HMAC_KEYS_CACHE_TTL_SECONDS = 300
_hmac_keys_cache = None
_hmac_keys_cached_at = 0.0
class RetryableDeliveryError(Exception):
"""Delivery failed in a way the ESM should retry in order (shard-blocking).
``status_code`` carries the HTTP status when one exists (429/5xx); it is
None for timeouts, connection errors, and secret-fetch failures.
"""
def __init__(self, message: str, status_code: int | None = None):
super().__init__(message)
self.status_code = status_code
def _get_hmac_keys() -> list[dict]:
"""Fetch the signing-key list from Secrets Manager, TTL-cached.
An empty/missing ``keys`` list is the bootstrap state before the first
rotation has run -- retryable, not a crash: the shard blocks until the
rotator populates the secret. Fetch failures (throttle, transient IAM/KMS
denial) are likewise retryable so a Secrets Manager blip blocks in order
instead of failing the whole batch. Key material is never logged.
"""
global _hmac_keys_cache, _hmac_keys_cached_at
now = time.monotonic()
if (
_hmac_keys_cache is not None
and now - _hmac_keys_cached_at < _HMAC_KEYS_CACHE_TTL_SECONDS
):
return _hmac_keys_cache
secrets = boto3.client("secretsmanager")
try:
secret = secrets.get_secret_value(SecretId=HMAC_SECRET_ARN)
keys = json.loads(secret["SecretString"]).get("keys") or []
except Exception as exc:
raise RetryableDeliveryError(f"hmac secret fetch failed: {exc}") from exc
if not keys:
raise RetryableDeliveryError("hmac secret not yet rotated")
_hmac_keys_cache = keys
_hmac_keys_cached_at = now
return keys
def sign_body(secret_hex: str, timestamp: int, raw_body: bytes) -> str:
"""HMAC-SHA256 hex digest over ``f"{timestamp}.{raw_body}"`` (raw bytes).
The key is the UTF-8 bytes of the secret string exactly as stored in the
secret's ``secret`` field (the receiver reads the same JSON field -- no
hex-decoding on either side). Pure function; the golden-vector tests and
Luby's receiver both pin against this definition.
"""
string_to_sign = f"{timestamp}.".encode() + raw_body
return hmac.new(
secret_hex.encode("utf-8"), string_to_sign, hashlib.sha256
).hexdigest()
def deliver(envelope: dict) -> tuple[str, int]:
"""POST one envelope to SHOC. Returns ("delivered"|"rejected", status).
Raises RetryableDeliveryError for 429/5xx/timeout/connection failures so
the caller can block the shard (in-order retry, contract section 7).
"""
raw_body = json.dumps(envelope).encode("utf-8")
timestamp = int(time.time())
signing_key = _get_hmac_keys()[0]
signature = sign_body(signing_key["secret"], timestamp, raw_body)
request = urllib.request.Request(
SHOC_WEBHOOK_URL,
data=raw_body,
headers={
"Content-Type": "application/json; charset=utf-8",
"User-Agent": USER_AGENT,
"X-SH-Timestamp": str(timestamp),
"X-SH-Key-Id": signing_key["kid"],
"X-SH-Signature": f"v1={signature}",
},
method="POST",
)
try:
with urllib.request.urlopen(request, timeout=POST_TIMEOUT_SECONDS) as response:
status_code = response.status
except urllib.error.HTTPError as exc:
status_code = exc.code
if (
status_code == HTTPStatus.TOO_MANY_REQUESTS
or status_code >= _HTTP_SERVER_ERROR_MIN
):
raise RetryableDeliveryError(
f"receiver returned {status_code}", status_code=status_code
) from exc
return ("rejected", status_code)
except (TimeoutError, urllib.error.URLError, OSError) as exc:
raise RetryableDeliveryError(f"connection error: {exc}") from exc
if status_code in _HTTP_SUCCESS_RANGE:
return ("delivered", status_code)
# A non-2xx that urlopen did not raise for (e.g. an unfollowed 3xx):
# contract bug on someone's side -- park it, never block the shard.
return ("rejected", status_code)