feat: 串联 Brain 独立纵切与内部事件 (#16) #130

Merged
ila merged 1 commits from feature/16-brain-local-events into dev 2026-08-28 23:25:39 +08:00
11 changed files with 426 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()
+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]]
}
]
}
}