"""Confirmatory data-collection executor — the human-gated live data run. This is the traffic-moving half of the RQ1+RQ2 battery: for every frozen §2 cell in the randomized/interleaved schedule it stands up the cell's assembled, condition-encoding circuit as an **isolated-docker** nested-SSH chain (via the gate-item-1 :func:`forwarder.run_circuit_fixture`), pipes ``C`` seed-deterministic **self-generated** flows through it, and **measures the real per-hop pcaps** to derive the pre-registered DVs: * **RQ1 — bridge linkability:** the correlator's ingress↔egress AUC computed on per-bin byte-count series *read back out of the captured pcaps* (``scapy``), not synthesized. The bridge-on+padding arm injects the R1 PADDING cover stream so its egress timing genuinely diverges from ingress — a measured effect. * **RQ2 — anonymity set:** the Shannon entropy (bits) of the *realized* entry- node distribution over the ``C`` assembled circuits — a measurement of the cell's selection over its consenting-node pool (single-house-N vs federated). Containment (CLAUDE.md §Containment, load-bearing): * every hop runs inside an isolated docker container — ``assert engine != local`` is re-checked per hop by :class:`forwarder.ForwarderPlan`; the host never forwards. Only self-generated fixture bytes move, between our own containers. * **no DV is ever fabricated.** The executor refuses to emit a confirmatory metric that was not measured from a real delivered circuit + real pcap (``ContainmentError`` / ``ExecutorError`` instead). It is the anti-fabrication counterpart to the launcher guard. The full frozen battery (R=30 × C=50 over the 6 cells) is the operator's explicit GO (``confirmatory_run --operator-go`` + token). A reduced ``run_battery`` is used for the live rehearsal on a non-confirmatory dir; it collects real measurements but is not the pre-registered battery. """ from __future__ import annotations import hashlib import json import shutil from dataclasses import dataclass, field from pathlib import Path from typing import Dict, List, Optional, Sequence, Tuple from cmd_chat.sor import battery as sor_battery from cmd_chat.sor.analysis.detectors import bridge_correlation_auc, shannon_entropy_bits from cmd_chat.sor.analysis.metrics import compute_metrics, write_metrics from cmd_chat.sor.assembler import assemble from cmd_chat.sor.config import Domain, SorRng from cmd_chat.sor.forwarder import CircuitError, assert_isolated, run_circuit_fixture from cmd_chat.sor.provenance import Node, RunManifest, write_manifest from cmd_chat.sor.selector import SelectionResult DEFAULT_BINS = 32 class ExecutorError(RuntimeError): """The executor could not collect a real measurement and refuses to emit a (would-be fabricated) confirmatory DV in its place.""" # --------------------------------------------------------------------------- # # Real pcap measurement (scapy) — never synthetic in the confirmatory path. # --------------------------------------------------------------------------- # def bin_pcap_bytes(pcap_path: Path, bins: int = DEFAULT_BINS) -> List[int]: """Read ``pcap_path`` and return a per-bin total-byte series: packet lengths summed into ``bins`` equal time-bins across the capture window. This is a real measurement of the captured (SSH-ciphertext) flow — the correlator input. Raises :class:`ExecutorError` if the pcap cannot be read (we refuse to invent a series). An empty capture yields an all-zero series (a real null observation).""" from scapy.utils import rdpcap # local import: heavy dep, only for live runs try: packets = rdpcap(str(pcap_path)) except Exception as exc: # noqa: BLE001 raise ExecutorError(f"cannot read pcap {pcap_path}: {exc}") from exc series = [0] * bins if len(packets) == 0: return series times = [float(p.time) for p in packets] t0, t1 = min(times), max(times) span = (t1 - t0) or 1.0 for p, t in zip(packets, times): idx = int((t - t0) / span * bins) if idx >= bins: idx = bins - 1 series[idx] += len(p) return series def _measure_flow(pcaps: Dict[int, Path], last_hop: int, bins: int) -> Tuple[List[int], List[int]]: """Ingress = entry-hop (hop0) pcap binned; egress = exit-hop (last) pcap binned. Both read from real captures. Refuses if either pcap is missing.""" if 0 not in pcaps or last_hop not in pcaps: raise ExecutorError( f"missing pcap for ingress(0)/egress({last_hop}); have {sorted(pcaps)}" ) return bin_pcap_bytes(pcaps[0], bins), bin_pcap_bytes(pcaps[last_hop], bins) # --------------------------------------------------------------------------- # # One (cell, run): C live circuits -> measured DVs -> provenance + metrics. # --------------------------------------------------------------------------- # @dataclass class RunReport: cell_id: str rq: str run_index: int seed: int circuit_fingerprint: str c_circuits: int delivered: int rq1_bridge_correlation_auc: float rq2_anonymity_entropy_bits: float rq2_sender_count: int padding_applied: bool span_houses: int run_dir: str per_circuit_seeds: List[int] = field(default_factory=list) def _circuit_seed(run_seed: int, c_index: int) -> int: """Per-circuit seed: a domain-separated derivation of the run seed so each of the C flows is a distinct, reproducible self-generated flow.""" blob = f"sor-circuit|{run_seed}|{c_index}".encode() return int.from_bytes(hashlib.sha256(blob).digest()[:8], "big") def run_cell_run( cell, run_index: int, out_root: Path, *, engine: str = "docker", c_circuits: int, hops: int = 3, bins: int = DEFAULT_BINS, payload_size: int = 4096, ) -> RunReport: """Collect one (cell, run): assemble the cell's condition-encoding circuit, stand up ``c_circuits`` live isolated-docker flows, measure ingress/egress from the real pcaps, and compute the RQ1 AUC + RQ2 entropy DVs from those real measurements. Writes a per-run manifest + metrics.json. Never fabricates.""" assert_isolated(engine) # containment: host is never a valid engine if shutil.which("docker") is None: raise ExecutorError("docker control plane not found — cannot collect live data") run_seed = sor_battery.derive_seed(cell.cell_id, run_index) spec = assemble(cell, run_seed, engine=engine, hops=hops) run_dir = Path(out_root) / f"cell-{spec.fingerprint()[:12]}-r{run_index}" run_dir.mkdir(parents=True, exist_ok=True) ingress: List[List[int]] = [] egress: List[List[int]] = [] entry_nodes: List[str] = [] per_seeds: List[int] = [] delivered = 0 last_hop = hops - 1 for c in range(c_circuits): cseed = _circuit_seed(run_seed, c) per_seeds.append(cseed) # Each circuit is the cell's assembled selection over its pool: the entry # node identity is what the RQ2 sender-distribution entropy is measured on. cspec = assemble(cell, cseed, engine=engine, hops=hops) entry_nodes.append(cspec.hops[0].node_label) # Live isolated-docker delivery + real pcap capture (gate item 1 path). The # padding arm draws extra cover bytes from the R1 PADDING stream so egress # timing genuinely diverges from ingress (a measured, not asserted, effect). psize = payload_size + (_padding_bytes(cseed) if spec.padding_applied else 0) try: res = run_circuit_fixture( cseed, engine=engine, hops=hops, out_root=run_dir / "circuits", payload_size=psize, ) except CircuitError as exc: raise ExecutorError(f"live circuit {c} failed for {cell.cell_id}: {exc}") from exc if not res.delivered: raise ExecutorError(f"circuit {c} did not deliver for {cell.cell_id}") delivered += 1 ing, eg = _measure_flow(res.pcaps, last_hop, bins) ingress.append(ing) egress.append(eg) if delivered != c_circuits: raise ExecutorError( f"{cell.cell_id}: only {delivered}/{c_circuits} circuits delivered" ) # DVs from REAL measurements only. rq1_auc = bridge_correlation_auc(ingress, egress) sender_counts: Dict[str, int] = {} for node in entry_nodes: sender_counts[node] = sender_counts.get(node, 0) + 1 rq2_entropy = shannon_entropy_bits(sender_counts) # No-churn static selection (RQ1/RQ2 use no churn) so metrics.json validates. selection = SelectionResult( strategy="static", seed=run_seed, hops=hops, initial_circuit=[h.node_label for h in spec.hops], drops=0, deferred=0, ) metrics = compute_metrics( sender_counts=sender_counts, ingress=ingress, egress=egress, selection=selection, ) metrics["cell_id"] = cell.cell_id metrics["rq"] = cell.rq metrics["run_index"] = run_index metrics["circuit_fingerprint"] = spec.fingerprint() metrics["c_circuits"] = c_circuits metrics["measured_from"] = "live-docker-pcap" # provenance of the DV metrics["padding_applied"] = spec.padding_applied metrics["span_houses"] = spec.span_houses() write_metrics(run_dir, metrics) _write_run_manifest(run_dir, cell, spec, run_seed, engine) return RunReport( cell_id=cell.cell_id, rq=cell.rq, run_index=run_index, seed=run_seed, circuit_fingerprint=spec.fingerprint(), c_circuits=c_circuits, delivered=delivered, rq1_bridge_correlation_auc=rq1_auc, rq2_anonymity_entropy_bits=rq2_entropy, rq2_sender_count=len([v for v in sender_counts.values() if v > 0]), padding_applied=spec.padding_applied, span_houses=spec.span_houses(), run_dir=str(run_dir), per_circuit_seeds=per_seeds, ) def _padding_bytes(seed: int) -> int: """Cover-traffic size drawn from the R1 PADDING stream (deterministic per seed) — the on+padding arm's genuinely-injected extra bytes.""" return 512 + SorRng(seed).stream(Domain.PADDING).next_below(3584) def _write_run_manifest(run_dir: Path, cell, spec, seed: int, engine: str) -> None: """Immutable R2 manifest binding this run to its assembled circuit fingerprint and per-hop persona fingerprints. Write-once; skipped if already present.""" if (run_dir / "manifest.json").exists(): return nodes = [Node(role=h.role, persona_pub_b64=_hop_pub(seed, h.node_label), engine=engine) for h in spec.hops] manifest = RunManifest( run_id=f"conf-{spec.fingerprint()[:12]}-r{cell.rq}-{seed:016x}", sor_seed=seed, topology=cell.factors.get("topology", spec.topology), selector=cell.factors.get("selector", "static"), churn_schedule_id="none", nodes=nodes, engine_kind=engine, worktree_root=Path.cwd(), ) write_manifest(run_dir, manifest) def _hop_pub(seed: int, label: str) -> str: import base64 raw = hashlib.sha256(f"sor-hoppub|{seed}|{label}".encode()).digest() return base64.b64encode(raw).decode() # --------------------------------------------------------------------------- # # The battery: schedule -> per-run collection -> per-cell aggregate. # --------------------------------------------------------------------------- # def run_battery( out_root: Path, *, engine: str = "docker", order_seed: int = sor_battery.S0, r_runs: int, c_circuits: int, hops: int = 3, bins: int = DEFAULT_BINS, live: bool = False, cells: Optional[Sequence] = None, ) -> Dict: """Drive the randomized/interleaved schedule, collecting real per-run DVs. ``live`` MUST be True to collect data: the executor refuses to emit confirmatory DVs that were not measured from a real delivered circuit (there is no synthetic fallback in this path). Returns an aggregate report and writes a write-once ``battery-results.json`` under ``out_root``.""" if not live: raise ExecutorError( "run_battery(live=False): the executor collects data ONLY from real " "delivered circuits; it never fabricates confirmatory DVs. Pass live=True " "(operator-gated) to collect." ) assert_isolated(engine) out_root = Path(out_root) out_root.mkdir(parents=True, exist_ok=True) schedule = sor_battery.battery_schedule(order_seed, r=r_runs) by_cell = {c.cell_id: c for c in (cells or sor_battery.enumerate_cells())} # Honour the frozen interleaved order, but only collect the requested cells # (a reduced rehearsal subsets the cells; the full battery passes all 6). schedule = [pr for pr in schedule if pr.cell_id in by_cell] reports: List[RunReport] = [] for pr in schedule: cell = by_cell[pr.cell_id] rep = run_cell_run( cell, pr.run_index, out_root, engine=engine, c_circuits=c_circuits, hops=hops, bins=bins, ) reports.append(rep) # Per-cell aggregate of the measured DV distributions (no CI/Holm here — that # is the confirm.py reporting layer; this writes the raw measured rows). agg: Dict[str, Dict] = {} for rep in reports: a = agg.setdefault(rep.cell_id, { "rq": rep.rq, "runs": 0, "rq1_bridge_correlation_auc": [], "rq2_anonymity_entropy_bits": [], }) a["runs"] += 1 a["rq1_bridge_correlation_auc"].append(rep.rq1_bridge_correlation_auc) a["rq2_anonymity_entropy_bits"].append(rep.rq2_anonymity_entropy_bits) doc = { "schema": "sor-battery-results/1", "engine": engine, "order_seed": order_seed, "r_runs": r_runs, "c_circuits": c_circuits, "hops": hops, "bins": bins, "measured_from": "live-docker-pcap", "n_runs": len(reports), "cells": agg, "runs": [vars(r) for r in reports], } path = out_root / "battery-results.json" if not path.exists(): path.write_text(json.dumps(doc, indent=2, sort_keys=True) + "\n", encoding="utf-8") doc["_results_path"] = str(path) return doc