Record host read and write commands, controller busy time, error log entries, warning and critical temperature time, and thermal management transitions from the existing smartctl -a -j acquisition. Schema 7 adds nullable columns to samples so legacy rows read as unknown, and the Diagnostics panel shows each counter with its change.
295 lines
11 KiB
Python
295 lines
11 KiB
Python
"""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 SCHEMA_VERSION, 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(
|
|
minute=35, 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
|
|
|
|
monkeypatch.setattr(collector, "PENDING_PUBLICATION_LIMIT", 1)
|
|
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] == SCHEMA_VERSION
|
|
assert migrated.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0
|
|
finally:
|
|
migrated.close()
|