"""Deliver immutable outbox bytes to Audit Core with its idempotency contract.""" from urllib.parse import urlencode from .http_transport import JSONTransport, TransportError, fixed_origin from .store import bounded_window class AuditDeliveryError(RuntimeError): def __init__(self, code, permanent=False): super().__init__(code) self.code = code self.permanent = permanent class AuditCoreSink: def __init__(self, origin, token_provider, *, transport=None, allow_internal_http=False): self.origin = fixed_origin(origin, allow_internal_http=allow_internal_http) self.token_provider = token_provider self.transport = transport or JSONTransport(allow_internal_http=allow_internal_http) def _request(self, method, path, *, body=None, event_id=None): try: token = self.token_provider() except (OSError, ValueError): raise AuditDeliveryError("unauthorized", permanent=True) from None if not isinstance(token, str) or not token or not token.isascii() or len(token) > 8192 or any(c.isspace() for c in token): raise AuditDeliveryError("unauthorized", permanent=True) headers = {"Authorization": "Bearer " + token, "Content-Type": "application/json"} if event_id: headers["Idempotency-Key"] = event_id try: status, data = self.transport.request(method, self.origin + path, headers=headers, body=body) except TransportError: raise AuditDeliveryError("unavailable") from None if status in (401, 403): raise AuditDeliveryError("unauthorized", permanent=True) if status == 400: raise AuditDeliveryError("rejected", permanent=True) if status == 409: raise AuditDeliveryError("conflict", permanent=True) if status >= 500: raise AuditDeliveryError("unavailable") return status, data def deliver(self, event_id, envelope_json): status, data = self._request("POST", "/v1/events", body=envelope_json.encode(), event_id=event_id) if ((status, data.get("status")) not in ((202, "accepted"), (200, "duplicate")) or data.get("reference") != "audit:" + event_id): raise AuditDeliveryError("invalid_receipt") return data["reference"] def counts(self, since, until): since, until = bounded_window(since, until) query = {"source": "informed-decision", "tenant": "tenant:platform", "since": since, "until": until} status, data = self._request("GET", "/v1/reconciliation?" + urlencode(query)) try: same_window = bounded_window(data.get("since"), data.get("until")) == (since, until) except (TypeError, ValueError): same_window = False rows = data.get("counts") if (status != 200 or data.get("source") != query["source"] or data.get("tenant") != query["tenant"] or not same_window or not isinstance(rows, list)): raise AuditDeliveryError("invalid_receipt") counts = {} for row in rows: if (not isinstance(row, dict) or not isinstance(row.get("class"), str) or not row["class"] or row["class"] in counts or type(row.get("count")) is not int or row["count"] < 0): raise AuditDeliveryError("invalid_receipt") counts[row["class"]] = row["count"] return counts class OutboxWorker: def __init__(self, store, sink): self.store = store self.sink = sink def run_once(self, *, limit=100): if type(limit) is not int or not 1 <= limit <= 1000: raise ValueError("delivery batch limit must be 1..1000") result = {"delivered": 0, "retrying": 0, "blocked": 0} for _ in range(limit): delivery = self.store.claim_delivery() if delivery is None: break event_id, lease, body = delivery try: reference = self.sink.deliver(event_id, body) except AuditDeliveryError as error: self.store.finish_delivery(event_id, lease, error=error.code, permanent=error.permanent) result["blocked" if error.permanent else "retrying"] += 1 else: self.store.finish_delivery(event_id, lease, reference=reference) result["delivered"] += 1 return result def reconcile(self, since, until): since, until = bounded_window(since, until) local = self.store.counts_by_class(since, until) remote = self.sink.counts(since, until) classes = sorted(set(local) | set(remote)) return {"source": "informed-decision", "tenant": "tenant:platform", "since": since, "until": until, "counts": {k: {"source": local.get(k, 0), "receiver": remote.get(k, 0)} for k in classes}, "count_values_match": all(local.get(k, 0) == remote.get(k, 0) for k in classes), "source_time_basis": "occurred_at", "receiver_time_basis": "accepted_at", "automatic_loss_finding": False, "completeness_proven": False, "reconstructability_proven": False}