"""Normal-collection retention acceptance tests for issue #100.""" import sqlite3 import sys from datetime import datetime, timedelta, timezone from pathlib import Path import pytest sys.path.insert(0, str(Path(__file__).parent.parent / "src")) from fenris.collector import run_collection from fenris.local_day import query_local_day_summary from fenris.pruning import prune_old_samples from fenris.status import read_status from fenris.store import init_store @pytest.fixture def sysfs_controller(tmp_path): controller = tmp_path / "sys" / "class" / "nvme" / "nvme0" controller.mkdir(parents=True) (controller / "subsysnqn").write_text("nqn.example:drive-1\n") (controller / "model").write_text("Fenris Test Drive\n") (controller / "serial").write_text("drive-1\n") (controller / "firmware_rev").write_text("1.0\n") transport = controller / "transport" transport.mkdir() (transport / "address").write_text("0000:00:01.0\n") (transport / "trstring").write_text("pcie\n") return controller class FakeClock: def __init__(self, now): self.now = now def utcnow(self): return self.now def _smartctl(written, read): return { "nvme_smart_health_information_log": { "critical_warning": 0, "temperature": 35, "available_spare": 100, "percentage_used": 1, "data_units_written": written, "data_units_read": read, "power_on_hours": 100, "power_cycles": 1, "unsafe_shutdowns": 0, "media_errors": 0, }, "user_capacity": {"bytes": 1_000_000_000_000}, "model_name": "Fenris Test Drive", "serial_number": "drive-1", "firmware_version": "1.0", } def _collect(config, controller, now, written, read): return run_collection( smartctl_data=_smartctl(written, read), sysfs_path=controller, config=config, clock=FakeClock(now), ) def test_normal_collection_expires_detail_and_keeps_local_day_evidence( tmp_path, sysfs_controller, monkeypatch, ): monkeypatch.setenv("TZ", "UTC") store_path = tmp_path / "observations.db" config = {"device": "/dev/nvme0", "store_path": str(store_path)} first_time = datetime(2026, 9, 1, 0, 0, tzinfo=timezone.utc) second_time = first_time + timedelta(days=15) first = _collect(config, sysfs_controller, first_time, 1000, 2000) assert first["ok"] is True, first second = _collect(config, sysfs_controller, second_time, 1001, 2003) assert second["ok"] is True, second with sqlite3.connect(store_path) as conn: assert conn.execute("SELECT ts FROM samples ORDER BY id").fetchall() == [ (second_time.isoformat(),) ] assert conn.execute( "SELECT bytes_written, bytes_read, reason " "FROM local_day_unallocated_evidence" ).fetchall() == [(512_000, 1_536_000, "local_midnight")] assert conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0 assert conn.execute("SELECT COUNT(*) FROM controller_segments").fetchone()[0] == 1 assert conn.execute("SELECT COUNT(*) FROM day_aggregates").fetchone()[0] == 15 monkeypatch.setenv("TZ", "Asia/Kolkata") with read_status(store_path, second_time, query_services=False) as (reader, state): assert reader is not None assert state.store_fault is None old_day = query_local_day_summary( reader, "2026-09-01", "UTC", second_time, ) assert old_day["utc_start"] == "2026-09-01T00:00:00+00:00" assert old_day["utc_end"] == "2026-09-02T00:00:00+00:00" assert old_day["bytes_written"] is None assert old_day["shared_bytes_written"] == 512_000 assert old_day["shared_bytes_read"] == 1_536_000 assert old_day["shared_evidence_count"] == 1 def test_pruning_keeps_samples_when_replacement_evidence_is_missing(tmp_path): store_path = tmp_path / "incomplete.db" conn = init_store(store_path) now = datetime(2026, 9, 30, tzinfo=timezone.utc) old_time = now - timedelta(days=20) next_time = now - timedelta(days=1) conn.execute( "INSERT INTO controller_segments (id, opened_at, identity_key) " "VALUES (1, ?, 'test')", (old_time.isoformat(),), ) conn.execute( "INSERT INTO samples (ts, device, bytes_written, bytes_read, segment_id, local_tz) " "VALUES (?, '/dev/nvme0', 100, 100, 1, 'UTC')", (old_time.isoformat(),), ) conn.execute( "INSERT INTO samples (ts, device, bytes_written, bytes_read, segment_id, local_tz) " "VALUES (?, '/dev/nvme0', 200, 200, 1, 'UTC')", (next_time.isoformat(),), ) conn.commit() assert prune_old_samples(conn, now) == 0 assert conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0] == 2 conn.close() def test_later_same_day_activity_does_not_replace_gap_interval(tmp_path): store_path = tmp_path / "unrelated-local-activity.db" conn = init_store(store_path) gap_start = datetime(2026, 9, 1, 0, 0, tzinfo=timezone.utc) gap_end = gap_start + timedelta(minutes=10) later_end = gap_start + timedelta(minutes=20) now = datetime(2026, 9, 30, tzinfo=timezone.utc) conn.execute( "INSERT INTO controller_segments (id, opened_at, identity_key) " "VALUES (1, ?, 'test')", (gap_start.isoformat(),), ) for sample_time, written, read in ( (gap_start, 1000, 100), (gap_end, 2000, 200), (later_end, 2100, 210), ): conn.execute( "INSERT INTO samples (ts, device, bytes_written, bytes_read, segment_id, local_tz) " "VALUES (?, '/dev/nvme0', ?, ?, 1, 'UTC')", (sample_time.isoformat(), written, read), ) # UTC aggregates can replace the old counter pair, but the local-day # summary contains only the later B->C interval, after monitoring resumed. conn.execute( "INSERT INTO hour_observations (hour, bytes_written_delta, bytes_read_delta) " "VALUES ('2026-09-01T00:00:00+00:00', 1000, 100)" ) conn.execute( "INSERT INTO day_aggregates (day) VALUES ('2026-09-01')" ) conn.execute( "INSERT INTO monitoring_periods (started_at, ended_at) VALUES (?, ?)", (gap_end.isoformat(), later_end.isoformat()), ) last_sample_id = conn.execute("SELECT MAX(id) FROM samples").fetchone()[0] conn.execute( "INSERT INTO local_days " "(local_date, tz_name, tz_offset, utc_start, utc_end, bytes_written, " " bytes_read, activity_intervals, activity_precision, last_sample_id) " "VALUES ('2026-09-01', 'UTC', '+00:00', ?, ?, 100, 10, 1, 'measured', ?)", ( gap_start.isoformat(), (gap_start + timedelta(days=1)).isoformat(), last_sample_id, ), ) local_day_id = conn.execute("SELECT MAX(id) FROM local_days").fetchone()[0] conn.execute( "INSERT INTO local_day_segment_totals " "(local_day_id, segment_id, bytes_written, bytes_read, activity_intervals) " "VALUES (?, 1, 100, 10, 1)", (local_day_id,), ) conn.commit() assert prune_old_samples(conn, now) == 0 assert conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0] == 3 conn.close() def test_pruning_keeps_oldest_unpaired_sample_as_successor_anchor(tmp_path): store_path = tmp_path / "anchor.db" conn = init_store(store_path) now = datetime(2026, 9, 30, tzinfo=timezone.utc) old_time = now - timedelta(days=20) conn.execute( "INSERT INTO controller_segments (id, opened_at, identity_key) " "VALUES (1, ?, 'test')", (old_time.isoformat(),), ) conn.execute( "INSERT INTO samples (ts, device, bytes_written, bytes_read, segment_id, local_tz) " "VALUES (?, '/dev/nvme0', 100, 100, 1, 'UTC')", (old_time.isoformat(),), ) conn.commit() assert prune_old_samples(conn, now) == 0 assert conn.execute("SELECT ts FROM samples").fetchone()[0] == old_time.isoformat() conn.close() def test_pruning_never_expires_pending_publications(tmp_path): store_path = tmp_path / "pending.db" conn = init_store(store_path) old_time = datetime(2026, 9, 1, tzinfo=timezone.utc) conn.execute( "INSERT INTO pending_publications (sample_ts, payload) VALUES (?, '{}')", (old_time.isoformat(),), ) conn.commit() assert prune_old_samples( conn, datetime(2026, 9, 30, tzinfo=timezone.utc), ) == 0 assert conn.execute( "SELECT sample_ts FROM pending_publications" ).fetchone()[0] == old_time.isoformat() conn.close() def test_pruning_preserves_unparseable_legacy_timestamp(tmp_path): store_path = tmp_path / "legacy.db" conn = init_store(store_path) conn.execute( "INSERT INTO samples (ts, device, bytes_written, bytes_read) " "VALUES ('legacy timestamp', '/dev/nvme0', 100, 100)" ) conn.commit() assert prune_old_samples(conn, datetime(2026, 9, 30, tzinfo=timezone.utc)) == 0 assert conn.execute("SELECT ts FROM samples").fetchone()[0] == "legacy timestamp" conn.close() def test_collection_defers_interrupted_pruning_and_retries_next_run( tmp_path, sysfs_controller, monkeypatch, ): monkeypatch.setenv("TZ", "UTC") store_path = tmp_path / "interrupted.db" config = {"device": "/dev/nvme0", "store_path": str(store_path)} first_time = datetime(2026, 9, 1, 0, 0, tzinfo=timezone.utc) second_time = first_time + timedelta(days=15) assert _collect(config, sysfs_controller, first_time, 1000, 2000)["ok"] is True writer = sqlite3.connect(store_path) writer.execute( "CREATE TRIGGER stop_retention BEFORE DELETE ON samples " "BEGIN SELECT RAISE(ABORT, 'simulated retention interruption'); END" ) writer.commit() writer.close() second = _collect(config, sysfs_controller, second_time, 1001, 2001) assert second["ok"] is True, second assert "simulated retention interruption" in second["retention_error"] with sqlite3.connect(store_path) as conn: assert conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0] == 2 assert conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0 assert conn.execute("SELECT COUNT(*) FROM local_day_unallocated_evidence").fetchone()[0] == 1 writer = sqlite3.connect(store_path) writer.execute("DROP TRIGGER stop_retention") writer.commit() writer.close() third_time = second_time + timedelta(days=15) third = _collect(config, sysfs_controller, third_time, 1002, 2002) assert third["ok"] is True, third assert third["retention_error"] is None with sqlite3.connect(store_path) as conn: assert conn.execute("SELECT ts FROM samples ORDER BY id").fetchall() == [ (third_time.isoformat(),) ] def test_collection_recovers_old_pending_work_before_pruning( tmp_path, sysfs_controller, monkeypatch, ): monkeypatch.setenv("TZ", "UTC") store_path = tmp_path / "pending-recovery.db" config = {"device": "/dev/nvme0", "store_path": str(store_path)} first_time = datetime(2026, 9, 1, 0, 0, tzinfo=timezone.utc) second_time = first_time + timedelta(days=15) assert _collect(config, sysfs_controller, first_time, 1000, 2000)["ok"] is True writer = sqlite3.connect(store_path) writer.execute( "CREATE TRIGGER interrupt_publication " "BEFORE INSERT ON local_day_unallocated_evidence " "BEGIN SELECT RAISE(ABORT, 'simulated publication interruption'); END" ) writer.commit() writer.close() failed = _collect(config, sysfs_controller, second_time, 1001, 2001) assert failed["ok"] is False assert "simulated publication interruption" in failed["error"] with sqlite3.connect(store_path) as conn: assert conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0] == 1 assert conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 1 writer = sqlite3.connect(store_path) writer.execute("DROP TRIGGER interrupt_publication") writer.commit() writer.close() third_time = second_time + timedelta(minutes=5) recovered = _collect(config, sysfs_controller, third_time, 1002, 2002) assert recovered["ok"] is True, recovered assert recovered["retention_error"] is None with sqlite3.connect(store_path) as conn: assert conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0 assert conn.execute("SELECT ts FROM samples ORDER BY id").fetchall() == [ (second_time.isoformat(),), (third_time.isoformat(),), ] assert conn.execute( "SELECT COUNT(*) FROM local_day_unallocated_evidence" ).fetchone()[0] == 1