Compare commits
2
Commits
017a562566
...
5aeee7b6f1
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5aeee7b6f1 | ||
|
|
d69753690a |
@@ -122,6 +122,11 @@ def main() -> None:
|
|||||||
|
|
||||||
if result["ok"]:
|
if result["ok"]:
|
||||||
print(f"Collection successful: {result['sample_count']} sample(s)")
|
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)
|
sys.exit(0)
|
||||||
else:
|
else:
|
||||||
print(f"Collection failed: {result['error']}", file=sys.stderr)
|
print(f"Collection failed: {result['error']}", file=sys.stderr)
|
||||||
|
|||||||
@@ -477,9 +477,21 @@ def run_collection(
|
|||||||
break
|
break
|
||||||
|
|
||||||
published_count += _recover_pending_for_collection(conn)
|
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 {
|
return {
|
||||||
"ok": True,
|
"ok": True,
|
||||||
"sample_count": published_count,
|
"sample_count": published_count,
|
||||||
|
"retention_error": retention_error,
|
||||||
"store_path": str(store_path),
|
"store_path": str(store_path),
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+360
-145
@@ -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.
|
The 14-day cutoff never overrides local-day preservation, shared boundary
|
||||||
Hour observations and day aggregates are retained indefinitely.
|
evidence, publication recovery, or the newest sample's successor-anchor role.
|
||||||
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).
|
|
||||||
"""
|
"""
|
||||||
import sqlite3
|
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
|
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(
|
def needs_boundary_anchor(
|
||||||
conn: sqlite3.Connection,
|
conn: sqlite3.Connection,
|
||||||
sample_ts: str,
|
sample_ts: str,
|
||||||
now: datetime,
|
now: datetime,
|
||||||
) -> bool:
|
) -> bool:
|
||||||
"""Check if a sample is needed as a boundary anchor for derivation.
|
"""Return whether an old sample still carries unreplaced evidence."""
|
||||||
|
sample_time = _parse_sample_time(sample_ts)
|
||||||
A sample is a boundary anchor if:
|
cutoff = _utc_datetime(now) - timedelta(days=RAW_SAMPLE_RETENTION_DAYS)
|
||||||
1. It's older than retention_days (strictly before cutoff)
|
if sample_time is None:
|
||||||
2. It has a next sample that forms an interval spanning the retention boundary
|
return True
|
||||||
3. The interval hasn't been derived yet
|
if sample_time >= cutoff:
|
||||||
|
|
||||||
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 False
|
return False
|
||||||
|
|
||||||
# Check if this sample has a next sample
|
rows = _sample_rows(conn)
|
||||||
cursor = conn.execute(
|
matching = [index for index, row in enumerate(rows) if row[1] == sample_ts]
|
||||||
"""SELECT ts, segment_id FROM samples WHERE ts > ? ORDER BY ts LIMIT 1""",
|
if not matching:
|
||||||
(sample_ts,),
|
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(
|
def prune_old_samples(
|
||||||
@@ -133,36 +334,50 @@ def prune_old_samples(
|
|||||||
now: datetime,
|
now: datetime,
|
||||||
retention_days: int = RAW_SAMPLE_RETENTION_DAYS,
|
retention_days: int = RAW_SAMPLE_RETENTION_DAYS,
|
||||||
) -> int:
|
) -> 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.
|
Pending publications are stored separately and never selected here. A
|
||||||
|
sample stays when it is the newest successor anchor, has an unpublishable
|
||||||
Args:
|
neighbour interval, lacks UTC or local-day replacement evidence, or has a
|
||||||
conn: Connection to the observation store.
|
timestamp that cannot be safely interpreted.
|
||||||
now: Current UTC time.
|
|
||||||
retention_days: Number of days to retain (default 14).
|
|
||||||
|
|
||||||
Returns:
|
|
||||||
Number of samples removed.
|
|
||||||
"""
|
"""
|
||||||
cutoff = (now - timedelta(days=retention_days)).isoformat()
|
cutoff = _utc_datetime(now) - timedelta(days=retention_days)
|
||||||
|
owns_transaction = not conn.in_transaction
|
||||||
# Get all samples older than cutoff
|
savepoint = "fenris_sample_retention"
|
||||||
cursor = conn.execute(
|
if owns_transaction:
|
||||||
"SELECT id, ts FROM samples WHERE ts < ? ORDER BY ts",
|
conn.execute("BEGIN IMMEDIATE")
|
||||||
(cutoff,),
|
else:
|
||||||
)
|
conn.execute(f"SAVEPOINT {savepoint}")
|
||||||
old_samples = cursor.fetchall()
|
|
||||||
|
try:
|
||||||
removed = 0
|
rows = _sample_rows(conn)
|
||||||
for sample_id, sample_ts in old_samples:
|
if not rows:
|
||||||
# Check if this sample is a boundary anchor
|
if owns_transaction:
|
||||||
if needs_boundary_anchor(conn, sample_ts, now):
|
conn.commit()
|
||||||
continue # Skip - it's a boundary anchor
|
else:
|
||||||
|
conn.execute(f"RELEASE SAVEPOINT {savepoint}")
|
||||||
# Remove the sample
|
return 0
|
||||||
conn.execute("DELETE FROM samples WHERE id = ?", (sample_id,))
|
|
||||||
removed += 1
|
newest_id = rows[-1][0]
|
||||||
|
expired = []
|
||||||
conn.commit()
|
for index, row in enumerate(rows):
|
||||||
return removed
|
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
|
||||||
|
|||||||
@@ -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
|
||||||
@@ -151,8 +151,8 @@ class TestDetailExpiresSummariesRemain:
|
|||||||
# Run pruning
|
# Run pruning
|
||||||
pruned = prune_old_samples(conn, now, retention_days=14)
|
pruned = prune_old_samples(conn, now, retention_days=14)
|
||||||
|
|
||||||
# Old samples should be pruned, recent ones retained
|
# The handcrafted totals lack replacement provenance, so retain samples.
|
||||||
assert pruned >= 1
|
assert pruned == 0
|
||||||
|
|
||||||
# Local-day summaries must still be queryable
|
# Local-day summaries must still be queryable
|
||||||
old_result = query_local_day_summary(conn, "2026-09-16")
|
old_result = query_local_day_summary(conn, "2026-09-16")
|
||||||
@@ -434,10 +434,10 @@ class TestExistingEntryPoints:
|
|||||||
_insert_sample(conn, ts, bw=days_ago * 100)
|
_insert_sample(conn, ts, bw=days_ago * 100)
|
||||||
|
|
||||||
pruned = prune_old_samples(conn, now, retention_days=14)
|
pruned = prune_old_samples(conn, now, retention_days=14)
|
||||||
assert pruned == 15
|
assert pruned == 0
|
||||||
|
|
||||||
cursor = conn.execute("SELECT COUNT(*) FROM samples")
|
cursor = conn.execute("SELECT COUNT(*) FROM samples")
|
||||||
assert cursor.fetchone()[0] == 14
|
assert cursor.fetchone()[0] == 29
|
||||||
conn.close()
|
conn.close()
|
||||||
|
|
||||||
def test_repair_with_real_temp_store(self, tmp_path):
|
def test_repair_with_real_temp_store(self, tmp_path):
|
||||||
|
|||||||
+17
-15
@@ -27,9 +27,14 @@ def store_conn(tmp_path: Path):
|
|||||||
|
|
||||||
|
|
||||||
def _insert_sample(conn, ts_iso, device="/dev/nvme0"):
|
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(
|
conn.execute(
|
||||||
"INSERT INTO samples (ts, device, data_units_written, data_units_read, "
|
"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),
|
(ts_iso, device),
|
||||||
)
|
)
|
||||||
conn.commit()
|
conn.commit()
|
||||||
@@ -50,33 +55,30 @@ class TestPruneOldSamples:
|
|||||||
cursor = store_conn.execute("SELECT COUNT(*) FROM samples")
|
cursor = store_conn.execute("SELECT COUNT(*) FROM samples")
|
||||||
assert cursor.fetchone()[0] == 1
|
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)
|
now = datetime(2026, 9, 15, 12, 0, 0, tzinfo=timezone.utc)
|
||||||
# Insert samples at 10, 14, and 15 days ago
|
# Old samples stay until UTC and local-day evidence replace them.
|
||||||
# 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
|
|
||||||
for days_ago in [10, 14, 15]:
|
for days_ago in [10, 14, 15]:
|
||||||
ts = (now - timedelta(days=days_ago)).isoformat()
|
ts = (now - timedelta(days=days_ago)).isoformat()
|
||||||
_insert_sample(store_conn, ts)
|
_insert_sample(store_conn, ts)
|
||||||
|
|
||||||
pruned = prune_old_samples(store_conn, now, retention_days=14)
|
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")
|
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)
|
now = datetime(2026, 9, 15, 12, 0, 0, tzinfo=timezone.utc)
|
||||||
for days_ago in range(1, 30):
|
for days_ago in range(1, 30):
|
||||||
ts = (now - timedelta(days=days_ago)).isoformat()
|
ts = (now - timedelta(days=days_ago)).isoformat()
|
||||||
_insert_sample(store_conn, ts)
|
_insert_sample(store_conn, ts)
|
||||||
|
|
||||||
pruned = prune_old_samples(store_conn, now, retention_days=14)
|
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")
|
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):
|
def test_empty_store_no_error(self, store_conn):
|
||||||
now = datetime(2026, 9, 15, 12, 0, 0, tzinfo=timezone.utc)
|
now = datetime(2026, 9, 15, 12, 0, 0, tzinfo=timezone.utc)
|
||||||
@@ -130,8 +132,8 @@ class TestPruneOldSamples:
|
|||||||
)
|
)
|
||||||
assert cursor.fetchone()[0] == 1
|
assert cursor.fetchone()[0] == 1
|
||||||
|
|
||||||
def test_old_sample_with_derived_interval_removed(self, store_conn):
|
def test_hour_rows_alone_do_not_replace_local_day_evidence(self, store_conn):
|
||||||
"""Old samples with fully derived intervals are removed."""
|
"""An hour row alone cannot authorize source-sample deletion."""
|
||||||
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
|
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
|
||||||
|
|
||||||
# Old sample with derived interval
|
# Old sample with derived interval
|
||||||
@@ -148,8 +150,8 @@ class TestPruneOldSamples:
|
|||||||
# Run pruning
|
# Run pruning
|
||||||
pruned = prune_old_samples(store_conn, now, retention_days=14)
|
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(
|
cursor = store_conn.execute(
|
||||||
"SELECT COUNT(*) FROM samples WHERE ts = '2026-09-10T10:00:00+00:00'"
|
"SELECT COUNT(*) FROM samples WHERE ts = '2026-09-10T10:00:00+00:00'"
|
||||||
)
|
)
|
||||||
assert cursor.fetchone()[0] == 0
|
assert cursor.fetchone()[0] == 1
|
||||||
|
|||||||
@@ -235,8 +235,8 @@ class TestBoundaryAnchorRetention:
|
|||||||
# the interval spans the retention boundary
|
# the interval spans the retention boundary
|
||||||
assert needs_boundary_anchor(store_conn, "2026-09-15T23:55:00+00:00", now)
|
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):
|
def test_hour_only_derivation_does_not_replace_local_day_evidence(self, store_conn):
|
||||||
"""Sample that's fully derived is not a boundary anchor."""
|
"""UTC hour evidence alone does not make a sample safe to prune."""
|
||||||
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
|
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
|
||||||
|
|
||||||
# Insert sample and fully derive its interval
|
# Insert sample and fully derive its interval
|
||||||
@@ -247,8 +247,7 @@ class TestBoundaryAnchorRetention:
|
|||||||
_insert_hour(store_conn, "2026-09-14T10:00:00+00:00",
|
_insert_hour(store_conn, "2026-09-14T10:00:00+00:00",
|
||||||
bytes_written_delta=1000000)
|
bytes_written_delta=1000000)
|
||||||
|
|
||||||
# Not a boundary anchor
|
assert needs_boundary_anchor(store_conn, "2026-09-14T10:00:00+00:00", now)
|
||||||
assert not needs_boundary_anchor(store_conn, "2026-09-14T10:00:00+00:00", now)
|
|
||||||
|
|
||||||
def test_pruning_retains_boundary_anchors(self, store_conn):
|
def test_pruning_retains_boundary_anchors(self, store_conn):
|
||||||
"""Pruning keeps samples needed as boundary anchors."""
|
"""Pruning keeps samples needed as boundary anchors."""
|
||||||
@@ -274,8 +273,8 @@ class TestBoundaryAnchorRetention:
|
|||||||
)
|
)
|
||||||
assert cursor.fetchone()[0] == 1
|
assert cursor.fetchone()[0] == 1
|
||||||
|
|
||||||
def test_pruning_removes_old_sample_with_derived_interval(self, store_conn):
|
def test_pruning_keeps_samples_without_local_day_replacement(self, store_conn):
|
||||||
"""Pruning removes old samples when interval is fully derived."""
|
"""Pruning keeps samples when only the UTC-hour row exists."""
|
||||||
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
|
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
|
||||||
|
|
||||||
# Old sample with derived interval
|
# Old sample with derived interval
|
||||||
@@ -289,11 +288,11 @@ class TestBoundaryAnchorRetention:
|
|||||||
# Run pruning
|
# Run pruning
|
||||||
pruned = prune_old_samples(store_conn, now, retention_days=14)
|
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(
|
cursor = store_conn.execute(
|
||||||
"SELECT COUNT(*) FROM samples WHERE ts = '2026-09-10T10:00:00+00:00'"
|
"SELECT COUNT(*) FROM samples WHERE ts = '2026-09-10T10:00:00+00:00'"
|
||||||
)
|
)
|
||||||
assert cursor.fetchone()[0] == 0
|
assert cursor.fetchone()[0] == 1
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|||||||
Reference in New Issue
Block a user