"""Historical repair and retention (issue #74, #99). Re-derives hour observations and day aggregates from surviving raw samples, rebuilds legacy local-day evidence from surviving samples or trustworthy UTC hours, and preserves import markers, boundary anchors, and valid summaries. Repair is idempotent and interruption-safe. 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"), ) 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)), ) 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() result = RepairResult(ok=True) try: # A committed in-progress marker may be stale after process death. # SQLite's write lock serializes active repairs; retry idempotently. if conn.in_transaction: conn.commit() conn.execute("BEGIN IMMEDIATE") _set_repair_in_progress(conn, True) conn.commit() 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 # Rebuild only legacy or unavailable local-day activity. This shares # the migration derivation and leaves trustworthy published totals intact. from .local_day import repair_legacy_local_day_evidence repair_legacy_local_day_evidence(conn) # 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] # Publish derived rows and completion status together. A process # interruption rolls back the in-progress marker with the repair. _update_repair_status(conn, result) _set_repair_in_progress(conn, False) conn.commit() except Exception as e: conn.rollback() try: _set_repair_in_progress(conn, False) conn.commit() except sqlite3.Error: conn.rollback() logger.error("Repair failed: %s", e) return RepairResult( ok=False, error=str(e), ) return result