diff --git a/solstone/observe/depict.py b/solstone/observe/depict.py index 7e6cc2ac4..a5a02add8 100644 --- a/solstone/observe/depict.py +++ b/solstone/observe/depict.py @@ -4,6 +4,7 @@ """Segment still-image description handler.""" import argparse +import io import json import logging import os @@ -11,6 +12,7 @@ from pathlib import Path from PIL import Image +from solstone.observe.detect import detect_objects, detections_block from solstone.observe.utils import get_segment_key, resize_for_vlm from solstone.think.journal_io import write_jsonl from solstone.think.models import generate @@ -51,6 +53,9 @@ def run(image_path: Path, *, redo: bool = False) -> Path | None: with Image.open(image_path) as img: img.load() + buf = io.BytesIO() + img.save(buf, format="PNG") + detect_png = buf.getvalue() prepared = resize_for_vlm(img) description = generate( contents=[_DESCRIBE_PROMPT, prepared], context="observe.depict" @@ -58,6 +63,9 @@ def run(image_path: Path, *, redo: bool = False) -> Path | None: header = _build_header(image_path.name, "image") entry = {"start": "00:00:00", "text": description} + cli = detect_objects(detect_png) + if cli is not None: + entry["detections"] = detections_block(cli, source="still", gate="still") write_jsonl(output_path, [header, entry]) return output_path diff --git a/solstone/observe/describe.py b/solstone/observe/describe.py index 9b818ceff..7238b12af 100644 --- a/solstone/observe/describe.py +++ b/solstone/observe/describe.py @@ -32,6 +32,7 @@ from typing import List, Optional from PIL import Image +from solstone.observe.detect import detect_objects, detections_block, screen_gate from solstone.observe.exit_codes import EXIT_PROVIDER_BLOCKED from solstone.observe.extract import ( DEFAULT_MAX_EXTRACTIONS, @@ -968,6 +969,13 @@ class VideoProcessor: if req.json_analysis: result["analysis"] = req.json_analysis + gate = screen_gate(req.json_analysis) + if gate is not None: + cli = await asyncio.to_thread(detect_objects, req.frame_bytes) + if cli is not None: + result["detections"] = detections_block( + cli, source="screen", gate=gate + ) # Check if this frame is selected for extraction if frame_id not in selected_ids or req.json_analysis is None: diff --git a/solstone/observe/detect.py b/solstone/observe/detect.py new file mode 100644 index 000000000..ef7374ce7 --- /dev/null +++ b/solstone/observe/detect.py @@ -0,0 +1,122 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""local rf-detr.cpp object detection + read-side qualification policy.""" + +import json +import logging +import subprocess +import tempfile +from pathlib import Path + +from solstone.think.providers.rfdetr_install import ( + ENGINE_REF, + MODEL_NAME, + rfdetr_paths, +) + +LOG = logging.getLogger(__name__) + +ENGINE_NAME = "rf-detr.cpp" +THRESHOLD = 0.25 +_TIMEOUT_S = 120.0 +_GATE_CATEGORIES = frozenset({"media", "social"}) + +_disabled = False + + +def _disable(reason: str) -> None: + global _disabled + if _disabled: + return + _disabled = True + LOG.warning("object detection disabled: %s", reason) + + +def detect_objects(image_bytes: bytes, *, threads: int = 4) -> dict | None: + if _disabled: + return None + try: + paths = rfdetr_paths() + if paths.status != "installed": + _disable(f"rf-detr provider {paths.status}") + return None + with tempfile.TemporaryDirectory(prefix="rfdetr_") as td: + input_png = Path(td) / "input.png" + input_png.write_bytes(image_bytes) + output_json = Path(td) / "output.json" + subprocess.run( + [ + str(paths.binary_path), + "detect", + "--model", + str(paths.model_path), + "--input", + str(input_png), + "--output", + str(output_json), + "--threshold", + str(THRESHOLD), + "--threads", + str(threads), + ], + timeout=_TIMEOUT_S, + capture_output=True, + check=True, + ) + parsed = json.loads(output_json.read_text(encoding="utf-8")) + return parsed + except Exception as exc: + _disable(str(exc)) + return None + + +def detections_block(result: dict, *, source: str, gate: str) -> dict: + return { + "engine": ENGINE_NAME, + "engine_ref": ENGINE_REF, + "model": MODEL_NAME, + "threshold": THRESHOLD, + "source": source, + "gate": gate, + "image": result["image"], + "objects": result["detections"], + } + + +def screen_gate(analysis: dict) -> str | None: + primary = analysis.get("primary") + secondary = analysis.get("secondary") + if primary in _GATE_CATEGORIES: + return f"primary:{primary}" + if secondary in _GATE_CATEGORIES: + return f"secondary:{secondary}" + return None + + +# Stored `detections` rows are UNFILTERED — raw CLI output at THRESHOLD. +# These constants filter at READ time only; the write side stores everything. +_DEVICE_CLASSES = frozenset({"laptop", "tv", "cell phone"}) +_PERSON_CLASS = "person" +_PERSON_MIN_SCORE = 0.4 +_ALLOWED_CLASSES = frozenset({"car", "bird", "tie", "bottle", "cup", "truck", "bowl"}) +_ALLOWED_MIN_SCORE = 0.4 + + +def qualified_objects(block: dict) -> list[dict]: + source = block.get("source") + kept: list[dict] = [] + for obj in block.get("objects", []): + name = obj.get("class_name") + score = obj.get("score", 0.0) + if source == "screen" and name in _DEVICE_CLASSES: + continue + if name == _PERSON_CLASS: + if score >= _PERSON_MIN_SCORE: + kept.append(obj) + continue + if name in _ALLOWED_CLASSES or name in _DEVICE_CLASSES: + if score >= _ALLOWED_MIN_SCORE: + kept.append(obj) + continue + return kept diff --git a/solstone/observe/screen.schema.json b/solstone/observe/screen.schema.json index 251266e7f..aa723afce 100644 --- a/solstone/observe/screen.schema.json +++ b/solstone/observe/screen.schema.json @@ -26,6 +26,7 @@ "frame_id": {"type": ["integer", "string"]}, "analysis": {"type": "object", "additionalProperties": true}, "content": {"type": "object", "additionalProperties": true}, + "detections": {"type": "object", "additionalProperties": true}, "aruco": {}, "enhanced": {"type": "boolean"}, "error": {"type": "string"}, diff --git a/solstone/talent/journal/contract/bundle.json b/solstone/talent/journal/contract/bundle.json index 6bc2987e6..9fba30176 100644 --- a/solstone/talent/journal/contract/bundle.json +++ b/solstone/talent/journal/contract/bundle.json @@ -344,6 +344,10 @@ "additionalProperties": true, "type": "object" }, + "detections": { + "additionalProperties": true, + "type": "object" + }, "enhanced": { "type": "boolean" }, diff --git a/tests/test_depict.py b/tests/test_depict.py index f5ca8060f..179dec1fe 100644 --- a/tests/test_depict.py +++ b/tests/test_depict.py @@ -49,6 +49,43 @@ def test_run_writes_image_jsonl_with_header_metadata(tmp_path, monkeypatch): assert entry == {"start": "00:00:00", "text": "A concise image description"} +def test_run_attaches_still_detection_block(tmp_path, monkeypatch): + image_path = _segment_image(tmp_path) + + def fake_generate(*, contents, context): + return "A concise image description" + + canned = { + "image": {"width": 4, "height": 4}, + "detections": [ + { + "class_id": 7, + "class_name": "bottle", + "score": 0.67, + "bbox": [0, 1, 2, 3], + } + ], + } + + monkeypatch.setattr(depict, "generate", fake_generate) + monkeypatch.setattr(depict, "detect_objects", lambda _image_bytes: canned) + + output_path = depict.run(image_path) + + lines = output_path.read_text(encoding="utf-8").splitlines() + entry = json.loads(lines[1]) + assert entry["detections"] == { + "engine": "rf-detr.cpp", + "engine_ref": "65c0ffcc", + "model": "rfdetr-nano-f16", + "threshold": 0.25, + "source": "still", + "gate": "still", + "image": canned["image"], + "objects": canned["detections"], + } + + def test_run_writes_image_jsonl_bytes_unchanged(tmp_path, monkeypatch): image_path = _segment_image(tmp_path) monkeypatch.setenv("OBSERVER_NAME", "camera") diff --git a/tests/test_describe_promote.py b/tests/test_describe_promote.py index d7678d3c0..37590cfed 100644 --- a/tests/test_describe_promote.py +++ b/tests/test_describe_promote.py @@ -3,6 +3,7 @@ import io import json +import logging from pathlib import Path from types import SimpleNamespace @@ -10,7 +11,9 @@ import pytest from PIL import Image from solstone.observe import describe as describe_module +from solstone.observe import detect as detect_module from solstone.observe import processing_record as processing_record_module +from solstone.think.providers.rfdetr_install import RfdetrPaths def _video_path(tmp_path: Path) -> Path: @@ -21,9 +24,9 @@ def _video_path(tmp_path: Path) -> Path: return video_path -def _png_bytes() -> bytes: +def _png_bytes(size: tuple[int, int] = (8, 8)) -> bytes: image_bytes = io.BytesIO() - Image.new("RGB", (8, 8), "white").save(image_bytes, format="PNG") + Image.new("RGB", size, "white").save(image_bytes, format="PNG") return image_bytes.getvalue() @@ -54,6 +57,28 @@ def _assert_no_describe_temp(directory: Path) -> None: ) +def _jsonl_rows(path: Path) -> list[dict]: + return [ + json.loads(line) + for line in path.read_text(encoding="utf-8").splitlines() + if line + ] + + +def _canned_detection() -> dict: + return { + "image": {"width": 8, "height": 8}, + "detections": [ + { + "class_id": 42, + "class_name": "cup", + "score": 0.7, + "bbox": [1, 2, 3, 4], + } + ], + } + + def _install_fakes(monkeypatch, outcomes: dict[int, dict]) -> list[tuple]: from solstone.think import batch as batch_module from solstone.think import models @@ -257,6 +282,347 @@ async def test_success_with_mixed_results_promotes_byte_identical_jsonl( _assert_no_describe_temp(output_path.parent) +@pytest.mark.asyncio +async def test_detection_blocks_attach_to_media_and_social_frames( + tmp_path, monkeypatch +): + video_path = _video_path(tmp_path) + output_path = video_path.with_suffix(".jsonl") + frame_bytes = _png_bytes() + processor = _processor( + video_path, + [ + _frame(1, 0.0, frame_bytes), + _frame(2, 1.25, frame_bytes), + ], + monkeypatch, + ) + canned = _canned_detection() + calls = [] + _install_fakes( + monkeypatch, + { + 1: { + "response": json.dumps( + {"primary": "media", "secondary": "none", "overlap": True} + ) + }, + 2: { + "response": json.dumps( + {"primary": "code", "secondary": "social", "overlap": True} + ) + }, + }, + ) + monkeypatch.setattr( + describe_module, + "select_frames_for_extraction", + lambda *_args, **_kwargs: [], + ) + + def fake_detect(image_bytes): + calls.append(image_bytes) + return canned + + monkeypatch.setattr(describe_module, "detect_objects", fake_detect) + + await processor.process_with_vision( + max_concurrent=1, + output_path=output_path, + work_key="20250101/143022_300/screen", + ) + + frame1, frame2 = _jsonl_rows(output_path)[1:] + assert frame1["detections"] == { + "engine": "rf-detr.cpp", + "engine_ref": "65c0ffcc", + "model": "rfdetr-nano-f16", + "threshold": 0.25, + "source": "screen", + "gate": "primary:media", + "image": canned["image"], + "objects": canned["detections"], + } + assert frame2["detections"] == { + "engine": "rf-detr.cpp", + "engine_ref": "65c0ffcc", + "model": "rfdetr-nano-f16", + "threshold": 0.25, + "source": "screen", + "gate": "secondary:social", + "image": canned["image"], + "objects": canned["detections"], + } + assert calls == [frame_bytes, frame_bytes] + + +@pytest.mark.asyncio +async def test_detection_gate_off_never_invokes_detector(tmp_path, monkeypatch): + video_path = _video_path(tmp_path) + output_path = video_path.with_suffix(".jsonl") + frame_bytes = _png_bytes() + processor = _processor( + video_path, + [ + _frame(1, 0.0, frame_bytes), + _frame(2, 1.25, frame_bytes), + _frame(3, 2.5, frame_bytes), + ], + monkeypatch, + ) + calls = 0 + _install_fakes( + monkeypatch, + { + 1: { + "response": json.dumps( + {"primary": "code", "secondary": "none", "overlap": True} + ) + }, + 2: { + "response": json.dumps( + {"primary": "terminal", "secondary": "none", "overlap": True} + ) + }, + 3: { + "response": json.dumps( + {"primary": "browsing", "secondary": "none", "overlap": True} + ) + }, + }, + ) + monkeypatch.setattr( + describe_module, + "select_frames_for_extraction", + lambda *_args, **_kwargs: [], + ) + + def fake_detect(_image_bytes): + nonlocal calls + calls += 1 + return _canned_detection() + + monkeypatch.setattr(describe_module, "detect_objects", fake_detect) + + await processor.process_with_vision( + max_concurrent=1, + output_path=output_path, + work_key="20250101/143022_300/screen", + ) + + rows = _jsonl_rows(output_path)[1:] + assert all("detections" not in row for row in rows) + assert calls == 0 + + +@pytest.mark.asyncio +async def test_detection_skips_categorization_failed_frame(tmp_path, monkeypatch): + video_path = _video_path(tmp_path) + output_path = video_path.with_suffix(".jsonl") + gated_frame_bytes = _png_bytes((8, 8)) + failed_frame_bytes = _png_bytes((10, 10)) + processor = _processor( + video_path, + [ + _frame(1, 0.0, gated_frame_bytes), + _frame(2, 1.25, failed_frame_bytes), + ], + monkeypatch, + ) + calls = [] + _install_fakes( + monkeypatch, + { + 1: { + "response": json.dumps( + {"primary": "media", "secondary": "none", "overlap": True} + ) + }, + 2: {"fail": True, "error": "boom"}, + }, + ) + monkeypatch.setattr( + describe_module, + "select_frames_for_extraction", + lambda *_args, **_kwargs: [], + ) + + def fake_detect(image_bytes): + calls.append(image_bytes) + return _canned_detection() + + monkeypatch.setattr(describe_module, "detect_objects", fake_detect) + + await processor.process_with_vision( + max_concurrent=1, + output_path=output_path, + work_key="20250101/143022_300/screen", + ) + + frame1, frame2 = _jsonl_rows(output_path)[1:] + assert "detections" in frame1 + assert "detections" not in frame2 + assert calls == [gated_frame_bytes] + + +@pytest.mark.asyncio +async def test_detection_provider_absence_latches_across_gated_frames( + tmp_path, monkeypatch, caplog +): + video_path = _video_path(tmp_path) + output_path = video_path.with_suffix(".jsonl") + frame_bytes = _png_bytes() + processor = _processor( + video_path, + [ + _frame(1, 0.0, frame_bytes), + _frame(2, 1.25, frame_bytes), + _frame(3, 2.5, frame_bytes), + ], + monkeypatch, + ) + calls = 0 + monkeypatch.setattr(detect_module, "_disabled", False) + monkeypatch.setattr(describe_module, "detect_objects", detect_module.detect_objects) + + def fake_paths(): + nonlocal calls + calls += 1 + return RfdetrPaths(status="not_installed") + + monkeypatch.setattr(detect_module, "rfdetr_paths", fake_paths) + caplog.set_level(logging.WARNING, logger=detect_module.LOG.name) + _install_fakes( + monkeypatch, + { + 1: { + "response": json.dumps( + {"primary": "media", "secondary": "none", "overlap": True} + ) + }, + 2: { + "response": json.dumps( + {"primary": "social", "secondary": "none", "overlap": True} + ) + }, + 3: { + "response": json.dumps( + {"primary": "code", "secondary": "social", "overlap": True} + ) + }, + }, + ) + monkeypatch.setattr( + describe_module, + "select_frames_for_extraction", + lambda *_args, **_kwargs: [], + ) + + await processor.process_with_vision( + max_concurrent=1, + output_path=output_path, + work_key="20250101/143022_300/screen", + ) + + rows = _jsonl_rows(output_path)[1:] + warnings = [ + record + for record in caplog.records + if record.name == detect_module.LOG.name and record.levelno == logging.WARNING + ] + assert all("detections" not in row for row in rows) + assert len(warnings) == 1 + assert warnings[0].getMessage() == ( + "object detection disabled: rf-detr provider not_installed" + ) + assert calls == 1 + + +@pytest.mark.asyncio +async def test_detection_empty_result_stores_empty_objects(tmp_path, monkeypatch): + video_path = _video_path(tmp_path) + output_path = video_path.with_suffix(".jsonl") + frame_bytes = _png_bytes() + processor = _processor( + video_path, + [_frame(1, 0.0, frame_bytes)], + monkeypatch, + ) + _install_fakes( + monkeypatch, + { + 1: { + "response": json.dumps( + {"primary": "media", "secondary": "none", "overlap": True} + ) + }, + }, + ) + monkeypatch.setattr( + describe_module, + "select_frames_for_extraction", + lambda *_args, **_kwargs: [], + ) + monkeypatch.setattr( + describe_module, + "detect_objects", + lambda _image_bytes: {"image": {"width": 8, "height": 8}, "detections": []}, + ) + + await processor.process_with_vision( + max_concurrent=1, + output_path=output_path, + work_key="20250101/143022_300/screen", + ) + + frame = _jsonl_rows(output_path)[1] + assert frame["detections"]["objects"] == [] + + +@pytest.mark.asyncio +async def test_detection_uses_full_resolution_frame_bytes(tmp_path, monkeypatch): + video_path = _video_path(tmp_path) + output_path = video_path.with_suffix(".jsonl") + frame_bytes = _png_bytes((2100, 20)) + processor = _processor( + video_path, + [_frame(1, 0.0, frame_bytes)], + monkeypatch, + ) + observed_sizes = [] + _install_fakes( + monkeypatch, + { + 1: { + "response": json.dumps( + {"primary": "media", "secondary": "none", "overlap": True} + ) + }, + }, + ) + monkeypatch.setattr( + describe_module, + "select_frames_for_extraction", + lambda *_args, **_kwargs: [], + ) + + def fake_detect(image_bytes): + with Image.open(io.BytesIO(image_bytes)) as img: + observed_sizes.append(img.size) + return _canned_detection() + + monkeypatch.setattr(describe_module, "detect_objects", fake_detect) + + await processor.process_with_vision( + max_concurrent=1, + output_path=output_path, + work_key="20250101/143022_300/screen", + ) + + assert observed_sizes == [(2100, 20)] + assert observed_sizes[0][0] > 1920 + + @pytest.mark.asyncio async def test_empty_run_promotes_header_only_file_for_event_precondition( tmp_path, monkeypatch diff --git a/tests/test_detect_objects.py b/tests/test_detect_objects.py new file mode 100644 index 000000000..8dfd5943f --- /dev/null +++ b/tests/test_detect_objects.py @@ -0,0 +1,158 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +import json +import logging +import subprocess +from pathlib import Path + +import pytest + +from solstone.observe import detect +from solstone.think.providers.rfdetr_install import ( + ENGINE_REF, + MODEL_NAME, + RfdetrPaths, +) + + +@pytest.fixture(autouse=True) +def _reset_detect_state(): + detect._disabled = False + yield + detect._disabled = False + + +def _canned_cli_json() -> dict: + return { + "image": {"width": 100, "height": 50}, + "detections": [ + { + "class_id": 1, + "class_name": "cup", + "score": 0.72, + "bbox": [1, 2, 3, 4], + } + ], + } + + +def test_detect_objects_missing_provider_latches_without_reattempt(monkeypatch, caplog): + calls = 0 + + def fake_paths(): + nonlocal calls + calls += 1 + return RfdetrPaths(status="not_installed") + + monkeypatch.setattr(detect, "rfdetr_paths", fake_paths) + caplog.set_level(logging.WARNING, logger=detect.LOG.name) + + assert detect.detect_objects(b"x") is None + assert detect.detect_objects(b"x") is None + + warnings = [ + record + for record in caplog.records + if record.name == detect.LOG.name and record.levelno == logging.WARNING + ] + assert detect._disabled is True + assert len(warnings) == 1 + assert warnings[0].getMessage() == ( + "object detection disabled: rf-detr provider not_installed" + ) + assert calls == 1 + + +def test_detect_objects_returns_parsed_cli_json(monkeypatch, tmp_path): + canned = _canned_cli_json() + argv_seen = [] + + monkeypatch.setattr( + detect, + "rfdetr_paths", + lambda: RfdetrPaths( + status="installed", + binary_path=tmp_path / "rfdetr-cli", + model_path=tmp_path / "model.gguf", + ), + ) + + def fake_run(argv, **kwargs): + argv_seen.append((argv, kwargs)) + output_path = Path(argv[argv.index("--output") + 1]) + output_path.write_text(json.dumps(canned), encoding="utf-8") + return subprocess.CompletedProcess(argv, 0) + + monkeypatch.setattr(detect.subprocess, "run", fake_run) + + assert detect.detect_objects(b"png-bytes") == canned + assert len(argv_seen) == 1 + argv, kwargs = argv_seen[0] + assert argv[1] == "detect" + assert argv[argv.index("--threshold") + 1] == str(detect.THRESHOLD) + assert kwargs["timeout"] == detect._TIMEOUT_S + assert kwargs["capture_output"] is True + assert kwargs["check"] is True + + +def test_detect_objects_subprocess_failure_latches(monkeypatch, tmp_path): + monkeypatch.setattr( + detect, + "rfdetr_paths", + lambda: RfdetrPaths( + status="installed", + binary_path=tmp_path / "rfdetr-cli", + model_path=tmp_path / "model.gguf", + ), + ) + + def fake_run(argv, **_kwargs): + raise subprocess.CalledProcessError(1, argv, stderr=b"bad") + + monkeypatch.setattr(detect.subprocess, "run", fake_run) + + assert detect.detect_objects(b"png-bytes") is None + assert detect._disabled is True + + +def test_detections_block_renames_and_preserves_provenance(): + canned = _canned_cli_json() + + assert detect.detections_block(canned, source="screen", gate="primary:media") == { + "engine": detect.ENGINE_NAME, + "engine_ref": ENGINE_REF, + "model": MODEL_NAME, + "threshold": detect.THRESHOLD, + "source": "screen", + "gate": "primary:media", + "image": canned["image"], + "objects": canned["detections"], + } + + +def test_screen_gate_prefers_primary_gate(): + assert detect.screen_gate({"primary": "media", "secondary": "none"}) == ( + "primary:media" + ) + assert detect.screen_gate({"primary": "code", "secondary": "social"}) == ( + "secondary:social" + ) + assert detect.screen_gate({"primary": "media", "secondary": "social"}) == ( + "primary:media" + ) + assert detect.screen_gate({"primary": "code", "secondary": "terminal"}) is None + + +def test_qualified_objects_filters_at_read_time_only(): + tv = {"class_name": "tv", "score": 0.7} + weak_person = {"class_name": "person", "score": 0.39} + person = {"class_name": "person", "score": 0.41} + cup = {"class_name": "cup", "score": 0.41} + sandwich = {"class_name": "sandwich", "score": 0.9} + + assert detect.qualified_objects({"source": "screen", "objects": [tv]}) == [] + assert detect.qualified_objects({"source": "still", "objects": [tv]}) == [tv] + assert detect.qualified_objects( + {"source": "still", "objects": [weak_person, person, cup, sandwich]} + ) == [person, cup]