Files
yovision/Brain/tests/test_ingress.py
T
QiuSW bd964e8831
Harness governance / validate (pull_request) Has been cancelled
feat: implement T-019 reliable event ingress
2026-08-11 15:41:07 +08:00

292 lines
12 KiB
Python

from __future__ import annotations
import base64
import hashlib
import hmac
import json
import tempfile
import threading
import unittest
from datetime import datetime, timezone
from http.server import BaseHTTPRequestHandler, HTTPServer
from pathlib import Path
from Brain.yovision_brain.domain import EventCandidate
from Brain.yovision_brain.ingress import (
BrainEventIngress,
DeliveryResult,
EventIngressClient,
EventOutbox,
EventRelayWorker,
IngressConfig,
KeyMaterial,
PermanentDelivery,
RetryableDelivery,
load_config,
map_candidate,
)
EVENT_ID = "evt_01J8XQ2K7M3P5R9T0V4W6Y8Z2B"
def candidate(source_event_id: str = "BRN-test-0001", fixture: bool = True) -> EventCandidate:
return EventCandidate(
source_event_id=source_event_id,
kind="zone_entry",
source_ref="must-not-leave-brain",
track_id="P-1",
zone_id="zone-1",
zone_version=3,
occurred_at=datetime(2026, 8, 11, tzinfo=timezone.utc).isoformat().replace("+00:00", "Z"),
confidence=None,
fixture=fixture,
)
def config(root: Path, bell_url: str = "http://127.0.0.1:8081/internal/v1/event-candidates") -> IngressConfig:
return IngressConfig(
producer_id="brain-main",
tenant_id=1,
site_id=2,
device_id=3,
modality="video",
severity="medium",
config_version="brain-demo-v1",
bell_url=bell_url,
key_id="brain-a",
key_file=root / "keys.json",
outbox_path=root / "outbox.sqlite3",
)
class ConfigAndMapperTests(unittest.TestCase):
def test_external_config_binds_key_and_rejects_remote_plaintext(self) -> None:
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
secret = bytes(range(32))
key_file = root / "keys.json"
key_file.write_text(
json.dumps(
{
"version": 1,
"keys": [
{
"key_id": "brain-a",
"producer_id": "brain-main",
"secret_base64url": base64.urlsafe_b64encode(secret).rstrip(b"=").decode("ascii"),
}
],
}
),
encoding="utf-8",
)
config_file = root / "config.json"
document = {
"version": 1,
"producer_id": "brain-main",
"tenant_id": 1,
"site_id": 2,
"device_id": 3,
"modality": "video",
"severity": "medium",
"config_version": "brain-demo-v1",
"bell_url": "http://127.0.0.1:8081/internal/v1/event-candidates",
"key_id": "brain-a",
"key_file": str(key_file.resolve()),
"outbox_path": str((root / "outbox.sqlite3").resolve()),
}
config_file.write_text(json.dumps(document), encoding="utf-8")
loaded, key = load_config(str(config_file.resolve()))
self.assertEqual(("brain-main", secret), (loaded.producer_id, key.secret))
document["bell_url"] = "http://camera.example/internal/v1/event-candidates"
config_file.write_text(json.dumps(document), encoding="utf-8")
with self.assertRaisesRegex(ValueError, "HTTPS"):
load_config(str(config_file.resolve()))
def test_mapper_produces_complete_candidate_without_sensitive_source(self) -> None:
with tempfile.TemporaryDirectory() as directory:
payload = json.loads(map_candidate(candidate(), config(Path(directory))))
self.assertNotIn("id", payload)
self.assertNotIn("source_ref", json.dumps(payload))
self.assertEqual("test", payload["outcome"])
self.assertEqual("auto", payload["outcome_source"])
self.assertIsNone(payload["confidence"])
self.assertEqual(payload["occurred_at"], payload["detected_at"])
self.assertEqual(0.0, payload["latency_seconds"])
self.assertEqual([{"device_id": 3, "modality": "video", "role": "primary"}], payload["sensors"])
expected = {
"schema_version", "source_event_id", "tenant_id", "site_id", "device_id", "sensors", "kind",
"severity", "confidence", "occurred_at", "detected_at", "latency_seconds", "config_version",
"rule", "subject", "observation", "evidence", "dedup_key", "aggregated_into", "outcome",
"outcome_source", "outcome_reason", "diagnostics", "ext",
}
self.assertEqual(expected, set(payload))
class OutboxTests(unittest.TestCase):
def test_crash_recovery_retry_and_terminal_records_are_durable(self) -> None:
clock = [1_000.0]
now = lambda: clock[0]
with tempfile.TemporaryDirectory() as directory:
path = Path(directory) / "outbox.sqlite3"
payload = map_candidate(candidate(), config(Path(directory)))
outbox = EventOutbox(path, now=now)
self.assertTrue(outbox.enqueue(payload))
self.assertFalse(outbox.enqueue(payload))
leased = outbox.lease_one()
self.assertIsNotNone(leased)
outbox.close() # Simulate a process crash while the lease is held.
clock[0] += 31.0
outbox = EventOutbox(path, now=now)
recovered = outbox.lease_one()
self.assertIsNotNone(recovered)
self.assertEqual(2, recovered.attempt_count)
outbox.mark_delivered(recovered, DeliveryResult("duplicate", EVENT_ID))
self.assertEqual(1, outbox.status()["delivered"])
outbox.close()
outbox = EventOutbox(path, now=now)
status = outbox.status()
self.assertEqual((1, 0), (status["delivered"], status["queued"]))
outbox.close()
def test_source_id_conflict_is_not_silently_overwritten(self) -> None:
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
outbox = EventOutbox(root / "outbox.sqlite3")
outbox.enqueue(map_candidate(candidate(), config(root)))
changed = json.loads(map_candidate(candidate(), config(root)))
changed["config_version"] = "different"
with self.assertRaisesRegex(ValueError, "different candidate"):
outbox.enqueue(json.dumps(changed, sort_keys=True, separators=(",", ":")).encode())
outbox.close()
def test_crash_on_final_attempt_becomes_dead_letter(self) -> None:
clock = [1_000.0]
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
path = root / "outbox.sqlite3"
outbox = EventOutbox(path, now=lambda: clock[0])
outbox.enqueue(map_candidate(candidate(), config(root)))
for _ in range(99):
item = outbox.lease_one()
self.assertIsNotNone(item)
outbox.mark_retry(item, "network_error")
clock[0] += 301.0
final = outbox.lease_one()
self.assertEqual(100, final.attempt_count)
outbox.close()
clock[0] += 31.0
outbox = EventOutbox(path, now=lambda: clock[0])
self.assertIsNone(outbox.lease_one())
self.assertEqual(1, outbox.status()["dead_letter"])
outbox.close()
def test_outbox_cannot_be_reused_for_another_producer_identity(self) -> None:
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
path = root / "outbox.sqlite3"
outbox = EventOutbox(path)
original = config(root)
outbox.bind_identity(original)
outbox.close()
outbox = EventOutbox(path)
changed = IngressConfig(**{**original.__dict__, "producer_id": "brain-other"})
with self.assertRaisesRegex(ValueError, "different producer"):
outbox.bind_identity(changed)
outbox.close()
class _BellHandler(BaseHTTPRequestHandler):
secret = bytes(range(32))
verified = False
failure: str | None = None
def log_message(self, _format: str, *_args: object) -> None:
return
def do_POST(self) -> None: # noqa: N802
try:
length = int(self.headers["Content-Length"])
body = self.rfile.read(length)
timestamp = self.headers["X-YoVision-Timestamp"]
nonce = self.headers["X-YoVision-Nonce"]
canonical = "\n".join(("POST", self.path, timestamp, nonce, hashlib.sha256(body).hexdigest())).encode()
encoded_signature = self.headers["X-YoVision-Signature"]
supplied = base64.urlsafe_b64decode(encoded_signature + "=" * (-len(encoded_signature) % 4))
type(self).verified = hmac.compare_digest(supplied, hmac.new(self.secret, canonical, hashlib.sha256).digest())
envelope = json.loads(body)
response = json.dumps(
{
"schema_version": 1,
"producer_id": envelope["producer_id"],
"source_event_id": envelope["candidate"]["source_event_id"],
"event_id": EVENT_ID,
"status": "accepted",
}
).encode()
self.send_response(201)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(response)))
self.end_headers()
self.wfile.write(response)
except Exception as exc: # pragma: no cover - only improves test diagnostics
type(self).failure = repr(exc)
self.send_response(500)
self.end_headers()
class ClientAndWorkerTests(unittest.TestCase):
def test_client_signs_and_accepts_bell_response(self) -> None:
server = HTTPServer(("127.0.0.1", 0), _BellHandler)
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
try:
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
client = EventIngressClient(
config(root, f"http://127.0.0.1:{server.server_port}/internal/v1/event-candidates"),
KeyMaterial("brain-a", "brain-main", _BellHandler.secret),
now=lambda: 1_800_000_000.0,
)
try:
result = client.deliver(map_candidate(candidate(), config(root)))
except RetryableDelivery as exc:
self.fail(f"test Bell handler failed: {_BellHandler.failure}; client={exc.code}")
self.assertEqual(("accepted", EVENT_ID), (result.status, result.event_id))
self.assertTrue(_BellHandler.verified)
finally:
server.shutdown()
server.server_close()
thread.join(timeout=2.0)
def test_worker_distinguishes_retryable_and_permanent_failures(self) -> None:
class FailingClient:
def __init__(self, failure: Exception) -> None:
self.failure = failure
def deliver(self, _payload: bytes) -> DeliveryResult:
raise self.failure
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
payload = map_candidate(candidate(), config(root))
outbox = EventOutbox(root / "retry.sqlite3", now=lambda: 100.0)
outbox.enqueue(payload)
EventRelayWorker(outbox, FailingClient(RetryableDelivery("unauthorized"))).run_once()
self.assertEqual((1, 0), (outbox.status()["queued"], outbox.status()["dead_letter"]))
outbox.close()
outbox = EventOutbox(root / "dead.sqlite3", now=lambda: 100.0)
outbox.enqueue(payload)
EventRelayWorker(outbox, FailingClient(PermanentDelivery("source_event_conflict"))).run_once()
self.assertEqual(1, outbox.status()["dead_letter"])
outbox.close()
if __name__ == "__main__":
unittest.main()