personal memory agent
1# SPDX-License-Identifier: AGPL-3.0-only
2# Copyright (c) 2026 sol pbc
3
4from __future__ import annotations
5
6import argparse
7import inspect
8import json
9import re
10import sqlite3
11from pathlib import Path
12from types import SimpleNamespace
13
14import pytest
15
16from solstone.apps.observer.prune import format_result, result_exit_code, run_prune
17from solstone.apps.observer.utils import (
18 append_history_record,
19 list_observers,
20 load_history,
21 revoke_observer_record,
22 save_observer,
23)
24from solstone.observe import observer_cli
25from solstone.observe.processing_record import (
26 HANDLER_TRANSCRIBE,
27 SCHEMA,
28 STATE_EMPTY,
29)
30from solstone.think.indexer.journal import get_journal_index
31from solstone.think.link.auth import AuthorizedClients
32from solstone.think.link.paths import authorized_clients_path
33from solstone.think.streams import (
34 get_stream_state,
35 read_segment_stream,
36 write_segment_stream,
37)
38from solstone.think.utils import iter_segments
39
40
41@pytest.fixture
42def observer_cli_env(tmp_path: Path, monkeypatch: pytest.MonkeyPatch):
43 home = tmp_path / "home"
44 journal = tmp_path / "journal"
45 home.mkdir()
46 journal.mkdir()
47 monkeypatch.setattr(Path, "home", lambda: home)
48 monkeypatch.setenv("HOME", str(home))
49 monkeypatch.setenv("SOLSTONE_JOURNAL", str(journal))
50
51 import solstone.convey.state as convey_state
52
53 convey_state.journal_root = ""
54 return SimpleNamespace(home=home, journal=journal)
55
56
57def _observer(
58 name: str = "archon",
59 key: str = "existing-key-abcdef",
60 *,
61 last_seen: int | None = None,
62 last_segment: str | None = None,
63 last_segment_received_at: object = None,
64 last_segment_day: object = None,
65 include_last_segment_freshness: bool = True,
66) -> dict:
67 record = {
68 "key": key,
69 "name": name,
70 "created_at": 1,
71 "last_seen": last_seen,
72 "last_segment": last_segment,
73 "enabled": True,
74 "stats": {"segments_received": 0, "bytes_received": 0},
75 }
76 if include_last_segment_freshness:
77 record["last_segment_received_at"] = last_segment_received_at
78 record["last_segment_day"] = last_segment_day
79 return record
80
81
82PRUNE_DAY = "20250103"
83PRUNE_STREAM = "field"
84PRUNE_AUDIO = b"identical legacy audio bytes"
85PRUNE_SCREEN = b"identical legacy screen bytes"
86
87
88def _observer_for_stream(
89 stream: str = PRUNE_STREAM, key: str = "field-key-abcdef"
90) -> dict:
91 record = _observer(name=stream, key=key)
92 record["stream"] = stream
93 return record
94
95
96def _processing_row(size: int) -> str:
97 return (
98 json.dumps(
99 {
100 "raw": "audio.flac",
101 "_solstone_processing": {
102 "schema": SCHEMA,
103 "state": STATE_EMPTY,
104 "handler": HANDLER_TRANSCRIBE,
105 "input_size": size,
106 },
107 }
108 )
109 + "\n"
110 )
111
112
113def _write_prune_segment(
114 journal: Path,
115 *,
116 segment: str,
117 seq: int,
118 prev_segment: str | None,
119 audio: bytes = PRUNE_AUDIO,
120 screen: bytes = PRUNE_SCREEN,
121 manifest: bool = False,
122 marker: bool = True,
123 unknown_file: bool = False,
124 extra_manifest_content: bool = False,
125 proof_only_audio: bool = False,
126) -> Path:
127 seg_dir = journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / segment
128 seg_dir.mkdir(parents=True)
129 if marker:
130 write_segment_stream(
131 seg_dir,
132 PRUNE_STREAM,
133 PRUNE_DAY if prev_segment else None,
134 prev_segment,
135 seq,
136 )
137 if not proof_only_audio:
138 (seg_dir / "audio.flac").write_bytes(audio)
139 (seg_dir / "audio.jsonl").write_text(
140 json.dumps({"segment": segment, "text": f"transcript {segment}"}) + "\n",
141 encoding="utf-8",
142 )
143 else:
144 (seg_dir / "audio.jsonl").write_text(
145 _processing_row(len(audio)),
146 encoding="utf-8",
147 )
148 (seg_dir / "screen.mp4").write_bytes(screen)
149 (seg_dir / "screen.jsonl").write_text(
150 json.dumps({"segment": segment, "text": f"description {segment}"}) + "\n",
151 encoding="utf-8",
152 )
153 (seg_dir / "events.jsonl").write_text(
154 json.dumps({"event": segment}) + "\n",
155 encoding="utf-8",
156 )
157 talents = seg_dir / "talents"
158 talents.mkdir()
159 (talents / "sense.json").write_text(
160 json.dumps({"segment": segment}),
161 encoding="utf-8",
162 )
163 if unknown_file:
164 (seg_dir / "notes.txt").write_text("operator note", encoding="utf-8")
165 manifest_files = None
166 if manifest:
167 manifest_files = {
168 "audio.flac": {"sha256": _sha(audio), "size": len(audio)},
169 "screen.mp4": {"sha256": _sha(screen), "size": len(screen)},
170 }
171 if extra_manifest_content:
172 extra = b"unique manifest content"
173 (seg_dir / "capture.bin").write_bytes(extra)
174 manifest_files["capture.bin"] = {"sha256": _sha(extra), "size": len(extra)}
175 (seg_dir / "ingest.json").write_text(
176 json.dumps(
177 {
178 "schema_version": 1,
179 "requested_segment": segment,
180 "files": manifest_files,
181 }
182 )
183 + "\n",
184 encoding="utf-8",
185 )
186 return seg_dir
187
188
189def _sha(data: bytes) -> str:
190 import hashlib
191
192 return hashlib.sha256(data).hexdigest()
193
194
195def _append_upload(prefix: str, segment: str, *, sha: str | None = None) -> None:
196 append_history_record(
197 prefix,
198 PRUNE_DAY,
199 {
200 "ts": 1,
201 "segment": segment,
202 "stream": PRUNE_STREAM,
203 "files": [
204 {
205 "submitted": "audio.flac",
206 "written": "audio.flac",
207 "size": len(PRUNE_AUDIO),
208 "sha256": sha or _sha(PRUNE_AUDIO),
209 }
210 ],
211 },
212 )
213
214
215def _write_stream_state(journal: Path, *, last_segment: str, seq: int = 9) -> None:
216 streams = journal / "streams"
217 streams.mkdir()
218 (streams / f"{PRUNE_STREAM}.json").write_text(
219 json.dumps(
220 {
221 "name": PRUNE_STREAM,
222 "type": "observer",
223 "host": "field-host",
224 "platform": "linux",
225 "created_at": 123,
226 "last_day": PRUNE_DAY,
227 "last_segment": last_segment,
228 "seq": seq,
229 }
230 )
231 + "\n",
232 encoding="utf-8",
233 )
234
235
236def _insert_index_rows(journal: Path, *segments: str) -> None:
237 conn, _db_path = get_journal_index(str(journal))
238 try:
239 for segment in segments:
240 rel = f"{PRUNE_DAY}/{PRUNE_STREAM}/{segment}"
241 conn.execute(
242 "INSERT INTO files(path, mtime) VALUES (?, ?)", (f"{rel}/audio.flac", 1)
243 )
244 conn.execute(
245 "INSERT INTO chunks(content, path, day, facet, agent, stream, idx, time_bucket) "
246 "VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
247 ("content", rel, PRUNE_DAY, "", "segment", PRUNE_STREAM, 0, "morning"),
248 )
249 conn.commit()
250 finally:
251 conn.close()
252
253
254def _index_count(journal: Path, segment: str) -> int:
255 db_path = journal / "indexer" / "journal.sqlite"
256 conn = sqlite3.connect(db_path)
257 try:
258 rel = f"{PRUNE_DAY}/{PRUNE_STREAM}/{segment}"
259 return conn.execute(
260 "SELECT count(*) FROM chunks WHERE path = ? OR path LIKE ?",
261 (rel, f"{rel}/%"),
262 ).fetchone()[0]
263 finally:
264 conn.close()
265
266
267def _journal_snapshot(journal: Path) -> list[tuple[str, str, bytes]]:
268 entries = []
269 for path in sorted(journal.rglob("*")):
270 rel = path.relative_to(journal).as_posix()
271 if path.is_dir():
272 entries.append((rel, "dir", b""))
273 elif path.is_file():
274 entries.append((rel, "file", path.read_bytes()))
275 return entries
276
277
278def _stream_marker_snapshot(journal: Path) -> dict[str, dict]:
279 snapshot = {}
280 for stream, segment, path in iter_segments(PRUNE_DAY):
281 if stream != PRUNE_STREAM:
282 continue
283 snapshot[segment] = read_segment_stream(path)
284 return snapshot
285
286
287def _pruned_history(prefix: str, segment: str) -> list[dict]:
288 return [
289 record
290 for record in load_history(prefix, PRUNE_DAY)
291 if record.get("type") == "pruned" and record.get("segment") == segment
292 ]
293
294
295def _build_legacy_prune_lattice(journal: Path, prefix: str) -> None:
296 _write_prune_segment(
297 journal, segment="080000_300", seq=1, prev_segment=None, audio=b"before"
298 )
299 _write_prune_segment(
300 journal, segment="090000_300", seq=2, prev_segment="080000_300"
301 )
302 _write_prune_segment(
303 journal, segment="090000_301", seq=3, prev_segment="090000_300"
304 )
305 _write_prune_segment(
306 journal, segment="090000_302", seq=4, prev_segment="090000_301"
307 )
308 _write_prune_segment(
309 journal,
310 segment="090000_303",
311 seq=5,
312 prev_segment="090000_302",
313 audio=b"near duplicate audio",
314 )
315 _write_prune_segment(
316 journal,
317 segment="090000_304",
318 seq=6,
319 prev_segment="090000_303",
320 unknown_file=True,
321 )
322 _write_prune_segment(
323 journal, segment="100000_300", seq=7, prev_segment="090000_304"
324 )
325 _write_prune_segment(
326 journal,
327 segment="100500_300",
328 seq=8,
329 prev_segment="100000_300",
330 audio=PRUNE_AUDIO,
331 screen=PRUNE_SCREEN,
332 )
333 for segment in (
334 "090000_300",
335 "090000_301",
336 "090000_302",
337 "090000_303",
338 "090000_304",
339 "110000_300",
340 ):
341 _append_upload(prefix, segment)
342 _write_stream_state(journal, last_segment="100500_300")
343 _insert_index_rows(journal, "090000_301", "090000_302")
344
345
346def _observer_with_stats(
347 *,
348 name: str,
349 key: str,
350 created_at: int,
351 segments_received: int,
352 bytes_received: int,
353 duplicates_rejected: int = 0,
354 last_segment_received_at: object = None,
355 last_segment_day: object = None,
356 include_last_segment_freshness: bool = True,
357) -> dict:
358 record = _observer(
359 name=name,
360 key=key,
361 last_segment_received_at=last_segment_received_at,
362 last_segment_day=last_segment_day,
363 include_last_segment_freshness=include_last_segment_freshness,
364 )
365 record["created_at"] = created_at
366 record["stats"] = {
367 "segments_received": segments_received,
368 "bytes_received": bytes_received,
369 "duplicates_rejected": duplicates_rejected,
370 }
371 return record
372
373
374def _table_row(output: str, name: str) -> str:
375 return next(line for line in output.splitlines() if line.startswith(f"{name:<20}"))
376
377
378def _table_header(output: str) -> str:
379 return next(
380 line
381 for line in output.splitlines()
382 if line.startswith("Name") and "Last Segment" in line
383 )
384
385
386def _table_cell(output: str, row: str, column: str) -> str:
387 header = _table_header(output)
388 starts = sorted(
389 header.index(name)
390 for name in (
391 "Name",
392 "Prefix",
393 "Status",
394 "Last Seen",
395 "Last Segment",
396 "Segments",
397 "Bytes",
398 )
399 if name in header
400 )
401 start = header.index(column)
402 end = next((pos for pos in starts if pos > start), len(row))
403 return row[start:end].strip()
404
405
406def _last_segment_cell(output: str, row: str) -> str:
407 return _table_cell(output, row, "Last Segment")
408
409
410def test_create_observer_record_reuses_existing_without_create_side_effects(
411 observer_cli_env,
412 monkeypatch: pytest.MonkeyPatch,
413) -> None:
414 existing = _observer()
415 assert save_observer(existing)
416 monkeypatch.setattr(
417 observer_cli,
418 "_generate_key",
419 lambda: pytest.fail("reuse must not generate a new key"),
420 )
421 monkeypatch.setattr(
422 observer_cli,
423 "save_observer",
424 lambda _data: pytest.fail("reuse must not save"),
425 )
426 monkeypatch.setattr(
427 observer_cli,
428 "log_app_action",
429 lambda **_kwargs: pytest.fail("reuse must not log observer_create"),
430 )
431
432 record, key, reused = observer_cli.create_observer_record(
433 "archon", reuse_existing=True
434 )
435
436 assert record["key"] == existing["key"]
437 assert record["name"] == existing["name"]
438 assert record["filename_prefix"] == "existing"
439 assert key == "existing-key-abcdef"
440 assert reused is True
441 assert list_observers() == [record]
442
443
444def test_create_observer_record_fresh_create_returns_reused_false_and_logs(
445 observer_cli_env,
446 monkeypatch: pytest.MonkeyPatch,
447) -> None:
448 logs = []
449 monkeypatch.setattr(observer_cli, "_generate_key", lambda: "fresh-key-abcdef")
450 monkeypatch.setattr(
451 observer_cli, "log_app_action", lambda **kwargs: logs.append(kwargs)
452 )
453
454 record, key, reused = observer_cli.create_observer_record("archon")
455
456 assert key == "fresh-key-abcdef"
457 assert reused is False
458 assert record["name"] == "archon"
459 assert list_observers()[0]["key"] == "fresh-key-abcdef"
460 assert logs == [
461 {
462 "app": "observer",
463 "facet": None,
464 "action": "observer_create",
465 "params": {"name": "archon", "key_prefix": "fresh-ke"},
466 }
467 ]
468
469
470def test_create_observer_record_duplicate_without_reuse_still_fails(
471 observer_cli_env,
472) -> None:
473 assert save_observer(_observer())
474
475 with pytest.raises(ValueError, match="observer already exists: archon"):
476 observer_cli.create_observer_record("archon")
477
478
479def test_cmd_create_duplicate_without_reuse_exits_one(
480 observer_cli_env,
481 capsys: pytest.CaptureFixture[str],
482) -> None:
483 assert save_observer(_observer())
484 args = argparse.Namespace(
485 name="archon",
486 json_output=False,
487 reuse_existing=False,
488 )
489
490 rc = observer_cli.cmd_create(args)
491
492 captured = capsys.readouterr()
493 assert rc == 1
494 assert captured.out == ""
495 assert captured.err == "Error: observer 'archon' already exists\n"
496
497
498def test_cmd_create_reuse_existing_json_shape(
499 observer_cli_env,
500 capsys: pytest.CaptureFixture[str],
501) -> None:
502 existing = _observer()
503 assert save_observer(existing)
504 args = argparse.Namespace(
505 name="archon",
506 json_output=True,
507 reuse_existing=True,
508 )
509
510 rc = observer_cli.cmd_create(args)
511
512 captured = capsys.readouterr()
513 assert rc == 0
514 assert captured.err == ""
515 assert captured.out == (
516 json.dumps(
517 {
518 "name": "archon",
519 "key": "existing-key-abcdef",
520 "prefix": "existing",
521 }
522 )
523 + "\n"
524 )
525
526
527def test_cmd_create_reuse_existing_human_header(
528 observer_cli_env,
529 capsys: pytest.CaptureFixture[str],
530) -> None:
531 existing = _observer()
532 assert save_observer(existing)
533 args = argparse.Namespace(
534 name="archon",
535 json_output=False,
536 reuse_existing=True,
537 )
538
539 rc = observer_cli.cmd_create(args)
540
541 captured = capsys.readouterr()
542 assert rc == 0
543 assert captured.err == ""
544 assert "Reusing existing observer:" in captured.out
545 assert "Observer created:" not in captured.out
546 assert " api key: existing-key-abcdef" in captured.out
547
548
549def test_cmd_create_reuse_existing_creates_normally_when_absent(
550 observer_cli_env,
551 monkeypatch: pytest.MonkeyPatch,
552 capsys: pytest.CaptureFixture[str],
553) -> None:
554 logs = []
555 monkeypatch.setattr(observer_cli, "_generate_key", lambda: "fresh-key-abcdef")
556 monkeypatch.setattr(
557 observer_cli, "log_app_action", lambda **kwargs: logs.append(kwargs)
558 )
559 args = argparse.Namespace(
560 name="archon",
561 json_output=False,
562 reuse_existing=True,
563 )
564
565 rc = observer_cli.cmd_create(args)
566
567 captured = capsys.readouterr()
568 assert rc == 0
569 assert captured.err == ""
570 assert "Observer created:" in captured.out
571 assert "Reusing existing observer:" not in captured.out
572 assert " api key: fresh-key-abcdef" in captured.out
573 assert list_observers()[0]["key"] == "fresh-key-abcdef"
574 assert logs == [
575 {
576 "app": "observer",
577 "facet": None,
578 "action": "observer_create",
579 "params": {"name": "archon", "key_prefix": "fresh-ke"},
580 }
581 ]
582
583
584def test_reconcile_collapses_duplicates_oldest_survives(observer_cli_env) -> None:
585 assert save_observer(
586 _observer_with_stats(
587 name="fedora.tmux",
588 key="newest03-key",
589 created_at=3,
590 segments_received=5,
591 bytes_received=100,
592 duplicates_rejected=1,
593 )
594 )
595 assert save_observer(
596 _observer_with_stats(
597 name="fedora.tmux",
598 key="oldest01-key",
599 created_at=1,
600 segments_received=7,
601 bytes_received=200,
602 duplicates_rejected=2,
603 )
604 )
605 assert save_observer(
606 _observer_with_stats(
607 name="fedora.tmux",
608 key="middle02-key",
609 created_at=2,
610 segments_received=11,
611 bytes_received=300,
612 )
613 )
614 lone = _observer_with_stats(
615 name="fedora",
616 key="desktop1-key",
617 created_at=4,
618 segments_received=13,
619 bytes_received=400,
620 duplicates_rejected=5,
621 )
622 assert save_observer(lone)
623
624 plan = observer_cli.reconcile_observers(dry_run=False)
625
626 assert plan == [
627 {
628 "name": "fedora.tmux",
629 "survivor_prefix": "oldest01",
630 "revoked_prefixes": ["newest03", "middle02"],
631 "stats": {
632 "segments_received": 23,
633 "bytes_received": 600,
634 "duplicates_rejected": 3,
635 },
636 }
637 ]
638 records = list_observers()
639 tmux_records = [record for record in records if record["name"] == "fedora.tmux"]
640 unrevoked_tmux = [
641 record for record in tmux_records if not record.get("revoked", False)
642 ]
643 assert len(unrevoked_tmux) == 1
644 assert unrevoked_tmux[0]["created_at"] == 1
645 assert unrevoked_tmux[0]["stats"] == {
646 "segments_received": 23,
647 "bytes_received": 600,
648 "duplicates_rejected": 3,
649 }
650 revoked_tmux = [record for record in tmux_records if record.get("revoked", False)]
651 assert {record["created_at"] for record in revoked_tmux} == {2, 3}
652 lone_record = next(record for record in records if record["name"] == "fedora")
653 assert lone_record.get("revoked", False) is False
654 assert lone_record["stats"] == lone["stats"]
655
656
657def test_reconcile_dry_run_mutates_nothing(observer_cli_env) -> None:
658 assert save_observer(
659 _observer_with_stats(
660 name="fedora.tmux",
661 key="newest03-key",
662 created_at=3,
663 segments_received=5,
664 bytes_received=100,
665 duplicates_rejected=1,
666 )
667 )
668 assert save_observer(
669 _observer_with_stats(
670 name="fedora.tmux",
671 key="oldest01-key",
672 created_at=1,
673 segments_received=7,
674 bytes_received=200,
675 duplicates_rejected=2,
676 )
677 )
678 assert save_observer(
679 _observer_with_stats(
680 name="fedora.tmux",
681 key="middle02-key",
682 created_at=2,
683 segments_received=11,
684 bytes_received=300,
685 )
686 )
687 observers_dir = observer_cli_env.journal / "apps" / "observer" / "observers"
688 before = {path.name: path.read_bytes() for path in observers_dir.glob("*.json")}
689
690 plan = observer_cli.reconcile_observers(dry_run=True)
691
692 assert plan == [
693 {
694 "name": "fedora.tmux",
695 "survivor_prefix": "oldest01",
696 "revoked_prefixes": ["newest03", "middle02"],
697 "stats": {
698 "segments_received": 23,
699 "bytes_received": 600,
700 "duplicates_rejected": 3,
701 },
702 }
703 ]
704 after = {path.name: path.read_bytes() for path in observers_dir.glob("*.json")}
705 assert after == before
706
707
708def test_reconcile_lone_stream_returns_empty_plan(observer_cli_env) -> None:
709 lone = _observer_with_stats(
710 name="fedora",
711 key="desktop1-key",
712 created_at=1,
713 segments_received=13,
714 bytes_received=400,
715 duplicates_rejected=5,
716 )
717 assert save_observer(lone)
718
719 plan = observer_cli.reconcile_observers(dry_run=False)
720
721 assert plan == []
722 records = list_observers()
723 assert len(records) == 1
724 assert records[0].get("revoked", False) is False
725 assert records[0]["stats"] == lone["stats"]
726
727
728def test_cmd_reconcile_reports_plan(
729 observer_cli_env,
730 capsys: pytest.CaptureFixture[str],
731) -> None:
732 assert save_observer(
733 _observer_with_stats(
734 name="fedora.tmux",
735 key="newest03-key",
736 created_at=3,
737 segments_received=5,
738 bytes_received=100,
739 )
740 )
741 assert save_observer(
742 _observer_with_stats(
743 name="fedora.tmux",
744 key="oldest01-key",
745 created_at=1,
746 segments_received=7,
747 bytes_received=200,
748 )
749 )
750
751 rc = observer_cli.cmd_reconcile(
752 argparse.Namespace(dry_run=False, json_output=False)
753 )
754
755 captured = capsys.readouterr()
756 assert rc == 0
757 assert captured.err == ""
758 assert "Reconciled stream 'fedora.tmux':" in captured.out
759 assert " survivor: oldest01" in captured.out
760 assert " revoking: newest03" in captured.out
761
762
763def test_cmd_reconcile_no_duplicates(
764 observer_cli_env,
765 capsys: pytest.CaptureFixture[str],
766) -> None:
767 assert save_observer(_observer(name="fedora", key="desktop1-key"))
768
769 rc = observer_cli.cmd_reconcile(
770 argparse.Namespace(dry_run=False, json_output=False)
771 )
772
773 captured = capsys.readouterr()
774 assert rc == 0
775 assert captured.err == ""
776 assert captured.out == "No duplicate observer streams to reconcile.\n"
777
778
779def test_cmd_list_json_includes_prefix_and_status(
780 observer_cli_env,
781 capsys: pytest.CaptureFixture[str],
782) -> None:
783 assert save_observer(_observer(name="desktop", key="abcdefgh12345678"))
784 args = argparse.Namespace(json_output=True)
785
786 rc = observer_cli.cmd_list(args)
787
788 captured = capsys.readouterr()
789 assert rc == 0
790 rows = {row["name"]: row for row in json.loads(captured.out)}
791 assert rows["desktop"]["prefix"] == "abcdefgh"
792 assert rows["desktop"]["status"] == "disconnected"
793 assert "mode" not in rows["desktop"]
794
795
796def test_fmt_compact_age_units_and_guards(
797 monkeypatch: pytest.MonkeyPatch,
798) -> None:
799 monkeypatch.setattr(observer_cli, "now_ms", lambda: 0)
800 assert observer_cli._fmt_compact_age(0) == "0s"
801
802 now = 2_000_000_000_000
803 monkeypatch.setattr(observer_cli, "now_ms", lambda: now)
804 assert observer_cli._fmt_compact_age(None) == "—"
805 assert observer_cli._fmt_compact_age("bad") == "—"
806 assert observer_cli._fmt_compact_age(True) == "—"
807 assert observer_cli._fmt_compact_age(-1) == "—"
808 assert observer_cli._fmt_compact_age(now + 1) == "—"
809 assert observer_cli._fmt_compact_age(now) == "0s"
810 assert observer_cli._fmt_compact_age(now - 30_000) == "30s"
811 assert observer_cli._fmt_compact_age(now - ((59 * 60 + 59) * 1000)) == "59m"
812 assert observer_cli._fmt_compact_age(now - (60 * 60 * 1000)) == "1h"
813 assert observer_cli._fmt_compact_age(now - ((23 * 60 + 59) * 60 * 1000)) == "23h"
814 assert observer_cli._fmt_compact_age(now - (24 * 60 * 60 * 1000)) == "1d"
815 assert observer_cli._fmt_compact_age(now - int(19.5 * 60 * 60 * 1000)) == "19h"
816
817
818def test_cmd_list_shows_last_segment_column_and_json(
819 observer_cli_env,
820 monkeypatch: pytest.MonkeyPatch,
821 capsys: pytest.CaptureFixture[str],
822) -> None:
823 now = 2_000_000_000_000
824 monkeypatch.setattr(observer_cli, "now_ms", lambda: now)
825 assert save_observer(
826 _observer(
827 name="desktop",
828 key="abcdefgh12345678",
829 last_segment_received_at=now - 2 * 60 * 1000,
830 last_segment_day="20260724",
831 )
832 )
833
834 rc = observer_cli.cmd_list(argparse.Namespace(json_output=False))
835
836 captured = capsys.readouterr()
837 assert rc == 0
838 assert "Last Seen Last Segment" in captured.out
839 assert "-" * 107 in captured.out
840 assert _last_segment_cell(captured.out, _table_row(captured.out, "desktop")) == "2m"
841
842 rc = observer_cli.cmd_list(argparse.Namespace(json_output=True))
843
844 captured = capsys.readouterr()
845 assert rc == 0
846 rows = {row["name"]: row for row in json.loads(captured.out)}
847 assert rows["desktop"]["last_segment_received_at"] == now - 2 * 60 * 1000
848 assert rows["desktop"]["last_segment_day"] == "20260724"
849
850
851def test_cmd_list_human_shows_prefix_column(
852 observer_cli_env,
853 capsys: pytest.CaptureFixture[str],
854) -> None:
855 assert save_observer(_observer(name="desktop", key="abcdefgh12345678"))
856 args = argparse.Namespace(json_output=False)
857
858 rc = observer_cli.cmd_list(args)
859
860 captured = capsys.readouterr()
861 assert rc == 0
862 assert "Name Prefix" in captured.out
863 assert "Mode" not in captured.out
864 assert "desktop abcdefgh" in captured.out
865
866
867def test_cmd_status_single_reports_prefix(
868 observer_cli_env,
869 capsys: pytest.CaptureFixture[str],
870) -> None:
871 assert save_observer(_observer(name="desktop", key="cdefghij12345678"))
872
873 rc = observer_cli.cmd_status(
874 argparse.Namespace(identifier="desktop", json_output=True)
875 )
876
877 captured = capsys.readouterr()
878 assert rc == 0
879 payload = json.loads(captured.out)
880 assert payload["prefix"] == "cdefghij"
881 assert payload["status"] == "disconnected"
882 assert "mode" not in payload
883
884
885def test_cmd_status_single_last_segment_age_uses_receipt_time_with_day_context(
886 observer_cli_env,
887 monkeypatch: pytest.MonkeyPatch,
888 capsys: pytest.CaptureFixture[str],
889) -> None:
890 now = 2_000_000_000_000
891 received_at = now - 2 * 60 * 1000
892 monkeypatch.setattr(observer_cli, "now_ms", lambda: now)
893 assert save_observer(
894 _observer(
895 name="desktop",
896 key="cdefghij12345678",
897 last_seen=now - 1_000,
898 last_segment="120000_300",
899 last_segment_received_at=received_at,
900 last_segment_day="20260722",
901 )
902 )
903
904 rc = observer_cli.cmd_status(
905 argparse.Namespace(identifier="desktop", json_output=False)
906 )
907
908 captured = capsys.readouterr()
909 assert rc == 0
910 assert " Last segment: 2m (20260722)\n" in captured.out
911
912 rc = observer_cli.cmd_status(
913 argparse.Namespace(identifier="desktop", json_output=True)
914 )
915
916 captured = capsys.readouterr()
917 assert rc == 0
918 payload = json.loads(captured.out)
919 assert payload["last_segment_received_at"] == received_at
920 assert payload["last_segment_day"] == "20260722"
921
922
923def test_cmd_status_all_table_shows_prefix(
924 observer_cli_env,
925 capsys: pytest.CaptureFixture[str],
926) -> None:
927 assert save_observer(_observer(name="desktop", key="abcdefgh12345678"))
928
929 rc = observer_cli.cmd_status(argparse.Namespace(identifier=None, json_output=False))
930
931 captured = capsys.readouterr()
932 assert rc == 0
933 assert "Name Prefix" in captured.out
934 assert "Mode" not in captured.out
935 assert "desktop abcdefgh" in captured.out
936
937
938def test_cmd_status_all_shows_last_segment_column_and_json(
939 observer_cli_env,
940 monkeypatch: pytest.MonkeyPatch,
941 capsys: pytest.CaptureFixture[str],
942) -> None:
943 now = 2_000_000_000_000
944 monkeypatch.setattr(observer_cli, "now_ms", lambda: now)
945 assert save_observer(
946 _observer(
947 name="desktop",
948 key="abcdefgh12345678",
949 last_segment_received_at=now - 2 * 60 * 1000,
950 last_segment_day="20260724",
951 )
952 )
953
954 rc = observer_cli.cmd_status(argparse.Namespace(identifier=None, json_output=False))
955
956 captured = capsys.readouterr()
957 assert rc == 0
958 assert "Last Seen Last Segment" in captured.out
959 assert "-" * 87 in captured.out
960 assert _last_segment_cell(captured.out, _table_row(captured.out, "desktop")) == "2m"
961
962 rc = observer_cli.cmd_status(argparse.Namespace(identifier=None, json_output=True))
963
964 captured = capsys.readouterr()
965 assert rc == 0
966 payload = json.loads(captured.out)
967 row = payload["observers"][0]
968 assert row["last_segment_received_at"] == now - 2 * 60 * 1000
969 assert row["last_segment_day"] == "20260724"
970
971
972def test_last_segment_freshness_does_not_change_connection_status(
973 observer_cli_env,
974 monkeypatch: pytest.MonkeyPatch,
975 capsys: pytest.CaptureFixture[str],
976) -> None:
977 now = 2_000_000_000_000
978 monkeypatch.setattr(observer_cli, "now_ms", lambda: now)
979 assert observer_cli.CONNECTED_THRESHOLD_MS == 2 * 60 * 1000
980 assert save_observer(
981 _observer(
982 name="desktop",
983 key="abcdefgh12345678",
984 last_seen=now - 1_000,
985 last_segment_received_at=now - 41 * 24 * 60 * 60 * 1000,
986 last_segment_day="20260613",
987 )
988 )
989
990 rc = observer_cli.cmd_status(
991 argparse.Namespace(identifier="desktop", json_output=False)
992 )
993
994 captured = capsys.readouterr()
995 assert rc == 0
996 assert " Status: connected\n" in captured.out
997 assert " Last segment: 41d (20260613)\n" in captured.out
998
999
1000def test_fleet_views_show_fresh_last_seen_with_unknown_last_segment(
1001 observer_cli_env,
1002 monkeypatch: pytest.MonkeyPatch,
1003 capsys: pytest.CaptureFixture[str],
1004) -> None:
1005 now = 2_000_000_000_000
1006 monkeypatch.setattr(observer_cli, "now_ms", lambda: now)
1007 assert save_observer(
1008 _observer(
1009 name="desktop",
1010 key="abcdefgh12345678",
1011 last_seen=now - 30_000,
1012 include_last_segment_freshness=False,
1013 )
1014 )
1015
1016 def assert_divergent_row(output: str) -> None:
1017 row = _table_row(output, "desktop")
1018 assert _table_cell(output, row, "Status") == "connected"
1019 assert _table_cell(output, row, "Last Seen") != "never"
1020 assert _last_segment_cell(output, row) == "—"
1021
1022 rc = observer_cli.cmd_list(argparse.Namespace(json_output=False))
1023
1024 captured = capsys.readouterr()
1025 assert rc == 0
1026 assert_divergent_row(captured.out)
1027
1028 rc = observer_cli.cmd_status(argparse.Namespace(identifier=None, json_output=False))
1029
1030 captured = capsys.readouterr()
1031 assert rc == 0
1032 assert_divergent_row(captured.out)
1033
1034
1035def test_observer_cli_last_segment_rendering_has_no_classification() -> None:
1036 source = Path(observer_cli.__file__).read_text(encoding="utf-8")
1037 thresholds = re.findall(
1038 r"^([A-Z][A-Z0-9_]*THRESHOLD[A-Z0-9_]*)\s*=", source, flags=re.MULTILINE
1039 )
1040 assert thresholds == ["CONNECTED_THRESHOLD_MS"]
1041
1042 render_source = "\n".join(
1043 inspect.getsource(obj)
1044 for obj in (
1045 observer_cli.cmd_list,
1046 observer_cli._status_single,
1047 observer_cli._status_all,
1048 observer_cli._fmt_compact_age,
1049 )
1050 )
1051 assert "\\x1b" not in render_source
1052 assert "\\033" not in render_source
1053 assert "color" not in render_source.lower()
1054 assert "colour" not in render_source.lower()
1055 assert "stale" not in render_source.lower()
1056 for glyph in ("▲", "▼", "●", "○", "✕", "✓"):
1057 assert glyph not in render_source
1058
1059
1060def test_fleet_last_segment_bad_rows_are_isolated(
1061 observer_cli_env,
1062 monkeypatch: pytest.MonkeyPatch,
1063 capsys: pytest.CaptureFixture[str],
1064) -> None:
1065 now = 2_000_000_000_000
1066 monkeypatch.setattr(observer_cli, "now_ms", lambda: now)
1067 records = [
1068 _observer(
1069 name="good",
1070 key="good000012345678",
1071 last_segment_received_at=now - 2 * 60 * 1000,
1072 ),
1073 _observer(
1074 name="malformed",
1075 key="malform12345678",
1076 last_segment_received_at="bad",
1077 ),
1078 _observer(
1079 name="negative",
1080 key="negative12345678",
1081 last_segment_received_at=-1,
1082 ),
1083 _observer(
1084 name="future",
1085 key="future0012345678",
1086 last_segment_received_at=now + 1,
1087 ),
1088 ]
1089 for record in records:
1090 assert save_observer(record)
1091
1092 rc = observer_cli.cmd_list(argparse.Namespace(json_output=False))
1093
1094 captured = capsys.readouterr()
1095 assert rc == 0
1096 assert _last_segment_cell(captured.out, _table_row(captured.out, "good")) == "2m"
1097 assert (
1098 _last_segment_cell(captured.out, _table_row(captured.out, "malformed")) == "—"
1099 )
1100 assert _last_segment_cell(captured.out, _table_row(captured.out, "negative")) == "—"
1101 assert _last_segment_cell(captured.out, _table_row(captured.out, "future")) == "—"
1102
1103 rc = observer_cli.cmd_status(argparse.Namespace(identifier=None, json_output=False))
1104
1105 captured = capsys.readouterr()
1106 assert rc == 0
1107 assert _last_segment_cell(captured.out, _table_row(captured.out, "good")) == "2m"
1108 assert (
1109 _last_segment_cell(captured.out, _table_row(captured.out, "malformed")) == "—"
1110 )
1111 assert _last_segment_cell(captured.out, _table_row(captured.out, "negative")) == "—"
1112 assert _last_segment_cell(captured.out, _table_row(captured.out, "future")) == "—"
1113
1114
1115def test_prechange_record_uses_status_single_history_fallback_only(
1116 observer_cli_env,
1117 monkeypatch: pytest.MonkeyPatch,
1118 capsys: pytest.CaptureFixture[str],
1119) -> None:
1120 now = 2_000_000_000_000
1121 today = observer_cli.datetime.date.today().strftime("%Y%m%d")
1122 monkeypatch.setattr(observer_cli, "now_ms", lambda: now)
1123 assert save_observer(
1124 _observer(
1125 name="desktop",
1126 key="abcdefgh12345678",
1127 last_segment="120000_300",
1128 include_last_segment_freshness=False,
1129 )
1130 )
1131 append_history_record(
1132 "abcdefgh",
1133 today,
1134 {
1135 "ts": now - 2 * 60 * 1000,
1136 "segment": "120000_300",
1137 "stream": "desktop",
1138 "files": [],
1139 },
1140 )
1141 append_history_record(
1142 "abcdefgh",
1143 today,
1144 {"type": "observed", "ts": now - 10_000, "segment": "ignored"},
1145 )
1146
1147 rc = observer_cli.cmd_list(argparse.Namespace(json_output=False))
1148
1149 captured = capsys.readouterr()
1150 assert rc == 0
1151 assert _last_segment_cell(captured.out, _table_row(captured.out, "desktop")) == "—"
1152
1153 rc = observer_cli.cmd_status(argparse.Namespace(identifier=None, json_output=False))
1154
1155 captured = capsys.readouterr()
1156 assert rc == 0
1157 assert _last_segment_cell(captured.out, _table_row(captured.out, "desktop")) == "—"
1158
1159 rc = observer_cli.cmd_status(
1160 argparse.Namespace(identifier="desktop", json_output=False)
1161 )
1162
1163 captured = capsys.readouterr()
1164 assert rc == 0
1165 assert f" Last segment: 2m ({today})\n" in captured.out
1166
1167 rc = observer_cli.cmd_status(
1168 argparse.Namespace(identifier="desktop", json_output=True)
1169 )
1170
1171 captured = capsys.readouterr()
1172 assert rc == 0
1173 payload = json.loads(captured.out)
1174 assert payload["last_segment_received_at"] is None
1175 assert payload["last_segment_day"] is None
1176
1177
1178def test_revoke_dl_observer_leaves_authorized_clients_untouched(
1179 observer_cli_env,
1180) -> None:
1181 assert save_observer(_observer(name="desktop", key="abcdefgh12345678"))
1182 fingerprint = "sha256:" + ("f" * 64)
1183 authorized = AuthorizedClients(authorized_clients_path())
1184 authorized.add(
1185 fingerprint,
1186 "phone",
1187 "inst-1",
1188 paired_at="2026-05-20T00:00:00Z",
1189 )
1190 before = authorized_clients_path().read_bytes()
1191
1192 record = revoke_observer_record("desktop")
1193
1194 assert record["revoked"] is True
1195 assert authorized_clients_path().read_bytes() == before
1196 assert (
1197 AuthorizedClients(authorized_clients_path()).is_authorized(fingerprint) is True
1198 )
1199
1200
1201def test_prune_dry_run_lists_legacy_duplicates_and_writes_nothing(
1202 observer_cli_env,
1203) -> None:
1204 observer = _observer_for_stream()
1205 assert save_observer(observer)
1206 prefix = observer["key"][:8]
1207 _build_legacy_prune_lattice(observer_cli_env.journal, prefix)
1208
1209 before_history = load_history(prefix, PRUNE_DAY)
1210 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=False)
1211
1212 assert len(result.groups) == 1
1213 assert [c.analysis.segment for c in result.groups[0].candidates] == [
1214 "090000_301",
1215 "090000_302",
1216 ]
1217 assert {refusal.gate for refusal in result.refusals} == {
1218 "content-identity",
1219 "derived-output",
1220 }
1221 assert (
1222 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "090000_301"
1223 ).is_dir()
1224 assert load_history(prefix, PRUNE_DAY) == before_history
1225 assert not (observer_cli_env.journal / "chronicle" / PRUNE_DAY / "health").exists()
1226 assert _index_count(observer_cli_env.journal, "090000_301") == 1
1227
1228
1229def test_prune_dry_run_without_observer_storage_writes_nothing(
1230 observer_cli_env,
1231) -> None:
1232 _write_prune_segment(
1233 observer_cli_env.journal,
1234 segment="081000_300",
1235 seq=1,
1236 prev_segment=None,
1237 )
1238 _write_prune_segment(
1239 observer_cli_env.journal,
1240 segment="081000_301",
1241 seq=2,
1242 prev_segment="081000_300",
1243 )
1244 _write_stream_state(observer_cli_env.journal, last_segment="081000_301", seq=2)
1245 assert not (observer_cli_env.journal / "apps").exists()
1246 before = _journal_snapshot(observer_cli_env.journal)
1247
1248 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=False)
1249
1250 assert any(refusal.gate == "observer-attribution" for refusal in result.refusals)
1251 assert result.deleted == []
1252 assert _journal_snapshot(observer_cli_env.journal) == before
1253
1254
1255def test_prune_execute_deletes_duplicates_repairs_chain_history_and_index(
1256 observer_cli_env,
1257) -> None:
1258 observer = _observer_for_stream()
1259 assert save_observer(observer)
1260 prefix = observer["key"][:8]
1261 _build_legacy_prune_lattice(observer_cli_env.journal, prefix)
1262 canonical = (
1263 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "090000_300"
1264 )
1265 canonical_hash = _sha((canonical / "audio.flac").read_bytes())
1266
1267 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1268
1269 assert {candidate.analysis.segment for candidate in result.deleted} == {
1270 "090000_301",
1271 "090000_302",
1272 }
1273 assert result.crash_repaired == 0
1274 assert result.chain_repaired == 1
1275 output = format_result(result)
1276 assert "chain-repaired: 1" in output
1277 assert "crash-repaired:" not in output
1278 assert result.refusals
1279 assert (canonical / "audio.flac").is_file()
1280 assert _sha((canonical / "audio.flac").read_bytes()) == canonical_hash
1281 assert not (
1282 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "090000_301"
1283 ).exists()
1284 assert not (
1285 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "090000_302"
1286 ).exists()
1287 history = load_history(prefix, PRUNE_DAY)
1288 pruned = [row for row in history if row.get("type") == "pruned"]
1289 assert {row["segment"] for row in pruned} == {"090000_301", "090000_302"}
1290 assert all(row["duplicate_of"] == "090000_300" for row in pruned)
1291 assert (
1292 read_segment_stream(
1293 observer_cli_env.journal
1294 / "chronicle"
1295 / PRUNE_DAY
1296 / PRUNE_STREAM
1297 / "090000_303"
1298 )["prev_segment"]
1299 == "090000_300"
1300 )
1301 assert _index_count(observer_cli_env.journal, "090000_301") == 0
1302 assert _index_count(observer_cli_env.journal, "090000_302") == 0
1303 assert (
1304 observer_cli_env.journal / "chronicle" / PRUNE_DAY / "health" / "stream.updated"
1305 ).exists()
1306
1307 state = get_stream_state(PRUNE_STREAM)
1308 assert state["type"] == "observer"
1309 assert state["host"] == "field-host"
1310 assert state["platform"] == "linux"
1311 assert state["created_at"] == 123
1312 assert state["last_segment"] == "100500_300"
1313 assert state["seq"] == 9
1314
1315 existing = {
1316 segment
1317 for _stream, segment, _path in iter_segments(PRUNE_DAY)
1318 if _stream == PRUNE_STREAM
1319 }
1320 for _stream, _segment, path in iter_segments(PRUNE_DAY):
1321 marker = read_segment_stream(path)
1322 if marker and marker.get("prev_segment"):
1323 assert marker["prev_segment"] in existing
1324
1325
1326def test_prune_same_start_grouping_does_not_delete_different_start_identical_bytes(
1327 observer_cli_env,
1328) -> None:
1329 observer = _observer_for_stream()
1330 assert save_observer(observer)
1331 prefix = observer["key"][:8]
1332 _write_prune_segment(
1333 observer_cli_env.journal,
1334 segment="120000_300",
1335 seq=1,
1336 prev_segment=None,
1337 )
1338 _write_prune_segment(
1339 observer_cli_env.journal,
1340 segment="120500_300",
1341 seq=2,
1342 prev_segment="120000_300",
1343 )
1344 _append_upload(prefix, "120000_300")
1345 _append_upload(prefix, "120500_300")
1346 _write_stream_state(observer_cli_env.journal, last_segment="120500_300", seq=2)
1347
1348 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1349
1350 assert result.groups == []
1351 assert result.deleted == []
1352 assert result.refusals == []
1353 assert (
1354 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "120500_300"
1355 ).is_dir()
1356
1357
1358def test_prune_near_duplicate_media_bytes_refuse_content_identity(
1359 observer_cli_env,
1360) -> None:
1361 observer = _observer_for_stream()
1362 assert save_observer(observer)
1363 prefix = observer["key"][:8]
1364 _write_prune_segment(
1365 observer_cli_env.journal,
1366 segment="125000_300",
1367 seq=1,
1368 prev_segment=None,
1369 )
1370 _write_prune_segment(
1371 observer_cli_env.journal,
1372 segment="125000_301",
1373 seq=2,
1374 prev_segment="125000_300",
1375 )
1376 _write_prune_segment(
1377 observer_cli_env.journal,
1378 segment="125000_302",
1379 seq=3,
1380 prev_segment="125000_301",
1381 audio=b"near duplicate audio bytes",
1382 )
1383 _append_upload(prefix, "125000_300")
1384 _append_upload(prefix, "125000_301")
1385 _append_upload(prefix, "125000_302", sha=_sha(b"near duplicate audio bytes"))
1386 _write_stream_state(observer_cli_env.journal, last_segment="125000_302", seq=3)
1387
1388 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1389
1390 assert [candidate.analysis.segment for candidate in result.deleted] == [
1391 "125000_301"
1392 ]
1393 refusal = next(
1394 refusal
1395 for refusal in result.refusals
1396 if refusal.subject.endswith("/125000_302")
1397 )
1398 assert refusal.gate == "content-identity"
1399 assert refusal.file == "audio.flac"
1400 assert "compared to canonical 125000_300" in refusal.resolution
1401 assert (
1402 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "125000_302"
1403 ).is_dir()
1404
1405
1406def test_prune_markerless_candidate_refuses(
1407 observer_cli_env,
1408) -> None:
1409 observer = _observer_for_stream()
1410 assert save_observer(observer)
1411 prefix = observer["key"][:8]
1412 _write_prune_segment(
1413 observer_cli_env.journal,
1414 segment="130000_300",
1415 seq=1,
1416 prev_segment=None,
1417 )
1418 _write_prune_segment(
1419 observer_cli_env.journal,
1420 segment="130000_301",
1421 seq=2,
1422 prev_segment="130000_300",
1423 marker=False,
1424 )
1425 _append_upload(prefix, "130000_300")
1426 _append_upload(prefix, "130000_301")
1427 _write_stream_state(observer_cli_env.journal, last_segment="130000_301", seq=2)
1428
1429 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1430
1431 assert result.deleted == []
1432 assert any(refusal.gate == "chain-identity" for refusal in result.refusals)
1433 assert (
1434 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "130000_301"
1435 ).is_dir()
1436
1437
1438def test_prune_unrecognized_derived_file_refuses_legacy_and_manifest_candidates(
1439 observer_cli_env,
1440) -> None:
1441 observer = _observer_for_stream()
1442 assert save_observer(observer)
1443 prefix = observer["key"][:8]
1444 _write_prune_segment(
1445 observer_cli_env.journal,
1446 segment="132000_300",
1447 seq=1,
1448 prev_segment=None,
1449 )
1450 _write_prune_segment(
1451 observer_cli_env.journal,
1452 segment="132000_301",
1453 seq=2,
1454 prev_segment="132000_300",
1455 unknown_file=True,
1456 )
1457 _write_prune_segment(
1458 observer_cli_env.journal,
1459 segment="132500_300",
1460 seq=3,
1461 prev_segment="132000_301",
1462 manifest=True,
1463 )
1464 _write_prune_segment(
1465 observer_cli_env.journal,
1466 segment="132500_301",
1467 seq=4,
1468 prev_segment="132500_300",
1469 manifest=True,
1470 unknown_file=True,
1471 )
1472 for segment in ("132000_300", "132000_301", "132500_300", "132500_301"):
1473 _append_upload(prefix, segment)
1474 _write_stream_state(observer_cli_env.journal, last_segment="132500_301", seq=4)
1475
1476 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1477
1478 refusals = {
1479 refusal.subject.rsplit("/", 1)[1]: refusal
1480 for refusal in result.refusals
1481 if refusal.gate == "derived-output"
1482 }
1483 assert set(refusals) == {"132000_301", "132500_301"}
1484 for refusal in refusals.values():
1485 assert refusal.file == "notes.txt"
1486 assert "remove the file" in refusal.resolution
1487 assert result.deleted == []
1488 assert (
1489 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "132000_301"
1490 ).is_dir()
1491 assert (
1492 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "132500_301"
1493 ).is_dir()
1494
1495
1496def test_prune_refuses_manifest_content_name_outside_segment(
1497 observer_cli_env,
1498) -> None:
1499 observer = _observer_for_stream()
1500 assert save_observer(observer)
1501 prefix = observer["key"][:8]
1502 _write_prune_segment(
1503 observer_cli_env.journal,
1504 segment="133000_300",
1505 seq=1,
1506 prev_segment=None,
1507 manifest=True,
1508 )
1509 candidate = _write_prune_segment(
1510 observer_cli_env.journal,
1511 segment="133000_301",
1512 seq=2,
1513 prev_segment="133000_300",
1514 manifest=True,
1515 )
1516 (candidate.parent / "outside.flac").write_bytes(PRUNE_AUDIO)
1517 manifest_path = candidate / "ingest.json"
1518 manifest = json.loads(manifest_path.read_text(encoding="utf-8"))
1519 manifest["files"]["../outside.flac"] = {
1520 "sha256": _sha(PRUNE_AUDIO),
1521 "size": len(PRUNE_AUDIO),
1522 }
1523 manifest_path.write_text(json.dumps(manifest) + "\n", encoding="utf-8")
1524 _append_upload(prefix, "133000_300")
1525 _append_upload(prefix, "133000_301")
1526 _write_stream_state(observer_cli_env.journal, last_segment="133000_301", seq=2)
1527
1528 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1529
1530 refusal = next(
1531 refusal
1532 for refusal in result.refusals
1533 if refusal.subject.endswith("/133000_301")
1534 )
1535 assert refusal.gate == "canonical-heldness"
1536 assert refusal.file == "../outside.flac"
1537 assert "plain in-segment names" in refusal.resolution
1538 assert candidate.is_dir()
1539
1540
1541def test_prune_manifest_extra_content_refuses_as_content_identity(
1542 observer_cli_env,
1543) -> None:
1544 observer = _observer_for_stream()
1545 assert save_observer(observer)
1546 prefix = observer["key"][:8]
1547 _write_prune_segment(
1548 observer_cli_env.journal,
1549 segment="135000_300",
1550 seq=1,
1551 prev_segment=None,
1552 manifest=True,
1553 )
1554 _write_prune_segment(
1555 observer_cli_env.journal,
1556 segment="135000_301",
1557 seq=2,
1558 prev_segment="135000_300",
1559 manifest=True,
1560 )
1561 _write_prune_segment(
1562 observer_cli_env.journal,
1563 segment="135000_302",
1564 seq=3,
1565 prev_segment="135000_301",
1566 manifest=True,
1567 extra_manifest_content=True,
1568 )
1569 _append_upload(prefix, "135000_300")
1570 _append_upload(prefix, "135000_301")
1571 _append_upload(prefix, "135000_302")
1572 _write_stream_state(observer_cli_env.journal, last_segment="135000_302", seq=3)
1573
1574 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1575
1576 assert [candidate.analysis.segment for candidate in result.deleted] == [
1577 "135000_301"
1578 ]
1579 refusal = next(
1580 refusal
1581 for refusal in result.refusals
1582 if refusal.subject.endswith("/135000_302")
1583 )
1584 assert refusal.gate == "content-identity"
1585 assert refusal.file == "capture.bin"
1586 assert (
1587 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "135000_302"
1588 ).is_dir()
1589
1590
1591def test_prune_refuses_when_canonical_heldness_becomes_unverifiable(
1592 observer_cli_env,
1593) -> None:
1594 observer = _observer_for_stream()
1595 assert save_observer(observer)
1596 prefix = observer["key"][:8]
1597 _write_prune_segment(
1598 observer_cli_env.journal,
1599 segment="140000_300",
1600 seq=1,
1601 prev_segment=None,
1602 manifest=True,
1603 )
1604 _write_prune_segment(
1605 observer_cli_env.journal,
1606 segment="140000_301",
1607 seq=2,
1608 prev_segment="140000_300",
1609 manifest=True,
1610 )
1611 _append_upload(prefix, "140000_300")
1612 _append_upload(prefix, "140000_301")
1613 _write_stream_state(observer_cli_env.journal, last_segment="140000_301", seq=2)
1614 dry = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=False)
1615 assert len(dry.groups) == 1
1616 assert dry.refusals == []
1617 assert [candidate.analysis.segment for candidate in dry.groups[0].candidates] == [
1618 "140000_301"
1619 ]
1620 canonical_audio = (
1621 observer_cli_env.journal
1622 / "chronicle"
1623 / PRUNE_DAY
1624 / PRUNE_STREAM
1625 / "140000_300"
1626 / "audio.flac"
1627 )
1628 canonical_audio.unlink()
1629
1630 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1631
1632 assert result.deleted == []
1633 assert result_exit_code(result) == 2
1634 assert any(refusal.gate == "canonical-heldness" for refusal in result.refusals)
1635 assert (
1636 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "140000_301"
1637 ).is_dir()
1638
1639
1640def test_prune_last_physical_copy_is_labeled_in_dry_run_and_execute(
1641 observer_cli_env,
1642) -> None:
1643 observer = _observer_for_stream()
1644 assert save_observer(observer)
1645 prefix = observer["key"][:8]
1646 _write_prune_segment(
1647 observer_cli_env.journal,
1648 segment="150000_300",
1649 seq=1,
1650 prev_segment=None,
1651 manifest=True,
1652 proof_only_audio=True,
1653 )
1654 _write_prune_segment(
1655 observer_cli_env.journal,
1656 segment="150000_301",
1657 seq=2,
1658 prev_segment="150000_300",
1659 manifest=True,
1660 )
1661 _append_upload(prefix, "150000_300")
1662 _append_upload(prefix, "150000_301")
1663 _write_stream_state(observer_cli_env.journal, last_segment="150000_301", seq=2)
1664
1665 dry = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=False)
1666 assert dry.last_physical_copy_count == 1
1667 assert dry.groups[0].candidates[0].last_physical_copy is True
1668
1669 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1670
1671 assert result.last_physical_copy_count == 1
1672 assert result.deleted[0].last_physical_copy is True
1673
1674
1675def test_prune_crash_after_pruned_record_before_delete_converges(
1676 observer_cli_env,
1677 monkeypatch: pytest.MonkeyPatch,
1678) -> None:
1679 observer = _observer_for_stream()
1680 assert save_observer(observer)
1681 prefix = observer["key"][:8]
1682 _write_prune_segment(
1683 observer_cli_env.journal,
1684 segment="155000_300",
1685 seq=1,
1686 prev_segment=None,
1687 )
1688 candidate = _write_prune_segment(
1689 observer_cli_env.journal,
1690 segment="155000_301",
1691 seq=2,
1692 prev_segment="155000_300",
1693 )
1694 _append_upload(prefix, "155000_300")
1695 _append_upload(prefix, "155000_301")
1696 _write_stream_state(observer_cli_env.journal, last_segment="155000_301", seq=2)
1697
1698 from solstone.apps.observer import prune
1699
1700 def crash_before_delete(_path: Path) -> None:
1701 raise RuntimeError("interrupt before delete")
1702
1703 with monkeypatch.context() as patch:
1704 patch.setattr(prune.shutil, "rmtree", crash_before_delete)
1705 with pytest.raises(RuntimeError, match="interrupt before delete"):
1706 run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1707
1708 assert candidate.is_dir()
1709 assert (candidate / "audio.flac").read_bytes() == PRUNE_AUDIO
1710 assert len(_pruned_history(prefix, "155000_301")) == 1
1711
1712 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1713
1714 assert result.refusals == []
1715 assert result_exit_code(result) == 0
1716 assert not candidate.exists()
1717 assert len(_pruned_history(prefix, "155000_301")) == 1
1718
1719
1720def test_prune_crash_rerun_repairs_pruned_dangling_prev_and_refuses_unexplained(
1721 observer_cli_env,
1722 monkeypatch: pytest.MonkeyPatch,
1723) -> None:
1724 observer = _observer_for_stream()
1725 assert save_observer(observer)
1726 prefix = observer["key"][:8]
1727 _write_prune_segment(
1728 observer_cli_env.journal,
1729 segment="160000_300",
1730 seq=1,
1731 prev_segment=None,
1732 )
1733 _write_prune_segment(
1734 observer_cli_env.journal,
1735 segment="160000_301",
1736 seq=2,
1737 prev_segment="160000_300",
1738 )
1739 _write_prune_segment(
1740 observer_cli_env.journal,
1741 segment="160500_300",
1742 seq=3,
1743 prev_segment="160000_301",
1744 audio=b"after",
1745 )
1746 _append_upload(prefix, "160000_300")
1747 _append_upload(prefix, "160000_301")
1748 _write_stream_state(observer_cli_env.journal, last_segment="160500_300", seq=3)
1749
1750 from solstone.apps.observer import prune
1751
1752 def fail_repair(*_args, **_kwargs):
1753 raise RuntimeError("interrupt after delete")
1754
1755 with monkeypatch.context() as patch:
1756 patch.setattr(
1757 prune, "repair_crash_leftovers", lambda *_args, **_kwargs: ([], 0)
1758 )
1759 patch.setattr(prune, "repair_stream_chain", fail_repair)
1760 with pytest.raises(RuntimeError, match="interrupt after delete"):
1761 run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1762 assert not (
1763 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "160000_301"
1764 ).exists()
1765 assert any(
1766 row.get("type") == "pruned" and row.get("segment") == "160000_301"
1767 for row in load_history(prefix, PRUNE_DAY)
1768 )
1769
1770 result = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1771 assert result.refusals == []
1772 assert result.crash_repaired == 1
1773 assert result.chain_repaired == 0
1774 output = format_result(result)
1775 assert "chain-repaired: 0" in output
1776 assert "crash-repaired: 1" in output
1777 assert (
1778 read_segment_stream(
1779 observer_cli_env.journal
1780 / "chronicle"
1781 / PRUNE_DAY
1782 / PRUNE_STREAM
1783 / "160500_300"
1784 )["prev_segment"]
1785 == "160000_300"
1786 )
1787
1788 unexplained = (
1789 observer_cli_env.journal / "chronicle" / PRUNE_DAY / PRUNE_STREAM / "160000_300"
1790 )
1791 __import__("shutil").rmtree(unexplained)
1792 bad = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1793 assert any(refusal.gate == "chain-repair" for refusal in bad.refusals)
1794
1795
1796def test_prune_execute_twice_is_clean_noop_second_run(
1797 observer_cli_env,
1798) -> None:
1799 observer = _observer_for_stream()
1800 assert save_observer(observer)
1801 prefix = observer["key"][:8]
1802 _write_prune_segment(
1803 observer_cli_env.journal,
1804 segment="165000_300",
1805 seq=1,
1806 prev_segment=None,
1807 )
1808 _write_prune_segment(
1809 observer_cli_env.journal,
1810 segment="165000_301",
1811 seq=2,
1812 prev_segment="165000_300",
1813 )
1814 _write_prune_segment(
1815 observer_cli_env.journal,
1816 segment="165500_300",
1817 seq=3,
1818 prev_segment="165000_301",
1819 audio=b"after",
1820 )
1821 _append_upload(prefix, "165000_300")
1822 _append_upload(prefix, "165000_301")
1823 _write_stream_state(observer_cli_env.journal, last_segment="165500_300", seq=3)
1824
1825 first = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1826 assert result_exit_code(first) == 0
1827 assert [candidate.analysis.segment for candidate in first.deleted] == ["165000_301"]
1828 state_after_first = get_stream_state(PRUNE_STREAM)
1829 markers_after_first = _stream_marker_snapshot(observer_cli_env.journal)
1830 pruned_after_first = _pruned_history(prefix, "165000_301")
1831
1832 second = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1833
1834 assert result_exit_code(second) == 0
1835 assert second.refusals == []
1836 assert second.deleted == []
1837 assert _pruned_history(prefix, "165000_301") == pruned_after_first
1838 assert _stream_marker_snapshot(observer_cli_env.journal) == markers_after_first
1839 assert get_stream_state(PRUNE_STREAM) == state_after_first
1840
1841
1842def test_prune_observer_attribution_refuses_no_owner_and_ambiguous_owner(
1843 observer_cli_env,
1844) -> None:
1845 _write_prune_segment(
1846 observer_cli_env.journal,
1847 segment="170000_300",
1848 seq=1,
1849 prev_segment=None,
1850 )
1851 _write_prune_segment(
1852 observer_cli_env.journal,
1853 segment="170000_301",
1854 seq=2,
1855 prev_segment="170000_300",
1856 )
1857 _write_stream_state(observer_cli_env.journal, last_segment="170000_301", seq=2)
1858
1859 no_owner = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1860 assert any(refusal.gate == "observer-attribution" for refusal in no_owner.refusals)
1861
1862 first = _observer_for_stream(key="field-one-key")
1863 second = _observer_for_stream(key="field-two-key")
1864 assert save_observer(first)
1865 assert save_observer(second)
1866 ambiguous = run_prune(days=[PRUNE_DAY], stream=PRUNE_STREAM, execute=True)
1867 assert any(
1868 "multiple active observers" in refusal.resolution
1869 for refusal in ambiguous.refusals
1870 )