342 lines
13 KiB
Python
342 lines
13 KiB
Python
"""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
|