feat: segment observation history by controller identity (closes #23)

This commit is contained in:
xavierk
2026-09-01 23:03:42 +05:30
parent 7d219c4697
commit 6a1841c447
3 changed files with 652 additions and 3 deletions
+40 -3
View File
@@ -124,6 +124,20 @@ def acquire_from_sysfs(sysfs_path: Path) -> Dict[str, Any]:
else: else:
identity["transport"] = None 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 return identity
@@ -180,18 +194,34 @@ def write_sample(
identity: Dict[str, Any], identity: Dict[str, Any],
conn: sqlite3.Connection, conn: sqlite3.Connection,
clock, clock,
) -> None: ) -> Dict[str, Any]:
"""Write one sample to the observation store. """Write one sample to the observation store.
Identity normalization happens exactly once here. 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 # Normalize identity exactly once at write time
identity_key = normalize_identity(identity) identity_key = normalize_identity(identity)
identity_degraded = compute_identity_degraded(identity) identity_degraded = compute_identity_degraded(identity)
# TODO: Implement full sample writing with controller segment handling # Find current segment
# For now, just insert a basic sample with identity fields 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( cursor = conn.execute(
""" """
INSERT INTO samples ( INSERT INTO samples (
@@ -226,6 +256,13 @@ def write_sample(
) )
conn.commit() conn.commit()
return {
"segment_opened": segment_opened,
"segment_reason": reason,
"identity_key": identity_key,
"identity_degraded": identity_degraded,
}
def run_collection( def run_collection(
+147
View File
@@ -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,
}
+465
View File
@@ -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