Compare commits

..
3 Commits
7 changed files with 519 additions and 81 deletions
+9 -10
View File
@@ -27,7 +27,6 @@ if VENV_DIR.exists():
if site_packages: if site_packages:
sys.path.insert(0, str(site_packages)) sys.path.insert(0, str(site_packages))
from fenris.store import init_store, get_store_path
from fenris.collector import run_collection from fenris.collector import run_collection
@@ -103,14 +102,6 @@ def main() -> None:
config = load_config() config = load_config()
device = config["device"] 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 # Inject a simple clock
class SimpleClock: class SimpleClock:
def utcnow(self): def utcnow(self):
@@ -118,8 +109,16 @@ def main() -> None:
clock = SimpleClock() 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 # Run collection
result = run_collection(smartctl_data, sysfs_path, config, clock) result = run_collection(config=config, clock=clock, acquire=acquire)
if result["ok"]: if result["ok"]:
print(f"Collection successful: {result['sample_count']} sample(s)") print(f"Collection successful: {result['sample_count']} sample(s)")
+167 -55
View File
@@ -5,19 +5,21 @@ This module implements the thinnest complete write path:
- Acquire controller identity from sysfs - Acquire controller identity from sysfs
- Normalize identity exactly once at write time - Normalize identity exactly once at write time
- Validate every row against store invariants - Validate every row against store invariants
- Stage valid observations privately before derivation
- Publish the sample and derived evidence in one collection-owned transaction - Publish the sample and derived evidence in one collection-owned transaction
No code path outside the collector interrogates the device. No code path outside the collector interrogates the device.
""" """
import json import json
import sqlite3 import sqlite3
from collections.abc import Callable
from datetime import datetime, timezone from datetime import datetime, timezone
from pathlib import Path 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 .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): class AcquisitionError(Exception):
@@ -30,6 +32,9 @@ class InvariantViolationError(Exception):
pass pass
PENDING_PUBLICATION_LIMIT = 6720
def acquire_from_smartctl(smartctl_data: Dict[str, Any]) -> Dict[str, Any]: def acquire_from_smartctl(smartctl_data: Dict[str, Any]) -> Dict[str, Any]:
"""Acquire counters and thermal evidence from smartctl -a -j data. """Acquire counters and thermal evidence from smartctl -a -j data.
@@ -195,7 +200,7 @@ def write_sample(
sample: Dict[str, Any], sample: Dict[str, Any],
identity: Dict[str, Any], identity: Dict[str, Any],
conn: sqlite3.Connection, conn: sqlite3.Connection,
clock, observed_at: datetime,
) -> Dict[str, Any]: ) -> Dict[str, Any]:
"""Write one sample to the observation store. """Write one sample to the observation store.
@@ -219,8 +224,7 @@ def write_sample(
# Open new segment if needed # Open new segment if needed
segment_opened = False segment_opened = False
if should_open: if should_open:
now = clock.utcnow() open_segment(conn, observed_at, identity, identity_key, identity_degraded)
open_segment(conn, now, identity, identity_key, identity_degraded)
segment_opened = True segment_opened = True
# Get current segment_id for provenance # Get current segment_id for provenance
@@ -271,48 +275,25 @@ def write_sample(
} }
def run_collection( def _observation_time(sample: Dict[str, Any]) -> datetime:
smartctl_data: Dict[str, Any], """Read an observation's original timestamp for ordered recovery."""
sysfs_path: Path, observed_at = datetime.fromisoformat(sample["ts"])
config: Dict[str, Any], if observed_at.tzinfo is None:
clock, return observed_at.replace(tzinfo=timezone.utc)
) -> Dict[str, Any]: return observed_at.astimezone(timezone.utc)
"""Run one collection run.
This is the main entry point for the collector.
Returns the run outcome.
"""
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
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())
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) validate_sample_invariants(sample, conn)
seg_info = write_sample(sample, identity, conn, clock) 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") cursor = conn.execute("SELECT id FROM samples ORDER BY id DESC LIMIT 1")
current_id = cursor.fetchone()[0] current_id = cursor.fetchone()[0]
@@ -331,31 +312,162 @@ def run_collection(
derive_hours_from_interval(conn, prev, current) derive_hours_from_interval(conn, prev, current)
from .day_aggregate import derive_all_days, persist_day_aggregate from .day_aggregate import derive_all_days, persist_day_aggregate
for agg in derive_all_days(conn): for aggregate in derive_all_days(conn):
persist_day_aggregate(conn, agg) persist_day_aggregate(conn, aggregate)
from .tz_util import detect_system_tz
from .local_day import derive_local_day_summary, persist_local_day 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, observed_at)
local_summary = derive_local_day_summary(conn, tz_name, clock.utcnow())
if local_summary is not None: if local_summary is not None:
persist_local_day(conn, local_summary) persist_local_day(conn, local_summary)
conn.commit()
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: 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]:
"""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:
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)
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)
published_count = _recover_pending_for_collection(conn)
# 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:
conn.rollback()
published_count += _recover_pending_for_collection(conn)
continue
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,
}
validate_sample_invariants(sample, conn)
_stage_observation(conn, sample, identity, tz_name)
conn.commit()
break
published_count += _recover_pending_for_collection(conn)
return { return {
"ok": True, "ok": True,
"sample_count": 1, "sample_count": published_count,
"store_path": str(store_path), "store_path": str(store_path),
} }
except Exception as e: except Exception as exc:
if conn is not None: if conn is not None:
conn.rollback() conn.rollback()
return { return {
"ok": False, "ok": False,
"error": str(e), "error": str(exc),
"error_type": type(e).__name__, "error_type": type(exc).__name__,
} }
finally: finally:
if conn is not None: if conn is not None:
+17
View File
@@ -115,6 +115,7 @@ class StatusComposition:
# Sample counts for waiting explanations # Sample counts for waiting explanations
sample_count: int = 0 sample_count: int = 0
day_count: int = 0 day_count: int = 0
pending_publication_count: int = 0
# Whether the status dot should blink (only Monitoring) # Whether the status dot should blink (only Monitoring)
should_blink: bool = False should_blink: bool = False
@@ -286,6 +287,7 @@ def compose_status(
freshness_age_s = None freshness_age_s = None
sample_count = 0 sample_count = 0
day_count = 0 day_count = 0
pending_publication_count = 0
deliberately_paused = False deliberately_paused = False
if conn is None and store_fault is None and newer_schema is None: 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] sample_count = cursor.fetchone()[0]
cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates") cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates")
day_count = cursor.fetchone()[0] 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: except sqlite3.Error as exc:
store_fault = str(exc) store_fault = str(exc)
@@ -328,6 +334,7 @@ def compose_status(
freshness = "unknown" freshness = "unknown"
freshness_age_s = None freshness_age_s = None
sample_count = day_count = 0 sample_count = day_count = 0
pending_publication_count = 0
deliberately_paused = False deliberately_paused = False
# --- External stop detection --- # --- External stop detection ---
@@ -417,6 +424,7 @@ def compose_status(
overlay_base=overlay_base, overlay_base=overlay_base,
sample_count=sample_count, sample_count=sample_count,
day_count=day_count, day_count=day_count,
pending_publication_count=pending_publication_count,
continuity=continuity, continuity=continuity,
paused_lines=paused_lines, paused_lines=paused_lines,
) )
@@ -459,6 +467,8 @@ def render_status_cli(comp: StatusComposition) -> str:
facts.append("timer: %s" % ( facts.append("timer: %s" % (
"unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive" "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: if facts:
lines.append(" · ".join(facts)) lines.append(" · ".join(facts))
@@ -524,6 +534,8 @@ def render_status_tui(comp: StatusComposition) -> str:
facts.append("Timer: %s" % ( facts.append("Timer: %s" % (
"unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive" "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: if facts:
lines.append(" · ".join(facts)) lines.append(" · ".join(facts))
@@ -539,3 +551,8 @@ def render_status_tui(comp: StatusComposition) -> str:
lines.append(pl[:1].upper() + pl[1:]) lines.append(pl[:1].upper() + pl[1:])
return "\n".join(lines) 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"
+22 -4
View File
@@ -3,8 +3,7 @@
This module handles: This module handles:
- Store initialization with WAL mode - Store initialization with WAL mode
- Schema versioning with PRAGMA user_version - Schema versioning with PRAGMA user_version
- The six entities: samples, hour_observations, day_aggregates, - Observation records, derived activity, publication state, and metadata
monitoring_periods, controller_segments, endurance_baseline
""" """
import sqlite3 import sqlite3
from pathlib import Path from pathlib import Path
@@ -12,7 +11,7 @@ from typing import Optional
# Schema version - increment on each migration # 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 # Packaged default placement (spec §8.3). The config may override it, but a
@@ -33,7 +32,7 @@ def get_store_path(config: dict) -> Path:
def init_store(store_path: Path) -> sqlite3.Connection: def init_store(store_path: Path) -> sqlite3.Connection:
"""Initialize the observation store if not present. """Initialize the observation store if not present.
Creates the database with WAL mode and all six entities. Creates the observation schema, including private pending-publication storage.
Returns a connection to the store. Returns a connection to the store.
""" """
conn = sqlite3.connect(str(store_path)) conn = sqlite3.connect(str(store_path))
@@ -221,6 +220,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): def _apply_migrations(conn: sqlite3.Connection, current_version: int):
"""Apply forward-only migrations from current_version to SCHEMA_VERSION. """Apply forward-only migrations from current_version to SCHEMA_VERSION.
@@ -275,6 +287,12 @@ def _apply_migrations(conn: sqlite3.Connection, current_version: int):
""") """)
current_version = 3 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: def migrate_to_latest(store_path: Path) -> int:
"""Apply forward-only migrations to bring the store to SCHEMA_VERSION. """Apply forward-only migrations to bring the store to SCHEMA_VERSION.
+1 -1
View File
@@ -154,7 +154,7 @@ class TestSchemaMigration:
# Migrate # Migrate
steps = migrate_to_latest(db) steps = migrate_to_latest(db)
assert steps == 2 # v1→v2→v3 assert steps == 3 # v1→v2→v3→v4
# Verify data preserved # Verify data preserved
conn = sqlite3.connect(str(db)) conn = sqlite3.connect(str(db))
+2 -1
View File
@@ -232,7 +232,7 @@ def test_store_initialization(config_fixture: Dict[str, Any]):
# Initialize store # Initialize store
conn = init_store(store_path) conn = init_store(store_path)
# Verify all six entities exist # Verify the observation, derived-history, and metadata tables exist.
cursor = conn.execute("SELECT name FROM sqlite_master WHERE type='table'") cursor = conn.execute("SELECT name FROM sqlite_master WHERE type='table'")
tables = {row[0] for row in cursor.fetchall()} tables = {row[0] for row in cursor.fetchall()}
@@ -249,6 +249,7 @@ def test_store_initialization(config_fixture: Dict[str, Any]):
expected_tables.add("sqlite_sequence") expected_tables.add("sqlite_sequence")
expected_tables.add("store_metadata") expected_tables.add("store_metadata")
expected_tables.add("local_days") expected_tables.add("local_days")
expected_tables.add("pending_publications")
assert expected_tables == tables assert expected_tables == tables
conn.close() conn.close()
+291
View File
@@ -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()