diff --git a/src/fenris/pruning.py b/src/fenris/pruning.py index 1e5cbde..0db762e 100644 --- a/src/fenris/pruning.py +++ b/src/fenris/pruning.py @@ -143,6 +143,24 @@ def _contains_instant( 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, @@ -157,7 +175,12 @@ def _has_local_replacement( 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: + 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 " diff --git a/tests/test_issue_100.py b/tests/test_issue_100.py index 18f06eb..8314c11 100644 --- a/tests/test_issue_100.py +++ b/tests/test_issue_100.py @@ -137,6 +137,68 @@ def test_pruning_keeps_samples_when_replacement_evidence_is_missing(tmp_path): 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)