From 017a56256641f8cb0418158d695a671ec399cd7d Mon Sep 17 00:00:00 2001 From: xavierk Date: Mon, 28 Sep 2026 13:09:00 +0530 Subject: [PATCH] fix: rebuild trustworthy local-day evidence (#99) --- docs/adr/0010-local-day-activity-history.md | 11 + src/fenris/local_day.py | 576 +++++++++++++++++++- src/fenris/projection.py | 12 +- src/fenris/repair.py | 52 +- src/fenris/store.py | 206 ++++--- tests/test_acceptance_sweep.py | 4 +- tests/test_collector_history_tracer.py | 6 +- tests/test_complete_observation_day_gate.py | 45 +- tests/test_issue_72_edge_cases.py | 4 +- tests/test_issue_77.py | 4 +- tests/test_issue_97.py | 6 +- tests/test_issue_99.py | 520 ++++++++++++++++++ tests/test_local_day.py | 2 +- tests/test_projection.py | 4 +- tests/test_tui.py | 4 +- 15 files changed, 1305 insertions(+), 151 deletions(-) create mode 100644 docs/adr/0010-local-day-activity-history.md create mode 100644 tests/test_issue_99.py diff --git a/docs/adr/0010-local-day-activity-history.md b/docs/adr/0010-local-day-activity-history.md new file mode 100644 index 0000000..3deadbb --- /dev/null +++ b/docs/adr/0010-local-day-activity-history.md @@ -0,0 +1,11 @@ +# 10. Preserve local-day activity history with its recorded timezone + +Status: Accepted; legacy-history migration and repair implemented in issue #99. + +Fenris will retain local-day read/write summaries and the boundary evidence needed to interpret them in the existing observation store before three-minute detail expires after 14 days. Each historical summary retains its recorded timezone and day boundaries; this extends [ADR 0001](0001-observation-store-sqlite.md) while preserving the UTC hour/day evidence used for endurance projections. Collection derives new measured intervals. Ordered schema migration and explicit repair rebuild legacy summaries from surviving evidence; TUI and CLI readers remain read-only. + +UTC hourly summaries alone cannot recover a local day whose midnight falls inside a UTC hour. Preserving labelled local summaries trades arbitrary future timezone reinterpretation for bounded detailed-history retention; retaining all fine-grained intervals indefinitely or estimating a split from UTC summaries would violate the agreed retention or evidence semantics. + +Only surviving evidence may be used to derive older local summaries. Dates without enough evidence remain incomplete or unavailable; a measured interval crossing midnight is retained once as shared boundary evidence, never prorated or counted in full on both days. A later timezone change does not silently rewrite historical day boundaries. + +The agreed presentation and validation requirements are in the [live drive activity specification](../spec/live-drive-activity.md). diff --git a/src/fenris/local_day.py b/src/fenris/local_day.py index 52c6892..d9e67a5 100644 --- a/src/fenris/local_day.py +++ b/src/fenris/local_day.py @@ -1,9 +1,10 @@ """Local-day activity derivation from measured sample intervals. Computes durable local-day read/write summaries using recorded local -midnight boundaries. UTC hour/day aggregates remain separate inputs to -endurance projections. The collector owns this derivation, preserving the -read-only TUI boundary (ADR 0010). +midnight boundaries. Collection derives new measured intervals; migration +and repair rebuild legacy summaries from surviving evidence. UTC hour/day +aggregates remain separate inputs to endurance projections, and readers stay +read-only (ADR 0010). Key contracts: - A compatible interval wholly inside a local day contributes its volume once @@ -28,6 +29,8 @@ Retention policy (issue #93): import sqlite3 from dataclasses import dataclass from datetime import date, datetime, time, timedelta, timezone +from itertools import pairwise +from typing import Any from zoneinfo import ZoneInfo @@ -50,6 +53,55 @@ class LocalDaySummary: activity_precision: str = "measured" +@dataclass(frozen=True) +class CoarseLocalDayActivity: + """Rebuilt byte evidence and its coverage limits for one local day.""" + bytes_written: int + bytes_read: int + activity_seconds: int + activity_intervals: int + incomplete: bool + segment_totals: dict[int, "LocalDaySegmentTotals"] + + +@dataclass +class LocalDaySegmentTotals: + """Activity totals attributed to one controller segment.""" + bytes_written: int = 0 + bytes_read: int = 0 + activity_seconds: int = 0 + activity_intervals: int = 0 + + def add( + self, + *, + bytes_written: int, + bytes_read: int, + activity_seconds: int, + activity_intervals: int = 1, + ) -> None: + """Accumulate one measured interval or coarse hour.""" + self.bytes_written += bytes_written + self.bytes_read += bytes_read + self.activity_seconds += activity_seconds + self.activity_intervals += activity_intervals + + +@dataclass(frozen=True) +class ReconstructedLocalDayInterval: + """One counter delta with source and timezone provenance.""" + start: datetime + end: datetime + bytes_written: int + bytes_read: int + activity_seconds: int + segment_id: int + start_sample_id: int + end_sample_id: int + start_tz: str | None + end_tz: str | None + + @dataclass(frozen=True) class LocalDayHistoryEntry: """A local-day summary with evidence-limit metadata for history readout. @@ -119,6 +171,8 @@ def local_day_boundaries(local_date: str, tz_name: str) -> tuple[datetime, datet start_utc = start_local.astimezone(timezone.utc) end_utc = end_local.astimezone(timezone.utc) offset = start_local.utcoffset() + if offset is None: + raise ValueError("Local midnight has no UTC offset") offset_seconds = int(offset.total_seconds()) sign = "+" if offset_seconds >= 0 else "-" offset_seconds = abs(offset_seconds) @@ -277,6 +331,472 @@ def _same_monitoring_period( return False +def repair_legacy_local_day_evidence(conn: sqlite3.Connection) -> int: + """Rebuild untrusted local-day totals from surviving evidence. + + Sample intervals take precedence over UTC-hour summaries. An interval is + usable only when its counters, controller segment, monitoring period, and + recorded local-day boundaries all agree. Old byte totals remain stored but + hidden when no such evidence survives. + """ + tables = { + row[0] for row in conn.execute( + "SELECT name FROM sqlite_master WHERE type='table'" + ) + } + if "local_days" not in tables: + return 0 + + recorded_days = conn.execute( + "SELECT id, local_date, tz_name, utc_start, utc_end, " + " activity_incomplete, activity_precision " + "FROM local_days ORDER BY utc_start, id" + ).fetchall() + local_days = [ + row[:6] for row in recorded_days + if row[6] in ("legacy", "unavailable") + ] + if not local_days: + return 0 + + required_sample_columns = {"id", "ts", "bytes_written", "bytes_read", "segment_id"} + periods = [] + segment_boundaries = [] + segment_windows = {} + if {"monitoring_periods", "controller_segments"}.issubset(tables): + periods = [ + (_utc_datetime(start), _utc_datetime(end) if end else None) + for start, end in conn.execute( + "SELECT started_at, ended_at FROM monitoring_periods" + ) + ] + segment_boundaries = [ + (segment_id, _utc_datetime(opened_at)) + for segment_id, opened_at in conn.execute( + "SELECT id, opened_at FROM controller_segments ORDER BY opened_at, id" + ) + ] + segment_windows = { + segment_id: ( + opened_at, + segment_boundaries[index + 1][1] + if index + 1 < len(segment_boundaries) else None, + ) + for index, (segment_id, opened_at) in enumerate(segment_boundaries) + } + sample_columns = ( + {row[1] for row in conn.execute("PRAGMA table_info(samples)")} + if "samples" in tables else set() + ) + + intervals: list[ReconstructedLocalDayInterval] = [] + samples = [] + if required_sample_columns.issubset(sample_columns): + sample_select = ( + "SELECT id, ts, bytes_written, bytes_read, segment_id, " + + ("local_tz " if "local_tz" in sample_columns else "NULL AS local_tz ") + + "FROM samples ORDER BY ts, id" + ) + samples = [ + { + "id": row[0], "ts": _utc_datetime(row[1]), + "bytes_written": row[2], "bytes_read": row[3], + "segment_id": row[4], "local_tz": row[5], + } + for row in conn.execute(sample_select) + ] + for previous, current in pairwise(samples): + start = previous["ts"] + end = current["ts"] + segment_id = current["segment_id"] + segment_window = segment_windows.get(segment_id) + if ( + end <= start + or segment_id is None + or previous["segment_id"] != segment_id + or segment_window is None + or segment_window[0] > start + or (segment_window[1] is not None and end >= segment_window[1]) + or bool(previous["local_tz"]) != bool(current["local_tz"]) + or previous["bytes_written"] is None + or current["bytes_written"] is None + or previous["bytes_read"] is None + or current["bytes_read"] is None + ): + continue + bytes_written = current["bytes_written"] - previous["bytes_written"] + bytes_read = current["bytes_read"] - previous["bytes_read"] + if bytes_written < 0 or bytes_read < 0: + continue + if not any( + period_start <= start and (period_end is None or end <= period_end) + for period_start, period_end in periods + ): + continue + intervals.append(ReconstructedLocalDayInterval( + start=start, + end=end, + bytes_written=bytes_written, + bytes_read=bytes_read, + activity_seconds=int((end - start).total_seconds()), + segment_id=segment_id, + start_sample_id=previous["id"], + end_sample_id=current["id"], + start_tz=previous["local_tz"], + end_tz=current["local_tz"], + )) + + _record_reconstructed_shared_intervals(conn, intervals, recorded_days) + + day_bounds = [ + { + "id": row[0], + "local_date": row[1], + "tz_name": row[2], + "start": _utc_datetime(row[3]), + "end": _utc_datetime(row[4]), + "precision": row[6], + } + for row in recorded_days + ] + rebuilding_ids = {row[0] for row in local_days} + matching_by_day: dict[int, list[ReconstructedLocalDayInterval]] = { + local_day_id: [] for local_day_id in rebuilding_ids + } + for interval in intervals: + candidates = [ + day for day in day_bounds + if day["start"] <= interval.start + and interval.end <= day["end"] + and all( + zone is None or zone == day["tz_name"] + for zone in (interval.start_tz, interval.end_tz) + ) + ] + if len(candidates) == 1 and candidates[0]["id"] in rebuilding_ids: + matching_by_day[candidates[0]["id"]].append(interval) + + rebuilt = 0 + for row in local_days: + ( + local_day_id, local_date, tz_name, utc_start, utc_end, + prior_incomplete, + ) = row + start = _utc_datetime(utc_start) + end = _utc_datetime(utc_end) + matching = matching_by_day[local_day_id] + + if not matching: + has_samples = any( + start <= sample["ts"] < end for sample in samples + ) + coarse = ( + None if has_samples + else _coarse_local_day_activity( + conn, local_day_id, start, end, periods, segment_boundaries, + ) + ) + if coarse is not None: + incomplete = bool(prior_incomplete) or coarse.incomplete + conn.execute( + "UPDATE local_days SET bytes_written = ?, bytes_read = ?, " + "activity_seconds = ?, activity_intervals = ?, " + "activity_incomplete = ?, activity_precision = 'coarse', " + "last_sample_id = NULL WHERE id = ?", + (coarse.bytes_written, coarse.bytes_read, + coarse.activity_seconds, coarse.activity_intervals, + incomplete, local_day_id), + ) + _replace_local_day_segment_totals( + conn, local_day_id, coarse.segment_totals, + ) + rebuilt += 1 + continue + conn.execute( + "UPDATE local_days SET activity_precision = 'unavailable', " + "activity_incomplete = CASE WHEN EXISTS (" + " SELECT 1 FROM local_day_unallocated_evidence " + " WHERE (start_local_date = ? AND start_tz_name = ?) " + " OR (end_local_date = ? AND end_tz_name = ?)" + ") THEN 1 ELSE 0 END WHERE id = ?", + (local_date, tz_name, local_date, tz_name, local_day_id), + ) + continue + + totals: dict[int, LocalDaySegmentTotals] = {} + bytes_written = 0 + bytes_read = 0 + activity_seconds = 0 + for interval in matching: + segment_total = totals.setdefault( + interval.segment_id, LocalDaySegmentTotals(), + ) + segment_total.add( + bytes_written=interval.bytes_written, + bytes_read=interval.bytes_read, + activity_seconds=interval.activity_seconds, + ) + bytes_written += interval.bytes_written + bytes_read += interval.bytes_read + activity_seconds += interval.activity_seconds + + day_seconds = int((end - start).total_seconds()) + incomplete = bool(prior_incomplete) or activity_seconds < day_seconds + last_sample_id = max(interval.end_sample_id for interval in matching) + _replace_local_day_segment_totals(conn, local_day_id, totals) + conn.execute( + "UPDATE local_days SET bytes_written = ?, bytes_read = ?, " + "activity_seconds = ?, activity_intervals = ?, " + "activity_incomplete = ?, activity_precision = 'measured', " + "last_sample_id = ? WHERE id = ?", + (bytes_written, bytes_read, activity_seconds, len(matching), + incomplete, last_sample_id, local_day_id), + ) + rebuilt += 1 + + return rebuilt + + +def _replace_local_day_segment_totals( + conn: sqlite3.Connection, + local_day_id: int, + totals: dict[int, LocalDaySegmentTotals], +) -> None: + """Replace segment provenance rows for one rebuilt local day.""" + conn.execute( + "DELETE FROM local_day_segment_totals WHERE local_day_id = ?", + (local_day_id,), + ) + conn.executemany( + "INSERT INTO local_day_segment_totals " + "(local_day_id, segment_id, bytes_written, bytes_read, " + " activity_seconds, activity_intervals) " + "VALUES (?, ?, ?, ?, ?, ?)", + ( + ( + local_day_id, + segment_id, + values.bytes_written, + values.bytes_read, + values.activity_seconds, + values.activity_intervals, + ) + for segment_id, values in totals.items() + ), + ) + + +def _record_reconstructed_shared_intervals( + conn: sqlite3.Connection, + intervals: list[ReconstructedLocalDayInterval], + recorded_days: list[tuple], +) -> None: + """Retain compatible counter deltas that cross stored local boundaries.""" + day_rows = [ + { + "id": row[0], "local_date": row[1], "tz_name": row[2], + "utc_start": _utc_datetime(row[3]), + "utc_end": _utc_datetime(row[4]), + "precision": row[6], + } + for row in recorded_days + ] + + def find_day(instant: datetime, zone: str | None, *, end: bool): + candidates = [ + row for row in day_rows + if (row["utc_start"] < instant <= row["utc_end"] if end + else row["utc_start"] <= instant < row["utc_end"]) + and (zone is None or row["tz_name"] == zone) + ] + return candidates[0] if len(candidates) == 1 else None + + for interval in intervals: + start_day = find_day(interval.start, interval.start_tz, end=False) + end_day = find_day(interval.end, interval.end_tz, end=True) + if start_day is None or end_day is None or start_day["id"] == end_day["id"]: + continue + if start_day["precision"] not in ("legacy", "unavailable") and ( + end_day["precision"] not in ("legacy", "unavailable") + ): + continue + + if start_day["tz_name"] == end_day["tz_name"]: + zone_days = sorted( + (row for row in day_rows if row["tz_name"] == start_day["tz_name"]), + key=lambda row: row["utc_start"], + ) + start_index = next( + (i for i, row in enumerate(zone_days) if row["id"] == start_day["id"]), + None, + ) + end_index = next( + (i for i, row in enumerate(zone_days) if row["id"] == end_day["id"]), + None, + ) + if start_index is None or end_index is None or end_index <= start_index: + continue + chain = zone_days[start_index:end_index + 1] + if any( + left["utc_end"] != right["utc_start"] + for left, right in pairwise(chain) + ): + continue + reason = "local_midnight" + else: + reason = "timezone_change" + + conn.execute( + "INSERT OR IGNORE INTO local_day_unallocated_evidence " + "(start_sample_id, end_sample_id, start_local_date, end_local_date, " + " start_tz_name, end_tz_name, started_at, ended_at, bytes_written, " + " bytes_read, reason, segment_id) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + ( + interval.start_sample_id, interval.end_sample_id, + start_day["local_date"], end_day["local_date"], + start_day["tz_name"], end_day["tz_name"], + interval.start.isoformat(), interval.end.isoformat(), + interval.bytes_written, interval.bytes_read, reason, + interval.segment_id, + ), + ) + conn.execute( + "UPDATE local_days SET activity_incomplete = 1 " + "WHERE id IN (?, ?)", + (start_day["id"], end_day["id"]), + ) + + +def _coarse_local_day_activity( + conn: sqlite3.Connection, + local_day_id: int, + day_start: datetime, + day_end: datetime, + periods: list[tuple[datetime, datetime | None]], + segments: list[tuple[int, datetime]], +) -> CoarseLocalDayActivity | None: + """Return UTC hours unique to these recorded local-day bounds.""" + if not periods or not segments or conn.execute( + "SELECT 1 FROM sqlite_master WHERE type='table' AND name='hour_observations'" + ).fetchone() is None: + return None + hour = timedelta(hours=1) + full_hours: dict[datetime, tuple] = {} + boundary_hour_found = False + other_day_bounds = [ + (_utc_datetime(row[0]), _utc_datetime(row[1])) + for row in conn.execute( + "SELECT utc_start, utc_end FROM local_days WHERE id != ?", + (local_day_id,), + ) + ] + first_hour = day_start.replace(minute=0, second=0, microsecond=0) + for row in conn.execute( + "SELECT hour, bytes_written_delta, bytes_read_delta, sample_count, " + " active_seconds, idle_seconds, powered_off_seconds " + "FROM hour_observations WHERE hour >= ? AND hour < ? ORDER BY hour", + (first_hour.isoformat(), day_end.isoformat()), + ): + hour_start = _utc_datetime(row[0]) + hour_end = hour_start + hour + if hour_start >= day_end or hour_end <= day_start: + continue + if day_start <= hour_start and hour_end <= day_end: + if any( + other_start < hour_end and hour_start < other_end + for other_start, other_end in other_day_bounds + ): + # Another historical local-day row claims this whole or + # partial UTC hour too. Keep the byte delta unattributed. + boundary_hour_found = True + continue + full_hours[hour_start] = row + else: + boundary_hour_found = True + + expected_hours = [] + expected_start = day_start.replace(minute=0, second=0, microsecond=0) + if expected_start < day_start: + expected_start += hour + while expected_start + hour <= day_end: + expected_hours.append(expected_start) + expected_start += hour + + def hour_context(hour_start: datetime) -> int | None: + hour_end = hour_start + hour + period_matches = any( + period_start <= hour_start + and (period_end is None or hour_end <= period_end) + for period_start, period_end in periods + ) + segment_matches = [] + for index, (segment_id, opened_at) in enumerate(segments): + next_segment = segments[index + 1][1] if index + 1 < len(segments) else None + if opened_at <= hour_start and (next_segment is None or hour_end <= next_segment): + segment_matches.append(segment_id) + if not period_matches or len(segment_matches) != 1: + return None + return segment_matches[0] + + totals: dict[int, LocalDaySegmentTotals] = {} + bytes_written = 0 + bytes_read = 0 + activity_seconds = 0 + activity_intervals = 0 + context_complete = True + for hour_start, row in full_hours.items(): + segment_id = hour_context(hour_start) + if segment_id is None: + context_complete = False + continue + hour_seconds = (row[4] or 0) + (row[5] or 0) + (row[6] or 0) + activity_seconds += hour_seconds + if (row[3] or 0) <= 0 and (row[1] or 0) == 0 and (row[2] or 0) == 0: + continue + segment_total = totals.setdefault(segment_id, LocalDaySegmentTotals()) + segment_total.add( + bytes_written=row[1] or 0, + bytes_read=row[2] or 0, + activity_seconds=hour_seconds, + ) + bytes_written += row[1] or 0 + bytes_read += row[2] or 0 + activity_intervals += 1 + if not activity_intervals: + return None + all_hours_present = all( + hour_start in full_hours for hour_start in expected_hours + ) + has_unallocated_bytes = False + if conn.execute( + "SELECT 1 FROM sqlite_master WHERE type='table' AND name='day_aggregates'" + ).fetchone() is not None: + last_utc_day = (day_end - timedelta(microseconds=1)).date() + has_unallocated_bytes = conn.execute( + "SELECT 1 FROM day_aggregates WHERE day >= ? AND day <= ? " + "AND (COALESCE(unattributed_bytes_written, 0) > 0 " + " OR COALESCE(unattributed_bytes_read, 0) > 0) LIMIT 1", + (day_start.date().isoformat(), last_utc_day.isoformat()), + ).fetchone() is not None + + complete = ( + all_hours_present + and activity_seconds >= int((day_end - day_start).total_seconds()) + and not boundary_hour_found + and not has_unallocated_bytes + and context_complete + ) + return CoarseLocalDayActivity( + bytes_written=bytes_written, + bytes_read=bytes_read, + activity_seconds=activity_seconds, + activity_intervals=activity_intervals, + incomplete=not complete, + segment_totals=totals, + ) + + def _upsert_activity_day( conn: sqlite3.Connection, local_date: str, @@ -298,17 +818,20 @@ def _upsert_activity_day( "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'measured', ?) " "ON CONFLICT(local_date, tz_name) DO UPDATE SET " "bytes_written = CASE WHEN local_days.activity_precision IN ('legacy', 'unavailable') " + " AND excluded.activity_intervals > 0 " " THEN excluded.bytes_written ELSE local_days.bytes_written + excluded.bytes_written END, " "bytes_read = CASE WHEN local_days.activity_precision IN ('legacy', 'unavailable') " + " AND excluded.activity_intervals > 0 " " THEN excluded.bytes_read ELSE local_days.bytes_read + excluded.bytes_read END, " "activity_seconds = CASE WHEN local_days.activity_precision IN ('legacy', 'unavailable') " + " AND excluded.activity_intervals > 0 " " THEN excluded.activity_seconds ELSE local_days.activity_seconds + excluded.activity_seconds END, " "activity_intervals = CASE WHEN local_days.activity_precision IN ('legacy', 'unavailable') " + " AND excluded.activity_intervals > 0 " " THEN excluded.activity_intervals ELSE local_days.activity_intervals + excluded.activity_intervals END, " - "activity_incomplete = CASE WHEN local_days.activity_precision IN ('legacy', 'unavailable') " - " THEN excluded.activity_incomplete " - " ELSE MAX(local_days.activity_incomplete, excluded.activity_incomplete) END, " - "activity_precision = 'measured', " + "activity_incomplete = MAX(local_days.activity_incomplete, excluded.activity_incomplete), " + "activity_precision = CASE WHEN excluded.activity_intervals > 0 " + " THEN 'measured' ELSE local_days.activity_precision END, " "last_sample_id = COALESCE(excluded.last_sample_id, local_days.last_sample_id)", ( local_date, tz_name, offset, start.isoformat(), end.isoformat(), @@ -320,8 +843,8 @@ def _upsert_activity_day( def record_local_activity_interval( conn: sqlite3.Connection, - previous: dict, - current: dict, + previous: dict[str, Any], + current: dict[str, Any], *, start_sample_id: int, end_sample_id: int, @@ -354,7 +877,12 @@ def record_local_activity_interval( prior_written = previous.get("bytes_written") read = current.get("bytes_read") prior_read = previous.get("bytes_read") - if None in (written, prior_written, read, prior_read): + if ( + written is None + or prior_written is None + or read is None + or prior_read is None + ): return None bytes_written = written - prior_written bytes_read = read - prior_read @@ -379,7 +907,7 @@ def record_local_activity_interval( reason = "local_midnight" else: already_accounted = conn.execute( - "SELECT last_sample_id FROM local_days " + "SELECT last_sample_id, activity_precision FROM local_days " "WHERE local_date = ? AND tz_name = ?", (end_date, current_tz), ).fetchone() @@ -389,6 +917,26 @@ def record_local_activity_interval( and already_accounted[0] >= end_sample_id ): return "known" + replacing_coarse = ( + already_accounted is not None + and already_accounted[1] == "coarse" + ) + if replacing_coarse: + local_day_id = conn.execute( + "SELECT id FROM local_days WHERE local_date = ? AND tz_name = ?", + (end_date, current_tz), + ).fetchone()[0] + conn.execute( + "DELETE FROM local_day_segment_totals WHERE local_day_id = ?", + (local_day_id,), + ) + conn.execute( + "UPDATE local_days SET bytes_written = 0, bytes_read = 0, " + "activity_seconds = 0, activity_intervals = 0, " + "activity_precision = 'unavailable', activity_incomplete = 1, " + "last_sample_id = NULL WHERE id = ?", + (local_day_id,), + ) seconds = int((end - start).total_seconds()) _upsert_activity_day( conn, @@ -398,6 +946,7 @@ def record_local_activity_interval( bytes_read=bytes_read, seconds=seconds, interval_count=1, + incomplete=replacing_coarse, last_sample_id=end_sample_id, ) local_day_id = conn.execute( @@ -567,7 +1116,8 @@ def query_local_day_summary( interval_count = activity_intervals or 0 has_unallocated = evidence_count > 0 - has_known = activity_precision == "measured" and interval_count > 0 + trusted_precision = activity_precision in ("measured", "coarse") + has_known = trusted_precision and interval_count > 0 if not has_known and not has_unallocated: activity_state = "incomplete" if bool(activity_incomplete) else "unavailable" else: @@ -595,7 +1145,7 @@ def query_local_day_summary( else: activity_state = "complete" - exact = activity_precision == "measured" and has_known + exact = trusted_precision and has_known return { "local_date": recorded_date, "tz_name": recorded_tz, diff --git a/src/fenris/projection.py b/src/fenris/projection.py index ee441aa..122671a 100644 --- a/src/fenris/projection.py +++ b/src/fenris/projection.py @@ -187,7 +187,7 @@ def _wall_clock_in_range(conn, start, end): def _resolve_baseline(conn, current_segment): baseline = _get_baseline(conn) - facts = [] + facts: list[str] = [] if baseline is None: return BaselineTier.NONE, None, "no baseline", facts @@ -494,16 +494,16 @@ def _has_complete_local_day(conn): within a monitoring period that has usable observation evidence. This is the prerequisite for showing an endurance outlook. """ - row = conn.execute( - "SELECT 1 FROM local_days WHERE complete = 1 LIMIT 1" - ).fetchone() - return row is not None + return _count_complete_local_days(conn) > 0 def _count_complete_local_days(conn): """Count the number of complete local observation days.""" row = conn.execute( - "SELECT COUNT(*) FROM local_days WHERE complete = 1" + "SELECT COUNT(*) FROM local_days " + "WHERE complete = 1 " + "AND activity_precision IN ('measured', 'coarse') " + "AND activity_intervals > 0" ).fetchone() return row[0] if row else 0 diff --git a/src/fenris/repair.py b/src/fenris/repair.py index 75a0678..bbec034 100644 --- a/src/fenris/repair.py +++ b/src/fenris/repair.py @@ -1,8 +1,9 @@ -"""Historical repair and retention (issue #74). +"""Historical repair and retention (issue #74, #99). 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. +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 @@ -71,7 +72,6 @@ def _set_repair_in_progress(conn: sqlite3.Connection, in_progress: bool) -> None "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: @@ -100,8 +100,6 @@ def _update_repair_status(conn: sqlite3.Connection, result: RepairResult) -> Non "INSERT OR REPLACE INTO store_metadata (key, value) VALUES (?, ?)", ("days_derived", str(days)), ) - - conn.commit() def get_repair_status(conn: sqlite3.Connection) -> RepairStatus: @@ -364,20 +362,16 @@ def repair_derivation( 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 + # 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 @@ -407,7 +401,12 @@ def repair_derivation( 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") @@ -428,21 +427,22 @@ def repair_derivation( ) result.legacy_summaries_preserved = cursor.fetchone()[0] - # Commit transaction - conn.commit() - - # Update repair status + # 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), ) - finally: - # Mark repair as complete - _set_repair_in_progress(conn, False) - return result diff --git a/src/fenris/store.py b/src/fenris/store.py index 806beea..fd7b89d 100644 --- a/src/fenris/store.py +++ b/src/fenris/store.py @@ -11,7 +11,7 @@ from typing import Optional # Schema version - increment on each migration -SCHEMA_VERSION = 5 +SCHEMA_VERSION = 6 # Packaged default placement (spec §8.3). The config may override it, but a @@ -72,9 +72,11 @@ def init_store(store_path: Path) -> sqlite3.Connection: ) elif current_version < SCHEMA_VERSION: # Older version - apply migrations - _apply_migrations(conn, current_version) - conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}") - conn.commit() + try: + _apply_migrations(conn, current_version) + except Exception: + conn.close() + raise return conn @@ -291,95 +293,119 @@ def _create_local_day_segment_totals(conn: sqlite3.Connection) -> None: def _apply_migrations(conn: sqlite3.Connection, current_version: int): """Apply forward-only migrations from current_version to SCHEMA_VERSION. - - Each migration step is a transactional block. Add new steps as sequential - elif branches when SCHEMA_VERSION increases. - - Spec: §3.6, §10.2 + + Commit each version transition independently. A failed step rolls back in + full while earlier successful steps remain versioned and retryable. """ - # Migration 1→2: add segment_id provenance to samples, - # unattributed byte tracking to day_aggregates (issue #73) - if current_version < 2: - # Defensive: only ALTER if table exists (handles minimal v1 stores) - tables = {row[0] for row in conn.execute( - "SELECT name FROM sqlite_master WHERE type='table'" - ).fetchall()} - if "samples" in tables: - # Check if column already exists (idempotent) - cols = {row[1] for row in conn.execute("PRAGMA table_info(samples)").fetchall()} - if "segment_id" not in cols: - conn.execute("ALTER TABLE samples ADD COLUMN segment_id INTEGER") - if "day_aggregates" in tables: - cols = {row[1] for row in conn.execute("PRAGMA table_info(day_aggregates)").fetchall()} - if "unattributed_bytes_written" not in cols: - conn.execute("ALTER TABLE day_aggregates ADD COLUMN unattributed_bytes_written INTEGER DEFAULT 0") - if "unattributed_bytes_read" not in cols: - conn.execute("ALTER TABLE day_aggregates ADD COLUMN unattributed_bytes_read INTEGER DEFAULT 0") - current_version = 2 + migrations = { + 2: _migrate_1_to_2, + 3: _migrate_2_to_3, + 4: _migrate_3_to_4, + 5: _migrate_4_to_5, + 6: _migrate_5_to_6, + } + while current_version < SCHEMA_VERSION: + target_version = current_version + 1 + migration = migrations.get(target_version) + if migration is None: + raise ValueError(f"No migration registered for schema {target_version}") + conn.execute("BEGIN IMMEDIATE") + try: + migration(conn) + conn.execute(f"PRAGMA user_version={target_version}") + conn.commit() + except Exception: + conn.rollback() + raise + current_version = target_version - # Migration 2→3: add local_days table for local-day activity totals - # (issue #90, ADR 0010). Pure addition — no existing rows touched. - if current_version < 3: - tables = {row[0] for row in conn.execute( - "SELECT name FROM sqlite_master WHERE type='table'" + +def _migrate_1_to_2(conn: sqlite3.Connection) -> None: + """Add segment provenance and unattributed UTC byte tracking.""" + tables = {row[0] for row in conn.execute( + "SELECT name FROM sqlite_master WHERE type='table'" + ).fetchall()} + if "samples" in tables: + cols = {row[1] for row in conn.execute( + "PRAGMA table_info(samples)" ).fetchall()} - if "local_days" not in tables: - conn.execute(""" - CREATE TABLE IF NOT EXISTS local_days ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - local_date TEXT NOT NULL, - tz_name TEXT NOT NULL, - tz_offset TEXT NOT NULL, - utc_start TEXT NOT NULL, - utc_end TEXT NOT NULL, - bytes_written INTEGER DEFAULT 0, - bytes_read INTEGER DEFAULT 0, - coverage REAL DEFAULT 0.0, - sample_count INTEGER DEFAULT 0, - complete BOOLEAN DEFAULT 0, - UNIQUE(local_date, tz_name) + if "segment_id" not in cols: + conn.execute("ALTER TABLE samples ADD COLUMN segment_id INTEGER") + if "day_aggregates" in tables: + cols = {row[1] for row in conn.execute( + "PRAGMA table_info(day_aggregates)" + ).fetchall()} + for column in ("unattributed_bytes_written", "unattributed_bytes_read"): + if column not in cols: + conn.execute( + f"ALTER TABLE day_aggregates ADD COLUMN {column} INTEGER DEFAULT 0" ) - """) - current_version = 3 - # Migration 3→4: retain acquired observations until derived evidence can - # be published atomically (issue #97, ADR 0011). - if current_version < 4: - _create_pending_publications(conn) - current_version = 4 - # Migration 4→5: retain precise measured local-day activity and once-only - # unallocated intervals. Existing local-day totals came from UTC-hour - # aggregates, so preserve their rows but mark their local precision legacy. - if current_version < 5: - tables = {row[0] for row in conn.execute( - "SELECT name FROM sqlite_master WHERE type='table'" +def _migrate_2_to_3(conn: sqlite3.Connection) -> None: + """Add local-day activity summaries.""" + tables = {row[0] for row in conn.execute( + "SELECT name FROM sqlite_master WHERE type='table'" + ).fetchall()} + if "local_days" not in tables: + conn.execute(""" + CREATE TABLE local_days ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + local_date TEXT NOT NULL, + tz_name TEXT NOT NULL, + tz_offset TEXT NOT NULL, + utc_start TEXT NOT NULL, + utc_end TEXT NOT NULL, + bytes_written INTEGER DEFAULT 0, + bytes_read INTEGER DEFAULT 0, + coverage REAL DEFAULT 0.0, + sample_count INTEGER DEFAULT 0, + complete BOOLEAN DEFAULT 0, + UNIQUE(local_date, tz_name) + ) + """) + + +def _migrate_3_to_4(conn: sqlite3.Connection) -> None: + """Add private publication staging for acquired observations.""" + _create_pending_publications(conn) + + +def _migrate_4_to_5(conn: sqlite3.Connection) -> None: + """Add measured local-day evidence storage and mark old totals legacy.""" + tables = {row[0] for row in conn.execute( + "SELECT name FROM sqlite_master WHERE type='table'" + ).fetchall()} + if "samples" in tables: + sample_cols = {row[1] for row in conn.execute( + "PRAGMA table_info(samples)" ).fetchall()} - if "samples" in tables: - sample_cols = {row[1] for row in conn.execute( - "PRAGMA table_info(samples)" - ).fetchall()} - if "local_tz" not in sample_cols: - conn.execute("ALTER TABLE samples ADD COLUMN local_tz TEXT") - if "local_days" in tables: - local_cols = {row[1] for row in conn.execute( - "PRAGMA table_info(local_days)" - ).fetchall()} - for column, declaration in ( - ("activity_seconds", "INTEGER NOT NULL DEFAULT 0"), - ("activity_intervals", "INTEGER NOT NULL DEFAULT 0"), - ("activity_incomplete", "BOOLEAN NOT NULL DEFAULT 0"), - ("activity_precision", "TEXT NOT NULL DEFAULT 'legacy'"), - ("last_sample_id", "INTEGER"), - ): - if column not in local_cols: - conn.execute( - "ALTER TABLE local_days ADD COLUMN %s %s" % - (column, declaration) - ) - _create_local_day_shared_evidence(conn) - _create_local_day_segment_totals(conn) - current_version = 5 + if "local_tz" not in sample_cols: + conn.execute("ALTER TABLE samples ADD COLUMN local_tz TEXT") + if "local_days" in tables: + local_cols = {row[1] for row in conn.execute( + "PRAGMA table_info(local_days)" + ).fetchall()} + for column, declaration in ( + ("activity_seconds", "INTEGER NOT NULL DEFAULT 0"), + ("activity_intervals", "INTEGER NOT NULL DEFAULT 0"), + ("activity_incomplete", "BOOLEAN NOT NULL DEFAULT 0"), + ("activity_precision", "TEXT NOT NULL DEFAULT 'legacy'"), + ("last_sample_id", "INTEGER"), + ): + if column not in local_cols: + conn.execute( + f"ALTER TABLE local_days ADD COLUMN {column} {declaration}" + ) + _create_local_day_shared_evidence(conn) + _create_local_day_segment_totals(conn) + + +def _migrate_5_to_6(conn: sqlite3.Connection) -> None: + """Rebuild local-day summaries from surviving trustworthy evidence.""" + from .local_day import repair_legacy_local_day_evidence + + repair_legacy_local_day_evidence(conn) def migrate_to_latest(store_path: Path) -> int: @@ -416,9 +442,11 @@ def migrate_to_latest(store_path: Path) -> int: return SCHEMA_VERSION steps = SCHEMA_VERSION - current_version - _apply_migrations(conn, current_version) - conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}") - conn.commit() + try: + _apply_migrations(conn, current_version) + except Exception: + conn.close() + raise conn.close() return steps diff --git a/tests/test_acceptance_sweep.py b/tests/test_acceptance_sweep.py index 739c362..e4b8911 100644 --- a/tests/test_acceptance_sweep.py +++ b/tests/test_acceptance_sweep.py @@ -121,8 +121,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00", conn.execute( "INSERT INTO local_days " "(local_date, tz_name, tz_offset, utc_start, utc_end, " - " bytes_written, bytes_read, coverage, sample_count, complete) " - "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + " bytes_written, bytes_read, coverage, sample_count, complete, activity_intervals) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)", (local_date, tz_name, tz_offset, utc_start, utc_end, bw, br, coverage, samples, complete), ) diff --git a/tests/test_collector_history_tracer.py b/tests/test_collector_history_tracer.py index 5678228..f8fb8fd 100644 --- a/tests/test_collector_history_tracer.py +++ b/tests/test_collector_history_tracer.py @@ -154,7 +154,7 @@ class TestSchemaMigration: # Migrate steps = migrate_to_latest(db) - assert steps == 4 # v1→v2→v3→v4→v5 + assert steps == 5 # v1→v2→v3→v4→v5→v6 # Verify data preserved conn = sqlite3.connect(str(db)) @@ -375,7 +375,9 @@ class TestCollectionAtomicity: } monkeypatch.setattr("fenris.status.query_service_state", lambda: service) - now = datetime.now(timezone.utc).replace(second=0, microsecond=0) + now = datetime.now(timezone.utc).replace( + minute=35, second=0, microsecond=0, + ) first = { **smartctl_fixture, "nvme_smart_health_information_log": { diff --git a/tests/test_complete_observation_day_gate.py b/tests/test_complete_observation_day_gate.py index 02366bf..574d083 100644 --- a/tests/test_complete_observation_day_gate.py +++ b/tests/test_complete_observation_day_gate.py @@ -98,8 +98,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00", conn.execute( "INSERT INTO local_days " "(local_date, tz_name, tz_offset, utc_start, utc_end, " - " bytes_written, bytes_read, coverage, sample_count, complete) " - "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + " bytes_written, bytes_read, coverage, sample_count, complete, activity_intervals) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)", (local_date, tz_name, tz_offset, utc_start, utc_end, bw, br, coverage, samples, complete), ) @@ -166,6 +166,47 @@ class TestGateNoCompleteDay: # The key assertion: gate message should NOT appear when gate IS met assert not any("full local observation day" in f for f in result.contributing_facts) + def test_complete_legacy_day_without_trusted_activity_does_not_open_gate( + self, store, + ): + """A complete flag cannot make unavailable local activity qualify.""" + _insert_baseline(store) + _insert_segment(store) + _open_period(store) + for i in range(3): + day = (datetime(2026, 9, 27) + timedelta(days=i)).strftime("%Y-%m-%d") + _insert_day(store, day, bw=1024 * 1024 * 100) + _insert_sample(store, "2026-09-30T10:00:00+00:00", pu=5) + _insert_local_day(store, "2026-09-29", complete=True) + store.execute( + "UPDATE local_days SET activity_precision = 'measured', activity_intervals = 0 " + "WHERE local_date = '2026-09-29'" + ) + result = compute_projection(store, _clock()) + assert result.confidence_state == ConfidenceState.UNSUPPORTED + assert any("full local observation day" in fact + for fact in result.contributing_facts) + + def test_partial_but_trusted_local_activity_opens_gate(self, store): + """Known local intervals can coexist with an incomplete day total.""" + _insert_baseline(store) + _insert_segment(store) + _open_period(store) + for i in range(3): + day = (datetime(2026, 9, 27) + timedelta(days=i)).strftime("%Y-%m-%d") + _insert_day(store, day, bw=1024 * 1024 * 100) + _insert_sample(store, "2026-09-30T10:00:00+00:00", pu=5) + _insert_local_day(store, "2026-09-29", complete=True) + store.execute( + "UPDATE local_days SET activity_incomplete = 1 " + "WHERE local_date = '2026-09-29'" + ) + + result = compute_projection(store, _clock()) + + assert not any("full local observation day" in fact + for fact in result.contributing_facts) + # --------------------------------------------------------------------------- # Gate-2: One complete local day → Limited confidence diff --git a/tests/test_issue_72_edge_cases.py b/tests/test_issue_72_edge_cases.py index 4aaf2c8..1a496bf 100644 --- a/tests/test_issue_72_edge_cases.py +++ b/tests/test_issue_72_edge_cases.py @@ -90,8 +90,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00", conn.execute( "INSERT INTO local_days " "(local_date, tz_name, tz_offset, utc_start, utc_end, " - " bytes_written, bytes_read, coverage, sample_count, complete) " - "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + " bytes_written, bytes_read, coverage, sample_count, complete, activity_intervals) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)", (local_date, tz_name, tz_offset, utc_start, utc_end, bw, br, coverage, samples, complete), ) diff --git a/tests/test_issue_77.py b/tests/test_issue_77.py index 7256c98..fbe7397 100644 --- a/tests/test_issue_77.py +++ b/tests/test_issue_77.py @@ -98,8 +98,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00", conn.execute( "INSERT INTO local_days " "(local_date, tz_name, tz_offset, utc_start, utc_end, " - " bytes_written, bytes_read, coverage, sample_count, complete) " - "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + " bytes_written, bytes_read, coverage, sample_count, complete, activity_intervals) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)", (local_date, tz_name, tz_offset, utc_start, utc_end, bw, br, coverage, samples, complete), ) diff --git a/tests/test_issue_97.py b/tests/test_issue_97.py index 16e5730..f3b6643 100644 --- a/tests/test_issue_97.py +++ b/tests/test_issue_97.py @@ -103,7 +103,9 @@ def test_failed_publication_survives_restart_and_readers_keep_last_consistent_vi monkeypatch.setenv("TZ", "UTC") service = {"boot_enabled": True, "timer_active": True, "last_collect_ok": True} monkeypatch.setattr("fenris.status.query_service_state", lambda: service) - now = datetime.now(timezone.utc).replace(second=0, microsecond=0) + now = datetime.now(timezone.utc).replace( + minute=35, second=0, microsecond=0, + ) store_path = tmp_path / "observations.db" config = {"device": "/dev/nvme0", "store_path": str(store_path)} sysfs_path = sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0" @@ -285,7 +287,7 @@ def test_v3_readers_ignore_pending_table_until_store_migrates(tmp_path): migrated = init_store(store_path) try: - assert migrated.execute("PRAGMA user_version").fetchone()[0] == 5 + assert migrated.execute("PRAGMA user_version").fetchone()[0] == 6 assert migrated.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0 finally: migrated.close() diff --git a/tests/test_issue_99.py b/tests/test_issue_99.py new file mode 100644 index 0000000..08a99f9 --- /dev/null +++ b/tests/test_issue_99.py @@ -0,0 +1,520 @@ +"""Historical local-day migration and repair acceptance tests for issue #99.""" +import sqlite3 +import sys +from datetime import datetime, timezone +from pathlib import Path + +import pytest + +sys.path.insert(0, str(Path(__file__).parent.parent / "src")) + +from fenris.local_day import ( + query_local_day_summary, + record_local_activity_interval, +) +from fenris.status import read_status +from fenris.store import init_store + +SOURCE_SCHEMA_VERSION = 5 +LOCAL_DATE = "2026-09-02" +TIMEZONE = "Asia/Kolkata" +UTC_START = "2026-09-01T18:30:00+00:00" +UTC_END = "2026-09-02T18:30:00+00:00" + + +def _legacy_store(path: Path) -> sqlite3.Connection: + """Build a v5 store with one old coarse local-day summary.""" + conn = init_store(path) + conn.execute( + "INSERT INTO local_days " + "(local_date, tz_name, tz_offset, utc_start, utc_end, bytes_written, " + " bytes_read, coverage, sample_count, complete, activity_incomplete, " + " activity_precision) " + "VALUES (?, ?, '+05:30', ?, ?, 999999, 888888, 1.0, 24, 1, 0, 'legacy')", + (LOCAL_DATE, TIMEZONE, UTC_START, UTC_END), + ) + conn.execute( + "INSERT INTO monitoring_periods (started_at, ended_at) VALUES (?, ?)", + (UTC_START, "2026-09-03T00:00:00+00:00"), + ) + conn.execute( + "INSERT INTO controller_segments " + "(opened_at, identity_key, identity_degraded) VALUES (?, ?, 0)", + (UTC_START, "controller-a"), + ) + conn.execute( + "INSERT INTO hour_observations " + "(hour, bytes_written_delta, bytes_read_delta, sample_count, coverage) " + "VALUES ('2026-09-02T00:00:00+00:00', 123, 45, 2, 1.0)" + ) + conn.execute( + "INSERT INTO day_aggregates " + "(day, bytes_written_delta, bytes_read_delta, coverage) " + "VALUES ('2026-09-02', 678, 90, 1.0)" + ) + conn.execute( + "INSERT INTO samples (ts, device, bytes_written, bytes_read, power_on_hours, " + " segment_id, local_tz) VALUES (?, '/dev/nvme0', 1000, 2000, 100, 1, ?)", + ("2026-09-01T19:00:00+00:00", TIMEZONE), + ) + conn.execute( + "INSERT INTO samples (ts, device, bytes_written, bytes_read, power_on_hours, " + " segment_id, local_tz) VALUES (?, '/dev/nvme0', 11000, 5000, 115, 1, ?)", + ("2026-09-02T10:00:00+00:00", TIMEZONE), + ) + conn.execute(f"PRAGMA user_version={SOURCE_SCHEMA_VERSION}") + conn.commit() + return conn + + +def test_upgrade_rebuilds_legacy_day_from_samples_and_keeps_recorded_bounds( + tmp_path, monkeypatch, +): + db = tmp_path / "legacy.db" + conn = _legacy_store(db) + conn.close() + monkeypatch.setenv("TZ", "UTC") + + conn = init_store(db) + conn.close() + + with read_status( + db, datetime(2026, 9, 4, tzinfo=timezone.utc), query_services=False, + ) as (reader, composition): + assert reader is not None + assert composition.store_fault is None + summary = query_local_day_summary( + reader, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary is not None + assert summary["local_date"] == LOCAL_DATE + assert summary["tz_name"] == TIMEZONE + assert summary["tz_offset"] == "+05:30" + assert summary["utc_start"] == UTC_START + assert summary["utc_end"] == UTC_END + assert summary["bytes_written"] == 10_000 + assert summary["bytes_read"] == 3_000 + assert summary["activity_precision"] == "measured" + assert summary["segments"] == [{ + "segment_id": 1, + "bytes_written": 10_000, + "bytes_read": 3_000, + "activity_seconds": 54_000, + "activity_intervals": 1, + }] + assert summary["activity_state"] == "incomplete" + assert reader.execute( + "SELECT bytes_written_delta FROM hour_observations" + ).fetchone()[0] == 123 + assert reader.execute( + "SELECT bytes_written_delta FROM day_aggregates" + ).fetchone()[0] == 678 + + conn = init_store(db) + conn.close() + with read_status( + db, datetime(2026, 9, 4, tzinfo=timezone.utc), query_services=False, + ) as (reader, _composition): + assert reader is not None + summary = query_local_day_summary( + reader, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] == 10_000 + assert summary["bytes_read"] == 3_000 + + +def test_repair_rebuilds_legacy_local_day_once_across_restarts(tmp_path): + from fenris.repair import repair_derivation + + db = tmp_path / "repair.db" + conn = _legacy_store(db) + + first = repair_derivation(conn) + assert first.ok + summary = query_local_day_summary( + conn, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] == 10_000 + assert summary["bytes_read"] == 3_000 + + conn.close() + conn = sqlite3.connect(db) + second = repair_derivation(conn) + assert second.ok + summary = query_local_day_summary( + conn, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] == 10_000 + assert summary["bytes_read"] == 3_000 + assert conn.execute( + "SELECT bytes_written, bytes_read, activity_intervals " + "FROM local_day_segment_totals" + ).fetchall() == [(10_000, 3_000, 1)] + conn.close() + + +def test_interrupted_repair_keeps_local_day_rebuild_retryable(tmp_path): + from fenris.repair import repair_derivation + + db = tmp_path / "repair-interrupted.db" + conn = _legacy_store(db) + conn.execute( + "CREATE TRIGGER interrupt_repair_local_rebuild " + "BEFORE INSERT ON local_day_segment_totals " + "BEGIN SELECT RAISE(ABORT, 'simulated repair interruption'); END" + ) + conn.commit() + + failed = repair_derivation(conn) + assert not failed.ok + assert query_local_day_summary( + conn, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + )["bytes_written"] is None + assert conn.execute( + "SELECT activity_precision, bytes_written FROM local_days" + ).fetchone() == ("legacy", 999_999) + conn.execute("DROP TRIGGER interrupt_repair_local_rebuild") + conn.commit() + conn.close() + + conn = sqlite3.connect(db) + retried = repair_derivation(conn) + assert retried.ok + summary = query_local_day_summary( + conn, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] == 10_000 + assert summary["bytes_read"] == 3_000 + conn.close() + + +def test_stale_repair_marker_does_not_block_crash_retry(tmp_path): + from fenris.repair import repair_derivation + + conn = _legacy_store(tmp_path / "stale-repair.db") + conn.execute( + "INSERT OR REPLACE INTO store_metadata (key, value) " + "VALUES ('repair_in_progress', 'true')" + ) + conn.commit() + + result = repair_derivation(conn) + + assert result.ok + assert conn.execute( + "SELECT value FROM store_metadata WHERE key = 'repair_in_progress'" + ).fetchone() == ("false",) + assert query_local_day_summary( + conn, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + )["bytes_written"] == 10_000 + conn.close() + + +def test_upgrade_uses_whole_utc_hours_but_skips_midnight_straddling_hours( + tmp_path, +): + db = tmp_path / "coarse.db" + conn = _legacy_store(db) + conn.execute("DELETE FROM samples") + conn.execute( + "INSERT INTO hour_observations " + "(hour, bytes_written_delta, bytes_read_delta, sample_count, " + " active_seconds, coverage) " + "VALUES ('2026-09-01T18:00:00+00:00', 99999, 9999, 2, 3600, 1.0)" + ) + conn.execute( + "INSERT INTO hour_observations " + "(hour, bytes_written_delta, bytes_read_delta, sample_count, " + " active_seconds, coverage) " + "VALUES ('2026-09-02T18:00:00+00:00', 88888, 8888, 2, 3600, 1.0)" + ) + conn.commit() + conn.close() + + conn = init_store(db) + summary = query_local_day_summary( + conn, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] == 123 + assert summary["bytes_read"] == 45 + assert summary["activity_precision"] == "coarse" + assert summary["activity_state"] == "incomplete" + assert summary["utc_start"] == UTC_START + assert summary["utc_end"] == UTC_END + + conn.execute( + "INSERT INTO samples " + "(id, ts, device, bytes_written, bytes_read, segment_id, local_tz) " + "VALUES (10, '2026-09-02T12:00:00+00:00', '/dev/nvme0', 1000, 2000, 1, ?), " + " (11, '2026-09-02T12:05:00+00:00', '/dev/nvme0', 2000, 2500, 1, ?)", + (TIMEZONE, TIMEZONE), + ) + conn.commit() + assert record_local_activity_interval( + conn, + { + "id": 10, "ts": "2026-09-02T12:00:00+00:00", + "bytes_written": 1_000, "bytes_read": 2_000, + "local_tz": TIMEZONE, + }, + { + "id": 11, "ts": "2026-09-02T12:05:00+00:00", + "bytes_written": 2_000, "bytes_read": 2_500, + "local_tz": TIMEZONE, + }, + start_sample_id=10, end_sample_id=11, segment_id=1, + ) == "known" + summary = query_local_day_summary( + conn, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] == 1_000 + assert summary["bytes_read"] == 500 + assert summary["activity_precision"] == "measured" + assert summary["segments"] == [{ + "segment_id": 1, + "bytes_written": 1_000, + "bytes_read": 500, + "activity_seconds": 300, + "activity_intervals": 1, + }] + conn.close() + + +def test_upgrade_does_not_double_count_coarse_hours_in_overlapping_days(tmp_path): + db = tmp_path / "overlapping-coarse-days.db" + conn = _legacy_store(db) + conn.execute("DELETE FROM samples") + conn.execute( + "INSERT INTO local_days " + "(local_date, tz_name, tz_offset, utc_start, utc_end, bytes_written, " + " bytes_read, coverage, sample_count, complete, activity_precision) " + "VALUES ('2026-09-02', 'UTC', '+00:00', " + "'2026-09-02T00:00:00+00:00', '2026-09-03T00:00:00+00:00', " + "999999, 888888, 1.0, 24, 1, 'legacy')" + ) + conn.commit() + conn.close() + + conn = init_store(db) + + for tz_name in (TIMEZONE, "UTC"): + summary = query_local_day_summary( + conn, LOCAL_DATE, tz_name, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] is None + assert summary["bytes_read"] is None + assert summary["activity_precision"] == "unavailable" + assert summary["activity_state"] == "unavailable" + assert conn.execute( + "SELECT bytes_written_delta, bytes_read_delta FROM hour_observations " + "WHERE hour = '2026-09-02T00:00:00+00:00'" + ).fetchone() == (123, 45) + conn.close() + + +def test_upgrade_retains_local_midnight_interval_once_as_shared_evidence(tmp_path): + db = tmp_path / "shared.db" + conn = _legacy_store(db) + conn.execute( + "INSERT INTO local_days " + "(local_date, tz_name, tz_offset, utc_start, utc_end, bytes_written, " + " bytes_read, coverage, sample_count, complete, activity_precision) " + "VALUES ('2026-09-01', ?, '+05:30', " + "'2026-08-31T18:30:00+00:00', '2026-09-01T18:30:00+00:00', " + "500, 600, 1.0, 24, 1, 'legacy')", + (TIMEZONE,), + ) + conn.execute( + "UPDATE monitoring_periods SET started_at = '2026-08-31T18:00:00+00:00'" + ) + conn.execute( + "UPDATE controller_segments SET opened_at = '2026-08-31T18:00:00+00:00'" + ) + conn.execute( + "UPDATE samples SET ts = '2026-09-01T18:25:00+00:00', " + "bytes_written = 1000, bytes_read = 2000 WHERE id = 1" + ) + conn.execute( + "UPDATE samples SET ts = '2026-09-01T18:35:00+00:00', " + "bytes_written = 1100, bytes_read = 2050 WHERE id = 2" + ) + conn.commit() + conn.close() + + conn = init_store(db) + shared = conn.execute( + "SELECT COUNT(*), SUM(bytes_written), SUM(bytes_read) " + "FROM local_day_unallocated_evidence" + ).fetchone() + assert shared == (1, 100, 50) + for local_date in ("2026-09-01", "2026-09-02"): + summary = query_local_day_summary( + conn, local_date, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] is None + assert summary["shared_bytes_written"] == 100 + assert summary["shared_bytes_read"] == 50 + assert summary["activity_state"] == "incomplete" + from fenris.repair import repair_derivation + assert repair_derivation(conn).ok + assert conn.execute( + "SELECT COUNT(*), SUM(bytes_written), SUM(bytes_read) " + "FROM local_day_unallocated_evidence" + ).fetchone() == (1, 100, 50) + conn.close() + + +def test_interrupted_upgrade_keeps_legacy_day_retryable(tmp_path): + db = tmp_path / "interrupted.db" + conn = _legacy_store(db) + conn.execute( + "CREATE TRIGGER interrupt_local_rebuild " + "BEFORE INSERT ON local_day_segment_totals " + "BEGIN SELECT RAISE(ABORT, 'simulated interruption'); END" + ) + conn.commit() + conn.close() + + with pytest.raises(sqlite3.IntegrityError, match="simulated interruption"): + init_store(db) + + conn = sqlite3.connect(db) + assert conn.execute("PRAGMA user_version").fetchone() == (SOURCE_SCHEMA_VERSION,) + assert conn.execute( + "SELECT activity_precision, bytes_written FROM local_days" + ).fetchone() == ("legacy", 999_999) + conn.execute("DROP TRIGGER interrupt_local_rebuild") + conn.commit() + conn.close() + + conn = init_store(db) + summary = query_local_day_summary( + conn, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] == 10_000 + assert summary["bytes_read"] == 3_000 + conn.close() + + +def test_prior_migration_step_commits_before_local_day_repair(tmp_path): + db = tmp_path / "migration-steps.db" + conn = _legacy_store(db) + conn.execute("PRAGMA user_version=4") + conn.execute( + "CREATE TRIGGER interrupt_local_rebuild BEFORE UPDATE OF bytes_written " + "ON local_days BEGIN SELECT RAISE(ABORT, 'simulated interruption'); END" + ) + conn.commit() + conn.close() + + with pytest.raises(sqlite3.IntegrityError, match="simulated interruption"): + init_store(db) + + conn = sqlite3.connect(db) + assert conn.execute("PRAGMA user_version").fetchone() == (5,) + assert conn.execute( + "SELECT 1 FROM sqlite_master WHERE type='table' " + "AND name='local_day_segment_totals'" + ).fetchone() is not None + assert conn.execute( + "SELECT bytes_written FROM local_days WHERE local_date=?", + (LOCAL_DATE,), + ).fetchone() == (999_999,) + conn.close() + + +def test_upgrade_does_not_allocate_unzoned_interval_to_overlapping_timezones( + tmp_path, +): + db = tmp_path / "overlapping-timezones.db" + conn = _legacy_store(db) + conn.execute( + "INSERT INTO local_days " + "(local_date, tz_name, tz_offset, utc_start, utc_end, bytes_written, " + " bytes_read, coverage, sample_count, complete, activity_precision) " + "VALUES ('2026-09-02', 'UTC', '+00:00', " + "'2026-09-02T00:00:00+00:00', '2026-09-03T00:00:00+00:00', " + "999999, 888888, 1.0, 24, 1, 'legacy')" + ) + conn.execute( + "UPDATE samples SET ts = '2026-09-02T01:00:00+00:00', local_tz = NULL " + "WHERE id = 1" + ) + conn.execute( + "UPDATE samples SET ts = '2026-09-02T10:00:00+00:00', local_tz = NULL " + "WHERE id = 2" + ) + conn.commit() + conn.close() + + conn = init_store(db) + for tz_name in (TIMEZONE, "UTC"): + summary = query_local_day_summary( + conn, LOCAL_DATE, tz_name, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] is None + assert summary["bytes_read"] is None + assert summary["activity_precision"] == "unavailable" + assert summary["activity_state"] == "unavailable" + conn.close() + + +def test_upgrade_rejects_interval_before_recorded_controller_segment( + tmp_path, +): + db = tmp_path / "segment-boundary.db" + conn = _legacy_store(db) + conn.execute( + "UPDATE controller_segments " + "SET opened_at = '2026-09-01T20:00:00+00:00' WHERE id = 1" + ) + conn.commit() + conn.close() + + conn = init_store(db) + summary = query_local_day_summary( + conn, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] is None + assert summary["bytes_read"] is None + assert summary["activity_precision"] == "unavailable" + conn.close() + + +def test_sample_at_local_midnight_does_not_hide_coarse_reconstruction(tmp_path): + db = tmp_path / "midnight-sample.db" + conn = _legacy_store(db) + conn.execute( + "UPDATE samples SET ts = '2026-09-01T17:00:00+00:00' WHERE id = 1" + ) + conn.execute( + "UPDATE samples SET ts = ?, local_tz = ? WHERE id = 2", + (UTC_END, TIMEZONE), + ) + conn.commit() + conn.close() + + conn = init_store(db) + summary = query_local_day_summary( + conn, LOCAL_DATE, TIMEZONE, + datetime(2026, 9, 4, tzinfo=timezone.utc), + ) + assert summary["bytes_written"] == 123 + assert summary["bytes_read"] == 45 + assert summary["activity_precision"] == "coarse" + assert summary["activity_state"] == "incomplete" + conn.close() diff --git a/tests/test_local_day.py b/tests/test_local_day.py index 35b6e09..2dede83 100644 --- a/tests/test_local_day.py +++ b/tests/test_local_day.py @@ -170,7 +170,7 @@ class TestSchemaMigration: stored = conn.execute( "SELECT bytes_written, activity_precision FROM local_days" ).fetchone() - assert stored == (5_120_000, "legacy") + assert stored == (5_120_000, "unavailable") visible = query_local_day_summary(conn, "2026-09-01", "UTC") assert visible["bytes_written"] is None assert visible["activity_state"] == "unavailable" diff --git a/tests/test_projection.py b/tests/test_projection.py index dbcc1f2..cd135e9 100644 --- a/tests/test_projection.py +++ b/tests/test_projection.py @@ -100,8 +100,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00", conn.execute( "INSERT INTO local_days " "(local_date, tz_name, tz_offset, utc_start, utc_end, " - " bytes_written, bytes_read, coverage, sample_count, complete) " - "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + " bytes_written, bytes_read, coverage, sample_count, complete, activity_intervals) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)", (local_date, tz_name, tz_offset, utc_start, utc_end, bw, br, coverage, samples, complete), ) diff --git a/tests/test_tui.py b/tests/test_tui.py index d7e3023..50a9ed6 100644 --- a/tests/test_tui.py +++ b/tests/test_tui.py @@ -124,8 +124,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00", conn.execute( "INSERT INTO local_days " "(local_date, tz_name, tz_offset, utc_start, utc_end, " - " bytes_written, bytes_read, coverage, sample_count, complete) " - "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + " bytes_written, bytes_read, coverage, sample_count, complete, activity_intervals) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)", (local_date, tz_name, tz_offset, utc_start, utc_end, bw, br, coverage, samples, complete), )