from __future__ import annotations import hashlib import json import os import subprocess import sys import time from pathlib import Path import pytest import k1link.web.camera_archive as camera_archive_module from k1link.device_plugins.xgrids_k1.archive import discover_legacy_viewer_sessions from k1link.web.camera_archive import ( CAMERA_ARCHIVE_SCHEMA, CAMERA_COMMIT_POLICY, CameraArchiveError, CameraArchiveWriter, recover_incomplete_camera_archives, ) def test_camera_archive_writes_canonical_segments_and_seals_summary(tmp_path: Path) -> None: session = tmp_path / "session" session.mkdir() writer = CameraArchiveWriter( session, "sensor.camera.left", 7, commit_interval_seconds=60, commit_bytes=1024, ) init = b"init-segment" first = b"first-fragment" second = b"second-fragment" writer.append("init", init, host_epoch_ns=100, host_monotonic_ns=200) writer.append("media", first, host_epoch_ns=110, host_monotonic_ns=210) writer.append("media", second, host_epoch_ns=120, host_monotonic_ns=220) # The tuning values cannot weaken the segment RPO: every append has already # published its fragment, index record, and interrupted crash checkpoint. checkpoint = json.loads(writer.summary_path.read_text(encoding="utf-8")) assert checkpoint["status"] == "interrupted" assert checkpoint["segment_count"] == 2 assert checkpoint["commit_policy"] == CAMERA_COMMIT_POLICY summary = writer.close() entries = [json.loads(line) for line in writer.index_path.read_text().splitlines()] assert writer.archive_dir.name == "epoch-7" assert writer.init_path.read_bytes() == init assert [path.name for path in sorted(writer.segments_dir.iterdir())] == [ "1.m4s", "2.m4s", ] assert [entry["kind"] for entry in entries] == ["media", "media"] assert [entry["sequence"] for entry in entries] == [1, 2] assert [entry["path"] for entry in entries] == [ "segments/1.m4s", "segments/2.m4s", ] for entry, expected in zip(entries, (first, second), strict=True): payload = (writer.archive_dir / entry["path"]).read_bytes() assert payload == expected assert entry["length"] == len(expected) assert entry["sha256"] == hashlib.sha256(expected).hexdigest() assert summary["schema_version"] == CAMERA_ARCHIVE_SCHEMA assert summary["status"] == "complete" assert summary["segment_count"] == 2 assert summary["entry_count"] == 2 assert summary["media_segment_count"] == 2 assert summary["valid_bytes"] == len(init) + len(first) + len(second) assert summary["stream_sha256"] == hashlib.sha256(init + first + second).hexdigest() assert summary["synchronization"] == "host-arrival-best-effort" assert summary["artifacts"] == { "init": "init.mp4", "segments": "segments", "index": "index.jsonl", } assert "192.168" not in json.dumps(summary) assert writer.close() == summary def test_recovery_accepts_sealed_index_larger_than_old_duration_cap_without_loading_video( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ) -> None: sessions_root = tmp_path / "sessions" epoch = ( sessions_root / "20260720T065719Z_viewer_live" / "media" / "sensor.camera.right" / "epoch-1" ) segments = epoch / "segments" segments.mkdir(parents=True) init = b"sealed-init" (epoch / "init.mp4").write_bytes(init) segment_count = 560 segment = b"x" segment_sha256 = hashlib.sha256(segment).hexdigest() index_sha256 = hashlib.sha256() with (epoch / "index.jsonl").open("wb") as stream: for sequence in range(1, segment_count + 1): (segments / f"{sequence}.m4s").write_bytes(segment) line = ( json.dumps( { "schema_version": "missioncore.camera-recording-index/v1", "sequence": sequence, "kind": "media", "path": f"segments/{sequence}.m4s", "length": 1, "sha256": segment_sha256, "padding": "p" * 60_000, }, separators=(",", ":"), ) + "\n" ).encode() index_sha256.update(line) stream.write(line) assert (epoch / "index.jsonl").stat().st_size > 32 * 1024 * 1024 (epoch / "summary.json").write_text( json.dumps( { "schema_version": CAMERA_ARCHIVE_SCHEMA, "source_id": "sensor.camera.right", "codec_epoch": 1, "status": "complete", "segment_count": segment_count, "entry_count": segment_count, "media_segment_count": segment_count, "valid_bytes": len(init) + segment_count, "init_sha256": hashlib.sha256(init).hexdigest(), "stream_sha256": hashlib.sha256(init + segment * segment_count).hexdigest(), "index_sha256": index_sha256.hexdigest(), "synchronization": "host-arrival-best-effort", "commit_policy": CAMERA_COMMIT_POLICY, "artifacts": { "init": "init.mp4", "segments": "segments", "index": "index.jsonl", }, } ), encoding="utf-8", ) def fail_if_recovery_reads_video(_segments_fd: int) -> object: raise AssertionError("a clean sealed archive must not enter payload recovery") monkeypatch.setattr( camera_archive_module, "_read_recovery_segment_catalog", fail_if_recovery_reads_video, ) assert recover_incomplete_camera_archives(sessions_root) == () def test_camera_archive_rejects_media_without_init_and_existing_epoch(tmp_path: Path) -> None: session = tmp_path / "session" session.mkdir() writer = CameraArchiveWriter(session, "sensor.camera.right", 1) with pytest.raises(CameraArchiveError, match="before"): writer.append("media", b"fragment") writer.append("init", b"init") writer.close(status="interrupted", failure_code="source-ended") with pytest.raises(FileExistsError): CameraArchiveWriter(session, "sensor.camera.right", 1) with pytest.raises(ValueError, match="safe storage"): CameraArchiveWriter(session, "../camera", 2) @pytest.mark.parametrize("symlink_component", ["media", "source"]) def test_camera_writer_never_follows_precreated_storage_symlinks( tmp_path: Path, symlink_component: str, ) -> None: session = tmp_path / "session" session.mkdir() outside = tmp_path / "outside" outside.mkdir() if symlink_component == "media": (session / "media").symlink_to(outside, target_is_directory=True) else: media = session / "media" media.mkdir() (media / "sensor.camera.left").symlink_to(outside, target_is_directory=True) with pytest.raises(CameraArchiveError, match="no-follow"): CameraArchiveWriter(session, "sensor.camera.left", 1) assert list(outside.iterdir()) == [] def test_camera_writer_keeps_directory_fds_across_path_swap(tmp_path: Path) -> None: session = tmp_path / "session" session.mkdir() writer = CameraArchiveWriter(session, "sensor.camera.left", 1) writer.append("init", b"init") original_epoch = writer.archive_dir moved_epoch = original_epoch.with_name("epoch-1-moved") original_epoch.rename(moved_epoch) outside_epoch = tmp_path / "outside-epoch" outside_epoch.mkdir() original_epoch.symlink_to(outside_epoch, target_is_directory=True) original_segments = moved_epoch / "segments" moved_segments = moved_epoch / "segments-held-by-fd" original_segments.rename(moved_segments) outside_segments = tmp_path / "outside-segments" outside_segments.mkdir() original_segments.symlink_to(outside_segments, target_is_directory=True) writer.append("media", b"frame-after-path-swap") summary = writer.close() assert (moved_segments / "1.m4s").read_bytes() == b"frame-after-path-swap" assert json.loads((moved_epoch / "summary.json").read_text(encoding="utf-8")) == summary assert list(outside_epoch.iterdir()) == [] assert list(outside_segments.iterdir()) == [] def test_recovery_seals_unindexed_durable_tail_and_preserves_orphans(tmp_path: Path) -> None: sessions_root = tmp_path / "sessions" session = sessions_root / "20260717T011050Z_viewer_live" session.mkdir(parents=True) writer = CameraArchiveWriter(session, "sensor.camera.left", 3) writer.append("init", b"init") writer.append("media", b"frame-1", host_epoch_ns=100, host_monotonic_ns=200) # Catalog refresh must not hash or rewrite a writer that is still active in # this server process. assert recover_incomplete_camera_archives(sessions_root) == () writer.close(status="interrupted", failure_code="synthetic-process-crash") epoch = writer.archive_dir writer.summary_path.unlink() # Simulate a process dying after atomically publishing the next fragment but # before its JSONL entry, plus a non-contiguous fragment that cannot be put on # the trusted timeline. Recovery salvages #2 and quarantines (never deletes) #4. (writer.segments_dir / "2.m4s").write_bytes(b"frame-2") (writer.segments_dir / "4.m4s").write_bytes(b"orphan-frame") recovered = recover_incomplete_camera_archives(sessions_root) assert len(recovered) == 1 summary = recovered[0] assert summary["status"] == "interrupted" assert summary["failure_code"] == "server-process-interrupted" assert summary["segment_count"] == 2 entries = [ json.loads(line) for line in writer.index_path.read_text(encoding="utf-8").splitlines() ] assert [entry["sequence"] for entry in entries] == [1, 2] assert entries[0]["host_epoch_ns"] == 100 assert entries[1]["recovered"] is True assert (writer.segments_dir / "2.m4s").read_bytes() == b"frame-2" assert not (writer.segments_dir / "4.m4s").exists() assert any( path.read_bytes() == b"orphan-frame" for path in (epoch / "recovery-orphans").iterdir() if path.is_file() ) # Recovery output is exactly the layout consumed by legacy media discovery, # and a second scan is an idempotent no-op. assert recover_incomplete_camera_archives(sessions_root) == () candidates = discover_legacy_viewer_sessions(sessions_root) assert len(candidates) == 1 assert [source.source_id for source in candidates[0].media_sources] == [ "sensor.camera.left" ] def test_recovery_replaces_index_and_summary_symlinks_without_following_targets( tmp_path: Path, ) -> None: sessions_root = tmp_path / "sessions" session = sessions_root / "20260717T011051Z_viewer_live" session.mkdir(parents=True) writer = CameraArchiveWriter(session, "sensor.camera.right", 1) writer.append("init", b"init") writer.append("media", b"frame") writer.close(status="interrupted") outside_index = tmp_path / "outside-index" outside_summary = tmp_path / "outside-summary" outside_index.write_bytes(b"do-not-read-or-overwrite-index") outside_summary.write_bytes(b"do-not-read-or-overwrite-summary") writer.index_path.unlink() writer.summary_path.unlink() writer.index_path.symlink_to(outside_index) writer.summary_path.symlink_to(outside_summary) recovered = recover_incomplete_camera_archives(sessions_root) assert len(recovered) == 1 assert recovered[0]["segment_count"] == 1 assert outside_index.read_bytes() == b"do-not-read-or-overwrite-index" assert outside_summary.read_bytes() == b"do-not-read-or-overwrite-summary" assert writer.index_path.is_symlink() is False assert writer.summary_path.is_symlink() is False assert json.loads(writer.index_path.read_text(encoding="utf-8"))["recovered"] is True def test_recovery_fails_closed_on_symlinked_quarantine_directory(tmp_path: Path) -> None: sessions_root = tmp_path / "sessions" session = sessions_root / "20260717T011052Z_viewer_live" session.mkdir(parents=True) writer = CameraArchiveWriter(session, "sensor.camera.left", 1) writer.append("init", b"init") writer.append("media", b"frame-1") writer.close(status="interrupted") writer.summary_path.unlink() orphan = writer.segments_dir / "3.m4s" orphan.write_bytes(b"orphan") outside = tmp_path / "outside-quarantine" outside.mkdir() (writer.archive_dir / "recovery-orphans").symlink_to(outside, target_is_directory=True) with pytest.raises(CameraArchiveError, match="real directory"): recover_incomplete_camera_archives(sessions_root) assert list(outside.iterdir()) == [] assert orphan.read_bytes() == b"orphan" def test_recovery_quarantines_segment_symlink_without_reading_target(tmp_path: Path) -> None: sessions_root = tmp_path / "sessions" session = sessions_root / "20260717T011053Z_viewer_live" session.mkdir(parents=True) writer = CameraArchiveWriter(session, "sensor.camera.left", 1) writer.append("init", b"init") writer.append("media", b"frame-1") writer.close(status="interrupted") writer.summary_path.unlink() outside = tmp_path / "outside-segment" outside.write_bytes(b"external-evidence-must-not-be-read-or-moved") (writer.segments_dir / "2.m4s").symlink_to(outside) recovered = recover_incomplete_camera_archives(sessions_root) assert len(recovered) == 1 assert recovered[0]["segment_count"] == 1 assert outside.read_bytes() == b"external-evidence-must-not-be-read-or-moved" quarantined = list((writer.archive_dir / "recovery-orphans").glob("2.m4s*")) assert len(quarantined) == 1 assert quarantined[0].is_symlink() assert quarantined[0].resolve() == outside.resolve() def test_recovery_lock_symlink_is_rejected_without_touching_target(tmp_path: Path) -> None: sessions_root = tmp_path / "sessions" sessions_root.mkdir() outside = tmp_path / "outside-lock" outside.write_bytes(b"external-lock-target") (sessions_root / ".camera-recovery.lock").symlink_to(outside) with pytest.raises(CameraArchiveError, match="lock failed no-follow"): recover_incomplete_camera_archives(sessions_root) assert outside.read_bytes() == b"external-lock-target" @pytest.mark.skipif(os.name != "posix", reason="uses POSIX flock to assert worker serialization") def test_recovery_startup_lease_serializes_another_process(tmp_path: Path) -> None: import fcntl sessions_root = tmp_path / "sessions" sessions_root.mkdir() lock_path = sessions_root / ".camera-recovery.lock" lock_path.touch(mode=0o600) descriptor = os.open(lock_path, os.O_RDWR) process: subprocess.Popen[str] | None = None locked = False try: fcntl.flock(descriptor, fcntl.LOCK_EX) locked = True process = subprocess.Popen( [ sys.executable, "-c", ( "from pathlib import Path; " "from k1link.web.camera_archive import " "recover_incomplete_camera_archives; " "recover_incomplete_camera_archives(Path(__import__('sys').argv[1])); " "print('recovered')" ), str(sessions_root), ], cwd=Path(__file__).parents[1], text=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE, ) time.sleep(0.15) assert process.poll() is None fcntl.flock(descriptor, fcntl.LOCK_UN) locked = False stdout, stderr = process.communicate(timeout=5) assert process.returncode == 0, stderr assert stdout.strip() == "recovered" finally: if locked: fcntl.flock(descriptor, fcntl.LOCK_UN) os.close(descriptor) if process is not None and process.poll() is None: process.kill() process.wait(timeout=5)