feat: implement repair and retention for observation history (issue #74)

This commit is contained in:
xavierk
2026-09-14 03:54:46 +05:30
parent d790ff84c5
commit 6917a658cb
4 changed files with 1033 additions and 5 deletions
+448
View File
@@ -0,0 +1,448 @@
"""Historical repair and retention (issue #74).
Re-derives hour observations and day aggregates from surviving raw samples,
with idempotent and interruption-safe guarantees. Preserves import markers,
boundary anchors, and valid historical summaries.
Contracts:
- Idempotent: running repair multiple times produces no duplicates
- Interruption-safe: partial repair preserves prior valid history
- Preserves existing valid data: never overwrites valid derived data
- Surfaces failures explicitly for retry
"""
import logging
import sqlite3
from dataclasses import dataclass, field
from datetime import datetime, timedelta, timezone
from typing import Optional, List, Tuple
from .derive import find_previous_sample, derive_hours_from_interval, _parse_ts
logger = logging.getLogger(__name__)
@dataclass
class RepairResult:
"""Result of a repair operation."""
ok: bool
hours_created: int = 0
hours_updated: int = 0
days_created: int = 0
days_updated: int = 0
intervals_derived: int = 0
boundary_anchors_retained: int = 0
legacy_summaries_preserved: int = 0
error: Optional[str] = None
@dataclass
class RepairStatus:
"""Current repair status for read-only views."""
last_repair: Optional[str] = None # ISO timestamp of last successful repair
repair_in_progress: bool = False
hours_derived: int = 0
days_derived: int = 0
def _ensure_repair_metadata(conn: sqlite3.Connection) -> None:
"""Ensure metadata table exists for tracking repair state."""
conn.execute("""
CREATE TABLE IF NOT EXISTS store_metadata (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
)
""")
def is_repair_in_progress(conn: sqlite3.Connection) -> bool:
"""Check if a repair operation is currently in progress."""
_ensure_repair_metadata(conn)
cursor = conn.execute(
"SELECT value FROM store_metadata WHERE key = 'repair_in_progress'"
)
row = cursor.fetchone()
return row is not None and row[0] == "true"
def _set_repair_in_progress(conn: sqlite3.Connection, in_progress: bool) -> None:
"""Mark repair as in progress or complete."""
_ensure_repair_metadata(conn)
conn.execute(
"INSERT OR REPLACE INTO store_metadata (key, value) VALUES (?, ?)",
("repair_in_progress", "true" if in_progress else "false"),
)
conn.commit()
def _update_repair_status(conn: sqlite3.Connection, result: RepairResult) -> None:
"""Update repair status after successful completion."""
_ensure_repair_metadata(conn)
now = datetime.now(timezone.utc).isoformat()
# Update last repair timestamp
conn.execute(
"INSERT OR REPLACE INTO store_metadata (key, value) VALUES (?, ?)",
("last_repair", now),
)
# Update derived counts
cursor = conn.execute("SELECT COUNT(*) FROM hour_observations")
hours = cursor.fetchone()[0]
cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates")
days = cursor.fetchone()[0]
conn.execute(
"INSERT OR REPLACE INTO store_metadata (key, value) VALUES (?, ?)",
("hours_derived", str(hours)),
)
conn.execute(
"INSERT OR REPLACE INTO store_metadata (key, value) VALUES (?, ?)",
("days_derived", str(days)),
)
conn.commit()
def get_repair_status(conn: sqlite3.Connection) -> RepairStatus:
"""Get current repair status for read-only views."""
_ensure_repair_metadata(conn)
last_repair = None
cursor = conn.execute(
"SELECT value FROM store_metadata WHERE key = 'last_repair'"
)
row = cursor.fetchone()
if row:
last_repair = row[0]
in_progress = is_repair_in_progress(conn)
hours_derived = 0
cursor = conn.execute(
"SELECT value FROM store_metadata WHERE key = 'hours_derived'"
)
row = cursor.fetchone()
if row:
hours_derived = int(row[0])
days_derived = 0
cursor = conn.execute(
"SELECT value FROM store_metadata WHERE key = 'days_derived'"
)
row = cursor.fetchone()
if row:
days_derived = int(row[0])
return RepairStatus(
last_repair=last_repair,
repair_in_progress=in_progress,
hours_derived=hours_derived,
days_derived=days_derived,
)
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
2. It has no derived hour observation for its hour
3. It's the last sample before a gap that needs derivation (gap > 24 hours)
"""
from .pruning import RAW_SAMPLE_RETENTION_DAYS
sample_dt = _parse_ts(sample_ts)
retention_cutoff = now - timedelta(days=RAW_SAMPLE_RETENTION_DAYS)
# If sample is within retention, not an anchor (will be kept anyway)
if sample_dt >= retention_cutoff:
return False
# Check if this sample's hour already has a derived observation
hour_start = sample_dt.replace(minute=0, second=0, microsecond=0).isoformat()
hour_end = (sample_dt + timedelta(hours=1)).replace(minute=0, second=0, microsecond=0).isoformat()
cursor = conn.execute(
"""SELECT COUNT(*) FROM hour_observations
WHERE hour >= ? AND hour < ?""",
(hour_start, hour_end),
)
# If there's an hour observation in this sample's hour, it's been derived
if cursor.fetchone()[0] > 0:
return False
# Check if this sample is the last sample before a gap
# (i.e., the next sample is significantly later)
cursor = conn.execute(
"""SELECT ts 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, might be needed
# But if it's old and fully derived, it's not needed
return False
next_ts = _parse_ts(next_row[0])
gap = (next_ts - sample_dt).total_seconds()
# If gap > 24 hours, this sample is a boundary anchor
# (needed to derive the interval spanning the gap)
return gap > 24 * 3600
def _get_unlinked_intervals(conn: sqlite3.Connection) -> List[Tuple[dict, dict]]:
"""Find sample pairs that form intervals but have no hour observations."""
cursor = conn.execute(
"""SELECT id, ts, bytes_written, bytes_read, power_on_hours,
temperature_c, data_units_written, data_units_read, segment_id
FROM samples ORDER BY ts"""
)
all_samples = []
for row in cursor.fetchall():
all_samples.append({
"id": row[0], "ts": row[1], "bytes_written": row[2],
"bytes_read": row[3], "power_on_hours": row[4],
"temperature_c": row[5], "data_units_written": row[6],
"data_units_read": row[7], "segment_id": row[8],
})
intervals = []
for i in range(len(all_samples) - 1):
prev = all_samples[i]
next_s = all_samples[i + 1]
# Skip if different segments
if prev["segment_id"] != next_s["segment_id"]:
continue
# Check if the interval spans hours that need derivation
prev_dt = _parse_ts(prev["ts"])
next_dt = _parse_ts(next_s["ts"])
# Check if any hour in the span lacks an observation
current = prev_dt.replace(minute=0, second=0, microsecond=0)
end = next_dt.replace(minute=0, second=0, microsecond=0)
needs_derivation = False
while current <= end:
cursor2 = conn.execute(
"SELECT id FROM hour_observations WHERE hour = ?",
(current.isoformat(),),
)
if cursor2.fetchone() is None:
needs_derivation = True
break
current += timedelta(hours=1)
if needs_derivation:
intervals.append((prev, next_s))
return intervals
def _derive_day_aggregate_from_hours(
conn: sqlite3.Connection,
day: str,
) -> Optional[dict]:
"""Derive a day aggregate from its hour observations."""
cursor = conn.execute(
"""SELECT SUM(active_seconds), SUM(idle_seconds),
SUM(powered_off_seconds), SUM(unknown_seconds),
SUM(bytes_written_delta), SUM(bytes_read_delta),
SUM(sample_count)
FROM hour_observations WHERE hour LIKE ?""",
(day + "T%",),
)
row = cursor.fetchone()
if row is None or row[0] is None:
return None
return {
"day": day,
"active_seconds": row[0] or 0,
"idle_seconds": row[1] or 0,
"powered_off_seconds": row[2] or 0,
"unknown_seconds": row[3] or 0,
"bytes_written_delta": row[4] or 0,
"bytes_read_delta": row[5] or 0,
"sample_count": row[6] or 0,
}
def _upsert_day_aggregate(conn: sqlite3.Connection, day_data: dict) -> bool:
"""Insert or update a day aggregate. Returns True if created."""
existing = conn.execute(
"SELECT id FROM day_aggregates WHERE day = ?",
(day_data["day"],),
).fetchone()
if existing is None:
# Calculate coverage
total_seconds = (day_data["active_seconds"] + day_data["idle_seconds"] +
day_data["powered_off_seconds"] + day_data["unknown_seconds"])
coverage = (day_data["active_seconds"] + day_data["idle_seconds"] +
day_data["powered_off_seconds"]) / total_seconds if total_seconds > 0 else 0.0
conn.execute(
"""INSERT INTO day_aggregates
(day, active_seconds, idle_seconds, powered_off_seconds, unknown_seconds,
bytes_written_delta, bytes_read_delta, sample_count, coverage)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)""",
(day_data["day"], day_data["active_seconds"], day_data["idle_seconds"],
day_data["powered_off_seconds"], day_data["unknown_seconds"],
day_data["bytes_written_delta"], day_data["bytes_read_delta"],
day_data["sample_count"], coverage),
)
return True
else:
# Update existing (but only if new data is more complete)
# This implements the "don't overwrite valid older history" rule
cursor = conn.execute(
"""SELECT bytes_written_delta, sample_count
FROM day_aggregates WHERE day = ?""",
(day_data["day"],),
)
existing_data = cursor.fetchone()
# Only update if new data has more samples or more bytes
if (day_data["sample_count"] > existing_data[1] or
day_data["bytes_written_delta"] > existing_data[0]):
total_seconds = (day_data["active_seconds"] + day_data["idle_seconds"] +
day_data["powered_off_seconds"] + day_data["unknown_seconds"])
coverage = (day_data["active_seconds"] + day_data["idle_seconds"] +
day_data["powered_off_seconds"]) / total_seconds if total_seconds > 0 else 0.0
conn.execute(
"""UPDATE day_aggregates SET
active_seconds = ?, idle_seconds = ?, powered_off_seconds = ?,
unknown_seconds = ?, bytes_written_delta = ?, bytes_read_delta = ?,
sample_count = ?, coverage = ?
WHERE day = ?""",
(day_data["active_seconds"], day_data["idle_seconds"],
day_data["powered_off_seconds"], day_data["unknown_seconds"],
day_data["bytes_written_delta"], day_data["bytes_read_delta"],
day_data["sample_count"], coverage, day_data["day"]),
)
return False # Updated, not created
else:
return False # No update needed
def repair_derivation(
conn: sqlite3.Connection,
clock=None,
) -> RepairResult:
"""Repair hour observations and day aggregates from surviving samples.
This is the main entry point for historical repair. It:
1. Finds sample pairs that need interval derivation
2. Derives hour observations from those intervals
3. Updates day aggregates from the hour observations
4. Preserves existing valid data
5. Is idempotent and interruption-safe
Args:
conn: Connection to the observation store
clock: Injected clock (for testing)
Returns:
RepairResult with operation details
"""
if clock is None:
clock = datetime.now(timezone.utc)
elif hasattr(clock, 'utcnow'):
clock = clock.utcnow()
# Check if repair is already in progress
if is_repair_in_progress(conn):
return RepairResult(
ok=False,
error="Repair already in progress",
)
# Mark repair as in progress
_set_repair_in_progress(conn, True)
result = RepairResult(ok=True)
try:
# Begin transaction
conn.execute("BEGIN IMMEDIATE")
# 1. Find and derive intervals from sample pairs
intervals = _get_unlinked_intervals(conn)
for prev, next_s in intervals:
try:
derived_hours = derive_hours_from_interval(conn, prev, next_s)
result.intervals_derived += 1
result.hours_created += len([h for h in derived_hours if h.get("attributed", True)])
except Exception as e:
logger.warning("Failed to derive interval %s -> %s: %s",
prev["ts"], next_s["ts"], e)
# Continue with other intervals (resilient)
# 2. Update day aggregates from hour observations
cursor = conn.execute(
"SELECT DISTINCT substr(hour, 1, 10) as day FROM hour_observations ORDER BY day"
)
days = [row[0] for row in cursor.fetchall()]
for day in days:
day_data = _derive_day_aggregate_from_hours(conn, day)
if day_data is not None:
created = _upsert_day_aggregate(conn, day_data)
if created:
result.days_created += 1
else:
result.days_updated += 1
# 3. Count boundary anchors retained
now = clock if isinstance(clock, datetime) else datetime.now(timezone.utc)
cursor = conn.execute("SELECT ts FROM samples ORDER BY ts")
anchor_count = 0
for row in cursor.fetchall():
if needs_boundary_anchor(conn, row[0], now):
anchor_count += 1
result.boundary_anchors_retained = anchor_count
# 4. Count preserved legacy summaries
# Legacy summaries are day aggregates without corresponding hour observations
cursor = conn.execute(
"""SELECT COUNT(*) FROM day_aggregates d
WHERE NOT EXISTS (
SELECT 1 FROM hour_observations h
WHERE h.hour LIKE d.day || 'T%'
)"""
)
result.legacy_summaries_preserved = cursor.fetchone()[0]
# Commit transaction
conn.commit()
# Update repair status
_update_repair_status(conn, result)
except Exception as e:
conn.rollback()
logger.error("Repair failed: %s", e)
return RepairResult(
ok=False,
error=str(e),
)
finally:
# Mark repair as complete
_set_repair_in_progress(conn, False)
return result