From 0d45f263d250823a1661b7dd577494f44f8d9e4a Mon Sep 17 00:00:00 2001 From: xavierk Date: Mon, 28 Sep 2026 03:08:00 +0530 Subject: [PATCH] feat: retain pending publications (#97) --- src/fenris/collect.py | 19 +- src/fenris/collector.py | 246 +++++++++++++++------ src/fenris/status_composition.py | 17 ++ src/fenris/store.py | 21 +- tests/test_collector_history_tracer.py | 2 +- tests/test_issue_97.py | 291 +++++++++++++++++++++++++ 6 files changed, 519 insertions(+), 77 deletions(-) create mode 100644 tests/test_issue_97.py diff --git a/src/fenris/collect.py b/src/fenris/collect.py index 840e3f7..b2a23a7 100644 --- a/src/fenris/collect.py +++ b/src/fenris/collect.py @@ -27,7 +27,6 @@ if VENV_DIR.exists(): if site_packages: sys.path.insert(0, str(site_packages)) -from fenris.store import init_store, get_store_path from fenris.collector import run_collection @@ -103,14 +102,6 @@ def main() -> None: config = load_config() device = config["device"] - # Interrogate the drive - smartctl_data = interrogate_drive(device) - - # Find sysfs path - sysfs_path = find_nvme_sysfs() - if sysfs_path is None: - raise RuntimeError("No NVMe controller found in sysfs") - # Inject a simple clock class SimpleClock: def utcnow(self): @@ -118,8 +109,16 @@ def main() -> None: clock = SimpleClock() + def acquire(): + """Interrogate the drive after pending-work preflight.""" + smartctl_data = interrogate_drive(device) + sysfs_path = find_nvme_sysfs() + if sysfs_path is None: + raise RuntimeError("No NVMe controller found in sysfs") + return smartctl_data, sysfs_path + # Run collection - result = run_collection(smartctl_data, sysfs_path, config, clock) + result = run_collection(config=config, clock=clock, acquire=acquire) if result["ok"]: print(f"Collection successful: {result['sample_count']} sample(s)") diff --git a/src/fenris/collector.py b/src/fenris/collector.py index 6c549c4..21f829f 100644 --- a/src/fenris/collector.py +++ b/src/fenris/collector.py @@ -5,19 +5,21 @@ This module implements the thinnest complete write path: - Acquire controller identity from sysfs - Normalize identity exactly once at write time - Validate every row against store invariants +- Stage valid observations privately before derivation - Publish the sample and derived evidence in one collection-owned transaction No code path outside the collector interrogates the device. """ import json import sqlite3 +from collections.abc import Callable from datetime import datetime, timezone from pathlib import Path -from typing import Any, Dict, Optional, Tuple +from typing import Any, Dict, Optional -from .store import init_store, get_store_path +from .derive import derive_hours_from_interval, find_previous_sample from .monitoring_periods import ensure_period_open -from .derive import find_previous_sample, derive_hours_from_interval +from .store import init_store, get_store_path class AcquisitionError(Exception): @@ -30,6 +32,9 @@ class InvariantViolationError(Exception): pass +PENDING_PUBLICATION_LIMIT = 6720 + + def acquire_from_smartctl(smartctl_data: Dict[str, Any]) -> Dict[str, Any]: """Acquire counters and thermal evidence from smartctl -a -j data. @@ -195,7 +200,7 @@ def write_sample( sample: Dict[str, Any], identity: Dict[str, Any], conn: sqlite3.Connection, - clock, + observed_at: datetime, ) -> Dict[str, Any]: """Write one sample to the observation store. @@ -219,8 +224,7 @@ def write_sample( # Open new segment if needed segment_opened = False if should_open: - now = clock.utcnow() - open_segment(conn, now, identity, identity_key, identity_degraded) + open_segment(conn, observed_at, identity, identity_key, identity_degraded) segment_opened = True # Get current segment_id for provenance @@ -271,91 +275,203 @@ def write_sample( } +def _observation_time(sample: Dict[str, Any]) -> datetime: + """Read an observation's original timestamp for ordered recovery.""" + observed_at = datetime.fromisoformat(sample["ts"]) + if observed_at.tzinfo is None: + return observed_at.replace(tzinfo=timezone.utc) + return observed_at.astimezone(timezone.utc) + + +def _publish_observation( + conn: sqlite3.Connection, + sample: Dict[str, Any], + identity: Dict[str, Any], + tz_name: str, +) -> None: + """Publish one staged sample and all dependent evidence in caller transaction.""" + observed_at = _observation_time(sample) + validate_sample_invariants(sample, conn) + ensure_period_open(conn, observed_at) + seg_info = write_sample(sample, identity, conn, observed_at) + + cursor = conn.execute("SELECT id FROM samples ORDER BY id DESC LIMIT 1") + current_id = cursor.fetchone()[0] + prev = find_previous_sample(conn, seg_info.get("segment_id"), current_id) + if prev is not None: + current = { + "id": current_id, + "ts": sample["ts"], + "bytes_written": sample["bytes_written"], + "bytes_read": sample["bytes_read"], + "power_on_hours": sample["power_on_hours"], + "temperature_c": sample["temperature_c"], + "data_units_written": sample["data_units_written"], + "data_units_read": sample["data_units_read"], + } + derive_hours_from_interval(conn, prev, current) + + from .day_aggregate import derive_all_days, persist_day_aggregate + for aggregate in derive_all_days(conn): + persist_day_aggregate(conn, aggregate) + + from .local_day import derive_local_day_summary, persist_local_day + local_summary = derive_local_day_summary(conn, tz_name, observed_at) + if local_summary is not None: + persist_local_day(conn, local_summary) + + +def _pending_count(conn: sqlite3.Connection) -> int: + return conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] + + +def _stage_observation( + conn: sqlite3.Connection, + sample: Dict[str, Any], + identity: Dict[str, Any], + tz_name: str, +) -> int: + """Durably stage valid acquired evidence before attempting publication.""" + payload = json.dumps( + {"sample": sample, "identity": identity, "tz_name": tz_name}, + separators=(",", ":"), + sort_keys=True, + ) + cursor = conn.execute( + "INSERT INTO pending_publications (sample_ts, payload) VALUES (?, ?)", + (sample["ts"], payload), + ) + return cursor.lastrowid + + +def _recover_pending(conn: sqlite3.Connection) -> int: + """Publish pending observations oldest first; stop at first failure.""" + recovered = 0 + while True: + conn.execute("BEGIN IMMEDIATE") + try: + row = conn.execute( + "SELECT id, payload FROM pending_publications ORDER BY id LIMIT 1" + ).fetchone() + if row is None: + conn.rollback() + return recovered + + pending_id, payload = row + observation = json.loads(payload) + _publish_observation( + conn, + observation["sample"], + observation["identity"], + observation["tz_name"], + ) + conn.execute("DELETE FROM pending_publications WHERE id = ?", (pending_id,)) + conn.commit() + recovered += 1 + except Exception: + conn.rollback() + raise + + +def _recover_pending_for_collection(conn: sqlite3.Connection) -> int: + """Retry queued work and report exhausted capacity without masking store faults.""" + try: + return _recover_pending(conn) + except Exception as exc: + if isinstance(exc, sqlite3.Error) and not isinstance(exc, sqlite3.IntegrityError): + raise + if isinstance(exc, InvariantViolationError): + raise + try: + capacity_full = _pending_count(conn) >= PENDING_PUBLICATION_LIMIT + except sqlite3.Error: + raise exc + if capacity_full: + raise RuntimeError( + "pending publication capacity full " + f"({PENDING_PUBLICATION_LIMIT} observations); no new observation acquired; " + f"recovery failed: {exc}" + ) from exc + raise + + def run_collection( - smartctl_data: Dict[str, Any], - sysfs_path: Path, - config: Dict[str, Any], - clock, + smartctl_data: Optional[Dict[str, Any]] = None, + sysfs_path: Optional[Path] = None, + config: Optional[Dict[str, Any]] = None, + clock=None, + *, + acquire: Optional[Callable[[], tuple[Dict[str, Any], Path]]] = None, ) -> Dict[str, Any]: - """Run one collection run. - - This is the main entry point for the collector. - Returns the run outcome. + """Recover old work, acquire one observation, then publish it atomically. + + Production callers pass ``acquire`` so recovery and capacity checks run + before device interrogation. Direct sample arguments remain useful for + deterministic collector tests. """ conn = None try: - # Acquire counters and thermal evidence - counters = acquire_from_smartctl(smartctl_data) - - # Acquire controller identity - identity = acquire_from_sysfs(sysfs_path) - - # Build sample with injected clock - sample = { - "ts": clock.utcnow().isoformat(), - "device": config["device"], - **counters, - **identity, - } - - # Initialize store if needed + if config is None or clock is None: + raise ValueError("config and clock are required") + store_path = get_store_path(config) conn = init_store(store_path) - # Run legacy import if needed (idempotent) from .legacy import import_legacy_history history_path = Path(config.get("data_dir", ".")) / "history.jsonl" if history_path.exists(): import_legacy_history(conn, history_path, clock=clock) - # Collection owns one transaction for the sample and its evidence. - ensure_period_open(conn, clock.utcnow()) + recovered_count = _recover_pending_for_collection(conn) - validate_sample_invariants(sample, conn) - seg_info = write_sample(sample, identity, conn, clock) + # Hold the writer reservation across the capacity check and acquisition. + # A concurrent collector will recheck pending work before it acquires. + while True: + conn.execute("BEGIN IMMEDIATE") + waiting = _pending_count(conn) + if waiting >= PENDING_PUBLICATION_LIMIT: + conn.rollback() + recovered_count += _recover_pending_for_collection(conn) + continue + if waiting: + conn.rollback() + recovered_count += _recover_pending_for_collection(conn) + continue - cursor = conn.execute("SELECT id FROM samples ORDER BY id DESC LIMIT 1") - current_id = cursor.fetchone()[0] - prev = find_previous_sample(conn, seg_info.get("segment_id"), current_id) - if prev is not None: - current = { - "id": current_id, - "ts": sample["ts"], - "bytes_written": sample["bytes_written"], - "bytes_read": sample["bytes_read"], - "power_on_hours": sample["power_on_hours"], - "temperature_c": sample["temperature_c"], - "data_units_written": sample["data_units_written"], - "data_units_read": sample["data_units_read"], + from .tz_util import detect_system_tz + tz_name = detect_system_tz() + if acquire is not None: + smartctl_data, sysfs_path = acquire() + if smartctl_data is None or sysfs_path is None: + raise AcquisitionError("No acquired SMART data or sysfs identity") + + counters = acquire_from_smartctl(smartctl_data) + identity = acquire_from_sysfs(sysfs_path) + sample = { + "ts": clock.utcnow().isoformat(), + "device": config["device"], + **counters, + **identity, } - derive_hours_from_interval(conn, prev, current) - - from .day_aggregate import derive_all_days, persist_day_aggregate - for agg in derive_all_days(conn): - persist_day_aggregate(conn, agg) - - from .tz_util import detect_system_tz - from .local_day import derive_local_day_summary, persist_local_day - tz_name = detect_system_tz() - local_summary = derive_local_day_summary(conn, tz_name, clock.utcnow()) - if local_summary is not None: - persist_local_day(conn, local_summary) - - conn.commit() + validate_sample_invariants(sample, conn) + _stage_observation(conn, sample, identity, tz_name) + conn.commit() + break + recovered_count += _recover_pending_for_collection(conn) return { "ok": True, - "sample_count": 1, + "sample_count": recovered_count, "store_path": str(store_path), } - except Exception as e: + except Exception as exc: if conn is not None: conn.rollback() return { "ok": False, - "error": str(e), - "error_type": type(e).__name__, + "error": str(exc), + "error_type": type(exc).__name__, } finally: if conn is not None: diff --git a/src/fenris/status_composition.py b/src/fenris/status_composition.py index bb39452..f3e632e 100644 --- a/src/fenris/status_composition.py +++ b/src/fenris/status_composition.py @@ -115,6 +115,7 @@ class StatusComposition: # Sample counts for waiting explanations sample_count: int = 0 day_count: int = 0 + pending_publication_count: int = 0 # Whether the status dot should blink (only Monitoring) should_blink: bool = False @@ -286,6 +287,7 @@ def compose_status( freshness_age_s = None sample_count = 0 day_count = 0 + pending_publication_count = 0 deliberately_paused = False if conn is None and store_fault is None and newer_schema is None: @@ -315,6 +317,10 @@ def compose_status( sample_count = cursor.fetchone()[0] cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates") day_count = cursor.fetchone()[0] + schema_version = conn.execute("PRAGMA user_version").fetchone()[0] + if schema_version >= 4: + cursor = conn.execute("SELECT COUNT(*) FROM pending_publications") + pending_publication_count = cursor.fetchone()[0] except sqlite3.Error as exc: store_fault = str(exc) @@ -328,6 +334,7 @@ def compose_status( freshness = "unknown" freshness_age_s = None sample_count = day_count = 0 + pending_publication_count = 0 deliberately_paused = False # --- External stop detection --- @@ -417,6 +424,7 @@ def compose_status( overlay_base=overlay_base, sample_count=sample_count, day_count=day_count, + pending_publication_count=pending_publication_count, continuity=continuity, paused_lines=paused_lines, ) @@ -459,6 +467,8 @@ def render_status_cli(comp: StatusComposition) -> str: facts.append("timer: %s" % ( "unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive" )) + if comp.pending_publication_count: + facts.append(_pending_publication_text(comp.pending_publication_count)) if facts: lines.append(" · ".join(facts)) @@ -524,6 +534,8 @@ def render_status_tui(comp: StatusComposition) -> str: facts.append("Timer: %s" % ( "unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive" )) + if comp.pending_publication_count: + facts.append(_pending_publication_text(comp.pending_publication_count).capitalize()) if facts: lines.append(" · ".join(facts)) @@ -539,3 +551,8 @@ def render_status_tui(comp: StatusComposition) -> str: lines.append(pl[:1].upper() + pl[1:]) return "\n".join(lines) + + +def _pending_publication_text(count: int) -> str: + noun = "observation" if count == 1 else "observations" + return f"pending publication: {count} {noun} retained for retry" diff --git a/src/fenris/store.py b/src/fenris/store.py index 796b0e3..64a64ab 100644 --- a/src/fenris/store.py +++ b/src/fenris/store.py @@ -12,7 +12,7 @@ from typing import Optional # Schema version - increment on each migration -SCHEMA_VERSION = 3 +SCHEMA_VERSION = 4 # Packaged default placement (spec §8.3). The config may override it, but a @@ -221,6 +221,19 @@ def _create_schema(conn: sqlite3.Connection): ) """) + _create_pending_publications(conn) + + +def _create_pending_publications(conn: sqlite3.Connection) -> None: + """Create private staging for valid observations awaiting derivation.""" + conn.execute(""" + CREATE TABLE IF NOT EXISTS pending_publications ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + sample_ts TEXT NOT NULL, + payload TEXT NOT NULL + ) + """) + def _apply_migrations(conn: sqlite3.Connection, current_version: int): """Apply forward-only migrations from current_version to SCHEMA_VERSION. @@ -275,6 +288,12 @@ def _apply_migrations(conn: sqlite3.Connection, current_version: int): """) current_version = 3 + # Migration 3→4: retain acquired observations until derived evidence can + # be published atomically (issue #97, ADR 0011). + if current_version < 4: + _create_pending_publications(conn) + current_version = 4 + def migrate_to_latest(store_path: Path) -> int: """Apply forward-only migrations to bring the store to SCHEMA_VERSION. diff --git a/tests/test_collector_history_tracer.py b/tests/test_collector_history_tracer.py index 13a3a10..b6d3479 100644 --- a/tests/test_collector_history_tracer.py +++ b/tests/test_collector_history_tracer.py @@ -154,7 +154,7 @@ class TestSchemaMigration: # Migrate steps = migrate_to_latest(db) - assert steps == 2 # v1→v2→v3 + assert steps == 3 # v1→v2→v3→v4 # Verify data preserved conn = sqlite3.connect(str(db)) diff --git a/tests/test_issue_97.py b/tests/test_issue_97.py new file mode 100644 index 0000000..ef08b29 --- /dev/null +++ b/tests/test_issue_97.py @@ -0,0 +1,291 @@ +"""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] == 4 + assert migrated.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0 + finally: + migrated.close()