Files
Fenris/src/fenris/collector.py
T

282 lines
8.6 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
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
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,
) -> None:
"""Write one sample to the observation store.
Identity normalization happens exactly once here.
"""
# Normalize identity exactly once at write time
identity_key = normalize_identity(identity)
identity_degraded = compute_identity_degraded(identity)
# TODO: Implement full sample writing with controller segment handling
# For now, just insert a basic sample with identity fields
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
) 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"],
),
)
conn.commit()
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)
try:
# Validate invariants
validate_sample_invariants(sample, conn)
# Write sample
write_sample(sample, identity, conn, clock)
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__,
}