"""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//subsysnqn (primary) - /sys/class/nvme//model - /sys/class/nvme//serial - /sys/class/nvme//firmware_rev - /sys/class/nvme//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 # Insert sample 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() return { "segment_opened": segment_opened, "segment_reason": reason, "identity_key": identity_key, "identity_degraded": identity_degraded, } 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: # 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__, }