"""Stream-record -> SHOC webhook envelope mapping (contract sections 3-4). Pure and total: no boto3 clients, no environment reads, no network. Every function either returns a value or returns None (skip) -- a malformed stream record must never raise here; classification gaps fall through to None so the shard is never blocked by an unmappable record. ``build_event`` is the single entry point the handler calls per record. Classification (docs/shoc-webhook-contract.md section 4): WorkOrders INSERT -> work_order.created WorkOrders MODIFY -> work_order.cancelled iff OLD wo_status was not "cancelled" AND NEW wo_status is "cancelled", else work_order.updated WorkOrderComments INSERT -> work_order.comment_added REMOVE (either table) -> skip (no deletes in this pipeline) NEW write_origin == "shoc-write-api" -> skip (echo guard; a no-op branch today -- nothing writes that attribute yet -- kept so the phase-2 write-back API never echoes SHOC's own writes back at it) Anything else (unknown table, comment MODIFYs) -> skip. """ from datetime import datetime, timezone from decimal import Decimal from boto3.dynamodb.types import TypeDeserializer SCHEMA_VERSION = 1 SOURCE = "procurement-ingest/workorder-shoc-emitter" CANCELLED_STATUS = "cancelled" SHOC_WRITE_ORIGIN = "shoc-write-api" WORK_ORDERS_TABLE = "WorkOrders" COMMENTS_TABLE = "WorkOrderComments" # data payload field lists (contract sections 4.1 / 4.2). Absent attributes are # null-filled; source_email_s3_key is deliberately EXCLUDED (internal key, not # part of the contract). 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", ) _deserializer = TypeDeserializer() # arn:aws:dynamodb:...:table/NAME[/stream/LABEL] -- slash-split needs at least # the ":table" head plus the name segment for the parse to be meaningful. _MIN_ARN_SLASH_SEGMENTS = 2 def _plain(value): """Convert TypeDeserializer output to json.dumps-safe plain Python. Mirrors lambdas/api/serialization.py: integral Decimals become JSON integers (exact), non-integral Decimals become floats, sets become sorted lists. Local copy on purpose -- the emitter bundles flat and must not grow a cross-package import for twelve lines of logic. """ if isinstance(value, Decimal): if value == value.to_integral_value(): return int(value) return float(value) if isinstance(value, dict): return {key: _plain(item) for key, item in value.items()} if isinstance(value, list): return [_plain(item) for item in value] if isinstance(value, set): return sorted(_plain(item) for item in value) return value def _deserialize_image(image: dict) -> dict: """DynamoDB stream image (attribute-value encoded) -> plain dict.""" return {key: _plain(_deserializer.deserialize(av)) for key, av in image.items()} def _table_name(event_source_arn: str) -> str | None: """Extract the exact table-name segment from a stream eventSourceARN. ARN format: arn:aws:dynamodb:region:acct:table/NAME/stream/LABEL. Parsing the segment exactly (not a loose substring) matters: "WorkOrders" is a prefix of nothing, but a substring test for it WOULD match "WorkOrderComments"-adjacent names -- segment equality can't. """ parts = event_source_arn.split("/") if len(parts) >= _MIN_ARN_SLASH_SEGMENTS and parts[0].endswith(":table"): return parts[1] return None def _classify(table, event_name, new_image, old_image) -> str | None: if table == WORK_ORDERS_TABLE: if event_name == "INSERT": return "work_order.created" if event_name == "MODIFY": was_cancelled = old_image.get("wo_status") == CANCELLED_STATUS if not was_cancelled and new_image.get("wo_status") == CANCELLED_STATUS: return "work_order.cancelled" return "work_order.updated" return None if table == COMMENTS_TABLE and event_name == "INSERT": return "work_order.comment_added" return None def build_event(record: dict) -> dict | None: """Map one DynamoDB stream record to a webhook envelope, or None to skip.""" event_name = record.get("eventName") if event_name == "REMOVE": return None stream = record.get("dynamodb") or {} new_image = _deserialize_image(stream.get("NewImage") or {}) if new_image.get("write_origin") == SHOC_WRITE_ORIGIN: return None # echo guard (forward-compat, contract section 3) table = _table_name(record.get("eventSourceARN") or "") old_image = _deserialize_image(stream.get("OldImage") or {}) event_type = _classify(table, event_name, new_image, old_image) if event_type is None: return None if event_type == "work_order.comment_added": fields = COMMENT_DATA_FIELDS else: fields = WO_DATA_FIELDS # eventID and ApproximateCreationDateTime are present on every real # DynamoDB stream record; guarding them keeps this function total (the # "never raise" contract above) rather than trusting a hard subscript. event_id = record.get("eventID") approx_creation = stream.get("ApproximateCreationDateTime") if event_id is None or approx_creation is None: return None # ApproximateCreationDateTime arrives as epoch seconds (float/Decimal). occurred_at = datetime.fromtimestamp( float(approx_creation), tz=timezone.utc ).isoformat() return { "schema_version": SCHEMA_VERSION, "delivery_id": event_id, "event_type": event_type, "occurred_at": occurred_at, "source": SOURCE, "replay": False, "data": {field: new_image.get(field) for field in fields}, }