diff --git a/README.md b/README.md index af6019a..26b7d27 100644 --- a/README.md +++ b/README.md @@ -263,6 +263,12 @@ measured read/write volumes, not transfer speed; missing evidence breaks the trace. `?` marks a gap, `~` a partial total, and `u` unallocated daily volume. Use the selected-point readout for exact values and evidence state. +Local-day totals use measured intervals inside each recorded local date. An +interval that crosses local midnight appears once as shared evidence and stays +outside both known totals. The dashboard labels current totals as so far and +shows incomplete or unavailable dates without treating them as zero. Older +UTC-only summaries cannot establish exact local-day totals. + - `v` cycles Live / Day / History; the tabs are also clickable. - `←` / `→` inspect points; `w` switches read/write volume in every view. - `[` / `]` browse dates, `g` enters a date, and `t` returns to today/live. diff --git a/src/fenris/collector.py b/src/fenris/collector.py index a09e317..5fc59f3 100644 --- a/src/fenris/collector.py +++ b/src/fenris/collector.py @@ -239,8 +239,8 @@ def write_sample( percentage_used, available_spare, media_errors, power_on_hours, power_cycles, unsafe_shutdowns, temperature_c, data_units_written, data_units_read, bytes_written, bytes_read, - critical_warning, segment_id - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + critical_warning, segment_id, local_tz + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( sample["ts"], @@ -263,6 +263,7 @@ def write_sample( sample["bytes_read"], sample["critical_warning"], segment_id, + sample.get("local_tz"), ), ) @@ -272,6 +273,7 @@ def write_sample( "identity_key": identity_key, "identity_degraded": identity_degraded, "segment_id": segment_id, + "sample_id": cursor.lastrowid, } @@ -291,12 +293,12 @@ def _publish_observation( ) -> None: """Publish one staged sample and all dependent evidence in caller transaction.""" observed_at = _observation_time(sample) + sample["local_tz"] = tz_name validate_sample_invariants(sample, conn) ensure_period_open(conn, observed_at) seg_info = write_sample(sample, identity, conn, observed_at) - cursor = conn.execute("SELECT id FROM samples ORDER BY id DESC LIMIT 1") - current_id = cursor.fetchone()[0] + current_id = seg_info["sample_id"] prev = find_previous_sample(conn, seg_info.get("segment_id"), current_id) if prev is not None: current = { @@ -308,8 +310,28 @@ def _publish_observation( "temperature_c": sample["temperature_c"], "data_units_written": sample["data_units_written"], "data_units_read": sample["data_units_read"], + "local_tz": tz_name, } derive_hours_from_interval(conn, prev, current) + from .local_day import record_local_activity_interval + record_local_activity_interval( + conn, + prev, + current, + start_sample_id=prev["id"], + end_sample_id=current_id, + segment_id=seg_info.get("segment_id"), + ) + else: + previous_any_segment = find_previous_sample(conn, None, current_id) + if previous_any_segment is not None: + from .local_day import mark_local_activity_gap + mark_local_activity_gap( + conn, + sample["ts"], + tz_name, + previous_any_segment, + ) from .day_aggregate import derive_all_days, persist_day_aggregate for aggregate in derive_all_days(conn): diff --git a/src/fenris/derive.py b/src/fenris/derive.py index 6c87da3..a592b97 100644 --- a/src/fenris/derive.py +++ b/src/fenris/derive.py @@ -29,7 +29,7 @@ def find_previous_sample( if segment_id is not None: cursor = conn.execute( "SELECT id, ts, bytes_written, bytes_read, power_on_hours, " - " temperature_c, data_units_written, data_units_read " + " temperature_c, data_units_written, data_units_read, local_tz " "FROM samples WHERE id < ? AND segment_id = ? " "ORDER BY id DESC LIMIT 1", (current_sample_id, segment_id), @@ -37,7 +37,7 @@ def find_previous_sample( else: cursor = conn.execute( "SELECT id, ts, bytes_written, bytes_read, power_on_hours, " - " temperature_c, data_units_written, data_units_read " + " temperature_c, data_units_written, data_units_read, local_tz " "FROM samples WHERE id < ? " "ORDER BY id DESC LIMIT 1", (current_sample_id,), @@ -49,7 +49,7 @@ def find_previous_sample( "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], + "data_units_read": row[7], "local_tz": row[8], } diff --git a/src/fenris/local_day.py b/src/fenris/local_day.py index 55be0a0..52c6892 100644 --- a/src/fenris/local_day.py +++ b/src/fenris/local_day.py @@ -1,22 +1,21 @@ -"""Local-day activity derivation from UTC hour observations. +"""Local-day activity derivation from measured sample intervals. -Computes durable local-day read/write summaries using the actual local -midnight boundaries, retaining UTC hour/day aggregates for endurance -projections. The collector owns this derivation, preserving the -existing read-only TUI boundary (ADR 0010). +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). Key contracts: -- An interval wholly attributable to a local day contributes its volume once -- Midnight-spanning intervals are retained once as shared/unallocated evidence +- A compatible interval wholly inside a local day contributes its volume once +- Midnight-spanning and incompatible intervals persist once as unallocated evidence - UTC hour/day aggregates are never modified or deleted - Migration cannot manufacture local precision from historical UTC data Retention policy (issue #93): - Raw three-minute samples are pruned after 14 days (spec §3.4, ST-5). - Hour observations and day aggregates are retained indefinitely. -- Local-day summaries are persisted at collection time and retained - indefinitely. They survive raw-sample pruning because they depend on - hour observations, not on raw samples. +- Local-day summaries and unallocated boundaries are persisted at collection + time and retained indefinitely. They survive raw-sample pruning. - After detail expires, aged local summaries remain queryable with their recorded timezone and UTC boundaries. The evidence-availability flag on each history entry indicates whether the underlying raw detail is @@ -28,7 +27,8 @@ Retention policy (issue #93): """ import sqlite3 from dataclasses import dataclass -from datetime import datetime, timedelta, timezone +from datetime import date, datetime, time, timedelta, timezone +from zoneinfo import ZoneInfo @dataclass(frozen=True) @@ -39,11 +39,15 @@ class LocalDaySummary: tz_offset: str # e.g. "+05:30" utc_start: str # ISO 8601 UTC: the local midnight that starts this day utc_end: str # ISO 8601 UTC: the local midnight that ends this day - bytes_written: int - bytes_read: int + bytes_written: int | None + bytes_read: int | None coverage: float # known / (known + unknown) in UTC hours sample_count: int complete: bool # day's UTC range fully covered by hour observations + activity_seconds: int = 0 + activity_intervals: int = 1 + activity_incomplete: bool = False + activity_precision: str = "measured" @dataclass(frozen=True) @@ -60,13 +64,20 @@ class LocalDayHistoryEntry: tz_offset: str utc_start: str utc_end: str - bytes_written: int - bytes_read: int + bytes_written: int | None + bytes_read: int | None coverage: float sample_count: int complete: bool detail_available: bool # True if raw samples for this day are within 14-day retention derived_from_surviving: bool # True if derived from hour observations, not raw samples + shared_bytes_written: int = 0 + shared_bytes_read: int = 0 + unallocated_bytes_written: int = 0 + unallocated_bytes_read: int = 0 + shared_evidence_count: int = 0 + unallocated_evidence_count: int = 0 + activity_state: str = "unavailable" @classmethod def from_summary( @@ -89,21 +100,32 @@ class LocalDayHistoryEntry: complete=summary["complete"], detail_available=detail_available, derived_from_surviving=derived_from_surviving, + shared_bytes_written=summary.get("shared_bytes_written", 0), + shared_bytes_read=summary.get("shared_bytes_read", 0), + unallocated_bytes_written=summary.get("unallocated_bytes_written", 0), + unallocated_bytes_read=summary.get("unallocated_bytes_read", 0), + shared_evidence_count=summary.get("shared_evidence_count", 0), + unallocated_evidence_count=summary.get("unallocated_evidence_count", 0), + activity_state=summary.get("activity_state", "unavailable"), ) -def _local_midnight_utc(dt: datetime, tz_name: str) -> datetime: - """Compute the UTC time of the local midnight that contains *dt*. - - Returns the most recent local midnight in UTC. For example, if *dt* - is 2026-09-01T17:35:00+05:30 (i.e. 12:05 UTC) this returns - 2026-08-31T18:30:00+00:00 (2026-09-01 00:00 in +05:30). - """ - from zoneinfo import ZoneInfo - local_tz = ZoneInfo(tz_name) - local_dt = dt.astimezone(local_tz) - local_midnight = local_dt.replace(hour=0, minute=0, second=0, microsecond=0) - return local_midnight.astimezone(timezone.utc) +def local_day_boundaries(local_date: str, tz_name: str) -> tuple[datetime, datetime, str]: + """Return recorded UTC boundaries and midnight offset for one local date.""" + day = date.fromisoformat(local_date) + zone = ZoneInfo(tz_name) + start_local = datetime.combine(day, time.min, tzinfo=zone) + end_local = datetime.combine(day + timedelta(days=1), time.min, tzinfo=zone) + start_utc = start_local.astimezone(timezone.utc) + end_utc = end_local.astimezone(timezone.utc) + offset = start_local.utcoffset() + offset_seconds = int(offset.total_seconds()) + sign = "+" if offset_seconds >= 0 else "-" + offset_seconds = abs(offset_seconds) + offset_text = "%s%02d:%02d" % ( + sign, offset_seconds // 3600, (offset_seconds % 3600) // 60 + ) + return start_utc, end_utc, offset_text def derive_local_day_summary( @@ -111,108 +133,78 @@ def derive_local_day_summary( tz_name: str, clock_now: datetime, ) -> LocalDaySummary | None: - """Derive a local-day summary from UTC hour observations. + """Build local-day coverage while preserving measured interval totals. - Computes the UTC boundaries of the current local day, queries the - UTC hours overlapping that range, and aggregates read/write volumes. - Midnight-spanning hours are retained once as shared evidence. - - Returns None if no hours exist for the local day. + Hour rows support the established usage-habit coverage contract. They + cannot allocate a read/write delta to a half-hour local boundary, so + local activity totals come only from measured sample intervals. """ - from zoneinfo import ZoneInfo - from .tz_util import get_tz_offset_str - - local_tz = ZoneInfo(tz_name) - local_dt = clock_now.astimezone(local_tz) - local_date = local_dt.strftime("%Y-%m-%d") - tz_offset_str = get_tz_offset_str(clock_now, tz_name) - - # UTC boundaries of this local day - utc_start = _local_midnight_utc(clock_now, tz_name) - next_local = local_dt + timedelta(days=1) - utc_end = next_local.replace(hour=0, minute=0, second=0, microsecond=0).astimezone(timezone.utc) - + local_dt = clock_now.astimezone(ZoneInfo(tz_name)) + local_date = local_dt.date().isoformat() + utc_start, utc_end, tz_offset_str = local_day_boundaries(local_date, tz_name) utc_start_iso = utc_start.isoformat() utc_end_iso = utc_end.isoformat() - # Query UTC hours overlapping the local day - cursor = conn.execute( - "SELECT hour, bytes_written_delta, bytes_read_delta, sample_count, " - " active_seconds, idle_seconds, powered_off_seconds, unknown_seconds " - "FROM hour_observations " - "WHERE hour >= ? AND hour < ? " - "ORDER BY hour", + rows = conn.execute( + "SELECT hour, sample_count, active_seconds, idle_seconds, " + " powered_off_seconds, unknown_seconds " + "FROM hour_observations WHERE hour >= ? AND hour < ? ORDER BY hour", (utc_start_iso, utc_end_iso), - ) - rows = cursor.fetchall() + ).fetchall() - # Check for a midnight-spanning UTC hour before utc_start. - # When the local midnight falls inside a UTC hour (e.g. UTC+5:30 - # where local midnight is 18:30 UTC), the hour 18:00 straddles the - # boundary. The main query (hour >= utc_start) excludes it because - # 18:00 < 18:30, so we must include it separately. This does NOT - # double-count: the hour falls outside the query range by - # construction (issue #93). + # A UTC hour can straddle local midnight. Count its existence for the + # coverage summary, but never use its bytes as local-day activity. midnight_hour = utc_start.replace(minute=0, second=0, microsecond=0) midnight_hour_iso = midnight_hour.strftime("%Y-%m-%dT%H:00:00+00:00") - prev_row = conn.execute( - "SELECT hour, bytes_written_delta, bytes_read_delta, sample_count " - "FROM hour_observations WHERE hour = ?", + previous_hour = conn.execute( + "SELECT hour, sample_count FROM hour_observations WHERE hour = ?", (midnight_hour_iso,), ).fetchone() - total_bw = 0 - total_br = 0 total_samples = 0 known_seconds = 0 - unknown_seconds = 0 hour_count = 0 - for row in rows: - total_bw += row[1] or 0 - total_br += row[2] or 0 - total_samples += row[3] or 0 - known_seconds += (row[4] or 0) + (row[5] or 0) + (row[6] or 0) - unknown_seconds += row[7] or 0 + total_samples += row[1] or 0 + known_seconds += (row[2] or 0) + (row[3] or 0) + (row[4] or 0) + hour_count += 1 + if previous_hour is not None and all(row[0] != previous_hour[0] for row in rows): + total_samples += previous_hour[1] or 0 hour_count += 1 - # Include midnight-spanning hour if it exists and is not already - # in the main query results (it won't be, since hour < utc_start). - if prev_row is not None: - # Verify this hour is NOT already counted in the main query - prev_hour_iso = prev_row[0] - already_counted = any(r[0] == prev_hour_iso for r in rows) - if not already_counted: - total_bw += prev_row[1] or 0 - total_br += prev_row[2] or 0 - total_samples += prev_row[3] or 0 - hour_count += 1 - - total_evidenced = known_seconds + unknown_seconds - # Coverage is known seconds as a share of the full local day, - # not just the hours present — gaps reduce coverage - day_seconds = int((utc_end - utc_start).total_seconds()) - coverage = known_seconds / day_seconds if day_seconds > 0 else 0.0 - - # Day is complete only if known (usable) seconds cover the full local day. - # Unknown seconds represent gaps without usable observation evidence. - complete = known_seconds >= day_seconds and hour_count > 0 - - # Return None if no hours exist for this local day if hour_count == 0: return None + day_seconds = int((utc_end - utc_start).total_seconds()) + coverage = known_seconds / day_seconds if day_seconds > 0 else 0.0 + complete = known_seconds >= day_seconds and hour_count > 0 + + existing = conn.execute( + "SELECT bytes_written, bytes_read, activity_seconds, activity_intervals, " + " activity_incomplete, activity_precision " + "FROM local_days WHERE local_date = ? AND tz_name = ?", + (local_date, tz_name), + ).fetchone() + if existing is None: + activity = (0, 0, 0, 0, False, "unavailable") + else: + activity = existing + return LocalDaySummary( local_date=local_date, tz_name=tz_name, tz_offset=tz_offset_str, utc_start=utc_start_iso, utc_end=utc_end_iso, - bytes_written=total_bw, - bytes_read=total_br, + bytes_written=activity[0] or 0, + bytes_read=activity[1] or 0, coverage=coverage, sample_count=total_samples, complete=complete, + activity_seconds=activity[2] or 0, + activity_intervals=activity[3] or 0, + activity_incomplete=bool(activity[4]), + activity_precision=activity[5], ) @@ -230,12 +222,15 @@ def persist_local_day( 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_seconds, activity_intervals, activity_incomplete, activity_precision) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", (summary.local_date, summary.tz_name, summary.tz_offset, summary.utc_start, summary.utc_end, summary.bytes_written, summary.bytes_read, - summary.coverage, summary.sample_count, summary.complete), + summary.coverage, summary.sample_count, summary.complete, + summary.activity_seconds, summary.activity_intervals, + summary.activity_incomplete, summary.activity_precision), ) created = True else: @@ -243,11 +238,15 @@ def persist_local_day( "UPDATE local_days " "SET tz_offset = ?, utc_start = ?, utc_end = ?, " " bytes_written = ?, bytes_read = ?, " - " coverage = ?, sample_count = ?, complete = ? " + " coverage = ?, sample_count = ?, complete = ?, " + " activity_seconds = ?, activity_intervals = ?, " + " activity_incomplete = ?, activity_precision = ? " "WHERE id = ?", (summary.tz_offset, summary.utc_start, summary.utc_end, summary.bytes_written, summary.bytes_read, summary.coverage, summary.sample_count, summary.complete, + summary.activity_seconds, summary.activity_intervals, + summary.activity_incomplete, summary.activity_precision, existing[0]), ) created = False @@ -255,34 +254,371 @@ def persist_local_day( return created +def _utc_datetime(value: str | datetime) -> datetime: + parsed = value if isinstance(value, datetime) else datetime.fromisoformat(value) + if parsed.tzinfo is None: + parsed = parsed.replace(tzinfo=timezone.utc) + return parsed.astimezone(timezone.utc) + + +def _same_monitoring_period( + conn: sqlite3.Connection, + start: datetime, + end: datetime, +) -> bool: + """Return whether one recorded monitoring period contains the interval.""" + for started_at, ended_at in conn.execute( + "SELECT started_at, ended_at FROM monitoring_periods" + ): + period_start = _utc_datetime(started_at) + period_end = _utc_datetime(ended_at) if ended_at else None + if period_start <= start and (period_end is None or end <= period_end): + return True + return False + + +def _upsert_activity_day( + conn: sqlite3.Connection, + local_date: str, + tz_name: str, + *, + bytes_written: int = 0, + bytes_read: int = 0, + seconds: int = 0, + interval_count: int = 0, + incomplete: bool = False, + last_sample_id: int | None = None, +) -> None: + start, end, offset = local_day_boundaries(local_date, tz_name) + conn.execute( + "INSERT INTO local_days " + "(local_date, tz_name, tz_offset, utc_start, utc_end, bytes_written, " + " bytes_read, activity_seconds, activity_intervals, activity_incomplete, " + " activity_precision, last_sample_id) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'measured', ?) " + "ON CONFLICT(local_date, tz_name) DO UPDATE SET " + "bytes_written = CASE WHEN local_days.activity_precision IN ('legacy', 'unavailable') " + " THEN excluded.bytes_written ELSE local_days.bytes_written + excluded.bytes_written END, " + "bytes_read = CASE WHEN local_days.activity_precision IN ('legacy', 'unavailable') " + " THEN excluded.bytes_read ELSE local_days.bytes_read + excluded.bytes_read END, " + "activity_seconds = CASE WHEN local_days.activity_precision IN ('legacy', 'unavailable') " + " THEN excluded.activity_seconds ELSE local_days.activity_seconds + excluded.activity_seconds END, " + "activity_intervals = CASE WHEN local_days.activity_precision IN ('legacy', 'unavailable') " + " 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', " + "last_sample_id = COALESCE(excluded.last_sample_id, local_days.last_sample_id)", + ( + local_date, tz_name, offset, start.isoformat(), end.isoformat(), + bytes_written, bytes_read, seconds, interval_count, incomplete, + last_sample_id, + ), + ) + + +def record_local_activity_interval( + conn: sqlite3.Connection, + previous: dict, + current: dict, + *, + start_sample_id: int, + end_sample_id: int, + segment_id: int | None, +) -> str | None: + """Persist one measured interval as known local activity or once-only evidence. + + Counter differences are allocated only when both sample timezones match, + both samples share one monitoring period, and the interval stays inside one + recorded local date. Every other valid difference remains one unallocated row. + """ + start = _utc_datetime(previous["ts"]) + end = _utc_datetime(current["ts"]) + if end <= start: + return None + + previous_tz = previous.get("local_tz") + current_tz = current.get("local_tz") + if not current_tz: + return None + + from_zone = ZoneInfo(previous_tz) if previous_tz else None + to_zone = ZoneInfo(current_tz) + start_date = start.astimezone(from_zone).date().isoformat() if from_zone else None + end_date = end.astimezone(to_zone).date().isoformat() + if start_date is None: + start_date = end_date + + written = current.get("bytes_written") + 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): + return None + bytes_written = written - prior_written + bytes_read = read - prior_read + if bytes_written < 0 or bytes_read < 0: + for local_date, zone_name in ( + (start_date, previous_tz), (end_date, current_tz), + ): + if zone_name is not None: + _upsert_activity_day( + conn, local_date, zone_name, incomplete=True, + ) + return "counter_discontinuity" + elif previous_tz is None: + reason = "legacy_timezone_unknown" + elif previous_tz != current_tz: + reason = "timezone_change" + elif segment_id is None: + reason = "segment_unknown" + elif not _same_monitoring_period(conn, start, end): + reason = "monitoring_period" + elif start_date != end_date: + reason = "local_midnight" + else: + already_accounted = conn.execute( + "SELECT last_sample_id FROM local_days " + "WHERE local_date = ? AND tz_name = ?", + (end_date, current_tz), + ).fetchone() + if ( + already_accounted is not None + and already_accounted[0] is not None + and already_accounted[0] >= end_sample_id + ): + return "known" + seconds = int((end - start).total_seconds()) + _upsert_activity_day( + conn, + end_date, + current_tz, + bytes_written=bytes_written, + bytes_read=bytes_read, + seconds=seconds, + interval_count=1, + last_sample_id=end_sample_id, + ) + local_day_id = conn.execute( + "SELECT id FROM local_days WHERE local_date = ? AND tz_name = ?", + (end_date, current_tz), + ).fetchone()[0] + conn.execute( + "INSERT INTO local_day_segment_totals " + "(local_day_id, segment_id, bytes_written, bytes_read, " + " activity_seconds, activity_intervals) " + "VALUES (?, ?, ?, ?, ?, 1) " + "ON CONFLICT(local_day_id, segment_id) DO UPDATE SET " + "bytes_written = local_day_segment_totals.bytes_written + excluded.bytes_written, " + "bytes_read = local_day_segment_totals.bytes_read + excluded.bytes_read, " + "activity_seconds = local_day_segment_totals.activity_seconds + excluded.activity_seconds, " + "activity_intervals = local_day_segment_totals.activity_intervals + 1", + (local_day_id, segment_id, bytes_written, bytes_read, seconds), + ) + return "known" + + 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 (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + ( + start_sample_id, end_sample_id, start_date, end_date, + previous_tz, current_tz, start.isoformat(), end.isoformat(), + bytes_written, bytes_read, reason, segment_id, + ), + ) + + affected_days: set[tuple[str, str]] = set() + if previous_tz is not None: + affected_days.add((start_date, previous_tz)) + affected_days.add((end_date, current_tz)) + if previous_tz is not None and previous_tz == current_tz and start_date < end_date: + day = date.fromisoformat(start_date) + timedelta(days=1) + last = date.fromisoformat(end_date) + while day < last: + affected_days.add((day.isoformat(), current_tz)) + day += timedelta(days=1) + for local_date, zone_name in affected_days: + _upsert_activity_day( + conn, local_date, zone_name, incomplete=True, + ) + return reason + + +def mark_local_activity_gap( + conn: sqlite3.Connection, + current_ts: str, + current_tz: str, + previous: dict | None = None, +) -> None: + """Mark dates around a controller boundary as missing local activity.""" + current = _utc_datetime(current_ts) + current_date = current.astimezone(ZoneInfo(current_tz)).date().isoformat() + _upsert_activity_day(conn, current_date, current_tz, incomplete=True) + if previous is None or not previous.get("local_tz"): + return + previous_at = _utc_datetime(previous["ts"]) + previous_tz = previous["local_tz"] + previous_date = previous_at.astimezone(ZoneInfo(previous_tz)).date().isoformat() + _upsert_activity_day(conn, previous_date, previous_tz, incomplete=True) + + def query_local_day_summary( conn: sqlite3.Connection, local_date: str, + tz_name: str | None = None, + clock_now: datetime | None = None, ) -> dict | None: """Query a stored local-day summary by date. Returns a dict with the summary fields, or None if no entry exists. """ + tz_filter = " AND tz_name = ?" if tz_name else "" + params = (local_date, tz_name) if tz_name else (local_date,) row = conn.execute( "SELECT local_date, tz_name, tz_offset, utc_start, utc_end, " - " bytes_written, bytes_read, coverage, sample_count, complete " - "FROM local_days WHERE local_date = ? " - "ORDER BY id DESC LIMIT 1", - (local_date,), + " bytes_written, bytes_read, coverage, sample_count, complete, " + " activity_seconds, activity_intervals, activity_incomplete, " + " activity_precision " + "FROM local_days WHERE local_date = ?" + tz_filter + + " ORDER BY id DESC LIMIT 1", + params, ).fetchone() if row is None: return None + ( + recorded_date, + recorded_tz, + tz_offset, + utc_start, + utc_end, + bytes_written, + bytes_read, + coverage, + sample_count, + complete, + activity_seconds, + activity_intervals, + activity_incomplete, + activity_precision, + ) = row + + evidence = conn.execute( + "SELECT COUNT(*), " + " COALESCE(SUM(CASE WHEN reason = 'local_midnight' THEN 1 ELSE 0 END), 0), " + " COALESCE(SUM(CASE WHEN reason != 'local_midnight' THEN 1 ELSE 0 END), 0), " + " COALESCE(SUM(CASE WHEN reason = 'local_midnight' " + " THEN bytes_written ELSE 0 END), 0), " + " COALESCE(SUM(CASE WHEN reason = 'local_midnight' " + " THEN bytes_read ELSE 0 END), 0), " + " COALESCE(SUM(CASE WHEN reason != 'local_midnight' " + " THEN bytes_written ELSE 0 END), 0), " + " COALESCE(SUM(CASE WHEN reason != 'local_midnight' " + " THEN bytes_read ELSE 0 END), 0) " + "FROM local_day_unallocated_evidence " + "WHERE (start_local_date = ? AND start_tz_name = ?) " + " OR (end_local_date = ? AND end_tz_name = ?) " + " OR (start_tz_name = ? AND end_tz_name = ? " + " AND start_local_date < ? AND end_local_date > ?)", + ( + recorded_date, recorded_tz, recorded_date, recorded_tz, + recorded_tz, recorded_tz, recorded_date, recorded_date, + ), + ).fetchone() + ( + evidence_count, + shared_evidence_count, + unallocated_evidence_count, + shared_bytes_written, + shared_bytes_read, + unallocated_bytes_written, + unallocated_bytes_read, + ) = evidence + + local_day_id = conn.execute( + "SELECT id FROM local_days WHERE local_date = ? AND tz_name = ?", + (recorded_date, recorded_tz), + ).fetchone()[0] + segments = [ + { + "segment_id": segment_id, + "bytes_written": segment_written, + "bytes_read": segment_read, + "activity_seconds": segment_seconds, + "activity_intervals": segment_intervals, + } + for ( + segment_id, + segment_written, + segment_read, + segment_seconds, + segment_intervals, + ) in conn.execute( + "SELECT segment_id, bytes_written, bytes_read, activity_seconds, " + " activity_intervals " + "FROM local_day_segment_totals WHERE local_day_id = ? " + "ORDER BY segment_id", + (local_day_id,), + ) + ] + + interval_count = activity_intervals or 0 + has_unallocated = evidence_count > 0 + has_known = activity_precision == "measured" and interval_count > 0 + if not has_known and not has_unallocated: + activity_state = "incomplete" if bool(activity_incomplete) else "unavailable" + else: + now = clock_now or datetime.now(timezone.utc) + is_current_day = ( + now.astimezone(ZoneInfo(recorded_tz)).date().isoformat() + == recorded_date + ) + day_seconds = int( + (_utc_datetime(utc_end) - _utc_datetime(utc_start)).total_seconds() + ) + interval_seconds = activity_seconds or 0 + full_activity_day = interval_seconds >= day_seconds or ( + interval_count == 0 and bool(complete) + ) + incomplete = ( + bool(activity_incomplete) or has_unallocated or not full_activity_day + ) + if is_current_day: + activity_state = "so_far" + elif incomplete: + activity_state = "incomplete" + elif has_known and (bytes_written or 0) == 0 and (bytes_read or 0) == 0: + activity_state = "zero" + else: + activity_state = "complete" + + exact = activity_precision == "measured" and has_known return { - "local_date": row[0], - "tz_name": row[1], - "tz_offset": row[2], - "utc_start": row[3], - "utc_end": row[4], - "bytes_written": row[5], - "bytes_read": row[6], - "coverage": row[7], - "sample_count": row[8], - "complete": bool(row[9]), + "local_date": recorded_date, + "tz_name": recorded_tz, + "tz_offset": tz_offset, + "utc_start": utc_start, + "utc_end": utc_end, + "bytes_written": bytes_written if exact else None, + "bytes_read": bytes_read if exact else None, + "shared_evidence_count": shared_evidence_count, + "unallocated_evidence_count": unallocated_evidence_count, + "shared_bytes_written": shared_bytes_written, + "shared_bytes_read": shared_bytes_read, + "unallocated_bytes_written": unallocated_bytes_written, + "unallocated_bytes_read": unallocated_bytes_read, + "coverage": coverage, + "sample_count": sample_count, + "complete": bool(complete), + "activity_seconds": activity_seconds or 0, + "activity_intervals": interval_count, + "activity_incomplete": bool(activity_incomplete), + "activity_precision": activity_precision, + "activity_state": activity_state, + "segments": segments, } @@ -296,7 +632,7 @@ def query_current_local_day( local_tz = ZoneInfo(tz_name) local_dt = clock_now.astimezone(local_tz) local_date = local_dt.strftime("%Y-%m-%d") - return query_local_day_summary(conn, local_date) + return query_local_day_summary(conn, local_date, tz_name, clock_now) def _is_detail_available( @@ -347,27 +683,20 @@ def query_local_day_history( """ rows = conn.execute( "SELECT local_date, tz_name, tz_offset, utc_start, utc_end, " - " bytes_written, bytes_read, coverage, sample_count, complete " + " bytes_written, bytes_read, coverage, sample_count, complete, " + " activity_seconds, activity_intervals, activity_incomplete, " + " activity_precision " "FROM local_days " "WHERE local_date >= ? AND local_date <= ? " - "ORDER BY local_date", + "ORDER BY local_date, id", (start_date, end_date), ).fetchall() entries = [] for row in rows: - summary = { - "local_date": row[0], - "tz_name": row[1], - "tz_offset": row[2], - "utc_start": row[3], - "utc_end": row[4], - "bytes_written": row[5], - "bytes_read": row[6], - "coverage": row[7], - "sample_count": row[8], - "complete": row[9], - } + summary = query_local_day_summary(conn, row[0], row[1], now) + if summary is None: + continue detail_available = _is_detail_available(conn, row[3], row[4], now) entries.append( LocalDayHistoryEntry.from_summary( diff --git a/src/fenris/store.py b/src/fenris/store.py index 1c70aa1..806beea 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 = 4 +SCHEMA_VERSION = 5 # Packaged default placement (spec §8.3). The config may override it, but a @@ -106,7 +106,8 @@ def _create_schema(conn: sqlite3.Connection): bytes_written INTEGER, bytes_read INTEGER, critical_warning INTEGER, - segment_id INTEGER + segment_id INTEGER, + local_tz TEXT ) """) @@ -208,10 +209,18 @@ def _create_schema(conn: sqlite3.Connection): coverage REAL DEFAULT 0.0, sample_count INTEGER DEFAULT 0, complete BOOLEAN DEFAULT 0, + 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 'measured', + last_sample_id INTEGER, UNIQUE(local_date, tz_name) ) """) + _create_local_day_shared_evidence(conn) + _create_local_day_segment_totals(conn) + # Metadata table for store state (e.g., legacy import marker) conn.execute(""" CREATE TABLE IF NOT EXISTS store_metadata ( @@ -234,6 +243,52 @@ def _create_pending_publications(conn: sqlite3.Connection) -> None: """) +def _create_local_day_shared_evidence(conn: sqlite3.Connection) -> None: + """Create once-only local activity evidence that cannot be day-allocated.""" + conn.execute(""" + CREATE TABLE IF NOT EXISTS local_day_unallocated_evidence ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + start_sample_id INTEGER NOT NULL, + end_sample_id INTEGER NOT NULL, + start_local_date TEXT NOT NULL, + end_local_date TEXT NOT NULL, + start_tz_name TEXT, + end_tz_name TEXT, + started_at TEXT NOT NULL, + ended_at TEXT NOT NULL, + bytes_written INTEGER NOT NULL DEFAULT 0, + bytes_read INTEGER NOT NULL DEFAULT 0, + reason TEXT NOT NULL, + segment_id INTEGER, + UNIQUE(start_sample_id, end_sample_id) + ) + """) + conn.execute( + "CREATE INDEX IF NOT EXISTS local_day_evidence_start " + "ON local_day_unallocated_evidence(start_local_date, start_tz_name)" + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS local_day_evidence_end " + "ON local_day_unallocated_evidence(end_local_date, end_tz_name)" + ) + + +def _create_local_day_segment_totals(conn: sqlite3.Connection) -> None: + """Retain the controller-segment provenance behind known day totals.""" + conn.execute(""" + CREATE TABLE IF NOT EXISTS local_day_segment_totals ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + local_day_id INTEGER NOT NULL, + segment_id INTEGER NOT NULL, + bytes_written INTEGER NOT NULL DEFAULT 0, + bytes_read INTEGER NOT NULL DEFAULT 0, + activity_seconds INTEGER NOT NULL DEFAULT 0, + activity_intervals INTEGER NOT NULL DEFAULT 0, + UNIQUE(local_day_id, segment_id) + ) + """) + + def _apply_migrations(conn: sqlite3.Connection, current_version: int): """Apply forward-only migrations from current_version to SCHEMA_VERSION. @@ -293,6 +348,39 @@ def _apply_migrations(conn: sqlite3.Connection, current_version: int): _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'" + ).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 + def migrate_to_latest(store_path: Path) -> int: """Apply forward-only migrations to bring the store to SCHEMA_VERSION. diff --git a/src/fenris/tui.py b/src/fenris/tui.py index 1f59daa..722730c 100644 --- a/src/fenris/tui.py +++ b/src/fenris/tui.py @@ -1709,7 +1709,9 @@ class FenrisTuiApp(App): ) tz_name = detect_system_tz() if self._browse_date is not None: - local = query_local_day_summary(conn, self._browse_date) + local = query_local_day_summary( + conn, self._browse_date, clock_now=self._clock_now + ) else: local = query_current_local_day(conn, self._clock_now, tz_name) except Exception: @@ -1720,7 +1722,10 @@ class FenrisTuiApp(App): if local is None: date = self._browse_date or self._clock_now.astimezone().date().isoformat() - widget.update("%s · local-day evidence unavailable\nW unavailable · R unavailable" % date) + widget.update( + "%s · local-day evidence unavailable\nW unavailable · R unavailable" + % date + ) widget.display = True main_grid.remove_class("local-day") return @@ -1728,23 +1733,45 @@ class FenrisTuiApp(App): main_grid.add_class("local-day") widget.styles.display = "block" + tz_display = "%s %s" % (local["tz_name"], local["tz_offset"]) + state = local["activity_state"] + state_labels = { + "so_far": "totals so far", + "incomplete": "incomplete", + "complete": "complete", + "zero": "measured zero", + "unavailable": "local activity unavailable", + } + label = state_labels.get(state, "local activity unavailable") bw = local["bytes_written"] br = local["bytes_read"] - partial = "" if local["complete"] else " · totals so far" - tz_display = "%s %s" % (local["tz_name"], local["tz_offset"]) - - text = ( - "[bold]%s[/bold] · %s%s\n" - " W %.3f GB · R %.3f GB · %.0f%% coverage" - % ( - local["local_date"], - tz_display, - partial, - bw / 1e9, - br / 1e9, - local["coverage"] * 100, + if bw is None or br is None: + totals = "W unavailable · R unavailable" + else: + totals = "W %.3f GB known · R %.3f GB known" % ( + bw / 1e9, br / 1e9 ) - ) + lines = [ + "[bold]%s[/bold] · %s · %s" % ( + local["local_date"], tz_display, label, + ), + totals, + ] + if local["shared_evidence_count"]: + lines.append( + "shared at midnight W %.3f GB · R %.3f GB" % ( + local["shared_bytes_written"] / 1e9, + local["shared_bytes_read"] / 1e9, + ) + ) + if local["unallocated_evidence_count"]: + lines.append( + "unallocated W %.3f GB · R %.3f GB" % ( + local["unallocated_bytes_written"] / 1e9, + local["unallocated_bytes_read"] / 1e9, + ) + ) + text = "\n".join(lines) widget.update(text) def _on_graph_drill(self, day: str) -> None: diff --git a/tests/test_collector_history_tracer.py b/tests/test_collector_history_tracer.py index b6d3479..5678228 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 == 3 # v1→v2→v3→v4 + assert steps == 4 # v1→v2→v3→v4→v5 # Verify data preserved conn = sqlite3.connect(str(db)) diff --git a/tests/test_collector_tracer.py b/tests/test_collector_tracer.py index f41e385..cb2412d 100644 --- a/tests/test_collector_tracer.py +++ b/tests/test_collector_tracer.py @@ -250,6 +250,8 @@ def test_store_initialization(config_fixture: Dict[str, Any]): expected_tables.add("store_metadata") expected_tables.add("local_days") expected_tables.add("pending_publications") + expected_tables.add("local_day_unallocated_evidence") + expected_tables.add("local_day_segment_totals") assert expected_tables == tables conn.close() diff --git a/tests/test_issue_93.py b/tests/test_issue_93.py index c6e5ddd..105a66c 100644 --- a/tests/test_issue_93.py +++ b/tests/test_issue_93.py @@ -277,7 +277,7 @@ class TestTimezonePreservation: now = datetime(2026, 10, 1, 12, 0, 0, tzinfo=timezone.utc) old = LocalDaySummary( - local_date="2026-09-16", tz_name="US/Eastern", tz_offset="-04:00", + local_date="2026-09-16", tz_name="America/New_York", tz_offset="-04:00", utc_start="2026-09-16T04:00:00+00:00", utc_end="2026-09-17T04:00:00+00:00", bytes_written=5000, bytes_read=2000, @@ -288,7 +288,7 @@ class TestTimezonePreservation: entries = query_local_day_history(conn, "2026-09-16", "2026-09-16", now) assert len(entries) == 1 entry = entries[0] - assert entry.tz_name == "US/Eastern" + assert entry.tz_name == "America/New_York" assert entry.tz_offset == "-04:00" assert entry.utc_start == "2026-09-16T04:00:00+00:00" assert entry.detail_available is False @@ -389,10 +389,10 @@ class TestIdempotencySafety: # --------------------------------------------------------------------------- class TestVolumeConservation: - """AC6: Summary volume plus unallocated evidence conserves each delta once.""" + """AC6: UTC-hour totals never become fabricated local-day totals.""" - def test_midnight_spanning_not_double_counted(self, tmp_path): - """Midnight-spanning interval appears once, not in both days.""" + def test_hourly_bytes_are_unavailable_at_local_precision(self, tmp_path): + """UTC-day hour bytes cannot be split across local-day boundaries.""" conn = init_store(tmp_path / "obs.db") # Insert hour observations that span midnight @@ -409,13 +409,10 @@ class TestVolumeConservation: summary_17 = derive_local_day_summary(conn, "UTC", clock_17) assert summary_17 is not None - # Each day should only have its own hour's bytes - # Day 16 has the 23:00 hour (100 bw) - # Day 17 has the 00:00 hour (200 bw) - assert summary_16.bytes_written == 100 - assert summary_17.bytes_written == 200 - # Total is conserved: 100 + 200 = 300 - assert summary_16.bytes_written + summary_17.bytes_written == 300 + assert summary_16.bytes_written == 0 + assert summary_17.bytes_written == 0 + assert summary_16.activity_precision == "unavailable" + assert summary_17.activity_precision == "unavailable" conn.close() diff --git a/tests/test_issue_97.py b/tests/test_issue_97.py index ef08b29..16e5730 100644 --- a/tests/test_issue_97.py +++ b/tests/test_issue_97.py @@ -285,7 +285,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] == 4 + assert migrated.execute("PRAGMA user_version").fetchone()[0] == 5 assert migrated.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0 finally: migrated.close() diff --git a/tests/test_local_day.py b/tests/test_local_day.py index 2aa2a1d..35b6e09 100644 --- a/tests/test_local_day.py +++ b/tests/test_local_day.py @@ -25,6 +25,7 @@ from fenris.store import init_store, SCHEMA_VERSION from fenris.local_day import ( derive_local_day_summary, persist_local_day, + record_local_activity_interval, query_local_day_summary, query_current_local_day, LocalDaySummary, @@ -135,13 +136,56 @@ class TestSchemaMigration: assert count == 0 conn.close() + def test_v4_totals_keep_values_but_are_not_treated_as_local_precision( + self, tmp_path, + ): + db = tmp_path / "v4.db" + conn = sqlite3.connect(db) + conn.execute( + "CREATE TABLE samples (id INTEGER PRIMARY KEY, ts TEXT, device TEXT)" + ) + conn.execute( + "CREATE TABLE pending_publications " + "(id INTEGER PRIMARY KEY, sample_ts TEXT, payload TEXT)" + ) + conn.execute( + "CREATE TABLE local_days (" + "id INTEGER PRIMARY KEY, 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, sample_count INTEGER DEFAULT 0, " + "complete BOOLEAN DEFAULT 0, UNIQUE(local_date, tz_name))" + ) + conn.execute( + "INSERT INTO local_days VALUES " + "(1, '2026-09-01', 'UTC', '+00:00', " + "'2026-09-01T00:00:00+00:00', '2026-09-02T00:00:00+00:00', " + "5120000, 2048000, 0.5, 10, 0)" + ) + conn.execute("PRAGMA user_version=4") + conn.commit() + conn.close() + + conn = init_store(db) + stored = conn.execute( + "SELECT bytes_written, activity_precision FROM local_days" + ).fetchone() + assert stored == (5_120_000, "legacy") + visible = query_local_day_summary(conn, "2026-09-01", "UTC") + assert visible["bytes_written"] is None + assert visible["activity_state"] == "unavailable" + assert "local_tz" in { + row[1] for row in conn.execute("PRAGMA table_info(samples)") + } + conn.close() + # --------------------------------------------------------------------------- -# Local-day derivation from UTC hours +# Local-day coverage derivation from UTC hours # --------------------------------------------------------------------------- class TestDeriveLocalDay: - """Derive local-day summaries from UTC hour observations.""" + """Derive coverage from UTC hours without allocating local-day bytes.""" def test_utc_plus_zero(self, tmp_path): """UTC timezone: local day boundaries = UTC day boundaries.""" @@ -158,8 +202,9 @@ class TestDeriveLocalDay: assert summary.local_date == "2026-09-01" assert summary.tz_name == "UTC" assert summary.tz_offset == "+00:00" - assert summary.bytes_written == 300 - assert summary.bytes_read == 130 + assert summary.bytes_written == 0 + assert summary.bytes_read == 0 + assert summary.activity_precision == "unavailable" assert summary.complete is False # not all 24 hours covered conn.close() @@ -179,12 +224,13 @@ class TestDeriveLocalDay: assert summary.local_date == "2026-09-01" assert summary.tz_name == "Asia/Kolkata" assert summary.tz_offset == "+05:30" - assert summary.bytes_written == 300 - assert summary.bytes_read == 130 + assert summary.bytes_written == 0 + assert summary.bytes_read == 0 + assert summary.activity_precision == "unavailable" conn.close() - def test_midnight_spanning_hour_included(self, tmp_path): - """UTC hour straddling local midnight is included in the local day.""" + def test_hourly_bytes_cannot_be_allocated_to_local_midnight(self, tmp_path): + """UTC-hour totals do not establish exact local-day byte totals.""" conn = init_store(tmp_path / "obs.db") # For UTC+5:30, local day 2026-09-01 starts at 2026-08-31T18:30 UTC @@ -196,9 +242,9 @@ class TestDeriveLocalDay: summary = derive_local_day_summary(conn, "Asia/Kolkata", clock) assert summary is not None - # The midnight-spanning hour (18:00) is included - assert summary.bytes_written == 300 - assert summary.bytes_read == 130 + assert summary.bytes_written == 0 + assert summary.bytes_read == 0 + assert summary.activity_precision == "unavailable" conn.close() def test_no_hours_returns_none(self, tmp_path): @@ -387,6 +433,411 @@ class TestCollectorIntegration: assert row[1] == 30 * 512000 # reads conn.close() + def test_interval_crossing_utc_hour_stays_wholly_in_local_day( + self, tmp_path, sysfs_tree, monkeypatch, + ): + """A measured interval inside one local date keeps its full delta.""" + from fenris.collector import run_collection + + monkeypatch.setenv("TZ", "Asia/Kolkata") + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + first_at = datetime(2026, 9, 1, 12, 59, tzinfo=timezone.utc) + second_at = datetime(2026, 9, 1, 13, 1, tzinfo=timezone.utc) + + first = run_collection( + _make_smartctl(10_000_000, 8_000_000), + sysfs_nvme, + cfg, + _Clock(first_at), + ) + second = run_collection( + _make_smartctl(10_000_010, 8_000_004), + sysfs_nvme, + cfg, + _Clock(second_at), + ) + + assert first["ok"] and second["ok"] + conn = sqlite3.connect(store) + summary = query_local_day_summary(conn, "2026-09-01") + + assert summary is not None + assert summary["bytes_written"] == 5_120_000 + assert summary["bytes_read"] == 2_048_000 + assert summary["segments"] == [{ + "segment_id": 1, + "bytes_written": 5_120_000, + "bytes_read": 2_048_000, + "activity_seconds": 120, + "activity_intervals": 1, + }] + segment_totals = conn.execute( + "SELECT segment_id, bytes_written, bytes_read, activity_intervals " + "FROM local_day_segment_totals" + ).fetchone() + assert segment_totals == (1, 5_120_000, 2_048_000, 1) + + samples = conn.execute( + "SELECT id, ts, bytes_written, bytes_read, local_tz " + "FROM samples ORDER BY id" + ).fetchall() + first_sample, second_sample = samples + assert record_local_activity_interval( + conn, + dict(zip(("id", "ts", "bytes_written", "bytes_read", "local_tz"), first_sample)), + dict(zip(("id", "ts", "bytes_written", "bytes_read", "local_tz"), second_sample)), + start_sample_id=first_sample[0], + end_sample_id=second_sample[0], + segment_id=1, + ) == "known" + repeated = query_local_day_summary(conn, "2026-09-01", "Asia/Kolkata") + assert repeated["bytes_written"] == 5_120_000 + assert repeated["bytes_read"] == 2_048_000 + assert conn.execute( + "SELECT SUM(bytes_written), SUM(bytes_read), SUM(activity_intervals) " + "FROM local_day_segment_totals" + ).fetchone() == (5_120_000, 2_048_000, 1) + conn.close() + + @pytest.mark.parametrize( + ( + "first_at", "second_at", "local_date", "expected_hours", + "expected_start", "expected_end", "expected_offset", + ), + [ + ( + datetime(2026, 3, 8, 6, 30, tzinfo=timezone.utc), + datetime(2026, 3, 8, 7, 30, tzinfo=timezone.utc), + "2026-03-08", 23, + "2026-03-08T05:00:00+00:00", + "2026-03-09T04:00:00+00:00", "-05:00", + ), + ( + datetime(2026, 11, 1, 5, 30, tzinfo=timezone.utc), + datetime(2026, 11, 1, 6, 30, tzinfo=timezone.utc), + "2026-11-01", 25, + "2026-11-01T04:00:00+00:00", + "2026-11-02T05:00:00+00:00", "-04:00", + ), + ], + ids=("spring-forward", "fall-back-repeated-hour"), + ) + def test_collector_retains_exact_interval_on_dst_local_day( + self, + tmp_path, + sysfs_tree, + monkeypatch, + first_at, + second_at, + local_date, + expected_hours, + expected_start, + expected_end, + expected_offset, + ): + """Collection preserves exact activity on short and long local days.""" + from fenris.collector import run_collection + + zone = "America/New_York" + monkeypatch.setenv("TZ", zone) + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + assert run_collection( + _make_smartctl(10_000_000, 8_000_000), + sysfs_nvme, + cfg, + _Clock(first_at), + )["ok"] + assert run_collection( + _make_smartctl(10_000_001, 8_000_002), + sysfs_nvme, + cfg, + _Clock(second_at), + )["ok"] + + conn = sqlite3.connect(store) + summary = query_local_day_summary( + conn, local_date, zone, second_at + timedelta(minutes=1) + ) + assert summary is not None + assert summary["bytes_written"] == 512_000 + assert summary["bytes_read"] == 1_024_000 + assert summary["activity_seconds"] == 3_600 + assert summary["utc_start"] == expected_start + assert summary["utc_end"] == expected_end + assert summary["tz_offset"] == expected_offset + assert ( + datetime.fromisoformat(summary["utc_end"]) + - datetime.fromisoformat(summary["utc_start"]) + ).total_seconds() == expected_hours * 3_600 + conn.close() + + @pytest.mark.asyncio + async def test_local_midnight_interval_is_shared_once( + self, tmp_path, sysfs_tree, monkeypatch, + ): + """A real counter interval crossing local midnight stays unallocated.""" + from fenris.collector import run_collection + + monkeypatch.setenv("TZ", "Asia/Kolkata") + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + first_at = datetime(2026, 9, 1, 18, 25, tzinfo=timezone.utc) + second_at = datetime(2026, 9, 1, 18, 35, tzinfo=timezone.utc) + + first = run_collection( + _make_smartctl(10_000_000, 8_000_000), + sysfs_nvme, + cfg, + _Clock(first_at), + ) + second = run_collection( + _make_smartctl(10_000_100, 8_000_060), + sysfs_nvme, + cfg, + _Clock(second_at), + ) + assert first["ok"] and second["ok"] + + conn = sqlite3.connect(store) + evidence = conn.execute( + "SELECT COUNT(*), SUM(bytes_written), SUM(bytes_read) " + "FROM local_day_unallocated_evidence" + ).fetchone() + assert evidence == (1, 51_200_000, 30_720_000) + assert conn.execute( + "SELECT SUM(bytes_written), SUM(bytes_read) FROM local_days" + ).fetchone() == (0, 0) + + for local_date in ("2026-09-01", "2026-09-02"): + summary = query_local_day_summary( + conn, local_date, "Asia/Kolkata", second_at + ) + assert summary is not None + assert summary["bytes_written"] is None + assert summary["shared_bytes_written"] == 51_200_000 + assert summary["shared_bytes_read"] == 30_720_000 + assert summary["activity_state"] in ("incomplete", "so_far") + conn.close() + + from fenris.status import read_status + from fenris.tui import FenrisTuiApp + + app = FenrisTuiApp(store_path=Path(store), refresh_interval_s=999) + async with app.run_test(size=(100, 32)): + app._browse_date = "2026-09-02" + with read_status( + Path(store), second_at, query_services=False + ) as (reader, _): + assert reader is not None + app._render_local_day(reader) + visible = str(app.query_one("#local-day").render()) + assert "2026-09-02" in visible + assert "Asia/Kolkata +05:30" in visible + assert "W unavailable" in visible + assert "shared at midnight W 0.051 GB" in visible + + def test_synthetic_shared_bytes_conserve_one_hundred_bytes(self, tmp_path): + """Shared evidence conserves byte values before display rounding.""" + conn = init_store(tmp_path / "obs.db") + first_at = datetime(2026, 9, 1, 18, 25, tzinfo=timezone.utc) + second_at = datetime(2026, 9, 3, 18, 35, tzinfo=timezone.utc) + ensure_period_open(conn, first_at) + conn.execute( + "INSERT INTO samples (ts, device, bytes_written, bytes_read, local_tz) " + "VALUES (?, ?, ?, ?, ?)", + (first_at.isoformat(), "/dev/nvme0", 1_000, 2_000, "Asia/Kolkata"), + ) + first_id = conn.execute("SELECT last_insert_rowid()").fetchone()[0] + conn.execute( + "INSERT INTO samples (ts, device, bytes_written, bytes_read, local_tz) " + "VALUES (?, ?, ?, ?, ?)", + (second_at.isoformat(), "/dev/nvme0", 1_100, 2_050, "Asia/Kolkata"), + ) + second_id = conn.execute("SELECT last_insert_rowid()").fetchone()[0] + + result = record_local_activity_interval( + conn, + { + "id": first_id, + "ts": first_at.isoformat(), + "bytes_written": 1_000, + "bytes_read": 2_000, + "local_tz": "Asia/Kolkata", + }, + { + "id": second_id, + "ts": second_at.isoformat(), + "bytes_written": 1_100, + "bytes_read": 2_050, + "local_tz": "Asia/Kolkata", + }, + start_sample_id=first_id, + end_sample_id=second_id, + segment_id=1, + ) + + assert result == "local_midnight" + stored = conn.execute( + "SELECT COUNT(*), SUM(bytes_written), SUM(bytes_read) " + "FROM local_day_unallocated_evidence" + ).fetchone() + assert stored == (1, 100, 50) + record_local_activity_interval( + conn, + { + "id": first_id, + "ts": first_at.isoformat(), + "bytes_written": 1_000, + "bytes_read": 2_000, + "local_tz": "Asia/Kolkata", + }, + { + "id": second_id, + "ts": second_at.isoformat(), + "bytes_written": 1_100, + "bytes_read": 2_050, + "local_tz": "Asia/Kolkata", + }, + start_sample_id=first_id, + end_sample_id=second_id, + segment_id=1, + ) + assert conn.execute( + "SELECT COUNT(*), SUM(bytes_written), SUM(bytes_read) " + "FROM local_day_unallocated_evidence" + ).fetchone() == (1, 100, 50) + assert conn.execute( + "SELECT SUM(bytes_written), SUM(bytes_read) FROM local_days" + ).fetchone() == (0, 0) + middle_day = query_local_day_summary( + conn, "2026-09-03", "Asia/Kolkata", second_at + timedelta(days=1) + ) + assert middle_day["shared_bytes_written"] == 100 + assert middle_day["shared_bytes_read"] == 50 + assert middle_day["shared_evidence_count"] == 1 + conn.close() + + def test_timezone_change_keeps_interval_unallocated( + self, tmp_path, sysfs_tree, monkeypatch, + ): + from fenris.collector import run_collection + + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + first_at = datetime(2026, 9, 1, 22, 0, tzinfo=timezone.utc) + second_at = datetime(2026, 9, 1, 22, 5, tzinfo=timezone.utc) + + monkeypatch.setenv("TZ", "UTC") + assert run_collection( + _make_smartctl(10_000_000, 8_000_000), + sysfs_nvme, + cfg, + _Clock(first_at), + )["ok"] + monkeypatch.setenv("TZ", "Asia/Kolkata") + assert run_collection( + _make_smartctl(10_000_010, 8_000_004), + sysfs_nvme, + cfg, + _Clock(second_at), + )["ok"] + + conn = sqlite3.connect(store) + event = conn.execute( + "SELECT COUNT(*), SUM(bytes_written), MIN(reason) " + "FROM local_day_unallocated_evidence" + ).fetchone() + assert event == (1, 5_120_000, "timezone_change") + old_zone = query_local_day_summary(conn, "2026-09-01", "UTC") + new_zone = query_local_day_summary(conn, "2026-09-02", "Asia/Kolkata") + assert old_zone["bytes_written"] is None + assert new_zone["bytes_written"] is None + assert old_zone["unallocated_bytes_written"] == 5_120_000 + assert new_zone["unallocated_bytes_written"] == 5_120_000 + conn.close() + + def test_pause_boundary_keeps_interval_unallocated( + self, tmp_path, sysfs_tree, monkeypatch, + ): + from fenris.collector import run_collection + + monkeypatch.setenv("TZ", "UTC") + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + first_at = datetime(2026, 9, 1, 12, 0, tzinfo=timezone.utc) + second_at = datetime(2026, 9, 1, 12, 10, tzinfo=timezone.utc) + assert run_collection( + _make_smartctl(10_000_000, 8_000_000), + sysfs_nvme, + cfg, + _Clock(first_at), + )["ok"] + + conn = sqlite3.connect(store) + close_period(conn, first_at + timedelta(minutes=5), "user_disabled") + conn.close() + + assert run_collection( + _make_smartctl(10_000_010, 8_000_004), + sysfs_nvme, + cfg, + _Clock(second_at), + )["ok"] + + conn = sqlite3.connect(store) + reason = conn.execute( + "SELECT reason FROM local_day_unallocated_evidence" + ).fetchone() + summary = query_local_day_summary(conn, "2026-09-01", "UTC") + assert reason == ("monitoring_period",) + assert summary["bytes_written"] is None + assert summary["unallocated_bytes_written"] == 5_120_000 + assert summary["shared_bytes_written"] == 0 + conn.close() + + def test_controller_segment_boundary_marks_local_activity_incomplete( + self, tmp_path, sysfs_tree, monkeypatch, + ): + from fenris.collector import run_collection + + monkeypatch.setenv("TZ", "UTC") + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + first_at = datetime(2026, 9, 1, 12, 0, tzinfo=timezone.utc) + second_at = datetime(2026, 9, 1, 12, 5, tzinfo=timezone.utc) + assert run_collection( + _make_smartctl(10_000_000, 8_000_000), + sysfs_nvme, + cfg, + _Clock(first_at), + )["ok"] + assert run_collection( + _make_smartctl(10, 8_000_004), + sysfs_nvme, + cfg, + _Clock(second_at), + )["ok"] + + conn = sqlite3.connect(store) + summary = query_local_day_summary( + conn, "2026-09-01", "UTC", second_at + ) + assert summary is not None + assert summary["bytes_written"] is None + assert summary["activity_state"] == "incomplete" + assert conn.execute( + "SELECT COUNT(*) FROM local_day_unallocated_evidence" + ).fetchone()[0] == 0 + conn.close() + # --------------------------------------------------------------------------- # DST handling (23/25-hour days) @@ -402,17 +853,72 @@ class TestDSTHandling: # Simulate a 23-hour day in a timezone with DST # For simplicity, just verify the summary records the correct UTC range clock = datetime(2026, 3, 8, 12, 0, 0, tzinfo=timezone.utc) - # US/Eastern springs forward on 2026-03-08 + # New York springs forward on 2026-03-08 # Local day 2026-03-08 is 23 hours: UTC [07:00, 06:00+1d) _insert_hour(conn, "2026-03-08T08:00:00+00:00", bw=100, active=3600) _insert_hour(conn, "2026-03-08T12:00:00+00:00", bw=200, active=3600) - summary = derive_local_day_summary(conn, "US/Eastern", clock) + summary = derive_local_day_summary(conn, "America/New_York", clock) assert summary is not None assert summary.local_date == "2026-03-08" # UTC range should be approximately 23 hours utc_start = datetime.fromisoformat(summary.utc_start) utc_end = datetime.fromisoformat(summary.utc_end) duration = (utc_end - utc_start).total_seconds() - assert 22 * 3600 <= duration <= 24 * 3600 # ~23h ± tolerance + assert duration == 23 * 3600 + conn.close() + + def test_long_day_25_hours(self): + from fenris.local_day import local_day_boundaries + + start, end, offset = local_day_boundaries( + "2026-11-01", "America/New_York" + ) + assert (end - start).total_seconds() == 25 * 3600 + assert offset == "-04:00" + + def test_half_hour_offset_records_exact_local_boundaries(self): + from fenris.local_day import local_day_boundaries + + start, end, offset = local_day_boundaries( + "2026-09-01", "Asia/Kolkata" + ) + assert start.isoformat() == "2026-08-31T18:30:00+00:00" + assert end.isoformat() == "2026-09-01T18:30:00+00:00" + assert offset == "+05:30" + + def test_repeated_local_clock_labels_remain_one_recorded_day( + self, tmp_path, sysfs_tree, monkeypatch, + ): + from fenris.collector import run_collection + + monkeypatch.setenv("TZ", "America/New_York") + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + # Both endpoints display as 01:30 locally, on opposite UTC offsets. + first_at = datetime(2026, 11, 1, 5, 30, tzinfo=timezone.utc) + second_at = datetime(2026, 11, 1, 6, 30, tzinfo=timezone.utc) + assert run_collection( + _make_smartctl(10_000_000, 8_000_000), + sysfs_nvme, + cfg, + _Clock(first_at), + )["ok"] + assert run_collection( + _make_smartctl(10_000_001, 8_000_002), + sysfs_nvme, + cfg, + _Clock(second_at), + )["ok"] + + conn = sqlite3.connect(store) + summary = query_local_day_summary( + conn, "2026-11-01", "America/New_York", second_at + ) + assert summary is not None + assert summary["bytes_written"] == 512_000 + assert summary["bytes_read"] == 1_024_000 + assert summary["utc_start"] == "2026-11-01T04:00:00+00:00" + assert summary["utc_end"] == "2026-11-02T05:00:00+00:00" conn.close()