diff --git a/AGENTS.md b/AGENTS.md index a2956c186..24e2fba53 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -187,7 +187,7 @@ Each domain has exactly **one** write-owning module (or one tightly-scoped famil | Observations (`observations.jsonl`) | `solstone/think/entities/observations.py` | | Activities (`facets/*/activities/*.jsonl`) | `solstone/think/activities.py` | | Timeline (`chronicle//timeline.json`, `chronicle/**//timeline.json`, root `timeline.json`) | `solstone/apps/timeline/maintenance.py` + `solstone/apps/timeline/talent/segment_summary.py` | -| Per-segment sense outputs (`chronicle/**//talents/{sense.json,facets.json,speakers.json,density.json,activity.md,sense.md}`) | `solstone/think/sense_splitter.py` | +| Per-segment sense outputs (`chronicle/**//talents/{sense.json,facets.json,speakers.json,density.json,change.json,activity.md,sense.md}`) | `solstone/think/sense_splitter.py` | | Awareness (`awareness/current.json`, `awareness/YYYYMMDD.jsonl`) | `solstone/think/awareness.py` | | Awareness activity state (`awareness/activity_state.json`) | `solstone/think/thinking.py` | | Identity (`identity/*.md`, `identity/history.jsonl` audit log) | `solstone/think/identity.py` | diff --git a/solstone/observe/describe.py b/solstone/observe/describe.py index 76ffd1e4e..54582d8b2 100644 --- a/solstone/observe/describe.py +++ b/solstone/observe/describe.py @@ -323,6 +323,9 @@ class VideoProcessor: DHASH_THRESHOLD = 8 # Skip frame if Convey UI covers more than this fraction of the frame MASK_SKIP_THRESHOLD = 0.8 + first_hash: Optional[int] = None + last_hash: Optional[int] = None + qualified_count: int = 0 def __init__(self, video_path: Path): self.video_path = video_path @@ -343,6 +346,9 @@ class VideoProcessor: """ # Cache for the last qualified frame hash last_hash: Optional[int] = None + self.first_hash = None + self.last_hash = None + self.qualified_count = 0 # Imports deferred: av (PyAV) and cv2 (via observe.aruco) bundle # mismatched libavdevice majors. Keeping them out of module scope @@ -422,6 +428,8 @@ class VideoProcessor: if last_hash is None: frame_data["frame_bytes"] = self._frame_to_bytes(pil_img) last_hash = self._dhash(pil_img) + self.first_hash = last_hash + self.last_hash = last_hash pil_img.close() self.qualified_frames.append(frame_data) @@ -446,11 +454,13 @@ class VideoProcessor: # Update cached frame hash last_hash = current_hash + self.last_hash = current_hash logger.debug( f"Qualified frame at {timestamp:.2f}s (hamming: {distance})" ) + self.qualified_count = len(self.qualified_frames) logger.info( f"Processed {frame_count} frames from {self.video_path.name}, " f"{len(self.qualified_frames)} qualified" @@ -470,6 +480,36 @@ class VideoProcessor: raise return self.qualified_frames + @staticmethod + def _format_dhash(hash_value: int | None) -> str | None: + if hash_value is None: + return None + return f"{hash_value:016x}" + + def _build_metadata_header(self) -> dict: + # Files are in segment directories, filename is simple (e.g., center_DP-3_screen.webm) + metadata = {"raw": self.video_path.name} + + # Add observer origin if set (from sense.py for observer uploads) + observer = os.getenv("OBSERVER_NAME") + if observer: + metadata["observer"] = observer + + # Add segment metadata (from sense.py via SEGMENT_META env var) + segment_meta_str = os.getenv("SEGMENT_META") + if segment_meta_str: + try: + segment_meta = json.loads(segment_meta_str) + for key, value in segment_meta.items(): + metadata[key] = value + except json.JSONDecodeError: + logger.warning(f"Invalid SEGMENT_META JSON: {segment_meta_str[:100]}") + + metadata["first_hash"] = self._format_dhash(self.first_hash) + metadata["last_hash"] = self._format_dhash(self.last_hash) + metadata["qualified_count"] = self.qualified_count + return metadata + def _dhash(self, img: Image.Image) -> int: """Compute 64-bit dHash (difference hash) for perceptual comparison.""" small = img.resize(self.DHASH_SIZE, Image.BILINEAR).convert("L") @@ -599,27 +639,7 @@ class VideoProcessor: try: # Write metadata header to JSONL file with actual video filename if output_file: - # Files are in segment directories, filename is simple (e.g., center_DP-3_screen.webm) - metadata = {"raw": self.video_path.name} - - # Add observer origin if set (from sense.py for observer uploads) - observer = os.getenv("OBSERVER_NAME") - if observer: - metadata["observer"] = observer - - # Add segment metadata (from sense.py via SEGMENT_META env var) - segment_meta_str = os.getenv("SEGMENT_META") - if segment_meta_str: - try: - segment_meta = json.loads(segment_meta_str) - for key, value in segment_meta.items(): - metadata[key] = value - except json.JSONDecodeError: - logger.warning( - f"Invalid SEGMENT_META JSON: {segment_meta_str[:100]}" - ) - - output_file.write(json.dumps(metadata) + "\n") + output_file.write(json.dumps(self._build_metadata_header()) + "\n") # Resolve model for frame description (tier from describe.md frontmatter) _, frame_model = resolve_provider(FRAME_CONTEXT, "generate") diff --git a/solstone/think/change_detection.py b/solstone/think/change_detection.py new file mode 100644 index 000000000..80a8fa42b --- /dev/null +++ b/solstone/think/change_detection.py @@ -0,0 +1,338 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Read-only segment change detection for Sense shadow mode.""" + +from __future__ import annotations + +import hashlib +import json +import logging +from datetime import datetime +from pathlib import Path +from typing import Any + +from solstone.observe.hear import load_transcript +from solstone.observe.utils import parse_screen_filename +from solstone.think.utils import DEFAULT_STREAM, iter_segments, segment_parse + +# Mirrors VideoProcessor.DHASH_THRESHOLD in solstone.observe.describe. +SCREEN_DHASH_THRESHOLD = 8 +# Bias to active: false redundant silently drops signal; false active only costs work. +TRANSCRIPT_WORD_DELTA_FLOOR = 5 +# Mirrors activity_state_machine.py's current gap threshold, but remains independent. +GAP_THRESHOLD_SECONDS = 600 + + +def assemble_sensor_state(seg_dir: Path) -> dict: + """Assemble current per-sensor state for a segment without writing anything.""" + screen_monitors: dict[str, dict[str, str | int | None]] = {} + for screen_path in sorted(seg_dir.glob("*screen.jsonl")): + if not screen_path.is_file(): + continue + position, connector = parse_screen_filename(screen_path.stem) + header = _read_header(screen_path) + screen_monitors[f"{position}:{connector}"] = { + "first_hash": _normalize_hash(header.get("first_hash") if header else None), + "last_hash": _normalize_hash(header.get("last_hash") if header else None), + "qualified_count": _normalize_count( + header.get("qualified_count") if header else None + ), + } + + transcript_text = _load_transcript_text(seg_dir) + normalized_text = _normalize_transcript(transcript_text) + transcript_state = { + "present": bool(normalized_text), + "word_count": len(normalized_text.split()) if normalized_text else 0, + "content_hash": ( + "sha256:" + hashlib.sha256(normalized_text.encode()).hexdigest() + if normalized_text + else None + ), + } + + return { + "screen": {"monitors": screen_monitors}, + "transcript": transcript_state, + } + + +def resolve_predecessor( + day: str, stream: str | None, segment: str +) -> dict[str, str] | None: + """Return the chronological prior same-stream segment ref, if comparable. + + This must derive from ``iter_segments(day)`` only. Do not read + ``last_segment_key`` or ``awareness/activity_state.json`` here: those track + last-processed state, not the chronological predecessor, and are wrong under + backfill or reprocess. + """ + stream_name = stream if stream is not None else DEFAULT_STREAM + same_stream = [ + (seg_key, seg_path) + for seg_stream, seg_key, seg_path in iter_segments(day) + if seg_stream == stream_name + ] + segment_index = next( + ( + index + for index, (seg_key, _seg_path) in enumerate(same_stream) + if seg_key == segment + ), + None, + ) + if segment_index is None or segment_index == 0: + return None + + predecessor_segment, _predecessor_path = same_stream[segment_index - 1] + predecessor_start, predecessor_end = segment_parse(predecessor_segment) + current_start, _current_end = segment_parse(segment) + if predecessor_start is None or predecessor_end is None or current_start is None: + return None + + day_date = datetime.strptime(day, "%Y%m%d").date() + predecessor_end_dt = datetime.combine(day_date, predecessor_end) + current_start_dt = datetime.combine(day_date, current_start) + gap_seconds = (current_start_dt - predecessor_end_dt).total_seconds() + if gap_seconds > GAP_THRESHOLD_SECONDS: + return None + + return {"day": day, "stream": stream_name, "segment": predecessor_segment} + + +def read_predecessor_state( + day: str, predecessor_ref: dict[str, str] | None +) -> dict | None: + """Read predecessor sensors from ``talents/change.json`` if available.""" + if predecessor_ref is None: + return None + + target_stream = predecessor_ref["stream"] + target_segment = predecessor_ref["segment"] + predecessor_dir = next( + ( + seg_path + for seg_stream, seg_key, seg_path in iter_segments(day) + if seg_stream == target_stream and seg_key == target_segment + ), + None, + ) + if predecessor_dir is None: + return None + + change_path = predecessor_dir / "talents" / "change.json" + try: + data = json.loads(change_path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError) as exc: + logging.debug( + "Failed to read predecessor change state %s: %s", change_path, exc + ) + return None + + sensors = data.get("sensors") if isinstance(data, dict) else None + if not isinstance(sensors, dict): + return None + if not isinstance(sensors.get("screen"), dict): + return None + if not isinstance(sensors.get("transcript"), dict): + return None + return sensors + + +def compare_screen(prev: dict, curr: dict) -> dict[str, bool]: + """Compare previous screen last-hash boundaries to current first hashes.""" + prev_monitors = _screen_monitors(prev) + curr_monitors = _screen_monitors(curr) + present = bool(prev_monitors or curr_monitors) + if not present: + return {"present": False, "changed": False} + + if set(prev_monitors) != set(curr_monitors): + return {"present": True, "changed": True} + + for monitor_key in sorted(curr_monitors): + prev_last = _hash_to_int(prev_monitors[monitor_key].get("last_hash")) + curr_first = _hash_to_int(curr_monitors[monitor_key].get("first_hash")) + if prev_last is None or curr_first is None: + return {"present": True, "changed": True} + if (prev_last ^ curr_first).bit_count() >= SCREEN_DHASH_THRESHOLD: + return {"present": True, "changed": True} + + return {"present": True, "changed": False} + + +def compare_transcript(prev: dict, curr: dict) -> dict[str, bool]: + """Compare transcript hashes with a tight word-delta noise gate.""" + if not prev.get("present") or not curr.get("present"): + return {"present": False, "changed": False} + + prev_hash = prev.get("content_hash") + curr_hash = curr.get("content_hash") + if not isinstance(prev_hash, str) or not isinstance(curr_hash, str): + return {"present": True, "changed": True} + if prev_hash == curr_hash: + return {"present": True, "changed": False} + + word_delta = abs(_word_count(curr) - _word_count(prev)) + return { + "present": True, + "changed": word_delta > TRANSCRIPT_WORD_DELTA_FLOOR, + } + + +def classify(vectors: dict) -> tuple[str, list[str]]: + """Classify segment change state from per-sensor vectors.""" + present = any(bool(vector.get("present")) for vector in vectors.values()) + if not present: + return "idle", [] + + changed_sensors = sorted( + name + for name, vector in vectors.items() + if vector.get("present") and vector.get("changed") + ) + if changed_sensors: + return "active", changed_sensors + return "redundant", [] + + +def detect_segment_change( + day: str, + stream: str | None, + segment: str, + seg_dir: Path, + *, + predecessor: dict | None, + timestamp: str, +) -> dict: + """Return the full change-detection result for persistence.""" + current_state = assemble_sensor_state(seg_dir) + predecessor_state = read_predecessor_state(day, predecessor) + + if predecessor is None or predecessor_state is None: + vectors = _missing_predecessor_vectors(current_state) + else: + vectors = { + "screen": compare_screen( + predecessor_state["screen"], current_state["screen"] + ), + "transcript": compare_transcript( + predecessor_state["transcript"], current_state["transcript"] + ), + } + + change_class, changed_sensors = classify(vectors) + return { + "timestamp": timestamp, + "predecessor": predecessor, + "change_class": change_class, + "changed_sensors": changed_sensors, + "sensors": current_state, + } + + +def _read_header(path: Path) -> dict[str, Any] | None: + try: + with path.open(encoding="utf-8") as handle: + first_line = handle.readline() + if not first_line.strip(): + return None + header = json.loads(first_line) + except (OSError, json.JSONDecodeError) as exc: + logging.debug("Failed to read screen header %s: %s", path, exc) + return None + if not isinstance(header, dict): + return None + return header + + +def _normalize_hash(value: object) -> str | None: + parsed = _hash_to_int(value) + if parsed is None: + return None + return f"{parsed:016x}" + + +def _hash_to_int(value: object) -> int | None: + if isinstance(value, bool): + return None + if isinstance(value, int): + parsed = value + elif isinstance(value, str): + try: + parsed = int(value, 16) + except ValueError: + return None + else: + return None + if not 0 <= parsed < 2**64: + return None + return parsed + + +def _normalize_count(value: object) -> int: + if isinstance(value, bool) or not isinstance(value, int): + return 0 + return max(value, 0) + + +def _load_transcript_text(seg_dir: Path) -> str: + parts: list[str] = [] + + jsonl_files = set() + for pattern in ("*audio.jsonl", "*_transcript.jsonl"): + jsonl_files.update(path for path in seg_dir.glob(pattern) if path.is_file()) + for jsonl_path in sorted(jsonl_files): + metadata, transcript_entries, formatted_text = load_transcript(jsonl_path) + if transcript_entries is None: + logging.debug( + "Skipping unreadable transcript %s: %s", + jsonl_path, + metadata.get("error"), + ) + continue + if formatted_text.strip(): + parts.append(formatted_text) + + md_files = set() + for pattern in ("*_transcript.md", "imported.md"): + md_files.update(path for path in seg_dir.glob(pattern) if path.is_file()) + for md_path in sorted(md_files): + try: + content = md_path.read_text(encoding="utf-8") + except OSError as exc: + logging.debug("Skipping unreadable transcript %s: %s", md_path, exc) + continue + if content.strip(): + parts.append(content) + + return "\n".join(parts) + + +def _normalize_transcript(text: str) -> str: + return " ".join(text.lower().split()).strip() + + +def _screen_monitors(state: dict) -> dict: + monitors = state.get("monitors") if isinstance(state, dict) else None + return monitors if isinstance(monitors, dict) else {} + + +def _word_count(state: dict) -> int: + value = state.get("word_count") + if isinstance(value, bool) or not isinstance(value, int): + return 0 + return max(value, 0) + + +def _missing_predecessor_vectors(current_state: dict) -> dict[str, dict[str, bool]]: + screen_present = bool(_screen_monitors(current_state["screen"])) + transcript_present = bool(current_state["transcript"].get("present")) + return { + "screen": {"present": screen_present, "changed": screen_present}, + "transcript": { + "present": transcript_present, + "changed": transcript_present, + }, + } diff --git a/solstone/think/sense_splitter.py b/solstone/think/sense_splitter.py index 66107b329..c5290a021 100644 --- a/solstone/think/sense_splitter.py +++ b/solstone/think/sense_splitter.py @@ -30,8 +30,6 @@ def write_sense_outputs( json.dumps( { "classification": density, - "transcript_lines": 0, - "screen_frames": 0, "timestamp": datetime.now(tz=timezone.utc).isoformat(), } ), @@ -68,9 +66,12 @@ def write_idle_stubs(seg_dir: Path) -> None: json.dumps( { "classification": "idle", - "transcript_lines": 0, - "screen_frames": 0, "timestamp": datetime.now(tz=timezone.utc).isoformat(), } ), ) + + +def write_change_detection(seg_dir: Path, result: dict) -> None: + """Write segment change-detection state.""" + atomic_replace(seg_dir / "talents" / "change.json", json.dumps(result)) diff --git a/solstone/think/thinking.py b/solstone/think/thinking.py index fa7517d32..26354f9d2 100644 --- a/solstone/think/thinking.py +++ b/solstone/think/thinking.py @@ -17,7 +17,7 @@ import logging import sys import threading import time -from datetime import date, datetime, timedelta +from datetime import date, datetime, timedelta, timezone from pathlib import Path from solstone.think.activities import ( @@ -28,6 +28,7 @@ from solstone.think.activities import ( ) from solstone.think.activity_state_machine import ActivityStateMachine from solstone.think.callosum import CallosumConnection +from solstone.think.change_detection import detect_segment_change, resolve_predecessor from solstone.think.cluster import cluster_segments from solstone.think.cogitate_policy import DETERMINISTIC_FAILURE_THRESHOLD from solstone.think.cortex_client import ( @@ -53,7 +54,11 @@ from solstone.think.pipeline_health import ( read_segment_progress, ) from solstone.think.runner import run_task -from solstone.think.sense_splitter import write_idle_stubs, write_sense_outputs +from solstone.think.sense_splitter import ( + write_change_detection, + write_idle_stubs, + write_sense_outputs, +) from solstone.think.talent import get_output_path, get_talent_configs from solstone.think.talent_provenance import ( compute_activity_input_hash, @@ -691,6 +696,7 @@ def run_segment_sense( skip_activity_prompts: bool = False, skip_talents: frozenset[str] = frozenset(), live: bool = False, + predecessor: dict | None = None, ) -> tuple[int, int, list[str]]: """Run Sense-first linear orchestrator for a single segment. @@ -910,6 +916,25 @@ def run_segment_sense( recommend=sense_json.get("recommend") or {}, **({"stream": stream} if stream else {}), ) + change_result = detect_segment_change( + day, + stream, + segment, + seg_dir, + predecessor=predecessor, + timestamp=datetime.now(tz=timezone.utc).isoformat(), + ) + write_change_detection(seg_dir, change_result) + _jsonl_log( + "sense.change_detect", + mode=target_schedule, + day=day, + segment=segment, + change_class=change_result["change_class"], + changed_sensors=change_result["changed_sensors"], + predecessor=change_result["predecessor"], + **({"stream": stream} if stream else {}), + ) if density == "idle" and not refresh: write_idle_stubs(seg_dir) @@ -3436,6 +3461,7 @@ def main() -> None: skip_activity_prompts=args.no_activity_prompts, skip_talents=skip_talents, live=False, + predecessor=resolve_predecessor(day, seg_stream, seg_key), ) batch_success += success batch_failed += failed @@ -3593,6 +3619,7 @@ def main() -> None: skip_activity_prompts=args.no_activity_prompts, skip_talents=skip_talents, live=args.live, + predecessor=resolve_predecessor(day, resolved_stream, args.segment), ) else: success_count, fail_count, failed_names, applicable_units = ( diff --git a/tests/test_change_detection.py b/tests/test_change_detection.py new file mode 100644 index 000000000..094d30191 --- /dev/null +++ b/tests/test_change_detection.py @@ -0,0 +1,309 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Tests for segment change detection.""" + +import json +from pathlib import Path + +from solstone.think.change_detection import ( + assemble_sensor_state, + classify, + compare_screen, + compare_transcript, + detect_segment_change, + resolve_predecessor, +) + + +def _screen_state( + *, + key: str = "center:DP-3", + first_hash: str | None = "0000000000000000", + last_hash: str | None = "0000000000000000", + qualified_count: int = 1, +) -> dict: + return { + "monitors": { + key: { + "first_hash": first_hash, + "last_hash": last_hash, + "qualified_count": qualified_count, + } + } + } + + +def _transcript_state( + *, + present: bool = True, + word_count: int = 10, + content_hash: str | None = "sha256:aaa", +) -> dict: + return { + "present": present, + "word_count": word_count, + "content_hash": content_hash, + } + + +def _make_segment(journal: Path, day: str, stream: str, segment: str) -> Path: + seg_dir = journal / "chronicle" / day / stream / segment + seg_dir.mkdir(parents=True) + return seg_dir + + +def _write_screen(seg_dir: Path, *, first: str, last: str, count: int = 1) -> None: + (seg_dir / "center_DP-3_screen.jsonl").write_text( + json.dumps( + { + "raw": "center_DP-3_screen.webm", + "first_hash": first, + "last_hash": last, + "qualified_count": count, + } + ) + + "\n", + encoding="utf-8", + ) + + +class TestCompareScreen: + def test_hamming_below_threshold_is_unchanged(self): + result = compare_screen( + _screen_state(last_hash="0000000000000000"), + _screen_state(first_hash="0000000000000001"), + ) + + assert result == {"present": True, "changed": False} + + def test_hamming_at_threshold_is_changed(self): + result = compare_screen( + _screen_state(last_hash="0000000000000000"), + _screen_state(first_hash="00000000000000ff"), + ) + + assert result == {"present": True, "changed": True} + + def test_monitor_present_on_one_side_is_changed(self): + result = compare_screen({"monitors": {}}, _screen_state()) + + assert result == {"present": True, "changed": True} + + def test_missing_hash_is_changed(self): + result = compare_screen( + _screen_state(last_hash=None), + _screen_state(first_hash="0000000000000000"), + ) + + assert result == {"present": True, "changed": True} + + def test_single_frame_monitor_uses_same_boundary_hash(self): + result = compare_screen( + _screen_state( + first_hash="0000000000000001", + last_hash="0000000000000001", + ), + _screen_state( + first_hash="0000000000000001", + last_hash="0000000000000001", + ), + ) + + assert result == {"present": True, "changed": False} + + +class TestCompareTranscript: + def test_identical_hash_is_unchanged(self): + result = compare_transcript( + _transcript_state(word_count=12, content_hash="sha256:same"), + _transcript_state(word_count=18, content_hash="sha256:same"), + ) + + assert result == {"present": True, "changed": False} + + def test_differing_hash_with_word_delta_above_floor_is_changed(self): + result = compare_transcript( + _transcript_state(word_count=10, content_hash="sha256:prev"), + _transcript_state(word_count=16, content_hash="sha256:curr"), + ) + + assert result == {"present": True, "changed": True} + + def test_differing_hash_with_word_delta_at_floor_is_unchanged(self): + result = compare_transcript( + _transcript_state(word_count=10, content_hash="sha256:prev"), + _transcript_state(word_count=15, content_hash="sha256:curr"), + ) + + assert result == {"present": True, "changed": False} + + def test_absent_on_either_side_is_not_present(self): + result = compare_transcript( + _transcript_state(present=False, word_count=0, content_hash=None), + _transcript_state(word_count=20, content_hash="sha256:curr"), + ) + + assert result == {"present": False, "changed": False} + + +class TestClassify: + def test_all_absent_is_idle(self): + assert classify( + { + "screen": {"present": False, "changed": False}, + "transcript": {"present": False, "changed": False}, + } + ) == ("idle", []) + + def test_present_and_none_changed_is_redundant(self): + assert classify( + { + "screen": {"present": True, "changed": False}, + "transcript": {"present": False, "changed": False}, + } + ) == ("redundant", []) + + def test_any_changed_is_active(self): + assert classify( + { + "transcript": {"present": True, "changed": True}, + "screen": {"present": True, "changed": True}, + } + ) == ("active", ["screen", "transcript"]) + + +class TestResolvePredecessor: + def test_resolution_uses_chronological_prior_even_before_state_exists( + self, tmp_path, monkeypatch + ): + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + day = "20260102" + _make_segment(tmp_path, day, "A", "101000_300") + _make_segment(tmp_path, day, "A", "100000_300") + + assert resolve_predecessor(day, "A", "101000_300") == { + "day": day, + "stream": "A", + "segment": "100000_300", + } + + def test_interleaved_streams_are_filtered(self, tmp_path, monkeypatch): + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + day = "20260102" + _make_segment(tmp_path, day, "A", "100000_300") + _make_segment(tmp_path, day, "B", "100500_300") + _make_segment(tmp_path, day, "A", "101000_300") + + assert resolve_predecessor(day, "A", "101000_300") == { + "day": day, + "stream": "A", + "segment": "100000_300", + } + + def test_first_of_stream_has_no_predecessor(self, tmp_path, monkeypatch): + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + day = "20260102" + _make_segment(tmp_path, day, "A", "100000_300") + + assert resolve_predecessor(day, "A", "100000_300") is None + + def test_gap_above_threshold_has_no_predecessor(self, tmp_path, monkeypatch): + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + day = "20260102" + _make_segment(tmp_path, day, "A", "100000_300") + _make_segment(tmp_path, day, "A", "102000_300") + + assert resolve_predecessor(day, "A", "102000_300") is None + + +class TestDetectSegmentChange: + def test_no_predecessor_with_present_sensor_is_active(self, tmp_path, monkeypatch): + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + day = "20260102" + seg_dir = _make_segment(tmp_path, day, "A", "100000_300") + _write_screen( + seg_dir, + first="0000000000000000", + last="0000000000000000", + ) + + result = detect_segment_change( + day, + "A", + "100000_300", + seg_dir, + predecessor=None, + timestamp="2026-01-02T10:00:00+00:00", + ) + + assert result["change_class"] == "active" + assert result["changed_sensors"] == ["screen"] + + def test_missing_predecessor_state_with_present_sensor_is_active( + self, tmp_path, monkeypatch + ): + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + day = "20260102" + _make_segment(tmp_path, day, "A", "100000_300") + seg_dir = _make_segment(tmp_path, day, "A", "101000_300") + _write_screen( + seg_dir, + first="0000000000000000", + last="0000000000000000", + ) + predecessor = {"day": day, "stream": "A", "segment": "100000_300"} + + result = detect_segment_change( + day, + "A", + "101000_300", + seg_dir, + predecessor=predecessor, + timestamp="2026-01-02T10:10:00+00:00", + ) + + assert result["predecessor"] == predecessor + assert result["change_class"] == "active" + assert result["changed_sensors"] == ["screen"] + + def test_empty_segment_without_predecessor_is_idle(self, tmp_path, monkeypatch): + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + day = "20260102" + seg_dir = _make_segment(tmp_path, day, "A", "100000_300") + + result = detect_segment_change( + day, + "A", + "100000_300", + seg_dir, + predecessor=None, + timestamp="2026-01-02T10:00:00+00:00", + ) + + assert result["change_class"] == "idle" + assert result["changed_sensors"] == [] + + +def test_assemble_transcript_excludes_document_extraction_jsonl(tmp_path): + seg_dir = tmp_path / "segment" + seg_dir.mkdir() + (seg_dir / "audio.jsonl").write_text( + json.dumps({"raw": "audio.flac"}) + + "\n" + + json.dumps({"start": "00:00:01", "text": "spoken words"}) + + "\n", + encoding="utf-8", + ) + (seg_dir / "document.jsonl").write_text( + json.dumps({"kind": "document"}) + + "\n" + + json.dumps({"start": "00:00:01", "text": "document words " * 20}) + + "\n", + encoding="utf-8", + ) + + state = assemble_sensor_state(seg_dir) + + assert state["transcript"]["present"] is True + assert state["transcript"]["word_count"] < 20 diff --git a/tests/test_describe_promote.py b/tests/test_describe_promote.py index e9a22b340..918b576fe 100644 --- a/tests/test_describe_promote.py +++ b/tests/test_describe_promote.py @@ -38,6 +38,9 @@ def _frame(frame_id: int, timestamp: float, frame_bytes: bytes) -> dict: def _processor(video_path: Path, frames: list[dict], monkeypatch) -> object: processor = describe_module.VideoProcessor.__new__(describe_module.VideoProcessor) processor.video_path = video_path + processor.first_hash = None + processor.last_hash = None + processor.qualified_count = len(frames) processor.qualified_frames = [] monkeypatch.setattr(processor, "process", lambda: frames) return processor @@ -73,6 +76,26 @@ def _install_fakes(monkeypatch, outcomes: dict[int, dict]) -> list[tuple]: return emitted +def test_build_metadata_header_includes_static_single_frame_hash(tmp_path, monkeypatch): + video_path = _video_path(tmp_path) + processor = describe_module.VideoProcessor.__new__(describe_module.VideoProcessor) + processor.video_path = video_path + processor.first_hash = 0x1234 + processor.last_hash = 0x1234 + processor.qualified_count = 1 + monkeypatch.setenv("OBSERVER_NAME", "desk") + monkeypatch.setenv("SEGMENT_META", json.dumps({"stream": "default"})) + + assert processor._build_metadata_header() == { + "raw": video_path.name, + "observer": "desk", + "stream": "default", + "first_hash": "0000000000001234", + "last_hash": "0000000000001234", + "qualified_count": 1, + } + + class FakeBatch: instances = [] outcomes = {} @@ -162,7 +185,14 @@ async def test_success_with_mixed_results_promotes_byte_identical_jsonl( ) request_type = describe_module.RequestType.DESCRIBE.value - header = json.dumps({"raw": video_path.name}) + header = json.dumps( + { + "raw": video_path.name, + "first_hash": None, + "last_hash": None, + "qualified_count": 2, + } + ) frame1 = json.dumps( { "frame_id": 1, @@ -219,7 +249,17 @@ async def test_empty_run_promotes_header_only_file_for_event_precondition( assert output_path.exists() # async_main's completion event branch is unchanged and gated on this exists(). - assert output_path.read_text() == json.dumps({"raw": video_path.name}) + "\n" + assert output_path.read_text() == ( + json.dumps( + { + "raw": video_path.name, + "first_hash": None, + "last_hash": None, + "qualified_count": 0, + } + ) + + "\n" + ) _assert_no_describe_temp(output_path.parent) @@ -240,7 +280,17 @@ async def test_all_frames_failed_promotes_header_only_then_raises( ) assert output_path.exists() - assert output_path.read_text() == json.dumps({"raw": video_path.name}) + "\n" + assert output_path.read_text() == ( + json.dumps( + { + "raw": video_path.name, + "first_hash": None, + "last_hash": None, + "qualified_count": 1, + } + ) + + "\n" + ) _assert_no_describe_temp(output_path.parent) diff --git a/tests/test_sense_splitter.py b/tests/test_sense_splitter.py index c9800dcd8..8da855ebc 100644 --- a/tests/test_sense_splitter.py +++ b/tests/test_sense_splitter.py @@ -63,13 +63,9 @@ class TestWriteSenseOutputs: density = json.loads((agents_dir / "density.json").read_text(encoding="utf-8")) assert set(density.keys()) == { "classification", - "transcript_lines", - "screen_frames", "timestamp", } assert density["classification"] == "active" - assert density["transcript_lines"] == 0 - assert density["screen_frames"] == 0 assert json.loads((agents_dir / "sense.json").read_text(encoding="utf-8")) == ( sense_json @@ -208,10 +204,45 @@ class TestWriteIdleStubs: agents_dir = seg_dir / "talents" assert (agents_dir / "density.json").exists() density = json.loads((agents_dir / "density.json").read_text(encoding="utf-8")) + assert set(density.keys()) == {"classification", "timestamp"} assert density["classification"] == "idle" - assert density["transcript_lines"] == 0 - assert density["screen_frames"] == 0 assert not (agents_dir / "activity.md").exists() assert not (agents_dir / "facets.json").exists() assert not (agents_dir / "speakers.json").exists() assert not (agents_dir / "sense.json").exists() + + +class TestWriteChangeDetection: + def test_round_trip(self, tmp_path): + from solstone.think.sense_splitter import write_change_detection + + seg_dir = Path(tmp_path) / "20260304" / "default" / "090000_300" + result = { + "timestamp": "2026-03-04T09:00:00+00:00", + "predecessor": None, + "change_class": "active", + "changed_sensors": ["screen"], + "sensors": { + "screen": { + "monitors": { + "center:DP-3": { + "first_hash": "0000000000000000", + "last_hash": "0000000000000001", + "qualified_count": 2, + } + } + }, + "transcript": { + "present": False, + "word_count": 0, + "content_hash": None, + }, + }, + } + + write_change_detection(seg_dir, result) + + stored = json.loads( + (seg_dir / "talents" / "change.json").read_text(encoding="utf-8") + ) + assert stored == result diff --git a/tests/test_think_segment.py b/tests/test_think_segment.py index 344f8ce38..d7c8b6848 100644 --- a/tests/test_think_segment.py +++ b/tests/test_think_segment.py @@ -36,6 +36,11 @@ def _segment_configs(*names: str) -> dict[str, dict]: "type": "cogitate", "schedule": "segment", }, + "documents": { + "priority": 20, + "type": "cogitate", + "schedule": "segment", + }, "timeline:segment_summary": { "priority": 41, "type": "generate", @@ -345,6 +350,66 @@ class TestRunSegmentSense: assert spawned == expected + def test_non_idle_dispatch_set_unchanged_with_change_detection( + self, segment_dir, monkeypatch + ): + from solstone.think import thinking as think + + spawned = [] + (segment_dir / "audio.npz").write_bytes(b"npz") + _write_sense_output( + segment_dir, + { + "density": "active", + "recommend": { + "screen_record": True, + "speaker_attribution": True, + }, + "facets": [], + }, + ) + + monkeypatch.setattr( + think, + "get_talent_configs", + lambda schedule=None, **kwargs: _segment_configs( + "sense", + "entities", + "documents", + "timeline:segment_summary", + "screen", + "speaker_attribution", + ), + ) + monkeypatch.setattr( + think, + "cortex_request", + lambda prompt, name, config=None: spawned.append(name) or f"agent-{name}", + ) + monkeypatch.setattr( + think, + "wait_for_uses", + lambda agent_ids, timeout=600: ({aid: "finish" for aid in agent_ids}, []), + ) + monkeypatch.setattr(think, "_callosum", None) + + think.run_segment_sense( + "20240115", + "120000_300", + refresh=False, + verbose=False, + stream="default", + ) + + assert spawned == [ + "sense", + "entities", + "documents", + "timeline:segment_summary", + "screen", + "speaker_attribution", + ] + def test_refresh_bypasses_idle(self, segment_dir, monkeypatch): from solstone.think import thinking as think @@ -1032,8 +1097,16 @@ class TestThinkJSONLEvents: for line in jsonl_path.read_text(encoding="utf-8").strip().splitlines() ] skips = [event for event in events if event["event"] == "talent.skip"] + dispatches = [event for event in events if event["event"] == "talent.dispatch"] + change_events = [ + event for event in events if event["event"] == "sense.change_detect" + ] assert any(skip["reason"] == "density_idle" for skip in skips) + assert [event["name"] for event in dispatches] == ["sense"] + assert len(change_events) == 1 + assert change_events[0]["change_class"] == "idle" + assert "changed_sensors" in change_events[0] def test_sense_complete_and_skip_events(self, segment_dir, monkeypatch): from solstone.think import thinking as think