Files
LegacyHUB/app/integrations/dispatch_inbox_consumer.py
Vadim Malanov 272cbf6d92
Some checks failed
CI / Backend (lint + tests + compose) (push) Has been cancelled
CI / Frontend (lint + type-check + build) (push) Has been cancelled
feat(dispatch): consume QMS asset events from the service inbox (C2)
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 <noreply@anthropic.com>
2026-07-12 12:39:44 +03:00

192 lines
7.5 KiB
Python

"""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,
)