Compare commits

...
7 Commits
28 changed files with 3148 additions and 403 deletions
+6
View File
@@ -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. trace. `?` marks a gap, `~` a partial total, and `u` unallocated daily volume.
Use the selected-point readout for exact values and evidence state. 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. - `v` cycles Live / Day / History; the tabs are also clickable.
- `←` / `→` inspect points; `w` switches read/write volume in every view. - `←` / `→` inspect points; `w` switches read/write volume in every view.
- `[` / `]` browse dates, `g` enters a date, and `t` returns to today/live. - `[` / `]` browse dates, `g` enters a date, and `t` returns to today/live.
@@ -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).
+9 -10
View File
@@ -27,7 +27,6 @@ if VENV_DIR.exists():
if site_packages: if site_packages:
sys.path.insert(0, str(site_packages)) sys.path.insert(0, str(site_packages))
from fenris.store import init_store, get_store_path
from fenris.collector import run_collection from fenris.collector import run_collection
@@ -103,14 +102,6 @@ def main() -> None:
config = load_config() config = load_config()
device = config["device"] device = config["device"]
# Interrogate the drive
smartctl_data = interrogate_drive(device)
# Find sysfs path
sysfs_path = find_nvme_sysfs()
if sysfs_path is None:
raise RuntimeError("No NVMe controller found in sysfs")
# Inject a simple clock # Inject a simple clock
class SimpleClock: class SimpleClock:
def utcnow(self): def utcnow(self):
@@ -118,8 +109,16 @@ def main() -> None:
clock = SimpleClock() clock = SimpleClock()
def acquire():
"""Interrogate the drive after pending-work preflight."""
smartctl_data = interrogate_drive(device)
sysfs_path = find_nvme_sysfs()
if sysfs_path is None:
raise RuntimeError("No NVMe controller found in sysfs")
return smartctl_data, sysfs_path
# Run collection # Run collection
result = run_collection(smartctl_data, sysfs_path, config, clock) result = run_collection(config=config, clock=clock, acquire=acquire)
if result["ok"]: if result["ok"]:
print(f"Collection successful: {result['sample_count']} sample(s)") print(f"Collection successful: {result['sample_count']} sample(s)")
+212 -96
View File
@@ -5,19 +5,21 @@ This module implements the thinnest complete write path:
- Acquire controller identity from sysfs - Acquire controller identity from sysfs
- Normalize identity exactly once at write time - Normalize identity exactly once at write time
- Validate every row against store invariants - Validate every row against store invariants
- Commit one well-formed sample - Stage valid observations privately before derivation
- Publish the sample and derived evidence in one collection-owned transaction
No code path outside the collector interrogates the device. No code path outside the collector interrogates the device.
""" """
import json import json
import sqlite3 import sqlite3
from collections.abc import Callable
from datetime import datetime, timezone from datetime import datetime, timezone
from pathlib import Path from pathlib import Path
from typing import Any, Dict, Optional, Tuple from typing import Any, Dict, Optional
from .store import init_store, get_store_path from .derive import derive_hours_from_interval, find_previous_sample
from .monitoring_periods import ensure_period_open from .monitoring_periods import ensure_period_open
from .derive import find_previous_sample, derive_hours_from_interval from .store import init_store, get_store_path
class AcquisitionError(Exception): class AcquisitionError(Exception):
@@ -30,6 +32,9 @@ class InvariantViolationError(Exception):
pass pass
PENDING_PUBLICATION_LIMIT = 6720
def acquire_from_smartctl(smartctl_data: Dict[str, Any]) -> Dict[str, Any]: def acquire_from_smartctl(smartctl_data: Dict[str, Any]) -> Dict[str, Any]:
"""Acquire counters and thermal evidence from smartctl -a -j data. """Acquire counters and thermal evidence from smartctl -a -j data.
@@ -195,12 +200,12 @@ def write_sample(
sample: Dict[str, Any], sample: Dict[str, Any],
identity: Dict[str, Any], identity: Dict[str, Any],
conn: sqlite3.Connection, conn: sqlite3.Connection,
clock, observed_at: datetime,
) -> Dict[str, Any]: ) -> Dict[str, Any]:
"""Write one sample to the observation store. """Write one sample to the observation store.
Identity normalization happens exactly once here. Identity normalization happens exactly once here.
Returns segment info for the caller. Returns segment info for the caller. Caller owns the transaction.
""" """
from .segment import find_current_segment, should_open_new_segment, open_segment from .segment import find_current_segment, should_open_new_segment, open_segment
@@ -219,8 +224,7 @@ def write_sample(
# Open new segment if needed # Open new segment if needed
segment_opened = False segment_opened = False
if should_open: if should_open:
now = clock.utcnow() open_segment(conn, observed_at, identity, identity_key, identity_degraded)
open_segment(conn, now, identity, identity_key, identity_degraded)
segment_opened = True segment_opened = True
# Get current segment_id for provenance # Get current segment_id for provenance
@@ -235,8 +239,8 @@ def write_sample(
percentage_used, available_spare, media_errors, power_on_hours, percentage_used, available_spare, media_errors, power_on_hours,
power_cycles, unsafe_shutdowns, temperature_c, power_cycles, unsafe_shutdowns, temperature_c,
data_units_written, data_units_read, bytes_written, bytes_read, data_units_written, data_units_read, bytes_written, bytes_read,
critical_warning, segment_id critical_warning, segment_id, local_tz
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""", """,
( (
sample["ts"], sample["ts"],
@@ -259,122 +263,234 @@ def write_sample(
sample["bytes_read"], sample["bytes_read"],
sample["critical_warning"], sample["critical_warning"],
segment_id, segment_id,
sample.get("local_tz"),
), ),
) )
conn.commit()
return { return {
"segment_opened": segment_opened, "segment_opened": segment_opened,
"segment_reason": reason, "segment_reason": reason,
"identity_key": identity_key, "identity_key": identity_key,
"identity_degraded": identity_degraded, "identity_degraded": identity_degraded,
"segment_id": segment_id, "segment_id": segment_id,
"sample_id": cursor.lastrowid,
} }
def run_collection( def _observation_time(sample: Dict[str, Any]) -> datetime:
smartctl_data: Dict[str, Any], """Read an observation's original timestamp for ordered recovery."""
sysfs_path: Path, observed_at = datetime.fromisoformat(sample["ts"])
config: Dict[str, Any], if observed_at.tzinfo is None:
clock, return observed_at.replace(tzinfo=timezone.utc)
) -> Dict[str, Any]: return observed_at.astimezone(timezone.utc)
"""Run one collection run.
This is the main entry point for the collector.
Returns the run outcome.
"""
try:
# Acquire counters and thermal evidence
counters = acquire_from_smartctl(smartctl_data)
# Acquire controller identity def _publish_observation(
identity = acquire_from_sysfs(sysfs_path) conn: sqlite3.Connection,
sample: Dict[str, Any],
identity: Dict[str, Any],
tz_name: str,
) -> 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)
# Build sample with injected clock current_id = seg_info["sample_id"]
sample = { prev = find_previous_sample(conn, seg_info.get("segment_id"), current_id)
"ts": clock.utcnow().isoformat(), if prev is not None:
"device": config["device"], current = {
**counters, "id": current_id,
**identity, "ts": sample["ts"],
"bytes_written": sample["bytes_written"],
"bytes_read": sample["bytes_read"],
"power_on_hours": sample["power_on_hours"],
"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):
persist_day_aggregate(conn, aggregate)
from .local_day import derive_local_day_summary, persist_local_day
local_summary = derive_local_day_summary(conn, tz_name, observed_at)
if local_summary is not None:
persist_local_day(conn, local_summary)
def _pending_count(conn: sqlite3.Connection) -> int:
return conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0]
def _stage_observation(
conn: sqlite3.Connection,
sample: Dict[str, Any],
identity: Dict[str, Any],
tz_name: str,
) -> int:
"""Durably stage valid acquired evidence before attempting publication."""
payload = json.dumps(
{"sample": sample, "identity": identity, "tz_name": tz_name},
separators=(",", ":"),
sort_keys=True,
)
cursor = conn.execute(
"INSERT INTO pending_publications (sample_ts, payload) VALUES (?, ?)",
(sample["ts"], payload),
)
return cursor.lastrowid
def _recover_pending(conn: sqlite3.Connection) -> int:
"""Publish pending observations oldest first; stop at first failure."""
recovered = 0
while True:
conn.execute("BEGIN IMMEDIATE")
try:
row = conn.execute(
"SELECT id, payload FROM pending_publications ORDER BY id LIMIT 1"
).fetchone()
if row is None:
conn.rollback()
return recovered
pending_id, payload = row
observation = json.loads(payload)
_publish_observation(
conn,
observation["sample"],
observation["identity"],
observation["tz_name"],
)
conn.execute("DELETE FROM pending_publications WHERE id = ?", (pending_id,))
conn.commit()
recovered += 1
except Exception:
conn.rollback()
raise
def _recover_pending_for_collection(conn: sqlite3.Connection) -> int:
"""Retry queued work and report exhausted capacity without masking store faults."""
try:
return _recover_pending(conn)
except Exception as exc:
if isinstance(exc, sqlite3.Error) and not isinstance(exc, sqlite3.IntegrityError):
raise
if isinstance(exc, InvariantViolationError):
raise
try:
capacity_full = _pending_count(conn) >= PENDING_PUBLICATION_LIMIT
except sqlite3.Error:
raise exc
if capacity_full:
raise RuntimeError(
"pending publication capacity full "
f"({PENDING_PUBLICATION_LIMIT} observations); no new observation acquired; "
f"recovery failed: {exc}"
) from exc
raise
def run_collection(
smartctl_data: Optional[Dict[str, Any]] = None,
sysfs_path: Optional[Path] = None,
config: Optional[Dict[str, Any]] = None,
clock=None,
*,
acquire: Optional[Callable[[], tuple[Dict[str, Any], Path]]] = None,
) -> Dict[str, Any]:
"""Recover old work, acquire one observation, then publish it atomically.
Production callers pass ``acquire`` so recovery and capacity checks run
before device interrogation. Direct sample arguments remain useful for
deterministic collector tests.
"""
conn = None
try:
if config is None or clock is None:
raise ValueError("config and clock are required")
# Initialize store if needed
store_path = get_store_path(config) store_path = get_store_path(config)
conn = init_store(store_path) conn = init_store(store_path)
# Run legacy import if needed (idempotent)
from .legacy import import_legacy_history from .legacy import import_legacy_history
history_path = Path(config.get("data_dir", ".")) / "history.jsonl" history_path = Path(config.get("data_dir", ".")) / "history.jsonl"
if history_path.exists(): if history_path.exists():
import_legacy_history(conn, history_path, clock=clock) import_legacy_history(conn, history_path, clock=clock)
published_count = _recover_pending_for_collection(conn)
try: # Hold the writer reservation across the capacity check and acquisition.
# Ensure monitoring period is open (issue #73 AC2) # A concurrent collector will recheck pending work before it acquires.
ensure_period_open(conn, clock.utcnow()) while True:
conn.execute("BEGIN IMMEDIATE")
waiting = _pending_count(conn)
if waiting:
conn.rollback()
published_count += _recover_pending_for_collection(conn)
continue
# Validate invariants from .tz_util import detect_system_tz
validate_sample_invariants(sample, conn) tz_name = detect_system_tz()
if acquire is not None:
smartctl_data, sysfs_path = acquire()
if smartctl_data is None or sysfs_path is None:
raise AcquisitionError("No acquired SMART data or sysfs identity")
# Write sample and get segment info counters = acquire_from_smartctl(smartctl_data)
seg_info = write_sample(sample, identity, conn, clock) identity = acquire_from_sysfs(sysfs_path)
sample = {
# Derive hour observations from interval with previous sample "ts": clock.utcnow().isoformat(),
try: "device": config["device"],
# Find the sample we just wrote **counters,
cursor = conn.execute("SELECT id FROM samples ORDER BY id DESC LIMIT 1") **identity,
current_id = cursor.fetchone()[0]
prev = find_previous_sample(conn, seg_info.get("segment_id"), current_id)
if prev is not None:
# Build current sample dict for derivation
current = {
"id": current_id,
"ts": sample["ts"],
"bytes_written": sample["bytes_written"],
"bytes_read": sample["bytes_read"],
"power_on_hours": sample["power_on_hours"],
"temperature_c": sample["temperature_c"],
"data_units_written": sample["data_units_written"],
"data_units_read": sample["data_units_read"],
}
derive_hours_from_interval(conn, prev, current)
# Rebuild day aggregates from hour observations
from .day_aggregate import derive_all_days, persist_day_aggregate
for agg in derive_all_days(conn):
persist_day_aggregate(conn, agg)
# Derive local-day summary using system timezone (issue #90)
try:
from .tz_util import detect_system_tz
from .local_day import derive_local_day_summary, persist_local_day
tz_name = detect_system_tz()
clock_now = clock.utcnow()
local_summary = derive_local_day_summary(conn, tz_name, clock_now)
if local_summary is not None:
persist_local_day(conn, local_summary)
except Exception:
# Local-day derivation failure must not prevent publication
pass
conn.commit()
except Exception:
# Derivation failure must not prevent sample persistence (issue #73 AC6)
pass
return {
"ok": True,
"sample_count": 1,
"store_path": str(store_path),
} }
finally: validate_sample_invariants(sample, conn)
conn.close() _stage_observation(conn, sample, identity, tz_name)
conn.commit()
break
except (AcquisitionError, InvariantViolationError) as e: published_count += _recover_pending_for_collection(conn)
return {
"ok": True,
"sample_count": published_count,
"store_path": str(store_path),
}
except Exception as exc:
if conn is not None:
conn.rollback()
return { return {
"ok": False, "ok": False,
"error": str(e), "error": str(exc),
"error_type": type(e).__name__, "error_type": type(exc).__name__,
} }
finally:
if conn is not None:
conn.close()
+4 -7
View File
@@ -29,7 +29,7 @@ def find_previous_sample(
if segment_id is not None: if segment_id is not None:
cursor = conn.execute( cursor = conn.execute(
"SELECT id, ts, bytes_written, bytes_read, power_on_hours, " "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 = ? " "FROM samples WHERE id < ? AND segment_id = ? "
"ORDER BY id DESC LIMIT 1", "ORDER BY id DESC LIMIT 1",
(current_sample_id, segment_id), (current_sample_id, segment_id),
@@ -37,7 +37,7 @@ def find_previous_sample(
else: else:
cursor = conn.execute( cursor = conn.execute(
"SELECT id, ts, bytes_written, bytes_read, power_on_hours, " "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 < ? " "FROM samples WHERE id < ? "
"ORDER BY id DESC LIMIT 1", "ORDER BY id DESC LIMIT 1",
(current_sample_id,), (current_sample_id,),
@@ -49,7 +49,7 @@ def find_previous_sample(
"id": row[0], "ts": row[1], "bytes_written": row[2], "id": row[0], "ts": row[1], "bytes_written": row[2],
"bytes_read": row[3], "power_on_hours": row[4], "bytes_read": row[3], "power_on_hours": row[4],
"temperature_c": row[5], "data_units_written": row[6], "temperature_c": row[5], "data_units_written": row[6],
"data_units_read": row[7], "data_units_read": row[7], "local_tz": row[8],
} }
@@ -96,6 +96,7 @@ def derive_hours_from_interval(
"""Derive hour observations from a sample pair interval. """Derive hour observations from a sample pair interval.
Returns list of hour observation dicts that were written/updated. Returns list of hour observation dicts that were written/updated.
Caller owns the transaction.
""" """
prev_ts = _parse_ts(prev_sample["ts"]) prev_ts = _parse_ts(prev_sample["ts"])
next_ts = _parse_ts(next_sample["ts"]) next_ts = _parse_ts(next_sample["ts"])
@@ -207,9 +208,6 @@ def _upsert_hour_observation(
"WHERE id = ?", "WHERE id = ?",
(new_bw, new_br, new_samples, existing[0]), (new_bw, new_br, new_samples, existing[0]),
) )
conn.commit()
def _add_unattributed_bytes( def _add_unattributed_bytes(
conn: sqlite3.Connection, conn: sqlite3.Connection,
prev_ts: datetime, prev_ts: datetime,
@@ -241,4 +239,3 @@ def _add_unattributed_bytes(
"unattributed_bytes_read = unattributed_bytes_read + ? WHERE day = ?", "unattributed_bytes_read = unattributed_bytes_read + ? WHERE day = ?",
(bw_delta, br_delta, day), (bw_delta, br_delta, day),
) )
conn.commit()
+1011 -132
View File
File diff suppressed because it is too large Load Diff
+1
View File
@@ -69,6 +69,7 @@ def cmd_enable(args: argparse.Namespace) -> None:
open_period = get_open_period(conn) open_period = get_open_period(conn)
if open_period is None: if open_period is None:
ensure_period_open(conn, now) ensure_period_open(conn, now)
conn.commit()
print("Monitoring period opened at", now.isoformat()) print("Monitoring period opened at", now.isoformat())
else: else:
print("Monitoring period already open (id=%d)" % open_period["id"]) print("Monitoring period already open (id=%d)" % open_period["id"])
+1 -2
View File
@@ -16,7 +16,7 @@ def ensure_period_open(conn: sqlite3.Connection, run_time: datetime) -> None:
"""Ensure a monitoring period is open. If none exists, open one at run_time. """Ensure a monitoring period is open. If none exists, open one at run_time.
Spec §9.8: A collection run finding no open monitoring period opens one Spec §9.8: A collection run finding no open monitoring period opens one
at the run moment, never backdated. at the run moment, never backdated. Caller owns the transaction.
""" """
if get_open_period(conn) is not None: if get_open_period(conn) is not None:
return # Already open — no-op return # Already open — no-op
@@ -26,7 +26,6 @@ def ensure_period_open(conn: sqlite3.Connection, run_time: datetime) -> None:
"INSERT INTO monitoring_periods (started_at) VALUES (?)", "INSERT INTO monitoring_periods (started_at) VALUES (?)",
(ts,), (ts,),
) )
conn.commit()
def close_period( def close_period(
+6 -6
View File
@@ -187,7 +187,7 @@ def _wall_clock_in_range(conn, start, end):
def _resolve_baseline(conn, current_segment): def _resolve_baseline(conn, current_segment):
baseline = _get_baseline(conn) baseline = _get_baseline(conn)
facts = [] facts: list[str] = []
if baseline is None: if baseline is None:
return BaselineTier.NONE, None, "no baseline", facts 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. within a monitoring period that has usable observation evidence.
This is the prerequisite for showing an endurance outlook. This is the prerequisite for showing an endurance outlook.
""" """
row = conn.execute( return _count_complete_local_days(conn) > 0
"SELECT 1 FROM local_days WHERE complete = 1 LIMIT 1"
).fetchone()
return row is not None
def _count_complete_local_days(conn): def _count_complete_local_days(conn):
"""Count the number of complete local observation days.""" """Count the number of complete local observation days."""
row = conn.execute( 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() ).fetchone()
return row[0] if row else 0 return row[0] if row else 0
+25 -25
View File
@@ -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, Re-derives hour observations and day aggregates from surviving raw samples,
with idempotent and interruption-safe guarantees. Preserves import markers, rebuilds legacy local-day evidence from surviving samples or trustworthy UTC
boundary anchors, and valid historical summaries. hours, and preserves import markers, boundary anchors, and valid summaries.
Repair is idempotent and interruption-safe.
Contracts: Contracts:
- Idempotent: running repair multiple times produces no duplicates - 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 (?, ?)", "INSERT OR REPLACE INTO store_metadata (key, value) VALUES (?, ?)",
("repair_in_progress", "true" if in_progress else "false"), ("repair_in_progress", "true" if in_progress else "false"),
) )
conn.commit()
def _update_repair_status(conn: sqlite3.Connection, result: RepairResult) -> None: def _update_repair_status(conn: sqlite3.Connection, result: RepairResult) -> None:
@@ -101,8 +101,6 @@ def _update_repair_status(conn: sqlite3.Connection, result: RepairResult) -> Non
("days_derived", str(days)), ("days_derived", str(days)),
) )
conn.commit()
def get_repair_status(conn: sqlite3.Connection) -> RepairStatus: def get_repair_status(conn: sqlite3.Connection) -> RepairStatus:
"""Get current repair status for read-only views.""" """Get current repair status for read-only views."""
@@ -364,20 +362,16 @@ def repair_derivation(
elif hasattr(clock, 'utcnow'): elif hasattr(clock, 'utcnow'):
clock = 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) result = RepairResult(ok=True)
try: 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") conn.execute("BEGIN IMMEDIATE")
# 1. Find and derive intervals from sample pairs # 1. Find and derive intervals from sample pairs
@@ -408,6 +402,11 @@ def repair_derivation(
else: else:
result.days_updated += 1 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 # 3. Count boundary anchors retained
now = clock if isinstance(clock, datetime) else datetime.now(timezone.utc) now = clock if isinstance(clock, datetime) else datetime.now(timezone.utc)
cursor = conn.execute("SELECT ts FROM samples ORDER BY ts") cursor = conn.execute("SELECT ts FROM samples ORDER BY ts")
@@ -428,21 +427,22 @@ def repair_derivation(
) )
result.legacy_summaries_preserved = cursor.fetchone()[0] result.legacy_summaries_preserved = cursor.fetchone()[0]
# Commit transaction # Publish derived rows and completion status together. A process
conn.commit() # interruption rolls back the in-progress marker with the repair.
# Update repair status
_update_repair_status(conn, result) _update_repair_status(conn, result)
_set_repair_in_progress(conn, False)
conn.commit()
except Exception as e: except Exception as e:
conn.rollback() conn.rollback()
try:
_set_repair_in_progress(conn, False)
conn.commit()
except sqlite3.Error:
conn.rollback()
logger.error("Repair failed: %s", e) logger.error("Repair failed: %s", e)
return RepairResult( return RepairResult(
ok=False, ok=False,
error=str(e), error=str(e),
) )
finally:
# Mark repair as complete
_set_repair_in_progress(conn, False)
return result return result
+17
View File
@@ -115,6 +115,7 @@ class StatusComposition:
# Sample counts for waiting explanations # Sample counts for waiting explanations
sample_count: int = 0 sample_count: int = 0
day_count: int = 0 day_count: int = 0
pending_publication_count: int = 0
# Whether the status dot should blink (only Monitoring) # Whether the status dot should blink (only Monitoring)
should_blink: bool = False should_blink: bool = False
@@ -286,6 +287,7 @@ def compose_status(
freshness_age_s = None freshness_age_s = None
sample_count = 0 sample_count = 0
day_count = 0 day_count = 0
pending_publication_count = 0
deliberately_paused = False deliberately_paused = False
if conn is None and store_fault is None and newer_schema is None: if conn is None and store_fault is None and newer_schema is None:
@@ -315,6 +317,10 @@ def compose_status(
sample_count = cursor.fetchone()[0] sample_count = cursor.fetchone()[0]
cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates") cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates")
day_count = cursor.fetchone()[0] day_count = cursor.fetchone()[0]
schema_version = conn.execute("PRAGMA user_version").fetchone()[0]
if schema_version >= 4:
cursor = conn.execute("SELECT COUNT(*) FROM pending_publications")
pending_publication_count = cursor.fetchone()[0]
except sqlite3.Error as exc: except sqlite3.Error as exc:
store_fault = str(exc) store_fault = str(exc)
@@ -328,6 +334,7 @@ def compose_status(
freshness = "unknown" freshness = "unknown"
freshness_age_s = None freshness_age_s = None
sample_count = day_count = 0 sample_count = day_count = 0
pending_publication_count = 0
deliberately_paused = False deliberately_paused = False
# --- External stop detection --- # --- External stop detection ---
@@ -417,6 +424,7 @@ def compose_status(
overlay_base=overlay_base, overlay_base=overlay_base,
sample_count=sample_count, sample_count=sample_count,
day_count=day_count, day_count=day_count,
pending_publication_count=pending_publication_count,
continuity=continuity, continuity=continuity,
paused_lines=paused_lines, paused_lines=paused_lines,
) )
@@ -459,6 +467,8 @@ def render_status_cli(comp: StatusComposition) -> str:
facts.append("timer: %s" % ( facts.append("timer: %s" % (
"unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive" "unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive"
)) ))
if comp.pending_publication_count:
facts.append(_pending_publication_text(comp.pending_publication_count))
if facts: if facts:
lines.append(" · ".join(facts)) lines.append(" · ".join(facts))
@@ -524,6 +534,8 @@ def render_status_tui(comp: StatusComposition) -> str:
facts.append("Timer: %s" % ( facts.append("Timer: %s" % (
"unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive" "unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive"
)) ))
if comp.pending_publication_count:
facts.append(_pending_publication_text(comp.pending_publication_count).capitalize())
if facts: if facts:
lines.append(" · ".join(facts)) lines.append(" · ".join(facts))
@@ -539,3 +551,8 @@ def render_status_tui(comp: StatusComposition) -> str:
lines.append(pl[:1].upper() + pl[1:]) lines.append(pl[:1].upper() + pl[1:])
return "\n".join(lines) return "\n".join(lines)
def _pending_publication_text(count: int) -> str:
noun = "observation" if count == 1 else "observations"
return f"pending publication: {count} {noun} retained for retry"
+190 -56
View File
@@ -3,8 +3,7 @@
This module handles: This module handles:
- Store initialization with WAL mode - Store initialization with WAL mode
- Schema versioning with PRAGMA user_version - Schema versioning with PRAGMA user_version
- The six entities: samples, hour_observations, day_aggregates, - Observation records, derived activity, publication state, and metadata
monitoring_periods, controller_segments, endurance_baseline
""" """
import sqlite3 import sqlite3
from pathlib import Path from pathlib import Path
@@ -12,7 +11,7 @@ from typing import Optional
# Schema version - increment on each migration # Schema version - increment on each migration
SCHEMA_VERSION = 3 SCHEMA_VERSION = 6
# Packaged default placement (spec §8.3). The config may override it, but a # Packaged default placement (spec §8.3). The config may override it, but a
@@ -33,7 +32,7 @@ def get_store_path(config: dict) -> Path:
def init_store(store_path: Path) -> sqlite3.Connection: def init_store(store_path: Path) -> sqlite3.Connection:
"""Initialize the observation store if not present. """Initialize the observation store if not present.
Creates the database with WAL mode and all six entities. Creates the observation schema, including private pending-publication storage.
Returns a connection to the store. Returns a connection to the store.
""" """
conn = sqlite3.connect(str(store_path)) conn = sqlite3.connect(str(store_path))
@@ -73,9 +72,11 @@ def init_store(store_path: Path) -> sqlite3.Connection:
) )
elif current_version < SCHEMA_VERSION: elif current_version < SCHEMA_VERSION:
# Older version - apply migrations # Older version - apply migrations
_apply_migrations(conn, current_version) try:
conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}") _apply_migrations(conn, current_version)
conn.commit() except Exception:
conn.close()
raise
return conn return conn
@@ -107,7 +108,8 @@ def _create_schema(conn: sqlite3.Connection):
bytes_written INTEGER, bytes_written INTEGER,
bytes_read INTEGER, bytes_read INTEGER,
critical_warning INTEGER, critical_warning INTEGER,
segment_id INTEGER segment_id INTEGER,
local_tz TEXT
) )
""") """)
@@ -209,10 +211,18 @@ def _create_schema(conn: sqlite3.Connection):
coverage REAL DEFAULT 0.0, coverage REAL DEFAULT 0.0,
sample_count INTEGER DEFAULT 0, sample_count INTEGER DEFAULT 0,
complete BOOLEAN 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) 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) # Metadata table for store state (e.g., legacy import marker)
conn.execute(""" conn.execute("""
CREATE TABLE IF NOT EXISTS store_metadata ( CREATE TABLE IF NOT EXISTS store_metadata (
@@ -221,59 +231,181 @@ def _create_schema(conn: sqlite3.Connection):
) )
""") """)
_create_pending_publications(conn)
def _create_pending_publications(conn: sqlite3.Connection) -> None:
"""Create private staging for valid observations awaiting derivation."""
conn.execute("""
CREATE TABLE IF NOT EXISTS pending_publications (
id INTEGER PRIMARY KEY AUTOINCREMENT,
sample_ts TEXT NOT NULL,
payload TEXT NOT NULL
)
""")
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): def _apply_migrations(conn: sqlite3.Connection, current_version: int):
"""Apply forward-only migrations from current_version to SCHEMA_VERSION. """Apply forward-only migrations from current_version to SCHEMA_VERSION.
Each migration step is a transactional block. Add new steps as sequential Commit each version transition independently. A failed step rolls back in
elif branches when SCHEMA_VERSION increases. full while earlier successful steps remain versioned and retryable.
Spec: §3.6, §10.2
""" """
# Migration 1→2: add segment_id provenance to samples, migrations = {
# unattributed byte tracking to day_aggregates (issue #73) 2: _migrate_1_to_2,
if current_version < 2: 3: _migrate_2_to_3,
# Defensive: only ALTER if table exists (handles minimal v1 stores) 4: _migrate_3_to_4,
tables = {row[0] for row in conn.execute( 5: _migrate_4_to_5,
"SELECT name FROM sqlite_master WHERE type='table'" 6: _migrate_5_to_6,
).fetchall()} }
if "samples" in tables: while current_version < SCHEMA_VERSION:
# Check if column already exists (idempotent) target_version = current_version + 1
cols = {row[1] for row in conn.execute("PRAGMA table_info(samples)").fetchall()} migration = migrations.get(target_version)
if "segment_id" not in cols: if migration is None:
conn.execute("ALTER TABLE samples ADD COLUMN segment_id INTEGER") raise ValueError(f"No migration registered for schema {target_version}")
if "day_aggregates" in tables: conn.execute("BEGIN IMMEDIATE")
cols = {row[1] for row in conn.execute("PRAGMA table_info(day_aggregates)").fetchall()} try:
if "unattributed_bytes_written" not in cols: migration(conn)
conn.execute("ALTER TABLE day_aggregates ADD COLUMN unattributed_bytes_written INTEGER DEFAULT 0") conn.execute(f"PRAGMA user_version={target_version}")
if "unattributed_bytes_read" not in cols: conn.commit()
conn.execute("ALTER TABLE day_aggregates ADD COLUMN unattributed_bytes_read INTEGER DEFAULT 0") except Exception:
current_version = 2 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. def _migrate_1_to_2(conn: sqlite3.Connection) -> None:
if current_version < 3: """Add segment provenance and unattributed UTC byte tracking."""
tables = {row[0] for row in conn.execute( tables = {row[0] for row in conn.execute(
"SELECT name FROM sqlite_master WHERE type='table'" "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()} ).fetchall()}
if "local_days" not in tables: if "segment_id" not in cols:
conn.execute(""" conn.execute("ALTER TABLE samples ADD COLUMN segment_id INTEGER")
CREATE TABLE IF NOT EXISTS local_days ( if "day_aggregates" in tables:
id INTEGER PRIMARY KEY AUTOINCREMENT, cols = {row[1] for row in conn.execute(
local_date TEXT NOT NULL, "PRAGMA table_info(day_aggregates)"
tz_name TEXT NOT NULL, ).fetchall()}
tz_offset TEXT NOT NULL, for column in ("unattributed_bytes_written", "unattributed_bytes_read"):
utc_start TEXT NOT NULL, if column not in cols:
utc_end TEXT NOT NULL, conn.execute(
bytes_written INTEGER DEFAULT 0, f"ALTER TABLE day_aggregates ADD COLUMN {column} 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)
) )
""")
current_version = 3
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 "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: def migrate_to_latest(store_path: Path) -> int:
@@ -310,9 +442,11 @@ def migrate_to_latest(store_path: Path) -> int:
return SCHEMA_VERSION return SCHEMA_VERSION
steps = SCHEMA_VERSION - current_version steps = SCHEMA_VERSION - current_version
_apply_migrations(conn, current_version) try:
conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}") _apply_migrations(conn, current_version)
conn.commit() except Exception:
conn.close()
raise
conn.close() conn.close()
return steps return steps
+43 -16
View File
@@ -1709,7 +1709,9 @@ class FenrisTuiApp(App):
) )
tz_name = detect_system_tz() tz_name = detect_system_tz()
if self._browse_date is not None: 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: else:
local = query_current_local_day(conn, self._clock_now, tz_name) local = query_current_local_day(conn, self._clock_now, tz_name)
except Exception: except Exception:
@@ -1720,7 +1722,10 @@ class FenrisTuiApp(App):
if local is None: if local is None:
date = self._browse_date or self._clock_now.astimezone().date().isoformat() 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 widget.display = True
main_grid.remove_class("local-day") main_grid.remove_class("local-day")
return return
@@ -1728,23 +1733,45 @@ class FenrisTuiApp(App):
main_grid.add_class("local-day") main_grid.add_class("local-day")
widget.styles.display = "block" 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"] bw = local["bytes_written"]
br = local["bytes_read"] br = local["bytes_read"]
partial = "" if local["complete"] else " · totals so far" if bw is None or br is None:
tz_display = "%s %s" % (local["tz_name"], local["tz_offset"]) totals = "W unavailable · R unavailable"
else:
text = ( totals = "W %.3f GB known · R %.3f GB known" % (
"[bold]%s[/bold] · %s%s\n" bw / 1e9, br / 1e9
" W %.3f GB · R %.3f GB · %.0f%% coverage"
% (
local["local_date"],
tz_display,
partial,
bw / 1e9,
br / 1e9,
local["coverage"] * 100,
) )
) 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) widget.update(text)
def _on_graph_drill(self, day: str) -> None: def _on_graph_drill(self, day: str) -> None:
+8 -5
View File
@@ -25,11 +25,14 @@ def detect_system_tz() -> str:
localtime = Path("/etc/localtime") localtime = Path("/etc/localtime")
if localtime.is_symlink(): if localtime.is_symlink():
target = os.readlink(str(localtime)) target_path = Path(os.readlink(str(localtime)))
# Strip common prefixes: /usr/share/zoneinfo/, /usr/lib/zoneinfo/ if not target_path.is_absolute():
for prefix in ("/usr/share/zoneinfo/", "/usr/lib/zoneinfo/"): target_path = localtime.parent / target_path
if target.startswith(prefix): target = target_path.resolve().as_posix()
return target[len(prefix):] marker = "/zoneinfo/"
marker_index = target.find(marker)
if marker_index >= 0:
return target[marker_index + len(marker):]
return target return target
return "UTC" return "UTC"
+2 -2
View File
@@ -121,8 +121,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00",
conn.execute( conn.execute(
"INSERT INTO local_days " "INSERT INTO local_days "
"(local_date, tz_name, tz_offset, utc_start, utc_end, " "(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_intervals) "
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)",
(local_date, tz_name, tz_offset, utc_start, utc_end, (local_date, tz_name, tz_offset, utc_start, utc_end,
bw, br, coverage, samples, complete), bw, br, coverage, samples, complete),
) )
+162 -1
View File
@@ -154,7 +154,7 @@ class TestSchemaMigration:
# Migrate # Migrate
steps = migrate_to_latest(db) steps = migrate_to_latest(db)
assert steps == 2 # v1→v2→v3 assert steps == 5 # v1→v2→v3→v4→v5→v6
# Verify data preserved # Verify data preserved
conn = sqlite3.connect(str(db)) conn = sqlite3.connect(str(db))
@@ -356,6 +356,167 @@ class TestCollectorDerivation:
assert count_after == 2 assert count_after == 2
class TestCollectionAtomicity:
"""Collection exposes sample and derived evidence as one publication."""
@pytest.mark.asyncio
async def test_collection_is_visible_through_cli_and_tui_readers(
self, config_fixture, smartctl_fixture, sysfs_fixture_tree, monkeypatch
):
"""Ordinary readers observe matching published sample and activity evidence."""
from fenris.status import get_status, read_status
from fenris.tui import FenrisTuiApp
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(
minute=35, second=0, microsecond=0,
)
first = {
**smartctl_fixture,
"nvme_smart_health_information_log": {
**smartctl_fixture["nvme_smart_health_information_log"],
"data_units_written": 12345678,
"data_units_read": 9876543,
},
}
second = {
**smartctl_fixture,
"nvme_smart_health_information_log": {
**smartctl_fixture["nvme_smart_health_information_log"],
"data_units_written": 12345698,
"data_units_read": 9876553,
},
}
sysfs_path = sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0"
assert run_collection(
first, sysfs_path, config_fixture, FakeClock(now - timedelta(minutes=5))
)["ok"] is True
assert run_collection(
second, sysfs_path, config_fixture, FakeClock(now)
)["ok"] is True
store_path = Path(config_fixture["store_path"])
with read_status(store_path, now, query_services=False) as (reader, composition):
assert reader is not None
assert composition.sample_count == 2
assert composition.day_count == 1
utc_day = reader.execute(
"SELECT bytes_written_delta, bytes_read_delta FROM day_aggregates"
).fetchone()
local_day = reader.execute(
"SELECT bytes_written, bytes_read FROM local_days"
).fetchone()
assert tuple(utc_day) == (10_240_000, 5_120_000)
assert tuple(local_day) == (10_240_000, 5_120_000)
cli_output = get_status(
store_path, now, query_services=True, query_journal=False
)
assert "Monitoring" in cli_output
app = FenrisTuiApp(store_path=store_path)
async with app.run_test(size=(100, 30)) as pilot:
await pilot.pause()
live_readout = str(app.query_one("#live-readout").render())
assert "W 0.010 GB" in live_readout
assert "R 0.005 GB" in live_readout
repeated_result = run_collection(
second, sysfs_path, config_fixture, FakeClock(now + timedelta(minutes=5))
)
assert repeated_result["ok"] is True
with read_status(
store_path, now + timedelta(minutes=5), query_services=False
) as (reader, composition):
assert reader is not None
assert composition.sample_count == 3
utc_day = reader.execute(
"SELECT bytes_written_delta, bytes_read_delta FROM day_aggregates"
).fetchone()
local_day = reader.execute(
"SELECT bytes_written, bytes_read FROM local_days"
).fetchone()
assert tuple(utc_day) == (10_240_000, 5_120_000)
assert tuple(local_day) == (10_240_000, 5_120_000)
def test_failed_local_day_publication_keeps_previous_publication(
self, config_fixture, smartctl_fixture, sysfs_fixture_tree, monkeypatch
):
"""A failed final derivation step leaves all prior reader state intact."""
from fenris.status import read_status
monkeypatch.setenv("TZ", "UTC")
now = datetime.now(timezone.utc).replace(second=0, microsecond=0)
first_clock = FakeClock(now)
first = {
**smartctl_fixture,
"nvme_smart_health_information_log": {
**smartctl_fixture["nvme_smart_health_information_log"],
"data_units_written": 12345678,
"data_units_read": 9876543,
},
}
first_result = run_collection(
first,
sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture,
first_clock,
)
assert first_result["ok"] is True, first_result
writer = sqlite3.connect(config_fixture["store_path"])
writer.execute(
"CREATE TRIGGER fail_local_day_publication "
"BEFORE INSERT ON local_days "
"BEGIN SELECT RAISE(ABORT, 'injected local-day publication failure'); END"
)
writer.commit()
writer.close()
next_time = now + timedelta(minutes=5)
second = {
**smartctl_fixture,
"nvme_smart_health_information_log": {
**smartctl_fixture["nvme_smart_health_information_log"],
"data_units_written": 12345698,
"data_units_read": 9876548,
},
}
failed_result = run_collection(
second,
sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture,
FakeClock(next_time),
)
assert failed_result["ok"] is False
assert "injected local-day publication failure" in failed_result["error"]
with read_status(
Path(config_fixture["store_path"]), next_time, query_services=False
) as (reader, composition):
assert reader is not None
assert composition.sample_count == 1
public_counts = reader.execute(
"SELECT (SELECT COUNT(*) FROM samples), "
"(SELECT COUNT(*) FROM controller_segments), "
"(SELECT COUNT(*) FROM monitoring_periods), "
"(SELECT COUNT(*) FROM hour_observations), "
"(SELECT COUNT(*) FROM day_aggregates), "
"(SELECT COUNT(*) FROM local_days)"
).fetchone()
assert tuple(public_counts) == (1, 1, 1, 0, 0, 0)
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Display states: awaiting first sample, awaiting another sample # Display states: awaiting first sample, awaiting another sample
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
+4 -1
View File
@@ -232,7 +232,7 @@ def test_store_initialization(config_fixture: Dict[str, Any]):
# Initialize store # Initialize store
conn = init_store(store_path) conn = init_store(store_path)
# Verify all six entities exist # Verify the observation, derived-history, and metadata tables exist.
cursor = conn.execute("SELECT name FROM sqlite_master WHERE type='table'") cursor = conn.execute("SELECT name FROM sqlite_master WHERE type='table'")
tables = {row[0] for row in cursor.fetchall()} tables = {row[0] for row in cursor.fetchall()}
@@ -249,6 +249,9 @@ def test_store_initialization(config_fixture: Dict[str, Any]):
expected_tables.add("sqlite_sequence") expected_tables.add("sqlite_sequence")
expected_tables.add("store_metadata") expected_tables.add("store_metadata")
expected_tables.add("local_days") 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 assert expected_tables == tables
conn.close() conn.close()
+43 -2
View File
@@ -98,8 +98,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00",
conn.execute( conn.execute(
"INSERT INTO local_days " "INSERT INTO local_days "
"(local_date, tz_name, tz_offset, utc_start, utc_end, " "(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_intervals) "
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)",
(local_date, tz_name, tz_offset, utc_start, utc_end, (local_date, tz_name, tz_offset, utc_start, utc_end,
bw, br, coverage, samples, complete), bw, br, coverage, samples, complete),
) )
@@ -166,6 +166,47 @@ class TestGateNoCompleteDay:
# The key assertion: gate message should NOT appear when gate IS met # 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) 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 # Gate-2: One complete local day → Limited confidence
+8
View File
@@ -634,6 +634,14 @@ class TestMonitorIntegration:
cmd_enable(args) cmd_enable(args)
mock_enable.assert_called_once_with(True) mock_enable.assert_called_once_with(True)
conn = sqlite3.connect(store_path)
try:
assert conn.execute(
"SELECT COUNT(*) FROM monitoring_periods WHERE ended_at IS NULL"
).fetchone()[0] == 1
finally:
conn.close()
def test_cmd_disable_uses_init_system(self, tmp_path): def test_cmd_disable_uses_init_system(self, tmp_path):
"""cmd_disable uses init_system.disable_timer.""" """cmd_disable uses init_system.disable_timer."""
from fenris.monitor import cmd_disable from fenris.monitor import cmd_disable
+2 -2
View File
@@ -90,8 +90,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00",
conn.execute( conn.execute(
"INSERT INTO local_days " "INSERT INTO local_days "
"(local_date, tz_name, tz_offset, utc_start, utc_end, " "(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_intervals) "
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)",
(local_date, tz_name, tz_offset, utc_start, utc_end, (local_date, tz_name, tz_offset, utc_start, utc_end,
bw, br, coverage, samples, complete), bw, br, coverage, samples, complete),
) )
+2 -2
View File
@@ -98,8 +98,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00",
conn.execute( conn.execute(
"INSERT INTO local_days " "INSERT INTO local_days "
"(local_date, tz_name, tz_offset, utc_start, utc_end, " "(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_intervals) "
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)",
(local_date, tz_name, tz_offset, utc_start, utc_end, (local_date, tz_name, tz_offset, utc_start, utc_end,
bw, br, coverage, samples, complete), bw, br, coverage, samples, complete),
) )
+9 -12
View File
@@ -277,7 +277,7 @@ class TestTimezonePreservation:
now = datetime(2026, 10, 1, 12, 0, 0, tzinfo=timezone.utc) now = datetime(2026, 10, 1, 12, 0, 0, tzinfo=timezone.utc)
old = LocalDaySummary( 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_start="2026-09-16T04:00:00+00:00",
utc_end="2026-09-17T04:00:00+00:00", utc_end="2026-09-17T04:00:00+00:00",
bytes_written=5000, bytes_read=2000, 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) entries = query_local_day_history(conn, "2026-09-16", "2026-09-16", now)
assert len(entries) == 1 assert len(entries) == 1
entry = entries[0] entry = entries[0]
assert entry.tz_name == "US/Eastern" assert entry.tz_name == "America/New_York"
assert entry.tz_offset == "-04:00" assert entry.tz_offset == "-04:00"
assert entry.utc_start == "2026-09-16T04:00:00+00:00" assert entry.utc_start == "2026-09-16T04:00:00+00:00"
assert entry.detail_available is False assert entry.detail_available is False
@@ -389,10 +389,10 @@ class TestIdempotencySafety:
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
class TestVolumeConservation: 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): def test_hourly_bytes_are_unavailable_at_local_precision(self, tmp_path):
"""Midnight-spanning interval appears once, not in both days.""" """UTC-day hour bytes cannot be split across local-day boundaries."""
conn = init_store(tmp_path / "obs.db") conn = init_store(tmp_path / "obs.db")
# Insert hour observations that span midnight # Insert hour observations that span midnight
@@ -409,13 +409,10 @@ class TestVolumeConservation:
summary_17 = derive_local_day_summary(conn, "UTC", clock_17) summary_17 = derive_local_day_summary(conn, "UTC", clock_17)
assert summary_17 is not None assert summary_17 is not None
# Each day should only have its own hour's bytes assert summary_16.bytes_written == 0
# Day 16 has the 23:00 hour (100 bw) assert summary_17.bytes_written == 0
# Day 17 has the 00:00 hour (200 bw) assert summary_16.activity_precision == "unavailable"
assert summary_16.bytes_written == 100 assert summary_17.activity_precision == "unavailable"
assert summary_17.bytes_written == 200
# Total is conserved: 100 + 200 = 300
assert summary_16.bytes_written + summary_17.bytes_written == 300
conn.close() conn.close()
+293
View File
@@ -0,0 +1,293 @@
"""Pending publication behavior from issue #97."""
import sqlite3
import sys
from datetime import datetime, timedelta, timezone
from pathlib import Path
import pytest
sys.path.insert(0, str(Path(__file__).parent.parent / "src"))
from fenris import collector
from fenris.projection import compute_projection
from fenris.status import _query_drive_facts, get_status, read_status
from fenris.status_composition import render_status_tui
from fenris.store import init_store
class FakeClock:
def __init__(self, now):
self.now = now
def utcnow(self):
return self.now
@pytest.fixture
def sysfs_fixture_tree(tmp_path):
ctrl_dir = tmp_path / "sys" / "class" / "nvme" / "nvme0"
ctrl_dir.mkdir(parents=True)
(ctrl_dir / "subsysnqn").write_text("nqn.test:drive\n")
(ctrl_dir / "model").write_text("Test NVMe\n")
(ctrl_dir / "serial").write_text("test-serial\n")
(ctrl_dir / "firmware_rev").write_text("1.0\n")
transport_dir = ctrl_dir / "transport"
transport_dir.mkdir()
(transport_dir / "trstring").write_text("pcie\n")
return tmp_path
def _smartctl(unit_written, unit_read=500):
return {
"nvme_smart_health_information_log": {
"critical_warning": 0,
"temperature": 35,
"available_spare": 100,
"percentage_used": 5,
"data_units_written": unit_written,
"data_units_read": unit_read,
"power_on_hours": 100,
},
"user_capacity": {"bytes": 1_024_000_000_000},
"model_name": "Test NVMe",
"serial_number": "test-serial",
"firmware_version": "1.0",
}
def _acquirer(sysfs_path, unit_written, unit_read=500):
return lambda: (_smartctl(unit_written, unit_read), sysfs_path)
def _public_snapshot(store_path, clock_now):
with read_status(store_path, clock_now) as (conn, comp):
assert conn is not None
return {
"state": comp.state,
"freshness": comp.freshness,
"sample_count": comp.sample_count,
"day_count": comp.day_count,
"boot_enabled": comp.boot_enabled,
"timer_active": comp.timer_active,
"last_collect_ok": comp.last_collect_ok,
"latest_sample": tuple(conn.execute(
"SELECT ts, subnqn, mn, sn, fr, capacity_bytes, temperature_c, "
"available_spare, data_units_written, bytes_written, critical_warning, "
"media_errors, unsafe_shutdowns "
"FROM samples ORDER BY id DESC LIMIT 1"
).fetchone()),
"drive_facts": tuple(_query_drive_facts(conn)),
"projection": repr(compute_projection(conn, clock_now)),
"segments": conn.execute("SELECT COUNT(*) FROM controller_segments").fetchone()[0],
"periods": conn.execute("SELECT COUNT(*) FROM monitoring_periods").fetchone()[0],
"hours": tuple(conn.execute(
"SELECT hour, bytes_written_delta, bytes_read_delta, sample_count "
"FROM hour_observations ORDER BY hour"
).fetchall()),
"days": tuple(conn.execute(
"SELECT day, bytes_written_delta, bytes_read_delta "
"FROM day_aggregates ORDER BY day"
).fetchall()),
"local_days": tuple(conn.execute(
"SELECT local_date, bytes_written, bytes_read, sample_count "
"FROM local_days ORDER BY local_date, tz_name"
).fetchall()),
"pending": comp.pending_publication_count,
}
def test_failed_publication_survives_restart_and_readers_keep_last_consistent_view(
tmp_path, sysfs_fixture_tree, monkeypatch
):
"""A failed derivation stays private, visible as pending, and replays in order."""
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(
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"
for units, at in ((1000, now - timedelta(minutes=10)), (1010, now - timedelta(minutes=5))):
result = collector.run_collection(
config=config,
clock=FakeClock(at),
acquire=_acquirer(sysfs_path, units),
)
assert result["ok"] is True, result
previous = _public_snapshot(store_path, now)
baseline_cli = get_status(store_path, now, query_services=False, query_journal=False)
original_derive = collector.derive_hours_from_interval
reader_during_write = {}
def fail_derivation(conn, previous_sample, next_sample):
# Read while writer has inserted the unpublished sample but not committed.
with read_status(store_path, now, query_services=False) as (reader, comp):
assert reader is not None
reader_during_write["samples"] = reader.execute(
"SELECT COUNT(*) FROM samples"
).fetchone()[0]
reader_during_write["pending"] = comp.pending_publication_count
reader_during_write["latest"] = tuple(reader.execute(
"SELECT data_units_written FROM samples ORDER BY id DESC LIMIT 1"
).fetchone())
raise RuntimeError("injected non-invariant derivation failure")
monkeypatch.setattr(collector, "derive_hours_from_interval", fail_derivation)
failed = collector.run_collection(
config=config,
clock=FakeClock(now),
acquire=_acquirer(sysfs_path, 1025),
)
assert failed["ok"] is False
assert "injected non-invariant derivation failure" in failed["error"]
assert reader_during_write == {"samples": 2, "pending": 1, "latest": (1010,)}
after_failure = _public_snapshot(store_path, now)
assert {**after_failure, "pending": previous["pending"]} == previous
assert after_failure["pending"] == 1
cli_text = get_status(store_path, now, query_services=False, query_journal=False)
pending_fact = " · pending publication: 1 observation retained for retry"
assert pending_fact in cli_text
assert cli_text.replace(pending_fact, "", 1) == baseline_cli
with read_status(store_path, now) as (_, comp):
assert comp.state == previous["state"]
assert comp.boot_enabled is True
assert comp.timer_active is True
assert "pending publication: 1 observation retained for retry" in render_status_tui(comp).lower()
called = False
def should_wait_for_recovery():
nonlocal called
called = True
return _smartctl(1030), sysfs_path
blocked = collector.run_collection(
config=config,
clock=FakeClock(now + timedelta(minutes=5)),
acquire=should_wait_for_recovery,
)
assert blocked["ok"] is False
assert called is False
assert _public_snapshot(store_path, now)["pending"] == 1
monkeypatch.setattr(collector, "derive_hours_from_interval", original_derive)
monkeypatch.setenv("TZ", "Asia/Kolkata")
recovered = collector.run_collection(
config=config,
clock=FakeClock(now + timedelta(minutes=5)),
acquire=_acquirer(sysfs_path, 1040),
)
assert recovered["ok"] is True, recovered
conn = sqlite3.connect(store_path)
try:
assert conn.execute(
"SELECT data_units_written FROM samples ORDER BY id"
).fetchall() == [(1000,), (1010,), (1025,), (1040,)]
assert conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0
assert conn.execute(
"SELECT SUM(bytes_written_delta) FROM day_aggregates"
).fetchone()[0] == (1040 - 1000) * 512_000
assert conn.execute(
"SELECT sample_count FROM local_days WHERE tz_name = 'UTC'"
).fetchone()[0] == 4
assert conn.execute(
"SELECT COUNT(*) FROM local_days WHERE tz_name = 'Asia/Kolkata'"
).fetchone()[0] == 1
finally:
conn.close()
def test_invalid_sample_writes_no_sample_segment_period_or_pending_row(tmp_path, sysfs_fixture_tree):
"""Sample invariants run before staging or public monitoring changes."""
now = datetime.now(timezone.utc).replace(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"
result = collector.run_collection(
config=config,
clock=FakeClock(now),
acquire=_acquirer(sysfs_path, -1),
)
assert result["ok"] is False
assert result["error_type"] == "InvariantViolationError"
conn = sqlite3.connect(store_path)
try:
assert conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0] == 0
assert conn.execute("SELECT COUNT(*) FROM controller_segments").fetchone()[0] == 0
assert conn.execute("SELECT COUNT(*) FROM monitoring_periods").fetchone()[0] == 0
assert conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0
finally:
conn.close()
def test_pending_capacity_retries_before_skipping_acquisition(tmp_path, sysfs_fixture_tree, monkeypatch):
now = datetime.now(timezone.utc).replace(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"
first = collector.run_collection(
config=config,
clock=FakeClock(now - timedelta(minutes=5)),
acquire=_acquirer(sysfs_path, 1000),
)
assert first["ok"] is True, first
monkeypatch.setattr(collector, "PENDING_PUBLICATION_LIMIT", 1)
def fail_derivation(*_):
raise RuntimeError("derivation still unavailable")
monkeypatch.setattr(collector, "derive_hours_from_interval", fail_derivation)
failed = collector.run_collection(
config=config,
clock=FakeClock(now),
acquire=_acquirer(sysfs_path, 1010),
)
assert failed["ok"] is False
assert "capacity full (1 observations)" in failed["error"]
called = False
def must_not_acquire():
nonlocal called
called = True
return _smartctl(1020), sysfs_path
retry = collector.run_collection(
config=config,
clock=FakeClock(now + timedelta(minutes=5)),
acquire=must_not_acquire,
)
assert retry["ok"] is False
assert "capacity full (1 observations)" in retry["error"]
assert called is False
def test_v3_readers_ignore_pending_table_until_store_migrates(tmp_path):
store_path = tmp_path / "observations.db"
conn = init_store(store_path)
conn.execute("DROP TABLE pending_publications")
conn.execute("PRAGMA user_version=3")
conn.commit()
conn.close()
with read_status(store_path, datetime.now(timezone.utc), query_services=False) as (reader, comp):
assert reader is not None
assert comp.store_fault is None
assert comp.pending_publication_count == 0
migrated = init_store(store_path)
try:
assert migrated.execute("PRAGMA user_version").fetchone()[0] == 6
assert migrated.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0
finally:
migrated.close()
+520
View File
@@ -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()
+520 -14
View File
@@ -25,6 +25,7 @@ from fenris.store import init_store, SCHEMA_VERSION
from fenris.local_day import ( from fenris.local_day import (
derive_local_day_summary, derive_local_day_summary,
persist_local_day, persist_local_day,
record_local_activity_interval,
query_local_day_summary, query_local_day_summary,
query_current_local_day, query_current_local_day,
LocalDaySummary, LocalDaySummary,
@@ -135,13 +136,56 @@ class TestSchemaMigration:
assert count == 0 assert count == 0
conn.close() 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, "unavailable")
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: 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): def test_utc_plus_zero(self, tmp_path):
"""UTC timezone: local day boundaries = UTC day boundaries.""" """UTC timezone: local day boundaries = UTC day boundaries."""
@@ -158,8 +202,9 @@ class TestDeriveLocalDay:
assert summary.local_date == "2026-09-01" assert summary.local_date == "2026-09-01"
assert summary.tz_name == "UTC" assert summary.tz_name == "UTC"
assert summary.tz_offset == "+00:00" assert summary.tz_offset == "+00:00"
assert summary.bytes_written == 300 assert summary.bytes_written == 0
assert summary.bytes_read == 130 assert summary.bytes_read == 0
assert summary.activity_precision == "unavailable"
assert summary.complete is False # not all 24 hours covered assert summary.complete is False # not all 24 hours covered
conn.close() conn.close()
@@ -179,12 +224,13 @@ class TestDeriveLocalDay:
assert summary.local_date == "2026-09-01" assert summary.local_date == "2026-09-01"
assert summary.tz_name == "Asia/Kolkata" assert summary.tz_name == "Asia/Kolkata"
assert summary.tz_offset == "+05:30" assert summary.tz_offset == "+05:30"
assert summary.bytes_written == 300 assert summary.bytes_written == 0
assert summary.bytes_read == 130 assert summary.bytes_read == 0
assert summary.activity_precision == "unavailable"
conn.close() conn.close()
def test_midnight_spanning_hour_included(self, tmp_path): def test_hourly_bytes_cannot_be_allocated_to_local_midnight(self, tmp_path):
"""UTC hour straddling local midnight is included in the local day.""" """UTC-hour totals do not establish exact local-day byte totals."""
conn = init_store(tmp_path / "obs.db") conn = init_store(tmp_path / "obs.db")
# For UTC+5:30, local day 2026-09-01 starts at 2026-08-31T18:30 UTC # 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) summary = derive_local_day_summary(conn, "Asia/Kolkata", clock)
assert summary is not None assert summary is not None
# The midnight-spanning hour (18:00) is included assert summary.bytes_written == 0
assert summary.bytes_written == 300 assert summary.bytes_read == 0
assert summary.bytes_read == 130 assert summary.activity_precision == "unavailable"
conn.close() conn.close()
def test_no_hours_returns_none(self, tmp_path): def test_no_hours_returns_none(self, tmp_path):
@@ -387,6 +433,411 @@ class TestCollectorIntegration:
assert row[1] == 30 * 512000 # reads assert row[1] == 30 * 512000 # reads
conn.close() 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) # DST handling (23/25-hour days)
@@ -402,17 +853,72 @@ class TestDSTHandling:
# Simulate a 23-hour day in a timezone with DST # Simulate a 23-hour day in a timezone with DST
# For simplicity, just verify the summary records the correct UTC range # For simplicity, just verify the summary records the correct UTC range
clock = datetime(2026, 3, 8, 12, 0, 0, tzinfo=timezone.utc) 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) # 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-08T08:00:00+00:00", bw=100, active=3600)
_insert_hour(conn, "2026-03-08T12:00:00+00:00", bw=200, 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 is not None
assert summary.local_date == "2026-03-08" assert summary.local_date == "2026-03-08"
# UTC range should be approximately 23 hours # UTC range should be approximately 23 hours
utc_start = datetime.fromisoformat(summary.utc_start) utc_start = datetime.fromisoformat(summary.utc_start)
utc_end = datetime.fromisoformat(summary.utc_end) utc_end = datetime.fromisoformat(summary.utc_end)
duration = (utc_end - utc_start).total_seconds() 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() conn.close()
+2 -2
View File
@@ -100,8 +100,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00",
conn.execute( conn.execute(
"INSERT INTO local_days " "INSERT INTO local_days "
"(local_date, tz_name, tz_offset, utc_start, utc_end, " "(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_intervals) "
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)",
(local_date, tz_name, tz_offset, utc_start, utc_end, (local_date, tz_name, tz_offset, utc_start, utc_end,
bw, br, coverage, samples, complete), bw, br, coverage, samples, complete),
) )
+3 -2
View File
@@ -109,6 +109,7 @@ def _insert_sample(conn, ts, pu=5, device="/dev/nvme0n1"):
def _open_period(conn, start="2026-09-01T00:00:00+00:00"): def _open_period(conn, start="2026-09-01T00:00:00+00:00"):
ensure_period_open(conn, datetime.fromisoformat(start)) ensure_period_open(conn, datetime.fromisoformat(start))
conn.commit()
def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00", def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00",
@@ -123,8 +124,8 @@ def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00",
conn.execute( conn.execute(
"INSERT INTO local_days " "INSERT INTO local_days "
"(local_date, tz_name, tz_offset, utc_start, utc_end, " "(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_intervals) "
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1)",
(local_date, tz_name, tz_offset, utc_start, utc_end, (local_date, tz_name, tz_offset, utc_start, utc_end,
bw, br, coverage, samples, complete), bw, br, coverage, samples, complete),
) )
+26
View File
@@ -0,0 +1,26 @@
"""Timezone detection tests."""
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).parent.parent / "src"))
from fenris import tz_util
def test_relative_localtime_symlink_returns_zoneinfo_key(tmp_path, monkeypatch):
zoneinfo = tmp_path / "usr" / "share" / "zoneinfo" / "Asia" / "Kolkata"
zoneinfo.parent.mkdir(parents=True)
zoneinfo.write_bytes(b"zoneinfo")
localtime = tmp_path / "etc" / "localtime"
localtime.parent.mkdir()
localtime.symlink_to("../usr/share/zoneinfo/Asia/Kolkata")
real_path = tz_util.Path
monkeypatch.setattr(
tz_util,
"Path",
lambda path: localtime if path == "/etc/localtime" else real_path(path),
)
monkeypatch.delenv("TZ", raising=False)
assert tz_util.detect_system_tz() == "Asia/Kolkata"