Compare commits

..
2 Commits
7 changed files with 746 additions and 172 deletions
+5
View File
@@ -122,6 +122,11 @@ def main() -> None:
if result["ok"]:
print(f"Collection successful: {result['sample_count']} sample(s)")
if result.get("retention_error"):
print(
f"Retention deferred: {result['retention_error']}",
file=sys.stderr,
)
sys.exit(0)
else:
print(f"Collection failed: {result['error']}", file=sys.stderr)
+12
View File
@@ -477,9 +477,21 @@ def run_collection(
break
published_count += _recover_pending_for_collection(conn)
retention_error = None
try:
if _pending_count(conn) == 0:
from .pruning import prune_old_samples
prune_old_samples(conn, clock.utcnow())
except Exception as exc: # noqa: BLE001 - publication already committed.
# Publication already committed. Keep collection successful and
# retry atomic retention on the next normal collection.
conn.rollback()
retention_error = str(exc)
return {
"ok": True,
"sample_count": published_count,
"retention_error": retention_error,
"store_path": str(store_path),
}
+360 -145
View File
@@ -1,131 +1,332 @@
"""Raw sample pruning per spec §3.4, ST-5.
"""Expire detail only after UTC and local_days replacement evidence is durable.
Raw samples are pruned opportunistically to 14 days.
Hour observations and day aggregates are retained indefinitely.
Boundary anchors required for successor evidence are retained.
Local-day summaries (local_days table) are never touched by pruning.
They are persisted at collection time from hour observations and survive
raw-sample pruning because they depend on hour observations, not on raw
samples. This is the mechanism that keeps local-day history trustworthy
after detail expires (issue #93).
The 14-day cutoff never overrides local-day preservation, shared boundary
evidence, publication recovery, or the newest sample's successor-anchor role.
"""
import sqlite3
from datetime import datetime, timedelta, timezone
from datetime import date, datetime, timedelta, timezone
# Spec §3.4: Raw-sample retention
RAW_SAMPLE_RETENTION_DAYS = 14
def _utc_datetime(value: datetime) -> datetime:
"""Normalize a clock value to aware UTC."""
if value.tzinfo is None:
value = value.replace(tzinfo=timezone.utc)
return value.astimezone(timezone.utc)
def _parse_sample_time(value: str) -> datetime | None:
"""Parse old timestamps conservatively; malformed legacy values stay."""
try:
parsed = datetime.fromisoformat(value)
except (TypeError, ValueError):
return None
return _utc_datetime(parsed)
def _sample_rows(conn: sqlite3.Connection) -> list[tuple]:
return conn.execute(
"SELECT id, ts, segment_id, bytes_written, bytes_read, local_tz "
"FROM samples ORDER BY id"
).fetchall()
def _required_hour_window(
start: datetime,
end: datetime,
) -> tuple[datetime, datetime, int] | None:
"""Return first hour, exclusive end and count for [start, end)."""
if end <= start:
return None
first_hour = start.replace(minute=0, second=0, microsecond=0)
last_hour = end.replace(minute=0, second=0, microsecond=0)
try:
exclusive_end = last_hour if end == last_hour else last_hour + timedelta(hours=1)
except OverflowError:
return None
count = int((exclusive_end - first_hour).total_seconds() // 3600)
return first_hour, exclusive_end, count
def _has_utc_replacement(
conn: sqlite3.Connection,
previous: tuple,
current: tuple,
start: datetime,
end: datetime,
) -> bool:
window = _required_hour_window(start, end)
if window is None:
return False
first_hour, exclusive_end, expected_hour_count = window
actual_hour_count = conn.execute(
"SELECT COUNT(*) FROM hour_observations WHERE hour >= ? AND hour < ?",
(
first_hour.strftime("%Y-%m-%dT%H:00:00+00:00"),
exclusive_end.strftime("%Y-%m-%dT%H:00:00+00:00"),
),
).fetchone()[0]
if actual_hour_count < expected_hour_count:
return False
start_written, end_written = previous[3], current[3]
start_read, end_read = previous[4], current[4]
if None in (start_written, end_written, start_read, end_read):
return False
if end_written < start_written or end_read < start_read:
return False
bytes_written = end_written - start_written
bytes_read = end_read - start_read
if expected_hour_count == 1:
totals = conn.execute(
"SELECT bytes_written_delta, bytes_read_delta FROM hour_observations "
"WHERE hour = ?",
(first_hour.strftime("%Y-%m-%dT%H:00:00+00:00"),),
).fetchone()
else:
totals = conn.execute(
"SELECT unattributed_bytes_written, unattributed_bytes_read "
"FROM day_aggregates WHERE day = ?",
(start.date().isoformat(),),
).fetchone()
if totals is None or totals[0] < bytes_written or totals[1] < bytes_read:
return False
last_day = (exclusive_end - timedelta(microseconds=1)).date()
expected_day_count = (last_day - first_hour.date()).days + 1
actual_day_count = conn.execute(
"SELECT COUNT(*) FROM day_aggregates WHERE day >= ? AND day <= ?",
(first_hour.date().isoformat(), last_day.isoformat()),
).fetchone()[0]
return actual_day_count >= expected_day_count
def _local_day_row_exists(
conn: sqlite3.Connection,
local_date: str,
tz_name: str,
) -> bool:
return conn.execute(
"SELECT 1 FROM local_days WHERE local_date = ? AND tz_name = ?",
(local_date, tz_name),
).fetchone() is not None
def _local_day_bounds(
conn: sqlite3.Connection,
local_date: str,
tz_name: str,
) -> tuple[datetime, datetime] | None:
row = conn.execute(
"SELECT utc_start, utc_end FROM local_days "
"WHERE local_date = ? AND tz_name = ?",
(local_date, tz_name),
).fetchone()
if row is None:
return None
start = _parse_sample_time(row[0])
end = _parse_sample_time(row[1])
if start is None or end is None or end <= start:
return None
return start, end
def _contains_instant(
conn: sqlite3.Connection,
local_date: str,
tz_name: str,
instant: datetime,
) -> bool:
bounds = _local_day_bounds(conn, local_date, tz_name)
return bounds is not None and bounds[0] <= instant < bounds[1]
def _monitoring_period_covers(
conn: sqlite3.Connection,
start: datetime,
end: datetime,
) -> bool:
"""Require continuous monitoring before a local-day total replaces detail."""
for started_at, ended_at in conn.execute(
"SELECT started_at, ended_at FROM monitoring_periods"
):
period_start = _parse_sample_time(started_at)
period_end = _parse_sample_time(ended_at) if ended_at is not None else None
if period_start is None or period_start > start:
continue
if ended_at is None or (period_end is not None and end <= period_end):
return True
return False
def _has_local_replacement(
conn: sqlite3.Connection,
previous: tuple,
current: tuple,
start: datetime,
end: datetime,
) -> bool:
start_tz = previous[5]
end_tz = current[5]
if not end_tz:
return False
start_id, end_id = previous[0], current[0]
segment_id = current[2]
if (
start_tz
and start_tz == end_tz
and segment_id is not None
and _monitoring_period_covers(conn, start, end)
):
known_days = conn.execute(
"SELECT local_days.utc_start, local_days.utc_end "
"FROM local_days JOIN local_day_segment_totals "
" ON local_day_segment_totals.local_day_id = local_days.id "
"WHERE local_days.tz_name = ? "
" AND local_days.activity_precision IN ('measured', 'coarse') "
" AND local_days.last_sample_id >= ? "
" AND local_days.activity_intervals > 0 "
" AND local_day_segment_totals.segment_id = ? "
" AND local_day_segment_totals.activity_intervals > 0",
(end_tz, end_id, segment_id),
).fetchall()
for row in known_days:
bounds_start = _parse_sample_time(row[0])
bounds_end = _parse_sample_time(row[1])
if (
bounds_start is not None
and bounds_end is not None
and bounds_start <= start
and end < bounds_end
):
return True
evidence = conn.execute(
"SELECT bytes_written, bytes_read, reason, start_local_date, end_local_date, "
" start_tz_name, end_tz_name, started_at, ended_at "
"FROM local_day_unallocated_evidence "
"WHERE start_sample_id = ? AND end_sample_id = ?",
(start_id, end_id),
).fetchone()
if evidence is None:
return False
start_written, end_written = previous[3], current[3]
start_read, end_read = previous[4], current[4]
if None in (start_written, end_written, start_read, end_read):
return False
if end_written < start_written or end_read < start_read:
return False
# The source IDs uniquely identify this interval. Confirm its durable
# deltas and retain every recorded local-day boundary it touches.
if (end_written - start_written, end_read - start_read) != evidence[:2]:
return False
if evidence[5] != start_tz or evidence[6] != end_tz:
return False
if evidence[2] not in {
"counter_discontinuity",
"legacy_timezone_unknown",
"timezone_change",
"monitoring_period",
"local_midnight",
"segment_unknown",
}:
return False
if _parse_sample_time(evidence[7]) != start or _parse_sample_time(evidence[8]) != end:
return False
affected_days: set[tuple[str, str]] = set()
if start_tz:
affected_days.add((evidence[3], start_tz))
affected_days.add((evidence[4], end_tz))
if start_tz and start_tz == end_tz and evidence[3] < evidence[4]:
day = date.fromisoformat(evidence[3]) + timedelta(days=1)
last_day = date.fromisoformat(evidence[4])
while day < last_day:
affected_days.add((day.isoformat(), end_tz))
day += timedelta(days=1)
if not all(
_local_day_row_exists(conn, local_date, zone)
for local_date, zone in affected_days
):
return False
if start_tz and not _contains_instant(conn, evidence[3], start_tz, start):
return False
return _contains_instant(conn, evidence[4], end_tz, end)
def _interval_has_replacement(
conn: sqlite3.Connection,
previous: tuple,
current: tuple,
) -> bool:
start = _parse_sample_time(previous[1])
end = _parse_sample_time(current[1])
if start is None or end is None or end <= start:
return False
try:
return (
_has_utc_replacement(conn, previous, current, start, end)
and _has_local_replacement(conn, previous, current, start, end)
)
except (sqlite3.Error, ValueError, OverflowError, TypeError):
# Missing or unreadable derived evidence never authorizes deletion.
return False
def _sample_needs_preservation(
conn: sqlite3.Connection,
rows: list[tuple],
index: int,
newest_id: int,
) -> bool:
sample = rows[index]
segment_id = sample[2]
if segment_id is None or conn.execute(
"SELECT 1 FROM controller_segments WHERE id = ?",
(segment_id,),
).fetchone() is None:
return True
if sample[0] == newest_id:
# Keep newest sample as the source anchor for the next collection.
return True
pairs = []
if index > 0 and rows[index - 1][2] == segment_id:
pairs.append((rows[index - 1], sample))
if index + 1 < len(rows) and rows[index + 1][2] == segment_id:
pairs.append((sample, rows[index + 1]))
return any(
not _interval_has_replacement(conn, previous, current)
for previous, current in pairs
)
def needs_boundary_anchor(
conn: sqlite3.Connection,
sample_ts: str,
now: datetime,
) -> bool:
"""Check if a sample is needed as a boundary anchor for derivation.
A sample is a boundary anchor if:
1. It's older than retention_days (strictly before cutoff)
2. It has a next sample that forms an interval spanning the retention boundary
3. The interval hasn't been derived yet
The interval spans the boundary if:
- The sample is before the cutoff, AND
- The next sample is strictly after the cutoff (or within retention)
"""
from .derive import _parse_ts
sample_dt = _parse_ts(sample_ts)
retention_cutoff = now - timedelta(days=RAW_SAMPLE_RETENTION_DAYS)
# If sample is within retention (strictly after cutoff), not an anchor
if sample_dt > retention_cutoff:
"""Return whether an old sample still carries unreplaced evidence."""
sample_time = _parse_sample_time(sample_ts)
cutoff = _utc_datetime(now) - timedelta(days=RAW_SAMPLE_RETENTION_DAYS)
if sample_time is None:
return True
if sample_time >= cutoff:
return False
# Check if this sample has a next sample
cursor = conn.execute(
"""SELECT ts, segment_id FROM samples WHERE ts > ? ORDER BY ts LIMIT 1""",
(sample_ts,),
rows = _sample_rows(conn)
matching = [index for index, row in enumerate(rows) if row[1] == sample_ts]
if not matching:
return False
newest_id = rows[-1][0] if rows else -1
return any(
_sample_needs_preservation(conn, rows, index, newest_id)
for index in matching
)
next_row = cursor.fetchone()
if next_row is None:
# No next sample - this is the last sample
# It's not needed for derivation (no interval to derive)
return False
next_ts_str = next_row[0]
next_segment_id = next_row[1]
next_dt = _parse_ts(next_ts_str)
# Check if the next sample is strictly after the cutoff (i.e., interval spans boundary)
if next_dt > retention_cutoff:
# The interval spans the retention boundary
# Check if the interval needs derivation
# Get current sample's segment_id
cursor = conn.execute(
"SELECT segment_id FROM samples WHERE ts = ?",
(sample_ts,),
)
current_segment_row = cursor.fetchone()
current_segment_id = current_segment_row[0] if current_segment_row else None
# If different segments, no interval to derive
if current_segment_id != next_segment_id:
return False
# Check if the interval [sample_ts, next_ts] needs derivation
# It needs derivation if any hour in the span lacks an observation
current_hour = sample_dt.replace(minute=0, second=0, microsecond=0)
end_hour = next_dt.replace(minute=0, second=0, microsecond=0)
while current_hour <= end_hour:
cursor = conn.execute(
"SELECT id FROM hour_observations WHERE hour = ?",
(current_hour.isoformat(),),
)
if cursor.fetchone() is None:
# This hour lacks an observation - interval needs derivation
return True
current_hour += timedelta(hours=1)
# All hours in the span have observations - interval is derived
return False
else:
# The interval doesn't span the boundary (both samples are old)
# Check if the interval needs derivation
# Get current sample's segment_id
cursor = conn.execute(
"SELECT segment_id FROM samples WHERE ts = ?",
(sample_ts,),
)
current_segment_row = cursor.fetchone()
current_segment_id = current_segment_row[0] if current_segment_row else None
# If different segments, no interval to derive
if current_segment_id != next_segment_id:
return False
# Check if the interval [sample_ts, next_ts] needs derivation
current_hour = sample_dt.replace(minute=0, second=0, microsecond=0)
end_hour = next_dt.replace(minute=0, second=0, microsecond=0)
while current_hour <= end_hour:
cursor = conn.execute(
"SELECT id FROM hour_observations WHERE hour = ?",
(current_hour.isoformat(),),
)
if cursor.fetchone() is None:
# This hour lacks an observation - interval needs derivation
# But only keep if the interval is significant (spans multiple hours)
# or if the next sample is the last sample before a gap
gap = (next_dt - sample_dt).total_seconds()
if gap > 24 * 3600: # Significant gap (> 24 hours)
return True
current_hour += timedelta(hours=1)
# All hours in the span have observations or gap is not significant
return False
def prune_old_samples(
@@ -133,36 +334,50 @@ def prune_old_samples(
now: datetime,
retention_days: int = RAW_SAMPLE_RETENTION_DAYS,
) -> int:
"""Remove raw samples older than retention_days.
"""Delete expired samples only after every dependent fact is durable.
Retains boundary anchors required for successor evidence.
Args:
conn: Connection to the observation store.
now: Current UTC time.
retention_days: Number of days to retain (default 14).
Returns:
Number of samples removed.
Pending publications are stored separately and never selected here. A
sample stays when it is the newest successor anchor, has an unpublishable
neighbour interval, lacks UTC or local-day replacement evidence, or has a
timestamp that cannot be safely interpreted.
"""
cutoff = (now - timedelta(days=retention_days)).isoformat()
# Get all samples older than cutoff
cursor = conn.execute(
"SELECT id, ts FROM samples WHERE ts < ? ORDER BY ts",
(cutoff,),
)
old_samples = cursor.fetchall()
removed = 0
for sample_id, sample_ts in old_samples:
# Check if this sample is a boundary anchor
if needs_boundary_anchor(conn, sample_ts, now):
continue # Skip - it's a boundary anchor
# Remove the sample
conn.execute("DELETE FROM samples WHERE id = ?", (sample_id,))
removed += 1
conn.commit()
return removed
cutoff = _utc_datetime(now) - timedelta(days=retention_days)
owns_transaction = not conn.in_transaction
savepoint = "fenris_sample_retention"
if owns_transaction:
conn.execute("BEGIN IMMEDIATE")
else:
conn.execute(f"SAVEPOINT {savepoint}")
try:
rows = _sample_rows(conn)
if not rows:
if owns_transaction:
conn.commit()
else:
conn.execute(f"RELEASE SAVEPOINT {savepoint}")
return 0
newest_id = rows[-1][0]
expired = []
for index, row in enumerate(rows):
sample_time = _parse_sample_time(row[1])
if sample_time is None or sample_time >= cutoff:
continue
if _sample_needs_preservation(conn, rows, index, newest_id):
continue
expired.append(row[0])
conn.executemany("DELETE FROM samples WHERE id = ?", ((sample_id,) for sample_id in expired))
if owns_transaction:
conn.commit()
else:
conn.execute(f"RELEASE SAVEPOINT {savepoint}")
return len(expired)
except Exception:
if owns_transaction:
conn.rollback()
else:
conn.execute(f"ROLLBACK TO SAVEPOINT {savepoint}")
conn.execute(f"RELEASE SAVEPOINT {savepoint}")
raise
+341
View File
@@ -0,0 +1,341 @@
"""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
+4 -4
View File
@@ -151,8 +151,8 @@ class TestDetailExpiresSummariesRemain:
# Run pruning
pruned = prune_old_samples(conn, now, retention_days=14)
# Old samples should be pruned, recent ones retained
assert pruned >= 1
# The handcrafted totals lack replacement provenance, so retain samples.
assert pruned == 0
# Local-day summaries must still be queryable
old_result = query_local_day_summary(conn, "2026-09-16")
@@ -434,10 +434,10 @@ class TestExistingEntryPoints:
_insert_sample(conn, ts, bw=days_ago * 100)
pruned = prune_old_samples(conn, now, retention_days=14)
assert pruned == 15
assert pruned == 0
cursor = conn.execute("SELECT COUNT(*) FROM samples")
assert cursor.fetchone()[0] == 14
assert cursor.fetchone()[0] == 29
conn.close()
def test_repair_with_real_temp_store(self, tmp_path):
+17 -15
View File
@@ -27,9 +27,14 @@ def store_conn(tmp_path: Path):
def _insert_sample(conn, ts_iso, device="/dev/nvme0"):
conn.execute(
"INSERT OR IGNORE INTO controller_segments "
"(id, opened_at, identity_key) VALUES (1, '2026-01-01T00:00:00+00:00', 'test')"
)
conn.execute(
"INSERT INTO samples (ts, device, data_units_written, data_units_read, "
" bytes_written, bytes_read, percentage_used) VALUES (?, ?, 0, 0, 0, 0, 0)",
" bytes_written, bytes_read, percentage_used, segment_id, local_tz) "
"VALUES (?, ?, 0, 0, 0, 0, 0, 1, 'UTC')",
(ts_iso, device),
)
conn.commit()
@@ -50,33 +55,30 @@ class TestPruneOldSamples:
cursor = store_conn.execute("SELECT COUNT(*) FROM samples")
assert cursor.fetchone()[0] == 1
def test_removes_old_samples(self, store_conn):
def test_keeps_old_samples_without_replacement_evidence(self, store_conn):
now = datetime(2026, 9, 15, 12, 0, 0, tzinfo=timezone.utc)
# Insert samples at 10, 14, and 15 days ago
# All three are before the cutoff (2026-09-01T12:00:00)
# The 15-day-old sample is not a boundary anchor because
# the next sample (14 days ago) is also before the cutoff
# Old samples stay until UTC and local-day evidence replace them.
for days_ago in [10, 14, 15]:
ts = (now - timedelta(days=days_ago)).isoformat()
_insert_sample(store_conn, ts)
pruned = prune_old_samples(store_conn, now, retention_days=14)
assert pruned == 1 # Only the 15-day-old sample removed
assert pruned == 0
cursor = store_conn.execute("SELECT COUNT(*) FROM samples")
assert cursor.fetchone()[0] == 2
assert cursor.fetchone()[0] == 3
def test_removes_many_old_samples(self, store_conn):
def test_keeps_many_old_samples_without_replacement_evidence(self, store_conn):
now = datetime(2026, 9, 15, 12, 0, 0, tzinfo=timezone.utc)
for days_ago in range(1, 30):
ts = (now - timedelta(days=days_ago)).isoformat()
_insert_sample(store_conn, ts)
pruned = prune_old_samples(store_conn, now, retention_days=14)
assert pruned == 15 # Days 15-29 removed
assert pruned == 0
cursor = store_conn.execute("SELECT COUNT(*) FROM samples")
assert cursor.fetchone()[0] == 14 # Days 1-14 kept
assert cursor.fetchone()[0] == 29
def test_empty_store_no_error(self, store_conn):
now = datetime(2026, 9, 15, 12, 0, 0, tzinfo=timezone.utc)
@@ -130,8 +132,8 @@ class TestPruneOldSamples:
)
assert cursor.fetchone()[0] == 1
def test_old_sample_with_derived_interval_removed(self, store_conn):
"""Old samples with fully derived intervals are removed."""
def test_hour_rows_alone_do_not_replace_local_day_evidence(self, store_conn):
"""An hour row alone cannot authorize source-sample deletion."""
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
# Old sample with derived interval
@@ -148,8 +150,8 @@ class TestPruneOldSamples:
# Run pruning
pruned = prune_old_samples(store_conn, now, retention_days=14)
# Old sample should be removed (interval is derived)
# Keep the source samples until local-day replacement evidence exists.
cursor = store_conn.execute(
"SELECT COUNT(*) FROM samples WHERE ts = '2026-09-10T10:00:00+00:00'"
)
assert cursor.fetchone()[0] == 0
assert cursor.fetchone()[0] == 1
+7 -8
View File
@@ -235,8 +235,8 @@ class TestBoundaryAnchorRetention:
# the interval spans the retention boundary
assert needs_boundary_anchor(store_conn, "2026-09-15T23:55:00+00:00", now)
def test_not_boundary_anchor_if_fully_derived(self, store_conn):
"""Sample that's fully derived is not a boundary anchor."""
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
@@ -247,8 +247,7 @@ class TestBoundaryAnchorRetention:
_insert_hour(store_conn, "2026-09-14T10:00:00+00:00",
bytes_written_delta=1000000)
# Not a boundary anchor
assert not needs_boundary_anchor(store_conn, "2026-09-14T10:00:00+00:00", now)
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."""
@@ -274,8 +273,8 @@ class TestBoundaryAnchorRetention:
)
assert cursor.fetchone()[0] == 1
def test_pruning_removes_old_sample_with_derived_interval(self, store_conn):
"""Pruning removes old samples when interval is fully derived."""
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
@@ -289,11 +288,11 @@ class TestBoundaryAnchorRetention:
# Run pruning
pruned = prune_old_samples(store_conn, now, retention_days=14)
# Old sample should be removed (interval is derived)
# 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] == 0
assert cursor.fetchone()[0] == 1
# ---------------------------------------------------------------------------