"""Stream-record -> SHOC webhook envelope mapping tests (plan Phase 5). Pins the emitter's classification and envelope-building contract (docs/shoc-webhook-contract.md sections 3-4) against realistic DynamoDB stream records: event-type classification incl. the cancelled transition and the already-cancelled no-false-cancel case, the write_origin echo guard, exact eventSourceARN table discrimination, delivery_id/occurred_at mapping, data null-fill + source_email_s3_key exclusion, and Decimal JSON-safety. Also the enum golden pin: the emitted event_type set must equal the four contract names AND the published spec's webhooks keys (lambdas/api/ openapi.json), and the classification's status/record_type value space must stay anchored to the WO pipeline's template_parser enums. Offline -- envelope.py has no boto3 clients and no network. Table names default to PascalCase; set WORK_ORDERS_TABLE / COMMENTS_TABLE for kebab. """ import json from datetime import datetime, timezone from decimal import Decimal from pathlib import Path import pytest from tests.support import REPO_ROOT, load_lambda_module # Realistic stream ARNs for BOTH tables: "WorkOrders" is a leading substring # of "WorkOrderComments", which is exactly the trap the exact-segment parse # in envelope._table_name exists to avoid. WO_STREAM_ARN = ( "arn:aws:dynamodb:us-east-1:011934824531:table/WorkOrders" "/stream/2026-07-23T00:00:00.000" ) COMMENTS_STREAM_ARN = ( "arn:aws:dynamodb:us-east-1:011934824531:table/WorkOrderComments" "/stream/2026-07-23T00:00:00.000" ) CONTRACT_EVENT_TYPES = { "work_order.created", "work_order.updated", "work_order.cancelled", "work_order.comment_added", } @pytest.fixture(scope="module") def envelope(): return load_lambda_module("wo", "shoc_emitter/envelope") def _record(arn, event_name, new_image=None, old_image=None, **meta): """Build a stream record; ``meta`` overrides event_id / creation.""" stream = { "ApproximateCreationDateTime": meta.get("creation", 1784642602.0), "SequenceNumber": meta.get("sequence_number", "100"), "StreamViewType": "NEW_AND_OLD_IMAGES", } if new_image is not None: stream["NewImage"] = new_image if old_image is not None: stream["OldImage"] = old_image return { "eventID": meta.get("event_id", "4b7c2f0e-0001"), "eventName": event_name, "eventSource": "aws:dynamodb", "eventSourceARN": arn, "dynamodb": stream, } def _wo_image(**overrides): image = { "work_order_id": {"S": "11144580730"}, "wo_status": {"S": "new"}, "description": {"S": "Dock door 14 won't close"}, "record_type": {"S": "new_work_order"}, "created_at": {"S": "2026-07-16T14:03:22.114208+00:00"}, "updated_at": {"S": "2026-07-16T14:03:22.114208+00:00"}, } image.update(overrides) return image def _comment_image(**overrides): image = { "work_order_id": {"S": "11144580730"}, "comment_id": {"S": "11144580730#2026-04-27T23:51:48#a1b2c3d4e5f6"}, "record_type": {"S": "comment"}, "commenter": {"S": "APM Technician"}, "text": {"S": "Vendor dispatched, ETA tomorrow AM."}, "created_at": {"S": "2026-04-27T23:51:48"}, "ingested_at": {"S": "2026-07-16T14:05:00+00:00"}, } image.update(overrides) return image # --- Classification ---------------------------------------------------------- def test_workorders_insert_is_created(envelope): event = envelope.build_event(_record(WO_STREAM_ARN, "INSERT", _wo_image())) assert event["event_type"] == "work_order.created" def test_workorders_modify_is_updated(envelope): event = envelope.build_event( _record( WO_STREAM_ARN, "MODIFY", _wo_image(wo_status={"S": "assigned"}), _wo_image(wo_status={"S": "new"}), ) ) assert event["event_type"] == "work_order.updated" def test_cancelled_transition_is_cancelled(envelope): event = envelope.build_event( _record( WO_STREAM_ARN, "MODIFY", _wo_image(wo_status={"S": "cancelled"}), _wo_image(wo_status={"S": "assigned"}), ) ) assert event["event_type"] == "work_order.cancelled" def test_already_cancelled_modify_is_updated_not_cancelled(envelope): # OLD already cancelled + NEW cancelled: any later touch to a cancelled # WO must NOT re-emit a cancelled event (no false cancels on retries). event = envelope.build_event( _record( WO_STREAM_ARN, "MODIFY", _wo_image(wo_status={"S": "cancelled"}), _wo_image(wo_status={"S": "cancelled"}), ) ) assert event["event_type"] == "work_order.updated" def test_comment_insert_is_comment_added(envelope): event = envelope.build_event( _record(COMMENTS_STREAM_ARN, "INSERT", _comment_image()) ) assert event["event_type"] == "work_order.comment_added" assert set(event["data"]) == set(envelope.COMMENT_DATA_FIELDS) @pytest.mark.parametrize( "overrides", [ {"text": {"S": ""}}, {"text": {"S": " \t\n"}}, {"_drop": "text"}, ], ids=["empty", "whitespace-only", "missing"], ) def test_blank_comment_text_is_skipped(envelope, overrides): image = _comment_image() if overrides.get("_drop") == "text": del image["text"] else: image.update(overrides) record = _record(COMMENTS_STREAM_ARN, "INSERT", image) assert envelope.build_event(record) is None assert envelope.is_blank_comment_skip(record) is True def test_non_blank_comment_is_not_blank_skip(envelope): record = _record(COMMENTS_STREAM_ARN, "INSERT", _comment_image()) assert envelope.build_event(record) is not None assert envelope.is_blank_comment_skip(record) is False def test_remove_is_skipped(envelope): assert envelope.build_event(_record(WO_STREAM_ARN, "REMOVE")) is None assert envelope.build_event(_record(COMMENTS_STREAM_ARN, "REMOVE")) is None assert ( envelope.is_blank_comment_skip(_record(COMMENTS_STREAM_ARN, "REMOVE")) is False ) def test_comment_modify_is_skipped(envelope): record = _record(COMMENTS_STREAM_ARN, "MODIFY", _comment_image(), _comment_image()) assert envelope.build_event(record) is None def test_shoc_write_origin_echo_guard_skips(envelope): record = _record( WO_STREAM_ARN, "INSERT", _wo_image(write_origin={"S": "shoc-write-api"}), ) assert envelope.build_event(record) is None # Any OTHER write_origin value is not an echo and must still deliver. record = _record( WO_STREAM_ARN, "INSERT", _wo_image(write_origin={"S": "email-processor"}), ) assert envelope.build_event(record) is not None # --- eventSourceARN table discrimination ------------------------------------- def test_table_discrimination_is_exact_segment_not_substring(envelope): # The WorkOrders stream must classify as a WO state event and the # WorkOrderComments stream as a comment event -- never cross-parsed even # though "WorkOrders" is a leading substring of "WorkOrderComments". wo_event = envelope.build_event(_record(WO_STREAM_ARN, "INSERT", _wo_image())) assert wo_event["event_type"] == "work_order.created" comment_event = envelope.build_event( _record(COMMENTS_STREAM_ARN, "INSERT", _comment_image()) ) assert comment_event["event_type"] == "work_order.comment_added" # Superstring table names must not substring-match either known table. for table in ("WorkOrdersLegacy", "WorkOrderCommentsArchive", "XWorkOrders"): arn = ( f"arn:aws:dynamodb:us-east-1:011934824531:table/{table}" "/stream/2026-07-23T00:00:00.000" ) assert envelope.build_event(_record(arn, "INSERT", _wo_image())) is None # A malformed / missing ARN skips rather than raises (total function). assert envelope.build_event(_record("", "INSERT", _wo_image())) is None def test_kebab_table_names_via_env(monkeypatch): """PLAT-11: classification follows WORK_ORDERS_TABLE / COMMENTS_TABLE.""" import sys monkeypatch.setenv("WORK_ORDERS_TABLE", "work-orders") monkeypatch.setenv("COMMENTS_TABLE", "work-order-comments") sys.modules.pop("wo_shoc_emitter_envelope", None) kebab_envelope = load_lambda_module("wo", "shoc_emitter/envelope") wo_arn = ( "arn:aws:dynamodb:us-east-1:011934824531:table/work-orders" "/stream/2026-08-07T00:00:00.000" ) comments_arn = ( "arn:aws:dynamodb:us-east-1:011934824531:table/work-order-comments" "/stream/2026-08-07T00:00:00.000" ) assert ( kebab_envelope.build_event(_record(wo_arn, "INSERT", _wo_image()))["event_type"] == "work_order.created" ) assert ( kebab_envelope.build_event(_record(comments_arn, "INSERT", _comment_image()))[ "event_type" ] == "work_order.comment_added" ) # Legacy PascalCase ARNs must not match kebab env names. assert ( kebab_envelope.build_event(_record(WO_STREAM_ARN, "INSERT", _wo_image())) is None ) # Restore cached module for later tests that expect PascalCase defaults. sys.modules.pop("wo_shoc_emitter_envelope", None) monkeypatch.delenv("WORK_ORDERS_TABLE", raising=False) monkeypatch.delenv("COMMENTS_TABLE", raising=False) load_lambda_module("wo", "shoc_emitter/envelope") # --- Envelope fields --------------------------------------------------------- def test_envelope_constant_fields_and_delivery_id(envelope): record = _record(WO_STREAM_ARN, "INSERT", _wo_image(), event_id="evt-abc-123") event = envelope.build_event(record) assert event["schema_version"] == 1 assert event["source"] == "procurement-ingest/workorder-shoc-emitter" assert event["replay"] is False # delivery_id is the stream record eventID: unique per source event and # stable across ESM retries (the receiver's idempotency key). assert event["delivery_id"] == "evt-abc-123" def test_occurred_at_from_decimal_creation_time(envelope): # ApproximateCreationDateTime arrives as epoch seconds; via the low-level # stream API it deserializes as Decimal -- must still map to an aware # ISO-8601 +00:00 instant. creation = Decimal("1784642602") record = _record(WO_STREAM_ARN, "INSERT", _wo_image(), creation=creation) event = envelope.build_event(record) expected = datetime.fromtimestamp(1784642602.0, tz=timezone.utc).isoformat() assert event["occurred_at"] == expected assert event["occurred_at"].endswith("+00:00") def test_data_null_fills_absent_fields_and_excludes_s3_key(envelope): image = { "work_order_id": {"S": "11144580730"}, "wo_status": {"S": "unknown"}, # Internal pointer: must never leave the account via the webhook. "source_email_s3_key": {"S": "inbound/2026/07/abc.eml"}, } event = envelope.build_event(_record(WO_STREAM_ARN, "INSERT", image)) assert set(event["data"]) == set(envelope.WO_DATA_FIELDS) assert "source_email_s3_key" not in event["data"] assert event["data"]["work_order_id"] == "11144580730" for field in set(envelope.WO_DATA_FIELDS) - {"work_order_id", "wo_status"}: assert event["data"][field] is None def test_decimal_image_values_become_json_safe(envelope): image = _wo_image( severity={"N": "3"}, priority={"N": "2.5"}, ) event = envelope.build_event(_record(WO_STREAM_ARN, "INSERT", image)) assert event["data"]["severity"] == 3 assert isinstance(event["data"]["severity"], int) assert event["data"]["priority"] == 2.5 assert isinstance(event["data"]["priority"], float) # The whole envelope must serialize -- this is the raw body that gets # signed and POSTed. json.dumps(event) # --- Enum golden pin --------------------------------------------------------- def test_event_type_set_matches_contract_and_openapi_webhooks(envelope): scenarios = [ _record(WO_STREAM_ARN, "INSERT", _wo_image()), _record( WO_STREAM_ARN, "MODIFY", _wo_image(wo_status={"S": "assigned"}), _wo_image(wo_status={"S": "new"}), ), _record( WO_STREAM_ARN, "MODIFY", _wo_image(wo_status={"S": "cancelled"}), _wo_image(wo_status={"S": "assigned"}), ), _record(COMMENTS_STREAM_ARN, "INSERT", _comment_image()), ] emitted = {envelope.build_event(record)["event_type"] for record in scenarios} assert emitted == CONTRACT_EVENT_TYPES # The published read-API spec documents the same four outbound webhook # events; emitter and spec must move together. spec_path = Path(REPO_ROOT) / "lambdas" / "api" / "openapi.json" spec = json.loads(spec_path.read_text(encoding="utf-8")) assert set(spec["webhooks"]) == emitted def test_status_and_record_type_value_space_anchored_to_template_parser(envelope): # The WO pipeline's template_parser enums are the wo_status/record_type # value space the stream carries (prompts.py mirrors them for the AI # path). The cancelled-transition trigger value must be a real status, # and the contract's record_type enum must be exactly the pipeline's # email-type enum -- widen either side only in a PR that updates both. template_parser = load_lambda_module("wo", "email_processor/template_parser") assert envelope.CANCELLED_STATUS in template_parser.VALID_STATUSES assert template_parser.VALID_EMAIL_TYPES == { "new_work_order", "update", "comment", "cancellation", }