Assistant: codex Assistant-Model: gpt-6-astra Assistant-Session: 01a07ff8-19d0-7820-b4d0-1353833cb7fc
108 lines
5.1 KiB
Python
108 lines
5.1 KiB
Python
"""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}
|