292 lines
12 KiB
Python
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()
|