diff --git a/Brain/src/yovision_brain/integration/sense_control/__init__.py b/Brain/src/yovision_brain/integration/sense_control/__init__.py new file mode 100644 index 0000000..060b161 --- /dev/null +++ b/Brain/src/yovision_brain/integration/sense_control/__init__.py @@ -0,0 +1,7 @@ +"""Credential-free Sense control-plane connector for Brain.""" + +from .consumer import ApplyResult, SourceConfigConsumer, SourceConfigError +from .replay import SQLiteReplayStore +from .status import RuntimeStatusPublisher + +__all__ = ["ApplyResult", "SourceConfigConsumer", "SourceConfigError", "SQLiteReplayStore", "RuntimeStatusPublisher"] diff --git a/Brain/src/yovision_brain/integration/sense_control/connector.py b/Brain/src/yovision_brain/integration/sense_control/connector.py new file mode 100644 index 0000000..afb7cf6 --- /dev/null +++ b/Brain/src/yovision_brain/integration/sense_control/connector.py @@ -0,0 +1,57 @@ +"""Machine-authenticated adapters with bounded timeout/backoff and a kill switch.""" + +from __future__ import annotations + +import time +import re +from dataclasses import dataclass +from typing import Callable + +from yovision_brain.integration.machine_identity.token import Signer, Verifier, bearer_token + +from .consumer import ApplyResult, SourceConfigConsumer + +_REQUEST_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:-]{15,127}$") + + +@dataclass(frozen=True, slots=True) +class ConnectorResponse: + status: int + body: bytes + correlation_id: str + + +class SourceConfigEndpoint: + def __init__(self, consumer: SourceConfigConsumer, verifier: Verifier, *, enabled: bool = True, max_body_bytes: int = 10 * 1024 * 1024) -> None: + self._consumer, self._verifier, self._enabled, self._max = consumer, verifier, enabled, max_body_bytes + + def receive(self, authorization: str, body: bytes, correlation_id: str) -> ApplyResult: + if not self._enabled: raise RuntimeError("CONNECTOR_DISABLED") + if not _REQUEST_ID.fullmatch(correlation_id): raise ValueError("INVALID_CORRELATION_ID") + if len(body) > self._max: raise ValueError("REQUEST_TOO_LARGE") + token=bearer_token(authorization) + self._verifier.verify(token,"yovision-brain","source-config:write","POST","/machine/v1/source-config",body) + return self._consumer.apply(body) + + +class StatusSender: + def __init__(self, signer: Signer, send: Callable[[str, bytes, str, float], ConnectorResponse], *, enabled: bool = True, timeout_seconds: float = 5.0, max_attempts: int = 4, sleeper: Callable[[float], None] = time.sleep) -> None: + if timeout_seconds <= 0 or max_attempts < 1: raise ValueError("invalid connector retry policy") + self._signer,self._send,self._enabled,self._timeout,self._attempts,self._sleep=signer,send,enabled,timeout_seconds,max_attempts,sleeper + + def publish(self, body: bytes, correlation_id: str) -> ConnectorResponse: + if not self._enabled: raise RuntimeError("CONNECTOR_DISABLED") + if not _REQUEST_ID.fullmatch(correlation_id): raise ValueError("INVALID_CORRELATION_ID") + last: Exception|None=None + for attempt in range(self._attempts): + try: + # A retry gets a fresh jti: the previous request may have been + # accepted even when its response was lost. + token=self._signer.mint("yovision-sense",("runtime-status:write",),"POST","/machine/v1/runtime-status",body) + response=self._send("Bearer "+token,body,correlation_id,self._timeout) + if 200<=response.status<300:return response + if response.status<500:raise RuntimeError(f"STATUS_REJECTED_{response.status}") + last=RuntimeError(f"STATUS_REMOTE_{response.status}") + except (TimeoutError,ConnectionError) as exc:last=exc + if attempt+1 None: + super().__init__(code) + self.code = code + + +@dataclass(frozen=True, slots=True) +class AppliedConfig: + config_id: str + revision: int + logical_device_id: str + media_ref: str + profile_encoding: str + frame_rate: float + rule_state: str + rules: RuleSet | None + + +@dataclass(frozen=True, slots=True) +class ApplyResult: + config_id: str + revision: int + state: str + config: AppliedConfig | None + + +class SourceConfigConsumer: + """Persists validated snapshots before atomically changing the active pointer.""" + + def __init__(self, state_path: str | Path, *, clock: Callable[[], float] | None = None) -> None: + self._path = str(state_path) + self._clock = clock or time.time + self._lock = threading.RLock() + with closing(self._connect()) as connection: + connection.executescript( + """ + CREATE TABLE IF NOT EXISTS source_snapshots ( + config_id TEXT NOT NULL, revision INTEGER NOT NULL, effective_at INTEGER NOT NULL, + state TEXT NOT NULL, payload TEXT NOT NULL, PRIMARY KEY(config_id, revision)); + CREATE TABLE IF NOT EXISTS source_active ( + config_id TEXT PRIMARY KEY, revision INTEGER NOT NULL, + FOREIGN KEY(config_id, revision) REFERENCES source_snapshots(config_id, revision)); + """ + ) + connection.commit() + + def _connect(self) -> sqlite3.Connection: + connection = sqlite3.connect(self._path, timeout=5) + connection.execute("PRAGMA foreign_keys=ON") + connection.execute("PRAGMA journal_mode=WAL") + return connection + + def apply(self, body: bytes) -> ApplyResult: + document = _parse_and_validate(body) + config_id, revision = document["config_id"], document["revision"] + mapped = _map(document) + effective_at = int(_timestamp(document["effective_at"])) + with self._lock, closing(self._connect()) as connection: + connection.execute("BEGIN IMMEDIATE") + latest = connection.execute( + "SELECT revision, payload, state, effective_at FROM source_snapshots WHERE config_id=? ORDER BY revision DESC LIMIT 1", + (config_id,), + ).fetchone() + if latest and revision < latest[0]: + connection.rollback() + raise SourceConfigError("STALE_REVISION") + canonical = body.decode("utf-8") + if latest and revision == latest[0]: + if json.loads(latest[1]) != document: + connection.rollback() + raise SourceConfigError("REVISION_CONFLICT") + connection.rollback() + return ApplyResult(config_id, revision, "idempotent", self.get_active(config_id)) + connection.execute( + "INSERT INTO source_snapshots(config_id, revision, effective_at, state, payload) VALUES (?, ?, ?, ?, ?)", + (config_id, revision, effective_at, document["rule_set"]["state"], canonical), + ) + if effective_at <= int(self._clock()): + connection.execute( + "INSERT INTO source_active(config_id, revision) VALUES (?, ?) ON CONFLICT(config_id) DO UPDATE SET revision=excluded.revision", + (config_id, revision), + ) + state = "applied" + else: + state = "scheduled" + connection.commit() + return ApplyResult(config_id, revision, state, mapped if state == "applied" else self.get_active(config_id)) + + def activate_due(self) -> tuple[AppliedConfig, ...]: + now = int(self._clock()) + activated: list[AppliedConfig] = [] + with self._lock, closing(self._connect()) as connection: + connection.execute("BEGIN IMMEDIATE") + rows = connection.execute( + "SELECT s.payload FROM source_snapshots s JOIN (SELECT config_id, MAX(revision) revision FROM source_snapshots WHERE effective_at<=? GROUP BY config_id) d ON d.config_id=s.config_id AND d.revision=s.revision", + (now,), + ).fetchall() + for (payload,) in rows: + document = json.loads(payload) + connection.execute( + "INSERT INTO source_active(config_id, revision) VALUES (?, ?) ON CONFLICT(config_id) DO UPDATE SET revision=excluded.revision", + (document["config_id"], document["revision"]), + ) + activated.append(_map(document)) + connection.commit() + return tuple(activated) + + def get_active(self, config_id: str) -> AppliedConfig | None: + with closing(self._connect()) as connection: + row = connection.execute( + "SELECT s.payload FROM source_active a JOIN source_snapshots s ON s.config_id=a.config_id AND s.revision=a.revision WHERE a.config_id=?", + (config_id,), + ).fetchone() + return _map(json.loads(row[0])) if row else None + + +def _parse_and_validate(body: bytes) -> dict[str, object]: + try: + document = json.loads(body.decode("utf-8")) + except (UnicodeDecodeError, json.JSONDecodeError): + raise SourceConfigError("CONFIG_INVALID") from None + if not isinstance(document, dict): + raise SourceConfigError("CONFIG_INVALID") + if document.get("schema_version") != VERSION: + raise SourceConfigError("UNSUPPORTED_SCHEMA_VERSION") + if _contains_secret(document): + raise SourceConfigError("CONFIG_INVALID") + required = {"schema_version", "config_id", "revision", "published_at", "effective_at", "site", "logical_device", "profile", "media", "rule_set", "integrity"} + if set(document) - (required | {"extensions"}) or not required <= set(document): + raise SourceConfigError("CONFIG_INVALID") + extensions = document.get("extensions", {}) + if not isinstance(extensions, dict) or any( + not isinstance(namespace, str) + or not _EXTENSION_NAMESPACE.fullmatch(namespace) + or not isinstance(value, dict) + for namespace, value in extensions.items() + ): + raise SourceConfigError("CONFIG_INVALID") + integrity = document.get("integrity") + if not isinstance(integrity, dict) or set(integrity) != {"algorithm", "value"} or integrity.get("algorithm") != "sha256": + raise SourceConfigError("CONFIG_INVALID") + unsigned = dict(document); unsigned.pop("integrity") + digest = hashlib.sha256(json.dumps(unsigned, ensure_ascii=False, separators=(",", ":"), sort_keys=True).encode()).hexdigest() + if not hmac.compare_digest(digest, str(integrity.get("value", ""))): + raise SourceConfigError("CONFIG_INVALID") + try: + if not _ID.fullmatch(document["config_id"]) or isinstance(document["revision"], bool) or document["revision"] < 1: + raise ValueError + published, effective = _timestamp(document["published_at"]), _timestamp(document["effective_at"]) + if effective < published: + raise ValueError + for field in ("site", "logical_device"): + if not isinstance(document[field], dict) or set(document[field]) != {"id"} or not _ID.fullmatch(document[field]["id"]): raise ValueError + _validate_profile(document["profile"]) + media = document["media"] + if not isinstance(media, dict) or set(media) != {"ref", "transport"} or media["transport"] != "rtsp" or not isinstance(media["ref"], str) or not media["ref"].startswith("media:") or any(marker in media["ref"] for marker in ("?", "#", "@", "\\", "://")): raise ValueError + _validate_rules(document["rule_set"], document["profile"]) + except (KeyError, TypeError, ValueError, AttributeError): + raise SourceConfigError("CONFIG_INVALID") from None + return document + + +def _validate_profile(profile: object) -> None: + if not isinstance(profile, dict) or set(profile) != {"id", "width", "height", "encoding", "frame_rate"}: raise ValueError + if not _ID.fullmatch(profile["id"]) or profile["encoding"] not in {"H264", "H265", "MJPEG"}: raise ValueError + for field in ("width", "height"): + if isinstance(profile[field], bool) or not isinstance(profile[field], int) or profile[field] < 1: raise ValueError + if isinstance(profile["frame_rate"], bool) or not isinstance(profile["frame_rate"], (int, float)) or profile["frame_rate"] <= 0: raise ValueError + + +def _validate_rules(rules: object, profile: Mapping[str, object]) -> None: + if not isinstance(rules, dict) or set(rules) != {"version", "state", "profile_binding", "areas", "directional_lines"}: raise ValueError + if not _ID.fullmatch(rules["version"]) or rules["state"] not in {"active", "disabled", "recalibration_required"}: raise ValueError + binding = rules["profile_binding"] + if binding != {"profile_id": profile["id"], "width": profile["width"], "height": profile["height"]}: raise ValueError + if not isinstance(rules["areas"], list) or not isinstance(rules["directional_lines"], list) or len(rules["areas"]) > 1024 or len(rules["directional_lines"]) > 1024: raise ValueError + identifiers: set[str] = set() + for area in rules["areas"]: + if not isinstance(area, dict) or set(area) != {"id", "version", "kind", "enabled", "points"} or area["kind"] != "danger_area" or not isinstance(area["enabled"], bool) or not 3 <= len(area["points"]) <= 256: raise ValueError + _rule_identity(area, identifiers); points = tuple(_point(value) for value in area["points"]) + polygon = sum(a[0]*b[1]-b[0]*a[1] for a,b in zip(points, points[1:]+points[:1])) / 2 + if abs(polygon) < 1e-9: raise ValueError + for line in rules["directional_lines"]: + if not isinstance(line, dict) or set(line) != {"id", "version", "kind", "enabled", "start", "end", "trigger_direction"} or line["kind"] != "directional_line" or not isinstance(line["enabled"], bool) or line["trigger_direction"] not in {"left_to_right", "right_to_left"}: raise ValueError + _rule_identity(line, identifiers) + if _point(line["start"]) == _point(line["end"]): raise ValueError + + +def _rule_identity(rule: Mapping[str, object], identifiers: set[str]) -> None: + if not isinstance(rule["id"], str) or not _ID.fullmatch(rule["id"]) or rule["id"] in identifiers or isinstance(rule["version"], bool) or not isinstance(rule["version"], int) or rule["version"] < 1: raise ValueError + identifiers.add(rule["id"]) + + +def _point(value: object) -> tuple[float, float]: + if not isinstance(value, dict) or set(value) != {"x", "y"}: raise ValueError + x, y = value["x"], value["y"] + if isinstance(x, bool) or isinstance(y, bool) or not isinstance(x, (int,float)) or not isinstance(y,(int,float)) or not 0 <= x <= 1 or not 0 <= y <= 1: raise ValueError + return float(x), float(y) + + +def _map(document: Mapping[str, object]) -> AppliedConfig: + profile, rules = document["profile"], document["rule_set"] + rule_set = None + if rules["state"] == "active": + rule_set = RuleSet( + version=rules["version"], profile_id=profile["id"], width=profile["width"], height=profile["height"], + areas=tuple(AreaDefinition(a["id"], tuple(NormalizedPoint(**p) for p in a["points"])) for a in rules["areas"] if a["enabled"]), + directional_lines=tuple(DirectionalLineDefinition(l["id"], NormalizedPoint(**l["start"]), NormalizedPoint(**l["end"]), l["trigger_direction"]) for l in rules["directional_lines"] if l["enabled"]), + ) + return AppliedConfig(document["config_id"], document["revision"], document["logical_device"]["id"], document["media"]["ref"], profile["encoding"], float(profile["frame_rate"]), rules["state"], rule_set) + + +def _timestamp(value: object) -> float: + if not isinstance(value, str) or not value.endswith("Z"): raise ValueError + return datetime.fromisoformat(value[:-1] + "+00:00").astimezone(timezone.utc).timestamp() + + +def _contains_secret(value: object) -> bool: + if isinstance(value, dict): return any(_SECRET.search(str(k)) or _contains_secret(v) for k,v in value.items()) + if isinstance(value, list): return any(_contains_secret(item) for item in value) + if isinstance(value, str): + split=urlsplit(value) + return bool(split.username or split.password or value.startswith("file:") or re.match(r"^[A-Za-z]:[\\/]", value)) + return False diff --git a/Brain/src/yovision_brain/integration/sense_control/replay.py b/Brain/src/yovision_brain/integration/sense_control/replay.py new file mode 100644 index 0000000..2fc3550 --- /dev/null +++ b/Brain/src/yovision_brain/integration/sense_control/replay.py @@ -0,0 +1,44 @@ +"""Brain-owned durable replay storage; never shared with Sense or Bell.""" + +from __future__ import annotations + +import sqlite3 +import threading +from contextlib import closing +from pathlib import Path + + +class SQLiteReplayStore: + def __init__(self, path: str | Path) -> None: + self._path = str(path) + self._lock = threading.Lock() + with closing(self._connect()) as connection: + connection.execute("PRAGMA journal_mode=WAL") + connection.execute( + "CREATE TABLE IF NOT EXISTS machine_replay (principal TEXT NOT NULL, token_id TEXT NOT NULL, expires_at INTEGER NOT NULL, PRIMARY KEY(principal, token_id))" + ) + + def _connect(self) -> sqlite3.Connection: + connection = sqlite3.connect(self._path, timeout=5, isolation_level=None) + connection.execute("PRAGMA busy_timeout=5000") + return connection + + def consume(self, principal: str, token_id: str, expires_at: int, now: int) -> bool: + if not principal or not token_id or expires_at <= now: + return False + with self._lock, closing(self._connect()) as connection: + try: + connection.execute("BEGIN IMMEDIATE") + connection.execute("DELETE FROM machine_replay WHERE expires_at <= ?", (now,)) + connection.execute( + "INSERT INTO machine_replay(principal, token_id, expires_at) VALUES (?, ?, ?)", + (principal, token_id, expires_at), + ) + connection.execute("COMMIT") + return True + except sqlite3.IntegrityError: + connection.execute("ROLLBACK") + return False + except Exception: + connection.execute("ROLLBACK") + return False diff --git a/Brain/src/yovision_brain/integration/sense_control/status.py b/Brain/src/yovision_brain/integration/sense_control/status.py new file mode 100644 index 0000000..f0210ee --- /dev/null +++ b/Brain/src/yovision_brain/integration/sense_control/status.py @@ -0,0 +1,77 @@ +"""Persistent sequence allocation and safe runtime-status/v1 production.""" + +from __future__ import annotations + +import json +import re +import sqlite3 +import threading +import uuid +from contextlib import closing +from datetime import datetime, timezone +from pathlib import Path +from typing import Mapping, Sequence + +_LOGICAL_REF = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$") +_VERSION = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._+-]{0,63}$") +_ERROR_CODE = re.compile(r"^[A-Z][A-Z0-9_]{2,63}$") +_RUNTIME_STATES = {"unconfigured", "starting", "running", "degraded", "failed", "stopped"} + + +class RuntimeStatusPublisher: + def __init__(self, state_path: str | Path, brain_instance_ref: str) -> None: + self._path, self._instance, self._lock = str(state_path), brain_instance_ref, threading.Lock() + with closing(self._connect()) as connection: + connection.execute("CREATE TABLE IF NOT EXISTS runtime_sequence(instance_ref TEXT PRIMARY KEY, sequence INTEGER NOT NULL)") + connection.commit() + + def _connect(self) -> sqlite3.Connection: + return sqlite3.connect(self._path, timeout=5) + + def build(self, *, runtime_state: str, runtime_version: str, started_at: datetime | None, model_ref: str, model_version: str, configurations: Sequence[Mapping[str, object]], health: Mapping[str, object], inputs: Sequence[Mapping[str, object]], observed_at: datetime | None = None) -> bytes: + _validate(runtime_state, runtime_version, self._instance, model_ref, model_version, configurations, health, inputs) + with self._lock, closing(self._connect()) as connection: + connection.execute("BEGIN IMMEDIATE") + row=connection.execute("SELECT sequence FROM runtime_sequence WHERE instance_ref=?",(self._instance,)).fetchone();sequence=(row[0]+1) if row else 0 + connection.execute("INSERT INTO runtime_sequence(instance_ref,sequence) VALUES(?,?) ON CONFLICT(instance_ref) DO UPDATE SET sequence=excluded.sequence",(self._instance,sequence));connection.commit() + observed=(observed_at or datetime.now(timezone.utc)).astimezone(timezone.utc) + document={"schema_version":"yovision.runtime-status/v1","status_id":str(uuid.uuid4()),"brain_instance_ref":self._instance,"sequence":sequence,"observed_at":_utc(observed),"runtime":{"state":runtime_state,"version":runtime_version,"started_at":_utc(started_at) if started_at else None},"model":{"model_ref":model_ref,"version":model_version},"configurations":list(configurations),"health":dict(health),"inputs":list(inputs)} + raw=json.dumps(document,separators=(",",":"),sort_keys=True).encode() + lowered=raw.lower(); + for marker in (b"password",b"credential",b"stream_uri",b"cookie",b"jwt",b"file://"): + if marker in lowered: raise ValueError("runtime status contains sensitive field") + return raw + + +def _utc(value: datetime) -> str: + if value.tzinfo is None: raise ValueError("runtime timestamp must be timezone-aware") + return value.astimezone(timezone.utc).isoformat(timespec="seconds").replace("+00:00","Z") + + +def _validate(runtime_state: str, runtime_version: str, instance: str, model_ref: str, model_version: str, configurations: Sequence[Mapping[str, object]], health: Mapping[str, object], inputs: Sequence[Mapping[str, object]]) -> None: + if runtime_state not in _RUNTIME_STATES or not _VERSION.fullmatch(runtime_version) or not _LOGICAL_REF.fullmatch(instance) or not _LOGICAL_REF.fullmatch(model_ref) or not _VERSION.fullmatch(model_version): + raise ValueError("invalid runtime identity or version") + if len(configurations) > 4096 or len(inputs) > 4096: + raise ValueError("runtime status collection too large") + seen: set[str] = set() + for item in configurations: + if set(item) != {"config_id", "apply_state", "applied_revision", "error_code"} or not isinstance(item["config_id"], str) or not _LOGICAL_REF.fullmatch(item["config_id"]) or item["config_id"] in seen: + raise ValueError("invalid configuration status") + seen.add(item["config_id"]); state=item["apply_state"]; revision=item["applied_revision"]; error=item["error_code"] + if state not in {"not_configured","applying","applied","rejected"} or (state=="not_configured" and revision is not None) or (state=="applied" and (isinstance(revision,bool) or not isinstance(revision,int) or revision<1)) or (state=="rejected" and (not isinstance(error,str) or not _ERROR_CODE.fullmatch(error))): + raise ValueError("invalid configuration status") + if set(health) != {"overall","error_codes","metrics"} or health["overall"] not in {"healthy","degraded","unhealthy"} or not _codes(health["error_codes"],32) or not _metrics(health["metrics"]): + raise ValueError("invalid health status") + for item in inputs: + if set(item) != {"input_ref","state","error_codes","metrics"} or not isinstance(item["input_ref"],str) or not _LOGICAL_REF.fullmatch(item["input_ref"]) or item["state"] not in _RUNTIME_STATES or not _codes(item["error_codes"],16) or not _metrics(item["metrics"]): + raise ValueError("invalid input status") + + +def _codes(value: object, limit: int) -> bool: + return isinstance(value,list) and len(value)<=limit and len(set(value))==len(value) and all(isinstance(code,str) and _ERROR_CODE.fullmatch(code) for code in value) + + +def _metrics(value: object) -> bool: + if not isinstance(value,Mapping) or set(value)!={"load_percent","queue_depth","latency_ms"}: return False + load,queue,latency=value["load_percent"],value["queue_depth"],value["latency_ms"] + return not isinstance(load,bool) and isinstance(load,(int,float)) and 0<=load<=100 and not isinstance(queue,bool) and isinstance(queue,int) and queue>=0 and not isinstance(latency,bool) and isinstance(latency,(int,float)) and latency>=0 diff --git a/Brain/tests/integration/sense_control/test_connector.py b/Brain/tests/integration/sense_control/test_connector.py new file mode 100644 index 0000000..222ef8e --- /dev/null +++ b/Brain/tests/integration/sense_control/test_connector.py @@ -0,0 +1,101 @@ +from __future__ import annotations + +import hashlib +import json +import threading +from datetime import datetime, timedelta, timezone + +import pytest +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey + +from yovision_brain.integration.machine_identity.token import KeyRecord, Registry, Signer, Verifier +from yovision_brain.integration.sense_control.connector import ConnectorResponse, SourceConfigEndpoint, StatusSender +from yovision_brain.integration.sense_control.consumer import SourceConfigConsumer, SourceConfigError +from yovision_brain.integration.sense_control.replay import SQLiteReplayStore +from yovision_brain.integration.sense_control.status import RuntimeStatusPublisher + + +def source_document(revision: int = 1, *, state: str = "active", effective: int = 0) -> bytes: + now=datetime(2026,8,31,tzinfo=timezone.utc) + document={"schema_version":"yovision.source-config/v1","config_id":"gate-primary","revision":revision,"published_at":now.isoformat().replace("+00:00","Z"),"effective_at":(now+timedelta(seconds=effective)).isoformat().replace("+00:00","Z"),"site":{"id":"site-east"},"logical_device":{"id":"camera-1"},"profile":{"id":"main","width":1920,"height":1080,"encoding":"H264","frame_rate":25},"media":{"ref":"media:site-east/camera-1/main","transport":"rtsp"},"rule_set":{"version":f"rules-{revision}","state":state,"profile_binding":{"profile_id":"main","width":1920,"height":1080},"areas":[{"id":"danger","version":1,"kind":"danger_area","enabled":True,"points":[{"x":.1,"y":.1},{"x":.8,"y":.1},{"x":.5,"y":.8}]}],"directional_lines":[]}} + digest=hashlib.sha256(json.dumps(document,separators=(",",":"),sort_keys=True).encode()).hexdigest();document["integrity"]={"algorithm":"sha256","value":digest} + return json.dumps(document,separators=(",",":"),sort_keys=True).encode() + + +def with_extension(body: bytes, extensions: object) -> bytes: + document=json.loads(body);document["extensions"]=extensions;unsigned=dict(document);unsigned.pop("integrity");document["integrity"]["value"]=hashlib.sha256(json.dumps(unsigned,separators=(",",":"),sort_keys=True).encode()).hexdigest();return json.dumps(document,separators=(",",":"),sort_keys=True).encode() + + +def identity(tmp_path, now: int): + private=Ed25519PrivateKey.generate();signer=Signer("yv:sense:east","sense-key-01",private,clock=lambda:now) + registry=Registry([KeyRecord("yv:sense:east","sense-key-01",private.public_key(),"yovision-brain",frozenset({"source-config:write"}))]) + replay=SQLiteReplayStore(tmp_path/"replay.sqlite") + return signer,Verifier(registry,replay,clock=lambda:now) + + +def test_authenticated_apply_is_idempotent_and_replay_survives_restart(tmp_path): + now=int(datetime(2026,8,31,tzinfo=timezone.utc).timestamp());signer,verifier=identity(tmp_path,now);body=source_document();consumer=SourceConfigConsumer(tmp_path/"state.sqlite",clock=lambda:now);endpoint=SourceConfigEndpoint(consumer,verifier) + token=signer.mint("yovision-brain",("source-config:write",),"POST","/machine/v1/source-config",body) + result=endpoint.receive("Bearer "+token,body,"corr-request-0001");assert result.state=="applied" and result.config.rules is not None + with pytest.raises(ValueError,match="machine_token_replayed"):endpoint.receive("Bearer "+token,body,"corr-request-0001") + restarted=SourceConfigEndpoint(SourceConfigConsumer(tmp_path/"state.sqlite",clock=lambda:now),Verifier(verifier._registry,SQLiteReplayStore(tmp_path/"replay.sqlite"),clock=lambda:now)) + with pytest.raises(ValueError,match="machine_token_replayed"):restarted.receive("Bearer "+token,body,"corr-request-0002") + new_token=signer.mint("yovision-brain",("source-config:write",),"POST","/machine/v1/source-config",body);assert restarted.receive("Bearer "+new_token,body,"corr-request-0003").state=="idempotent" + + +def test_atomic_replay_accepts_once_under_concurrency(tmp_path): + store=SQLiteReplayStore(tmp_path/"atomic.sqlite");results=[] + threads=[threading.Thread(target=lambda:results.append(store.consume("yv:sense:east","token-id",200,100))) for _ in range(12)] + for thread in threads:thread.start() + for thread in threads:thread.join() + assert results.count(True)==1 + + +def test_stale_unknown_profile_and_recalibration_are_safe(tmp_path): + now=int(datetime(2026,8,31,tzinfo=timezone.utc).timestamp());consumer=SourceConfigConsumer(tmp_path/"state.sqlite",clock=lambda:now) + assert consumer.apply(source_document(2)).config.rules is not None + with pytest.raises(SourceConfigError,match="STALE_REVISION"):consumer.apply(source_document(1)) + invalid=json.loads(source_document(3));invalid["profile"]["width"]=1280;invalid["integrity"]["value"]="0"*64 + with pytest.raises(SourceConfigError,match="CONFIG_INVALID"):consumer.apply(json.dumps(invalid).encode()) + unknown=json.loads(source_document(3));unknown["schema_version"]="yovision.source-config/v2" + with pytest.raises(SourceConfigError,match="UNSUPPORTED_SCHEMA_VERSION"):consumer.apply(json.dumps(unknown).encode()) + safe=consumer.apply(source_document(3,state="recalibration_required"));assert safe.config.rule_state=="recalibration_required" and safe.config.rules is None + + +def test_future_effective_snapshot_activates_atomically_after_restart(tmp_path): + base=int(datetime(2026,8,31,tzinfo=timezone.utc).timestamp());current=[base];path=tmp_path/"state.sqlite";consumer=SourceConfigConsumer(path,clock=lambda:current[0]) + assert consumer.apply(source_document(1)).state=="applied";scheduled=consumer.apply(source_document(2,effective=60));assert scheduled.state=="scheduled" and scheduled.config.revision==1 + current[0]+=61;restarted=SourceConfigConsumer(path,clock=lambda:current[0]);activated=restarted.activate_due();assert activated[0].revision==2 and restarted.get_active("gate-primary").revision==2 + + +def test_unknown_valid_extension_namespace_is_ignored(tmp_path): + now=int(datetime(2026,8,31,tzinfo=timezone.utc).timestamp());consumer=SourceConfigConsumer(tmp_path/"state.sqlite",clock=lambda:now) + result=consumer.apply(with_extension(source_document(),{"vendor.example":{"feature":"safe"}}));assert result.state=="applied" and result.config.revision==1 + + +@pytest.mark.parametrize("extensions", [[], {"1invalid":{}}, {"vendor_ok":{}}, {"vendor.example":"not-an-object"}]) +def test_invalid_extensions_are_rejected(tmp_path, extensions): + now=int(datetime(2026,8,31,tzinfo=timezone.utc).timestamp());consumer=SourceConfigConsumer(tmp_path/"state.sqlite",clock=lambda:now) + with pytest.raises(SourceConfigError,match="CONFIG_INVALID"):consumer.apply(with_extension(source_document(),extensions)) + + +def test_status_sequence_restart_retry_timeout_and_disable(tmp_path): + path=tmp_path/"status.sqlite";publisher=RuntimeStatusPublisher(path,"brain-east-01");health={"overall":"healthy","error_codes":[],"metrics":{"load_percent":1.0,"queue_depth":0,"latency_ms":2.0}} + one=json.loads(publisher.build(runtime_state="running",runtime_version="1.0.0",started_at=datetime.now(timezone.utc),model_ref="people-detection",model_version="1",configurations=[],health=health,inputs=[]));two=json.loads(RuntimeStatusPublisher(path,"brain-east-01").build(runtime_state="running",runtime_version="1.0.0",started_at=None,model_ref="people-detection",model_version="1",configurations=[],health=health,inputs=[]));assert (one["sequence"],two["sequence"])==(0,1) + private=Ed25519PrivateKey.generate();signer=Signer("yv:brain:east","brain-key-01",private,clock=lambda:1_787_000_000);attempts=[] + def send(_auth,_body,_corr,timeout):attempts.append(timeout);raise TimeoutError + sender=StatusSender(signer,send,max_attempts=3,sleeper=lambda _:None) + with pytest.raises(RuntimeError,match="STATUS_DELIVERY_EXHAUSTED"):sender.publish(b"{}","corr-request-0001") + assert attempts==[5.0,5.0,5.0] + disabled=StatusSender(signer,lambda *_:ConnectorResponse(204,b"","corr-request-0001"),enabled=False) + with pytest.raises(RuntimeError,match="CONNECTOR_DISABLED"):disabled.publish(b"{}","corr-request-0001") + + +@pytest.mark.parametrize("request_id", ["short", "0123456789abcde\n", "0123456789abcde!", "a"*129]) +def test_connector_rejects_unsafe_request_ids(tmp_path, request_id): + now=int(datetime(2026,8,31,tzinfo=timezone.utc).timestamp());signer,verifier=identity(tmp_path,now);body=source_document();endpoint=SourceConfigEndpoint(SourceConfigConsumer(tmp_path/"state.sqlite",clock=lambda:now),verifier) + token=signer.mint("yovision-brain",("source-config:write",),"POST","/machine/v1/source-config",body) + with pytest.raises(ValueError,match="INVALID_CORRELATION_ID"):endpoint.receive("Bearer "+token,body,request_id) + status_signer=Signer("yv:brain:east","brain-key-01",Ed25519PrivateKey.generate(),clock=lambda:now) + sender=StatusSender(status_signer,lambda *_:ConnectorResponse(204,b"",request_id)) + with pytest.raises(ValueError,match="INVALID_CORRELATION_ID"):sender.publish(b"{}",request_id) diff --git a/Sense/server/app/sense/integration/brain_control/connector.go b/Sense/server/app/sense/integration/brain_control/connector.go new file mode 100644 index 0000000..d3ca85d --- /dev/null +++ b/Sense/server/app/sense/integration/brain_control/connector.go @@ -0,0 +1,118 @@ +package brain_control + +import ( + "errors" + "fmt" + "regexp" + "time" + + mi "git.ilapage.cn/ila/yovision/Sense/server/app/sense/integration/machine_identity" +) + +const ( + sourceConfigPath = "/machine/v1/source-config" + runtimeStatusPath = "/machine/v1/runtime-status" +) + +var requestIDPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._:-]{15,127}$`) + +type SendResponse struct { + StatusCode int + Body []byte + CorrelationID string +} +type SendFunc func(authorization string, body []byte, correlationID string, timeout time.Duration) (SendResponse, error) + +type ConfigSender struct { + Signer mi.Signer + Send SendFunc + Enabled bool + Timeout time.Duration + MaxAttempts int + Sleep func(time.Duration) +} + +func (s ConfigSender) Publish(config SourceConfig, correlationID string) (SendResponse, error) { + if !s.Enabled { + return SendResponse{}, errors.New("CONNECTOR_DISABLED") + } + if s.Send == nil || !requestIDPattern.MatchString(correlationID) { + return SendResponse{}, errors.New("invalid connector configuration") + } + if s.Timeout <= 0 { + s.Timeout = 5 * time.Second + } + if s.MaxAttempts == 0 { + s.MaxAttempts = 4 + } + if s.MaxAttempts < 1 { + return SendResponse{}, errors.New("invalid connector retry policy") + } + if s.Sleep == nil { + s.Sleep = time.Sleep + } + if err := ValidateSourceConfig(config); err != nil { + return SendResponse{}, err + } + body, err := MarshalSourceConfig(config) + if err != nil { + return SendResponse{}, err + } + var last error + for attempt := 0; attempt < s.MaxAttempts; attempt++ { + token, mintErr := s.Signer.Mint("yovision-brain", []string{"source-config:write"}, "POST", sourceConfigPath, body) + if mintErr != nil { + return SendResponse{}, mintErr + } + response, sendErr := s.Send("Bearer "+token, body, correlationID, s.Timeout) + if sendErr == nil && response.StatusCode >= 200 && response.StatusCode < 300 { + return response, nil + } + if sendErr == nil && response.StatusCode < 500 { + return SendResponse{}, fmt.Errorf("source config rejected: %d", response.StatusCode) + } + if sendErr != nil { + last = sendErr + } else { + last = fmt.Errorf("source config remote status: %d", response.StatusCode) + } + if attempt+1 < s.MaxAttempts { + delay := time.Second << attempt + if delay > 30*time.Second { + delay = 30 * time.Second + } + s.Sleep(delay) + } + } + return SendResponse{}, fmt.Errorf("source config delivery exhausted: %w", last) +} + +type RuntimeStatusEndpoint struct { + Verifier mi.Verifier + Store ProjectionStore + Enabled bool + MaxBodyBytes int +} + +func (e RuntimeStatusEndpoint) Receive(authorization string, body []byte, correlationID string, expected map[string]int64) (ProjectionView, error) { + if !e.Enabled { + return ProjectionView{}, errors.New("CONNECTOR_DISABLED") + } + if !requestIDPattern.MatchString(correlationID) { + return ProjectionView{}, errors.New("INVALID_CORRELATION_ID") + } + if e.MaxBodyBytes == 0 { + e.MaxBodyBytes = 10 * 1024 * 1024 + } + if len(body) > e.MaxBodyBytes { + return ProjectionView{}, errors.New("REQUEST_TOO_LARGE") + } + token, err := mi.BearerToken(authorization) + if err != nil { + return ProjectionView{}, err + } + if _, err = e.Verifier.Verify(token, "yovision-sense", "runtime-status:write", "POST", runtimeStatusPath, body); err != nil { + return ProjectionView{}, err + } + return e.Store.Ingest(body, expected) +} diff --git a/Sense/server/app/sense/integration/brain_control/connector_test.go b/Sense/server/app/sense/integration/brain_control/connector_test.go new file mode 100644 index 0000000..dbec520 --- /dev/null +++ b/Sense/server/app/sense/integration/brain_control/connector_test.go @@ -0,0 +1,124 @@ +package brain_control + +import ( + "crypto/ed25519" + "crypto/rand" + "encoding/json" + "strings" + "sync" + "testing" + "time" + + mi "git.ilapage.cn/ila/yovision/Sense/server/app/sense/integration/machine_identity" + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +func memoryDB(t *testing.T, name string) *gorm.DB { + t.Helper() + db, err := gorm.Open(sqlite.Open("file:"+name+"?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err = db.AutoMigrate(&ReplayToken{}, &RuntimeProjection{}, &SourceRevision{}); err != nil { + t.Fatal(err) + } + return db +} + +func TestSourceMapperAndRecalibration(t *testing.T) { + now := time.Date(2026, 8, 31, 0, 0, 0, 0, time.UTC) + facts := SourceFacts{ConfigID: "gate-primary", SiteID: "site-east", LogicalDeviceID: "camera-1", MediaPath: "site/camera/main", Revision: 1, PublishedAt: now, EffectiveAt: now, Profile: Profile{ID: "main", Width: 1920, Height: 1080, Encoding: "h264", FrameRate: 25}, RuleSetVersion: "rules-1", Areas: []AreaRule{{ID: "danger", Version: 1, Kind: "danger_area", Enabled: true, Points: []Point{{.1, .1}, {.8, .1}, {.5, .8}}}}} + config, err := MapSourceConfig(facts) + if err != nil { + t.Fatal(err) + } + if err = ValidateSourceConfig(config); err != nil { + t.Fatal(err) + } + facts.Revision = 2 + facts.NeedsRecalibration = true + config, err = MapSourceConfig(facts) + if err != nil { + t.Fatal(err) + } + if config.RuleSet.State != "recalibration_required" || config.RuleSet.Areas[0].Enabled { + t.Fatal("recalibration did not disable rules") + } +} + +func TestReplayAtomicAndRestart(t *testing.T) { + db := memoryDB(t, "replay-package") + pub, priv, _ := ed25519.GenerateKey(rand.Reader) + now := time.Date(2026, 8, 31, 0, 0, 0, 0, time.UTC) + registry, _ := mi.NewRegistry(mi.KeyRecord{Principal: "yv:brain:east", KeyID: "brain-key-01", PublicKey: pub, Audience: "yovision-sense", Scopes: []string{"runtime-status:write"}, Enabled: true}) + signer := mi.Signer{Principal: "yv:brain:east", KeyID: "brain-key-01", PrivateKey: priv, Now: func() time.Time { return now }} + body := []byte("{}") + token, _ := signer.Mint("yovision-sense", []string{"runtime-status:write"}, "POST", "/machine/v1/runtime-status", body) + accepted := 0 + var mu sync.Mutex + var wg sync.WaitGroup + for range 8 { + wg.Add(1) + go func() { + defer wg.Done() + v := mi.Verifier{Registry: registry, Replay: GORMReplayStore{DB: db}, Now: func() time.Time { return now }} + if _, err := v.Verify(token, "yovision-sense", "runtime-status:write", "POST", "/machine/v1/runtime-status", body); err == nil { + mu.Lock() + accepted++ + mu.Unlock() + } + }() + } + wg.Wait() + if accepted != 1 { + t.Fatalf("accepted %d", accepted) + } + v := mi.Verifier{Registry: registry, Replay: GORMReplayStore{DB: db}, Now: func() time.Time { return now }} + if _, err := v.Verify(token, "yovision-sense", "runtime-status:write", "POST", "/machine/v1/runtime-status", body); err == nil { + t.Fatal("restart replay accepted") + } +} + +func TestProjectionStaleRecoveryAndMismatch(t *testing.T) { + db := memoryDB(t, "projection-package") + now := time.Date(2026, 8, 31, 0, 0, 0, 0, time.UTC) + store := ProjectionStore{DB: db, Clock: func() time.Time { return now }, StaleAfter: 90 * time.Second} + raw := runtimeFixture("018f4d6a-8d1b-4a25-8b37-9085f9c0d101", 1, now, 2) + view, err := store.Ingest(raw, map[string]int64{"gate": 3}) + if err != nil || !view.RevisionMismatch { + t.Fatalf("view %+v err %v", view, err) + } + now = now.Add(91 * time.Second) + view, _ = store.View("brain-east-01") + if !view.Offline { + t.Fatal("not offline") + } + raw = runtimeFixture("018f4d6a-8d1b-4a25-8b37-9085f9c0d102", 2, now, 3) + view, err = store.Ingest(raw, map[string]int64{"gate": 3}) + if err != nil || !view.Recovered { + t.Fatalf("recovery %+v err %v", view, err) + } +} + +func TestConnectorRejectsUnsafeRequestIDs(t *testing.T) { + for _, value := range []string{"short", "0123456789abcde\n", "0123456789abcde!", strings.Repeat("a", 129)} { + sender := ConfigSender{Enabled: true, Send: func(string, []byte, string, time.Duration) (SendResponse, error) { + t.Fatal("unsafe request id reached transport") + return SendResponse{}, nil + }} + if _, err := sender.Publish(SourceConfig{}, value); err == nil { + t.Fatalf("sender accepted request id %q", value) + } + endpoint := RuntimeStatusEndpoint{Enabled: true} + if _, err := endpoint.Receive("Bearer ignored", nil, value, nil); err == nil || err.Error() != "INVALID_CORRELATION_ID" { + t.Fatalf("endpoint accepted request id %q: %v", value, err) + } + } +} + +func runtimeFixture(id string, seq int64, observed time.Time, revision int64) []byte { + value := map[string]any{"schema_version": RuntimeStatusVersion, "status_id": id, "brain_instance_ref": "brain-east-01", "sequence": seq, "observed_at": observed.Format(time.RFC3339), "runtime": map[string]any{"state": "running", "version": "1.0.0", "started_at": observed.Format(time.RFC3339)}, "model": map[string]any{"model_ref": "people", "version": "1"}, "configurations": []any{map[string]any{"config_id": "gate", "apply_state": "applied", "applied_revision": revision, "error_code": nil}}, "health": map[string]any{"overall": "healthy", "error_codes": []any{}, "metrics": map[string]any{"load_percent": 1.0, "queue_depth": 0, "latency_ms": 1.0}}, "inputs": []any{}} + raw, _ := json.Marshal(value) + return raw +} diff --git a/Sense/server/app/sense/integration/brain_control/mapper.go b/Sense/server/app/sense/integration/brain_control/mapper.go new file mode 100644 index 0000000..9f9fef1 --- /dev/null +++ b/Sense/server/app/sense/integration/brain_control/mapper.go @@ -0,0 +1,135 @@ +package brain_control + +import ( + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "regexp" + "strings" + "time" +) + +var stableID = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._~-]{0,127}$`) + +func MapSourceConfig(f SourceFacts) (SourceConfig, error) { + if !stableID.MatchString(f.ConfigID) || !stableID.MatchString(f.SiteID) || !stableID.MatchString(f.LogicalDeviceID) || !stableID.MatchString(f.Profile.ID) || f.Revision < 1 || f.Profile.Width < 1 || f.Profile.Height < 1 || f.Profile.FrameRate <= 0 { + return SourceConfig{}, errors.New("invalid source configuration facts") + } + encoding := strings.ToUpper(f.Profile.Encoding) + if encoding != "H264" && encoding != "H265" && encoding != "MJPEG" { + return SourceConfig{}, errors.New("unsupported profile encoding") + } + if f.PublishedAt.IsZero() || f.EffectiveAt.Before(f.PublishedAt) { + return SourceConfig{}, errors.New("invalid source configuration time") + } + if !stableID.MatchString(f.RuleSetVersion) { + return SourceConfig{}, errors.New("invalid rule set version") + } + if strings.ContainsAny(f.MediaPath, "?#@\\") || strings.Contains(f.MediaPath, "://") || f.MediaPath == "" { + return SourceConfig{}, errors.New("media path must be opaque and credential-free") + } + if len(f.Areas) > 1024 || len(f.DirectionalLines) > 1024 { + return SourceConfig{}, errors.New("too many rules") + } + seen := map[string]bool{} + for _, a := range f.Areas { + if !stableID.MatchString(a.ID) || seen[a.ID] || a.Version < 1 || a.Kind != "danger_area" || len(a.Points) < 3 || len(a.Points) > 256 || !validPoints(a.Points) || polygonArea(a.Points) == 0 { + return SourceConfig{}, errors.New("invalid area rule") + } + seen[a.ID] = true + } + for _, l := range f.DirectionalLines { + if !stableID.MatchString(l.ID) || seen[l.ID] || l.Version < 1 || l.Kind != "directional_line" || (l.TriggerDirection != "left_to_right" && l.TriggerDirection != "right_to_left") || !validPoints([]Point{l.Start, l.End}) || l.Start == l.End { + return SourceConfig{}, errors.New("invalid directional line rule") + } + seen[l.ID] = true + } + var out SourceConfig + out.SchemaVersion, out.ConfigID, out.Revision = SourceConfigVersion, f.ConfigID, f.Revision + out.PublishedAt, out.EffectiveAt = f.PublishedAt.UTC(), f.EffectiveAt.UTC() + out.Site.ID, out.LogicalDevice.ID = f.SiteID, f.LogicalDeviceID + out.Profile = f.Profile + out.Profile.Encoding = encoding + out.Media.Ref, out.Media.Transport = "media:"+strings.TrimPrefix(f.MediaPath, "/"), "rtsp" + out.RuleSet.Version = f.RuleSetVersion + out.RuleSet.State = "active" + if f.Disabled { + out.RuleSet.State = "disabled" + } + if f.NeedsRecalibration { + out.RuleSet.State = "recalibration_required" + } + out.RuleSet.ProfileBinding.ProfileID, out.RuleSet.ProfileBinding.Width, out.RuleSet.ProfileBinding.Height = f.Profile.ID, f.Profile.Width, f.Profile.Height + out.RuleSet.Areas = append([]AreaRule(nil), f.Areas...) + out.RuleSet.DirectionalLines = append([]DirectionalLineRule(nil), f.DirectionalLines...) + if out.RuleSet.State != "active" { + for i := range out.RuleSet.Areas { + out.RuleSet.Areas[i].Enabled = false + } + for i := range out.RuleSet.DirectionalLines { + out.RuleSet.DirectionalLines[i].Enabled = false + } + } + digest, err := sourceDigest(out) + if err != nil { + return SourceConfig{}, err + } + out.Integrity.Algorithm, out.Integrity.Value = "sha256", digest + return out, nil +} + +func MarshalSourceConfig(config SourceConfig) ([]byte, error) { return json.Marshal(config) } + +func sourceDigest(config SourceConfig) (string, error) { + raw, err := json.Marshal(config) + if err != nil { + return "", err + } + var value map[string]any + if err = json.Unmarshal(raw, &value); err != nil { + return "", err + } + delete(value, "integrity") + canonical, err := json.Marshal(value) + if err != nil { + return "", err + } + sum := sha256.Sum256(canonical) + return hex.EncodeToString(sum[:]), nil +} + +func validPoints(points []Point) bool { + for _, p := range points { + if p.X < 0 || p.X > 1 || p.Y < 0 || p.Y > 1 { + return false + } + } + return true +} +func polygonArea(p []Point) float64 { + var a float64 + for i := range p { + n := p[(i+1)%len(p)] + a += p[i].X*n.Y - n.X*p[i].Y + } + if a < 0 { + a = -a + } + return a / 2 +} + +func ValidateSourceConfig(config SourceConfig) error { + if config.SchemaVersion != SourceConfigVersion { + return fmt.Errorf("unsupported source config version") + } + digest, err := sourceDigest(config) + if err != nil || config.Integrity.Algorithm != "sha256" || digest != config.Integrity.Value { + return errors.New("source config integrity mismatch") + } + _, err = MapSourceConfig(SourceFacts{ConfigID: config.ConfigID, SiteID: config.Site.ID, LogicalDeviceID: config.LogicalDevice.ID, MediaPath: strings.TrimPrefix(config.Media.Ref, "media:"), Revision: config.Revision, PublishedAt: config.PublishedAt, EffectiveAt: config.EffectiveAt, Profile: config.Profile, RuleSetVersion: config.RuleSet.Version, Disabled: config.RuleSet.State == "disabled", NeedsRecalibration: config.RuleSet.State == "recalibration_required", Areas: config.RuleSet.Areas, DirectionalLines: config.RuleSet.DirectionalLines}) + return err +} + +func UTCNow() time.Time { return time.Now().UTC() } diff --git a/Sense/server/app/sense/integration/brain_control/models.go b/Sense/server/app/sense/integration/brain_control/models.go new file mode 100644 index 0000000..1bb83f5 --- /dev/null +++ b/Sense/server/app/sense/integration/brain_control/models.go @@ -0,0 +1,122 @@ +package brain_control + +import "time" + +const ( + SourceConfigVersion = "yovision.source-config/v1" + RuntimeStatusVersion = "yovision.runtime-status/v1" +) + +type Point struct { + X float64 `json:"x"` + Y float64 `json:"y"` +} +type Profile struct { + ID string `json:"id"` + Width int `json:"width"` + Height int `json:"height"` + Encoding string `json:"encoding"` + FrameRate float64 `json:"frame_rate"` +} +type AreaRule struct { + ID string `json:"id"` + Version int64 `json:"version"` + Kind string `json:"kind"` + Enabled bool `json:"enabled"` + Points []Point `json:"points"` +} +type DirectionalLineRule struct { + ID string `json:"id"` + Version int64 `json:"version"` + Kind string `json:"kind"` + Enabled bool `json:"enabled"` + Start Point `json:"start"` + End Point `json:"end"` + TriggerDirection string `json:"trigger_direction"` +} + +type SourceConfig struct { + SchemaVersion string `json:"schema_version"` + ConfigID string `json:"config_id"` + Revision int64 `json:"revision"` + PublishedAt time.Time `json:"published_at"` + EffectiveAt time.Time `json:"effective_at"` + Site struct { + ID string `json:"id"` + } `json:"site"` + LogicalDevice struct { + ID string `json:"id"` + } `json:"logical_device"` + Profile Profile `json:"profile"` + Media struct { + Ref string `json:"ref"` + Transport string `json:"transport"` + } `json:"media"` + RuleSet struct { + Version string `json:"version"` + State string `json:"state"` + ProfileBinding struct { + ProfileID string `json:"profile_id"` + Width int `json:"width"` + Height int `json:"height"` + } `json:"profile_binding"` + Areas []AreaRule `json:"areas"` + DirectionalLines []DirectionalLineRule `json:"directional_lines"` + } `json:"rule_set"` + Integrity struct { + Algorithm string `json:"algorithm"` + Value string `json:"value"` + } `json:"integrity"` +} + +// SourceFacts is an explicit, credential-free boundary DTO. Callers map their +// GORM entities into it; database models are never serialized as a contract. +type SourceFacts struct { + ConfigID, SiteID, LogicalDeviceID, MediaRouteID, MediaPath string + Revision int64 + PublishedAt, EffectiveAt time.Time + Profile Profile + RuleSetVersion string + Disabled, NeedsRecalibration bool + Areas []AreaRule + DirectionalLines []DirectionalLineRule +} + +type ReplayToken struct { + Principal string `gorm:"size:128;primaryKey"` + TokenID string `gorm:"size:96;primaryKey"` + ExpiresAt time.Time `gorm:"not null;index"` + CreatedAt time.Time `gorm:"not null"` +} + +func (ReplayToken) TableName() string { return "sense_brain_runtime_replay_tokens" } + +type RuntimeProjection struct { + BrainInstanceRef string `gorm:"size:128;primaryKey"` + StatusID string `gorm:"size:36;not null;uniqueIndex"` + Sequence int64 `gorm:"not null"` + ObservedAt time.Time `gorm:"not null;index"` + ReceivedAt time.Time `gorm:"not null"` + RuntimeState string `gorm:"size:32;not null"` + RuntimeVersion string `gorm:"size:64;not null"` + ModelRef string `gorm:"size:128;not null"` + ModelVersion string `gorm:"size:64;not null"` + HealthOverall string `gorm:"size:32;not null"` + ExpectedRevisionsJSON string `gorm:"type:jsonb;not null"` + ConfigurationsJSON string `gorm:"type:jsonb;not null"` + HealthJSON string `gorm:"type:jsonb;not null"` + InputsJSON string `gorm:"type:jsonb;not null"` + WasOffline bool `gorm:"not null;default:false"` + CreatedAt time.Time + UpdatedAt time.Time +} + +func (RuntimeProjection) TableName() string { return "sense_brain_runtime_projections" } + +type SourceRevision struct { + ConfigID string `gorm:"size:128;primaryKey"` + Revision int64 `gorm:"not null"` + UpdatedAt time.Time +} + +func (SourceRevision) TableName() string { return "sense_brain_source_revisions" } diff --git a/Sense/server/app/sense/integration/brain_control/runtime.go b/Sense/server/app/sense/integration/brain_control/runtime.go new file mode 100644 index 0000000..fbdf60d --- /dev/null +++ b/Sense/server/app/sense/integration/brain_control/runtime.go @@ -0,0 +1,145 @@ +package brain_control + +import ( + "bytes" + "encoding/json" + "errors" + "io" + "regexp" + "time" +) + +type runtimeStatus struct { + SchemaVersion string `json:"schema_version"` + StatusID string `json:"status_id"` + BrainInstanceRef string `json:"brain_instance_ref"` + Sequence int64 `json:"sequence"` + ObservedAt time.Time `json:"observed_at"` + Runtime struct { + State string `json:"state"` + Version string `json:"version"` + StartedAt *time.Time `json:"started_at"` + } `json:"runtime"` + Model struct { + ModelRef string `json:"model_ref"` + Version string `json:"version"` + } `json:"model"` + Configurations []configurationStatus `json:"configurations"` + Health healthStatus `json:"health"` + Inputs []inputStatus `json:"inputs"` +} +type configurationStatus struct { + ConfigID string `json:"config_id"` + ApplyState string `json:"apply_state"` + AppliedRevision *int64 `json:"applied_revision"` + ErrorCode *string `json:"error_code"` +} +type metrics struct { + LoadPercent float64 `json:"load_percent"` + QueueDepth int64 `json:"queue_depth"` + LatencyMS float64 `json:"latency_ms"` +} +type healthStatus struct { + Overall string `json:"overall"` + ErrorCodes []string `json:"error_codes"` + Metrics metrics `json:"metrics"` +} +type inputStatus struct { + InputRef string `json:"input_ref"` + State string `json:"state"` + ErrorCodes []string `json:"error_codes"` + Metrics metrics `json:"metrics"` +} + +var uuid4 = regexp.MustCompile(`^[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$`) +var errCode = regexp.MustCompile(`^[A-Z][A-Z0-9_]{2,63}$`) + +func parseRuntimeStatus(raw []byte) (runtimeStatus, error) { + var s runtimeStatus + d := json.NewDecoder(bytes.NewReader(raw)) + d.DisallowUnknownFields() + if err := d.Decode(&s); err != nil { + return s, errors.New("CONFIG_INVALID") + } + if err := d.Decode(&struct{}{}); !errors.Is(err, io.EOF) { + return s, errors.New("CONFIG_INVALID") + } + if s.SchemaVersion != RuntimeStatusVersion { + return s, errors.New("UNSUPPORTED_SCHEMA_VERSION") + } + if !uuid4.MatchString(s.StatusID) || !stableID.MatchString(s.BrainInstanceRef) || s.Sequence < 0 || s.ObservedAt.IsZero() || !validState(s.Runtime.State) || s.Runtime.Version == "" || !stableID.MatchString(s.Model.ModelRef) || s.Model.Version == "" || len(s.Configurations) > 4096 || len(s.Inputs) > 4096 || !validHealth(s.Health) { + return s, errors.New("CONFIG_INVALID") + } + seen := map[string]bool{} + for _, c := range s.Configurations { + if !stableID.MatchString(c.ConfigID) || seen[c.ConfigID] || !validApply(c) { + return s, errors.New("CONFIG_INVALID") + } + seen[c.ConfigID] = true + } + for _, i := range s.Inputs { + if !stableID.MatchString(i.InputRef) || !validState(i.State) || !validMetrics(i.Metrics) || !validCodes(i.ErrorCodes, 16) { + return s, errors.New("CONFIG_INVALID") + } + } + return s, nil +} +func validState(v string) bool { + switch v { + case "unconfigured", "starting", "running", "degraded", "failed", "stopped": + return true + } + return false +} +func validApply(c configurationStatus) bool { + switch c.ApplyState { + case "not_configured": + return c.AppliedRevision == nil + case "applying": + return true + case "applied": + return c.AppliedRevision != nil && *c.AppliedRevision >= 1 + case "rejected": + return c.ErrorCode != nil && errCode.MatchString(*c.ErrorCode) + } + return false +} +func validMetrics(m metrics) bool { + return m.LoadPercent >= 0 && m.LoadPercent <= 100 && m.QueueDepth >= 0 && m.LatencyMS >= 0 +} +func validCodes(v []string, n int) bool { + if len(v) > n { + return false + } + seen := map[string]bool{} + for _, x := range v { + if seen[x] || !errCode.MatchString(x) { + return false + } + seen[x] = true + } + return true +} +func validHealth(h healthStatus) bool { + return (h.Overall == "healthy" || h.Overall == "degraded" || h.Overall == "unhealthy") && validCodes(h.ErrorCodes, 32) && validMetrics(h.Metrics) +} +func validRuntimeTransition(from, to string) bool { + if from == to { + return true + } + switch from { + case "unconfigured": + return to == "starting" || to == "stopped" + case "starting": + return to == "running" || to == "degraded" || to == "failed" || to == "stopped" + case "running": + return to == "degraded" || to == "failed" || to == "stopped" + case "degraded": + return to == "running" || to == "failed" || to == "stopped" + case "failed": + return to == "starting" || to == "stopped" + case "stopped": + return to == "starting" + } + return false +} diff --git a/Sense/server/app/sense/integration/brain_control/store.go b/Sense/server/app/sense/integration/brain_control/store.go new file mode 100644 index 0000000..f35ed40 --- /dev/null +++ b/Sense/server/app/sense/integration/brain_control/store.go @@ -0,0 +1,159 @@ +package brain_control + +import ( + "encoding/json" + "errors" + "time" + + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +type GORMReplayStore struct{ DB *gorm.DB } + +func (s GORMReplayStore) Consume(principal, tokenID string, expiresAt, now time.Time) bool { + if s.DB == nil || principal == "" || tokenID == "" || !expiresAt.After(now) { + return false + } + return s.DB.Transaction(func(tx *gorm.DB) error { + if err := tx.Where("expires_at <= ?", now).Delete(&ReplayToken{}).Error; err != nil { + return err + } + result := tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&ReplayToken{Principal: principal, TokenID: tokenID, ExpiresAt: expiresAt, CreatedAt: now}) + if result.Error != nil { + return result.Error + } + if result.RowsAffected != 1 { + return errors.New("replayed") + } + return nil + }) == nil +} + +type RevisionStore struct{ DB *gorm.DB } + +func (s RevisionStore) Next(configID string) (int64, error) { + if s.DB == nil || !stableID.MatchString(configID) { + return 0, errors.New("invalid revision store") + } + var next int64 + err := s.DB.Transaction(func(tx *gorm.DB) error { + var row SourceRevision + err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("config_id = ?", configID).First(&row).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + row = SourceRevision{ConfigID: configID, Revision: 1} + if err = tx.Create(&row).Error; err != nil { + return err + } + next = 1 + return nil + } + if err != nil { + return err + } + row.Revision++ + next = row.Revision + return tx.Save(&row).Error + }) + return next, err +} + +type ProjectionStore struct { + DB *gorm.DB + Clock func() time.Time + StaleAfter time.Duration + FutureSkew time.Duration +} +type ProjectionView struct { + Projection RuntimeProjection + Offline, Stale, Recovered, RevisionMismatch bool +} + +func (s ProjectionStore) Ingest(raw []byte, expected map[string]int64) (ProjectionView, error) { + if s.DB == nil { + return ProjectionView{}, errors.New("projection database required") + } + now := time.Now().UTC() + if s.Clock != nil { + now = s.Clock().UTC() + } + if s.StaleAfter == 0 { + s.StaleAfter = 90 * time.Second + } + if s.FutureSkew == 0 { + s.FutureSkew = 30 * time.Second + } + status, err := parseRuntimeStatus(raw) + if err != nil { + return ProjectionView{}, err + } + if status.ObservedAt.After(now.Add(s.FutureSkew)) { + return ProjectionView{}, errors.New("FUTURE_OBSERVATION") + } + var view ProjectionView + err = s.DB.Transaction(func(tx *gorm.DB) error { + var old RuntimeProjection + find := tx.Where("brain_instance_ref = ?", status.BrainInstanceRef).First(&old).Error + if find == nil { + if old.StatusID == status.StatusID { + view.Projection = old + return nil + } + if status.Sequence <= old.Sequence { + return errors.New("OUT_OF_ORDER_STATUS") + } + if !validRuntimeTransition(old.RuntimeState, status.Runtime.State) { + return errors.New("INVALID_STATUS_TRANSITION") + } + } + if find != nil && !errors.Is(find, gorm.ErrRecordNotFound) { + return find + } + expectedJSON, _ := json.Marshal(expected) + configs, _ := json.Marshal(status.Configurations) + health, _ := json.Marshal(status.Health) + inputs, _ := json.Marshal(status.Inputs) + p := RuntimeProjection{BrainInstanceRef: status.BrainInstanceRef, StatusID: status.StatusID, Sequence: status.Sequence, ObservedAt: status.ObservedAt, ReceivedAt: now, RuntimeState: status.Runtime.State, RuntimeVersion: status.Runtime.Version, ModelRef: status.Model.ModelRef, ModelVersion: status.Model.Version, HealthOverall: status.Health.Overall, ExpectedRevisionsJSON: string(expectedJSON), ConfigurationsJSON: string(configs), HealthJSON: string(health), InputsJSON: string(inputs), WasOffline: find == nil && now.Sub(old.ObservedAt) > s.StaleAfter} + if find == nil { + p.CreatedAt = old.CreatedAt + view.Recovered = p.WasOffline + } + if err := tx.Save(&p).Error; err != nil { + return err + } + view.Projection = p + return nil + }) + if err != nil { + return ProjectionView{}, err + } + view.Stale = now.Sub(view.Projection.ObservedAt) > s.StaleAfter + view.Offline = view.Stale + applied := make(map[string]*int64, len(status.Configurations)) + for _, c := range status.Configurations { + applied[c.ConfigID] = c.AppliedRevision + } + for configID, want := range expected { + got, ok := applied[configID] + if !ok || got == nil || *got != want { + view.RevisionMismatch = true + } + } + return view, nil +} + +func (s ProjectionStore) View(instance string) (ProjectionView, error) { + var p RuntimeProjection + if err := s.DB.First(&p, "brain_instance_ref = ?", instance).Error; err != nil { + return ProjectionView{}, err + } + now := time.Now().UTC() + if s.Clock != nil { + now = s.Clock().UTC() + } + stale := s.StaleAfter + if stale == 0 { + stale = 90 * time.Second + } + return ProjectionView{Projection: p, Offline: now.Sub(p.ObservedAt) > stale, Stale: now.Sub(p.ObservedAt) > stale}, nil +} diff --git a/Sense/server/cmd/migrate/migration/version/2026083112000_brain_runtime.go b/Sense/server/cmd/migrate/migration/version/2026083112000_brain_runtime.go new file mode 100644 index 0000000..7b18fe3 --- /dev/null +++ b/Sense/server/cmd/migrate/migration/version/2026083112000_brain_runtime.go @@ -0,0 +1,25 @@ +package version + +import ( + "runtime" + + "gorm.io/gorm" + + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/integration/brain_control" + "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration" + common "git.ilapage.cn/ila/yovision/Sense/server/common/models" +) + +func init() { + _, fileName, _, _ := runtime.Caller(0) + migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateBrainRuntime) +} + +func migrateBrainRuntime(db *gorm.DB, version string) error { + return db.Transaction(func(tx *gorm.DB) error { + if err := tx.AutoMigrate(&brain_control.ReplayToken{}, &brain_control.RuntimeProjection{}, &brain_control.SourceRevision{}); err != nil { + return err + } + return tx.Create(&common.Migration{Version: version}).Error + }) +} diff --git a/Sense/server/cmd/migrate/migration/version/2026083112000_brain_runtime_test.go b/Sense/server/cmd/migrate/migration/version/2026083112000_brain_runtime_test.go new file mode 100644 index 0000000..9b9fff6 --- /dev/null +++ b/Sense/server/cmd/migrate/migration/version/2026083112000_brain_runtime_test.go @@ -0,0 +1,28 @@ +package version + +import ( + "testing" + + "git.ilapage.cn/ila/yovision/Sense/server/app/sense/integration/brain_control" + common "git.ilapage.cn/ila/yovision/Sense/server/common/models" + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +func TestMigrateBrainRuntime(t *testing.T) { + db, err := gorm.Open(sqlite.Open("file:brain-runtime-migration?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err = db.AutoMigrate(&common.Migration{}); err != nil { + t.Fatal(err) + } + if err = migrateBrainRuntime(db, "2026083112000"); err != nil { + t.Fatal(err) + } + for _, model := range []any{&brain_control.ReplayToken{}, &brain_control.RuntimeProjection{}, &brain_control.SourceRevision{}} { + if !db.Migrator().HasTable(model) { + t.Fatalf("missing table for %T", model) + } + } +} diff --git a/Sense/tests/integration/brain_control/connector_test.go b/Sense/tests/integration/brain_control/connector_test.go new file mode 100644 index 0000000..d0039b4 --- /dev/null +++ b/Sense/tests/integration/brain_control/connector_test.go @@ -0,0 +1,138 @@ +package brain_control_test + +import ( + "crypto/ed25519" + "crypto/rand" + "encoding/json" + "strings" + "sync" + "testing" + "time" + + bc "git.ilapage.cn/ila/yovision/Sense/server/app/sense/integration/brain_control" + mi "git.ilapage.cn/ila/yovision/Sense/server/app/sense/integration/machine_identity" + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +func testDB(t *testing.T, name string) *gorm.DB { + t.Helper() + db, err := gorm.Open(sqlite.Open("file:"+name+"?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err = db.AutoMigrate(&bc.ReplayToken{}, &bc.RuntimeProjection{}, &bc.SourceRevision{}); err != nil { + t.Fatal(err) + } + return db +} + +func TestMapperProducesCredentialFreeFrozenContract(t *testing.T) { + now := time.Date(2026, 8, 31, 0, 0, 0, 0, time.UTC) + c, err := bc.MapSourceConfig(bc.SourceFacts{ConfigID: "gate-primary", SiteID: "site-east", LogicalDeviceID: "camera-1", MediaPath: "site-east/camera-1/main", Revision: 1, PublishedAt: now, EffectiveAt: now, Profile: bc.Profile{ID: "main", Width: 1920, Height: 1080, Encoding: "h264", FrameRate: 25}, RuleSetVersion: "rules-1", Areas: []bc.AreaRule{{ID: "danger", Version: 1, Kind: "danger_area", Enabled: true, Points: []bc.Point{{X: .1, Y: .1}, {X: .8, Y: .1}, {X: .5, Y: .8}}}}}) + if err != nil { + t.Fatal(err) + } + if err = bc.ValidateSourceConfig(c); err != nil { + t.Fatal(err) + } + raw, _ := json.Marshal(c) + text := strings.ToLower(string(raw)) + for _, secret := range []string{"password", "username", "rtsp://", "stream_uri", "credential"} { + if strings.Contains(text, secret) { + t.Fatalf("leaked %q", secret) + } + } + c2, err := bc.MapSourceConfig(bc.SourceFacts{ConfigID: "gate-primary", SiteID: "site-east", LogicalDeviceID: "camera-1", MediaPath: "site-east/camera-1/main", Revision: 2, PublishedAt: now, EffectiveAt: now, Profile: bc.Profile{ID: "main-v2", Width: 1280, Height: 720, Encoding: "H265", FrameRate: 20}, RuleSetVersion: "rules-2", NeedsRecalibration: true, Areas: []bc.AreaRule{{ID: "danger", Version: 2, Kind: "danger_area", Enabled: true, Points: []bc.Point{{X: .1, Y: .1}, {X: .8, Y: .1}, {X: .5, Y: .8}}}}}) + if err != nil { + t.Fatal(err) + } + if c2.RuleSet.State != "recalibration_required" || c2.RuleSet.Areas[0].Enabled { + t.Fatal("recalibration must disable rules") + } +} + +func TestDurableReplayIsAtomicAndSurvivesVerifierRestart(t *testing.T) { + db := testDB(t, "sense-replay") + pub, priv, _ := ed25519.GenerateKey(rand.Reader) + now := time.Date(2026, 8, 31, 0, 0, 0, 0, time.UTC) + registry, _ := mi.NewRegistry(mi.KeyRecord{Principal: "yv:brain:east", KeyID: "brain-key-01", PublicKey: pub, Audience: "yovision-sense", Scopes: []string{"runtime-status:write"}, Enabled: true}) + signer := mi.Signer{Principal: "yv:brain:east", KeyID: "brain-key-01", PrivateKey: priv, Now: func() time.Time { return now }} + body := []byte(`{"ok":true}`) + token, _ := signer.Mint("yovision-sense", []string{"runtime-status:write"}, "POST", "/machine/v1/runtime-status", body) + results := make(chan bool, 8) + var wg sync.WaitGroup + for i := 0; i < 8; i++ { + wg.Add(1) + go func() { + defer wg.Done() + v := mi.Verifier{Registry: registry, Replay: bc.GORMReplayStore{DB: db}, Now: func() time.Time { return now }} + _, err := v.Verify(token, "yovision-sense", "runtime-status:write", "POST", "/machine/v1/runtime-status", body) + results <- err == nil + }() + } + wg.Wait() + close(results) + accepted := 0 + for ok := range results { + if ok { + accepted++ + } + } + if accepted != 1 { + t.Fatalf("accepted=%d", accepted) + } + v2 := mi.Verifier{Registry: registry, Replay: bc.GORMReplayStore{DB: db}, Now: func() time.Time { return now }} + if _, err := v2.Verify(token, "yovision-sense", "runtime-status:write", "POST", "/machine/v1/runtime-status", body); err == nil { + t.Fatal("replay accepted after verifier restart") + } +} + +func TestRevisionStoreConcurrentAndRestart(t *testing.T) { + db := testDB(t, "sense-revisions") + store := bc.RevisionStore{DB: db} + for want := int64(1); want <= 3; want++ { + got, err := store.Next("gate-primary") + if err != nil || got != want { + t.Fatalf("got %d err %v", got, err) + } + } + restarted := bc.RevisionStore{DB: db} + got, err := restarted.Next("gate-primary") + if err != nil || got != 4 { + t.Fatalf("restart got %d err %v", got, err) + } +} + +func TestRuntimeProjectionStaleRecoveryMismatchAndOrdering(t *testing.T) { + db := testDB(t, "sense-projection") + now := time.Date(2026, 8, 31, 0, 0, 0, 0, time.UTC) + store := bc.ProjectionStore{DB: db, Clock: func() time.Time { return now }, StaleAfter: 90 * time.Second} + raw := statusJSON("018f4d6a-8d1b-4a25-8b37-9085f9c0d101", 41, now, "running", 20) + view, err := store.Ingest(raw, map[string]int64{"gate-primary": 21}) + if err != nil { + t.Fatal(err) + } + if !view.RevisionMismatch || view.Stale { + t.Fatalf("bad initial view %+v", view) + } + now = now.Add(91 * time.Second) + view, err = store.View("brain-east-01") + if err != nil || !view.Offline || !view.Stale { + t.Fatalf("offline %+v %v", view, err) + } + raw = statusJSON("018f4d6a-8d1b-4a25-8b37-9085f9c0d102", 42, now, "running", 21) + view, err = store.Ingest(raw, map[string]int64{"gate-primary": 21}) + if err != nil || !view.Recovered || view.RevisionMismatch { + t.Fatalf("recovery %+v %v", view, err) + } + if _, err = store.Ingest(statusJSON("018f4d6a-8d1b-4a25-8b37-9085f9c0d103", 41, now, "running", 21), nil); err == nil || err.Error() != "OUT_OF_ORDER_STATUS" { + t.Fatalf("expected ordering rejection: %v", err) + } +} + +func statusJSON(id string, seq int64, observed time.Time, state string, revision int64) []byte { + v := map[string]any{"schema_version": bc.RuntimeStatusVersion, "status_id": id, "brain_instance_ref": "brain-east-01", "sequence": seq, "observed_at": observed.Format(time.RFC3339), "runtime": map[string]any{"state": state, "version": "1.0.0", "started_at": observed.Add(-time.Minute).Format(time.RFC3339)}, "model": map[string]any{"model_ref": "people-detection", "version": "2026.08.1"}, "configurations": []any{map[string]any{"config_id": "gate-primary", "apply_state": "applied", "applied_revision": revision, "error_code": nil}}, "health": map[string]any{"overall": "healthy", "error_codes": []any{}, "metrics": map[string]any{"load_percent": 1.0, "queue_depth": 0, "latency_ms": 2.0}}, "inputs": []any{}} + raw, _ := json.Marshal(v) + return raw +} diff --git a/Sense/tests/integration/brain_control/go.mod b/Sense/tests/integration/brain_control/go.mod new file mode 100644 index 0000000..455114d --- /dev/null +++ b/Sense/tests/integration/brain_control/go.mod @@ -0,0 +1,18 @@ +module git.ilapage.cn/ila/yovision/Sense/tests/integration/brain_control + +go 1.26.5 + +require ( + git.ilapage.cn/ila/yovision/Sense/server v0.0.0 + gorm.io/driver/sqlite v1.6.0 + gorm.io/gorm v1.31.2 +) + +require ( + github.com/jinzhu/inflection v1.0.0 // indirect + github.com/jinzhu/now v1.1.5 // indirect + github.com/mattn/go-sqlite3 v1.14.49 // indirect + golang.org/x/text v0.40.0 // indirect +) + +replace git.ilapage.cn/ila/yovision/Sense/server => ../../../server diff --git a/Sense/tests/integration/brain_control/go.sum b/Sense/tests/integration/brain_control/go.sum new file mode 100644 index 0000000..6b181ed --- /dev/null +++ b/Sense/tests/integration/brain_control/go.sum @@ -0,0 +1,12 @@ +github.com/jinzhu/inflection v1.0.0 h1:K317FqzuhWc8YvSVlFMCCUb36O/S9MCKRDI7QkRKD/E= +github.com/jinzhu/inflection v1.0.0/go.mod h1:h+uFLlag+Qp1Va5pdKtLDYj+kHp5pxUVkryuEj+Srlc= +github.com/jinzhu/now v1.1.5 h1:/o9tlHleP7gOFmsnYNz3RGnqzefHA47wQpKrrdTIwXQ= +github.com/jinzhu/now v1.1.5/go.mod h1:d3SSVoowX0Lcu0IBviAWJpolVfI5UJVZZ7cO71lE/z8= +github.com/mattn/go-sqlite3 v1.14.49 h1:B8jBHC3xhxZgxztrgruTuLucebnULQnx4W7cF7SAE9w= +github.com/mattn/go-sqlite3 v1.14.49/go.mod h1:6JTjA44L93a0QCyJef5YvlPoKXntQPjzWv5gtm9sB6w= +golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs= +golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY= +gorm.io/driver/sqlite v1.6.0 h1:WHRRrIiulaPiPFmDcod6prc4l2VGVWHz80KspNsxSfQ= +gorm.io/driver/sqlite v1.6.0/go.mod h1:AO9V1qIQddBESngQUKWL9yoH93HIeA1X6V633rBwyT8= +gorm.io/gorm v1.31.2 h1:3o8FXNo9v9S858gil+3LlZA1LkCOzgb4g5BL64FgaCo= +gorm.io/gorm v1.31.2/go.mod h1:XyQVbO2k6YkOis7C2437jSit3SsDK72s7n7rsSHd+Gs=