diff --git a/src/fenris/collect.py b/src/fenris/collect.py index b2a23a7..31197b2 100644 --- a/src/fenris/collect.py +++ b/src/fenris/collect.py @@ -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) diff --git a/src/fenris/collector.py b/src/fenris/collector.py index 5fc59f3..460b5d5 100644 --- a/src/fenris/collector.py +++ b/src/fenris/collector.py @@ -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), } diff --git a/src/fenris/pruning.py b/src/fenris/pruning.py index ec42542..1e5cbde 100644 --- a/src/fenris/pruning.py +++ b/src/fenris/pruning.py @@ -1,131 +1,309 @@ -"""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 _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: + 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 +311,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 diff --git a/tests/test_issue_100.py b/tests/test_issue_100.py new file mode 100644 index 0000000..18f06eb --- /dev/null +++ b/tests/test_issue_100.py @@ -0,0 +1,279 @@ +"""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_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 diff --git a/tests/test_issue_93.py b/tests/test_issue_93.py index 105a66c..2629e03 100644 --- a/tests/test_issue_93.py +++ b/tests/test_issue_93.py @@ -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): diff --git a/tests/test_pruning.py b/tests/test_pruning.py index fbe0d3c..03af590 100644 --- a/tests/test_pruning.py +++ b/tests/test_pruning.py @@ -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 diff --git a/tests/test_repair.py b/tests/test_repair.py index 248ce60..4dbbd89 100644 --- a/tests/test_repair.py +++ b/tests/test_repair.py @@ -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 # ---------------------------------------------------------------------------