"""Pending publication behavior from issue #97.""" 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 import collector from fenris.projection import compute_projection from fenris.status import _query_drive_facts, get_status, read_status from fenris.status_composition import render_status_tui from fenris.store import init_store class FakeClock: def __init__(self, now): self.now = now def utcnow(self): return self.now @pytest.fixture def sysfs_fixture_tree(tmp_path): ctrl_dir = tmp_path / "sys" / "class" / "nvme" / "nvme0" ctrl_dir.mkdir(parents=True) (ctrl_dir / "subsysnqn").write_text("nqn.test:drive\n") (ctrl_dir / "model").write_text("Test NVMe\n") (ctrl_dir / "serial").write_text("test-serial\n") (ctrl_dir / "firmware_rev").write_text("1.0\n") transport_dir = ctrl_dir / "transport" transport_dir.mkdir() (transport_dir / "trstring").write_text("pcie\n") return tmp_path def _smartctl(unit_written, unit_read=500): return { "nvme_smart_health_information_log": { "critical_warning": 0, "temperature": 35, "available_spare": 100, "percentage_used": 5, "data_units_written": unit_written, "data_units_read": unit_read, "power_on_hours": 100, }, "user_capacity": {"bytes": 1_024_000_000_000}, "model_name": "Test NVMe", "serial_number": "test-serial", "firmware_version": "1.0", } def _acquirer(sysfs_path, unit_written, unit_read=500): return lambda: (_smartctl(unit_written, unit_read), sysfs_path) def _public_snapshot(store_path, clock_now): with read_status(store_path, clock_now) as (conn, comp): assert conn is not None return { "state": comp.state, "freshness": comp.freshness, "sample_count": comp.sample_count, "day_count": comp.day_count, "boot_enabled": comp.boot_enabled, "timer_active": comp.timer_active, "last_collect_ok": comp.last_collect_ok, "latest_sample": tuple(conn.execute( "SELECT ts, subnqn, mn, sn, fr, capacity_bytes, temperature_c, " "available_spare, data_units_written, bytes_written, critical_warning, " "media_errors, unsafe_shutdowns " "FROM samples ORDER BY id DESC LIMIT 1" ).fetchone()), "drive_facts": tuple(_query_drive_facts(conn)), "projection": repr(compute_projection(conn, clock_now)), "segments": conn.execute("SELECT COUNT(*) FROM controller_segments").fetchone()[0], "periods": conn.execute("SELECT COUNT(*) FROM monitoring_periods").fetchone()[0], "hours": tuple(conn.execute( "SELECT hour, bytes_written_delta, bytes_read_delta, sample_count " "FROM hour_observations ORDER BY hour" ).fetchall()), "days": tuple(conn.execute( "SELECT day, bytes_written_delta, bytes_read_delta " "FROM day_aggregates ORDER BY day" ).fetchall()), "local_days": tuple(conn.execute( "SELECT local_date, bytes_written, bytes_read, sample_count " "FROM local_days ORDER BY local_date, tz_name" ).fetchall()), "pending": comp.pending_publication_count, } def test_failed_publication_survives_restart_and_readers_keep_last_consistent_view( tmp_path, sysfs_fixture_tree, monkeypatch ): """A failed derivation stays private, visible as pending, and replays in order.""" monkeypatch.setenv("TZ", "UTC") service = {"boot_enabled": True, "timer_active": True, "last_collect_ok": True} monkeypatch.setattr("fenris.status.query_service_state", lambda: service) now = datetime.now(timezone.utc).replace(second=0, microsecond=0) store_path = tmp_path / "observations.db" config = {"device": "/dev/nvme0", "store_path": str(store_path)} sysfs_path = sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0" for units, at in ((1000, now - timedelta(minutes=10)), (1010, now - timedelta(minutes=5))): result = collector.run_collection( config=config, clock=FakeClock(at), acquire=_acquirer(sysfs_path, units), ) assert result["ok"] is True, result previous = _public_snapshot(store_path, now) baseline_cli = get_status(store_path, now, query_services=False, query_journal=False) original_derive = collector.derive_hours_from_interval reader_during_write = {} def fail_derivation(conn, previous_sample, next_sample): # Read while writer has inserted the unpublished sample but not committed. with read_status(store_path, now, query_services=False) as (reader, comp): assert reader is not None reader_during_write["samples"] = reader.execute( "SELECT COUNT(*) FROM samples" ).fetchone()[0] reader_during_write["pending"] = comp.pending_publication_count reader_during_write["latest"] = tuple(reader.execute( "SELECT data_units_written FROM samples ORDER BY id DESC LIMIT 1" ).fetchone()) raise RuntimeError("injected non-invariant derivation failure") monkeypatch.setattr(collector, "derive_hours_from_interval", fail_derivation) failed = collector.run_collection( config=config, clock=FakeClock(now), acquire=_acquirer(sysfs_path, 1025), ) assert failed["ok"] is False assert "injected non-invariant derivation failure" in failed["error"] assert reader_during_write == {"samples": 2, "pending": 1, "latest": (1010,)} after_failure = _public_snapshot(store_path, now) assert {**after_failure, "pending": previous["pending"]} == previous assert after_failure["pending"] == 1 cli_text = get_status(store_path, now, query_services=False, query_journal=False) pending_fact = " ยท pending publication: 1 observation retained for retry" assert pending_fact in cli_text assert cli_text.replace(pending_fact, "", 1) == baseline_cli with read_status(store_path, now) as (_, comp): assert comp.state == previous["state"] assert comp.boot_enabled is True assert comp.timer_active is True assert "pending publication: 1 observation retained for retry" in render_status_tui(comp).lower() called = False def should_wait_for_recovery(): nonlocal called called = True return _smartctl(1030), sysfs_path blocked = collector.run_collection( config=config, clock=FakeClock(now + timedelta(minutes=5)), acquire=should_wait_for_recovery, ) assert blocked["ok"] is False assert called is False assert _public_snapshot(store_path, now)["pending"] == 1 monkeypatch.setattr(collector, "derive_hours_from_interval", original_derive) monkeypatch.setenv("TZ", "Asia/Kolkata") recovered = collector.run_collection( config=config, clock=FakeClock(now + timedelta(minutes=5)), acquire=_acquirer(sysfs_path, 1040), ) assert recovered["ok"] is True, recovered conn = sqlite3.connect(store_path) try: assert conn.execute( "SELECT data_units_written FROM samples ORDER BY id" ).fetchall() == [(1000,), (1010,), (1025,), (1040,)] assert conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0 assert conn.execute( "SELECT SUM(bytes_written_delta) FROM day_aggregates" ).fetchone()[0] == (1040 - 1000) * 512_000 assert conn.execute( "SELECT sample_count FROM local_days WHERE tz_name = 'UTC'" ).fetchone()[0] == 4 assert conn.execute( "SELECT COUNT(*) FROM local_days WHERE tz_name = 'Asia/Kolkata'" ).fetchone()[0] == 1 finally: conn.close() def test_invalid_sample_writes_no_sample_segment_period_or_pending_row(tmp_path, sysfs_fixture_tree): """Sample invariants run before staging or public monitoring changes.""" now = datetime.now(timezone.utc).replace(second=0, microsecond=0) store_path = tmp_path / "observations.db" config = {"device": "/dev/nvme0", "store_path": str(store_path)} sysfs_path = sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0" result = collector.run_collection( config=config, clock=FakeClock(now), acquire=_acquirer(sysfs_path, -1), ) assert result["ok"] is False assert result["error_type"] == "InvariantViolationError" conn = sqlite3.connect(store_path) try: assert conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0] == 0 assert conn.execute("SELECT COUNT(*) FROM controller_segments").fetchone()[0] == 0 assert conn.execute("SELECT COUNT(*) FROM monitoring_periods").fetchone()[0] == 0 assert conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0 finally: conn.close() def test_pending_capacity_retries_before_skipping_acquisition(tmp_path, sysfs_fixture_tree, monkeypatch): now = datetime.now(timezone.utc).replace(second=0, microsecond=0) store_path = tmp_path / "observations.db" config = {"device": "/dev/nvme0", "store_path": str(store_path)} sysfs_path = sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0" first = collector.run_collection( config=config, clock=FakeClock(now - timedelta(minutes=5)), acquire=_acquirer(sysfs_path, 1000), ) assert first["ok"] is True, first monkeypatch.setattr(collector, "PENDING_PUBLICATION_LIMIT", 1) def fail_derivation(*_): raise RuntimeError("derivation still unavailable") monkeypatch.setattr(collector, "derive_hours_from_interval", fail_derivation) failed = collector.run_collection( config=config, clock=FakeClock(now), acquire=_acquirer(sysfs_path, 1010), ) assert failed["ok"] is False assert "capacity full (1 observations)" in failed["error"] called = False def must_not_acquire(): nonlocal called called = True return _smartctl(1020), sysfs_path retry = collector.run_collection( config=config, clock=FakeClock(now + timedelta(minutes=5)), acquire=must_not_acquire, ) assert retry["ok"] is False assert "capacity full (1 observations)" in retry["error"] assert called is False def test_v3_readers_ignore_pending_table_until_store_migrates(tmp_path): store_path = tmp_path / "observations.db" conn = init_store(store_path) conn.execute("DROP TABLE pending_publications") conn.execute("PRAGMA user_version=3") conn.commit() conn.close() with read_status(store_path, datetime.now(timezone.utc), query_services=False) as (reader, comp): assert reader is not None assert comp.store_fault is None assert comp.pending_publication_count == 0 migrated = init_store(store_path) try: assert migrated.execute("PRAGMA user_version").fetchone()[0] == 5 assert migrated.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0 finally: migrated.close()