163 lines
5.8 KiB
Python
163 lines
5.8 KiB
Python
"""Raw sample pruning per spec §3.4, ST-5.
|
|
|
|
Raw samples are pruned opportunistically to 14 days.
|
|
Hour observations and day aggregates are retained indefinitely.
|
|
Boundary anchors required for successor evidence are retained.
|
|
"""
|
|
import sqlite3
|
|
from datetime import datetime, timedelta, timezone
|
|
|
|
# Spec §3.4: Raw-sample retention
|
|
RAW_SAMPLE_RETENTION_DAYS = 14
|
|
|
|
|
|
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 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,),
|
|
)
|
|
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(
|
|
conn: sqlite3.Connection,
|
|
now: datetime,
|
|
retention_days: int = RAW_SAMPLE_RETENTION_DAYS,
|
|
) -> int:
|
|
"""Remove raw samples older than retention_days.
|
|
|
|
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.
|
|
"""
|
|
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
|