procurement-ingest/scripts/replay_shoc_webhooks.py
Adam Moussa 13efee9682
feat(infra): add kebab WO tables and PLAT-11 cutover tooling (PLAT-11) (#172)
* feat(infra): add kebab WO tables and PLAT-11 cutover tooling

Create empty work-orders/work-order-comments under Terraform while writers
stay on PascalCase; make emitter table classification env-driven and add
same-account migrate/verify plus cutover runbook.

* style(test): ruff-format shoc emitter envelope tests
2026-08-07 13:42:35 -04:00

473 lines
16 KiB
Python

"""
Rebuild SHOC work-order webhook events from DynamoDB and re-POST them to a
receiver endpoint (contract section 8 replay backstop, for the
workorder-shoc-emitter-failures / -rejected alarm runbook).
Dry-run is the DEFAULT in every mode; nothing is POSTed unless --execute is
passed. Validations (account gate, --since parsing, secret shape, unknown
--work-order-id) run on dry-run too, so the rehearsal surfaces the same
failures the real run would.
Semantics (deliberate divergences from the live emitter):
* Every envelope carries "replay": true.
* Every WorkOrders state row is emitted as work_order.updated -- a replay
reads current table state, not the stream, so it cannot distinguish the
original created (or cancelled transition). The receiver upserts
idempotently, so updated is always safe.
* delivery_id is deterministic: "replay-" + sha256 over
table#work_order_id#comment_id#occurred_at, truncated. Re-running the
same replay yields the SAME ids -- receivers dedupe on delivery_id, so
overlapping replay runs are harmless.
* occurred_at is the row's updated_at (WO) / ingested_at (comment), falling
back to now-UTC when absent.
Selection is exactly one of --work-order-id (repeatable; per-item get_item /
Query) or --since (ISO-8601 UTC; a FULL TABLE Scan on WorkOrders filtered on
updated_at >= since and on WorkOrderComments filtered on ingested_at >= since
-- a count banner is printed per table). --events narrows to state rows,
comments, or both (default).
Usage:
# dry-run two work orders (state + comments)
python scripts/replay_shoc_webhooks.py \
--url https://api.dev.seahaven.com/api/webhooks/work-orders \
--work-order-id 11144580730 --work-order-id 11144580731
# replay everything touched since a timestamp (full table Scan)
python scripts/replay_shoc_webhooks.py --url https://... \
--since 2026-07-24T02:00:00Z --execute
# comments only, for one work order
python scripts/replay_shoc_webhooks.py --url https://... \
--work-order-id 11144580730 --events comments --execute
--url is required with no default: replay must be a deliberate act against a
known receiver. Signing mirrors the emitter's delivery module byte-for-byte
(same string-to-sign, same header names; only the User-Agent gains a
"-replay" suffix). Non-2xx responses are counted as failures; the script
keeps going and exits 1 if any event failed.
"""
import argparse
import hashlib
import hmac
import json
import sys
import time
import urllib.error
import urllib.request
from datetime import datetime, timedelta, timezone
from decimal import Decimal
import os
import boto3
from boto3.dynamodb.conditions import Key
EXPECTED_ACCOUNT = "011934824531"
REGION = "us-east-1"
WO_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders")
COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments")
SECRET_NAME = "workorder-ingest/shoc-webhook-hmac"
SCHEMA_VERSION = 1
SOURCE = "procurement-ingest/workorder-shoc-emitter"
USER_AGENT = "workorder-shoc-emitter/1-replay"
POST_TIMEOUT_SECONDS = 10
DELIVERY_ID_HASH_CHARS = 32
HTTP_2XX_MIN = 200
HTTP_2XX_MAX = 300
HTTP_TOO_MANY_REQUESTS = 429
HTTP_5XX_MIN = 500
class _NoRedirectHandler(urllib.request.HTTPRedirectHandler):
"""Refuse to follow receiver redirects (matches the emitter's opener).
Following a 3xx would forward the live X-SH-* auth headers to a
receiver-chosen Location and could downgrade an http:// target past the
https-only --url check; a redirect is a receiver misconfiguration, so let
it surface as an HTTPError status.
"""
def redirect_request(self, *args, **kwargs):
return None
_opener = urllib.request.build_opener(_NoRedirectHandler)
# data payload field lists -- pinned to the emitter's envelope module. Absent
# attributes are emitted as null. source_email_s3_key is deliberately
# EXCLUDED (internal S3 pointer, never leaves the account).
WO_DATA_FIELDS = [
"work_order_id",
"wo_status",
"description",
"customer",
"site_code",
"building",
"address",
"severity",
"priority",
"assigned_to",
"date_reported",
"scheduled_start",
"due_date",
"record_type",
"created_at",
"updated_at",
]
COMMENT_DATA_FIELDS = [
"work_order_id",
"comment_id",
"record_type",
"commenter",
"text",
"created_at",
"ingested_at",
]
def parse_since(value):
"""Parse the operator-supplied --since and re-emit the canonical
zero-padded ``+00:00`` isoformat string that ``updated_at`` /
``ingested_at`` are lexicographically comparable against (both are
``datetime.now(timezone.utc).isoformat()`` strings). A raw operator
string would compare WRONG under DynamoDB's byte-wise ``>=``, so parse
strictly and require an aware UTC instant."""
raw = value[:-1] + "+00:00" if value.endswith("Z") else value
try:
dt = datetime.fromisoformat(raw)
except ValueError:
sys.exit(
f"ERROR: --since {value!r} is not ISO-8601 (e.g. 2026-07-24T02:00:00Z)."
)
if dt.tzinfo is None or dt.utcoffset() != timedelta(0):
sys.exit(f"ERROR: --since {value!r} must be UTC (trailing 'Z' or '+00:00').")
return dt.isoformat()
def _json_safe(value):
"""boto3's DynamoDB resource deserializes numbers as Decimal, which
json.dumps rejects; convert recursively (defensive -- these fields are
string-typed today)."""
if isinstance(value, Decimal):
return int(value) if value == value.to_integral_value() else float(value)
if isinstance(value, list):
return [_json_safe(v) for v in value]
if isinstance(value, dict):
return {k: _json_safe(v) for k, v in value.items()}
return value
def sign_body(secret_hex, timestamp, raw_body_bytes):
"""HMAC-SHA256 over ``f"{timestamp}.{raw_body}"`` (contract section 6),
computed over the raw UTF-8 body bytes exactly as the emitter's delivery
module does -- the shared test vectors pin both implementations to
identical output. Never log the secret."""
string_to_sign = f"{timestamp}.".encode("utf-8") + raw_body_bytes
mac = hmac.new(secret_hex.encode("utf-8"), string_to_sign, hashlib.sha256)
return mac.hexdigest()
def build_headers(kid, signature, timestamp):
return {
"Content-Type": "application/json; charset=utf-8",
"User-Agent": USER_AGENT,
"X-SH-Timestamp": str(timestamp),
"X-SH-Key-Id": kid,
"X-SH-Signature": f"v1={signature}",
}
def post_event(url, raw_body_bytes, kid, secret_hex):
"""POST one signed envelope. Returns the HTTP status code, or None on
timeout / connection error."""
timestamp = int(time.time())
signature = sign_body(secret_hex, timestamp, raw_body_bytes)
request = urllib.request.Request(
url,
data=raw_body_bytes,
headers=build_headers(kid, signature, timestamp),
method="POST",
)
try:
with _opener.open(request, timeout=POST_TIMEOUT_SECONDS) as resp:
return resp.status
except urllib.error.HTTPError as exc:
return exc.code
except (urllib.error.URLError, TimeoutError, OSError):
return None
def classify_response(status):
"""Contract section 7 classification, printed per event for the operator
(the replay tool itself never retries -- re-run it instead)."""
if status is not None and HTTP_2XX_MIN <= status < HTTP_2XX_MAX:
return "delivered"
if status is None or status == HTTP_TOO_MANY_REQUESTS or status >= HTTP_5XX_MIN:
return "retryable"
return "non-retryable"
def fetch_signing_key(session, secret_id):
"""Fetch the HMAC secret once and return (kid, secret_hex) from keys[0]
(the producer signing key). Secret material is never printed."""
client = session.client("secretsmanager")
value = json.loads(client.get_secret_value(SecretId=secret_id)["SecretString"])
keys = value.get("keys") or []
if not keys:
sys.exit(
f"ERROR: secret {secret_id!r} has no keys (bootstrap state); run the "
"workorder-shoc-hmac-rotator rotation first."
)
return keys[0]["kid"], keys[0]["secret"]
def replay_delivery_id(table, work_order_id, comment_id, occurred_at):
"""Deterministic delivery_id, stable across replay runs, so receivers
dedupe repeated replays on it."""
seed = f"{table}#{work_order_id}#{comment_id}#{occurred_at}"
digest = hashlib.sha256(seed.encode("utf-8")).hexdigest()
return f"replay-{digest[:DELIVERY_ID_HASH_CHARS]}"
def _envelope(event_type, delivery_id, occurred_at, data):
return {
"schema_version": SCHEMA_VERSION,
"delivery_id": delivery_id,
"event_type": event_type,
"occurred_at": occurred_at,
"source": SOURCE,
"replay": True,
"data": data,
}
def build_state_event(item):
# Always work_order.updated: current-state replay cannot distinguish the
# original created/cancelled, and the receiver upserts idempotently.
occurred_at = item.get("updated_at") or datetime.now(timezone.utc).isoformat()
return _envelope(
"work_order.updated",
replay_delivery_id(WO_TABLE, item["work_order_id"], "", occurred_at),
occurred_at,
{f: _json_safe(item.get(f)) for f in WO_DATA_FIELDS},
)
def build_comment_event(item):
# Skip blank/missing text the same way the live emitter does -- otherwise
# a --since / --work-order-id replay would re-POST #nocomment# rows and
# re-park them on the rejected queue.
text = item.get("text")
if text is None or (isinstance(text, str) and not text.strip()):
return None
occurred_at = item.get("ingested_at") or datetime.now(timezone.utc).isoformat()
delivery_id = replay_delivery_id(
COMMENTS_TABLE, item["work_order_id"], item["comment_id"], occurred_at
)
return _envelope(
"work_order.comment_added",
delivery_id,
occurred_at,
{f: _json_safe(item.get(f)) for f in COMMENT_DATA_FIELDS},
)
def query_comments(table, work_order_id):
kwargs = {"KeyConditionExpression": Key("work_order_id").eq(work_order_id)}
while True:
page = table.query(**kwargs)
yield from page.get("Items", [])
lek = page.get("LastEvaluatedKey")
if not lek:
return
kwargs["ExclusiveStartKey"] = lek
def scan_since(table, attr, since):
kwargs = {
"FilterExpression": "#a >= :since",
"ExpressionAttributeNames": {"#a": attr},
"ExpressionAttributeValues": {":since": since},
}
while True:
page = table.scan(**kwargs)
yield from page.get("Items", [])
lek = page.get("LastEvaluatedKey")
if not lek:
return
kwargs["ExclusiveStartKey"] = lek
def _append_comment_events(envelopes, items):
"""Build comment envelopes, dropping blank-text skips; return skip count."""
skipped = 0
for item in items:
event = build_comment_event(item)
if event is None:
skipped += 1
continue
envelopes.append(event)
return skipped
def select_by_ids(dynamodb, work_order_ids, want_state, want_comments):
envelopes = []
skipped_empty = 0
if want_state:
table = dynamodb.Table(WO_TABLE)
for wid in work_order_ids:
item = table.get_item(Key={"work_order_id": wid}).get("Item")
if item is None:
sys.exit(f"ERROR: work order {wid!r} not found in {WO_TABLE}.")
envelopes.append(build_state_event(item))
if want_comments:
table = dynamodb.Table(COMMENTS_TABLE)
for wid in work_order_ids:
skipped_empty += _append_comment_events(
envelopes, query_comments(table, wid)
)
if skipped_empty:
print(f"skipped {skipped_empty} empty-text comment(s)")
return envelopes
def select_since(dynamodb, since, want_state, want_comments):
envelopes = []
skipped_empty = 0
if want_state:
table = dynamodb.Table(WO_TABLE)
items = list(scan_since(table, "updated_at", since))
print(
f"[{WO_TABLE}] full table Scan (updated_at >= {since}): {len(items)} row(s)"
)
envelopes.extend(build_state_event(i) for i in items)
if want_comments:
table = dynamodb.Table(COMMENTS_TABLE)
items = list(scan_since(table, "ingested_at", since))
print(
f"[{COMMENTS_TABLE}] full table Scan (ingested_at >= {since}): "
f"{len(items)} row(s)"
)
skipped_empty += _append_comment_events(envelopes, items)
if skipped_empty:
print(f"skipped {skipped_empty} empty-text comment(s)")
return envelopes
def main():
ap = argparse.ArgumentParser(description=__doc__)
ap.add_argument(
"--work-order-id",
action="append",
help="Replay this work order (repeatable). Exactly one of "
"--work-order-id or --since is required.",
)
ap.add_argument(
"--since",
help="Replay every row with updated_at/ingested_at >= this ISO-8601 "
"UTC timestamp (FULL TABLE Scan on both tables).",
)
ap.add_argument(
"--events",
choices=["state", "comments", "both"],
default="both",
help="Which event families to replay (default: both).",
)
ap.add_argument(
"--url",
required=True,
help="Receiver endpoint URL. Required, no default -- replay must be "
"deliberate.",
)
ap.add_argument(
"--secret-arn",
default=SECRET_NAME,
help=f"HMAC secret to sign with (default: resolve {SECRET_NAME!r} by name).",
)
ap.add_argument("--profile", help="AWS profile (default: seahaven-prod).")
ap.add_argument(
"--execute",
action="store_true",
help="Actually POST. Default is dry-run (list only).",
)
args = ap.parse_args()
if bool(args.work_order_id) == bool(args.since):
sys.exit("ERROR: pass exactly one of --work-order-id or --since.")
if not args.url.startswith("https://"):
# urllib follows file:// and http:// too; the feed is HTTPS-only.
sys.exit("ERROR: --url must be an https:// URL.")
session = boto3.Session(
profile_name=args.profile or "seahaven-prod", region_name=REGION
)
acct = session.client("sts").get_caller_identity()["Account"]
if acct != EXPECTED_ACCOUNT:
sys.exit(
f"ERROR: profile resolves to account {acct}, expected "
f"{EXPECTED_ACCOUNT}. Aborting."
)
want_state = args.events in ("state", "both")
want_comments = args.events in ("comments", "both")
dynamodb = session.resource("dynamodb")
if args.work_order_id:
events = select_by_ids(dynamodb, args.work_order_id, want_state, want_comments)
else:
since = parse_since(args.since)
events = select_since(dynamodb, since, want_state, want_comments)
# Fetch on dry-run too: an empty/missing secret should surface on the
# rehearsal, not first appear when the operator adds --execute.
kid, secret_hex = fetch_signing_key(session, args.secret_arn)
if not events:
print("No matching events -- nothing to replay.")
return
# In-order per work order: state and comments interleave by occurred_at
# (receiver upserts + skeleton-upserts, so cross-family order is a
# nicety, not a requirement).
events.sort(
key=lambda e: (e["data"]["work_order_id"], e["occurred_at"], e["event_type"])
)
state_count = sum(1 for e in events if e["event_type"] == "work_order.updated")
print(
f"{len(events)} event(s) selected ({state_count} state, "
f"{len(events) - state_count} comment) -> {args.url}\n"
)
failed = 0
for event in events:
wid = event["data"]["work_order_id"]
if not args.execute:
print(
f"[{event['delivery_id']}] DRY-RUN: would POST {event['event_type']} wo={wid}"
)
continue
raw_body = json.dumps(event).encode("utf-8")
status = post_event(args.url, raw_body, kid, secret_hex)
outcome = classify_response(status)
status_str = "conn-error/timeout" if status is None else str(status)
print(
f"[{event['delivery_id']}] POST {event['event_type']} wo={wid} -> "
f"{status_str} ({outcome})"
)
if outcome != "delivered":
failed += 1
if not args.execute:
print("\nDry run complete. Re-run with --execute to POST.")
return
print(f"\n{len(events) - failed} delivered, {failed} failed.")
if failed:
sys.exit(1)
if __name__ == "__main__":
main()