""" SHOC work-order webhook emitter Lambda. Consumes the WorkOrders/WorkOrderComments DynamoDB streams and pushes each mutation to SHOC as an HMAC-signed HTTPS POST (docs/shoc-webhook-contract.md). This handler is the thin event loop; the work lives in flat sibling modules (bare-name imports resolve via the same flat-landing bundling as the email processor's siblings): envelope.py -- stream record -> envelope mapping + event classification delivery.py -- secret cache, signing, POST, response classification Ships DARK: both event-source mappings deploy with enabled=False, so this code runs zero deliveries until the activation PR flips them on after SHOC's receiver passes the shared HMAC test vectors. Ordering semantics: the ESMs run parallelization_factor=1 with bisect-on-error off, and this loop processes records strictly in order. A retryable failure (429/5xx/timeout/connection error) stops the batch immediately and reports that record via report_batch_item_failures -- earlier successes are not re-delivered, and the ESM blocks the shard and retries from the failed record, preserving per-work-order commit order. Non-retryable 4xx responses are a contract bug, not an availability blip: the full envelope is parked on the rejected queue (alarmed) and the loop continues, so a bad payload can never block the shard for 24 hours. No healthcheck branch on purpose: stream consumers are not smoke-gated (po-ingest-site-extractor precedent) -- there is no direct-invoke path to probe, and a synthetic stream record would be a real delivery. """ import json import logging import os import time import boto3 import delivery import envelope from botocore.config import Config logger = logging.getLogger() logger.setLevel(logging.INFO) REJECTED_QUEUE_URL = os.environ.get("REJECTED_QUEUE_URL") # Bound the SQS client's timeouts (Open SWE #17): a slow SQS response while # parking a rejected envelope must not hang the invocation toward its 60s # timeout and re-deliver the whole batch. _BOTO_CONFIG = Config( connect_timeout=3, read_timeout=5, retries={"max_attempts": 2, "mode": "standard"} ) # Lazy cached SQS client. Keeps the public attribute name ``sqs`` so the test # monkeypatch target changes module only, not attribute name. sqs = None def _get_sqs(): global sqs if sqs is None: sqs = boto3.client("sqs", config=_BOTO_CONFIG) return sqs def _delivery_log(event: dict, status_code, latency_ms: int, outcome: str) -> str: """One structured record per delivery attempt (Logs-Insights-queryable).""" return json.dumps( { "event": "shoc_delivery", "delivery_id": event["delivery_id"], "event_type": event["event_type"], "work_order_id": event["data"].get("work_order_id"), "status_code": status_code, "latency_ms": latency_ms, "outcome": outcome, } ) def _park_rejected(event: dict, status_code: int): """Send the full envelope to the rejected queue for operator replay.""" _get_sqs().send_message( QueueUrl=REJECTED_QUEUE_URL, MessageBody=json.dumps({"envelope": event, "response_status": status_code}), ) logger.warning( json.dumps( { "event": "shoc_delivery_rejected", "delivery_id": event["delivery_id"], "event_type": event["event_type"], "work_order_id": event["data"].get("work_order_id"), "response_status": status_code, } ) ) def _batch_summary(records: int, emitted: int, skipped: int) -> str: """Per-invocation emit/skip tally. Skips (REMOVE, echo guard, blank-text comments, comment MODIFYs, unmappable records) return no delivery and would otherwise be invisible: a systemic classification break -- e.g. a table rename desyncing the eventSourceARN parse -- would drop 100%% of events while every batch still reports success. Logging the ratio makes that queryable in Logs Insights and alarmable. Blank-text comment skips also emit a distinct ``shoc_skipped_empty_comment`` line so their volume stays observable apart from other skip reasons. """ return json.dumps( { "event": "shoc_batch_summary", "records": records, "emitted": emitted, "skipped": skipped, } ) def _log_empty_comment_skip(record: dict) -> None: """Structured log for a blank-text comment_added skip (Logs Insights).""" stream = record.get("dynamodb") or {} # Attribute-value encoded NewImage; only the S forms are read for the # log fields -- classification already decided this is a blank skip. new_image = stream.get("NewImage") or {} work_order_id = (new_image.get("work_order_id") or {}).get("S") comment_id = (new_image.get("comment_id") or {}).get("S") logger.info( json.dumps( { "event": "shoc_skipped_empty_comment", "work_order_id": work_order_id, "comment_id": comment_id, } ) ) def handler(event, context): """Lambda entry point. Triggered by the two WO-table stream ESMs.""" records = event.get("Records", []) emitted = 0 skipped = 0 for record in records: webhook_event = envelope.build_event(record) if webhook_event is None: skipped += 1 if envelope.is_blank_comment_skip(record): _log_empty_comment_skip(record) continue # Capture the sequence number BEFORE any delivery attempt so the # failure paths below can never KeyError on the subscript (Open SWE # #9/#26 -- the catch-all exists to keep an exception from re-delivering # the whole batch, so it must not itself raise). SequenceNumber is # present on every real DynamoDB stream record; a record lacking it is # unmappable-to-a-checkpoint, so log and skip rather than crash. seq = record.get("dynamodb", {}).get("SequenceNumber") if seq is None: logger.warning( json.dumps( { "event": "shoc_missing_sequence_number", "delivery_id": webhook_event["delivery_id"], } ) ) skipped += 1 continue emitted += 1 start = time.monotonic() try: outcome, status_code = delivery.deliver(webhook_event) latency_ms = int((time.monotonic() - start) * 1000) logger.info(_delivery_log(webhook_event, status_code, latency_ms, outcome)) if outcome == "rejected": _park_rejected(webhook_event, status_code) except delivery.RetryableDeliveryError as exc: latency_ms = int((time.monotonic() - start) * 1000) logger.warning( _delivery_log(webhook_event, exc.status_code, latency_ms, "retryable") ) logger.info(_batch_summary(len(records), emitted, skipped)) # Stop here: reporting this record's sequence number makes the ESM # retry from it in order; earlier successes are not re-delivered. return {"batchItemFailures": [{"itemIdentifier": seq}]} except Exception: # Catch-all so an unexpected error (e.g. an SQS park failure) # cannot escape and fail the WHOLE invocation -- that would make # the ESM re-deliver every earlier success in the batch for up to # 24h. Report only THIS record so the ESM retries from it in order. logger.exception( json.dumps( { "event": "shoc_delivery_error", "delivery_id": webhook_event["delivery_id"], "event_type": webhook_event["event_type"], } ) ) logger.info(_batch_summary(len(records), emitted, skipped)) return {"batchItemFailures": [{"itemIdentifier": seq}]} logger.info(_batch_summary(len(records), emitted, skipped)) return {"batchItemFailures": []}