mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-10-04 12:31:56 +00:00
117 lines
4.3 KiB
Python
117 lines
4.3 KiB
Python
|
|
"""
|
||
|
|
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
|
||
|
|
|
||
|
|
logger = logging.getLogger()
|
||
|
|
logger.setLevel(logging.INFO)
|
||
|
|
|
||
|
|
REJECTED_QUEUE_URL = os.environ.get("REJECTED_QUEUE_URL")
|
||
|
|
|
||
|
|
# 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")
|
||
|
|
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 handler(event, context):
|
||
|
|
"""Lambda entry point. Triggered by the two WO-table stream ESMs."""
|
||
|
|
for record in event.get("Records", []):
|
||
|
|
webhook_event = envelope.build_event(record)
|
||
|
|
if webhook_event is None:
|
||
|
|
continue
|
||
|
|
start = time.monotonic()
|
||
|
|
try:
|
||
|
|
outcome, status_code = delivery.deliver(webhook_event)
|
||
|
|
except delivery.RetryableDeliveryError as exc:
|
||
|
|
latency_ms = int((time.monotonic() - start) * 1000)
|
||
|
|
logger.warning(
|
||
|
|
_delivery_log(webhook_event, exc.status_code, latency_ms, "retryable")
|
||
|
|
)
|
||
|
|
# 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": record["dynamodb"]["SequenceNumber"]}
|
||
|
|
]
|
||
|
|
}
|
||
|
|
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)
|
||
|
|
return {"batchItemFailures": []}
|