646 lines
26 KiB
Python
646 lines
26 KiB
Python
from __future__ import annotations
|
|
|
|
import base64
|
|
import binascii
|
|
import hashlib
|
|
import hmac
|
|
import ipaddress
|
|
import json
|
|
import re
|
|
import secrets
|
|
import sqlite3
|
|
import threading
|
|
import time
|
|
from dataclasses import dataclass
|
|
from pathlib import Path
|
|
from typing import Callable
|
|
from urllib.error import HTTPError, URLError
|
|
from urllib.parse import urlsplit
|
|
from urllib.request import HTTPRedirectHandler, ProxyHandler, Request, build_opener
|
|
|
|
from .domain import EventCandidate
|
|
|
|
|
|
INGRESS_PATH = "/internal/v1/event-candidates"
|
|
MAX_PAYLOAD_BYTES = 1 << 20
|
|
MAX_QUEUED = 10_000
|
|
MAX_ATTEMPTS = 100
|
|
HTTP_TIMEOUT_SECONDS = 10.0
|
|
_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$")
|
|
_SOURCE_ID = re.compile(r"^[A-Za-z0-9_-]{1,128}$")
|
|
_EVENT_ID = re.compile(r"^evt_[0-9A-HJKMNP-TV-Z]{26}$")
|
|
_ERROR_CODE = re.compile(r"^[a-z][a-z0-9_]{0,63}$")
|
|
_REPO_ROOT = Path(__file__).resolve().parents[2]
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class IngressConfig:
|
|
producer_id: str
|
|
tenant_id: int
|
|
site_id: int
|
|
device_id: int
|
|
modality: str
|
|
severity: str
|
|
config_version: str
|
|
bell_url: str
|
|
key_id: str
|
|
key_file: Path
|
|
outbox_path: Path
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class KeyMaterial:
|
|
key_id: str
|
|
producer_id: str
|
|
secret: bytes
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class LeasedEvent:
|
|
source_event_id: str
|
|
payload: bytes
|
|
attempt_count: int
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class DeliveryResult:
|
|
status: str
|
|
event_id: str
|
|
|
|
|
|
class RetryableDelivery(RuntimeError):
|
|
def __init__(self, code: str) -> None:
|
|
super().__init__(code)
|
|
self.code = _safe_error(code)
|
|
|
|
|
|
class PermanentDelivery(RuntimeError):
|
|
def __init__(self, code: str) -> None:
|
|
super().__init__(code)
|
|
self.code = _safe_error(code)
|
|
|
|
|
|
class _NoRedirect(HTTPRedirectHandler):
|
|
def redirect_request(self, _request: Request, _file_pointer: object, _code: int, _message: str, _headers: object, _new_url: str) -> None:
|
|
return None
|
|
|
|
|
|
def _is_within(path: Path, parent: Path) -> bool:
|
|
try:
|
|
path.relative_to(parent)
|
|
return True
|
|
except ValueError:
|
|
return False
|
|
|
|
|
|
def _external_absolute_path(raw: object, field: str, *, must_exist: bool) -> Path:
|
|
if not isinstance(raw, str) or not raw:
|
|
raise ValueError(f"{field} must be an absolute external path")
|
|
path = Path(raw)
|
|
if not path.is_absolute():
|
|
raise ValueError(f"{field} must be an absolute external path")
|
|
resolved = path.resolve(strict=must_exist)
|
|
if _is_within(resolved, _REPO_ROOT):
|
|
raise ValueError(f"{field} must stay outside the repository")
|
|
return resolved
|
|
|
|
|
|
def _load_json_file(path: Path, maximum: int = 64 << 10) -> object:
|
|
raw = path.read_bytes()
|
|
if len(raw) > maximum:
|
|
raise ValueError("external configuration file is too large")
|
|
try:
|
|
return json.loads(raw.decode("utf-8"))
|
|
except (UnicodeError, json.JSONDecodeError) as exc:
|
|
raise ValueError("external configuration is not valid JSON") from exc
|
|
|
|
|
|
def _positive_int(value: object, field: str) -> int:
|
|
if isinstance(value, bool) or not isinstance(value, int) or value < 1:
|
|
raise ValueError(f"{field} must be a positive integer")
|
|
return value
|
|
|
|
|
|
def _validate_bell_url(raw: object) -> str:
|
|
if not isinstance(raw, str):
|
|
raise ValueError("bell_url must be a URL")
|
|
parsed = urlsplit(raw)
|
|
if parsed.scheme not in {"http", "https"} or parsed.hostname is None or parsed.path != INGRESS_PATH:
|
|
raise ValueError("bell_url must target the versioned event ingress path")
|
|
if parsed.username is not None or parsed.password is not None or parsed.query or parsed.fragment:
|
|
raise ValueError("bell_url cannot contain credentials, query or fragment")
|
|
try:
|
|
port = parsed.port
|
|
except ValueError as exc:
|
|
raise ValueError("bell_url has an invalid port") from exc
|
|
if port is not None and not 1 <= port <= 65535:
|
|
raise ValueError("bell_url has an invalid port")
|
|
if parsed.scheme == "http":
|
|
host = parsed.hostname
|
|
loopback = host == "localhost"
|
|
if not loopback:
|
|
try:
|
|
loopback = ipaddress.ip_address(host).is_loopback
|
|
except ValueError:
|
|
loopback = False
|
|
if not loopback:
|
|
raise ValueError("non-loopback event ingress requires HTTPS")
|
|
return raw
|
|
|
|
|
|
def load_config(path_value: str) -> tuple[IngressConfig, KeyMaterial]:
|
|
config_path = _external_absolute_path(path_value, "event ingress config", must_exist=True)
|
|
value = _load_json_file(config_path)
|
|
expected = {
|
|
"version",
|
|
"producer_id",
|
|
"tenant_id",
|
|
"site_id",
|
|
"device_id",
|
|
"modality",
|
|
"severity",
|
|
"config_version",
|
|
"bell_url",
|
|
"key_id",
|
|
"key_file",
|
|
"outbox_path",
|
|
}
|
|
if not isinstance(value, dict) or set(value) != expected or value.get("version") != 1:
|
|
raise ValueError("invalid event ingress configuration shape")
|
|
producer_id = value["producer_id"]
|
|
key_id = value["key_id"]
|
|
config_version = value["config_version"]
|
|
if not isinstance(producer_id, str) or not _ID.fullmatch(producer_id):
|
|
raise ValueError("invalid producer_id")
|
|
if not isinstance(key_id, str) or not _ID.fullmatch(key_id):
|
|
raise ValueError("invalid key_id")
|
|
if not isinstance(config_version, str) or not 1 <= len(config_version) <= 128:
|
|
raise ValueError("config_version must contain 1 to 128 characters")
|
|
if value["modality"] != "video":
|
|
raise ValueError("T-019 Brain ingress supports only video primary sensors")
|
|
if value["severity"] not in {"low", "medium", "high", "critical"}:
|
|
raise ValueError("invalid severity")
|
|
key_file = _external_absolute_path(value["key_file"], "key_file", must_exist=True)
|
|
outbox_path = _external_absolute_path(value["outbox_path"], "outbox_path", must_exist=False)
|
|
if not outbox_path.parent.is_dir():
|
|
raise ValueError("outbox_path parent directory must already exist")
|
|
config = IngressConfig(
|
|
producer_id=producer_id,
|
|
tenant_id=_positive_int(value["tenant_id"], "tenant_id"),
|
|
site_id=_positive_int(value["site_id"], "site_id"),
|
|
device_id=_positive_int(value["device_id"], "device_id"),
|
|
modality="video",
|
|
severity=value["severity"],
|
|
config_version=config_version,
|
|
bell_url=_validate_bell_url(value["bell_url"]),
|
|
key_id=key_id,
|
|
key_file=key_file,
|
|
outbox_path=outbox_path,
|
|
)
|
|
document = _load_json_file(key_file)
|
|
if not isinstance(document, dict) or set(document) != {"version", "keys"} or document.get("version") != 1:
|
|
raise ValueError("invalid event ingress key document")
|
|
keys = document.get("keys")
|
|
if not isinstance(keys, list) or not keys:
|
|
raise ValueError("event ingress key document has no keys")
|
|
selected: KeyMaterial | None = None
|
|
seen: set[str] = set()
|
|
for item in keys:
|
|
if not isinstance(item, dict) or set(item) != {"key_id", "producer_id", "secret_base64url"}:
|
|
raise ValueError("invalid event ingress key entry")
|
|
item_key = item["key_id"]
|
|
item_producer = item["producer_id"]
|
|
encoded = item["secret_base64url"]
|
|
if not isinstance(item_key, str) or not _ID.fullmatch(item_key) or item_key in seen:
|
|
raise ValueError("invalid or duplicate event ingress key ID")
|
|
if not isinstance(item_producer, str) or not _ID.fullmatch(item_producer) or not isinstance(encoded, str):
|
|
raise ValueError("invalid event ingress key entry")
|
|
seen.add(item_key)
|
|
try:
|
|
secret = base64.b64decode(encoded + "=" * (-len(encoded) % 4), altchars=b"-_", validate=True)
|
|
except (ValueError, binascii.Error) as exc:
|
|
raise ValueError("invalid event ingress secret") from exc
|
|
if len(secret) < 32:
|
|
raise ValueError("event ingress secret must contain at least 32 bytes")
|
|
if item_key == key_id:
|
|
selected = KeyMaterial(item_key, item_producer, secret)
|
|
if selected is None or selected.producer_id != producer_id:
|
|
raise ValueError("selected key is not bound to the configured producer")
|
|
return config, selected
|
|
|
|
|
|
def map_candidate(value: EventCandidate, config: IngressConfig) -> bytes:
|
|
if not _SOURCE_ID.fullmatch(value.source_event_id):
|
|
raise ValueError("invalid source_event_id")
|
|
payload = {
|
|
"schema_version": "0.1",
|
|
"source_event_id": value.source_event_id,
|
|
"tenant_id": config.tenant_id,
|
|
"site_id": config.site_id,
|
|
"device_id": config.device_id,
|
|
"sensors": [{"device_id": config.device_id, "modality": config.modality, "role": "primary"}],
|
|
"kind": value.kind,
|
|
"severity": config.severity,
|
|
"confidence": None,
|
|
"occurred_at": value.occurred_at,
|
|
"detected_at": value.occurred_at,
|
|
"latency_seconds": 0.0,
|
|
"config_version": config.config_version,
|
|
"rule": None,
|
|
"subject": {
|
|
"class": "person",
|
|
"track_id": value.track_id,
|
|
"attributes": {},
|
|
"anon_id": None,
|
|
"identity": None,
|
|
"identity_status": "not_enabled",
|
|
},
|
|
"observation": {
|
|
"zone": value.zone_id,
|
|
"dwell_sec": 0.0,
|
|
"bbox_seq_uri": None,
|
|
"keypoint_seq_uri": None,
|
|
"signal_seq_uri": None,
|
|
},
|
|
"evidence": {"snapshot_uris": [], "clip_uri": None, "clip_range": None},
|
|
"dedup_key": None,
|
|
"aggregated_into": None,
|
|
"outcome": "test" if value.fixture else "unknown",
|
|
"outcome_source": "auto" if value.fixture else None,
|
|
"outcome_reason": None,
|
|
"diagnostics": None,
|
|
"ext": {"brain_candidate_version": "brain-demo-v1", "zone_version": value.zone_version},
|
|
}
|
|
encoded = json.dumps(payload, ensure_ascii=False, sort_keys=True, separators=(",", ":")).encode("utf-8")
|
|
if len(encoded) > MAX_PAYLOAD_BYTES:
|
|
raise ValueError("event candidate exceeds 1 MiB")
|
|
return encoded
|
|
|
|
|
|
class EventOutbox:
|
|
def __init__(self, path: Path, *, now: Callable[[], float] = time.time) -> None:
|
|
if not path.is_absolute() or _is_within(path.resolve(), _REPO_ROOT):
|
|
raise ValueError("outbox path must be absolute and external")
|
|
self._now = now
|
|
self._lock = threading.RLock()
|
|
self._connection = sqlite3.connect(str(path), timeout=5.0, isolation_level=None, check_same_thread=False)
|
|
self._connection.row_factory = sqlite3.Row
|
|
self._connection.execute("PRAGMA journal_mode=WAL")
|
|
self._connection.execute("PRAGMA synchronous=FULL")
|
|
self._connection.execute("PRAGMA foreign_keys=ON")
|
|
self._connection.executescript(
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS event_outbox (
|
|
source_event_id TEXT PRIMARY KEY,
|
|
candidate_hash BLOB NOT NULL,
|
|
payload BLOB NOT NULL,
|
|
state TEXT NOT NULL CHECK(state IN ('queued','delivering','delivered','dead_letter')),
|
|
attempt_count INTEGER NOT NULL DEFAULT 0 CHECK(attempt_count BETWEEN 0 AND 100),
|
|
available_at REAL NOT NULL,
|
|
lease_until REAL,
|
|
bell_event_id TEXT,
|
|
last_error_code TEXT,
|
|
created_at REAL NOT NULL,
|
|
updated_at REAL NOT NULL,
|
|
delivered_at REAL,
|
|
dead_lettered_at REAL,
|
|
CHECK(length(candidate_hash)=32),
|
|
CHECK(length(payload)<=1048576),
|
|
CHECK((state='delivering')=(lease_until IS NOT NULL)),
|
|
CHECK(delivered_at IS NULL OR dead_lettered_at IS NULL)
|
|
);
|
|
CREATE INDEX IF NOT EXISTS event_outbox_due_idx
|
|
ON event_outbox(available_at, created_at, source_event_id)
|
|
WHERE state='queued';
|
|
CREATE TABLE IF NOT EXISTS event_outbox_identity (
|
|
singleton INTEGER PRIMARY KEY CHECK(singleton=1),
|
|
producer_id TEXT NOT NULL,
|
|
tenant_id INTEGER NOT NULL CHECK(tenant_id>=1),
|
|
site_id INTEGER NOT NULL CHECK(site_id>=1),
|
|
device_id INTEGER NOT NULL CHECK(device_id>=1)
|
|
);
|
|
"""
|
|
)
|
|
|
|
def bind_identity(self, config: IngressConfig) -> None:
|
|
identity = (config.producer_id, config.tenant_id, config.site_id, config.device_id)
|
|
with self._lock:
|
|
self._connection.execute("BEGIN IMMEDIATE")
|
|
try:
|
|
row = self._connection.execute(
|
|
"SELECT producer_id,tenant_id,site_id,device_id FROM event_outbox_identity WHERE singleton=1"
|
|
).fetchone()
|
|
if row is None:
|
|
self._connection.execute(
|
|
"""INSERT INTO event_outbox_identity(
|
|
singleton,producer_id,tenant_id,site_id,device_id
|
|
) VALUES (1,?,?,?,?)""",
|
|
identity,
|
|
)
|
|
elif tuple(row) != identity:
|
|
raise ValueError("event outbox is bound to a different producer or device identity")
|
|
self._connection.execute("COMMIT")
|
|
except Exception:
|
|
self._connection.execute("ROLLBACK")
|
|
raise
|
|
|
|
def enqueue(self, payload: bytes) -> bool:
|
|
if not 0 < len(payload) <= MAX_PAYLOAD_BYTES:
|
|
raise ValueError("event candidate payload must contain at most 1 MiB")
|
|
try:
|
|
candidate = json.loads(payload)
|
|
source_event_id = candidate["source_event_id"]
|
|
except (UnicodeError, json.JSONDecodeError, KeyError, TypeError) as exc:
|
|
raise ValueError("invalid event candidate payload") from exc
|
|
if not isinstance(source_event_id, str) or not _SOURCE_ID.fullmatch(source_event_id):
|
|
raise ValueError("invalid source_event_id")
|
|
digest = hashlib.sha256(payload).digest()
|
|
now = self._now()
|
|
with self._lock:
|
|
self._connection.execute("BEGIN IMMEDIATE")
|
|
try:
|
|
existing = self._connection.execute(
|
|
"SELECT candidate_hash FROM event_outbox WHERE source_event_id=?", (source_event_id,)
|
|
).fetchone()
|
|
if existing is not None:
|
|
if not hmac.compare_digest(existing["candidate_hash"], digest):
|
|
raise ValueError("source_event_id already has a different candidate")
|
|
self._connection.execute("COMMIT")
|
|
return False
|
|
active = self._connection.execute(
|
|
"SELECT count(*) FROM event_outbox WHERE state IN ('queued','delivering')"
|
|
).fetchone()[0]
|
|
if active >= MAX_QUEUED:
|
|
raise RuntimeError("event outbox capacity exceeded")
|
|
self._connection.execute(
|
|
"""INSERT INTO event_outbox(
|
|
source_event_id,candidate_hash,payload,state,available_at,created_at,updated_at
|
|
) VALUES (?,?,?,'queued',?,?,?)""",
|
|
(source_event_id, digest, payload, now, now, now),
|
|
)
|
|
self._connection.execute("COMMIT")
|
|
return True
|
|
except Exception:
|
|
self._connection.execute("ROLLBACK")
|
|
raise
|
|
|
|
def lease_one(self, lease_seconds: float = 30.0) -> LeasedEvent | None:
|
|
now = self._now()
|
|
with self._lock:
|
|
self._connection.execute("BEGIN IMMEDIATE")
|
|
try:
|
|
self._connection.execute(
|
|
"""UPDATE event_outbox SET state='dead_letter', lease_until=NULL,
|
|
last_error_code='retry_exhausted', dead_lettered_at=?, updated_at=?
|
|
WHERE attempt_count>=? AND (
|
|
state='queued' OR (state='delivering' AND lease_until<=?)
|
|
)""",
|
|
(now, now, MAX_ATTEMPTS, now),
|
|
)
|
|
self._connection.execute(
|
|
"""UPDATE event_outbox SET state='queued', lease_until=NULL,
|
|
available_at=?, last_error_code='lease_expired', updated_at=?
|
|
WHERE state='delivering' AND lease_until<=? AND attempt_count<?""",
|
|
(now, now, now, MAX_ATTEMPTS),
|
|
)
|
|
row = self._connection.execute(
|
|
"""SELECT source_event_id,payload,attempt_count FROM event_outbox
|
|
WHERE state='queued' AND available_at<=? AND attempt_count<?
|
|
ORDER BY available_at,created_at,source_event_id LIMIT 1""",
|
|
(now, MAX_ATTEMPTS),
|
|
).fetchone()
|
|
if row is None:
|
|
self._connection.execute("COMMIT")
|
|
return None
|
|
attempt = row["attempt_count"] + 1
|
|
self._connection.execute(
|
|
"""UPDATE event_outbox SET state='delivering',attempt_count=?,lease_until=?,updated_at=?
|
|
WHERE source_event_id=? AND state='queued'""",
|
|
(attempt, now + lease_seconds, now, row["source_event_id"]),
|
|
)
|
|
self._connection.execute("COMMIT")
|
|
return LeasedEvent(row["source_event_id"], bytes(row["payload"]), attempt)
|
|
except Exception:
|
|
self._connection.execute("ROLLBACK")
|
|
raise
|
|
|
|
def mark_delivered(self, item: LeasedEvent, result: DeliveryResult) -> None:
|
|
now = self._now()
|
|
with self._lock:
|
|
changed = self._connection.execute(
|
|
"""UPDATE event_outbox SET state='delivered',lease_until=NULL,bell_event_id=?,
|
|
last_error_code=NULL,delivered_at=?,updated_at=?
|
|
WHERE source_event_id=? AND state='delivering' AND attempt_count=?""",
|
|
(result.event_id, now, now, item.source_event_id, item.attempt_count),
|
|
).rowcount
|
|
if changed != 1:
|
|
raise RuntimeError("event outbox delivery lease was lost")
|
|
|
|
def mark_retry(self, item: LeasedEvent, code: str) -> None:
|
|
now = self._now()
|
|
if item.attempt_count >= MAX_ATTEMPTS:
|
|
self.mark_dead(item, "retry_exhausted")
|
|
return
|
|
delay = min(300.0, float(2 ** min(item.attempt_count - 1, 9)))
|
|
with self._lock:
|
|
changed = self._connection.execute(
|
|
"""UPDATE event_outbox SET state='queued',lease_until=NULL,available_at=?,
|
|
last_error_code=?,updated_at=?
|
|
WHERE source_event_id=? AND state='delivering' AND attempt_count=?""",
|
|
(now + delay, _safe_error(code), now, item.source_event_id, item.attempt_count),
|
|
).rowcount
|
|
if changed != 1:
|
|
raise RuntimeError("event outbox retry lease was lost")
|
|
|
|
def mark_dead(self, item: LeasedEvent, code: str) -> None:
|
|
now = self._now()
|
|
with self._lock:
|
|
changed = self._connection.execute(
|
|
"""UPDATE event_outbox SET state='dead_letter',lease_until=NULL,
|
|
last_error_code=?,dead_lettered_at=?,updated_at=?
|
|
WHERE source_event_id=? AND state='delivering' AND attempt_count=?""",
|
|
(_safe_error(code), now, now, item.source_event_id, item.attempt_count),
|
|
).rowcount
|
|
if changed != 1:
|
|
raise RuntimeError("event outbox dead-letter lease was lost")
|
|
|
|
def status(self) -> dict[str, object]:
|
|
with self._lock:
|
|
counts = {row["state"]: row["count"] for row in self._connection.execute(
|
|
"SELECT state,count(*) AS count FROM event_outbox GROUP BY state"
|
|
)}
|
|
error = self._connection.execute(
|
|
"""SELECT last_error_code FROM event_outbox WHERE last_error_code IS NOT NULL
|
|
ORDER BY updated_at DESC LIMIT 1"""
|
|
).fetchone()
|
|
return {
|
|
"enabled": True,
|
|
"queued": counts.get("queued", 0),
|
|
"delivering": counts.get("delivering", 0),
|
|
"delivered": counts.get("delivered", 0),
|
|
"dead_letter": counts.get("dead_letter", 0),
|
|
"last_error_code": None if error is None else error["last_error_code"],
|
|
}
|
|
|
|
def close(self) -> None:
|
|
with self._lock:
|
|
self._connection.close()
|
|
|
|
|
|
class EventIngressClient:
|
|
def __init__(self, config: IngressConfig, key: KeyMaterial, *, now: Callable[[], float] = time.time) -> None:
|
|
self._config = config
|
|
self._key = key
|
|
self._now = now
|
|
# Internal event payloads must never be redirected through ambient
|
|
# HTTP(S)_PROXY settings.
|
|
self._opener = build_opener(ProxyHandler({}), _NoRedirect())
|
|
|
|
def deliver(self, candidate: bytes) -> DeliveryResult:
|
|
try:
|
|
candidate_object = json.loads(candidate)
|
|
except (UnicodeError, json.JSONDecodeError) as exc:
|
|
raise PermanentDelivery("candidate_invalid") from exc
|
|
body = json.dumps(
|
|
{"schema_version": 1, "producer_id": self._config.producer_id, "candidate": candidate_object},
|
|
ensure_ascii=False,
|
|
sort_keys=True,
|
|
separators=(",", ":"),
|
|
).encode("utf-8")
|
|
if len(body) > MAX_PAYLOAD_BYTES:
|
|
raise PermanentDelivery("payload_too_large")
|
|
timestamp = str(int(self._now()))
|
|
nonce = base64.urlsafe_b64encode(secrets.token_bytes(16)).rstrip(b"=").decode("ascii")
|
|
digest = hashlib.sha256(body).hexdigest()
|
|
canonical = "\n".join(("POST", INGRESS_PATH, timestamp, nonce, digest)).encode("utf-8")
|
|
signature_value = base64.urlsafe_b64encode(hmac.new(self._key.secret, canonical, hashlib.sha256).digest()).rstrip(b"=").decode("ascii")
|
|
request = Request(self._config.bell_url, data=body, method="POST")
|
|
request.add_header("Content-Type", "application/json")
|
|
request.add_header("X-YoVision-Key-Id", self._key.key_id)
|
|
request.add_header("X-YoVision-Timestamp", timestamp)
|
|
request.add_header("X-YoVision-Nonce", nonce)
|
|
request.add_header("X-YoVision-Signature", signature_value)
|
|
try:
|
|
with self._opener.open(request, timeout=HTTP_TIMEOUT_SECONDS) as response:
|
|
status = response.status
|
|
response_body = response.read(64 << 10)
|
|
except HTTPError as exc:
|
|
status = exc.code
|
|
response_body = exc.read(64 << 10)
|
|
except (URLError, TimeoutError, OSError) as exc:
|
|
raise RetryableDelivery("network_error") from exc
|
|
if status in {200, 201}:
|
|
try:
|
|
value = json.loads(response_body)
|
|
response_status = value["status"]
|
|
event_id = value["event_id"]
|
|
if (
|
|
value["schema_version"] != 1
|
|
or value["producer_id"] != self._config.producer_id
|
|
or value["source_event_id"] != candidate_object["source_event_id"]
|
|
or response_status not in {"accepted", "duplicate"}
|
|
or not isinstance(event_id, str)
|
|
or not _EVENT_ID.fullmatch(event_id)
|
|
):
|
|
raise ValueError
|
|
except (KeyError, TypeError, ValueError, json.JSONDecodeError) as exc:
|
|
raise RetryableDelivery("invalid_response") from exc
|
|
return DeliveryResult(response_status, event_id)
|
|
code = "http_error"
|
|
try:
|
|
error_value = json.loads(response_body)
|
|
if isinstance(error_value, dict) and isinstance(error_value.get("error"), str):
|
|
code = error_value["error"]
|
|
except (UnicodeError, json.JSONDecodeError):
|
|
pass
|
|
if status == 401 or status >= 500:
|
|
raise RetryableDelivery(code)
|
|
raise PermanentDelivery(code)
|
|
|
|
|
|
class EventRelayWorker:
|
|
def __init__(self, outbox: EventOutbox, client: EventIngressClient) -> None:
|
|
self._outbox = outbox
|
|
self._client = client
|
|
self._stop = threading.Event()
|
|
self._thread: threading.Thread | None = None
|
|
self._last_error_code: str | None = None
|
|
|
|
def run_once(self) -> bool:
|
|
item = self._outbox.lease_one()
|
|
if item is None:
|
|
return False
|
|
try:
|
|
result = self._client.deliver(item.payload)
|
|
except PermanentDelivery as exc:
|
|
self._outbox.mark_dead(item, exc.code)
|
|
except RetryableDelivery as exc:
|
|
self._outbox.mark_retry(item, exc.code)
|
|
except Exception:
|
|
self._outbox.mark_retry(item, "client_error")
|
|
else:
|
|
self._outbox.mark_delivered(item, result)
|
|
return True
|
|
|
|
def start(self) -> None:
|
|
if self._thread is not None:
|
|
return
|
|
self._thread = threading.Thread(target=self._run, name="brain-event-relay", daemon=True)
|
|
self._thread.start()
|
|
|
|
def _run(self) -> None:
|
|
while not self._stop.is_set():
|
|
try:
|
|
worked = self.run_once()
|
|
self._last_error_code = None
|
|
except Exception:
|
|
self._last_error_code = "event_outbox_unavailable"
|
|
self._stop.wait(1.0)
|
|
continue
|
|
if not worked:
|
|
self._stop.wait(0.25)
|
|
|
|
def stop(self) -> None:
|
|
self._stop.set()
|
|
if self._thread is not None:
|
|
self._thread.join(timeout=HTTP_TIMEOUT_SECONDS + 2.0)
|
|
self._thread = None
|
|
|
|
def last_error_code(self) -> str | None:
|
|
return self._last_error_code
|
|
|
|
|
|
class BrainEventIngress:
|
|
def __init__(self, config: IngressConfig, key: KeyMaterial) -> None:
|
|
self._config = config
|
|
self._outbox = EventOutbox(config.outbox_path)
|
|
self._outbox.bind_identity(config)
|
|
self._worker = EventRelayWorker(self._outbox, EventIngressClient(config, key))
|
|
|
|
@classmethod
|
|
def from_file(cls, path: str) -> "BrainEventIngress":
|
|
config, key = load_config(path)
|
|
return cls(config, key)
|
|
|
|
def submit(self, value: EventCandidate) -> None:
|
|
self._outbox.enqueue(map_candidate(value, self._config))
|
|
|
|
def start(self) -> None:
|
|
self._worker.start()
|
|
|
|
def stop(self) -> None:
|
|
self._worker.stop()
|
|
self._outbox.close()
|
|
|
|
def status(self) -> dict[str, object]:
|
|
value = self._outbox.status()
|
|
if self._worker.last_error_code() is not None:
|
|
value["last_error_code"] = self._worker.last_error_code()
|
|
return value
|
|
|
|
|
|
def _safe_error(value: str) -> str:
|
|
if not isinstance(value, str) or not _ERROR_CODE.fullmatch(value):
|
|
return "unknown_error"
|
|
return value
|