"""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, }