fix: guard local-day pruning across monitoring gaps (#100)
This commit is contained in:
+24
-1
@@ -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 "
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user