Files
Fenris/src/fenris/collector.py
T
xavierk 4f884b4b73 Show trustworthy local-day activity totals (#90)
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).
2026-09-18 12:28:23 +05:30

381 lines
13 KiB
Python

"""Collector: acquires counters and identity, writes to observation store.
This module implements the thinnest complete write path:
- Acquire counters and thermal evidence from smartctl -a -j
- Acquire controller identity from sysfs
- Normalize identity exactly once at write time
- Validate every row against store invariants
- Commit one well-formed sample
No code path outside the collector interrogates the device.
"""
import json
import sqlite3
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Dict, Optional, Tuple
from .store import init_store, get_store_path
from .monitoring_periods import ensure_period_open
from .derive import find_previous_sample, derive_hours_from_interval
class AcquisitionError(Exception):
"""Raised when acquisition fails - whole run is refused."""
pass
class InvariantViolationError(Exception):
"""Raised when a row would violate store invariants - writes nothing."""
pass
def acquire_from_smartctl(smartctl_data: Dict[str, Any]) -> Dict[str, Any]:
"""Acquire counters and thermal evidence from smartctl -a -j data.
Validates that all required fields are present.
Raises AcquisitionError on any failure.
"""
required_fields = [
"nvme_smart_health_information_log",
"user_capacity",
"model_name",
"serial_number",
"firmware_version",
]
for field in required_fields:
if field not in smartctl_data:
raise AcquisitionError(f"Missing required field in smartctl data: {field}")
log = smartctl_data["nvme_smart_health_information_log"]
required_log_fields = [
"data_units_written",
"data_units_read",
"percentage_used",
"power_on_hours",
"temperature",
]
for field in required_log_fields:
if field not in log:
raise AcquisitionError(f"Missing required field in SMART log: {field}")
return {
"model": smartctl_data["model_name"],
"serial": smartctl_data["serial_number"],
"firmware_rev": smartctl_data["firmware_version"],
"capacity_bytes": smartctl_data["user_capacity"]["bytes"],
"percentage_used": log["percentage_used"],
"available_spare": log.get("available_spare"),
"media_errors": log.get("media_errors", 0),
"power_on_hours": log["power_on_hours"],
"power_cycles": log.get("power_cycles"),
"unsafe_shutdowns": log.get("unsafe_shutdowns"),
"temperature_c": log["temperature"],
"data_units_written": log["data_units_written"],
"data_units_read": log["data_units_read"],
"bytes_written": log["data_units_written"] * 512000,
"bytes_read": log["data_units_read"] * 512000,
"critical_warning": log.get("critical_warning", 0),
}
def acquire_from_sysfs(sysfs_path: Path) -> Dict[str, Any]:
"""Acquire controller identity from sysfs.
Reads identity from:
- /sys/class/nvme/<ctrl>/subsysnqn (primary)
- /sys/class/nvme/<ctrl>/model
- /sys/class/nvme/<ctrl>/serial
- /sys/class/nvme/<ctrl>/firmware_rev
- /sys/class/nvme/<ctrl>/transport/ (optional)
Raises AcquisitionError on any failure.
"""
identity_files = {
"subnqn": "subsysnqn",
"mn": "model",
"sn": "serial",
"fr": "firmware_rev",
}
identity = {}
for key, filename in identity_files.items():
filepath = sysfs_path / filename
if not filepath.exists():
raise AcquisitionError(f"Missing sysfs file: {filepath}")
try:
value = filepath.read_text().strip()
identity[key] = value if value else ""
except Exception as e:
raise AcquisitionError(f"Failed to read {filepath}: {e}")
# Transport info (optional)
transport_dir = sysfs_path / "transport"
if transport_dir.exists():
try:
transport_file = transport_dir / "trstring"
if transport_file.exists():
identity["transport"] = transport_file.read_text().strip()
else:
identity["transport"] = None
except Exception:
identity["transport"] = None
else:
identity["transport"] = None
# vid/ssvid from PCI node (optional, metadata only - never key components)
# PCI device directory is the sysfs_path itself (the controller dir is a symlink to PCI)
pci_device = sysfs_path
for attr, key in [("vendor", "vid"), ("subsystem_vendor", "ssvid")]:
filepath = pci_device / attr
if filepath.exists():
try:
value = filepath.read_text().strip()
identity[key] = value if value else None
except Exception:
identity[key] = None
else:
identity[key] = None
return identity
def normalize_identity(identity: Dict[str, Any]) -> str:
"""Normalize identity exactly once at write time.
Rules:
- Strip trailing spaces and newlines
- No case folding
- Empty-after-strip stored blank
Returns normalized identity key.
"""
# Primary key: normalized kernel-exposed subsystem NQN
key = identity.get("subnqn", "")
if key:
key = key.rstrip()
return key
# Fallback 1: kernel composite (not implemented yet)
# Fallback 2: model|serial
mn = identity.get("mn", "").rstrip()
sn = identity.get("sn", "").rstrip()
if mn or sn:
return f"{mn}|{sn}"
# All keys blank - degraded identity
return ""
def compute_identity_degraded(identity: Dict[str, Any]) -> bool:
"""Check if identity is degraded (all key rungs empty)."""
key = normalize_identity(identity)
return key == ""
def validate_sample_invariants(sample: Dict[str, Any], conn: sqlite3.Connection) -> None:
"""Validate sample against store invariants.
Raises InvariantViolationError if any invariant is violated.
"""
# TODO: Implement more complex invariants as needed
# For now, just check basic constraints
if sample.get("bytes_written", 0) < 0:
raise InvariantViolationError("Negative bytes_written")
if sample.get("bytes_read", 0) < 0:
raise InvariantViolationError("Negative bytes_read")
def write_sample(
sample: Dict[str, Any],
identity: Dict[str, Any],
conn: sqlite3.Connection,
clock,
) -> Dict[str, Any]:
"""Write one sample to the observation store.
Identity normalization happens exactly once here.
Returns segment info for the caller.
"""
from .segment import find_current_segment, should_open_new_segment, open_segment
# Normalize identity exactly once at write time
identity_key = normalize_identity(identity)
identity_degraded = compute_identity_degraded(identity)
# Find current segment
current_segment = find_current_segment(conn)
# Determine if we need a new segment
should_open, reason = should_open_new_segment(
current_segment, identity_key, sample["bytes_written"], conn
)
# Open new segment if needed
segment_opened = False
if should_open:
now = clock.utcnow()
open_segment(conn, now, identity, identity_key, identity_degraded)
segment_opened = True
# Get current segment_id for provenance
current_segment = find_current_segment(conn)
segment_id = current_segment["id"] if current_segment else None
# Insert sample with segment_id
cursor = conn.execute(
"""
INSERT INTO samples (
ts, device, subnqn, sn, mn, fr, capacity_bytes,
percentage_used, available_spare, media_errors, power_on_hours,
power_cycles, unsafe_shutdowns, temperature_c,
data_units_written, data_units_read, bytes_written, bytes_read,
critical_warning, segment_id
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
sample["ts"],
sample["device"],
identity.get("subnqn", ""),
identity.get("sn", ""),
identity.get("mn", ""),
identity.get("fr", ""),
sample["capacity_bytes"],
sample["percentage_used"],
sample["available_spare"],
sample["media_errors"],
sample["power_on_hours"],
sample["power_cycles"],
sample["unsafe_shutdowns"],
sample["temperature_c"],
sample["data_units_written"],
sample["data_units_read"],
sample["bytes_written"],
sample["bytes_read"],
sample["critical_warning"],
segment_id,
),
)
conn.commit()
return {
"segment_opened": segment_opened,
"segment_reason": reason,
"identity_key": identity_key,
"identity_degraded": identity_degraded,
"segment_id": segment_id,
}
def run_collection(
smartctl_data: Dict[str, Any],
sysfs_path: Path,
config: Dict[str, Any],
clock,
) -> Dict[str, Any]:
"""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
identity = acquire_from_sysfs(sysfs_path)
# Build sample with injected clock
sample = {
"ts": clock.utcnow().isoformat(),
"device": config["device"],
**counters,
**identity,
}
# Initialize store if needed
store_path = get_store_path(config)
conn = init_store(store_path)
# Run legacy import if needed (idempotent)
from .legacy import import_legacy_history
history_path = Path(config.get("data_dir", ".")) / "history.jsonl"
if history_path.exists():
import_legacy_history(conn, history_path, clock=clock)
try:
# Ensure monitoring period is open (issue #73 AC2)
ensure_period_open(conn, clock.utcnow())
# Validate invariants
validate_sample_invariants(sample, conn)
# Write sample and get segment info
seg_info = write_sample(sample, identity, conn, clock)
# Derive hour observations from interval with previous sample
try:
# Find the sample we just wrote
cursor = conn.execute("SELECT id FROM samples ORDER BY id DESC LIMIT 1")
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:
conn.close()
except (AcquisitionError, InvariantViolationError) as e:
return {
"ok": False,
"error": str(e),
"error_type": type(e).__name__,
}