Add local-day activity summaries derived from UTC hour observations using the system timezone, with durable storage in a new local_days table. The collector derives local-day read/write totals after UTC aggregation; the TUI displays them with timezone, completeness state, and coverage. Schema: bump SCHEMA_VERSION to 3, add local_days table (migration 2→3 is pure addition, idempotent, preserves newer-schema refusal).
332 lines
12 KiB
Python
332 lines
12 KiB
Python
"""Observation store: SQLite database for persisting observation history.
|
|
|
|
This module handles:
|
|
- Store initialization with WAL mode
|
|
- Schema versioning with PRAGMA user_version
|
|
- The six entities: samples, hour_observations, day_aggregates,
|
|
monitoring_periods, controller_segments, endurance_baseline
|
|
"""
|
|
import sqlite3
|
|
from pathlib import Path
|
|
from typing import Optional
|
|
|
|
|
|
# Schema version - increment on each migration
|
|
SCHEMA_VERSION = 3
|
|
|
|
|
|
# Packaged default placement (spec §8.3). The config may override it, but a
|
|
# fresh install that sets only the device selector must collect cleanly.
|
|
DEFAULT_STORE_PATH = Path("/var/lib/fenris/observations.db")
|
|
|
|
|
|
def get_store_path(config: dict) -> Path:
|
|
"""Get the store path from config.
|
|
|
|
Falls back to the packaged default when the config does not pin one,
|
|
so a fresh install whose config holds only the device selector works
|
|
instead of crashing with KeyError 'store_path' (issue #53).
|
|
"""
|
|
return Path(config.get("store_path", DEFAULT_STORE_PATH))
|
|
|
|
|
|
def init_store(store_path: Path) -> sqlite3.Connection:
|
|
"""Initialize the observation store if not present.
|
|
|
|
Creates the database with WAL mode and all six entities.
|
|
Returns a connection to the store.
|
|
"""
|
|
conn = sqlite3.connect(str(store_path))
|
|
|
|
# Enable WAL mode for concurrent reads during writes
|
|
conn.execute("PRAGMA journal_mode=WAL")
|
|
|
|
# Group members (fenris group) read the live store read-only, but SQLite
|
|
# in WAL mode needs write access to the db and its -wal/-shm sidecars even
|
|
# for readers. Best effort: root-created stores stay group-accessible
|
|
# without relying on the creating process's umask (issue #54).
|
|
import os as _os
|
|
for sidecar in (store_path,
|
|
store_path.with_name(store_path.name + "-wal"),
|
|
store_path.with_name(store_path.name + "-shm")):
|
|
try:
|
|
mode = _os.stat(sidecar).st_mode & 0o777
|
|
_os.chmod(sidecar, mode | 0o060)
|
|
except OSError:
|
|
pass
|
|
|
|
# Check if this is a new database
|
|
cursor = conn.execute("PRAGMA user_version")
|
|
current_version = cursor.fetchone()[0]
|
|
|
|
if current_version == 0:
|
|
# New database - create schema
|
|
_create_schema(conn)
|
|
conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}")
|
|
conn.commit()
|
|
elif current_version > SCHEMA_VERSION:
|
|
# Unknown newer version - refuse
|
|
conn.close()
|
|
raise ValueError(
|
|
f"Observation store written by a newer Fenris (version {current_version}) "
|
|
f"— upgrade Fenris"
|
|
)
|
|
elif current_version < SCHEMA_VERSION:
|
|
# Older version - apply migrations
|
|
_apply_migrations(conn, current_version)
|
|
conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}")
|
|
conn.commit()
|
|
|
|
return conn
|
|
|
|
|
|
def _create_schema(conn: sqlite3.Connection):
|
|
"""Create the initial schema with all six entities."""
|
|
|
|
# Samples: raw collection runs (14-day retention)
|
|
conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS samples (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
ts TEXT NOT NULL, -- ISO 8601 UTC timestamp
|
|
device TEXT NOT NULL,
|
|
-- Normalized controller-identity fields captured at acquisition
|
|
subnqn TEXT,
|
|
sn TEXT,
|
|
mn TEXT,
|
|
fr TEXT,
|
|
capacity_bytes INTEGER,
|
|
percentage_used INTEGER,
|
|
available_spare INTEGER,
|
|
media_errors INTEGER,
|
|
power_on_hours INTEGER,
|
|
power_cycles INTEGER,
|
|
unsafe_shutdowns INTEGER,
|
|
temperature_c INTEGER,
|
|
data_units_written INTEGER,
|
|
data_units_read INTEGER,
|
|
bytes_written INTEGER,
|
|
bytes_read INTEGER,
|
|
critical_warning INTEGER,
|
|
segment_id INTEGER
|
|
)
|
|
""")
|
|
|
|
# Hour observations: UTC-hour usage-habit split
|
|
conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS hour_observations (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
hour TEXT NOT NULL UNIQUE, -- ISO 8601 UTC hour (e.g., "2026-09-01T12:00:00Z")
|
|
active_seconds INTEGER DEFAULT 0,
|
|
idle_seconds INTEGER DEFAULT 0,
|
|
powered_off_seconds INTEGER DEFAULT 0,
|
|
unknown_seconds INTEGER DEFAULT 0,
|
|
bytes_written_delta INTEGER DEFAULT 0,
|
|
bytes_read_delta INTEGER DEFAULT 0,
|
|
temperature_min INTEGER,
|
|
temperature_avg REAL,
|
|
temperature_max INTEGER,
|
|
sample_count INTEGER DEFAULT 0,
|
|
coverage REAL DEFAULT 0.0
|
|
)
|
|
""")
|
|
|
|
# Day aggregates: derived from hour observations
|
|
conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS day_aggregates (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
day TEXT NOT NULL UNIQUE, -- ISO 8601 UTC day (e.g., "2026-09-01")
|
|
active_seconds INTEGER DEFAULT 0,
|
|
idle_seconds INTEGER DEFAULT 0,
|
|
powered_off_seconds INTEGER DEFAULT 0,
|
|
unknown_seconds INTEGER DEFAULT 0,
|
|
bytes_written_delta INTEGER DEFAULT 0,
|
|
bytes_read_delta INTEGER DEFAULT 0,
|
|
sample_count INTEGER DEFAULT 0,
|
|
coverage REAL DEFAULT 0.0,
|
|
unattributed_bytes_written INTEGER DEFAULT 0,
|
|
unattributed_bytes_read INTEGER DEFAULT 0
|
|
)
|
|
""")
|
|
|
|
# Monitoring periods: tracking when monitoring was enabled/disabled
|
|
conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS monitoring_periods (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
started_at TEXT NOT NULL, -- ISO 8601 UTC timestamp
|
|
ended_at TEXT, -- NULL if currently active
|
|
end_cause TEXT CHECK(end_cause IN ('user_disabled', 'migrated', 'unknown_gap'))
|
|
)
|
|
""")
|
|
|
|
# Controller segments: identity key plus metadata snapshot
|
|
conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS controller_segments (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
opened_at TEXT NOT NULL, -- ISO 8601 UTC timestamp
|
|
identity_key TEXT, -- Normalized identity key (NULL if degraded)
|
|
identity_degraded BOOLEAN DEFAULT 0,
|
|
subnqn TEXT,
|
|
sn TEXT,
|
|
mn TEXT,
|
|
fr TEXT,
|
|
vid TEXT,
|
|
ssvid TEXT,
|
|
transport TEXT
|
|
)
|
|
""")
|
|
|
|
# Endurance baseline: one active row, replaced on edit
|
|
conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS endurance_baseline (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
tbw_terabytes REAL NOT NULL,
|
|
source_url TEXT,
|
|
document_revision TEXT,
|
|
entry_date TEXT,
|
|
model_string TEXT,
|
|
nominal_capacity_bytes INTEGER,
|
|
validated_by TEXT, -- 'user' or 'machine_match'
|
|
verified BOOLEAN DEFAULT 0,
|
|
created_at TEXT NOT NULL,
|
|
updated_at TEXT NOT NULL
|
|
)
|
|
""")
|
|
|
|
|
|
# Local-day activity summaries derived from UTC hour observations.
|
|
# Each row retains its recorded timezone and UTC boundaries so that
|
|
# historical summaries survive a system-timezone change (ADR 0010).
|
|
conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS local_days (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
local_date TEXT NOT NULL, -- e.g. "2026-09-01" in the recorded tz
|
|
tz_name TEXT NOT NULL, -- POSIX tz name, e.g. "Asia/Kolkata"
|
|
tz_offset TEXT NOT NULL, -- e.g. "+05:30"
|
|
utc_start TEXT NOT NULL, -- ISO 8601 UTC: local midnight start
|
|
utc_end TEXT NOT NULL, -- ISO 8601 UTC: local midnight end
|
|
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)
|
|
)
|
|
""")
|
|
|
|
# Metadata table for store state (e.g., legacy import marker)
|
|
conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS store_metadata (
|
|
key TEXT PRIMARY KEY,
|
|
value TEXT NOT NULL
|
|
)
|
|
""")
|
|
|
|
|
|
def _apply_migrations(conn: sqlite3.Connection, current_version: int):
|
|
"""Apply forward-only migrations from current_version to SCHEMA_VERSION.
|
|
|
|
Each migration step is a transactional block. Add new steps as sequential
|
|
elif branches when SCHEMA_VERSION increases.
|
|
|
|
Spec: §3.6, §10.2
|
|
"""
|
|
# Migration 1→2: add segment_id provenance to samples,
|
|
# unattributed byte tracking to day_aggregates (issue #73)
|
|
if current_version < 2:
|
|
# Defensive: only ALTER if table exists (handles minimal v1 stores)
|
|
tables = {row[0] for row in conn.execute(
|
|
"SELECT name FROM sqlite_master WHERE type='table'"
|
|
).fetchall()}
|
|
if "samples" in tables:
|
|
# Check if column already exists (idempotent)
|
|
cols = {row[1] for row in conn.execute("PRAGMA table_info(samples)").fetchall()}
|
|
if "segment_id" not in cols:
|
|
conn.execute("ALTER TABLE samples ADD COLUMN segment_id INTEGER")
|
|
if "day_aggregates" in tables:
|
|
cols = {row[1] for row in conn.execute("PRAGMA table_info(day_aggregates)").fetchall()}
|
|
if "unattributed_bytes_written" not in cols:
|
|
conn.execute("ALTER TABLE day_aggregates ADD COLUMN unattributed_bytes_written INTEGER DEFAULT 0")
|
|
if "unattributed_bytes_read" not in cols:
|
|
conn.execute("ALTER TABLE day_aggregates ADD COLUMN unattributed_bytes_read INTEGER DEFAULT 0")
|
|
current_version = 2
|
|
|
|
# Migration 2→3: add local_days table for local-day activity totals
|
|
# (issue #90, ADR 0010). Pure addition — no existing rows touched.
|
|
if current_version < 3:
|
|
tables = {row[0] for row in conn.execute(
|
|
"SELECT name FROM sqlite_master WHERE type='table'"
|
|
).fetchall()}
|
|
if "local_days" not in tables:
|
|
conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS local_days (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
local_date TEXT NOT NULL,
|
|
tz_name TEXT NOT NULL,
|
|
tz_offset TEXT NOT NULL,
|
|
utc_start TEXT NOT NULL,
|
|
utc_end TEXT NOT NULL,
|
|
bytes_written INTEGER DEFAULT 0,
|
|
bytes_read INTEGER DEFAULT 0,
|
|
coverage REAL DEFAULT 0.0,
|
|
sample_count INTEGER DEFAULT 0,
|
|
complete BOOLEAN DEFAULT 0,
|
|
UNIQUE(local_date, tz_name)
|
|
)
|
|
""")
|
|
current_version = 3
|
|
|
|
|
|
def migrate_to_latest(store_path: Path) -> int:
|
|
"""Apply forward-only migrations to bring the store to SCHEMA_VERSION.
|
|
|
|
Called by the upgrade target (§10.2). Returns the number of migration
|
|
steps applied. Raises ValueError on newer-schema store (§3.6, §9.5).
|
|
|
|
Spec: §3.6, §10.2, §10.3
|
|
"""
|
|
conn = sqlite3.connect(str(store_path))
|
|
conn.execute("PRAGMA journal_mode=WAL")
|
|
|
|
cursor = conn.execute("PRAGMA user_version")
|
|
current_version = cursor.fetchone()[0]
|
|
|
|
if current_version > SCHEMA_VERSION:
|
|
conn.close()
|
|
raise ValueError(
|
|
f"Observation store written by a newer Fenris (version {current_version}) "
|
|
f"— upgrade Fenris"
|
|
)
|
|
|
|
if current_version == SCHEMA_VERSION:
|
|
conn.close()
|
|
return 0 # Already up to date
|
|
|
|
# Version 0 means no schema — create fresh (issue #73)
|
|
if current_version == 0:
|
|
_create_schema(conn)
|
|
conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}")
|
|
conn.commit()
|
|
conn.close()
|
|
return SCHEMA_VERSION
|
|
|
|
steps = SCHEMA_VERSION - current_version
|
|
_apply_migrations(conn, current_version)
|
|
conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}")
|
|
conn.commit()
|
|
conn.close()
|
|
return steps
|
|
|
|
|
|
def is_store_faulty(store_path: Path) -> bool:
|
|
"""Check if the store is present but cannot be read or trusted."""
|
|
if not store_path.exists():
|
|
return False
|
|
|
|
try:
|
|
conn = sqlite3.connect(f"file:{store_path}?mode=ro", uri=True)
|
|
conn.execute("PRAGMA user_version")
|
|
conn.close()
|
|
return False
|
|
except sqlite3.Error:
|
|
return True
|