"""Batch-loop tests for the SHOC webhook emitter handler (plan Phase 5). Pins the partial-batch ordering contract: on a retryable failure at record i the loop STOPS -- record i's DynamoDB SequenceNumber is reported via report_batch_item_failures (so the ESM retries from it, in order) and later records are never attempted; earlier in-batch successes are not re-delivered. Non-retryable rejections park the full envelope on the rejected queue and the loop CONTINUES (a contract bug must never block the shard). Skip records (REMOVE / echo guard / blank-text comments) produce no delivery attempts at all. delivery.deliver is monkeypatched on the handler's own sibling instance (the handler imports it by bare name); envelope.build_event runs for real, so these batches exercise the true record -> envelope -> deliver path. SQS is a plain fake on the handler's public ``sqs`` module global. """ import json import pytest from tests.support import load_lambda_module WO_STREAM_ARN = ( "arn:aws:dynamodb:us-east-1:011934824531:table/WorkOrders" "/stream/2026-07-23T00:00:00.000" ) REJECTED_QUEUE_URL = ( "https://sqs.us-east-1.amazonaws.com/011934824531/workorder-shoc-emitter-rejected" ) @pytest.fixture(scope="module") def handler_mod(): return load_lambda_module("wo", "shoc_emitter/handler") class _FakeSQS: def __init__(self): self.sent = [] def send_message(self, QueueUrl, MessageBody): # noqa: N803 (boto3 kwargs) self.sent.append({"QueueUrl": QueueUrl, "MessageBody": MessageBody}) def _wo_record(event_id, sequence_number, event_name="INSERT", **image_overrides): image = { "work_order_id": {"S": f"wo-{event_id}"}, "wo_status": {"S": "new"}, "record_type": {"S": "new_work_order"}, } image.update(image_overrides) stream = { "ApproximateCreationDateTime": 1784642602.0, "SequenceNumber": sequence_number, "StreamViewType": "NEW_AND_OLD_IMAGES", } if event_name != "REMOVE": stream["NewImage"] = image return { "eventID": event_id, "eventName": event_name, "eventSource": "aws:dynamodb", "eventSourceARN": WO_STREAM_ARN, "dynamodb": stream, } def _wire(monkeypatch, handler_mod, outcomes): """Patch delivery.deliver with a per-delivery_id outcome table. ``outcomes`` maps delivery_id (the record eventID) to either a ("delivered"|"rejected", status) tuple or the string "retryable" (raise). Returns the ordered list of attempted delivery_ids and the fake SQS. """ attempted = [] def _fake_deliver(webhook_event): attempted.append(webhook_event["delivery_id"]) outcome = outcomes[webhook_event["delivery_id"]] if outcome == "retryable": raise handler_mod.delivery.RetryableDeliveryError( "receiver returned 503", status_code=503 ) return outcome monkeypatch.setattr(handler_mod.delivery, "deliver", _fake_deliver) fake_sqs = _FakeSQS() monkeypatch.setattr(handler_mod, "sqs", fake_sqs) monkeypatch.setattr(handler_mod, "REJECTED_QUEUE_URL", REJECTED_QUEUE_URL) return attempted, fake_sqs def test_retryable_middle_record_stops_batch_and_reports_its_sequence( monkeypatch, handler_mod ): event = { "Records": [ _wo_record("evt-1", "101"), _wo_record("evt-2", "102"), _wo_record("evt-3", "103"), ] } attempted, fake_sqs = _wire( monkeypatch, handler_mod, {"evt-1": ("delivered", 200), "evt-2": "retryable"}, ) result = handler_mod.handler(event, None) # EXACTLY the failed record's SequenceNumber: the ESM retries from it in # order, and evt-1 (already delivered) is not re-delivered. assert result == {"batchItemFailures": [{"itemIdentifier": "102"}]} # The loop stopped at the failure: record 3 was never attempted. assert attempted == ["evt-1", "evt-2"] assert fake_sqs.sent == [] def test_rejected_record_parks_to_sqs_and_processing_continues( monkeypatch, handler_mod ): event = { "Records": [ _wo_record("evt-1", "201"), _wo_record("evt-2", "202"), _wo_record("evt-3", "203"), ] } attempted, fake_sqs = _wire( monkeypatch, handler_mod, { "evt-1": ("delivered", 200), "evt-2": ("rejected", 422), "evt-3": ("delivered", 204), }, ) result = handler_mod.handler(event, None) # A 4xx rejection must NOT block the shard: no batch item failures, and # the records after the rejection were still attempted. assert result == {"batchItemFailures": []} assert attempted == ["evt-1", "evt-2", "evt-3"] # The full envelope was parked for operator replay. assert len(fake_sqs.sent) == 1 assert fake_sqs.sent[0]["QueueUrl"] == REJECTED_QUEUE_URL parked = json.loads(fake_sqs.sent[0]["MessageBody"]) assert set(parked) == {"envelope", "response_status"} assert parked["response_status"] == 422 assert parked["envelope"]["delivery_id"] == "evt-2" assert parked["envelope"]["event_type"] == "work_order.created" assert parked["envelope"]["data"]["work_order_id"] == "wo-evt-2" def test_skip_records_produce_no_delivery_calls(monkeypatch, handler_mod): event = { "Records": [ _wo_record("evt-1", "301", event_name="REMOVE"), _wo_record("evt-2", "302", write_origin={"S": "shoc-write-api"}), ] } attempted, fake_sqs = _wire(monkeypatch, handler_mod, {}) result = handler_mod.handler(event, None) assert result == {"batchItemFailures": []} assert attempted == [] assert fake_sqs.sent == [] COMMENTS_STREAM_ARN = ( "arn:aws:dynamodb:us-east-1:011934824531:table/WorkOrderComments" "/stream/2026-07-23T00:00:00.000" ) def _comment_record(event_id, sequence_number, **image_overrides): image = { "work_order_id": {"S": f"wo-{event_id}"}, "comment_id": {"S": f"wo-{event_id}#nocomment#abc123"}, "record_type": {"S": "new_work_order"}, "commenter": {"S": ""}, "text": {"S": ""}, "created_at": {"S": "2026-07-16T14:05:00+00:00"}, "ingested_at": {"S": "2026-07-16T14:05:00+00:00"}, } image.update(image_overrides) return { "eventID": event_id, "eventName": "INSERT", "eventSource": "aws:dynamodb", "eventSourceARN": COMMENTS_STREAM_ARN, "dynamodb": { "ApproximateCreationDateTime": 1784642602.0, "SequenceNumber": sequence_number, "StreamViewType": "NEW_AND_OLD_IMAGES", "NewImage": image, }, } def test_blank_comment_skip_logs_and_produces_no_delivery( monkeypatch, handler_mod, caplog ): event = { "Records": [ _comment_record("evt-empty", "501"), _wo_record("evt-wo", "502"), ] } attempted, fake_sqs = _wire( monkeypatch, handler_mod, {"evt-wo": ("delivered", 200)} ) with caplog.at_level("INFO"): result = handler_mod.handler(event, None) assert result == {"batchItemFailures": []} assert attempted == ["evt-wo"] assert fake_sqs.sent == [] skip_logs = [ json.loads(rec.message) for rec in caplog.records if rec.message.startswith("{") and '"shoc_skipped_empty_comment"' in rec.message ] assert len(skip_logs) == 1 assert skip_logs[0]["event"] == "shoc_skipped_empty_comment" assert skip_logs[0]["work_order_id"] == "wo-evt-empty" assert skip_logs[0]["comment_id"] == "wo-evt-empty#nocomment#abc123" def test_unexpected_error_reports_record_not_whole_batch(monkeypatch, handler_mod): # An unexpected exception (here: SQS park failure on a 4xx rejection) must # NOT escape the loop -- that would fail the invocation and make the ESM # re-deliver every earlier success for 24h. The offending record is # reported so the ESM retries from it in order; evt-1 is not re-delivered. event = { "Records": [ _wo_record("evt-1", "401"), _wo_record("evt-2", "402"), _wo_record("evt-3", "403"), ] } attempted, fake_sqs = _wire( monkeypatch, handler_mod, { "evt-1": ("delivered", 200), "evt-2": ("rejected", 400), "evt-3": ("delivered", 200), }, ) def _boom(QueueUrl, MessageBody): # noqa: N803 (boto3 kwargs) raise RuntimeError("sqs unavailable") monkeypatch.setattr(fake_sqs, "send_message", _boom) result = handler_mod.handler(event, None) assert result == {"batchItemFailures": [{"itemIdentifier": "402"}]} # Stopped at the failing record; evt-3 not attempted. assert attempted == ["evt-1", "evt-2"] def test_record_missing_sequence_number_is_skipped_not_crashed( monkeypatch, handler_mod ): # A record that maps to an envelope but lacks SequenceNumber (unreachable # for real streams) must be skipped, not crash the failure paths (Open SWE # #9/#26). Craft one that build_event accepts but strip SequenceNumber. rec = _wo_record("evt-1", "999") del rec["dynamodb"]["SequenceNumber"] attempted, _ = _wire( monkeypatch, handler_mod, {"evt-1": "retryable", "evt-2": ("delivered", 200)} ) good = _wo_record("evt-2", "1000") result = handler_mod.handler({"Records": [rec, good]}, None) # The seq-less record is skipped (never delivered); the good one delivers. assert result == {"batchItemFailures": []} assert "evt-1" not in attempted assert "evt-2" in attempted def test_empty_batch_returns_no_failures(monkeypatch, handler_mod): attempted, _ = _wire(monkeypatch, handler_mod, {}) assert handler_mod.handler({"Records": []}, None) == {"batchItemFailures": []} assert attempted == []