Files
Fenris/src/fenris/store.py
T

466 lines
17 KiB
Python

"""Observation store: SQLite database for persisting observation history.
This module handles:
- Store initialization with WAL mode
- Schema versioning with PRAGMA user_version
- Observation records, derived activity, publication state, and metadata
"""
import sqlite3
from pathlib import Path
from typing import Optional
# Schema version - increment on each migration
SCHEMA_VERSION = 6
# 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 observation schema, including private pending-publication storage.
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
try:
_apply_migrations(conn, current_version)
except Exception:
conn.close()
raise
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,
local_tz TEXT
)
""")
# 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,
activity_seconds INTEGER NOT NULL DEFAULT 0,
activity_intervals INTEGER NOT NULL DEFAULT 0,
activity_incomplete BOOLEAN NOT NULL DEFAULT 0,
activity_precision TEXT NOT NULL DEFAULT 'measured',
last_sample_id INTEGER,
UNIQUE(local_date, tz_name)
)
""")
_create_local_day_shared_evidence(conn)
_create_local_day_segment_totals(conn)
# Metadata table for store state (e.g., legacy import marker)
conn.execute("""
CREATE TABLE IF NOT EXISTS store_metadata (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
)
""")
_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):
"""Apply forward-only migrations from current_version to SCHEMA_VERSION.
Commit each version transition independently. A failed step rolls back in
full while earlier successful steps remain versioned and retryable.
"""
migrations = {
2: _migrate_1_to_2,
3: _migrate_2_to_3,
4: _migrate_3_to_4,
5: _migrate_4_to_5,
6: _migrate_5_to_6,
}
while current_version < SCHEMA_VERSION:
target_version = current_version + 1
migration = migrations.get(target_version)
if migration is None:
raise ValueError(f"No migration registered for schema {target_version}")
conn.execute("BEGIN IMMEDIATE")
try:
migration(conn)
conn.execute(f"PRAGMA user_version={target_version}")
conn.commit()
except Exception:
conn.rollback()
raise
current_version = target_version
def _migrate_1_to_2(conn: sqlite3.Connection) -> None:
"""Add segment provenance and unattributed UTC byte tracking."""
tables = {row[0] for row in conn.execute(
"SELECT name FROM sqlite_master WHERE type='table'"
).fetchall()}
if "samples" in tables:
cols = {row[1] for row in conn.execute(
"PRAGMA table_info(samples)"
).fetchall()}
if "segment_id" not in cols:
conn.execute("ALTER TABLE samples ADD COLUMN segment_id INTEGER")
if "day_aggregates" in tables:
cols = {row[1] for row in conn.execute(
"PRAGMA table_info(day_aggregates)"
).fetchall()}
for column in ("unattributed_bytes_written", "unattributed_bytes_read"):
if column not in cols:
conn.execute(
f"ALTER TABLE day_aggregates ADD COLUMN {column} INTEGER DEFAULT 0"
)
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:
"""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
try:
_apply_migrations(conn, current_version)
except Exception:
conn.close()
raise
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