personal memory agent
0

Configure Feed

Select the types of activity you want to include in your feed.

solstone / tests / test_observer_cli.py
58 kB 1870 lines
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 )