fix: prune old detail after publication (#100)
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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),
|
||||
}
|
||||
|
||||
|
||||
+337
-145
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user