mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 06:03:14 +00:00
* 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
473 lines
16 KiB
Python
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()
|