diff --git a/src/fenris/collector.py b/src/fenris/collector.py index bc16f7a..0317c61 100644 --- a/src/fenris/collector.py +++ b/src/fenris/collector.py @@ -124,6 +124,20 @@ def acquire_from_sysfs(sysfs_path: Path) -> Dict[str, Any]: 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 @@ -180,18 +194,34 @@ def write_sample( identity: Dict[str, Any], conn: sqlite3.Connection, clock, -) -> None: +) -> 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) - # TODO: Implement full sample writing with controller segment handling - # For now, just insert a basic sample with identity fields + # 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 ( @@ -226,6 +256,13 @@ def write_sample( ) conn.commit() + + return { + "segment_opened": segment_opened, + "segment_reason": reason, + "identity_key": identity_key, + "identity_degraded": identity_degraded, + } def run_collection( diff --git a/src/fenris/segment.py b/src/fenris/segment.py new file mode 100644 index 0000000..b409358 --- /dev/null +++ b/src/fenris/segment.py @@ -0,0 +1,147 @@ +"""Controller segment management. + +Handles identity-based segmentation of observation history: +- Find current (most recent) segment +- Determine if a new segment should open +- Open new segments with frozen metadata snapshot + +Segmentation axes (independent): +- Identity key change → quarantines prior history +- DUW decrease with unchanged identity → new segment, prior history stays as habit evidence + +Blank-key semantics (PR-16): +- To/from blank is an identity change → quarantines +- Equal blanks continue the segment, segmented by DUW monotonicity alone +""" +import sqlite3 +from datetime import datetime +from typing import Any, Dict, Optional, Tuple + +from .collector import normalize_identity + + +def find_current_segment(conn: sqlite3.Connection) -> Optional[Dict[str, Any]]: + """Find the most recent (open) controller segment. + + Returns the segment dict or None if no segments exist. + """ + cursor = conn.execute( + "SELECT id, opened_at, identity_key, identity_degraded, " + "subnqn, sn, mn, fr, vid, ssvid, transport " + "FROM controller_segments ORDER BY id DESC LIMIT 1" + ) + row = cursor.fetchone() + if row is None: + return None + + return { + "id": row[0], + "opened_at": row[1], + "identity_key": row[2], + "identity_degraded": bool(row[3]), + "subnqn": row[4], + "sn": row[5], + "mn": row[6], + "fr": row[7], + "vid": row[8], + "ssvid": row[9], + "transport": row[10], + } + + +def get_last_duw(conn: sqlite3.Connection, segment_id: int) -> Optional[int]: + """Get the bytes_written from the most recent sample in a segment. + + Returns None if no samples exist in the segment. + """ + # Samples don't have a segment_id FK yet, so we need to find the + # latest sample before the segment's opened_at, or the latest sample + # if this is the first segment. + # + # For now, we'll use a simpler approach: get the latest sample's bytes_written. + # TODO: Add segment_id FK to samples table in next schema migration + cursor = conn.execute( + "SELECT bytes_written FROM samples ORDER BY id DESC LIMIT 1" + ) + row = cursor.fetchone() + return row[0] if row else None + + +def should_open_new_segment( + current_segment: Optional[Dict[str, Any]], + new_identity_key: str, + new_bytes_written: int, + conn: sqlite3.Connection, +) -> Tuple[bool, Optional[str]]: + """Determine if a new segment should open. + + Returns (should_open, reason). + reason is None if no new segment, or a string describing why. + """ + # No current segment → must open first segment + if current_segment is None: + return True, "first_segment" + + old_key = current_segment["identity_key"] or "" + + # Identity key change (including to/from blank) + if old_key != new_identity_key: + return True, "identity_change" + + # DUW decrease (counter reset or controller replacement with same identity) + last_duw = get_last_duw(conn, current_segment["id"]) + if last_duw is not None and new_bytes_written < last_duw: + return True, "duw_decrease" + + # Same identity, DUW non-decreasing → continue segment + return False, None + + + +def open_segment( + conn: sqlite3.Connection, + now: datetime, + identity: Dict[str, Any], + identity_key: str, + identity_degraded: bool, +) -> Dict[str, Any]: + """Open a new controller segment with frozen metadata snapshot. + + The metadata is immutable once frozen. + """ + cursor = conn.execute( + """ + INSERT INTO controller_segments ( + opened_at, identity_key, identity_degraded, + subnqn, sn, mn, fr, vid, ssvid, transport + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + now.isoformat(), + identity_key if identity_key else None, + identity_degraded, + identity.get("subnqn") or None, + identity.get("sn") or None, + identity.get("mn") or None, + identity.get("fr") or None, + identity.get("vid") or None, + identity.get("ssvid") or None, + identity.get("transport") or None, + ), + ) + + segment_id = cursor.lastrowid + + return { + "id": segment_id, + "opened_at": now.isoformat(), + "identity_key": identity_key if identity_key else None, + "identity_degraded": identity_degraded, + "subnqn": identity.get("subnqn") or None, + "sn": identity.get("sn") or None, + "mn": identity.get("mn") or None, + "fr": identity.get("fr") or None, + "vid": identity.get("vid") or None, + "ssvid": identity.get("ssvid") or None, + "transport": identity.get("transport") or None, + } diff --git a/tests/test_segmentation.py b/tests/test_segmentation.py new file mode 100644 index 0000000..8ec58e8 --- /dev/null +++ b/tests/test_segmentation.py @@ -0,0 +1,465 @@ +"""Controller segmentation tests. + +Covers acceptance criteria: +- ID-1: Identity key ladder, FR metadata only, independent axes +- ID-2: Frozen metadata snapshot at segment open +- ID-4: identity_degraded set exactly when key is blank +- AC-4: vid/ssvid from PCI node, stored null otherwise +- PR-16: Blank-key semantics (to/from blank quarantines, equal blanks continue) +""" +import sqlite3 +from datetime import datetime, timezone +from pathlib import Path +from typing import Any, Dict + +import pytest + +import sys +sys.path.insert(0, str(Path(__file__).parent.parent / "src")) + +from fenris.store import init_store +from fenris.collector import normalize_identity, compute_identity_degraded, acquire_from_sysfs +from fenris.segment import find_current_segment, should_open_new_segment, open_segment + + +# Fixtures + +@pytest.fixture +def store(tmp_path: Path) -> sqlite3.Connection: + conn = init_store(tmp_path / "observations.db") + yield conn + conn.close() + + +@pytest.fixture +def clock(): + class FakeClock: + def __init__(self): + self.now = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + def utcnow(self): + return self.now + return FakeClock() + + +@pytest.fixture +def identity_nqn() -> Dict[str, Any]: + return { + "subnqn": "nqn.2014-08.org.nvmexpress:uuid:12345678-1234-1234-1234-123456789abc", + "mn": "Samsung SSD 970 EVO Plus 1TB", + "sn": "S4EWNX0N123456", + "fr": "2B2QEXM7", + "vid": "0x144d", + "ssvid": "0x144d", + "transport": "pcie", + } + + +@pytest.fixture +def identity_model_serial() -> Dict[str, Any]: + return { + "subnqn": "", + "mn": "Samsung SSD 970 EVO Plus 1TB", + "sn": "S4EWNX0N123456", + "fr": "2B2QEXM7", + "vid": "0x144d", + "ssvid": "0x144d", + "transport": "pcie", + } + + +@pytest.fixture +def identity_degraded() -> Dict[str, Any]: + return { + "subnqn": "", + "mn": "", + "sn": "", + "fr": "2B2QEXM7", + "vid": "0x144d", + "ssvid": "0x144d", + "transport": "pcie", + } + + +# ID-1: Identity key ladder + +class TestIdentityKeyLadder: + def test_subnqn_primary(self, identity_nqn): + key = normalize_identity(identity_nqn) + assert key == "nqn.2014-08.org.nvmexpress:uuid:12345678-1234-1234-1234-123456789abc" + + def test_fallback_to_model_serial(self, identity_model_serial): + key = normalize_identity(identity_model_serial) + assert key == "Samsung SSD 970 EVO Plus 1TB|S4EWNX0N123456" + + def test_all_blank_degraded(self, identity_degraded): + key = normalize_identity(identity_degraded) + assert key == "" + + def test_firmware_is_metadata_only(self, identity_nqn): + key = normalize_identity(identity_nqn) + assert "2B2QEXM7" not in key + + +# ID-2: Frozen metadata snapshot + +class TestFrozenMetadataSnapshot: + def test_segment_stores_all_metadata_fields(self, store, clock, identity_nqn): + key = normalize_identity(identity_nqn) + degraded = compute_identity_degraded(identity_nqn) + + segment = open_segment(store, clock.now, identity_nqn, key, degraded) + + assert segment["identity_key"] == key + assert segment["identity_degraded"] is False + assert segment["subnqn"] == identity_nqn["subnqn"] + assert segment["sn"] == identity_nqn["sn"] + assert segment["mn"] == identity_nqn["mn"] + assert segment["fr"] == identity_nqn["fr"] + assert segment["vid"] == identity_nqn["vid"] + assert segment["ssvid"] == identity_nqn["ssvid"] + assert segment["transport"] == identity_nqn["transport"] + + def test_metadata_nullable(self, store, clock): + sparse_identity = { + "subnqn": "", + "mn": "Legacy Model", + "sn": "LEGACY123", + "fr": "1.0", + "vid": None, + "ssvid": None, + "transport": None, + } + + key = normalize_identity(sparse_identity) + degraded = compute_identity_degraded(sparse_identity) + segment = open_segment(store, clock.now, sparse_identity, key, degraded) + + assert segment["vid"] is None + assert segment["ssvid"] is None + assert segment["transport"] is None + + def test_cntlid_excluded(self, store, clock, identity_nqn): + key = normalize_identity(identity_nqn) + degraded = compute_identity_degraded(identity_nqn) + + segment = open_segment(store, clock.now, identity_nqn, key, degraded) + + assert "cntlid" not in segment + + def test_metadata_immutable_after_open(self, store, clock, identity_nqn): + key = normalize_identity(identity_nqn) + degraded = compute_identity_degraded(identity_nqn) + + open_segment(store, clock.now, identity_nqn, key, degraded) + + found = find_current_segment(store) + assert found["subnqn"] == identity_nqn["subnqn"] + assert found["vid"] == identity_nqn["vid"] + assert found["ssvid"] == identity_nqn["ssvid"] + + +# ID-4: identity_degraded + +class TestIdentityDegraded: + def test_degraded_when_blank_key(self, store, clock, identity_degraded): + key = normalize_identity(identity_degraded) + degraded = compute_identity_degraded(identity_degraded) + + assert key == "" + assert degraded is True + + segment = open_segment(store, clock.now, identity_degraded, key, degraded) + assert segment["identity_degraded"] is True + + def test_not_degraded_with_subnqn(self, store, clock, identity_nqn): + key = normalize_identity(identity_nqn) + degraded = compute_identity_degraded(identity_nqn) + + assert key != "" + assert degraded is False + + segment = open_segment(store, clock.now, identity_nqn, key, degraded) + assert segment["identity_degraded"] is False + + def test_not_degraded_with_model_serial(self, store, clock, identity_model_serial): + key = normalize_identity(identity_model_serial) + degraded = compute_identity_degraded(identity_model_serial) + + assert key != "" + assert degraded is False + + segment = open_segment(store, clock.now, identity_model_serial, key, degraded) + assert segment["identity_degraded"] is False + + +# AC-4: vid/ssvid from PCI node + +class TestVidSsvidAcquisition: + def test_vid_ssvid_from_pci_node(self, tmp_path: Path): + ctrl_dir = tmp_path / "nvme0" + ctrl_dir.mkdir() + + (ctrl_dir / "subsysnqn").write_text("nqn.test\n") + (ctrl_dir / "model").write_text("Test Model\n") + (ctrl_dir / "serial").write_text("TEST123\n") + (ctrl_dir / "firmware_rev").write_text("1.0\n") + (ctrl_dir / "vendor").write_text("0x144d\n") + (ctrl_dir / "subsystem_vendor").write_text("0x144d\n") + + identity = acquire_from_sysfs(ctrl_dir) + + assert identity["vid"] == "0x144d" + assert identity["ssvid"] == "0x144d" + + def test_vid_ssvid_null_when_absent(self, tmp_path: Path): + ctrl_dir = tmp_path / "nvme0" + ctrl_dir.mkdir() + + (ctrl_dir / "subsysnqn").write_text("nqn.test\n") + (ctrl_dir / "model").write_text("Test Model\n") + (ctrl_dir / "serial").write_text("TEST123\n") + (ctrl_dir / "firmware_rev").write_text("1.0\n") + + identity = acquire_from_sysfs(ctrl_dir) + + assert identity["vid"] is None + assert identity["ssvid"] is None + + def test_vid_ssvid_metadata_only(self, store, clock, identity_nqn): + identity_a = {**identity_nqn, "vid": "0x144d", "ssvid": "0x144d"} + identity_b = {**identity_nqn, "vid": "0xFFFF", "ssvid": "0xFFFF"} + + key_a = normalize_identity(identity_a) + key_b = normalize_identity(identity_b) + + assert key_a == key_b + + +# PR-16: Blank-key semantics + +class TestBlankKeySemantics: + def test_to_blank_quarantines(self, store, clock): + healthy = { + "subnqn": "nqn.healthy", + "mn": "Model A", "sn": "SN1", "fr": "1.0", + "vid": "0x1", "ssvid": "0x1", "transport": "pcie", + } + key_healthy = normalize_identity(healthy) + degraded_healthy = compute_identity_degraded(healthy) + + open_segment(store, clock.now, healthy, key_healthy, degraded_healthy) + current = find_current_segment(store) + + blank = { + "subnqn": "", "mn": "", "sn": "", + "fr": "1.0", "vid": "0x1", "ssvid": "0x1", "transport": "pcie", + } + key_blank = normalize_identity(blank) + + should_open, reason = should_open_new_segment( + current, key_blank, 1000, store + ) + + assert should_open is True + assert reason == "identity_change" + + def test_from_blank_quarantines(self, store, clock): + blank = { + "subnqn": "", "mn": "", "sn": "", + "fr": "1.0", "vid": "0x1", "ssvid": "0x1", "transport": "pcie", + } + key_blank = normalize_identity(blank) + degraded_blank = compute_identity_degraded(blank) + + open_segment(store, clock.now, blank, key_blank, degraded_blank) + current = find_current_segment(store) + + healthy = { + "subnqn": "nqn.healthy", + "mn": "Model A", "sn": "SN1", "fr": "1.0", + "vid": "0x1", "ssvid": "0x1", "transport": "pcie", + } + key_healthy = normalize_identity(healthy) + + should_open, reason = should_open_new_segment( + current, key_healthy, 1000, store + ) + + assert should_open is True + assert reason == "identity_change" + + def test_equal_blanks_continue_by_duw(self, store, clock): + blank = { + "subnqn": "", "mn": "", "sn": "", + "fr": "1.0", "vid": "0x1", "ssvid": "0x1", "transport": "pcie", + } + key_blank = normalize_identity(blank) + degraded_blank = compute_identity_degraded(blank) + + open_segment(store, clock.now, blank, key_blank, degraded_blank) + current = find_current_segment(store) + + store.execute( + "INSERT INTO samples (ts, device, bytes_written) VALUES (?, ?, ?)", + ("2026-09-01T12:00:00Z", "/dev/nvme0", 1000) + ) + store.commit() + + should_open, reason = should_open_new_segment( + current, key_blank, 2000, store + ) + + assert should_open is False + assert reason is None + + def test_equal_blanks_duw_decrease_opens(self, store, clock): + blank = { + "subnqn": "", "mn": "", "sn": "", + "fr": "1.0", "vid": "0x1", "ssvid": "0x1", "transport": "pcie", + } + key_blank = normalize_identity(blank) + degraded_blank = compute_identity_degraded(blank) + + open_segment(store, clock.now, blank, key_blank, degraded_blank) + current = find_current_segment(store) + + store.execute( + "INSERT INTO samples (ts, device, bytes_written) VALUES (?, ?, ?)", + ("2026-09-01T12:00:00Z", "/dev/nvme0", 2000) + ) + store.commit() + + should_open, reason = should_open_new_segment( + current, key_blank, 1000, store + ) + + assert should_open is True + assert reason == "duw_decrease" + + +# Independent segmentation axes + +class TestIndependentAxes: + def test_identity_change_with_duw_increase(self, store, clock): + identity_a = { + "subnqn": "nqn.drive-a", + "mn": "Model A", "sn": "SN1", "fr": "1.0", + "vid": "0x1", "ssvid": "0x1", "transport": "pcie", + } + key_a = normalize_identity(identity_a) + degraded_a = compute_identity_degraded(identity_a) + + open_segment(store, clock.now, identity_a, key_a, degraded_a) + current = find_current_segment(store) + + store.execute( + "INSERT INTO samples (ts, device, bytes_written) VALUES (?, ?, ?)", + ("2026-09-01T12:00:00Z", "/dev/nvme0", 2000) + ) + store.commit() + + identity_b = { + "subnqn": "nqn.drive-b", + "mn": "Model B", "sn": "SN2", "fr": "2.0", + "vid": "0x2", "ssvid": "0x2", "transport": "pcie", + } + key_b = normalize_identity(identity_b) + + should_open, reason = should_open_new_segment( + current, key_b, 3000, store + ) + + assert should_open is True + assert reason == "identity_change" + + def test_duw_decrease_same_identity(self, store, clock): + identity = { + "subnqn": "nqn.drive-a", + "mn": "Model A", "sn": "SN1", "fr": "1.0", + "vid": "0x1", "ssvid": "0x1", "transport": "pcie", + } + key = normalize_identity(identity) + degraded = compute_identity_degraded(identity) + + open_segment(store, clock.now, identity, key, degraded) + current = find_current_segment(store) + + store.execute( + "INSERT INTO samples (ts, device, bytes_written) VALUES (?, ?, ?)", + ("2026-09-01T12:00:00Z", "/dev/nvme0", 5000) + ) + store.commit() + + should_open, reason = should_open_new_segment( + current, key, 3000, store + ) + + assert should_open is True + assert reason == "duw_decrease" + + +# Segment lifecycle + +class TestSegmentLifecycle: + def test_first_segment_always_opens(self, store): + key = "nqn.test" + + should_open, reason = should_open_new_segment( + None, key, 1000, store + ) + + assert should_open is True + assert reason == "first_segment" + + def test_same_identity_duw_non_decreasing_continues(self, store, clock): + identity = { + "subnqn": "nqn.drive", + "mn": "Model", "sn": "SN1", "fr": "1.0", + "vid": "0x1", "ssvid": "0x1", "transport": "pcie", + } + key = normalize_identity(identity) + degraded = compute_identity_degraded(identity) + + open_segment(store, clock.now, identity, key, degraded) + current = find_current_segment(store) + + store.execute( + "INSERT INTO samples (ts, device, bytes_written) VALUES (?, ?, ?)", + ("2026-09-01T12:00:00Z", "/dev/nvme0", 1000) + ) + store.commit() + + should_open, reason = should_open_new_segment( + current, key, 1500, store + ) + + assert should_open is False + assert reason is None + + def test_multiple_segments(self, store, clock): + identity_a = { + "subnqn": "nqn.drive-a", + "mn": "Model A", "sn": "SN1", "fr": "1.0", + "vid": "0x1", "ssvid": "0x1", "transport": "pcie", + } + key_a = normalize_identity(identity_a) + degraded_a = compute_identity_degraded(identity_a) + + open_segment(store, clock.now, identity_a, key_a, degraded_a) + + identity_b = { + "subnqn": "nqn.drive-b", + "mn": "Model B", "sn": "SN2", "fr": "2.0", + "vid": "0x2", "ssvid": "0x2", "transport": "pcie", + } + key_b = normalize_identity(identity_b) + degraded_b = compute_identity_degraded(identity_b) + + open_segment(store, clock.now, identity_b, key_b, degraded_b) + + cursor = store.execute("SELECT COUNT(*) FROM controller_segments") + count = cursor.fetchone()[0] + assert count == 2 + + current = find_current_segment(store) + assert current["identity_key"] == key_b