398 lines
16 KiB
Python
398 lines
16 KiB
Python
"""Repair and retention tests (issue #74).
|
|
|
|
Tests the idempotent, safe repair of hour observations and day aggregates
|
|
from surviving raw samples, boundary anchor retention, and legacy summary
|
|
handling at actual precision.
|
|
"""
|
|
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.store import init_store
|
|
from fenris.repair import (
|
|
repair_derivation,
|
|
is_repair_in_progress,
|
|
get_repair_status,
|
|
)
|
|
from fenris.pruning import prune_old_samples, needs_boundary_anchor
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Fixtures
|
|
# ---------------------------------------------------------------------------
|
|
|
|
@pytest.fixture
|
|
def store_conn(tmp_path: Path):
|
|
"""Create a fresh store for each test."""
|
|
db_path = tmp_path / "test.db"
|
|
conn = init_store(db_path)
|
|
yield conn
|
|
conn.close()
|
|
|
|
|
|
def _insert_sample(conn, ts_iso, bytes_written, bytes_read=0, power_on_hours=100,
|
|
device="/dev/nvme0", segment_id=None):
|
|
"""Insert a raw sample."""
|
|
conn.execute(
|
|
"""INSERT INTO samples
|
|
(ts, device, bytes_written, bytes_read, power_on_hours,
|
|
data_units_written, data_units_read, segment_id)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?)""",
|
|
(ts_iso, device, bytes_written, bytes_read, power_on_hours,
|
|
bytes_written // 512000, bytes_read // 512000, segment_id),
|
|
)
|
|
conn.commit()
|
|
|
|
|
|
def _insert_hour(conn, hour_iso, bytes_written_delta=0, active_seconds=3600,
|
|
idle_seconds=0, powered_off_seconds=0, unknown_seconds=0,
|
|
sample_count=1, coverage=1.0):
|
|
"""Insert an hour observation."""
|
|
conn.execute(
|
|
"""INSERT INTO hour_observations
|
|
(hour, active_seconds, idle_seconds, powered_off_seconds, unknown_seconds,
|
|
bytes_written_delta, bytes_read_delta, sample_count, coverage)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)""",
|
|
(hour_iso, active_seconds, idle_seconds, powered_off_seconds, unknown_seconds,
|
|
bytes_written_delta, 0, sample_count, coverage),
|
|
)
|
|
conn.commit()
|
|
|
|
|
|
def _insert_day_aggregate(conn, day, bytes_written_delta=0, coverage=1.0):
|
|
"""Insert a day aggregate."""
|
|
conn.execute(
|
|
"""INSERT INTO day_aggregates
|
|
(day, active_seconds, bytes_written_delta, coverage, sample_count)
|
|
VALUES (?, 3600, ?, ?, 1)""",
|
|
(day, bytes_written_delta, coverage),
|
|
)
|
|
conn.commit()
|
|
|
|
|
|
def _open_period(conn, start_iso, end_iso=None, end_cause=None):
|
|
"""Insert a monitoring period."""
|
|
conn.execute(
|
|
"""INSERT INTO monitoring_periods (started_at, ended_at, end_cause)
|
|
VALUES (?, ?, ?)""",
|
|
(start_iso, end_iso, end_cause),
|
|
)
|
|
conn.commit()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# AC1: Normal collection, legacy import and recovery use same evidence rules
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class TestRepairUsesSameEvidenceRules:
|
|
"""AC1: Repair derives only surviving supported evidence, transactionally
|
|
and idempotently."""
|
|
|
|
def test_repair_idempotent_on_empty_store(self, store_conn):
|
|
"""Repair on empty store succeeds and does nothing."""
|
|
result = repair_derivation(store_conn)
|
|
assert result.ok is True
|
|
assert result.hours_created == 0
|
|
assert result.days_created == 0
|
|
|
|
def test_repair_idempotent_on_fully_derived(self, store_conn):
|
|
"""Repair on store with existing derived data does not duplicate."""
|
|
now = datetime(2026, 9, 15, 12, 0, 0, tzinfo=timezone.utc)
|
|
_open_period(store_conn, "2026-09-14T00:00:00+00:00")
|
|
|
|
# Insert samples that span an hour
|
|
_insert_sample(store_conn, "2026-09-14T10:00:00+00:00", 1000000)
|
|
_insert_sample(store_conn, "2026-09-14T10:30:00+00:00", 2000000)
|
|
|
|
# Pre-existing hour observation for the 10:00 hour
|
|
# This hour observation captures the interval [10:00, 10:30]
|
|
_insert_hour(store_conn, "2026-09-14T10:00:00+00:00",
|
|
bytes_written_delta=1000000)
|
|
_insert_day_aggregate(store_conn, "2026-09-14",
|
|
bytes_written_delta=1000000)
|
|
|
|
# Run repair
|
|
result = repair_derivation(store_conn)
|
|
|
|
# Should not create new observations (already exist)
|
|
assert result.hours_created == 0
|
|
assert result.days_created == 0
|
|
|
|
# Verify no duplicates
|
|
cursor = store_conn.execute(
|
|
"SELECT COUNT(*) FROM hour_observations WHERE hour = '2026-09-14T10:00:00+00:00'"
|
|
)
|
|
assert cursor.fetchone()[0] == 1
|
|
|
|
def test_repair_preserves_import_markers(self, store_conn):
|
|
"""Repair does not remove legacy import markers."""
|
|
# Set legacy import marker
|
|
store_conn.execute(
|
|
"INSERT INTO store_metadata (key, value) VALUES ('legacy_imported', 'true')"
|
|
)
|
|
store_conn.commit()
|
|
|
|
# Run repair
|
|
result = repair_derivation(store_conn)
|
|
|
|
# Verify marker preserved
|
|
cursor = store_conn.execute(
|
|
"SELECT value FROM store_metadata WHERE key = 'legacy_imported'"
|
|
)
|
|
assert cursor.fetchone()[0] == "true"
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# AC2: Rerunning or interrupting repair produces no duplicates
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class TestRepairIdempotency:
|
|
"""AC2: No duplicated intervals or totals from repeated repair."""
|
|
|
|
def test_repair_no_duplicate_hours(self, store_conn):
|
|
"""Running repair twice produces no duplicate hour observations."""
|
|
_open_period(store_conn, "2026-09-14T00:00:00+00:00")
|
|
_insert_sample(store_conn, "2026-09-14T10:00:00+00:00", 1000000)
|
|
_insert_sample(store_conn, "2026-09-14T10:30:00+00:00", 2000000)
|
|
|
|
# First repair
|
|
result1 = repair_derivation(store_conn)
|
|
assert result1.hours_created == 1
|
|
|
|
# Second repair
|
|
result2 = repair_derivation(store_conn)
|
|
assert result2.hours_created == 0 # No new hours
|
|
|
|
# Verify only one hour observation
|
|
cursor = store_conn.execute("SELECT COUNT(*) FROM hour_observations")
|
|
assert cursor.fetchone()[0] == 1
|
|
|
|
def test_repair_no_duplicate_days(self, store_conn):
|
|
"""Running repair twice produces no duplicate day aggregates."""
|
|
_open_period(store_conn, "2026-09-14T00:00:00+00:00")
|
|
_insert_sample(store_conn, "2026-09-14T10:00:00+00:00", 1000000)
|
|
_insert_sample(store_conn, "2026-09-14T10:30:00+00:00", 2000000)
|
|
|
|
# First repair
|
|
result1 = repair_derivation(store_conn)
|
|
assert result1.days_created == 1
|
|
|
|
# Second repair
|
|
result2 = repair_derivation(store_conn)
|
|
assert result2.days_created == 0 # No new days
|
|
|
|
# Verify only one day aggregate
|
|
cursor = store_conn.execute("SELECT COUNT(*) FROM day_aggregates")
|
|
assert cursor.fetchone()[0] == 1
|
|
|
|
def test_interrupted_repair_preserves_evidence(self, store_conn):
|
|
"""If repair fails, prior valid history is preserved."""
|
|
_open_period(store_conn, "2026-09-14T00:00:00+00:00")
|
|
|
|
# Insert valid existing data
|
|
_insert_hour(store_conn, "2026-09-14T10:00:00+00:00",
|
|
bytes_written_delta=500000)
|
|
_insert_day_aggregate(store_conn, "2026-09-14",
|
|
bytes_written_delta=500000)
|
|
|
|
# Run repair (should succeed but not modify existing valid data)
|
|
result = repair_derivation(store_conn)
|
|
assert result.ok is True
|
|
|
|
# Verify existing data preserved
|
|
cursor = store_conn.execute(
|
|
"SELECT bytes_written_delta FROM hour_observations "
|
|
"WHERE hour = '2026-09-14T10:00:00+00:00'"
|
|
)
|
|
assert cursor.fetchone()[0] == 500000
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# AC3: Boundary anchor retention
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class TestBoundaryAnchorRetention:
|
|
"""AC3: Prune samples only after durable derivation; retain anchors."""
|
|
|
|
def test_needs_boundary_anchor_sample(self, store_conn):
|
|
"""Sample before 14-day boundary is needed for derivation."""
|
|
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
|
|
|
|
# Sample just before 14-day boundary (2026-09-15T23:55:00)
|
|
# is 14 days and 0.75 hours old (before cutoff at 2026-09-16T12:00:00)
|
|
_insert_sample(store_conn, "2026-09-15T23:55:00+00:00", 1000000)
|
|
|
|
# Sample just after boundary (2026-09-16T12:30:00)
|
|
# is 13 days and 23.5 hours old (within retention)
|
|
_insert_sample(store_conn, "2026-09-16T12:30:00+00:00", 2000000)
|
|
|
|
# The sample at 2026-09-15 is a boundary anchor because
|
|
# the interval spans the retention boundary
|
|
assert needs_boundary_anchor(store_conn, "2026-09-15T23:55:00+00:00", now)
|
|
|
|
def test_hour_only_derivation_does_not_replace_local_day_evidence(self, store_conn):
|
|
"""UTC hour evidence alone does not make a sample safe to prune."""
|
|
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
|
|
|
|
# Insert sample and fully derive its interval
|
|
_insert_sample(store_conn, "2026-09-14T10:00:00+00:00", 1000000)
|
|
_insert_sample(store_conn, "2026-09-14T10:30:00+00:00", 2000000)
|
|
|
|
# Hour observation already exists for this interval
|
|
_insert_hour(store_conn, "2026-09-14T10:00:00+00:00",
|
|
bytes_written_delta=1000000)
|
|
|
|
assert needs_boundary_anchor(store_conn, "2026-09-14T10:00:00+00:00", now)
|
|
|
|
def test_pruning_retains_boundary_anchors(self, store_conn):
|
|
"""Pruning keeps samples needed as boundary anchors."""
|
|
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
|
|
|
|
# Old sample before boundary (2026-09-14T23:55:00)
|
|
# is 15 days and 0.75 hours old (before cutoff at 2026-09-16T12:00:00)
|
|
_insert_sample(store_conn, "2026-09-14T23:55:00+00:00", 1000000)
|
|
|
|
# Sample after boundary (2026-09-16T12:30:00)
|
|
# is 13 days and 23.5 hours old (within retention)
|
|
_insert_sample(store_conn, "2026-09-16T12:30:00+00:00", 2000000)
|
|
|
|
# Recent sample
|
|
_insert_sample(store_conn, "2026-09-29T12:00:00+00:00", 3000000)
|
|
|
|
# Run pruning
|
|
pruned = prune_old_samples(store_conn, now, retention_days=14)
|
|
|
|
# The boundary anchor should be retained
|
|
cursor = store_conn.execute(
|
|
"SELECT COUNT(*) FROM samples WHERE ts = '2026-09-14T23:55:00+00:00'"
|
|
)
|
|
assert cursor.fetchone()[0] == 1
|
|
|
|
def test_pruning_keeps_samples_without_local_day_replacement(self, store_conn):
|
|
"""Pruning keeps samples when only the UTC-hour row exists."""
|
|
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
|
|
|
|
# Old sample with derived interval
|
|
_insert_sample(store_conn, "2026-09-10T10:00:00+00:00", 1000000)
|
|
_insert_sample(store_conn, "2026-09-10T10:30:00+00:00", 2000000)
|
|
|
|
# Hour observation exists for the interval
|
|
_insert_hour(store_conn, "2026-09-10T10:00:00+00:00",
|
|
bytes_written_delta=1000000)
|
|
|
|
# Run pruning
|
|
pruned = prune_old_samples(store_conn, now, retention_days=14)
|
|
|
|
# The source remains until local-day evidence is also durable.
|
|
cursor = store_conn.execute(
|
|
"SELECT COUNT(*) FROM samples WHERE ts = '2026-09-10T10:00:00+00:00'"
|
|
)
|
|
assert cursor.fetchone()[0] == 1
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# AC4: Legacy day-only summaries retain actual precision
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class TestLegacySummaryPrecision:
|
|
"""AC4: Legacy summaries at actual precision, no interpolation."""
|
|
|
|
def test_legacy_summary_not_reconstructed(self, store_conn):
|
|
"""Legacy day-only summaries are not interpolated to hour detail."""
|
|
# Insert a legacy-style day aggregate without hour observations
|
|
_insert_day_aggregate(store_conn, "2026-08-01",
|
|
bytes_written_delta=5000000)
|
|
|
|
# Run repair
|
|
result = repair_derivation(store_conn)
|
|
|
|
# Should not create hour observations for legacy day
|
|
cursor = store_conn.execute(
|
|
"SELECT COUNT(*) FROM hour_observations WHERE hour LIKE '2026-08-01%'"
|
|
)
|
|
assert cursor.fetchone()[0] == 0
|
|
|
|
def test_legacy_summary_no_double_counting(self, store_conn):
|
|
"""Legacy summaries and derived intervals don't double-count."""
|
|
# Insert legacy day aggregate
|
|
_insert_day_aggregate(store_conn, "2026-08-01",
|
|
bytes_written_delta=5000000)
|
|
|
|
# Run repair
|
|
result = repair_derivation(store_conn)
|
|
|
|
# Day aggregate should not be modified
|
|
cursor = store_conn.execute(
|
|
"SELECT bytes_written_delta FROM day_aggregates WHERE day = '2026-08-01'"
|
|
)
|
|
assert cursor.fetchone()[0] == 5000000
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# AC5: Status distinguishes evidence states
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class TestStatusEvidenceDistingushing:
|
|
"""AC5: Read-only views distinguish evidence states."""
|
|
|
|
def test_repair_status_available(self, store_conn):
|
|
"""Repair status is available for read-only views."""
|
|
status = get_repair_status(store_conn)
|
|
assert hasattr(status, 'last_repair')
|
|
assert hasattr(status, 'repair_in_progress')
|
|
assert hasattr(status, 'hours_derived')
|
|
assert hasattr(status, 'days_derived')
|
|
|
|
def test_repair_in_progress_flag(self, store_conn):
|
|
"""Repair in progress flag is trackable."""
|
|
assert not is_repair_in_progress(store_conn)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# AC6: Migration then collection then reader consumption
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class TestMigrationCollectionReader:
|
|
"""AC6: End-to-end migration, collection, and reader consumption."""
|
|
|
|
def test_store_with_raw_evidence_and_summaries(self, store_conn):
|
|
"""Store with raw evidence and old summaries works correctly."""
|
|
# Set up store with mixed data
|
|
_open_period(store_conn, "2026-08-01T00:00:00+00:00")
|
|
|
|
# Old day-only summary (legacy)
|
|
_insert_day_aggregate(store_conn, "2026-08-01",
|
|
bytes_written_delta=5000000)
|
|
|
|
# Recent raw samples
|
|
_insert_sample(store_conn, "2026-09-14T10:00:00+00:00", 1000000)
|
|
_insert_sample(store_conn, "2026-09-14T10:30:00+00:00", 2000000)
|
|
|
|
# Run repair
|
|
result = repair_derivation(store_conn)
|
|
assert result.ok is True
|
|
|
|
# Verify legacy summary preserved
|
|
cursor = store_conn.execute(
|
|
"SELECT bytes_written_delta FROM day_aggregates WHERE day = '2026-08-01'"
|
|
)
|
|
assert cursor.fetchone()[0] == 5000000
|
|
|
|
# Verify new hour observation created
|
|
cursor = store_conn.execute(
|
|
"SELECT COUNT(*) FROM hour_observations WHERE hour LIKE '2026-09-14%'"
|
|
)
|
|
assert cursor.fetchone()[0] == 1
|
|
|
|
def test_concurrent_reader_consistency(self, store_conn):
|
|
"""Reader sees consistent snapshot during repair."""
|
|
# This is more of a documentation test - SQLite WAL mode handles this
|
|
# We verify the store is in WAL mode
|
|
cursor = store_conn.execute("PRAGMA journal_mode")
|
|
assert cursor.fetchone()[0] == "wal"
|