From 2c49fe00c356d0b2ccaa62adc31868e10aab458d Mon Sep 17 00:00:00 2001 From: epi13 Date: Sat, 8 Aug 2026 18:02:23 -0800 Subject: [PATCH 1/2] Add bounded diagnostic snapshot transport views --- schemas/mnel-diagnostic-snapshot.schema.json | 5 +- src/mnel/snapshots.py | 486 ++++++++++++++++++- tests/test_snapshots.py | 64 ++- 3 files changed, 545 insertions(+), 10 deletions(-) diff --git a/schemas/mnel-diagnostic-snapshot.schema.json b/schemas/mnel-diagnostic-snapshot.schema.json index 57040db..f5f0023 100644 --- a/schemas/mnel-diagnostic-snapshot.schema.json +++ b/schemas/mnel-diagnostic-snapshot.schema.json @@ -4,11 +4,12 @@ "title": "MNEL identity-bound diagnostic snapshot metadata", "type": "object", "additionalProperties": false, - "required": ["schema", "snapshot_type", "schema_version", "producer_identity", "source_identity", "dependency_identity", "feature_extractor_identity", "payload_identity", "payload_bytes", "snapshot_identity", "authority", "semantics"], + "required": ["schema", "snapshot_type", "schema_version", "schema_identity", "producer_identity", "source_identity", "dependency_identity", "feature_extractor_identity", "payload_identity", "payload_bytes", "snapshot_identity", "authority", "semantics"], "properties": { "schema": {"const": "mnel-diagnostic-snapshot/0.3"}, - "snapshot_type": {"enum": ["transition", "pair", "tabular"]}, + "snapshot_type": {"enum": ["transition", "pair", "tabular", "trace", "graph", "composite"]}, "schema_version": {"type": "integer", "minimum": 1}, + "schema_identity": {"type": "string", "minLength": 1}, "producer_identity": {"type": "string", "minLength": 1}, "source_identity": {"type": "string", "minLength": 1}, "dependency_identity": {"type": "string", "minLength": 1}, diff --git a/src/mnel/snapshots.py b/src/mnel/snapshots.py index bfd4a18..ee5307c 100644 --- a/src/mnel/snapshots.py +++ b/src/mnel/snapshots.py @@ -11,7 +11,7 @@ import math import struct from dataclasses import dataclass -from typing import Sequence +from typing import Sequence, TypeAlias from .core import canonical_digest @@ -20,6 +20,17 @@ class SnapshotError(ValueError): pass +MAX_SNAPSHOT_BYTES = 1024 * 1024 +MAX_TRANSITION_MEMBER_BYTES = 65535 +MAX_TRACE_EVENTS = 256 +MAX_TRACE_LABEL_BYTES = 64 +MAX_TRACE_PAYLOAD_BYTES = 4096 +MAX_GRAPH_NODES = 4096 +MAX_GRAPH_EDGES = 8192 +MAX_GRAPH_LABEL_BYTES = 64 +MAX_COMPOSITE_COMPONENTS = 64 + + @dataclass(frozen=True, slots=True) class DiagnosticSnapshot: snapshot_type: str @@ -31,6 +42,7 @@ class DiagnosticSnapshot: payload: bytes payload_identity: str snapshot_identity: str + schema_identity: str = "mnel-diagnostic-snapshot" @classmethod def build( @@ -43,23 +55,33 @@ def build( dependency_identity: str, feature_extractor_identity: str, payload: bytes, + schema_identity: str = "mnel-diagnostic-snapshot", ) -> "DiagnosticSnapshot": - if not snapshot_type.strip() or schema_version < 1: - raise SnapshotError("snapshot type and positive schema version are required") + if ( + not isinstance(snapshot_type, str) + or not snapshot_type.strip() + or not isinstance(schema_version, int) + or isinstance(schema_version, bool) + or schema_version < 1 + or not isinstance(schema_identity, str) + or not schema_identity.strip() + ): + raise SnapshotError("snapshot type, schema identity, and positive schema version are required") identities = ( producer_identity, source_identity, dependency_identity, feature_extractor_identity, ) - if any(not identity.strip() for identity in identities): + if any(not isinstance(identity, str) or not identity.strip() for identity in identities): raise SnapshotError("snapshot producer and dependency identities are required") - if not payload or len(payload) > 1024 * 1024: + if not payload or len(payload) > MAX_SNAPSHOT_BYTES: raise SnapshotError("snapshot payload must be non-empty and bounded") payload_identity = "sha256:" + hashlib.sha256(payload).hexdigest() identity_body = { "snapshot_type": snapshot_type, "schema_version": schema_version, + "schema_identity": schema_identity, "producer_identity": producer_identity, "source_identity": source_identity, "dependency_identity": dependency_identity, @@ -76,13 +98,31 @@ def build( bytes(payload), payload_identity, canonical_digest(identity_body), + schema_identity, ) + def validate_integrity(self) -> None: + """Reject tampered metadata or payload before a consumer receives a view.""" + + rebuilt = DiagnosticSnapshot.build( + snapshot_type=self.snapshot_type, + schema_version=self.schema_version, + schema_identity=self.schema_identity, + producer_identity=self.producer_identity, + source_identity=self.source_identity, + dependency_identity=self.dependency_identity, + feature_extractor_identity=self.feature_extractor_identity, + payload=self.payload, + ) + if rebuilt.payload_identity != self.payload_identity or rebuilt.snapshot_identity != self.snapshot_identity: + raise SnapshotError("snapshot identity or payload identity does not match its content") + def to_dict(self) -> dict[str, object]: return { "schema": "mnel-diagnostic-snapshot/0.3", "snapshot_type": self.snapshot_type, "schema_version": self.schema_version, + "schema_identity": self.schema_identity, "producer_identity": self.producer_identity, "source_identity": self.source_identity, "dependency_identity": self.dependency_identity, @@ -104,6 +144,7 @@ def transition_snapshot( dependency_identity: str, feature_extractor_identity: str, schema_version: int = 1, + schema_identity: str = "mnel-diagnostic-snapshot", ) -> DiagnosticSnapshot: return DiagnosticSnapshot.build( snapshot_type="transition", @@ -113,6 +154,7 @@ def transition_snapshot( dependency_identity=dependency_identity, feature_extractor_identity=feature_extractor_identity, payload=_pair_payload(b"MNEL-T1", previous_state, next_state), + schema_identity=schema_identity, ) @@ -125,6 +167,7 @@ def pair_snapshot( dependency_identity: str, feature_extractor_identity: str, schema_version: int = 1, + schema_identity: str = "mnel-diagnostic-snapshot", ) -> DiagnosticSnapshot: return DiagnosticSnapshot.build( snapshot_type="pair", @@ -134,6 +177,7 @@ def pair_snapshot( dependency_identity=dependency_identity, feature_extractor_identity=feature_extractor_identity, payload=_pair_payload(b"MNEL-P1", left, right), + schema_identity=schema_identity, ) @@ -145,16 +189,19 @@ def tabular_snapshot( dependency_identity: str, feature_extractor_identity: str, schema_version: int = 1, + schema_identity: str = "mnel-diagnostic-snapshot", ) -> DiagnosticSnapshot: if not rows or not rows[0]: raise SnapshotError("tabular snapshot requires non-empty rows and columns") + if any(not isinstance(row, Sequence) for row in rows): + raise SnapshotError("tabular rows must be sequences") column_count = len(rows[0]) if len(rows) > 65535 or column_count > 65535 or any(len(row) != column_count for row in rows): raise SnapshotError("tabular shape is invalid or exceeds the bounded format") values: list[float] = [] for row in rows: for value in row: - if not math.isfinite(value): + if not isinstance(value, (int, float)) or isinstance(value, bool) or not math.isfinite(value): raise SnapshotError("tabular values must be finite") values.append(float(value)) payload = struct.pack(">4sHH", b"MNET", len(rows), column_count) + struct.pack( @@ -168,10 +215,435 @@ def tabular_snapshot( dependency_identity=dependency_identity, feature_extractor_identity=feature_extractor_identity, payload=payload, + schema_identity=schema_identity, ) def _pair_payload(header: bytes, left: bytes, right: bytes) -> bytes: - if not left or not right or len(left) > 65535 or len(right) > 65535: + if ( + not isinstance(left, bytes) + or not isinstance(right, bytes) + or not left + or not right + or len(left) > MAX_TRANSITION_MEMBER_BYTES + or len(right) > MAX_TRANSITION_MEMBER_BYTES + ): raise SnapshotError("pair members must be non-empty and bounded") return header + struct.pack(">HH", len(left), len(right)) + left + right + + +@dataclass(frozen=True, slots=True) +class TransitionView: + previous_state: bytes + next_state: bytes + + +@dataclass(frozen=True, slots=True) +class PairView: + left: bytes + right: bytes + + +@dataclass(frozen=True, slots=True) +class TabularView: + rows: tuple[tuple[float, ...], ...] + + @property + def row_count(self) -> int: + return len(self.rows) + + @property + def column_count(self) -> int: + return len(self.rows[0]) if self.rows else 0 + + +@dataclass(frozen=True, slots=True) +class TraceEvent: + event_type: str + payload: bytes + + +@dataclass(frozen=True, slots=True) +class TraceView: + events: tuple[TraceEvent, ...] + + +@dataclass(frozen=True, slots=True) +class GraphNode: + node_id: int + label: str + + +@dataclass(frozen=True, slots=True) +class GraphEdge: + source: int + target: int + edge_type: str + + +@dataclass(frozen=True, slots=True) +class GraphView: + nodes: tuple[GraphNode, ...] + edges: tuple[GraphEdge, ...] + + +@dataclass(frozen=True, slots=True) +class CompositeView: + component_snapshot_identities: tuple[str, ...] + + +SnapshotView: TypeAlias = TransitionView | PairView | TabularView | TraceView | GraphView | CompositeView + + +def decode_transition(payload: bytes) -> TransitionView: + left, right = _decode_pair_payload(payload, b"MNEL-T1") + return TransitionView(left, right) + + +def decode_pair(payload: bytes) -> PairView: + left, right = _decode_pair_payload(payload, b"MNEL-P1") + return PairView(left, right) + + +def decode_tabular(payload: bytes) -> TabularView: + if len(payload) > MAX_SNAPSHOT_BYTES or len(payload) < 8 or payload[:4] != b"MNET": + raise SnapshotError("invalid tabular magic or truncated header") + rows, columns = struct.unpack(">HH", payload[4:8]) + if rows == 0 or columns == 0 or rows > 65535 or columns > 65535: + raise SnapshotError("invalid tabular dimensions") + expected = 8 + rows * columns * 8 + if len(payload) != expected: + raise SnapshotError("tabular payload is truncated or has trailing bytes") + values = struct.unpack(f">{rows * columns}d", payload[8:]) + if any(not math.isfinite(value) for value in values): + raise SnapshotError("tabular payload contains a non-finite value") + return TabularView(tuple(tuple(values[row * columns : (row + 1) * columns]) for row in range(rows))) + + +def decode_trace(payload: bytes) -> TraceView: + if len(payload) > MAX_SNAPSHOT_BYTES or len(payload) < 9 or payload[:7] != b"MNEL-R1": + raise SnapshotError("invalid trace magic or truncated header") + count = struct.unpack(">H", payload[7:9])[0] + if count == 0 or count > MAX_TRACE_EVENTS: + raise SnapshotError("invalid trace event count") + offset = 9 + events: list[TraceEvent] = [] + for _ in range(count): + if offset + 3 > len(payload): + raise SnapshotError("truncated trace event header") + label_length, body_length = struct.unpack(">BH", payload[offset : offset + 3]) + offset += 3 + if label_length == 0 or label_length > MAX_TRACE_LABEL_BYTES or body_length > MAX_TRACE_PAYLOAD_BYTES: + raise SnapshotError("trace event exceeds its label or payload ceiling") + end = offset + label_length + body_length + if end > len(payload): + raise SnapshotError("truncated trace event payload") + try: + event_type = payload[offset : offset + label_length].decode("utf-8") + except UnicodeDecodeError as error: + raise SnapshotError("trace event type is not UTF-8") from error + offset += label_length + events.append(TraceEvent(event_type, bytes(payload[offset : offset + body_length]))) + offset = end + if offset != len(payload): + raise SnapshotError("trace payload has trailing bytes") + return TraceView(tuple(events)) + + +def decode_graph(payload: bytes) -> GraphView: + if len(payload) > MAX_SNAPSHOT_BYTES or len(payload) < 11 or payload[:7] != b"MNEL-G1": + raise SnapshotError("invalid graph magic or truncated header") + node_count, edge_count = struct.unpack(">HH", payload[7:11]) + if node_count == 0 or node_count > MAX_GRAPH_NODES or edge_count > MAX_GRAPH_EDGES: + raise SnapshotError("invalid graph dimensions") + offset = 11 + nodes: list[GraphNode] = [] + seen_nodes: set[int] = set() + for _ in range(node_count): + if offset + 5 > len(payload): + raise SnapshotError("truncated graph node") + node_id, label_length = struct.unpack(">IB", payload[offset : offset + 5]) + offset += 5 + if label_length == 0 or label_length > MAX_GRAPH_LABEL_BYTES or node_id in seen_nodes: + raise SnapshotError("invalid or duplicate graph node") + if offset + label_length > len(payload): + raise SnapshotError("truncated graph node label") + try: + label = payload[offset : offset + label_length].decode("utf-8") + except UnicodeDecodeError as error: + raise SnapshotError("graph node label is not UTF-8") from error + offset += label_length + seen_nodes.add(node_id) + nodes.append(GraphNode(node_id, label)) + edges: list[GraphEdge] = [] + seen_edges: set[tuple[int, int, str]] = set() + for _ in range(edge_count): + if offset + 9 > len(payload): + raise SnapshotError("truncated graph edge") + source, target, label_length = struct.unpack(">IIB", payload[offset : offset + 9]) + offset += 9 + if source not in seen_nodes or target not in seen_nodes or label_length == 0 or label_length > MAX_GRAPH_LABEL_BYTES: + raise SnapshotError("graph edge references an unknown node or invalid type") + if offset + label_length > len(payload): + raise SnapshotError("truncated graph edge type") + try: + edge_type = payload[offset : offset + label_length].decode("utf-8") + except UnicodeDecodeError as error: + raise SnapshotError("graph edge type is not UTF-8") from error + offset += label_length + edge = (source, target, edge_type) + if edge in seen_edges: + raise SnapshotError("duplicate graph edge") + seen_edges.add(edge) + edges.append(GraphEdge(*edge)) + if offset != len(payload): + raise SnapshotError("graph payload has trailing bytes") + return GraphView(tuple(nodes), tuple(edges)) + + +def decode_composite(payload: bytes) -> CompositeView: + if len(payload) > MAX_SNAPSHOT_BYTES or len(payload) < 9 or payload[:7] != b"MNEL-C1": + raise SnapshotError("invalid composite magic or truncated header") + count = struct.unpack(">H", payload[7:9])[0] + if count == 0 or count > MAX_COMPOSITE_COMPONENTS: + raise SnapshotError("invalid composite component count") + offset = 9 + identities: list[str] = [] + for _ in range(count): + if offset + 2 > len(payload): + raise SnapshotError("truncated composite component") + length = struct.unpack(">H", payload[offset : offset + 2])[0] + offset += 2 + if length == 0 or length > 128 or offset + length > len(payload): + raise SnapshotError("invalid composite component identity") + try: + identity = payload[offset : offset + length].decode("ascii") + except UnicodeDecodeError as error: + raise SnapshotError("composite component identity is not ASCII") from error + if ( + len(identity) != 71 + or not identity.startswith("sha256:") + or any(character not in "0123456789abcdef" for character in identity[7:]) + or identity in identities + ): + raise SnapshotError("composite component identity is invalid or duplicated") + identities.append(identity) + offset += length + if offset != len(payload): + raise SnapshotError("composite payload has trailing bytes") + return CompositeView(tuple(identities)) + + +def decode_snapshot(snapshot: DiagnosticSnapshot) -> SnapshotView: + snapshot.validate_integrity() + decoders = { + "transition": decode_transition, + "pair": decode_pair, + "tabular": decode_tabular, + "trace": decode_trace, + "graph": decode_graph, + "composite": decode_composite, + } + try: + decoder = decoders[snapshot.snapshot_type] + except KeyError as error: + raise SnapshotError(f"unsupported snapshot type: {snapshot.snapshot_type}") from error + return decoder(snapshot.payload) + + +class SnapshotStore: + """Identity-keyed immutable snapshot store shared by probes and providers.""" + + def __init__(self) -> None: + self._snapshots: dict[str, DiagnosticSnapshot] = {} + + def register(self, snapshot: DiagnosticSnapshot) -> str: + snapshot.validate_integrity() + view = decode_snapshot(snapshot) + if isinstance(view, CompositeView): + for identity in view.component_snapshot_identities: + if identity not in self._snapshots: + raise SnapshotError("composite references an unregistered snapshot") + existing = self._snapshots.get(snapshot.snapshot_identity) + if existing is not None and existing != snapshot: + raise SnapshotError("snapshot identity collision") + self._snapshots[snapshot.snapshot_identity] = snapshot + return snapshot.snapshot_identity + + def get(self, identity: str) -> DiagnosticSnapshot: + try: + return self._snapshots[identity] + except KeyError as error: + raise SnapshotError(f"unknown snapshot identity: {identity}") from error + + def view(self, identity: str, *, accepted_types: Sequence[str] = (), schema_versions: Sequence[int] = ()) -> SnapshotView: + snapshot = self.get(identity) + if accepted_types and snapshot.snapshot_type not in accepted_types: + raise SnapshotError("snapshot type is incompatible with the requested consumer") + if schema_versions and snapshot.schema_version not in schema_versions: + raise SnapshotError("snapshot schema version is incompatible with the requested consumer") + return decode_snapshot(snapshot) + + def identities(self) -> tuple[str, ...]: + return tuple(sorted(self._snapshots)) + + +def trace_snapshot( + events: Sequence[tuple[str, bytes]], + *, + producer_identity: str, + source_identity: str, + dependency_identity: str, + feature_extractor_identity: str, + schema_version: int = 1, + schema_identity: str = "mnel-diagnostic-snapshot", +) -> DiagnosticSnapshot: + if not events or len(events) > MAX_TRACE_EVENTS: + raise SnapshotError("trace requires a bounded non-empty event sequence") + encoded = bytearray(b"MNEL-R1" + struct.pack(">H", len(events))) + for event_type, payload in events: + if not isinstance(event_type, str) or not isinstance(payload, bytes): + raise SnapshotError("trace events require a string type and byte payload") + label = event_type.encode("utf-8") + if not label or len(label) > MAX_TRACE_LABEL_BYTES or len(payload) > MAX_TRACE_PAYLOAD_BYTES: + raise SnapshotError("trace event exceeds its label or payload ceiling") + encoded.extend(struct.pack(">BH", len(label), len(payload))) + encoded.extend(label) + encoded.extend(payload) + return DiagnosticSnapshot.build( + snapshot_type="trace", + schema_version=schema_version, + producer_identity=producer_identity, + source_identity=source_identity, + dependency_identity=dependency_identity, + feature_extractor_identity=feature_extractor_identity, + payload=bytes(encoded), + schema_identity=schema_identity, + ) + + +def graph_snapshot( + nodes: Sequence[tuple[int, str]], + edges: Sequence[tuple[int, int, str]], + *, + producer_identity: str, + source_identity: str, + dependency_identity: str, + feature_extractor_identity: str, + schema_version: int = 1, + schema_identity: str = "mnel-diagnostic-snapshot", +) -> DiagnosticSnapshot: + if not nodes or len(nodes) > MAX_GRAPH_NODES or len(edges) > MAX_GRAPH_EDGES: + raise SnapshotError("graph exceeds its node or edge ceiling") + if any( + not isinstance(item, (tuple, list)) or len(item) != 2 for item in nodes + ) or any(not isinstance(item, (tuple, list)) or len(item) != 3 for item in edges): + raise SnapshotError("graph nodes and edges have invalid shapes") + normalized_nodes = list(nodes) + if any( + not isinstance(node_id, int) + or isinstance(node_id, bool) + or node_id < 0 + or node_id > 0xFFFFFFFF + or not isinstance(label, str) + or not label + for node_id, label in normalized_nodes + ): + raise SnapshotError("graph nodes require non-negative ids and labels") + if len({node_id for node_id, _ in normalized_nodes}) != len(normalized_nodes): + raise SnapshotError("graph node identities must be unique") + node_ids = {node_id for node_id, _ in normalized_nodes} + normalized_edges = list(edges) + if any( + not isinstance(source, int) + or not isinstance(target, int) + or isinstance(source, bool) + or isinstance(target, bool) + or source < 0 + or target < 0 + or source > 0xFFFFFFFF + or target > 0xFFFFFFFF + or source not in node_ids + or target not in node_ids + or not isinstance(edge_type, str) + or not edge_type + for source, target, edge_type in normalized_edges + ): + raise SnapshotError("graph edges must reference declared nodes") + if len(set(normalized_edges)) != len(normalized_edges): + raise SnapshotError("graph edges must be unique") + normalized_nodes.sort(key=lambda item: item[0]) + normalized_edges.sort(key=lambda item: (item[0], item[1], item[2])) + encoded = bytearray(b"MNEL-G1" + struct.pack(">HH", len(normalized_nodes), len(normalized_edges))) + for node_id, label in normalized_nodes: + label_bytes = label.encode("utf-8") + if len(label_bytes) > MAX_GRAPH_LABEL_BYTES: + raise SnapshotError("graph node label exceeds its ceiling") + encoded.extend(struct.pack(">IB", node_id, len(label_bytes))) + encoded.extend(label_bytes) + for source, target, edge_type in normalized_edges: + label_bytes = edge_type.encode("utf-8") + if len(label_bytes) > MAX_GRAPH_LABEL_BYTES: + raise SnapshotError("graph edge type exceeds its ceiling") + encoded.extend(struct.pack(">IIB", source, target, len(label_bytes))) + encoded.extend(label_bytes) + return DiagnosticSnapshot.build( + snapshot_type="graph", + schema_version=schema_version, + producer_identity=producer_identity, + source_identity=source_identity, + dependency_identity=dependency_identity, + feature_extractor_identity=feature_extractor_identity, + payload=bytes(encoded), + schema_identity=schema_identity, + ) + + +def composite_snapshot( + components: Sequence[DiagnosticSnapshot], + *, + producer_identity: str, + source_identity: str, + dependency_identity: str, + feature_extractor_identity: str, + schema_version: int = 1, + schema_identity: str = "mnel-diagnostic-snapshot", +) -> DiagnosticSnapshot: + if not components or len(components) > MAX_COMPOSITE_COMPONENTS: + raise SnapshotError("composite requires a bounded non-empty component set") + if any(not isinstance(component, DiagnosticSnapshot) for component in components): + raise SnapshotError("composite components must be diagnostic snapshots") + for component in components: + component.validate_integrity() + identities = [component.snapshot_identity for component in components] + if len(set(identities)) != len(identities): + raise SnapshotError("composite component identities must be unique") + encoded = bytearray(b"MNEL-C1" + struct.pack(">H", len(identities))) + for identity in identities: + value = identity.encode("ascii") + if len(value) > 128: + raise SnapshotError("composite component identity exceeds its ceiling") + encoded.extend(struct.pack(">H", len(value))) + encoded.extend(value) + return DiagnosticSnapshot.build( + snapshot_type="composite", + schema_version=schema_version, + producer_identity=producer_identity, + source_identity=source_identity, + dependency_identity=dependency_identity, + feature_extractor_identity=feature_extractor_identity, + payload=bytes(encoded), + schema_identity=schema_identity, + ) + + +def _decode_pair_payload(payload: bytes, magic: bytes) -> tuple[bytes, bytes]: + if len(payload) > MAX_SNAPSHOT_BYTES or len(payload) < 11 or payload[:7] != magic: + raise SnapshotError("invalid pair magic or truncated header") + left_length, right_length = struct.unpack(">HH", payload[7:11]) + if left_length == 0 or right_length == 0: + raise SnapshotError("pair members must be non-empty") + end = 11 + left_length + right_length + if end != len(payload): + raise SnapshotError("pair payload is truncated or has trailing bytes") + return bytes(payload[11 : 11 + left_length]), bytes(payload[11 + left_length : end]) diff --git a/tests/test_snapshots.py b/tests/test_snapshots.py index e0505f1..161bf90 100644 --- a/tests/test_snapshots.py +++ b/tests/test_snapshots.py @@ -1,6 +1,21 @@ import unittest -from mnel.snapshots import SnapshotError, pair_snapshot, tabular_snapshot, transition_snapshot +from mnel.snapshots import ( + SnapshotError, + SnapshotStore, + composite_snapshot, + decode_composite, + decode_graph, + decode_pair, + decode_snapshot, + decode_tabular, + decode_trace, + graph_snapshot, + pair_snapshot, + tabular_snapshot, + trace_snapshot, + transition_snapshot, +) class SnapshotTests(unittest.TestCase): @@ -34,9 +49,56 @@ def test_tabular_payload_is_binary_bounded_and_rejects_nonfinite_values(self) -> snapshot = tabular_snapshot(((1.0, 2.0), (3.0, 4.0)), **self._kwargs()) self.assertEqual(snapshot.payload[:4], b"MNET") self.assertGreater(len(snapshot.payload), 4) + self.assertEqual(decode_tabular(snapshot.payload).rows[1], (3.0, 4.0)) with self.assertRaises(SnapshotError): tabular_snapshot(((float("nan"),),), **self._kwargs()) + def test_all_views_round_trip_through_shared_store(self) -> None: + transition = transition_snapshot(b"a", b"b", **self._kwargs()) + pair = pair_snapshot(b"left", b"right", **self._kwargs()) + trace = trace_snapshot((("start", b"1"), ("stop", b"2")), **self._kwargs()) + graph = graph_snapshot( + ((2, "target"), (1, "source")), + ((1, 2, "calls"),), + **self._kwargs(), + ) + composite = composite_snapshot((transition, graph), **self._kwargs()) + self.assertEqual(decode_snapshot(transition).next_state, b"b") + self.assertEqual(decode_pair(pair.payload).right, b"right") + self.assertEqual(len(decode_trace(trace.payload).events), 2) + self.assertEqual(decode_graph(graph.payload).edges[0].edge_type, "calls") + self.assertEqual(decode_composite(composite.payload).component_snapshot_identities, (transition.snapshot_identity, graph.snapshot_identity)) + store = SnapshotStore() + for snapshot in (transition, pair, trace, graph, composite): + store.register(snapshot) + self.assertEqual(store.view(graph.snapshot_identity, accepted_types=("graph",)).nodes[0].node_id, 1) + with self.assertRaises(SnapshotError): + store.view(graph.snapshot_identity, accepted_types=("tabular",)) + + def test_malformed_views_fail_closed_and_store_rejects_tampering(self) -> None: + transition = transition_snapshot(b"a", b"b", **self._kwargs()) + with self.assertRaises(SnapshotError): + decode_pair(transition.payload) + with self.assertRaises(SnapshotError): + decode_trace(b"MNEL-R1\x00\x01\x05") + with self.assertRaises(SnapshotError): + decode_graph(b"MNEL-G1\x00\x01\x00\x00\x00") + with self.assertRaises(SnapshotError): + decode_composite(b"MNEL-C1\x00\x01\x00\x05abc") + tampered = transition.payload[:-1] + object.__setattr__(transition, "payload", tampered) + with self.assertRaises(SnapshotError): + SnapshotStore().register(transition) + + def test_graph_trace_and_composite_limits_are_explicit(self) -> None: + with self.assertRaises(SnapshotError): + trace_snapshot(tuple(("x", b"") for _ in range(257)), **self._kwargs()) + with self.assertRaises(SnapshotError): + graph_snapshot(((1, "one"),), ((1, 2, "missing"),), **self._kwargs()) + component = transition_snapshot(b"a", b"b", **self._kwargs()) + with self.assertRaises(SnapshotError): + composite_snapshot(tuple(component for _ in range(2)), **self._kwargs()) + if __name__ == "__main__": unittest.main() From a113b04b1647ae02c1429b04f086e52e37316d4d Mon Sep 17 00:00:00 2001 From: epi13 Date: Sat, 8 Aug 2026 18:02:30 -0800 Subject: [PATCH 2/2] Implement bounded Forge diagnostic lifecycle --- CHANGELOG.md | 11 + README.md | 22 +- docs/ARCHITECTURE.md | 16 + docs/INTEGRATIONS.md | 10 + docs/LEARNED_PROVIDER_RUNTIME.md | 9 +- docs/ROADMAP.md | 18 +- docs/THREAT_MODEL.md | 8 + schemas/mnel-forge-lifecycle.schema.json | 211 ++++ src/mnel/cli.py | 8 + src/mnel/forge_lifecycle.py | 1280 ++++++++++++++++++++++ src/mnel/integrations.py | 12 + tests/test_forge_lifecycle.py | 242 ++++ 12 files changed, 1834 insertions(+), 13 deletions(-) create mode 100644 schemas/mnel-forge-lifecycle.schema.json create mode 100644 src/mnel/forge_lifecycle.py create mode 100644 tests/test_forge_lifecycle.py diff --git a/CHANGELOG.md b/CHANGELOG.md index ef48c50..ef5c67f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,16 @@ # Changelog +## 0.3.0a0 — unreleased + +- Add bounded binary views and an identity-keyed shared store for transition, pair, + tabular, trace, graph, and composite diagnostic snapshots. +- Add an explicit MNEL-side Forge lifecycle surface: verifier registry and declarations, + preconditions, bounded probe requests, diagnostic-only witnesses, reference verifiers, + registered mutations, independent comparison, health/quarantine, coverage, learned + observation events, and omitted-question candidates. +- Add a deterministic `mnel forge-reference` study path and lifecycle schema. This is a + local diagnostic reference surface, not an implementation of external Forge authority. + ## 0.2.0a0 — unreleased - Add backend-neutral CPU, full-CUDA, and sequential CPU offload placement policy with diff --git a/README.md b/README.md index d5f9972..eb9c637 100644 --- a/README.md +++ b/README.md @@ -37,13 +37,15 @@ experience, negative memory, causal attribution, transfer-gated principles, reus strategies, and append-only candidate lineage rather than relying exclusively on conventional neural-weight training. -> **Current status:** functional `0.2.0a0` iteration. The repository now includes a +> **Current status:** functional `0.3.0a0` iteration. The repository now includes a > backend-neutral accelerator placement policy, optional Torch/Accelerate adapter, > process-local persistent Rust host, reusable identity-bound snapshots, bounded and > normalized diagnostic results, failure quarantine, an executable Rust HMM baseline, > deterministic runtime measurements, and bounded investigator context/workspace > contracts, an executable local-harness/worktree path, and a validated Rust v1 dynamic -> provider loader. It still does not provide process isolation, +> provider loader, and an executable bounded Forge-oriented diagnostic lifecycle with +> compact snapshot views, a verifier registry, reference probes, witnesses, mutations, +> comparison, health, coverage, and omitted-question candidates. It still does not provide process isolation, > unattended model execution, distributed scheduling, protected final custody, formal > MNCS/MNCDS conformance, or automatic RAVEL promotion. @@ -104,6 +106,14 @@ copy their authority or silently create substitute implementations. integration through the existing provider host; - initial immutable transition, tabular, and pair diagnostic snapshot producers with compact binary payloads and dependency-bound content identities; +- bounded transition, pair, tabular, trace, graph, and composite snapshot views backed by + one identity-keyed immutable snapshot store; +- explicit diagnostic verifier declarations and registry matching, bounded preconditions, + proposal-bound probe requests, diagnostic-only witnesses, and deterministic reference + transition/tabular/pair/trace/graph verifiers; +- registered-only mutation operators, independent witness comparison, verifier health and + quarantine state, coverage records, learned-provider observation normalization, and + proposal-only omitted-question candidates; - deterministic reference workflow, JSON schemas, mutation-oriented tests, and CI. ## Install @@ -170,6 +180,14 @@ The demo preregisters a bounded experiment, records an observation, evaluates ha gates, attributes the intervention, and creates a provisional principle proposal. It does not call a model or modify RAVEL. +Run the bounded 0.3 diagnostic lifecycle, including snapshot production, two independent +reference witnesses, a mutation, comparison, health/coverage records, and a next-question +candidate: + +```bash +mnel forge-reference --workspace build/forge-reference +``` + Verify and summarize the resulting ledger: ```bash diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index abc6807..020d559 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -72,6 +72,22 @@ sequential CPU offload. The latter keeps weights in system RAM and temporarily e modules on CUDA; it does not reload a provider for each query. Placement decisions are resource-accounted and diagnostic, never evaluator decisions. +### Bounded Forge-oriented diagnostic lifecycle + +The executable local lifecycle is deliberately split into independent records: + +```text +identified snapshot -> compatible verifier -> precondition report -> bounded probe + -> diagnostic witness -> optional registered mutation -> independent comparison + -> health/coverage -> proposal-only omitted-question candidate +``` + +`src/mnel/snapshots.py` provides compact immutable views shared by deterministic +verifiers and learned providers. `src/mnel/forge_lifecycle.py` provides the explicit +registry and reference execution surface. Witnesses characterize observations; they +contain no evaluator verdict field. Health describes execution reliability, not truth, +and comparisons preserve disagreement rather than voting it away. + ### Probe plane Forge supplies small, identity-bearing questions and witnesses. The stable interface is diff --git a/docs/INTEGRATIONS.md b/docs/INTEGRATIONS.md index 55db372..17a7bb9 100644 --- a/docs/INTEGRATIONS.md +++ b/docs/INTEGRATIONS.md @@ -27,6 +27,16 @@ subject identities, expected witness type, resource budget, mutation prohibition provider identity. Large prose scans should be decomposed into bounded witnesses where possible. +MNEL now includes a small local reference surface in `mnel.forge_lifecycle`. It is used +for deterministic tests and the `mnel forge-reference` command: it provides explicit +verifier declarations, bounded snapshot views, preconditions, witnesses, registered +mutations, independent comparison, health, and coverage. It is an adapter/test surface, +not a substitute Forge implementation and does not claim MNCS/MNCDS conformance. + +The external `mncs-forge-mcp` checkout is optional. The MNEL adapter contract is +identity-bound and provider-neutral; no developer-local Forge path is a runtime +dependency, and no hidden network or model service is invoked by the reference study. + ## MNCS Fabric Fabric distributes identified experiment bundles, captures node capabilities, and diff --git a/docs/LEARNED_PROVIDER_RUNTIME.md b/docs/LEARNED_PROVIDER_RUNTIME.md index a124203..7bcffe4 100644 --- a/docs/LEARNED_PROVIDER_RUNTIME.md +++ b/docs/LEARNED_PROVIDER_RUNTIME.md @@ -146,8 +146,8 @@ and a production accelerator backend remain future work. The C ABI v1 remains un ## Snapshot transport -The initial Python snapshot producers construct bounded transition, pair, or tabular -snapshots once. Each immutable payload is binary-friendly and carries producer, source, +The initial Python snapshot producers construct bounded transition, pair, tabular, trace, +graph, or composite snapshots once. Each immutable payload is binary-friendly and carries producer, source, dependency, feature-extractor, schema, and payload identities. Compatible deterministic probes and learned providers can consume the same payload boundary; changing a material dependency changes the content identity and prevents silent reuse. Forge or another @@ -158,6 +158,11 @@ compact binary bytes with explicit schema and feature-extractor identities. Any change to source, dependency, extractor, normalization, toolchain, or environment invalidates reuse unless the dependency envelope proves the snapshot unaffected. +The shared `SnapshotStore` exposes validated immutable views to compatible consumers, so +the same identified payload can be reused by a deterministic micro-verifier and a +learned provider without reparsing ad hoc JSON. Composite snapshots reference component +identities rather than duplicating their payloads. + ## Native-language exceptions A non-Rust provider may enter `native-trusted` only when its manifest includes: diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index 1f5f7f8..4413f22 100644 --- a/docs/ROADMAP.md +++ b/docs/ROADMAP.md @@ -40,18 +40,18 @@ ## 0.3 — Forge experiment lifecycle -- micro-verifier registry; -- probe preconditions and witness schemas; -- counterfactual and mutation probe support; -- independent-probe comparison; -- verifier health and coverage records; -- skeptic-driven omitted-question discovery; -- **Started:** identity-bound transition, tabular, and pair diagnostic snapshots with +- **Implemented:** explicit micro-verifier registry and identity-bound declarations; +- **Implemented:** probe preconditions, proposal-bound requests, and diagnostic-only witness schemas; +- **Implemented:** bounded counterfactual and registered mutation probe support; +- **Implemented:** independent-probe comparison preserving agreement, disagreement, and incomplete coverage; +- **Implemented:** verifier health, quarantine, and coverage records; +- **Started:** deterministic skeptic-driven omitted-question candidate discovery; +- **Implemented:** identity-bound transition, tabular, pair, trace, graph, and composite diagnostic snapshots with immutable compact binary payloads, producer/source/dependency/extractor identities, and deterministic content identities suitable for deterministic probes and learned micro-providers; -- compact binary snapshot views shared across compatible providers; -- learned observations normalized as diagnostic events without verifier status. +- **Implemented:** compact binary snapshot views shared across compatible providers; +- **Implemented:** learned observations normalized as diagnostic events without verifier status. ## 0.4 — verified distillation and learned-provider studies diff --git a/docs/THREAT_MODEL.md b/docs/THREAT_MODEL.md index c1f8784..999b407 100644 --- a/docs/THREAT_MODEL.md +++ b/docs/THREAT_MODEL.md @@ -63,6 +63,14 @@ dependency, extractor, producer, schema, and payload identities in the content i material dependency changes therefore invalidate reuse rather than silently transferring stale diagnostic context. +The 0.3 diagnostic lifecycle fails closed on malformed snapshot bytes, incompatible +verifiers, unavailable preconditions, malformed verifier output, budget exhaustion, and +repeated verifier errors. A verifier may be quarantined for runtime reliability without +being treated as truthful. Mutation operators are a fixed registered set; arbitrary +callbacks and in-place authoritative snapshot mutation are not accepted. Learned +provider observations and verifier witnesses remain distinct diagnostic records, and +neither can authorize conformance or promotion. + ### Apparent independence Multiple local machines run the same operator-controlled stack. This is replication, diff --git a/schemas/mnel-forge-lifecycle.schema.json b/schemas/mnel-forge-lifecycle.schema.json new file mode 100644 index 0000000..327d03b --- /dev/null +++ b/schemas/mnel-forge-lifecycle.schema.json @@ -0,0 +1,211 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "https://github.com/epi13/Machine-Native-Experimental-Learning/schemas/mnel-forge-lifecycle.schema.json", + "title": "MNEL bounded Forge-oriented diagnostic lifecycle records", + "$defs": { + "authority": {"const": "diagnostic-only"}, + "precondition": { + "type": "object", + "additionalProperties": false, + "required": ["kind", "value"], + "properties": { + "kind": {"type": "string", "minLength": 1}, + "value": {"type": ["string", "integer", "boolean"]} + } + }, + "limits": { + "type": "object", + "additionalProperties": false, + "required": ["operation_limit", "wall_time_ms", "output_bytes"], + "properties": { + "operation_limit": {"type": "integer", "minimum": 1}, + "wall_time_ms": {"type": "integer", "minimum": 1}, + "output_bytes": {"type": "integer", "minimum": 1} + } + }, + "verifierDeclaration": { + "type": "object", + "additionalProperties": false, + "required": ["schema", "verifier_id", "verifier_version", "implementation_identity", "accepted_snapshot_types", "accepted_schema_versions", "required_preconditions", "expected_witness_type", "resource_limits", "deterministic", "mutation_capability", "authority", "declaration_identity"], + "properties": { + "schema": {"const": "mnel-verifier-declaration/0.3"}, + "verifier_id": {"type": "string", "minLength": 1}, + "verifier_version": {"type": "string", "minLength": 1}, + "implementation_identity": {"type": "string", "minLength": 1}, + "accepted_snapshot_types": {"type": "array", "minItems": 1, "items": {"type": "string", "minLength": 1}}, + "accepted_schema_versions": {"type": "array", "minItems": 1, "items": {"type": "integer", "minimum": 1}}, + "required_preconditions": {"type": "array", "items": {"$ref": "#/$defs/precondition"}}, + "expected_witness_type": {"type": "string", "minLength": 1}, + "resource_limits": {"$ref": "#/$defs/limits"}, + "deterministic": {"type": "boolean"}, + "mutation_capability": {"type": "boolean"}, + "authority": {"$ref": "#/$defs/authority"}, + "declaration_identity": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"} + } + }, + "probeRequest": { + "type": "object", + "additionalProperties": false, + "required": ["schema", "probe_id", "question", "question_identity", "subject_identities", "verifier_id", "snapshot_identities", "expected_witness_type", "preconditions", "resource_budget", "mutation_policy", "runtime_identity", "lineage", "parameters", "authority", "request_identity"], + "properties": { + "schema": {"const": "mnel-forge-probe-request/0.3"}, + "probe_id": {"type": "string", "minLength": 1}, + "question": {"type": "string", "minLength": 1}, + "question_identity": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"}, + "subject_identities": {"type": "object", "minProperties": 1, "additionalProperties": {"type": "string", "minLength": 1}}, + "verifier_id": {"type": "string", "minLength": 1}, + "snapshot_identities": {"type": "array", "minItems": 1, "items": {"type": "string", "minLength": 1}}, + "expected_witness_type": {"type": "string", "minLength": 1}, + "preconditions": {"type": "array", "items": {"$ref": "#/$defs/precondition"}}, + "resource_budget": {"$ref": "#/$defs/limits"}, + "mutation_policy": {"enum": ["forbidden", "registered-only"]}, + "runtime_identity": {"type": "object", "additionalProperties": {"type": "string"}}, + "lineage": {"type": "object", "additionalProperties": {"type": "string"}}, + "parameters": {"type": "object"}, + "authority": {"const": "proposal-only"}, + "semantics": {"const": "diagnostic-request; not-a-verdict"}, + "request_identity": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"} + } + }, + "witness": { + "type": "object", + "additionalProperties": false, + "required": ["schema", "probe_identity", "question_identity", "verifier_id", "verifier_version", "implementation_identity", "snapshot_identities", "expected_witness_type", "precondition_report", "execution_status", "diagnostic_output", "resource_usage", "mutation_identity", "error", "authority", "semantics", "witness_identity"], + "properties": { + "schema": {"const": "mnel-diagnostic-witness/0.3"}, + "probe_identity": {"type": "string", "minLength": 1}, + "question_identity": {"type": "string", "minLength": 1}, + "verifier_id": {"type": "string", "minLength": 1}, + "verifier_version": {"type": "string", "minLength": 1}, + "implementation_identity": {"type": "string", "minLength": 1}, + "snapshot_identities": {"type": "array", "minItems": 1, "items": {"type": "string", "minLength": 1}}, + "expected_witness_type": {"type": "string", "minLength": 1}, + "precondition_report": {"type": "object"}, + "execution_status": {"enum": ["completed", "ineligible", "not-applicable", "abstained", "error", "budget-exceeded", "unavailable", "quarantined"]}, + "diagnostic_output": {"type": "object", "not": {"required": ["verdict"]}}, + "resource_usage": {"type": "object"}, + "mutation_identity": {"type": ["string", "null"]}, + "error": {"type": ["string", "null"]}, + "authority": {"$ref": "#/$defs/authority"}, + "semantics": {"const": "not-a-verdict"}, + "witness_identity": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"} + } + }, + "mutation": { + "type": "object", + "additionalProperties": false, + "required": ["schema", "original_snapshot_identity", "mutation_operator_identity", "parameters", "resulting_snapshot_identity", "scope", "authority", "semantics", "mutation_identity"], + "properties": { + "schema": {"const": "mnel-mutation-record/0.3"}, + "original_snapshot_identity": {"type": "string", "minLength": 1}, + "mutation_operator_identity": {"type": "string", "minLength": 1}, + "parameters": {"type": "object"}, + "resulting_snapshot_identity": {"type": "string", "minLength": 1}, + "scope": {"type": "string", "minLength": 1}, + "authority": {"$ref": "#/$defs/authority"}, + "semantics": {"const": "diagnostic-experiment; not-a-verdict"}, + "mutation_identity": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"} + } + }, + "comparison": { + "type": "object", + "additionalProperties": false, + "required": ["schema", "question_identity", "subject_identities", "witness_identities", "verifier_ids", "comparison_status", "disagreement_fields", "authority", "semantics", "comparison_identity"], + "properties": { + "schema": {"const": "mnel-witness-comparison/0.3"}, + "question_identity": {"type": "string", "minLength": 1}, + "subject_identities": {"type": "object"}, + "witness_identities": {"type": "array", "minItems": 2, "items": {"type": "string", "minLength": 1}}, + "verifier_ids": {"type": "array", "minItems": 2, "items": {"type": "string", "minLength": 1}}, + "comparison_status": {"enum": ["agreement", "disagreement", "incomplete"]}, + "disagreement_fields": {"type": "array", "items": {"type": "string"}}, + "authority": {"$ref": "#/$defs/authority"}, + "semantics": {"const": "evidence-characterization; not-a-verdict"}, + "comparison_identity": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"} + } + }, + "health": { + "type": "object", + "additionalProperties": false, + "required": ["schema", "verifier_id", "successful_executions", "execution_errors", "precondition_exclusions", "abstentions", "budget_violations", "malformed_outputs", "latency_ns_total", "latency_samples", "snapshot_types", "quarantined", "quarantine_reason", "authority", "semantics", "health_identity"], + "properties": { + "schema": {"const": "mnel-verifier-health/0.3"}, + "verifier_id": {"type": "string", "minLength": 1}, + "successful_executions": {"type": "integer", "minimum": 0}, + "execution_errors": {"type": "integer", "minimum": 0}, + "precondition_exclusions": {"type": "integer", "minimum": 0}, + "abstentions": {"type": "integer", "minimum": 0}, + "budget_violations": {"type": "integer", "minimum": 0}, + "malformed_outputs": {"type": "integer", "minimum": 0}, + "latency_ns_total": {"type": "integer", "minimum": 0}, + "latency_samples": {"type": "integer", "minimum": 0}, + "snapshot_types": {"type": "array", "items": {"type": "string"}}, + "quarantined": {"type": "boolean"}, + "quarantine_reason": {"type": ["string", "null"]}, + "authority": {"$ref": "#/$defs/authority"}, + "semantics": {"const": "health-is-not-truth"}, + "health_identity": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"} + } + }, + "coverage": { + "type": "object", + "additionalProperties": false, + "required": ["schema", "registered_snapshot_types", "exercised_snapshot_types", "exercised_verifier_ids", "uncovered_snapshot_types", "single_source_question_identities", "authority", "semantics", "coverage_identity"], + "properties": { + "schema": {"const": "mnel-verifier-coverage/0.3"}, + "registered_snapshot_types": {"type": "array", "items": {"type": "string"}}, + "exercised_snapshot_types": {"type": "array", "items": {"type": "string"}}, + "exercised_verifier_ids": {"type": "array", "items": {"type": "string"}}, + "uncovered_snapshot_types": {"type": "array", "items": {"type": "string"}}, + "single_source_question_identities": {"type": "array", "items": {"type": "string"}}, + "authority": {"$ref": "#/$defs/authority"}, + "semantics": {"const": "coverage-is-not-truth"}, + "coverage_identity": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"} + } + }, + "learnedObservation": { + "type": "object", + "additionalProperties": false, + "required": ["schema", "record_type", "provider_id", "provider_observation_identity", "snapshot_identities", "declaration_identity", "source_record_ids", "payload", "authority", "semantics", "event_identity"], + "properties": { + "schema": {"const": "mnel-learned-provider-observation-event/0.3"}, + "record_type": {"const": "learned-provider-observation"}, + "provider_id": {"type": "string", "minLength": 1}, + "provider_observation_identity": {"type": "string", "minLength": 1}, + "snapshot_identities": {"type": "array", "minItems": 1, "items": {"type": "string", "minLength": 1}}, + "declaration_identity": {"type": "string", "minLength": 1}, + "source_record_ids": {"type": "array", "items": {"type": "string"}}, + "payload": {"type": "object"}, + "authority": {"$ref": "#/$defs/authority"}, + "semantics": {"const": "not-a-verdict"}, + "event_identity": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"} + } + }, + "questionCandidate": { + "type": "object", + "additionalProperties": false, + "required": ["schema", "subject_identity", "reason", "target_snapshot_type", "supporting_record_ids", "authority", "semantics", "candidate_identity"], + "properties": { + "schema": {"const": "mnel-omitted-question-candidate/0.3"}, + "subject_identity": {"type": "string", "minLength": 1}, + "reason": {"type": "string", "minLength": 1}, + "target_snapshot_type": {"type": ["string", "null"]}, + "supporting_record_ids": {"type": "array", "items": {"type": "string"}}, + "authority": {"const": "proposal-only"}, + "semantics": {"const": "proposal-only; not-a-verdict"}, + "candidate_identity": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"} + } + } + }, + "oneOf": [ + {"$ref": "#/$defs/verifierDeclaration"}, + {"$ref": "#/$defs/probeRequest"}, + {"$ref": "#/$defs/witness"}, + {"$ref": "#/$defs/mutation"}, + {"$ref": "#/$defs/comparison"}, + {"$ref": "#/$defs/health"}, + {"$ref": "#/$defs/coverage"}, + {"$ref": "#/$defs/learnedObservation"}, + {"$ref": "#/$defs/questionCandidate"} + ] +} diff --git a/src/mnel/cli.py b/src/mnel/cli.py index 9871899..7a9f261 100644 --- a/src/mnel/cli.py +++ b/src/mnel/cli.py @@ -11,6 +11,7 @@ from . import __version__ from .core import EvidenceLedger, run_reference_study +from .forge_lifecycle import run_reference_forge_study from .investigators import DEFAULT_ROLE_CONTRACTS from .learned_providers import ( DEFAULT_LEARNED_PROVIDER_REGISTRY, @@ -58,6 +59,10 @@ def parser() -> argparse.ArgumentParser: command.add_argument("path") demo = commands.add_parser("demo") demo.add_argument("--workspace", default="build/demo") + forge_reference = commands.add_parser( + "forge-reference", description="Run the deterministic MNEL Forge diagnostic lifecycle" + ) + forge_reference.add_argument("--workspace", default=None) return root @@ -122,5 +127,8 @@ def main(argv: list[str] | None = None) -> int: ) print(json.dumps(payload, indent=2, sort_keys=True)) return 0 if result.valid else 1 + if args.command == "forge-reference": + print(json.dumps(run_reference_forge_study(args.workspace), indent=2, sort_keys=True)) + return 0 print(json.dumps(run_reference_study(args.workspace), indent=2, sort_keys=True)) return 0 diff --git a/src/mnel/forge_lifecycle.py b/src/mnel/forge_lifecycle.py new file mode 100644 index 0000000..55474e7 --- /dev/null +++ b/src/mnel/forge_lifecycle.py @@ -0,0 +1,1280 @@ +"""Bounded MNEL-side Forge diagnostic lifecycle. + +This module provides explicit registry, probe, witness, mutation, comparison, health, +coverage, and question-candidate contracts. It is a deterministic reference surface, not +an implementation of the external MNCS Forge authority. +""" + +from __future__ import annotations + +import math +import time +from dataclasses import dataclass, field +from enum import StrEnum +from pathlib import Path +from typing import Any, Protocol, Sequence + +from .core import EvidenceLedger, canonical_digest, canonical_json +from .snapshots import ( + DiagnosticSnapshot, + GraphView, + PairView, + SnapshotError, + SnapshotStore, + SnapshotView, + TabularView, + TraceView, + TransitionView, + decode_snapshot, + graph_snapshot, + pair_snapshot, + tabular_snapshot, + trace_snapshot, + transition_snapshot, +) + + +AUTHORITY_DIAGNOSTIC_ONLY = "diagnostic-only" +AUTHORITY_PROPOSAL_ONLY = "proposal-only" +SEMANTICS_NOT_A_VERDICT = "not-a-verdict" +FORBIDDEN_AUTHORITY_KEYS = frozenset( + { + "verdict", + "conformance", + "promotion_authorized", + "promotion", + "evaluator_verdict", + "evaluator_eligible", + "future_final", + "hidden_transfer", + "ravel_promotion", + "verifier_authority", + "evaluator_authority", + "mncs_verdict", + "mncds_verdict", + } +) + + +class ForgeLifecycleError(ValueError): + pass + + +class ProbeExecutionStatus(StrEnum): + COMPLETED = "completed" + INELIGIBLE = "ineligible" + NOT_APPLICABLE = "not-applicable" + ABSTAINED = "abstained" + ERROR = "error" + BUDGET_EXCEEDED = "budget-exceeded" + UNAVAILABLE = "unavailable" + QUARANTINED = "quarantined" + + +class MutationPolicy(StrEnum): + FORBIDDEN = "forbidden" + REGISTERED_ONLY = "registered-only" + + +class VerifierState(StrEnum): + ENABLED = "enabled" + DISABLED = "disabled" + QUARANTINED = "quarantined" + + +def _reject_authority(value: Any) -> None: + if isinstance(value, dict): + for key, child in value.items(): + if key in FORBIDDEN_AUTHORITY_KEYS: + raise ForgeLifecycleError(f"diagnostic record contains forbidden field: {key}") + if key == "authority" and child not in {AUTHORITY_DIAGNOSTIC_ONLY, AUTHORITY_PROPOSAL_ONLY}: + raise ForgeLifecycleError("diagnostic record attempted to expand authority") + _reject_authority(child) + elif isinstance(value, list): + for child in value: + _reject_authority(child) + + +def _nonempty(value: str, label: str) -> str: + if not isinstance(value, str) or not value.strip(): + raise ForgeLifecycleError(f"{label} is required") + return value + + +def _bounded_dict(value: dict[str, Any], label: str, max_bytes: int = 16 * 1024) -> dict[str, Any]: + _reject_authority(value) + try: + encoded = canonical_json(value) + except (TypeError, ValueError) as error: + raise ForgeLifecycleError(f"{label} is not canonical JSON") from error + if len(encoded) > max_bytes: + raise ForgeLifecycleError(f"{label} exceeds its byte ceiling") + return dict(value) + + +@dataclass(frozen=True, slots=True) +class Precondition: + kind: str + value: str | int | bool + + def __post_init__(self) -> None: + _nonempty(self.kind, "precondition kind") + if not isinstance(self.value, (str, int, bool)): + raise ForgeLifecycleError("precondition values must be scalar") + + def to_dict(self) -> dict[str, object]: + return {"kind": self.kind, "value": self.value} + + +@dataclass(frozen=True, slots=True) +class PreconditionOutcome: + kind: str + expected: str | int | bool + actual: object + satisfied: bool + availability: str = "available" + + def to_dict(self) -> dict[str, object]: + return { + "kind": self.kind, + "expected": self.expected, + "actual": self.actual, + "satisfied": self.satisfied, + "availability": self.availability, + } + + +@dataclass(frozen=True, slots=True) +class PreconditionReport: + status: str + outcomes: tuple[PreconditionOutcome, ...] + + @property + def satisfied(self) -> bool: + return self.status == "satisfied" + + def to_dict(self) -> dict[str, object]: + return {"status": self.status, "outcomes": [item.to_dict() for item in self.outcomes]} + + +@dataclass(frozen=True, slots=True) +class VerifierDeclaration: + verifier_id: str + verifier_version: str + implementation_identity: str + accepted_snapshot_types: tuple[str, ...] + accepted_schema_versions: tuple[int, ...] + required_preconditions: tuple[Precondition, ...] + expected_witness_type: str + resource_limits: dict[str, int] + deterministic: bool + mutation_capability: bool + authority: str = AUTHORITY_DIAGNOSTIC_ONLY + + def __post_init__(self) -> None: + for value, label in ( + (self.verifier_id, "verifier id"), + (self.verifier_version, "verifier version"), + (self.implementation_identity, "implementation identity"), + (self.expected_witness_type, "expected witness type"), + ): + _nonempty(value, label) + if not self.accepted_snapshot_types or not self.accepted_schema_versions: + raise ForgeLifecycleError("verifiers must declare accepted snapshot types and schema versions") + if any(not item.strip() for item in self.accepted_snapshot_types): + raise ForgeLifecycleError("verifier snapshot types must be non-empty") + if any(not isinstance(item, int) or item < 1 for item in self.accepted_schema_versions): + raise ForgeLifecycleError("verifier schema versions must be positive integers") + required = {"operation_limit", "wall_time_ms", "output_bytes"} + if set(self.resource_limits) != required or any( + not isinstance(value, int) or value < 1 for value in self.resource_limits.values() + ): + raise ForgeLifecycleError("verifier resource limits must declare positive operation, wall, and output bounds") + if self.authority != AUTHORITY_DIAGNOSTIC_ONLY: + raise ForgeLifecycleError("micro-verifiers are diagnostic-only") + + @property + def declaration_identity(self) -> str: + return canonical_digest(self.to_dict(include_identity=False)) + + @classmethod + def from_dict(cls, value: dict[str, Any]) -> "VerifierDeclaration": + if not isinstance(value, dict): + raise ForgeLifecycleError("verifier declaration must be an object") + if value.get("schema") != "mnel-verifier-declaration/0.3": + raise ForgeLifecycleError("unsupported verifier declaration schema") + raw_preconditions = value.get("required_preconditions", ()) + if not isinstance(raw_preconditions, (list, tuple)) or any( + not isinstance(item, dict) or "kind" not in item or "value" not in item + for item in raw_preconditions + ): + raise ForgeLifecycleError("verifier preconditions are malformed") + declaration = cls( + verifier_id=value.get("verifier_id"), + verifier_version=value.get("verifier_version"), + implementation_identity=value.get("implementation_identity"), + accepted_snapshot_types=tuple(value.get("accepted_snapshot_types", ())), + accepted_schema_versions=tuple(value.get("accepted_schema_versions", ())), + required_preconditions=tuple(Precondition(item["kind"], item["value"]) for item in raw_preconditions), + expected_witness_type=value.get("expected_witness_type"), + resource_limits=dict(value.get("resource_limits", {})), + deterministic=value.get("deterministic"), + mutation_capability=value.get("mutation_capability"), + authority=value.get("authority", ""), + ) + supplied_identity = value.get("declaration_identity") + if supplied_identity is not None and supplied_identity != declaration.declaration_identity: + raise ForgeLifecycleError("verifier declaration identity does not match its content") + return declaration + + def to_dict(self, *, include_identity: bool = True) -> dict[str, object]: + value: dict[str, object] = { + "schema": "mnel-verifier-declaration/0.3", + "verifier_id": self.verifier_id, + "verifier_version": self.verifier_version, + "implementation_identity": self.implementation_identity, + "accepted_snapshot_types": list(self.accepted_snapshot_types), + "accepted_schema_versions": list(self.accepted_schema_versions), + "required_preconditions": [item.to_dict() for item in self.required_preconditions], + "expected_witness_type": self.expected_witness_type, + "resource_limits": dict(self.resource_limits), + "deterministic": self.deterministic, + "mutation_capability": self.mutation_capability, + "authority": self.authority, + } + if include_identity: + value["declaration_identity"] = self.declaration_identity + return value + + +class ReferenceVerifier(Protocol): + def run(self, view: SnapshotView, parameters: dict[str, Any], operation_limit: int) -> dict[str, Any]: ... + + +class VerifierRegistry: + def __init__(self) -> None: + self._declarations: dict[str, VerifierDeclaration] = {} + self._implementations: dict[str, ReferenceVerifier] = {} + self._states: dict[str, VerifierState] = {} + self._state_reasons: dict[str, str] = {} + + def register(self, declaration: VerifierDeclaration, implementation: ReferenceVerifier) -> str: + identity = declaration.declaration_identity + existing = self._declarations.get(declaration.verifier_id) + if existing is not None: + if existing.declaration_identity != identity: + raise ForgeLifecycleError("verifier id collision with different declaration identity") + raise ForgeLifecycleError("duplicate verifier registration") + if not callable(getattr(implementation, "run", None)): + raise ForgeLifecycleError("verifier implementation must expose run") + self._declarations[declaration.verifier_id] = declaration + self._implementations[declaration.verifier_id] = implementation + self._states[declaration.verifier_id] = VerifierState.ENABLED + return identity + + def register_declaration(self, declaration: VerifierDeclaration) -> str: + """Register a declaration without an executable implementation.""" + + identity = declaration.declaration_identity + if declaration.verifier_id in self._declarations: + raise ForgeLifecycleError("duplicate verifier registration") + self._declarations[declaration.verifier_id] = declaration + self._states[declaration.verifier_id] = VerifierState.DISABLED + self._state_reasons[declaration.verifier_id] = "no local implementation attached" + return identity + + def load(self, declarations: Sequence[dict[str, Any]]) -> tuple[str, ...]: + """Load explicit declaration objects; executable implementations remain separate.""" + + identities = [] + for value in declarations: + declaration = VerifierDeclaration.from_dict(value) + identities.append(self.register_declaration(declaration)) + return tuple(identities) + + def lookup(self, verifier_id: str) -> VerifierDeclaration: + try: + return self._declarations[verifier_id] + except KeyError as error: + raise ForgeLifecycleError(f"unknown verifier: {verifier_id}") from error + + def implementation(self, verifier_id: str) -> ReferenceVerifier: + try: + return self._implementations[verifier_id] + except KeyError as error: + raise ForgeLifecycleError(f"verifier has no executable implementation: {verifier_id}") from error + + def set_state(self, verifier_id: str, state: VerifierState, reason: str = "") -> None: + self.lookup(verifier_id) + if state is VerifierState.QUARANTINED and not reason.strip(): + raise ForgeLifecycleError("quarantine requires a reason") + self._states[verifier_id] = state + if reason: + self._state_reasons[verifier_id] = reason + + def state(self, verifier_id: str) -> VerifierState: + self.lookup(verifier_id) + return self._states[verifier_id] + + def declarations(self) -> tuple[VerifierDeclaration, ...]: + return tuple(self._declarations[key] for key in sorted(self._declarations)) + + def match(self, snapshot: DiagnosticSnapshot, context: dict[str, Any] | None = None) -> tuple[VerifierDeclaration, ...]: + snapshot.validate_integrity() + decode_snapshot(snapshot) + matches: list[VerifierDeclaration] = [] + for declaration in self.declarations(): + if self._states[declaration.verifier_id] is not VerifierState.ENABLED: + continue + if snapshot.snapshot_type not in declaration.accepted_snapshot_types: + continue + if snapshot.schema_version not in declaration.accepted_schema_versions: + continue + report = evaluate_preconditions( + declaration.required_preconditions, + snapshot, + decode_snapshot(snapshot), + context or {}, + ) + if report.satisfied: + matches.append(declaration) + return tuple(matches) + + def to_dict(self) -> dict[str, object]: + return { + "schema": "mnel-verifier-registry/0.3", + "verifiers": [ + { + **declaration.to_dict(), + "state": self._states[declaration.verifier_id].value, + "state_reason": self._state_reasons.get(declaration.verifier_id), + } + for declaration in self.declarations() + ], + "authority": AUTHORITY_DIAGNOSTIC_ONLY, + "registry_identity": canonical_digest( + [declaration.to_dict() for declaration in self.declarations()] + ), + } + + +@dataclass(frozen=True, slots=True) +class ProbeRequest: + probe_id: str + question: str + subject_identities: dict[str, str] + verifier_id: str + snapshot_identities: tuple[str, ...] + expected_witness_type: str + preconditions: tuple[Precondition, ...] + resource_budget: dict[str, int] + mutation_policy: MutationPolicy + runtime_identity: dict[str, str] + lineage: dict[str, str] + parameters: dict[str, Any] = field(default_factory=dict) + authority: str = AUTHORITY_PROPOSAL_ONLY + + def __post_init__(self) -> None: + for value, label in ( + (self.probe_id, "probe id"), + (self.question, "probe question"), + (self.verifier_id, "verifier id"), + (self.expected_witness_type, "expected witness type"), + ): + _nonempty(value, label) + if not self.snapshot_identities or any( + not isinstance(item, str) or not item.strip() for item in self.snapshot_identities + ): + raise ForgeLifecycleError("probe requests require snapshot identities") + _bounded_dict(self.subject_identities, "subject identities") + _bounded_dict(self.runtime_identity, "runtime identity") + _bounded_dict(self.lineage, "probe lineage") + _bounded_dict(self.parameters, "probe parameters") + required = {"operation_limit", "wall_time_ms", "output_bytes"} + if set(self.resource_budget) != required or any( + not isinstance(value, int) or value < 1 for value in self.resource_budget.values() + ): + raise ForgeLifecycleError("probe budget must declare positive operation, wall, and output bounds") + if self.authority != AUTHORITY_PROPOSAL_ONLY: + raise ForgeLifecycleError("investigator probe requests are proposal-only") + + @property + def question_identity(self) -> str: + return canonical_digest( + { + "question": self.question, + "subject_identities": self.subject_identities, + "snapshot_identities": self.snapshot_identities, + "parameters": self.parameters, + } + ) + + @property + def request_identity(self) -> str: + return canonical_digest(self.to_dict(include_identity=False)) + + def to_dict(self, *, include_identity: bool = True) -> dict[str, object]: + value: dict[str, object] = { + "schema": "mnel-forge-probe-request/0.3", + "probe_id": self.probe_id, + "question": self.question, + "question_identity": self.question_identity, + "subject_identities": dict(self.subject_identities), + "verifier_id": self.verifier_id, + "snapshot_identities": list(self.snapshot_identities), + "expected_witness_type": self.expected_witness_type, + "preconditions": [item.to_dict() for item in self.preconditions], + "resource_budget": dict(self.resource_budget), + "mutation_policy": self.mutation_policy.value, + "runtime_identity": dict(self.runtime_identity), + "lineage": dict(self.lineage), + "parameters": dict(self.parameters), + "authority": self.authority, + "semantics": "diagnostic-request; not-a-verdict", + } + if include_identity: + value["request_identity"] = self.request_identity + return value + + +def evaluate_preconditions( + preconditions: Sequence[Precondition], + snapshot: DiagnosticSnapshot, + view: SnapshotView, + context: dict[str, Any], +) -> PreconditionReport: + outcomes: list[PreconditionOutcome] = [] + for condition in preconditions: + actual: object + available = "available" + if condition.kind == "required_snapshot_type": + actual = snapshot.snapshot_type + elif condition.kind == "required_schema_version": + actual = snapshot.schema_version + elif condition.kind == "dependency_identity": + actual = snapshot.dependency_identity + elif condition.kind == "source_language": + actual = context.get("source_language") + elif condition.kind == "feature_available": + features = context.get("features", ()) + actual = condition.value in features + elif condition.kind == "tool_available": + tools = context.get("tools", ()) + actual = condition.value in tools + elif condition.kind == "mutation_allowed": + actual = context.get("mutation_allowed") + elif condition.kind == "minimum_rows": + actual = len(view.rows) if isinstance(view, TabularView) else None + elif condition.kind == "minimum_nodes": + actual = len(view.nodes) if isinstance(view, GraphView) else None + elif condition.kind == "minimum_events": + actual = len(view.events) if isinstance(view, TraceView) else None + else: + actual = None + available = "unavailable" + satisfied = _precondition_matches(condition, actual, available) + outcomes.append(PreconditionOutcome(condition.kind, condition.value, actual, satisfied, available)) + if any(item.availability == "unavailable" for item in outcomes): + status = "unavailable" + elif all(item.satisfied for item in outcomes): + status = "satisfied" + else: + status = "failed" + return PreconditionReport(status, tuple(outcomes)) + + +def _precondition_matches(condition: Precondition, actual: object, availability: str) -> bool: + if availability != "available": + return False + if condition.kind.startswith("minimum_"): + return isinstance(actual, int) and isinstance(condition.value, int) and actual >= condition.value + return actual == condition.value + + +@dataclass(frozen=True, slots=True) +class Witness: + probe_identity: str + question_identity: str + verifier_id: str + verifier_version: str + implementation_identity: str + snapshot_identities: tuple[str, ...] + expected_witness_type: str + precondition_report: PreconditionReport + execution_status: ProbeExecutionStatus + diagnostic_output: dict[str, Any] + resource_usage: dict[str, int] + mutation_identity: str | None = None + error: str | None = None + authority: str = AUTHORITY_DIAGNOSTIC_ONLY + semantics: str = SEMANTICS_NOT_A_VERDICT + + def __post_init__(self) -> None: + _bounded_dict(self.diagnostic_output, "witness diagnostic output", 32 * 1024) + _bounded_dict(self.resource_usage, "witness resource usage", 4096) + if self.authority != AUTHORITY_DIAGNOSTIC_ONLY or self.semantics != SEMANTICS_NOT_A_VERDICT: + raise ForgeLifecycleError("witness authority is fixed to diagnostic-only") + + @property + def witness_identity(self) -> str: + identity_body = self.to_dict(include_identity=False) + # Timing is evidence, but wall-clock timing is not a stable content identity. + identity_body["resource_usage"] = { + key: value for key, value in self.resource_usage.items() if key != "elapsed_ns" + } + return canonical_digest(identity_body) + + def to_dict(self, *, include_identity: bool = True) -> dict[str, object]: + value: dict[str, object] = { + "schema": "mnel-diagnostic-witness/0.3", + "probe_identity": self.probe_identity, + "question_identity": self.question_identity, + "verifier_id": self.verifier_id, + "verifier_version": self.verifier_version, + "implementation_identity": self.implementation_identity, + "snapshot_identities": list(self.snapshot_identities), + "expected_witness_type": self.expected_witness_type, + "precondition_report": self.precondition_report.to_dict(), + "execution_status": self.execution_status.value, + "diagnostic_output": dict(self.diagnostic_output), + "resource_usage": dict(self.resource_usage), + "mutation_identity": self.mutation_identity, + "error": self.error, + "authority": self.authority, + "semantics": self.semantics, + } + if include_identity: + value["witness_identity"] = self.witness_identity + return value + + +class TransitionChangeVerifier: + def run(self, view: SnapshotView, parameters: dict[str, Any], operation_limit: int) -> dict[str, Any]: + if not isinstance(view, TransitionView): + raise ForgeLifecycleError("transition verifier received an incompatible view") + if operation_limit < 1: + raise ForgeLifecycleError("operation budget exhausted") + changed = view.previous_state != view.next_state + return {"condition_observed": changed, "changed": changed, "member_bytes": len(view.previous_state) + len(view.next_state)} + + +class TabularBoundsVerifier: + def run(self, view: SnapshotView, parameters: dict[str, Any], operation_limit: int) -> dict[str, Any]: + if not isinstance(view, TabularView): + raise ForgeLifecycleError("tabular verifier received an incompatible view") + minimum = parameters.get("minimum") + maximum = parameters.get("maximum") + if ( + not isinstance(minimum, (int, float)) + or isinstance(minimum, bool) + or not isinstance(maximum, (int, float)) + or isinstance(maximum, bool) + or not math.isfinite(minimum) + or not math.isfinite(maximum) + or minimum > maximum + ): + raise ForgeLifecycleError("tabular verifier requires finite minimum and maximum parameters") + values = [value for row in view.rows for value in row] + if len(values) > operation_limit: + raise ForgeLifecycleError("tabular verifier operation budget exceeded") + outliers = [value for value in values if value < minimum or value > maximum] + return {"condition_observed": not outliers, "outlier_count": len(outliers), "minimum": minimum, "maximum": maximum} + + +class PairRelationVerifier: + def run(self, view: SnapshotView, parameters: dict[str, Any], operation_limit: int) -> dict[str, Any]: + if not isinstance(view, PairView): + raise ForgeLifecycleError("pair verifier received an incompatible view") + relation = parameters.get("relation", "equal") + if relation not in {"equal", "different"}: + raise ForgeLifecycleError("pair relation must be equal or different") + equal = view.left == view.right + observed = equal if relation == "equal" else not equal + return {"condition_observed": observed, "equal": equal, "relation": relation} + + +class TraceOrderVerifier: + def run(self, view: SnapshotView, parameters: dict[str, Any], operation_limit: int) -> dict[str, Any]: + if not isinstance(view, TraceView): + raise ForgeLifecycleError("trace verifier received an incompatible view") + before = parameters.get("before") + after = parameters.get("after") + if not isinstance(before, str) or not isinstance(after, str): + raise ForgeLifecycleError("trace verifier requires before and after event types") + if len(view.events) > operation_limit: + raise ForgeLifecycleError("trace verifier operation budget exceeded") + before_positions = [index for index, event in enumerate(view.events) if event.event_type == before] + after_positions = [index for index, event in enumerate(view.events) if event.event_type == after] + observed = bool(before_positions and after_positions and min(before_positions) < max(after_positions)) + return {"condition_observed": observed, "before": before, "after": after, "event_count": len(view.events)} + + +class GraphEdgeVerifier: + def run(self, view: SnapshotView, parameters: dict[str, Any], operation_limit: int) -> dict[str, Any]: + if not isinstance(view, GraphView): + raise ForgeLifecycleError("graph verifier received an incompatible view") + source, target, edge_type = parameters.get("source"), parameters.get("target"), parameters.get("edge_type") + if not isinstance(source, int) or not isinstance(target, int) or not isinstance(edge_type, str): + raise ForgeLifecycleError("graph verifier requires source, target, and edge_type") + if len(view.edges) > operation_limit: + raise ForgeLifecycleError("graph verifier operation budget exceeded") + observed = any(edge.source == source and edge.target == target and edge.edge_type == edge_type for edge in view.edges) + return {"condition_observed": observed, "source": source, "target": target, "edge_type": edge_type} + + +def reference_verifier_registry() -> VerifierRegistry: + registry = VerifierRegistry() + common = {"operation_limit": 4096, "wall_time_ms": 100, "output_bytes": 8192} + registry.register( + VerifierDeclaration("transition-change", "0.3.0", "mnel-reference-transition/1", ("transition",), (1,), (), "transition-witness", common, True, False), + TransitionChangeVerifier(), + ) + registry.register( + VerifierDeclaration("transition-change-independent", "0.3.0", "mnel-reference-transition-independent/1", ("transition",), (1,), (), "transition-witness", common, True, False), + TransitionChangeVerifier(), + ) + registry.register( + VerifierDeclaration("tabular-bounds", "0.3.0", "mnel-reference-tabular-bounds/1", ("tabular",), (1,), (), "tabular-witness", common, True, False), + TabularBoundsVerifier(), + ) + registry.register( + VerifierDeclaration("pair-relation", "0.3.0", "mnel-reference-pair-relation/1", ("pair",), (1,), (), "pair-witness", common, True, False), + PairRelationVerifier(), + ) + registry.register( + VerifierDeclaration("trace-order", "0.3.0", "mnel-reference-trace-order/1", ("trace",), (1,), (), "trace-witness", common, True, False), + TraceOrderVerifier(), + ) + registry.register( + VerifierDeclaration("graph-edge", "0.3.0", "mnel-reference-graph-edge/1", ("graph",), (1,), (), "graph-witness", common, True, False), + GraphEdgeVerifier(), + ) + return registry + + +class VerifierHealthStore: + def __init__(self, quarantine_after_errors: int = 3) -> None: + if quarantine_after_errors < 1: + raise ForgeLifecycleError("quarantine_after_errors must be positive") + self.quarantine_after_errors = quarantine_after_errors + self._records: dict[str, dict[str, Any]] = {} + + def _record(self, verifier_id: str) -> dict[str, Any]: + return self._records.setdefault( + verifier_id, + { + "successful_executions": 0, + "execution_errors": 0, + "precondition_exclusions": 0, + "abstentions": 0, + "budget_violations": 0, + "malformed_outputs": 0, + "latency_ns_total": 0, + "latency_samples": 0, + "snapshot_types": set(), + "quarantined": False, + "quarantine_reason": None, + }, + ) + + def observe(self, verifier_id: str, snapshot_type: str, status: ProbeExecutionStatus, elapsed_ns: int) -> None: + record = self._record(verifier_id) + record["latency_ns_total"] += max(0, elapsed_ns) + record["latency_samples"] += 1 + record["snapshot_types"].add(snapshot_type) + if status is ProbeExecutionStatus.COMPLETED: + record["successful_executions"] += 1 + elif status in {ProbeExecutionStatus.INELIGIBLE, ProbeExecutionStatus.NOT_APPLICABLE}: + record["precondition_exclusions"] += 1 + elif status is ProbeExecutionStatus.ABSTAINED: + record["abstentions"] += 1 + elif status is ProbeExecutionStatus.BUDGET_EXCEEDED: + record["budget_violations"] += 1 + elif status is ProbeExecutionStatus.ERROR: + record["execution_errors"] += 1 + if record["execution_errors"] >= self.quarantine_after_errors: + record["quarantined"] = True + record["quarantine_reason"] = "repeated verifier execution errors" + + def quarantine(self, verifier_id: str, reason: str) -> None: + _nonempty(reason, "quarantine reason") + record = self._record(verifier_id) + record["quarantined"] = True + record["quarantine_reason"] = reason + + def is_quarantined(self, verifier_id: str) -> bool: + return bool(self._record(verifier_id)["quarantined"]) + + def malformed_output(self, verifier_id: str, reason: str) -> None: + _nonempty(reason, "malformed output reason") + record = self._record(verifier_id) + record["malformed_outputs"] += 1 + + def to_dict(self, verifier_id: str) -> dict[str, object]: + record = self._record(verifier_id) + value = { + "schema": "mnel-verifier-health/0.3", + "verifier_id": verifier_id, + **{key: (sorted(item) if isinstance(item, set) else item) for key, item in record.items()}, + "authority": AUTHORITY_DIAGNOSTIC_ONLY, + "semantics": "health-is-not-truth", + } + value["health_identity"] = canonical_digest(value) + return value + + +class ReferenceForgeRuntime: + def __init__(self, snapshots: SnapshotStore, verifiers: VerifierRegistry, health: VerifierHealthStore | None = None) -> None: + self.snapshots = snapshots + self.verifiers = verifiers + self.health = health or VerifierHealthStore() + + def execute(self, request: ProbeRequest, *, mutation_identity: str | None = None) -> Witness: + started = time.monotonic_ns() + if mutation_identity is not None and request.mutation_policy is MutationPolicy.FORBIDDEN: + return self._witness( + request, + None, + ProbeExecutionStatus.NOT_APPLICABLE, + {}, + {}, + "probe request forbids mutation execution", + started, + mutation_identity, + ) + try: + declaration = self.verifiers.lookup(request.verifier_id) + except ForgeLifecycleError as error: + return self._witness(request, None, ProbeExecutionStatus.UNAVAILABLE, {}, {}, str(error), started, mutation_identity) + state = self.verifiers.state(request.verifier_id) + if state is not VerifierState.ENABLED: + status = ProbeExecutionStatus.QUARANTINED if state is VerifierState.QUARANTINED else ProbeExecutionStatus.UNAVAILABLE + return self._witness(request, declaration, status, {}, {}, f"verifier state is {state.value}", started, mutation_identity) + if declaration.expected_witness_type != request.expected_witness_type: + return self._witness(request, declaration, ProbeExecutionStatus.NOT_APPLICABLE, {}, {}, "witness type mismatch", started, mutation_identity) + try: + snapshot = self.snapshots.get(request.snapshot_identities[0]) + if len(request.snapshot_identities) != 1: + raise ForgeLifecycleError("reference verifiers require exactly one snapshot") + view = self.snapshots.view( + snapshot.snapshot_identity, + accepted_types=declaration.accepted_snapshot_types, + schema_versions=declaration.accepted_schema_versions, + ) + report = evaluate_preconditions(request.preconditions + declaration.required_preconditions, snapshot, view, request.parameters) + if report.status != "satisfied": + status = ProbeExecutionStatus.INELIGIBLE if report.status == "failed" else ProbeExecutionStatus.NOT_APPLICABLE + return self._witness(request, declaration, status, {}, {"precondition_report": report.to_dict()}, None, started, mutation_identity, report) + operation_limit = min( + request.resource_budget["operation_limit"], declaration.resource_limits["operation_limit"] + ) + output = self.verifiers.implementation(request.verifier_id).run( + view, request.parameters, operation_limit + ) + if not isinstance(output, dict): + self.health.malformed_output(request.verifier_id, "verifier output is not an object") + raise ForgeLifecycleError("verifier output must be an object") + operation_count = output.get("operation_count", 0) + if ( + not isinstance(operation_count, int) + or isinstance(operation_count, bool) + or operation_count < 0 + or operation_count > operation_limit + ): + self.health.malformed_output(request.verifier_id, "verifier operation count is invalid") + raise ForgeLifecycleError("verifier operation count exceeds its budget") + encoded = canonical_json(output) + if len(encoded) > min( + request.resource_budget["output_bytes"], declaration.resource_limits["output_bytes"] + ): + self.health.malformed_output(request.verifier_id, "verifier output exceeds its byte budget") + raise ForgeLifecycleError("verifier output exceeds its byte budget") + elapsed = time.monotonic_ns() - started + wall_limit_ms = min(request.resource_budget["wall_time_ms"], declaration.resource_limits["wall_time_ms"]) + if elapsed > wall_limit_ms * 1_000_000: + status = ProbeExecutionStatus.BUDGET_EXCEEDED + else: + status = ProbeExecutionStatus.COMPLETED + witness = self._witness( + request, + declaration, + status, + output, + {"operation_count": operation_count, "output_bytes": len(encoded)}, + None, + started, + mutation_identity, + report, + ) + self.health.observe(request.verifier_id, snapshot.snapshot_type, status, elapsed) + if self.health.is_quarantined(request.verifier_id): + self.verifiers.set_state( + request.verifier_id, + VerifierState.QUARANTINED, + "repeated verifier execution errors", + ) + return witness + except (ForgeLifecycleError, SnapshotError, TypeError, ValueError) as error: + snapshot_type = snapshot.snapshot_type if "snapshot" in locals() else "unknown" + status = ProbeExecutionStatus.BUDGET_EXCEEDED if "budget" in str(error).lower() else ProbeExecutionStatus.ERROR + self.health.observe(request.verifier_id, snapshot_type, status, time.monotonic_ns() - started) + if self.health.is_quarantined(request.verifier_id): + self.verifiers.set_state( + request.verifier_id, + VerifierState.QUARANTINED, + "repeated verifier execution errors", + ) + return self._witness(request, declaration, status, {}, {}, str(error), started, mutation_identity) + + def _witness( + self, + request: ProbeRequest, + declaration: VerifierDeclaration | None, + status: ProbeExecutionStatus, + output: dict[str, Any], + usage: dict[str, Any], + error: str | None, + started: int, + mutation_identity: str | None, + report: PreconditionReport | None = None, + ) -> Witness: + report = report or PreconditionReport("unavailable", ()) + return Witness( + probe_identity=request.request_identity, + question_identity=request.question_identity, + verifier_id=request.verifier_id, + verifier_version=declaration.verifier_version if declaration else "unknown", + implementation_identity=declaration.implementation_identity if declaration else "unknown", + snapshot_identities=request.snapshot_identities, + expected_witness_type=request.expected_witness_type, + precondition_report=report, + execution_status=status, + diagnostic_output=output, + resource_usage={"elapsed_ns": time.monotonic_ns() - started, **usage}, + mutation_identity=mutation_identity, + error=error, + ) + + +@dataclass(frozen=True, slots=True) +class WitnessComparison: + question_identity: str + subject_identities: dict[str, str] + witness_identities: tuple[str, ...] + verifier_ids: tuple[str, ...] + comparison_status: str + disagreement_fields: tuple[str, ...] + authority: str = AUTHORITY_DIAGNOSTIC_ONLY + semantics: str = "evidence-characterization; not-a-verdict" + + def __post_init__(self) -> None: + if self.authority != AUTHORITY_DIAGNOSTIC_ONLY: + raise ForgeLifecycleError("witness comparisons are diagnostic-only") + if self.semantics != "evidence-characterization; not-a-verdict": + raise ForgeLifecycleError("witness comparison semantics are fixed") + + @property + def comparison_identity(self) -> str: + return canonical_digest(self.to_dict(include_identity=False)) + + def to_dict(self, *, include_identity: bool = True) -> dict[str, object]: + value: dict[str, object] = { + "schema": "mnel-witness-comparison/0.3", + "question_identity": self.question_identity, + "subject_identities": dict(self.subject_identities), + "witness_identities": list(self.witness_identities), + "verifier_ids": list(self.verifier_ids), + "comparison_status": self.comparison_status, + "disagreement_fields": list(self.disagreement_fields), + "authority": self.authority, + "semantics": self.semantics, + } + if include_identity: + value["comparison_identity"] = self.comparison_identity + return value + + +def compare_witnesses(witnesses: Sequence[Witness], subject_identities: dict[str, str]) -> WitnessComparison: + if len(witnesses) < 2: + raise ForgeLifecycleError("independent comparison requires at least two witnesses") + if len({item.question_identity for item in witnesses}) != 1: + raise ForgeLifecycleError("witnesses address different identified questions") + completed = [item for item in witnesses if item.execution_status is ProbeExecutionStatus.COMPLETED] + signatures = [ + canonical_json({key: item.diagnostic_output.get(key) for key in ("condition_observed", "relation", "equal")}) + for item in completed + ] + if len(completed) != len(witnesses): + status = "incomplete" + fields = ("execution_status",) + elif len(set(signatures)) == 1: + status = "agreement" + fields = () + else: + status = "disagreement" + fields = ("diagnostic_output",) + return WitnessComparison( + witnesses[0].question_identity, + dict(subject_identities), + tuple(item.witness_identity for item in witnesses), + tuple(sorted(item.verifier_id for item in witnesses)), + status, + fields, + ) + + +@dataclass(frozen=True, slots=True) +class MutationRecord: + original_snapshot_identity: str + mutation_operator_identity: str + parameters: dict[str, Any] + resulting_snapshot_identity: str + scope: str + authority: str = AUTHORITY_DIAGNOSTIC_ONLY + + def __post_init__(self) -> None: + for value, label in ( + (self.original_snapshot_identity, "original snapshot identity"), + (self.mutation_operator_identity, "mutation operator identity"), + (self.resulting_snapshot_identity, "resulting snapshot identity"), + (self.scope, "mutation scope"), + ): + _nonempty(value, label) + _bounded_dict(self.parameters, "mutation parameters") + if self.authority != AUTHORITY_DIAGNOSTIC_ONLY: + raise ForgeLifecycleError("mutations are diagnostic-only") + + @property + def mutation_identity(self) -> str: + return canonical_digest(self.to_dict(include_identity=False)) + + def to_dict(self, *, include_identity: bool = True) -> dict[str, object]: + value: dict[str, object] = { + "schema": "mnel-mutation-record/0.3", + "original_snapshot_identity": self.original_snapshot_identity, + "mutation_operator_identity": self.mutation_operator_identity, + "parameters": dict(self.parameters), + "resulting_snapshot_identity": self.resulting_snapshot_identity, + "scope": self.scope, + "authority": self.authority, + "semantics": "diagnostic-experiment; not-a-verdict", + } + if include_identity: + value["mutation_identity"] = self.mutation_identity + return value + + +class MutationRegistry: + OPERATORS = frozenset({"transition.swap", "trace.remove-event", "trace.swap-events", "tabular.perturb", "graph.remove-edge", "pair.swap"}) + + def apply(self, operator_id: str, snapshot: DiagnosticSnapshot, parameters: dict[str, Any], store: SnapshotStore) -> MutationRecord: + if operator_id not in self.OPERATORS: + raise ForgeLifecycleError("mutation operator is not registered") + snapshot.validate_integrity() + view = decode_snapshot(snapshot) + kwargs = { + "producer_identity": f"{snapshot.producer_identity}|mutation:{operator_id}", + "source_identity": snapshot.source_identity, + "dependency_identity": snapshot.dependency_identity, + "feature_extractor_identity": snapshot.feature_extractor_identity, + "schema_version": snapshot.schema_version, + "schema_identity": snapshot.schema_identity, + } + if operator_id == "transition.swap" and isinstance(view, TransitionView): + mutated = transition_snapshot(view.next_state, view.previous_state, **kwargs) + elif operator_id == "pair.swap" and isinstance(view, PairView): + mutated = pair_snapshot(view.right, view.left, **kwargs) + elif operator_id == "trace.remove-event" and isinstance(view, TraceView): + index = _parameter_index(parameters, len(view.events)) + events = view.events[:index] + view.events[index + 1 :] + if not events: + raise ForgeLifecycleError("trace mutation may not remove its final event") + mutated = trace_snapshot(tuple((event.event_type, event.payload) for event in events), **kwargs) + elif operator_id == "trace.swap-events" and isinstance(view, TraceView): + first = _parameter_index(parameters, len(view.events)) + second = _parameter_index(parameters, len(view.events), "second_index") + if first == second: + raise ForgeLifecycleError("trace swap requires two distinct event indexes") + events = list(view.events) + events[first], events[second] = events[second], events[first] + mutated = trace_snapshot(tuple((event.event_type, event.payload) for event in events), **kwargs) + elif operator_id == "tabular.perturb" and isinstance(view, TabularView): + row = _parameter_index(parameters, len(view.rows), "row") + column = _parameter_index(parameters, view.column_count, "column") + delta = parameters.get("delta") + if not isinstance(delta, (int, float)) or not math.isfinite(delta): + raise ForgeLifecycleError("tabular perturbation delta must be finite") + rows = [list(values) for values in view.rows] + rows[row][column] += float(delta) + mutated = tabular_snapshot(tuple(tuple(values) for values in rows), **kwargs) + elif operator_id == "graph.remove-edge" and isinstance(view, GraphView): + index = _parameter_index(parameters, len(view.edges)) + edges = view.edges[:index] + view.edges[index + 1 :] + mutated = graph_snapshot(tuple((node.node_id, node.label) for node in view.nodes), tuple((edge.source, edge.target, edge.edge_type) for edge in edges), **kwargs) + else: + raise ForgeLifecycleError("mutation operator is incompatible with the snapshot type") + store.register(mutated) + return MutationRecord(snapshot.snapshot_identity, f"mnel-mutation/{operator_id}/1", dict(parameters), mutated.snapshot_identity, "bounded registered mutation") + + +def _parameter_index(parameters: dict[str, Any], length: int, name: str = "index") -> int: + value = parameters.get(name) + if not isinstance(value, int) or isinstance(value, bool) or value < 0 or value >= length: + raise ForgeLifecycleError(f"mutation {name} is outside its bounded range") + return value + + +@dataclass(frozen=True, slots=True) +class LearnedDiagnosticEvent: + provider_id: str + provider_observation_identity: str + snapshot_identities: tuple[str, ...] + declaration_identity: str + source_record_ids: tuple[str, ...] + payload: dict[str, Any] + authority: str = AUTHORITY_DIAGNOSTIC_ONLY + semantics: str = SEMANTICS_NOT_A_VERDICT + + def __post_init__(self) -> None: + for value, label in ( + (self.provider_id, "provider id"), + (self.provider_observation_identity, "provider observation identity"), + (self.declaration_identity, "provider declaration identity"), + ): + _nonempty(value, label) + if not self.snapshot_identities or any( + not isinstance(item, str) or not item.strip() for item in self.snapshot_identities + ): + raise ForgeLifecycleError("learned observations require snapshot identities") + _bounded_dict(self.payload, "learned diagnostic payload") + if self.authority != AUTHORITY_DIAGNOSTIC_ONLY or self.semantics != SEMANTICS_NOT_A_VERDICT: + raise ForgeLifecycleError("learned provider observations are diagnostic-only") + + @property + def event_identity(self) -> str: + return canonical_digest(self.to_dict(include_identity=False)) + + @classmethod + def from_observation(cls, observation: dict[str, Any]) -> "LearnedDiagnosticEvent": + _reject_authority(observation) + required = ("provider_id", "observation_identity", "snapshot_ids", "declaration_identity") + if any(key not in observation for key in required): + raise ForgeLifecycleError("learned observation lacks identity-bound fields") + provider_id = observation["provider_id"] + observation_identity = observation["observation_identity"] + snapshot_ids = observation["snapshot_ids"] + declaration_identity = observation["declaration_identity"] + if not isinstance(snapshot_ids, (list, tuple)) or not snapshot_ids: + raise ForgeLifecycleError("learned observation requires snapshot identities") + for value, label in ( + (provider_id, "provider id"), + (observation_identity, "provider observation identity"), + (declaration_identity, "provider declaration identity"), + ): + _nonempty(value, label) + if any(not isinstance(item, str) or not item.strip() for item in snapshot_ids): + raise ForgeLifecycleError("learned observation snapshot identities must be non-empty strings") + payload = {key: value for key, value in observation.items() if key not in required} + _bounded_dict(payload, "learned diagnostic payload") + return cls( + provider_id, + observation_identity, + tuple(snapshot_ids), + declaration_identity, + tuple(str(item) for item in observation.get("source_record_ids", ())), + payload, + ) + + def to_dict(self, *, include_identity: bool = True) -> dict[str, object]: + value: dict[str, object] = { + "schema": "mnel-learned-provider-observation-event/0.3", + "record_type": "learned-provider-observation", + "provider_id": self.provider_id, + "provider_observation_identity": self.provider_observation_identity, + "snapshot_identities": list(self.snapshot_identities), + "declaration_identity": self.declaration_identity, + "source_record_ids": list(self.source_record_ids), + "payload": dict(self.payload), + "authority": self.authority, + "semantics": self.semantics, + } + if include_identity: + value["event_identity"] = self.event_identity + return value + + +@dataclass(frozen=True, slots=True) +class CoverageRecord: + registered_snapshot_types: tuple[str, ...] + exercised_snapshot_types: tuple[str, ...] + exercised_verifier_ids: tuple[str, ...] + uncovered_snapshot_types: tuple[str, ...] + single_source_question_identities: tuple[str, ...] + + def __post_init__(self) -> None: + if not self.registered_snapshot_types and self.exercised_snapshot_types: + raise ForgeLifecycleError("exercised coverage cannot exceed registered coverage") + + @property + def coverage_identity(self) -> str: + return canonical_digest(self.to_dict(include_identity=False)) + + def to_dict(self, *, include_identity: bool = True) -> dict[str, object]: + value: dict[str, object] = { + "schema": "mnel-verifier-coverage/0.3", + "registered_snapshot_types": list(self.registered_snapshot_types), + "exercised_snapshot_types": list(self.exercised_snapshot_types), + "exercised_verifier_ids": list(self.exercised_verifier_ids), + "uncovered_snapshot_types": list(self.uncovered_snapshot_types), + "single_source_question_identities": list(self.single_source_question_identities), + "authority": AUTHORITY_DIAGNOSTIC_ONLY, + "semantics": "coverage-is-not-truth", + } + if include_identity: + value["coverage_identity"] = self.coverage_identity + return value + + +def build_coverage(registry: VerifierRegistry, snapshots: SnapshotStore, witnesses: Sequence[Witness]) -> CoverageRecord: + registered = sorted( + { + item + for declaration in registry.declarations() + if registry.state(declaration.verifier_id) is VerifierState.ENABLED + for item in declaration.accepted_snapshot_types + } + ) + exercised = sorted({snapshots.get(identity).snapshot_type for witness in witnesses for identity in witness.snapshot_identities if witness.execution_status is ProbeExecutionStatus.COMPLETED}) + verifier_ids = sorted({witness.verifier_id for witness in witnesses if witness.execution_status is ProbeExecutionStatus.COMPLETED}) + question_verifiers: dict[str, set[str]] = {} + for witness in witnesses: + if witness.execution_status is ProbeExecutionStatus.COMPLETED: + question_verifiers.setdefault(witness.question_identity, set()).add(witness.verifier_id) + single = sorted(question for question, verifier_ids_for_question in question_verifiers.items() if len(verifier_ids_for_question) < 2) + return CoverageRecord(tuple(registered), tuple(exercised), tuple(verifier_ids), tuple(item for item in registered if item not in exercised), tuple(single)) + + +@dataclass(frozen=True, slots=True) +class QuestionCandidate: + subject_identity: str + reason: str + target_snapshot_type: str | None + supporting_record_ids: tuple[str, ...] + authority: str = AUTHORITY_PROPOSAL_ONLY + + def __post_init__(self) -> None: + _nonempty(self.subject_identity, "candidate subject identity") + _nonempty(self.reason, "candidate reason") + if self.authority != AUTHORITY_PROPOSAL_ONLY: + raise ForgeLifecycleError("question candidates are proposal-only") + + @property + def candidate_identity(self) -> str: + return canonical_digest(self.to_dict(include_identity=False)) + + def to_dict(self, *, include_identity: bool = True) -> dict[str, object]: + value: dict[str, object] = { + "schema": "mnel-omitted-question-candidate/0.3", + "subject_identity": self.subject_identity, + "reason": self.reason, + "target_snapshot_type": self.target_snapshot_type, + "supporting_record_ids": list(self.supporting_record_ids), + "authority": self.authority, + "semantics": "proposal-only; not-a-verdict", + } + if include_identity: + value["candidate_identity"] = self.candidate_identity + return value + + +def discover_question_candidates( + subject_identity: str, + snapshots: SnapshotStore, + coverage: CoverageRecord, + comparisons: Sequence[WitnessComparison] = (), +) -> tuple[QuestionCandidate, ...]: + candidates: list[QuestionCandidate] = [ + QuestionCandidate(subject_identity, "no compatible verifier exercised this snapshot type", snapshot_type, (coverage.coverage_identity,)) + for snapshot_type in coverage.uncovered_snapshot_types + ] + candidates.extend( + QuestionCandidate(subject_identity, "identified question has only one diagnostic verifier", None, (question_identity,)) + for question_identity in coverage.single_source_question_identities + ) + candidates.extend( + QuestionCandidate(subject_identity, "independent witnesses disagree", None, comparison.witness_identities) + for comparison in comparisons + if comparison.comparison_status == "disagreement" + ) + return tuple(sorted(candidates, key=lambda item: item.candidate_identity)) + + +def run_reference_forge_study(workspace: str | Path | None = None) -> dict[str, Any]: + identities = { + "producer_identity": "mnel-reference-producer/0.3", + "source_identity": "sha256:reference-source", + "dependency_identity": "sha256:reference-dependency", + "feature_extractor_identity": "sha256:reference-extractor", + } + store = SnapshotStore() + transition = transition_snapshot(b"cold", b"warm", **identities) + table = tabular_snapshot(((0.2, 0.4), (0.3, 0.5)), **identities) + trace = trace_snapshot((("start", b""), ("finish", b"")), **identities) + graph = graph_snapshot(((1, "source"), (2, "target")), ((1, 2, "calls"),), **identities) + for snapshot in (transition, table, trace, graph): + store.register(snapshot) + registry = reference_verifier_registry() + health = VerifierHealthStore() + runtime = ReferenceForgeRuntime(store, registry, health) + base = { + "subject_identities": {"source": identities["source_identity"]}, + "snapshot_identities": (transition.snapshot_identity,), + "preconditions": (), + "resource_budget": {"operation_limit": 100, "wall_time_ms": 1000, "output_bytes": 4096}, + "mutation_policy": MutationPolicy.REGISTERED_ONLY, + "runtime_identity": {"runtime": "mnel-reference-runtime/0.3"}, + "lineage": {"investigator_request": "sha256:reference-investigator-request"}, + "parameters": {}, + } + witness_a = runtime.execute(ProbeRequest("probe-transition-a", "did the identified transition change state?", verifier_id="transition-change", expected_witness_type="transition-witness", **base)) + witness_b = runtime.execute(ProbeRequest("probe-transition-b", "did the identified transition change state?", verifier_id="transition-change-independent", expected_witness_type="transition-witness", **base)) + table_witness = runtime.execute(ProbeRequest("probe-table", "are all table values bounded?", subject_identities=base["subject_identities"], verifier_id="tabular-bounds", snapshot_identities=(table.snapshot_identity,), expected_witness_type="tabular-witness", preconditions=(Precondition("minimum_rows", 2),), resource_budget=base["resource_budget"], mutation_policy=base["mutation_policy"], runtime_identity=base["runtime_identity"], lineage=base["lineage"], parameters={"minimum": 0.0, "maximum": 1.0})) + comparison = compare_witnesses((witness_a, witness_b), base["subject_identities"]) + mutation = MutationRegistry().apply("trace.swap-events", trace, {"index": 0, "second_index": 1}, store) + mutation_request = ProbeRequest( + "probe-trace-mutation", + "did the identified trace preserve the requested event order after mutation?", + subject_identities=base["subject_identities"], + verifier_id="trace-order", + snapshot_identities=(mutation.resulting_snapshot_identity,), + expected_witness_type="trace-witness", + preconditions=(), + resource_budget=base["resource_budget"], + mutation_policy=base["mutation_policy"], + runtime_identity=base["runtime_identity"], + lineage=base["lineage"], + parameters={"before": "start", "after": "finish"}, + ) + mutation_witness = runtime.execute(mutation_request, mutation_identity=mutation.mutation_identity) + learned = LearnedDiagnosticEvent.from_observation({"provider_id": "state.hidden-markov-model", "observation_identity": "sha256:learned-observation", "snapshot_ids": [transition.snapshot_identity], "declaration_identity": "sha256:learned-declaration", "value": 0.8, "out_of_distribution": False}) + witnesses = (witness_a, witness_b, table_witness, mutation_witness) + coverage = build_coverage(registry, store, witnesses) + candidates = discover_question_candidates(base["subject_identities"]["source"], store, coverage, (comparison,)) + records: list[dict[str, Any]] = [ + {"record_type": "diagnostic-snapshot", **snapshot.to_dict()} for snapshot in (transition, table, trace, graph, mutation_record_snapshot(store, mutation.resulting_snapshot_identity)) + ] + records.extend({"record_type": "verifier-witness", **witness.to_dict()} for witness in witnesses) + records.extend(({"record_type": "witness-comparison", **comparison.to_dict()}, {"record_type": "mutation", **mutation.to_dict()}, {"record_type": "learned-provider-observation", **learned.to_dict()}, {"record_type": "verifier-coverage", **coverage.to_dict()})) + records.extend({"record_type": "omitted-question-candidate", **candidate.to_dict()} for candidate in candidates) + result: dict[str, Any] = {"schema": "mnel-forge-reference-study/0.3", "records": records, "comparison": comparison.to_dict(), "mutation": mutation.to_dict(), "coverage": coverage.to_dict(), "health": [health.to_dict(item.verifier_id) for item in registry.declarations()], "question_candidates": [item.to_dict() for item in candidates], "authority": AUTHORITY_DIAGNOSTIC_ONLY, "study_identity": canonical_digest(records)} + if workspace is not None: + root = Path(workspace) + ledger = EvidenceLedger(root / "forge-evidence.jsonl") + for record in records: + ledger.append(record["record_type"], record, actor="mnel-reference-verifier") + result["ledger"] = ledger.summarize() + return result + + +def mutation_record_snapshot(store: SnapshotStore, identity: str) -> DiagnosticSnapshot: + return store.get(identity) diff --git a/src/mnel/integrations.py b/src/mnel/integrations.py index daad46e..a63bfa7 100644 --- a/src/mnel/integrations.py +++ b/src/mnel/integrations.py @@ -10,6 +10,7 @@ from typing import Any, Protocol, Sequence from .core import canonical_digest +from .forge_lifecycle import ProbeRequest as ForgeLifecycleProbeRequest from .investigator_harness import RuntimeIdentityEnvelope, WorkspaceAccess @@ -404,6 +405,17 @@ class ForgeProbeProvider(Protocol): def run_probe(self, request: ForgeProbeRequest) -> dict[str, Any]: ... +class ForgeLifecycleProvider(Protocol): + """MNEL-side boundary for an identity-bound diagnostic probe provider. + + Implementations may target the external Forge project or the local reference + surface. The returned record remains a witness/observation and never an evaluator + verdict. + """ + + def execute_probe(self, request: ForgeLifecycleProbeRequest) -> dict[str, Any]: ... + + @dataclass(frozen=True) class FabricExperimentRequest: experiment_id: str diff --git a/tests/test_forge_lifecycle.py b/tests/test_forge_lifecycle.py new file mode 100644 index 0000000..70a3378 --- /dev/null +++ b/tests/test_forge_lifecycle.py @@ -0,0 +1,242 @@ +import tempfile +import unittest +from dataclasses import replace +import json +from pathlib import Path + +from mnel.forge_lifecycle import ( + ForgeLifecycleError, + LearnedDiagnosticEvent, + MutationPolicy, + MutationRegistry, + Precondition, + ProbeExecutionStatus, + ProbeRequest, + ReferenceForgeRuntime, + VerifierDeclaration, + VerifierHealthStore, + VerifierRegistry, + VerifierState, + build_coverage, + compare_witnesses, + discover_question_candidates, + reference_verifier_registry, + run_reference_forge_study, +) +from mnel.snapshots import SnapshotStore, transition_snapshot + + +class _BrokenVerifier: + def run(self, view, parameters, operation_limit): + raise ForgeLifecycleError("fixture verifier failure") + + +class _MalformedVerifier: + def run(self, view, parameters, operation_limit): + return ["not", "an", "object"] + + +class ForgeLifecycleTests(unittest.TestCase): + def setUp(self) -> None: + self.identities = { + "producer_identity": "producer:v1", + "source_identity": "source:v1", + "dependency_identity": "dependency:v1", + "feature_extractor_identity": "extractor:v1", + } + self.snapshot = transition_snapshot(b"cold", b"warm", **self.identities) + self.store = SnapshotStore() + self.store.register(self.snapshot) + self.registry = reference_verifier_registry() + self.runtime = ReferenceForgeRuntime(self.store, self.registry) + + def _request(self, verifier_id="transition-change", **overrides): + values = { + "probe_id": "probe-1", + "question": "did this transition change state?", + "subject_identities": {"subject": "source:v1"}, + "verifier_id": verifier_id, + "snapshot_identities": (self.snapshot.snapshot_identity,), + "expected_witness_type": "transition-witness", + "preconditions": (), + "resource_budget": {"operation_limit": 100, "wall_time_ms": 1000, "output_bytes": 4096}, + "mutation_policy": MutationPolicy.REGISTERED_ONLY, + "runtime_identity": {"runtime": "fixture"}, + "lineage": {"request": "request:v1"}, + "parameters": {}, + } + values.update(overrides) + return ProbeRequest(**values) + + def _contains_forbidden_key(self, value) -> bool: + if isinstance(value, dict): + return any( + key in {"verdict", "conformance", "promotion_authorized"} + or self._contains_forbidden_key(child) + for key, child in value.items() + ) + if isinstance(value, list): + return any(self._contains_forbidden_key(child) for child in value) + return False + + def test_registry_rejects_collisions_and_matches_identity_bound_snapshot(self) -> None: + declaration = VerifierDeclaration( + "fixture", + "1", + "fixture-implementation", + ("transition",), + (1,), + (Precondition("required_snapshot_type", "transition"),), + "fixture-witness", + {"operation_limit": 10, "wall_time_ms": 100, "output_bytes": 1000}, + True, + False, + ) + registry = VerifierRegistry() + registry.register(declaration, _BrokenVerifier()) + with self.assertRaises(ForgeLifecycleError): + registry.register(declaration, _BrokenVerifier()) + self.assertEqual(registry.match(self.snapshot)[0].verifier_id, "fixture") + + changed = transition_snapshot(b"cold", b"warm", **{**self.identities, "dependency_identity": "dependency:v2"}) + self.assertNotEqual(changed.snapshot_identity, self.snapshot.snapshot_identity) + + loaded = VerifierRegistry() + loaded.load((declaration.to_dict(),)) + self.assertEqual(loaded.state("fixture"), VerifierState.DISABLED) + + def test_failed_precondition_is_ineligible_and_has_no_verdict(self) -> None: + request = self._request(preconditions=(Precondition("required_snapshot_type", "tabular"),)) + witness = self.runtime.execute(request) + self.assertEqual(witness.execution_status, ProbeExecutionStatus.INELIGIBLE) + self.assertNotIn("verdict", witness.to_dict()) + self.assertEqual(witness.precondition_report.status, "failed") + + def test_probe_parameters_cannot_expand_nested_authority(self) -> None: + with self.assertRaises(ForgeLifecycleError): + self._request(parameters={"nested": {"authority": "evaluator"}}) + + def test_reference_witness_is_bounded_and_identity_is_stable(self) -> None: + first = self.runtime.execute(self._request()) + second = self.runtime.execute(self._request()) + self.assertEqual(first.execution_status, ProbeExecutionStatus.COMPLETED) + self.assertEqual(first.witness_identity, second.witness_identity) + self.assertLessEqual(first.resource_usage["output_bytes"], 4096) + self.assertEqual(first.authority, "diagnostic-only") + self.assertEqual(first.semantics, "not-a-verdict") + + def test_mutation_preserves_original_and_changes_identity(self) -> None: + mutations = MutationRegistry() + record = mutations.apply("transition.swap", self.snapshot, {}, self.store) + self.assertNotEqual(record.original_snapshot_identity, record.resulting_snapshot_identity) + self.assertEqual(self.store.get(record.original_snapshot_identity).payload, self.snapshot.payload) + self.assertEqual(record.authority, "diagnostic-only") + with self.assertRaises(ForgeLifecycleError): + mutations.apply("arbitrary.python", self.snapshot, {}, self.store) + + def test_comparison_preserves_agreement_and_disagreement(self) -> None: + first = self.runtime.execute(self._request("transition-change")) + second = self.runtime.execute( + self._request("transition-change-independent", probe_id="probe-2") + ) + comparison = compare_witnesses((first, second), {"subject": "source:v1"}) + self.assertEqual(comparison.comparison_status, "agreement") + contrary = replace(second, diagnostic_output={"condition_observed": False}) + disagreement = compare_witnesses((first, contrary), {"subject": "source:v1"}) + self.assertEqual(disagreement.comparison_status, "disagreement") + self.assertNotIn("verdict", disagreement.to_dict()) + + def test_repeated_errors_quarantine_verifier(self) -> None: + declaration = VerifierDeclaration( + "broken", + "1", + "broken-implementation", + ("transition",), + (1,), + (), + "broken-witness", + {"operation_limit": 10, "wall_time_ms": 100, "output_bytes": 1000}, + True, + False, + ) + registry = VerifierRegistry() + registry.register(declaration, _BrokenVerifier()) + runtime = ReferenceForgeRuntime(self.store, registry, VerifierHealthStore(2)) + request = self._request("broken", expected_witness_type="broken-witness") + self.assertEqual(runtime.execute(request).execution_status, ProbeExecutionStatus.ERROR) + self.assertEqual(runtime.execute(request).execution_status, ProbeExecutionStatus.ERROR) + self.assertEqual(registry.state("broken"), VerifierState.QUARANTINED) + self.assertEqual(runtime.execute(request).execution_status, ProbeExecutionStatus.QUARANTINED) + + def test_malformed_verifier_output_fails_closed_and_is_counted(self) -> None: + declaration = VerifierDeclaration( + "malformed", + "1", + "malformed-implementation", + ("transition",), + (1,), + (), + "malformed-witness", + {"operation_limit": 10, "wall_time_ms": 100, "output_bytes": 1000}, + True, + False, + ) + registry = VerifierRegistry() + registry.register(declaration, _MalformedVerifier()) + health = VerifierHealthStore() + witness = ReferenceForgeRuntime(self.store, registry, health).execute( + self._request("malformed", expected_witness_type="malformed-witness") + ) + self.assertEqual(witness.execution_status, ProbeExecutionStatus.ERROR) + self.assertEqual(health.to_dict("malformed")["malformed_outputs"], 1) + + def test_learned_observation_is_not_a_verifier_witness(self) -> None: + event = LearnedDiagnosticEvent.from_observation( + { + "provider_id": "provider:v1", + "observation_identity": "observation:v1", + "snapshot_ids": [self.snapshot.snapshot_identity], + "declaration_identity": "declaration:v1", + "score": 0.75, + } + ) + value = event.to_dict() + self.assertEqual(value["record_type"], "learned-provider-observation") + self.assertNotIn("verifier_id", value) + with self.assertRaises(ForgeLifecycleError): + LearnedDiagnosticEvent.from_observation( + { + "provider_id": "provider:v1", + "observation_identity": "observation:v1", + "snapshot_ids": [self.snapshot.snapshot_identity], + "declaration_identity": "declaration:v1", + "verdict": "PASS", + } + ) + + def test_reference_study_writes_a_valid_append_only_ledger(self) -> None: + with tempfile.TemporaryDirectory() as directory: + result = run_reference_forge_study(Path(directory)) + self.assertEqual(result["comparison"]["comparison_status"], "agreement") + self.assertTrue(result["question_candidates"]) + self.assertTrue(Path(directory, "forge-evidence.jsonl").is_file()) + self.assertTrue(all(not self._contains_forbidden_key(record) for record in result["records"])) + + def test_lifecycle_schema_is_machine_readable(self) -> None: + root = Path(__file__).resolve().parents[1] + schema = json.loads((root / "schemas" / "mnel-forge-lifecycle.schema.json").read_text()) + self.assertEqual(schema["$schema"], "https://json-schema.org/draft/2020-12/schema") + self.assertIn("witness", schema["$defs"]) + self.assertIn("comparison", schema["$defs"]) + + def test_coverage_surfaces_single_source_questions(self) -> None: + witness = self.runtime.execute(self._request()) + coverage = build_coverage(self.registry, self.store, (witness,)) + self.assertIn("transition", coverage.exercised_snapshot_types) + self.assertIn(witness.question_identity, coverage.single_source_question_identities) + candidates = discover_question_candidates("source:v1", self.store, coverage) + self.assertTrue(candidates) + + +if __name__ == "__main__": + unittest.main()