From 272cbf6d9288c3e4b170e32b6edabdeeb59786ab Mon Sep 17 00:00:00 2001 From: Vadim Malanov Date: Sun, 12 Jul 2026 12:39:02 +0300 Subject: [PATCH] feat(dispatch): consume QMS asset events from the service inbox (C2) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit LegacyHUB is subscribed to AssetRegistered/AssetApproved from QMS but never polled its service inbox. New consumer maps the §7 payload to an AssetManifestEnvelope and runs the local knowledge-ingest, then confirms the delivery. Outcome policy avoids head-of-line blocking: Registered → ingest (accepted/duplicate/rejected all confirm), Approved/unknown → confirm without work, malformed → confirm as poison; only transient infrastructure failures (object_read_failed) leave the message pending and stop the run for the next scheduled retry. Thin CLI wrapper scripts/consume_dispatch_inbox.py (scripts/ ships in the image). Co-Authored-By: Claude Fable 5 --- app/integrations/dispatch_inbox_consumer.py | 191 ++++++++++++++++++++ scripts/consume_dispatch_inbox.py | 24 +++ tests/test_dispatch_inbox_consumer.py | 160 ++++++++++++++++ 3 files changed, 375 insertions(+) create mode 100644 app/integrations/dispatch_inbox_consumer.py create mode 100644 scripts/consume_dispatch_inbox.py create mode 100644 tests/test_dispatch_inbox_consumer.py diff --git a/app/integrations/dispatch_inbox_consumer.py b/app/integrations/dispatch_inbox_consumer.py new file mode 100644 index 0000000..e4a3143 --- /dev/null +++ b/app/integrations/dispatch_inbox_consumer.py @@ -0,0 +1,191 @@ +"""Consume QMS §7 asset events from the dispatch service inbox (C2). + +LegacyHUB is subscribed to ``AssetRegistered``/``AssetApproved`` emitted by +QMS for release/evidence objects in the shared MinIO. This consumer polls +``GET /dispatch/inbox/service/{participant}`` (the inbox returns the same +head message until it is confirmed), maps the §7 payload to an +``AssetManifestEnvelope`` and runs the local knowledge-ingest, then confirms +the delivery. + +Outcome policy: +- ``AssetRegistered`` → ingest; accepted/duplicate/rejected all confirm + (rejected is a permanent verdict, retrying cannot change it). +- ``AssetApproved`` / unknown types → confirm without work (no head-of-line + blocking on messages that carry nothing to do). +- Malformed payload → confirm as poison (permanently unmappable). +- Transient infrastructure failure (object storage / DB down) → NO confirm + and stop the run; the next scheduled run retries the same message. + +Secrets (X-API-Key) are never logged. +""" + +from __future__ import annotations + +from collections.abc import Callable +from typing import Any + +import httpx + +from app.config import settings +from app.logging_config import get_logger + +logger = get_logger(__name__) + +MANIFEST_VERSION = "1.0" +CONSUMED_EVENT_TYPES = {"AssetRegistered", "AssetApproved"} + + +class TransientIngestError(Exception): + """Ingest failed for a reason a later retry can fix (storage/DB outage).""" + + +def manifest_envelope_from_event(event: dict[str, Any]) -> dict[str, Any]: + """Map a §7 asset-event payload to an AssetManifestEnvelope dict (§8).""" + ref = event.get("manifest_ref") + if not isinstance(ref, dict): + raise ValueError("event has no manifest_ref") + for required in ("asset_id", "owner_module", "owner_record_type", + "owner_record_id", "content_type", "size_bytes", "sha256"): + if event.get(required) in (None, ""): + raise ValueError(f"event misses required field: {required}") + object_key = str(ref.get("object_key") or "") + filename = object_key.rsplit("/", 1)[-1] or None + return { + "manifest": { + "manifest_version": MANIFEST_VERSION, + "asset": { + "asset_id": str(event["asset_id"]), + "owner_module": str(event["owner_module"]), + "owner_record_type": str(event["owner_record_type"]), + "owner_record_id": str(event["owner_record_id"]), + "asset_kind": str(event.get("asset_kind") or "document"), + "title": None, + "original_filename": filename, + "content_type": str(event["content_type"]), + "size_bytes": int(event["size_bytes"]), + "sha256": str(event["sha256"]), + "created_at": None, + "created_by": event.get("actor"), + }, + "storage": { + "provider": str(ref.get("provider") or "s3"), + "bucket": str(ref.get("bucket") or ""), + "object_key": object_key, + "version_id": None, + }, + "security": { + "gate_status": str((event.get("security") or {}).get("gate_status") or ""), + }, + "retention": { + "policy_id": f"{event['owner_module']}-{event['owner_record_type']}", + "retain_until": None, + "legal_hold": False, + }, + "derivatives": [], + "links": [], + "source": {"source_system": str(event["owner_module"])}, + } + } + + +def _event_from_message(message: dict[str, Any]) -> dict[str, Any]: + envelope = message.get("envelope") or {} + card = envelope.get("card") or {} + event = card.get("event") + return event if isinstance(event, dict) else {} + + +def consume_inbox( + *, + fetch: Callable[[], dict[str, Any] | None], + ingest: Callable[[dict[str, Any]], str], + confirm: Callable[[str], Any], + limit: int = 25, +) -> dict[str, Any]: + """Drain up to ``limit`` messages; injected callables keep this pure.""" + summary: dict[str, Any] = {"ingested": 0, "confirmed": 0, "skipped": 0} + for _ in range(limit): + message = fetch() + if not message: + break + message_id = str(message.get("message_id") or "") + message_type = str(message.get("message_type") or "") + if message_type == "AssetRegistered": + try: + envelope = manifest_envelope_from_event(_event_from_message(message)) + except ValueError as exc: + logger.warning("dispatch_consume.poison", message_id=message_id, error=str(exc)) + confirm(message_id) + summary["confirmed"] += 1 + summary["skipped"] += 1 + continue + try: + status = ingest(envelope) + except TransientIngestError as exc: + logger.warning("dispatch_consume.transient", message_id=message_id, error=str(exc)) + summary["stopped"] = "transient_error" + break + summary["ingested"] += 1 + logger.info("dispatch_consume.ingested", message_id=message_id, status=str(status)) + else: + if message_type not in CONSUMED_EVENT_TYPES: + logger.info("dispatch_consume.unhandled_type", message_id=message_id, + message_type=message_type) + summary["skipped"] += 1 + confirm(message_id) + summary["confirmed"] += 1 + return summary + + +# ---------------- live wiring (HTTP fetch/confirm + local ingest) ---------------- + +def _http_fetch(client: httpx.Client) -> dict[str, Any] | None: + response = client.get( + f"{settings.dispatch_api_url.rstrip('/')}/dispatch/inbox/service/" + f"{settings.dispatch_participant_code}", + headers={"X-API-Key": settings.dispatch_api_key}, + ) + response.raise_for_status() + payload = response.json() + return payload if payload.get("envelope") else None + + +def _http_confirm(client: httpx.Client, message_id: str) -> None: + response = client.post( + f"{settings.dispatch_api_url.rstrip('/')}/dispatch/deliveries/confirm", + headers={"X-API-Key": settings.dispatch_api_key}, + json={"message_id": message_id, + "participant_code": settings.dispatch_participant_code}, + ) + response.raise_for_status() + + +def _local_ingest(envelope: dict[str, Any]) -> str: + from app.api.schemas import AssetManifestEnvelope + from app.ingestion.knowledge_ingest import KnowledgeIngestError, accept_knowledge_ingest + + request = AssetManifestEnvelope.model_validate(envelope) + idempotency_key = ( + f"dispatch-consume-{request.manifest.asset.asset_id}-{request.manifest.asset.sha256}" + ) + try: + response = accept_knowledge_ingest(request, idempotency_key=idempotency_key) + except KnowledgeIngestError as exc: + if exc.reason_code == "object_read_failed": + raise TransientIngestError(exc.message) from exc + logger.warning("dispatch_consume.ingest_rejected", reason_code=exc.reason_code) + return f"rejected:{exc.reason_code}" + return str(response.status) + + +def consume_once(limit: int = 25) -> dict[str, Any]: + if not (settings.dispatch_enabled and settings.dispatch_api_url + and settings.dispatch_api_key): + return {"skipped_run": "dispatch_disabled"} + with httpx.Client(timeout=settings.dispatch_timeout_seconds) as client: + return consume_inbox( + fetch=lambda: _http_fetch(client), + ingest=_local_ingest, + confirm=lambda message_id: _http_confirm(client, message_id), + limit=limit, + ) diff --git a/scripts/consume_dispatch_inbox.py b/scripts/consume_dispatch_inbox.py new file mode 100644 index 0000000..8608051 --- /dev/null +++ b/scripts/consume_dispatch_inbox.py @@ -0,0 +1,24 @@ +"""Drain the dispatch service inbox once (QMS §7 asset events → knowledge-ingest).""" +from __future__ import annotations + +import argparse +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +if str(ROOT) not in sys.path: + sys.path.insert(0, str(ROOT)) + +from app.integrations.dispatch_inbox_consumer import consume_once # noqa: E402 + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--limit", type=int, default=25) + args = parser.parse_args(argv) + print(consume_once(limit=args.limit)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_dispatch_inbox_consumer.py b/tests/test_dispatch_inbox_consumer.py new file mode 100644 index 0000000..37a02d2 --- /dev/null +++ b/tests/test_dispatch_inbox_consumer.py @@ -0,0 +1,160 @@ +"""C2: consume QMS §7 asset events from the dispatch service inbox. + +AssetRegistered → build AssetManifestEnvelope from the event payload and run +knowledge-ingest; accepted/duplicate/rejected → confirm. AssetApproved and +unknown types → confirm (they carry no new work). Transient ingest failures → +NO confirm and stop: the inbox returns the same head message until confirmed, +so continuing would spin; the next cron run retries. +""" + +from __future__ import annotations + +import pytest + +from app.api.schemas import AssetManifestEnvelope +from app.integrations.dispatch_inbox_consumer import ( + TransientIngestError, + consume_inbox, + manifest_envelope_from_event, +) + + +def _event(**overrides): + event = { + "event_type": "AssetRegistered", + "event_version": "1.0", + "owner_module": "qms", + "asset_id": "512d5f53-6edf-58f3-93aa-2a387f38fcc0", + "owner_record_type": "release_candidate", + "owner_record_id": "rel-2026-001", + "asset_kind": "document", + "content_type": "application/pdf", + "size_bytes": 619, + "sha256": "f" * 64, + "security": {"gate_status": "approved"}, + "manifest_ref": { + "provider": "s3", + "bucket": "teamhub-qmshub-releases", + "object_key": "qmshub/2026/07/12/512d5f53-6edf-58f3-93aa-2a387f38fcc0/original/" + + "f" * 64 + ".pdf", + "sha256": "f" * 64, + }, + "actor": "qa@x.com", + } + event.update(overrides) + return event + + +def _message(message_type="AssetRegistered", event=None): + return { + "message_id": "11111111-2222-3333-4444-555555555555", + "message_type": message_type, + "envelope": { + "message_type": message_type, + "participant_code": "qms", + "card": {"event": event if event is not None else _event(event_type=message_type)}, + }, + } + + +# --- payload -> AssetManifestEnvelope --- + +def test_manifest_envelope_from_event_is_valid_against_schema(): + envelope = manifest_envelope_from_event(_event()) + model = AssetManifestEnvelope.model_validate(envelope) + assert model.manifest.asset.asset_id == "512d5f53-6edf-58f3-93aa-2a387f38fcc0" + assert model.manifest.asset.owner_module == "qms" + assert model.manifest.asset.sha256 == "f" * 64 + assert model.manifest.storage.bucket == "teamhub-qmshub-releases" + assert model.manifest.security.gate_status == "approved" + assert model.manifest.retention.policy_id + + +def test_manifest_envelope_rejects_incomplete_event(): + broken = _event() + del broken["manifest_ref"] + with pytest.raises(ValueError): + manifest_envelope_from_event(broken) + + +# --- orchestration --- + +class Recorder: + def __init__(self, messages, ingest_results=None): + self._messages = list(messages) + self.ingested = [] + self.confirmed = [] + self._ingest_results = list(ingest_results or []) + + def fetch(self): + return self._messages.pop(0) if self._messages else None + + def ingest(self, envelope): + self.ingested.append(envelope) + if self._ingest_results: + result = self._ingest_results.pop(0) + if isinstance(result, Exception): + raise result + return result + return "accepted" + + def confirm(self, message_id): + self.confirmed.append(message_id) + + +def test_registered_is_ingested_and_confirmed(): + rec = Recorder([_message("AssetRegistered")]) + summary = consume_inbox(fetch=rec.fetch, ingest=rec.ingest, confirm=rec.confirm) + assert len(rec.ingested) == 1 and len(rec.confirmed) == 1 + assert summary == {"ingested": 1, "confirmed": 1, "skipped": 0} + + +def test_approved_is_confirmed_without_ingest(): + rec = Recorder([_message("AssetApproved")]) + summary = consume_inbox(fetch=rec.fetch, ingest=rec.ingest, confirm=rec.confirm) + assert rec.ingested == [] and len(rec.confirmed) == 1 + assert summary["skipped"] == 1 + + +def test_unknown_type_is_confirmed_to_avoid_head_of_line_blocking(): + rec = Recorder([_message("SomethingElse")]) + consume_inbox(fetch=rec.fetch, ingest=rec.ingest, confirm=rec.confirm) + assert rec.ingested == [] and len(rec.confirmed) == 1 + + +def test_duplicate_and_rejected_still_confirm(): + rec = Recorder( + [_message("AssetRegistered"), _message("AssetRegistered")], + ingest_results=["duplicate", "rejected"], + ) + summary = consume_inbox(fetch=rec.fetch, ingest=rec.ingest, confirm=rec.confirm) + assert len(rec.confirmed) == 2 + assert summary["ingested"] == 2 + + +def test_transient_ingest_error_stops_without_confirm(): + rec = Recorder( + [_message("AssetRegistered"), _message("AssetApproved")], + ingest_results=[TransientIngestError("minio down")], + ) + summary = consume_inbox(fetch=rec.fetch, ingest=rec.ingest, confirm=rec.confirm) + assert rec.confirmed == [] # head message left pending for retry + assert summary["ingested"] == 0 + assert summary.get("stopped") == "transient_error" + + +def test_malformed_registered_event_is_confirmed_as_poison(): + """A permanently unmappable event must not block the queue forever.""" + bad = _message("AssetRegistered", event={"event_type": "AssetRegistered"}) + rec = Recorder([bad]) + summary = consume_inbox(fetch=rec.fetch, ingest=rec.ingest, confirm=rec.confirm) + assert rec.ingested == [] + assert len(rec.confirmed) == 1 + assert summary["skipped"] == 1 + + +def test_empty_inbox_returns_zero_summary(): + rec = Recorder([]) + assert consume_inbox(fetch=rec.fetch, ingest=rec.ingest, confirm=rec.confirm) == { + "ingested": 0, "confirmed": 0, "skipped": 0, + }