Compare commits

...
Author SHA1 Message Date
QiuSW adbd1c6aba feat: 串联 Brain 独立纵切与内部事件 (#16) 2026-08-28 23:25:23 +08:00
ila 6ffdbcce84 feat: 实现 Brain 区域与方向越线判定 (#15)
合入 dev,#15 保持待验收。
2026-08-28 23:08:01 +08:00
QiuSW f6f561f2e2 feat: 实现 Brain 区域与方向越线判定 (#15) 2026-08-28 23:07:46 +08:00
ila ee9cfb0433 feat: 实现 Brain 匿名检测与单路跟踪 (#14)
合入 dev,#14 保持待验收。
2026-08-28 22:39:29 +08:00
QiuSW 407ffa17b2 feat: 实现 Brain 匿名检测与单路跟踪 (#14) 2026-08-28 22:39:12 +08:00
ila ba4ec28763 feat: 建立 Brain 可替换视频解码流水线 (#13)
合入 dev,#13 保持待验收。
2026-08-28 22:27:27 +08:00
22 changed files with 1074 additions and 0 deletions
+5
View File
@@ -0,0 +1,5 @@
"""Brain independent vertical-slice application."""
from .runner import RunSummary, run_pipeline
__all__ = ["RunSummary", "run_pipeline"]
+63
View File
@@ -0,0 +1,63 @@
"""CLI for the isolated Brain local-event vertical slice."""
from __future__ import annotations
import argparse
import json
import sys
from pathlib import Path
from yovision_brain.config import ConfigError
from yovision_brain.decode import DecoderError
from yovision_brain.events import JsonLinesSink
from yovision_brain.input import InputError
from yovision_brain.rules import RuleConfigError
from .runner import run_pipeline
def build_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(description="Run the isolated Brain internal-event pipeline")
parser.add_argument("--config", required=True, help="explicit Brain-internal JSON configuration")
parser.add_argument("--output", default="-", help="JSON Lines output file, or - for stdout")
return parser
def main(argv: list[str] | None = None) -> int:
args = build_parser().parse_args(argv)
config_path = Path(args.config)
try:
raw = json.loads(config_path.read_text(encoding="utf-8"))
except (OSError, UnicodeError, json.JSONDecodeError):
print(json.dumps({"status": "error", "message": "Brain internal config cannot be read"}), file=sys.stderr)
return 2
if not isinstance(raw, dict):
print(json.dumps({"status": "error", "message": "Brain internal config must be an object"}), file=sys.stderr)
return 2
stream = sys.stdout
owned_stream = None
try:
if args.output != "-":
try:
owned_stream = Path(args.output).open("w", encoding="utf-8", newline="\n")
except OSError:
print(json.dumps({"status": "error", "message": "event output cannot be opened"}), file=sys.stderr)
return 2
stream = owned_stream
summary = run_pipeline(raw, JsonLinesSink(stream), base_dir=config_path.parent)
except (ConfigError, DecoderError, InputError, RuleConfigError, RuntimeError, ValueError) as exc:
print(json.dumps({"status": "error", "message": str(exc)}, ensure_ascii=False), file=sys.stderr)
return 3
except KeyboardInterrupt:
print(json.dumps({"status": "cancelled"}), file=sys.stderr)
return 130
finally:
if owned_stream is not None:
owned_stream.close()
print(json.dumps({"status": summary.status, "frames": summary.frames, "detections": summary.detections, "events": summary.events}, sort_keys=True), file=sys.stderr)
return 0
if __name__ == "__main__":
raise SystemExit(main())
+110
View File
@@ -0,0 +1,110 @@
"""Compose input, decode, anonymous vision, rules and local event output."""
from __future__ import annotations
from dataclasses import dataclass
from pathlib import Path
from typing import Any, Mapping
from yovision_brain.config import BrainInputConfig, parse_input_config
from yovision_brain.decode import decode_packets
from yovision_brain.events import EventSink, candidate_from_decision
from yovision_brain.input import CancellationToken, build_input_source
from yovision_brain.rules import (
AreaDefinition,
DirectionalLineDefinition,
NormalizedPoint,
RuleEngine,
RuleSet,
)
from yovision_brain.vision import LumaBlobDetector, SingleStreamTracker, TorchLumaBlobDetector
@dataclass(frozen=True, slots=True)
class RunSummary:
status: str
frames: int
detections: int
events: int
def _rules(config: BrainInputConfig, version: str) -> RuleSet:
return RuleSet(
version=version,
profile_id=config.profile.profile_id,
width=config.profile.width,
height=config.profile.height,
areas=tuple(
AreaDefinition(area.rule_id, tuple(NormalizedPoint(point.x, point.y) for point in area.points))
for area in config.areas
),
directional_lines=tuple(
DirectionalLineDefinition(
line.rule_id,
NormalizedPoint(line.start.x, line.start.y),
NormalizedPoint(line.end.x, line.end.y),
line.trigger_direction,
)
for line in config.directional_lines
),
)
def run_pipeline(
raw: Mapping[str, Any],
sink: EventSink,
*,
base_dir: Path | None = None,
cancellation: CancellationToken | None = None,
) -> RunSummary:
input_raw = raw.get("input")
if not isinstance(input_raw, Mapping):
raise ValueError("input must be an object")
config = parse_input_config(input_raw, base_dir=base_dir)
version = raw.get("rules_version")
if not isinstance(version, str) or not version:
raise ValueError("rules_version must be a non-empty string")
detector_raw = raw.get("detector", {})
if not isinstance(detector_raw, Mapping):
raise ValueError("detector must be an object")
backend = detector_raw.get("backend", "python")
threshold = detector_raw.get("threshold", 200)
minimum_area = detector_raw.get("minimum_area", 1)
if not isinstance(threshold, int) or not isinstance(minimum_area, int):
raise ValueError("detector threshold and minimum_area must be integers")
if backend == "python":
detector = LumaBlobDetector(threshold=threshold, minimum_area=minimum_area)
elif backend == "torch_cpu":
detector = TorchLumaBlobDetector(threshold=threshold, minimum_area=minimum_area, device="cpu")
else:
raise ValueError("detector.backend must be python or torch_cpu")
source = build_input_source(config)
tracker = SingleStreamTracker()
engine = RuleEngine(_rules(config, version))
frame_count = detection_count = event_count = 0
for frame in decode_packets(source.packets(cancellation), cancellation):
frame_count += 1
detections = detector.detect(frame)
detection_count += len(detections)
tracks = tracker.update(detections, frame_sequence=frame.sequence, timestamp_ns=frame.timestamp_ns)
decisions = engine.evaluate(
tracks,
profile_id=frame.profile_id,
width=frame.width,
height=frame.height,
)
tracks_by_id = {track.track_id: track for track in tracks}
for decision in decisions:
if not decision.triggered:
continue
sink.write(candidate_from_decision(
decision,
tracks_by_id[decision.track_id],
logical_input_id=config.logical_device_id,
detector=detector.metadata,
))
event_count += 1
tracker.finish()
status = "cancelled" if cancellation is not None and cancellation.cancelled else "completed"
return RunSummary(status, frame_count, detection_count, event_count)
@@ -0,0 +1,13 @@
"""Brain-internal event candidates and local sinks."""
from .mapper import candidate_from_decision
from .models import INTERNAL_EVENT_SCHEMA, InternalEventCandidate
from .sink import EventSink, JsonLinesSink
__all__ = [
"EventSink",
"INTERNAL_EVENT_SCHEMA",
"InternalEventCandidate",
"JsonLinesSink",
"candidate_from_decision",
]
+53
View File
@@ -0,0 +1,53 @@
"""Stable mapping from an internal rule hit to an internal event candidate."""
from __future__ import annotations
import hashlib
import json
from yovision_brain.rules import RuleDecision
from yovision_brain.vision import DetectorMetadata, TrackedObject
from .models import INTERNAL_EVENT_SCHEMA, InternalEventCandidate
def candidate_from_decision(
decision: RuleDecision,
track: TrackedObject,
*,
logical_input_id: str,
detector: DetectorMetadata,
) -> InternalEventCandidate:
if not decision.triggered or decision.track_id != track.track_id:
raise ValueError("only a triggered decision for the same anonymous track can become an event")
event_type = "danger_area_entered" if decision.rule_type == "danger_area" else "directional_line_crossed"
observation: dict[str, object] = {
"category": track.category,
"confidence": track.confidence,
"box": {
"left": track.box.left,
"top": track.box.top,
"right": track.box.right,
"bottom": track.box.bottom,
},
"anchor": {"x": decision.anchor.x, "y": decision.anchor.y},
}
fact = {
"schema": INTERNAL_EVENT_SCHEMA,
"logical_input_id": logical_input_id,
"event_type": event_type,
"occurred_at_ns": decision.timestamp_ns,
"rule_id": decision.rule_id,
"rule_version": decision.config_version,
"model_name": detector.name,
"model_version": detector.version,
"profile_id": decision.profile_id,
"frame_width": decision.width,
"frame_height": decision.height,
"track_id": track.track_id,
"observation": observation,
"reason": decision.reason,
}
canonical = json.dumps(fact, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
event_id = "brain-local-" + hashlib.sha256(canonical.encode("utf-8")).hexdigest()
return InternalEventCandidate(event_id=event_id, **fact)
+29
View File
@@ -0,0 +1,29 @@
"""Explicitly internal event candidate model; not a shared contract."""
from __future__ import annotations
from dataclasses import asdict, dataclass
INTERNAL_EVENT_SCHEMA = "brain.internal.event-candidate/v1"
@dataclass(frozen=True, slots=True)
class InternalEventCandidate:
schema: str
event_id: str
logical_input_id: str
event_type: str
occurred_at_ns: int
rule_id: str
rule_version: str
model_name: str
model_version: str
profile_id: str
frame_width: int
frame_height: int
track_id: str
observation: dict[str, object]
reason: str
def to_dict(self) -> dict[str, object]:
return asdict(self)
+21
View File
@@ -0,0 +1,21 @@
"""Replaceable local event sinks."""
from __future__ import annotations
import json
from typing import Protocol, TextIO
from .models import InternalEventCandidate
class EventSink(Protocol):
def write(self, candidate: InternalEventCandidate) -> None: ...
class JsonLinesSink:
def __init__(self, stream: TextIO) -> None:
self._stream = stream
def write(self, candidate: InternalEventCandidate) -> None:
self._stream.write(json.dumps(candidate.to_dict(), ensure_ascii=False, sort_keys=True) + "\n")
self._stream.flush()
@@ -0,0 +1,21 @@
"""Brain-internal anonymous area and directional-line rules."""
from .engine import RuleEngine
from .models import (
AreaDefinition,
DirectionalLineDefinition,
NormalizedPoint,
RuleConfigError,
RuleDecision,
RuleSet,
)
__all__ = [
"AreaDefinition",
"DirectionalLineDefinition",
"NormalizedPoint",
"RuleConfigError",
"RuleDecision",
"RuleEngine",
"RuleSet",
]
+112
View File
@@ -0,0 +1,112 @@
"""Stateful, explainable area and directional-line evaluation."""
from __future__ import annotations
from yovision_brain.vision import TrackedObject
from .models import NormalizedPoint, RuleConfigError, RuleDecision, RuleSet
_EPSILON = 1e-9
def _anchor(track: TrackedObject, width: int, height: int) -> NormalizedPoint:
x = (track.box.left + track.box.right) / (2.0 * width)
y = track.box.bottom / height
try:
return NormalizedPoint(x, y)
except RuleConfigError as exc:
raise RuleConfigError(f"track {track.track_id!r} anchor is outside the configured frame") from exc
def _on_segment(point: NormalizedPoint, first: NormalizedPoint, second: NormalizedPoint) -> bool:
cross = (second.x - first.x) * (point.y - first.y) - (second.y - first.y) * (point.x - first.x)
return abs(cross) <= _EPSILON and min(first.x, second.x) - _EPSILON <= point.x <= max(first.x, second.x) + _EPSILON and min(first.y, second.y) - _EPSILON <= point.y <= max(first.y, second.y) + _EPSILON
def _inside(point: NormalizedPoint, polygon: tuple[NormalizedPoint, ...]) -> bool:
inside = False
previous = polygon[-1]
for current in polygon:
if _on_segment(point, previous, current):
return True
if (current.y > point.y) != (previous.y > point.y):
crossing_x = (previous.x - current.x) * (point.y - current.y) / (previous.y - current.y) + current.x
if point.x < crossing_x:
inside = not inside
previous = current
return inside
def _side(point: NormalizedPoint, start: NormalizedPoint, end: NormalizedPoint) -> float:
return (end.x - start.x) * (point.y - start.y) - (end.y - start.y) * (point.x - start.x)
class RuleEngine:
"""Evaluates one versioned rule set against one stream session."""
def __init__(self, rules: RuleSet) -> None:
self._rules = rules
self._area_inside: dict[tuple[str, str], bool] = {}
self._line_side: dict[tuple[str, str], int] = {}
def evaluate(
self,
tracks: tuple[TrackedObject, ...],
*,
profile_id: str,
width: int,
height: int,
) -> tuple[RuleDecision, ...]:
if (profile_id, width, height) != (self._rules.profile_id, self._rules.width, self._rules.height):
raise RuleConfigError("track Profile/resolution does not match the versioned rule configuration")
decisions: list[RuleDecision] = []
for track in tracks:
anchor = _anchor(track, width, height)
common = dict(
track_id=track.track_id,
config_version=self._rules.version,
profile_id=profile_id,
width=width,
height=height,
anchor=anchor,
timestamp_ns=track.timestamp_ns,
)
for area in self._rules.areas:
key = (track.track_id, area.rule_id)
current = _inside(anchor, area.points)
previous = self._area_inside.get(key, False)
state = "entered" if current and not previous else "inside" if current else "outside"
self._area_inside[key] = current
decisions.append(RuleDecision(
rule_id=area.rule_id,
rule_type="danger_area",
state=state,
triggered=state == "entered",
reason=f"bottom-center anchor is {state} the configured polygon",
**common,
))
for line in self._rules.directional_lines:
key = (track.track_id, line.rule_id)
value = _side(anchor, line.start, line.end)
if abs(value) <= line.deadband:
decisions.append(RuleDecision(
rule_id=line.rule_id, rule_type="directional_line", state="on_line",
triggered=False, reason="anchor is inside the line deadband; previous significant side is retained",
**common,
))
continue
current_side = 1 if value > 0 else -1
previous_side = self._line_side.get(key)
self._line_side[key] = current_side
wanted = (previous_side, current_side) == ((1, -1) if line.trigger_direction == "left_to_right" else (-1, 1))
crossed = previous_side is not None and previous_side != current_side
state = "triggered" if wanted else "reverse_crossing" if crossed else "same_side"
decisions.append(RuleDecision(
rule_id=line.rule_id,
rule_type="directional_line",
state=state,
triggered=wanted,
reason=f"directed side transition {previous_side!r}->{current_side}; expected {line.trigger_direction}",
**common,
))
return tuple(decisions)
+84
View File
@@ -0,0 +1,84 @@
"""Versioned Brain-internal rule configuration and decisions."""
from __future__ import annotations
from dataclasses import dataclass
class RuleConfigError(ValueError):
pass
@dataclass(frozen=True, slots=True)
class NormalizedPoint:
x: float
y: float
def __post_init__(self) -> None:
if not 0.0 <= self.x <= 1.0 or not 0.0 <= self.y <= 1.0:
raise RuleConfigError("rule coordinates must be normalized to 0..1")
@dataclass(frozen=True, slots=True)
class AreaDefinition:
rule_id: str
points: tuple[NormalizedPoint, ...]
@dataclass(frozen=True, slots=True)
class DirectionalLineDefinition:
rule_id: str
start: NormalizedPoint
end: NormalizedPoint
trigger_direction: str
deadband: float = 0.005
@dataclass(frozen=True, slots=True)
class RuleSet:
version: str
profile_id: str
width: int
height: int
areas: tuple[AreaDefinition, ...] = ()
directional_lines: tuple[DirectionalLineDefinition, ...] = ()
def __post_init__(self) -> None:
if not self.version or not self.profile_id or self.width <= 0 or self.height <= 0:
raise RuleConfigError("rule version, profile and dimensions are required")
identifiers = [rule.rule_id for rule in self.areas] + [rule.rule_id for rule in self.directional_lines]
if any(not identifier for identifier in identifiers) or len(set(identifiers)) != len(identifiers):
raise RuleConfigError("rule ids must be non-empty and unique")
for area in self.areas:
if len(area.points) < 3 or abs(_polygon_area(area.points)) < 1e-9:
raise RuleConfigError(f"area {area.rule_id!r} must be a non-degenerate polygon")
for line in self.directional_lines:
if line.start == line.end:
raise RuleConfigError(f"line {line.rule_id!r} must have distinct endpoints")
if line.trigger_direction not in {"left_to_right", "right_to_left"}:
raise RuleConfigError(f"line {line.rule_id!r} has invalid trigger direction")
if not 0.0 <= line.deadband < 0.5:
raise RuleConfigError(f"line {line.rule_id!r} has invalid deadband")
def _polygon_area(points: tuple[NormalizedPoint, ...]) -> float:
return sum(
first.x * second.y - second.x * first.y
for first, second in zip(points, points[1:] + points[:1])
) / 2.0
@dataclass(frozen=True, slots=True)
class RuleDecision:
rule_id: str
rule_type: str
track_id: str
state: str
triggered: bool
reason: str
config_version: str
profile_id: str
width: int
height: int
anchor: NormalizedPoint
timestamp_ns: int
@@ -0,0 +1,16 @@
"""Anonymous detection and single-stream tracking."""
from .detector import LumaBlobDetector, TorchLumaBlobDetector
from .models import BoundingBox, Detection, Detector, DetectorMetadata, TrackedObject
from .tracker import SingleStreamTracker
__all__ = [
"BoundingBox",
"Detection",
"Detector",
"DetectorMetadata",
"LumaBlobDetector",
"SingleStreamTracker",
"TorchLumaBlobDetector",
"TrackedObject",
]
+120
View File
@@ -0,0 +1,120 @@
"""Deterministic anonymous blob detectors with no biometric semantics."""
from __future__ import annotations
from collections.abc import Sequence
from yovision_brain.decode import DecodedFrame, DecoderError
from .models import BoundingBox, Detection, DetectorMetadata
_METADATA = DetectorMetadata(
name="yovision-luma-blob",
version="1.0.0",
source="YoVision Brain first-party deterministic algorithm",
license="No external model license; no learned weights are distributed",
weights="none",
)
def _components(mask: Sequence[Sequence[bool]], minimum_area: int) -> tuple[BoundingBox, ...]:
height = len(mask)
width = len(mask[0]) if height else 0
visited: set[tuple[int, int]] = set()
boxes: list[BoundingBox] = []
for y in range(height):
for x in range(width):
if not mask[y][x] or (x, y) in visited:
continue
pending = [(x, y)]
visited.add((x, y))
points: list[tuple[int, int]] = []
while pending:
current_x, current_y = pending.pop()
points.append((current_x, current_y))
for neighbor in (
(current_x - 1, current_y),
(current_x + 1, current_y),
(current_x, current_y - 1),
(current_x, current_y + 1),
):
nx, ny = neighbor
if 0 <= nx < width and 0 <= ny < height and mask[ny][nx] and neighbor not in visited:
visited.add(neighbor)
pending.append(neighbor)
if len(points) >= minimum_area:
xs, ys = zip(*points)
boxes.append(BoundingBox(min(xs), min(ys), max(xs) + 1, max(ys) + 1))
return tuple(sorted(boxes, key=lambda box: (box.top, box.left, box.bottom, box.right)))
def _validate_frame(frame: DecodedFrame) -> None:
if frame.pixel_format not in {"rgb24", "yuv444p"}:
raise DecoderError(f"anonymous detector does not support pixel format {frame.pixel_format!r}")
expected = frame.width * frame.height * 3
if len(frame.payload) != expected:
raise DecoderError(f"vision frame has {len(frame.payload)} bytes; expected {expected}")
class LumaBlobDetector:
"""Small CPU reference detector used for deterministic integration tests."""
metadata = _METADATA
def __init__(self, *, threshold: int = 200, minimum_area: int = 1) -> None:
if not 0 <= threshold <= 255 or minimum_area < 1:
raise ValueError("invalid luma detector threshold or minimum area")
self._threshold = threshold
self._minimum_area = minimum_area
def detect(self, frame: DecodedFrame) -> tuple[Detection, ...]:
_validate_frame(frame)
if frame.pixel_format == "rgb24":
pixels = [
max(frame.payload[index : index + 3])
for index in range(0, len(frame.payload), 3)
]
else:
pixels = list(frame.payload[: frame.width * frame.height])
mask = [
[pixels[y * frame.width + x] >= self._threshold for x in range(frame.width)]
for y in range(frame.height)
]
return tuple(
Detection(box=box, category="anonymous_target", confidence=1.0)
for box in _components(mask, self._minimum_area)
)
class TorchLumaBlobDetector:
"""PyTorch CPU/GPU smoke backend; it contains no external model weights."""
metadata = DetectorMetadata(
name="yovision-torch-luma-blob",
version="1.0.0",
source="YoVision Brain first-party PyTorch tensor implementation",
license="PyTorch BSD-3-Clause; no external model weights",
weights="none",
)
def __init__(self, *, threshold: int = 200, minimum_area: int = 1, device: str = "cpu") -> None:
self._threshold = threshold
self._minimum_area = minimum_area
self._device = device
def detect(self, frame: DecodedFrame) -> tuple[Detection, ...]:
_validate_frame(frame)
try:
import torch
except ImportError as exc:
raise RuntimeError("PyTorch runtime is required for TorchLumaBlobDetector") from exc
values = torch.tensor(list(frame.payload), dtype=torch.uint8, device=self._device)
if frame.pixel_format == "rgb24":
luma = values.reshape(frame.height, frame.width, 3).amax(dim=2)
else:
luma = values[: frame.width * frame.height].reshape(frame.height, frame.width)
mask = (luma >= self._threshold).cpu().tolist()
return tuple(
Detection(box=box, category="anonymous_target", confidence=1.0)
for box in _components(mask, self._minimum_area)
)
+52
View File
@@ -0,0 +1,52 @@
"""Privacy-preserving vision ports and observations."""
from __future__ import annotations
from dataclasses import dataclass
from typing import Protocol
from yovision_brain.decode import DecodedFrame
@dataclass(frozen=True, slots=True)
class DetectorMetadata:
name: str
version: str
source: str
license: str
weights: str
@dataclass(frozen=True, slots=True)
class BoundingBox:
left: int
top: int
right: int
bottom: int
@property
def area(self) -> int:
return max(0, self.right - self.left) * max(0, self.bottom - self.top)
@dataclass(frozen=True, slots=True)
class Detection:
box: BoundingBox
category: str
confidence: float
@dataclass(frozen=True, slots=True)
class TrackedObject:
track_id: str
box: BoundingBox
category: str
confidence: float
frame_sequence: int
timestamp_ns: int
class Detector(Protocol):
metadata: DetectorMetadata
def detect(self, frame: DecodedFrame) -> tuple[Detection, ...]: ...
@@ -0,0 +1,83 @@
"""Session-local single-stream IoU tracker."""
from __future__ import annotations
from dataclasses import dataclass
from .models import BoundingBox, Detection, TrackedObject
def _iou(first: BoundingBox, second: BoundingBox) -> float:
intersection = BoundingBox(
max(first.left, second.left),
max(first.top, second.top),
min(first.right, second.right),
min(first.bottom, second.bottom),
).area
union = first.area + second.area - intersection
return intersection / union if union else 0.0
@dataclass(slots=True)
class _Track:
track_id: str
detection: Detection
missed: int = 0
class SingleStreamTracker:
"""Tracks anonymous boxes only within one process session and one stream."""
def __init__(self, *, iou_threshold: float = 0.2, max_missed: int = 2) -> None:
if not 0.0 <= iou_threshold <= 1.0 or max_missed < 0:
raise ValueError("invalid tracker threshold or missed-frame limit")
self._iou_threshold = iou_threshold
self._max_missed = max_missed
self._tracks: dict[str, _Track] = {}
self._next_id = 1
def update(
self,
detections: tuple[Detection, ...],
*,
frame_sequence: int,
timestamp_ns: int,
) -> tuple[TrackedObject, ...]:
unmatched_tracks = set(self._tracks)
results: list[TrackedObject] = []
for detection in detections:
candidates = [
(track_id, _iou(self._tracks[track_id].detection.box, detection.box))
for track_id in unmatched_tracks
if self._tracks[track_id].detection.category == detection.category
]
track_id, score = max(candidates, key=lambda item: item[1], default=("", -1.0))
if score < self._iou_threshold:
track_id = f"track-{self._next_id:06d}"
self._next_id += 1
self._tracks[track_id] = _Track(track_id, detection)
else:
unmatched_tracks.remove(track_id)
self._tracks[track_id].detection = detection
self._tracks[track_id].missed = 0
results.append(
TrackedObject(
track_id=track_id,
box=detection.box,
category=detection.category,
confidence=detection.confidence,
frame_sequence=frame_sequence,
timestamp_ns=timestamp_ns,
)
)
for track_id in unmatched_tracks:
track = self._tracks[track_id]
track.missed += 1
if track.missed > self._max_missed:
del self._tracks[track_id]
return tuple(results)
def finish(self) -> tuple[str, ...]:
ended = tuple(sorted(self._tracks))
self._tracks.clear()
return ended
+68
View File
@@ -0,0 +1,68 @@
from __future__ import annotations
import io
import json
import subprocess
import sys
from pathlib import Path
from yovision_brain.app import run_pipeline
from yovision_brain.events import JsonLinesSink
from yovision_brain.input import CancellationToken
FIXTURE = Path(__file__).parents[1] / "fixtures" / "events" / "area.json"
def test_pipeline_generates_stable_internal_event() -> None:
raw = json.loads(FIXTURE.read_text(encoding="utf-8"))
first_stream, second_stream = io.StringIO(), io.StringIO()
first = run_pipeline(raw, JsonLinesSink(first_stream), base_dir=FIXTURE.parent)
second = run_pipeline(raw, JsonLinesSink(second_stream), base_dir=FIXTURE.parent)
assert first.events == second.events == 1
assert first_stream.getvalue() == second_stream.getvalue()
event = json.loads(first_stream.getvalue())
assert event["event_type"] == "danger_area_entered"
assert event["schema"] == "brain.internal.event-candidate/v1"
def test_no_hit_has_explicit_zero_event_summary() -> None:
raw = json.loads(FIXTURE.read_text(encoding="utf-8"))
raw["detector"]["threshold"] = 255
raw["detector"]["minimum_area"] = 1000
stream = io.StringIO()
summary = run_pipeline(raw, JsonLinesSink(stream))
assert (summary.status, summary.events, stream.getvalue()) == ("completed", 0, "")
def test_cli_runs_without_sense_or_bell() -> None:
result = subprocess.run(
[sys.executable, "-m", "yovision_brain.app", "--config", str(FIXTURE), "--output", "-"],
check=False,
capture_output=True,
text=True,
)
assert result.returncode == 0
assert json.loads(result.stdout)["schema"] == "brain.internal.event-candidate/v1"
assert json.loads(result.stderr)["events"] == 1
def test_pre_cancelled_run_is_explicit() -> None:
raw = json.loads(FIXTURE.read_text(encoding="utf-8"))
token = CancellationToken()
token.cancel()
summary = run_pipeline(raw, JsonLinesSink(io.StringIO()), cancellation=token)
assert (summary.status, summary.frames, summary.events) == ("cancelled", 0, 0)
def test_cli_config_failure_is_nonzero_and_does_not_echo_path(tmp_path: Path) -> None:
missing = tmp_path / "private-machine-path.json"
result = subprocess.run(
[sys.executable, "-m", "yovision_brain.app", "--config", str(missing)],
check=False,
capture_output=True,
text=True,
)
assert result.returncode == 2
assert json.loads(result.stderr)["status"] == "error"
assert str(tmp_path) not in result.stderr
+30
View File
@@ -0,0 +1,30 @@
from __future__ import annotations
import io
import json
from yovision_brain.events import JsonLinesSink, candidate_from_decision
from yovision_brain.rules import NormalizedPoint, RuleDecision
from yovision_brain.vision import BoundingBox, DetectorMetadata, TrackedObject
def test_internal_event_id_is_stable_and_payload_is_safe() -> None:
track = TrackedObject("track-000001", BoundingBox(1, 2, 3, 4), "anonymous_target", 1.0, 7, 123)
decision = RuleDecision("yard", "danger_area", track.track_id, "entered", True, "entered polygon", "rules-v1", "main", 10, 10, NormalizedPoint(0.2, 0.4), 123)
metadata = DetectorMetadata("detector", "1", "first-party", "no external weights", "none")
first = candidate_from_decision(decision, track, logical_input_id="camera-01", detector=metadata)
second = candidate_from_decision(decision, track, logical_input_id="camera-01", detector=metadata)
assert first == second
assert first.event_id.startswith("brain-local-")
payload = json.dumps(first.to_dict())
for forbidden in ("password", "rtsp://", "evidence", "face"):
assert forbidden not in payload.lower()
def test_json_lines_sink_writes_one_canonical_line() -> None:
track = TrackedObject("track-000001", BoundingBox(0, 0, 1, 1), "anonymous_target", 1.0, 0, 0)
decision = RuleDecision("yard", "danger_area", track.track_id, "entered", True, "entered", "v1", "main", 2, 2, NormalizedPoint(0.25, 0.5), 0)
candidate = candidate_from_decision(decision, track, logical_input_id="synthetic", detector=DetectorMetadata("d", "1", "first", "none", "none"))
stream = io.StringIO()
JsonLinesSink(stream).write(candidate)
assert json.loads(stream.getvalue())["schema"] == "brain.internal.event-candidate/v1"
+5
View File
@@ -0,0 +1,5 @@
# Brain internal-event fixtures
These fixtures are synthetic and explicitly internal. They are not the future
Brain-to-Bell event contract and must not contain evidence references, customer
media, identities, credentials, or machine-specific paths.
+29
View File
@@ -0,0 +1,29 @@
{
"rules_version": "fixture-rules-v1",
"detector": {
"backend": "python",
"threshold": 0,
"minimum_area": 1
},
"input": {
"schema": "brain.internal.input/v1",
"logical_device_id": "synthetic-camera-01",
"profile": {
"id": "main",
"width": 4,
"height": 3,
"fps": 5
},
"source": {
"kind": "synthetic",
"seed": 17,
"frame_count": 3
},
"areas": [
{
"id": "full-frame-area",
"points": [[0, 0], [1, 0], [1, 1], [0, 1]]
}
]
}
}
+5
View File
@@ -0,0 +1,5 @@
# Brain rule fixtures
Rule tests use normalized synthetic geometry and anonymous track IDs only. Do
not add customer site layouts, camera paths, identities, credentials, or a
copy of a future cross-project contract.
+5
View File
@@ -0,0 +1,5 @@
# Brain vision fixtures
Vision tests create anonymous geometric RGB frames in memory. Never add faces,
customer recordings, biometric templates, camera credentials, or unreviewed
model weights to this directory.
+81
View File
@@ -0,0 +1,81 @@
from __future__ import annotations
import pytest
from yovision_brain.rules import (
AreaDefinition,
DirectionalLineDefinition,
NormalizedPoint,
RuleConfigError,
RuleEngine,
RuleSet,
)
from yovision_brain.vision import BoundingBox, TrackedObject
def point(x: float, y: float) -> NormalizedPoint:
return NormalizedPoint(x, y)
def rules() -> RuleSet:
return RuleSet(
version="rules-v7",
profile_id="main",
width=100,
height=100,
areas=(AreaDefinition("yard", (point(0.2, 0.2), point(0.8, 0.2), point(0.8, 0.8), point(0.2, 0.8))),),
directional_lines=(DirectionalLineDefinition("gate", point(0.5, 0.1), point(0.5, 0.9), "left_to_right", 0.01),),
)
def track(track_id: str, anchor_x: int, anchor_y: int, sequence: int = 0) -> TrackedObject:
return TrackedObject(track_id, BoundingBox(anchor_x - 1, anchor_y - 2, anchor_x + 1, anchor_y), "anonymous_target", 1.0, sequence, sequence)
def decisions(engine: RuleEngine, item: TrackedObject):
return engine.evaluate((item,), profile_id="main", width=100, height=100)
def test_area_outside_entered_inside_and_boundary() -> None:
engine = RuleEngine(rules())
assert decisions(engine, track("one", 10, 50))[0].state == "outside"
entered = decisions(engine, track("one", 20, 50, 1))[0]
assert (entered.state, entered.triggered) == ("entered", True)
inside = decisions(engine, track("one", 50, 50, 2))[0]
assert (inside.state, inside.triggered) == ("inside", False)
assert inside.config_version == "rules-v7"
def test_direction_and_reverse_crossing_are_distinct() -> None:
engine = RuleEngine(rules())
decisions(engine, track("one", 40, 50))
forward = decisions(engine, track("one", 60, 50, 1))[1]
assert (forward.state, forward.triggered) == ("triggered", True)
reverse_engine = RuleEngine(rules())
decisions(reverse_engine, track("two", 60, 50))
reverse = decisions(reverse_engine, track("two", 40, 50, 1))[1]
assert (reverse.state, reverse.triggered) == ("reverse_crossing", False)
def test_line_deadband_prevents_jitter_trigger() -> None:
engine = RuleEngine(rules())
decisions(engine, track("one", 40, 50))
on_line = decisions(engine, track("one", 50, 50, 1))[1]
assert (on_line.state, on_line.triggered) == ("on_line", False)
triggered = decisions(engine, track("one", 60, 50, 2))[1]
assert triggered.triggered is True
def test_profile_resolution_mismatch_is_rejected() -> None:
with pytest.raises(RuleConfigError, match="Profile/resolution"):
RuleEngine(rules()).evaluate((track("one", 20, 20),), profile_id="sub", width=100, height=100)
def test_invalid_polygon_line_and_duplicate_ids_are_rejected() -> None:
with pytest.raises(RuleConfigError, match="non-degenerate"):
RuleSet("v", "main", 10, 10, areas=(AreaDefinition("bad", (point(0, 0), point(0.5, 0.5), point(1, 1))),))
with pytest.raises(RuleConfigError, match="distinct endpoints"):
RuleSet("v", "main", 10, 10, directional_lines=(DirectionalLineDefinition("bad", point(0, 0), point(0, 0), "left_to_right"),))
with pytest.raises(RuleConfigError, match="unique"):
RuleSet("v", "main", 10, 10, areas=(AreaDefinition("same", (point(0, 0), point(1, 0), point(0, 1))),), directional_lines=(DirectionalLineDefinition("same", point(0, 0), point(1, 1), "left_to_right"),))
@@ -0,0 +1,69 @@
from __future__ import annotations
import pytest
from yovision_brain.decode import DecodedFrame
from yovision_brain.vision import (
BoundingBox,
Detection,
LumaBlobDetector,
SingleStreamTracker,
TorchLumaBlobDetector,
)
def frame(payload: bytes, *, sequence: int = 0, width: int = 4, height: int = 3) -> DecodedFrame:
return DecodedFrame(sequence, sequence * 40_000_000, "camera", "main", width, height, "rgb24", payload)
def rgb(values: list[int]) -> bytes:
return b"".join(bytes((value, value, value)) for value in values)
def detection(left: int, top: int, right: int, bottom: int) -> Detection:
return Detection(BoundingBox(left, top, right, bottom), "anonymous_target", 0.9)
def test_detector_emits_only_anonymous_observations() -> None:
payload = rgb([0, 255, 255, 0, 0, 255, 255, 0, 0, 0, 0, 0])
result = LumaBlobDetector(minimum_area=2).detect(frame(payload))
assert result == (Detection(BoundingBox(1, 0, 3, 2), "anonymous_target", 1.0),)
assert LumaBlobDetector.metadata.weights == "none"
assert "external model license" in LumaBlobDetector.metadata.license
def test_empty_frame_has_no_detection() -> None:
assert LumaBlobDetector().detect(frame(rgb([0] * 12))) == ()
def test_tracker_keeps_session_id_across_motion_and_short_occlusion() -> None:
tracker = SingleStreamTracker(iou_threshold=0.1, max_missed=2)
first = tracker.update((detection(0, 0, 3, 3),), frame_sequence=0, timestamp_ns=0)
assert first[0].track_id == "track-000001"
assert tracker.update((), frame_sequence=1, timestamp_ns=1) == ()
resumed = tracker.update((detection(1, 0, 4, 3),), frame_sequence=2, timestamp_ns=2)
assert resumed[0].track_id == "track-000001"
assert tracker.finish() == ("track-000001",)
def test_disappeared_track_ends_and_new_target_gets_new_id() -> None:
tracker = SingleStreamTracker(max_missed=1)
first = tracker.update((detection(0, 0, 2, 2),), frame_sequence=0, timestamp_ns=0)
tracker.update((), frame_sequence=1, timestamp_ns=1)
tracker.update((), frame_sequence=2, timestamp_ns=2)
second = tracker.update((detection(0, 0, 2, 2),), frame_sequence=3, timestamp_ns=3)
assert first[0].track_id == "track-000001"
assert second[0].track_id == "track-000002"
def test_track_ids_are_session_local() -> None:
one = SingleStreamTracker().update((detection(0, 0, 1, 1),), frame_sequence=0, timestamp_ns=0)
two = SingleStreamTracker().update((detection(0, 0, 1, 1),), frame_sequence=0, timestamp_ns=0)
assert one[0].track_id == two[0].track_id == "track-000001"
def test_torch_backend_cpu_smoke_uses_no_external_weights() -> None:
pytest.importorskip("torch")
result = TorchLumaBlobDetector().detect(frame(rgb([0, 255] + [0] * 10)))
assert result[0].category == "anonymous_target"
assert TorchLumaBlobDetector.metadata.weights == "none"