From 5d0392ed59f9ef47bc21c634b8c53300012ef1e7 Mon Sep 17 00:00:00 2001 From: Leonardo Araujo Date: Sat, 26 Sep 2026 09:33:16 -0300 Subject: [PATCH 1/3] feat(eval): frozen M3 evaluator for the structural model (#209) The M3 protocol on #209 requires the evaluator to be committed and reviewed before the out-of-sample capture exists, so the measurement cannot be shaped by the data. This encodes it: - src/eval/structural_m3.py (pure): - three nested arms over the same inputs: M1 (sig), M2a (+edge), M2a+M2b (+util, the frozen model); - per-case attribution of which observable family witnessed the retained truth (sig / edge / util / topology); - healthy-window edge specificity at the unchanged cutoff; - series-identity coverage; - decide(), which applies the pre-registered thresholds (truth retained >= 0.50, median candidate fraction <= 0.33, healthy abstention >= 0.58; M2b validated / not_validated / untestable) and reports protocol deviations. - scripts/eval/m3_structural.py: - refuses a dirty tree and records the model commit + corpus sha256; - ingests once, runs the default path (ranker + rare_event) and the three arms on the persisted rows. - src/core/rca/structural.py: build_structural_view is split into load_structural_rows + resolve_structural, so the evaluator reads the same rows and runs the same resolve step as the product (behaviour-preserving; an integration test asserts parity). Co-Authored-By: Claude Opus 5.5 --- scripts/eval/m3_structural.py | 131 +++++++ src/core/rca/structural.py | 64 ++-- src/eval/structural_m3.py | 321 ++++++++++++++++++ .../integration/test_structural_m3_parity.py | 90 +++++ tests/unit/test_structural_m3.py | 194 +++++++++++ 5 files changed, 776 insertions(+), 24 deletions(-) create mode 100644 scripts/eval/m3_structural.py create mode 100644 src/eval/structural_m3.py create mode 100644 tests/integration/test_structural_m3_parity.py create mode 100644 tests/unit/test_structural_m3.py diff --git a/scripts/eval/m3_structural.py b/scripts/eval/m3_structural.py new file mode 100644 index 0000000..4fab832 --- /dev/null +++ b/scripts/eval/m3_structural.py @@ -0,0 +1,131 @@ +#!/usr/bin/env python3 +"""#209 M3 — run the frozen structural-model evaluation ONCE over a captured corpus. + +The protocol (pre-registered on #209 before capture) fixes the model, the metrics and the decision +thresholds; ``src/eval/structural_m3.py`` encodes them. This driver: + +1. refuses a dirty working tree (the frozen model must be an exact commit) unless --allow-dirty; +2. records provenance: model commit, corpus content hash, trigger mode, ranker/calibrator paths; +3. ingests every case once (logs + trace/metric sidecars, the harness ingest) and runs the default + explain path — learned ranker + rare_event trigger — i.e. ``raglogs eval``'s scoring; +4. reads each case's persisted rows over [baseline_start, window_end] exactly as the product's + structural view does and runs the three arms (M1, M2a, M2a+M2b); +5. writes the JSON report and the as-is markdown post for #209. + + python scripts/eval/m3_structural.py data/eval-cases/otel-m3 \\ + --json m3.json --md m3.md --wipe + +Needs Postgres (DB_URL). --wipe truncates the eval tables first (the local eval DB, never a shared one). +""" +from __future__ import annotations + +import argparse +import hashlib +import json +import os +import subprocess +import sys +from pathlib import Path + +_EVAL_TABLES = ("log_entries, ingestion_jobs, clusters, cluster_runs, cluster_members, cluster_embeddings, " + "log_embeddings, metric_samples, trace_spans, explanations, ingest_idempotency_keys") + + +def _git(*args: str) -> str: + return subprocess.run(["git", *args], capture_output=True, text=True, check=True).stdout.strip() + + +def corpus_hash(cases_dir: Path) -> str: + """sha256 over every file's relative path and bytes — identifies the exact corpus evaluated.""" + h = hashlib.sha256() + for path in sorted(p for p in cases_dir.rglob("*") if p.is_file()): + h.update(str(path.relative_to(cases_dir)).encode() + b"\0") + with open(path, "rb") as f: + for chunk in iter(lambda: f.read(1 << 20), b""): + h.update(chunk) + h.update(b"\0") + return h.hexdigest() + + +def main() -> int: + ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("cases", type=Path, help="captured case directory (one sub-dir per case)") + ap.add_argument("--json", type=Path, required=True, help="write the full report here") + ap.add_argument("--md", type=Path, required=True, help="write the #209 results post here") + ap.add_argument("--ranker", default="models/rca_ranker.json") + ap.add_argument("--calibrator", default="models/rca_calibrator.json") + ap.add_argument("--wipe", action="store_true", help="truncate the eval tables before ingesting") + ap.add_argument("--allow-dirty", action="store_true", help="run on an uncommitted tree (not a frozen run)") + args = ap.parse_args() + + dirty = bool(_git("status", "--porcelain")) + if dirty and not args.allow_dirty: + print("refusing: working tree is dirty — the frozen model must be an exact commit " + "(--allow-dirty for a non-canonical run)", file=sys.stderr) + return 2 + + # The default explain path as pre-registered: learned ranker + rare_event trigger. + os.environ["TRIGGER_MODE"] = "rare_event" + os.environ["RCA_RANKER_MODEL_PATH"] = args.ranker + os.environ["RCA_CALIBRATOR_MODEL_PATH"] = args.calibrator + from src.config import reload_settings + + settings = reload_settings() + + from sqlalchemy import text + + from src.core.rca.structural import load_structural_rows + from src.db.session import get_db + from src.eval.case import load_cases + from src.eval.report import build_report + from src.eval.runner import _scope_for, run_cases + from src.eval.structural_m3 import build_m3_report, evaluate_case, render_markdown + from src.utils.time import parse_duration + + cases = load_cases(args.cases) + if not cases: + print(f"no cases in {args.cases}", file=sys.stderr) + return 2 + provenance = { + "model_commit": _git("rev-parse", "HEAD") + (" (DIRTY — not a frozen run)" if dirty else ""), + "corpus": str(args.cases), + "corpus_sha256": corpus_hash(args.cases), + "n_cases": len(cases), + "trigger_mode": settings.trigger_mode, + "ranker": args.ranker, + "calibrator": args.calibrator, + } + + with get_db() as db: + if args.wipe: + db.execute(text(f"TRUNCATE {_EVAL_TABLES} RESTART IDENTITY CASCADE")) + db.commit() + results = run_cases(db, cases) # the one ingest + the default explain path + db.commit() + evals = [] + for case in cases: + baseline = parse_duration(case.baseline_window or settings.default_baseline_window) + metric_rows, span_rows = load_structural_rows(db, _scope_for(case), case.window_end, + case.window_start - baseline) + cause = case.root_cause.service if case.root_cause else None + evals.append(evaluate_case(case.id, cause, bool(case.expect_explanation and cause), span_rows, + metric_rows, case.window_start)) + + report = build_m3_report(evals) + default = build_report(results) + report["default_path"] = {k: default[k] for k in ("raglogs", "baseline", "lift_over_baseline")} + report["provenance"] = provenance + args.json.write_text(json.dumps(report, indent=1, sort_keys=True, default=str)) + + md = render_markdown(report, provenance=provenance) + r = default["raglogs"] + md += ("\n**Default path (learned ranker + rare_event):** " + + ", ".join(f"{k}={r[k]}" for k in sorted(r) if not isinstance(r[k], (dict, list))) + + f"\n\nLift over baseline: {default['lift_over_baseline']}\n") + args.md.write_text(md) + print(md) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/src/core/rca/structural.py b/src/core/rca/structural.py index 90ecbd1..893abe7 100644 --- a/src/core/rca/structural.py +++ b/src/core/rca/structural.py @@ -21,9 +21,10 @@ from sqlalchemy.orm import Session from src.core.rca.outcome import StructuralResult, resolve -from src.core.rca.partition import partition +from src.core.rca.partition import Partition, partition from src.core.rca.scoring import ScoredClass, rank_classes from src.core.rca.structural_model import ( + StructuralInputs, build_hypotheses, build_observables, structural_signals, @@ -31,20 +32,12 @@ from src.db.models import MetricSample, TraceSpan -def build_structural_view( - db: Session, - scope: str, - window_start: datetime, - window_end: datetime, - baseline_start: datetime, -) -> Optional[tuple[StructuralResult, tuple[ScoredClass, ...]]]: - """Build the deterministic structural result + Phase G ranking for ``scope`` over the incident - window ``[window_start, window_end]`` against the baseline from ``baseline_start``. Returns ``None`` - when the telemetry yields no candidate hypothesis (an honest abstention — no structural claim). - - The ``MetricSample`` / ``TraceSpan`` ORM rows duck-type directly into the shared structural-model - builders (``service`` / ``value`` / ``ts`` / ``metric`` and ``trace_id`` / ``span_id`` / - ``parent_span_id`` / ``service``), so the live view runs the exact same model as the shadow eval.""" +def load_structural_rows( + db: Session, scope: str, window_end: datetime, baseline_start: datetime +) -> tuple[list[MetricSample], list[TraceSpan]]: + """The ``metric_samples`` and ``trace_spans`` rows the structural model reads for ``scope`` over + ``[baseline_start, window_end]``. Shared by the live view and the M3 evaluator (#209), so both read + exactly the persisted rows — series identity included — that the product path sees.""" metric_rows = db.execute( select(MetricSample).where( MetricSample.scope == scope, @@ -62,17 +55,40 @@ def build_structural_view( TraceSpan.start_time <= window_end, ) ).scalars().all() + return list(metric_rows), list(span_rows) - # The one shared structural-model builder (#209 M1): span-derived sig (latency median + explicit - # OTLP status only) merged with the metric sig, and the INCIDENT call graph (edges whose child span - # starts at/after window_start — the baseline half of this load feeds latency ratios, never edges). - inp = structural_signals(span_rows, metric_rows, window_start) + +def resolve_structural(inp: StructuralInputs) -> Optional[tuple[StructuralResult, Partition]]: + """Hypotheses → observables → ``partition → resolve`` for one set of structural inputs, or ``None`` + when no candidate hypothesis exists (an honest abstention — no structural claim).""" hypotheses = build_hypotheses(inp.signals, inp.edges, inp.edge_signals, inp.util_signals) if not hypotheses: - return None # no candidate -> abstain, never a fabricated structural result - + return None observations = build_observables(inp.signals, inp.edge_signals, inp.util_signals) part = partition(hypotheses, observations) - result = resolve(part, observations=observations) - ranking = rank_classes(part) - return result, ranking + return resolve(part, observations=observations), part + + +def build_structural_view( + db: Session, + scope: str, + window_start: datetime, + window_end: datetime, + baseline_start: datetime, +) -> Optional[tuple[StructuralResult, tuple[ScoredClass, ...]]]: + """Build the deterministic structural result + Phase G ranking for ``scope`` over the incident + window ``[window_start, window_end]`` against the baseline from ``baseline_start``. Returns ``None`` + when the telemetry yields no candidate hypothesis (an honest abstention — no structural claim). + + The ``MetricSample`` / ``TraceSpan`` ORM rows duck-type directly into the shared structural-model + builders (``service`` / ``value`` / ``ts`` / ``metric`` and ``trace_id`` / ``span_id`` / + ``parent_span_id`` / ``service``), so the live view runs the exact same model as the shadow eval.""" + metric_rows, span_rows = load_structural_rows(db, scope, window_end, baseline_start) + # The one shared structural-model builder (#209 M1): span-derived sig (latency median + explicit + # OTLP status only) merged with the metric sig, and the INCIDENT call graph (edges whose child span + # starts at/after window_start — the baseline half of this load feeds latency ratios, never edges). + resolved = resolve_structural(structural_signals(span_rows, metric_rows, window_start)) + if resolved is None: + return None # no candidate -> abstain, never a fabricated structural result + result, part = resolved + return result, rank_classes(part) diff --git a/src/eval/structural_m3.py b/src/eval/structural_m3.py new file mode 100644 index 0000000..259d708 --- /dev/null +++ b/src/eval/structural_m3.py @@ -0,0 +1,321 @@ +"""#209 M3 — the frozen out-of-sample evaluator for the structural model. + +M3 is the generalization test for M1 (span ``sig``), M2a (``edge:`` observables) and M2b (``util:`` +observables), run **once** on a corpus captured after the model was frozen. The protocol — corpus, +metrics, decision thresholds — was pre-registered on #209 before any capture existed; this module +encodes it so the measurement cannot be shaped by the data: + +* Three nested arms over the *same* inputs: ``M1`` (sig only), ``M2a`` (+ edge), ``M2a+M2b`` + (+ util, the frozen model). Each arm is :func:`~src.core.rca.structural.resolve_structural` on a + trimmed :class:`~src.core.rca.structural_model.StructuralInputs` — the product's own resolve step. +* Inputs are the persisted rows the product path reads + (:func:`~src.core.rca.structural.load_structural_rows`), so series identity is exercised exactly + as it is in production, not on a hand-built sample list. +* The verdicts (:func:`decide`) apply the thresholds fixed in the protocol; none of them is a tunable. + +Pure: no DB. ``scripts/eval/m3_structural.py`` ingests a corpus and feeds rows in.""" +from __future__ import annotations + +import statistics +from collections import Counter +from dataclasses import dataclass, field +from datetime import datetime +from typing import Optional + +from src.core.rca.metric_series import instance_identity +from src.core.rca.observable import State +from src.core.rca.partition import discretize_rate +from src.core.rca.structural import resolve_structural +from src.core.rca.structural_model import ( + StructuralInputs, + service_universe, + structural_signals, + util_metric_class, +) + +ARMS = ("M1", "M2a", "M2a+M2b") +FULL_ARM = "M2a+M2b" + +# Pre-registered on #209 (M3 protocol), from the otel-fresh results; fixed before capture. +TRUTH_RETAINED_MIN = 0.50 # cause-has-telemetry truth retained (M2a's trace-only ceiling) +CANDIDATE_FRACTION_MAX = 0.33 # median |candidates| / |services| over positives that returned any +HEALTHY_ABSTENTION_MIN = 0.58 # healthy windows with no localization claim (7/12 on otel-fresh) +MIN_HEALTHY_NEGATIVES = 12 +# Positives whose cause is a resource fault — the cases M2b exists to recover. +RESOURCE_FAULT_SCENARIOS = ("adHighCpu", "recommendationCpuStress") + + +def arm_inputs(inp: StructuralInputs, arm: str) -> StructuralInputs: + """The inputs one arm sees: M1 drops edge + util observables, M2a drops util, the full arm keeps all. + Candidate generation from the incident call graph (``edges``) is M1 behaviour and stays in every arm.""" + if arm == "M1": + return inp._replace(edge_signals={}, util_signals={}) + if arm == "M2a": + return inp._replace(util_signals={}) + if arm == FULL_ARM: + return inp + raise ValueError(f"unknown arm {arm!r}") + + +@dataclass(frozen=True) +class ArmResult: + outcome: str # Outcome value, or "no_candidates" when no hypothesis was generated + localization: tuple[str, ...] + + @property + def claims(self) -> bool: + """Did this arm localize anything? An empty localization (no hypotheses, or none compatible) + is an abstention.""" + return bool(self.localization) + + +@dataclass(frozen=True) +class CaseEval: + case_id: str + cause: Optional[str] # ground-truth root-cause service; None for a healthy window + positive: bool + n_services: int + cause_has_telemetry: bool + arms: dict[str, ArmResult] + # PRESENT observable families *about the cause* under the full inputs: sig / edge (into it) / util + cause_families: tuple[str, ...] = () + util_measured_for_cause: bool = False + util_present: tuple[str, ...] = () + edges_measured: int = 0 + edges_present: int = 0 + edges_present_latency_only: int = 0 + util_samples: int = 0 # utilization-class metric samples in the rows + util_samples_named: int = 0 # ... of which name their reporting instance + notes: list[str] = field(default_factory=list) + + def retained(self, arm: str) -> bool: + return self.cause is not None and self.cause in self.arms[arm].localization + + @property + def retained_via(self) -> tuple[str, ...]: + """Families that witnessed the retained truth under the frozen model; ``topology`` when the + cause was retained with no PRESENT observable about it (a callee of an unexplained anomaly).""" + if not self.retained(FULL_ARM): + return () + return self.cause_families or ("topology",) + + +def scenario_of(case_id: str) -> str: + """``otel_adHighCpu`` / ``otel_chaos_recommendationCpuStress`` -> the scenario name.""" + return case_id.rsplit("_", 1)[-1] + + +def evaluate_case(case_id: str, cause: Optional[str], positive: bool, spans: list, metrics: list, + window_start: datetime) -> CaseEval: + """Run every arm on one case's rows (spans + metrics over ``[baseline_start, window_end]``).""" + inp = structural_signals(spans, metrics, window_start) + universe = service_universe(inp.signals, {sp.service for sp in spans if sp.service}) + arms: dict[str, ArmResult] = {} + for arm in ARMS: + resolved = resolve_structural(arm_inputs(inp, arm)) + if resolved is None: + arms[arm] = ArmResult("no_candidates", ()) + else: + result, _part = resolved + arms[arm] = ArmResult(result.outcome.value, tuple(result.localization)) + + families: list[str] = [] + if cause is not None: + sig = inp.signals.get(cause) + if sig is not None and sig.sig_state == State.PRESENT: + families.append("sig") + if any(e.callee == cause and e.sig_state == State.PRESENT for e in inp.edge_signals.values()): + families.append("edge") + if any(u.service == cause and u.sig_state == State.PRESENT for u in inp.util_signals.values()): + families.append("util") + + edges = list(inp.edge_signals.values()) + present_edges = [e for e in edges if e.sig_state == State.PRESENT] + latency_only = [e for e in present_edges + if not (e.error_measured and discretize_rate(e.error_rate) == State.PRESENT)] + util_rows = [m for m in metrics if util_metric_class(getattr(m, "metric", None)) is not None] + named = [m for m in util_rows if instance_identity(getattr(m, "attributes", None)) is not None] + return CaseEval( + case_id=case_id, cause=cause, positive=positive, + n_services=len(universe), cause_has_telemetry=cause is not None and cause in universe, + arms=arms, cause_families=tuple(families), + util_measured_for_cause=cause is not None and any( + u.service == cause and u.sig_state is not None for u in inp.util_signals.values()), + util_present=tuple(sorted(u.id for u in inp.util_signals.values() if u.sig_state == State.PRESENT)), + edges_measured=sum(e.sig_state is not None for e in edges), + edges_present=len(present_edges), + edges_present_latency_only=len(latency_only), + util_samples=len(util_rows), util_samples_named=len(named), + ) + + +def _rate(k: int, n: int) -> Optional[float]: + return k / n if n else None + + +def summarize_arm(cases: list[CaseEval], arm: str) -> dict: + """The pre-registered per-arm metrics.""" + pos = [c for c in cases if c.positive] + neg = [c for c in cases if not c.positive] + pos_tel = [c for c in pos if c.cause_has_telemetry] + claiming = [c for c in pos if c.arms[arm].claims] + fractions = [len(c.arms[arm].localization) / c.n_services for c in claiming if c.n_services] + identified = [c for c in cases if c.arms[arm].outcome == "identified"] + ident_correct = [c for c in identified if c.positive and c.arms[arm].localization == (c.cause,)] + return { + "positives": len(pos), + "generated": sum(c.arms[arm].outcome != "no_candidates" for c in pos), + "truth_retained_all": [sum(c.retained(arm) for c in pos), len(pos)], + "truth_retained_cause_has_telemetry": [sum(c.retained(arm) for c in pos_tel), len(pos_tel)], + "truth_retained_cause_has_telemetry_rate": _rate(sum(c.retained(arm) for c in pos_tel), len(pos_tel)), + "median_candidates": statistics.median([len(c.arms[arm].localization) for c in claiming]) + if claiming else None, + "median_candidate_fraction": statistics.median(fractions) if fractions else None, + "negatives": len(neg), + "healthy_abstained": [sum(not c.arms[arm].claims for c in neg), len(neg)], + "healthy_abstention_rate": _rate(sum(not c.arms[arm].claims for c in neg), len(neg)), + "outcomes": dict(sorted(Counter(c.arms[arm].outcome for c in cases).items())), + "identified_precision": [len(ident_correct), len(identified)], + "identified_on_healthy": sum(not c.positive for c in identified), + } + + +def edge_specificity(cases: list[CaseEval]) -> dict: + """Healthy-window flagged-edge rate at the unchanged M2a cutoff (reported, never recalibrated here).""" + neg = [c for c in cases if not c.positive] + measured = sum(c.edges_measured for c in neg) + present = sum(c.edges_present for c in neg) + return { + "healthy_windows_with_present_edge": [sum(c.edges_present > 0 for c in neg), len(neg)], + "healthy_present_over_measured_edges": [present, measured], + "healthy_present_edges_latency_only": [sum(c.edges_present_latency_only for c in neg), present], + } + + +def identity_coverage(cases: list[CaseEval]) -> dict: + """Share of utilization-class samples that name their reporting instance — an M2b prerequisite.""" + total = sum(c.util_samples for c in cases) + named = sum(c.util_samples_named for c in cases) + return {"util_samples": total, "named_instance": named, "rate": _rate(named, total)} + + +def decide(cases: list[CaseEval]) -> dict: + """Apply the pre-registered decision rules to the frozen model (the full arm).""" + full = summarize_arm(cases, FULL_ARM) + retained = full["truth_retained_cause_has_telemetry_rate"] + fraction = full["median_candidate_fraction"] + abstention = full["healthy_abstention_rate"] + criteria = { + "truth_retained_cause_has_telemetry": [retained, TRUTH_RETAINED_MIN, ">="], + "median_candidate_fraction": [fraction, CANDIDATE_FRACTION_MAX, "<="], + "healthy_abstention": [abstention, HEALTHY_ABSTENTION_MIN, ">="], + } + if retained is None or abstention is None: + generalizes = "untestable" # no positive with cause telemetry, or no healthy window + elif (retained >= TRUTH_RETAINED_MIN and abstention >= HEALTHY_ABSTENTION_MIN + and fraction is not None and fraction <= CANDIDATE_FRACTION_MAX): + generalizes = "generalizes" + else: + generalizes = "does_not_generalize" + + resource = [c for c in cases if c.positive and scenario_of(c.case_id) in RESOURCE_FAULT_SCENARIOS] + healthy_util = sorted(c.case_id for c in cases if not c.positive and c.util_present) + if healthy_util: + m2b = "not_validated" # util PRESENT on a healthy window + elif not any(c.util_measured_for_cause for c in resource): + m2b = "untestable" # no resource fault's cause has a measured util observable + elif any("util" in c.retained_via for c in resource): + m2b = "validated" + else: + m2b = "not_validated" + + neg = [c for c in cases if not c.positive] + deviations = [] + if len(neg) < MIN_HEALTHY_NEGATIVES: + deviations.append(f"{len(neg)} healthy negatives < {MIN_HEALTHY_NEGATIVES} pre-registered") + missing = [s for s in RESOURCE_FAULT_SCENARIOS if not any(scenario_of(c.case_id) == s for c in cases)] + if missing: + deviations.append(f"resource-fault scenarios absent from corpus: {missing}") + return { + "generalization": generalizes, + "generalization_criteria": criteria, + "m2b": m2b, + "m2b_resource_cases": {c.case_id: {"util_measured_for_cause": c.util_measured_for_cause, + "retained_via": list(c.retained_via)} for c in resource}, + "m2b_healthy_util_present": healthy_util, + "protocol_deviations": deviations, + } + + +def build_m3_report(cases: list[CaseEval]) -> dict: + return { + "arms": {arm: summarize_arm(cases, arm) for arm in ARMS}, + "edge_specificity": edge_specificity(cases), + "identity_coverage": identity_coverage(cases), + "verdicts": decide(cases), + "cases": [ + {"id": c.case_id, "cause": c.cause, "positive": c.positive, "n_services": c.n_services, + "cause_has_telemetry": c.cause_has_telemetry, + "arms": {a: {"outcome": r.outcome, "localization": list(r.localization), + "retained": c.retained(a)} for a, r in c.arms.items()}, + "retained_via": list(c.retained_via), "util_present": list(c.util_present)} + for c in cases + ], + } + + +def _frac(pair: list) -> str: + k, n = pair + return f"{k}/{n}" + (f" ({k / n:.0%})" if n else "") + + +def _num(v: Optional[float]) -> str: + return "n/a" if v is None else f"{v:.2f}" + + +def render_markdown(report: dict, *, provenance: dict) -> str: + """The as-is results post for #209.""" + lines = ["## M3 — frozen out-of-sample result", ""] + lines += [f"- {k}: `{v}`" for k, v in provenance.items()] + v = report["verdicts"] + lines += ["", f"**Generalization:** `{v['generalization']}` · **M2b:** `{v['m2b']}`", ""] + for name, (value, bound, op) in v["generalization_criteria"].items(): + lines.append(f"- {name}: {_num(value)} (pre-registered {op} {bound})") + if v["protocol_deviations"]: + lines += ["", "**Protocol deviations:** " + "; ".join(v["protocol_deviations"])] + lines += ["", "| metric | " + " | ".join(ARMS) + " |", "|---|" + "---|" * len(ARMS)] + rows = [ + ("generated / positives", lambda a: f"{a['generated']}/{a['positives']}"), + ("truth retained (all)", lambda a: _frac(a["truth_retained_all"])), + ("truth retained (cause has telemetry)", lambda a: _frac(a["truth_retained_cause_has_telemetry"])), + ("median candidates", lambda a: str(a["median_candidates"])), + ("median candidate fraction", lambda a: _num(a["median_candidate_fraction"])), + ("healthy abstention", lambda a: _frac(a["healthy_abstained"])), + ("IDENTIFIED precision", lambda a: f"{a['identified_precision'][0]}/{a['identified_precision'][1]}"), + ("IDENTIFIED on healthy", lambda a: str(a["identified_on_healthy"])), + ("outcomes", lambda a: ", ".join(f"{k}={n}" for k, n in a["outcomes"].items())), + ] + for label, fn in rows: + lines.append(f"| {label} | " + " | ".join(fn(report["arms"][arm]) for arm in ARMS) + " |") + es, ic = report["edge_specificity"], report["identity_coverage"] + lines += [ + "", + f"- Edge specificity (healthy): windows with a PRESENT edge " + f"{_frac(es['healthy_windows_with_present_edge'])}; PRESENT/measured edges " + f"{_frac(es['healthy_present_over_measured_edges'])}; latency-only " + f"{_frac(es['healthy_present_edges_latency_only'])}", + f"- Series-identity coverage: {ic['named_instance']}/{ic['util_samples']} utilization samples name " + f"their instance ({_num(ic['rate'])})", + "", + "| case | cause | " + " | ".join(ARMS) + " | retained via | util PRESENT |", + "|---|---|" + "---|" * len(ARMS) + "---|---|", + ] + for c in report["cases"]: + cells = [] + for arm in ARMS: + r = c["arms"][arm] + mark = ("✓ " if r["retained"] else "✗ ") if c["positive"] else "" + cells.append(f"{mark}{r['outcome']} ({len(r['localization'])})") + lines.append(f"| {c['id']} | {c['cause'] or '—'} | " + " | ".join(cells) + + f" | {', '.join(c['retained_via']) or '—'} | {', '.join(c['util_present']) or '—'} |") + return "\n".join(lines) + "\n" diff --git a/tests/integration/test_structural_m3_parity.py b/tests/integration/test_structural_m3_parity.py new file mode 100644 index 0000000..c19e6b1 --- /dev/null +++ b/tests/integration/test_structural_m3_parity.py @@ -0,0 +1,90 @@ +"""Integration: the #209 M3 evaluator's frozen arm IS the product's structural view. On the same persisted +rows, ``evaluate_case``'s M2a+M2b arm must return exactly ``build_structural_view``'s outcome and +localization — the evaluator measures the shipped model, not a copy. Skipped without Postgres.""" +from __future__ import annotations + +import os +from datetime import datetime, timedelta, timezone + +import pytest + +pytestmark = pytest.mark.skipif( + not os.getenv("DB_URL") and not os.getenv("INTEGRATION_TESTS"), + reason="Integration tests require DB_URL environment variable", +) + +W = datetime(2026, 1, 1, 12, 0, 0, tzinfo=timezone.utc) +BASE = W - timedelta(seconds=300) +END = W + timedelta(seconds=300) +SCOPE = "eval:m3-parity" + + +@pytest.fixture +def db_session(): + from sqlalchemy import text + + from src.db.models import Base + from src.db.session import check_connection, get_db, get_engine + + if not check_connection(): + pytest.skip("Cannot connect to database") + engine = get_engine() + with engine.connect() as conn: + conn.execute(text("CREATE EXTENSION IF NOT EXISTS vector")) + conn.commit() + Base.metadata.drop_all(engine) + Base.metadata.create_all(engine) + with get_db() as db: + yield db + + +def _spans(): + """frontend -> payment (Charge) and frontend -> ad (GetAds); in the incident, Charge fails with no + payment span (unreachable callee).""" + from src.core.ingestion.telemetry import ParsedSpan + + out, n = [], 0 + for i in range(1, 11): + for ts, incident in ((W - timedelta(seconds=i), False), (W + timedelta(seconds=i), True)): + for op, callee in (("oteldemo.PaymentService/Charge", "payment"), ("oteldemo.AdService/GetAds", "ad")): + n += 1 + failing = incident and callee == "payment" + out.append(ParsedSpan(trace_id=f"t{n}", span_id=f"p{n}", service="frontend", operation=op, + start_time=ts, duration_ms=10.0, status_code="2" if failing else "1")) + if not failing: + out.append(ParsedSpan(trace_id=f"t{n}", span_id=f"c{n}", parent_span_id=f"p{n}", + service=callee, operation="handle", + start_time=ts + timedelta(milliseconds=1), duration_ms=9.0, + status_code="1")) + return out + + +def _cpu(): + from src.core.ingestion.telemetry import ParsedMetricSample + + inst = {"service.instance.id": "ad-1"} + return [ParsedMetricSample(service="ad", metric="jvm.cpu.recent_utilization", value=v, + ts=W + timedelta(seconds=sign * 30 * i), metric_type="gauge", attributes=inst) + for i in range(1, 6) for sign, v in ((-1, 0.01), (1, 0.9))] + + +def test_full_arm_matches_the_product_structural_view(db_session): + from src.core.ingestion.telemetry import persist_metric_samples, persist_spans + from src.core.rca.structural import build_structural_view, load_structural_rows + from src.eval.structural_m3 import FULL_ARM, evaluate_case + + persist_spans(db_session, _spans(), scope=SCOPE) + persist_metric_samples(db_session, _cpu(), scope=SCOPE) + db_session.flush() + + built = build_structural_view(db_session, SCOPE, W, END, BASE) + assert built is not None + product, _ranking = built + + metric_rows, span_rows = load_structural_rows(db_session, SCOPE, END, BASE) + c = evaluate_case("otel_parity", "payment", True, span_rows, metric_rows, W) + assert c.arms[FULL_ARM].outcome == product.outcome.value + assert c.arms[FULL_ARM].localization == tuple(product.localization) + # and the evaluator attributes the families it saw: the edge into payment, util on ad + assert "edge" in c.retained_via + assert c.util_present == ("util:ad:cpu",) diff --git a/tests/unit/test_structural_m3.py b/tests/unit/test_structural_m3.py new file mode 100644 index 0000000..223ffe1 --- /dev/null +++ b/tests/unit/test_structural_m3.py @@ -0,0 +1,194 @@ +"""#209 M3 evaluator — arms, per-case attribution and the pre-registered decision rules. Pure, no DB. + +The arms are nested views of one set of inputs (M1 ⊂ M2a ⊂ M2a+M2b); each family must recover exactly +the fault it exists for, and the verdicts must apply the thresholds frozen in the M3 protocol.""" +from datetime import datetime, timedelta, timezone + +import pytest + +from src.eval.structural_m3 import ( + ARMS, + CANDIDATE_FRACTION_MAX, + FULL_ARM, + HEALTHY_ABSTENTION_MIN, + MIN_HEALTHY_NEGATIVES, + ArmResult, + CaseEval, + arm_inputs, + build_m3_report, + decide, + evaluate_case, + render_markdown, + scenario_of, +) + +_W = datetime(2026, 1, 1, 12, 0, 0, tzinfo=timezone.utc) +OK, ERR = "1", "2" +CHARGE = "oteldemo.PaymentService/Charge" +_ONE = {"service.instance.id": "pod-1"} + + +class _Span: + def __init__(self, service, op, ts, *, trace, sid, parent=None, dur=10.0, status=OK): + self.service, self.operation, self.start_time = service, op, ts + self.trace_id, self.span_id, self.parent_span_id = trace, sid, parent + self.duration_ms, self.status_code = dur, status + + +class _M: + def __init__(self, service, metric, value, ts, *, attributes=_ONE, metric_type="gauge"): + self.service, self.metric, self.value, self.ts = service, metric, value, ts + self.attributes, self.metric_type = attributes, metric_type + + +_n = 0 + + +def _call(ts, *, status=OK, child=True, callee="payment", op=CHARGE): + global _n + _n += 1 + out = [_Span("frontend", op, ts, trace=f"t{_n}", sid=f"p{_n}", status=status)] + if child: + out.append(_Span(callee, "handle", ts + timedelta(milliseconds=1), trace=f"t{_n}", sid=f"c{_n}", + parent=f"p{_n}", dur=9.0, status=OK)) + return out + + +def _traffic(n=10, *, incident_status=OK, incident_child=True, callee="payment", op=CHARGE): + spans = [] + for i in range(n): + spans += _call(_W - timedelta(seconds=i + 1), callee=callee, op=op) + spans += _call(_W + timedelta(seconds=i + 1), status=incident_status, child=incident_child, + callee=callee, op=op) + return spans + + +def _gauge(service, metric, base, inc, n=5, attributes=_ONE): + out = [_M(service, metric, base, _W - timedelta(seconds=10 * (i + 1)), attributes=attributes) + for i in range(n)] + return out + [_M(service, metric, inc, _W + timedelta(seconds=10 * (i + 1)), attributes=attributes) + for i in range(n)] + + +class TestArms: + def test_unknown_arm_is_rejected(self): + with pytest.raises(ValueError): + arm_inputs(None, "M4") # type: ignore[arg-type] + + def test_unreachable_callee_is_recovered_by_the_edge_family_only(self): + # frontend's Charge calls fail and payment emits nothing: only edge:frontend->payment sees it + spans = _traffic(incident_status=ERR, incident_child=False) + c = evaluate_case("otel_paymentUnreachable", "payment", True, spans, [], _W) + assert not c.retained("M1") + assert c.retained("M2a") and c.retained(FULL_ARM) + assert "edge" in c.retained_via + + def test_locally_silent_cpu_fault_is_recovered_by_the_util_family_only(self): + spans = _traffic() + _traffic(callee="ad", op="oteldemo.AdService/GetAds") + metrics = _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9) + c = evaluate_case("otel_adHighCpu", "ad", True, spans, metrics, _W) + assert not c.retained("M1") and not c.retained("M2a") + assert c.retained(FULL_ARM) and c.retained_via == ("util",) + assert c.util_measured_for_cause and c.util_present == ("util:ad:cpu",) + assert (c.util_samples, c.util_samples_named) == (10, 10) + + def test_unnamed_utilization_counts_against_identity_coverage(self): + metrics = _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9, attributes={}) + c = evaluate_case("otel_adHighCpu", "ad", True, _traffic(), metrics, _W) + assert (c.util_samples, c.util_samples_named) == (10, 0) + assert not c.util_measured_for_cause # unnamed -> UNKNOWN, never measured + + def test_healthy_window_abstains_in_every_arm(self): + c = evaluate_case("otel_healthy_1", None, False, _traffic(), + _gauge("ad", "jvm.cpu.recent_utilization", 0.3, 0.31), _W) + assert all(not c.arms[a].claims for a in ARMS) + assert c.edges_measured == 1 and c.edges_present == 0 and c.util_present == () + + +def _case(cid, cause, *, retained=True, n_loc=1, n_services=10, families=("sig",), util_measured=False, + util_present=(), outcome="uncertain"): + loc = ((cause,) + tuple(f"x{i}" for i in range(n_loc - 1))) if (cause and retained) else \ + tuple(f"x{i}" for i in range(n_loc)) + arm = ArmResult(outcome if loc else "no_candidates", loc) + return CaseEval(cid, cause, cause is not None, n_services, cause is not None, {a: arm for a in ARMS}, + cause_families=tuple(families) if cause else (), util_measured_for_cause=util_measured, + util_present=tuple(util_present)) + + +def _healthy(i, *, claims=False, util_present=()): + return _case(f"otel_healthy_{i}", None, n_loc=1 if claims else 0, util_present=util_present) + + +def _corpus(*, retained=(True, True), n_loc=2, healthy_claims=0, n_healthy=MIN_HEALTHY_NEGATIVES, extra=()): + pos = [_case("otel_a", "a", retained=retained[0], n_loc=n_loc), + _case("otel_b", "b", retained=retained[1], n_loc=n_loc)] + neg = [_healthy(i, claims=i < healthy_claims) for i in range(n_healthy)] + return pos + neg + list(extra) + + +class TestDecisionRules: + def test_passes_all_three_pre_registered_bars(self): + assert decide(_corpus())["generalization"] == "generalizes" + + def test_truth_retained_below_half_fails(self): + cases = _corpus(retained=(True, False)) + [_case("otel_c", "c", retained=False)] + assert decide(cases)["generalization"] == "does_not_generalize" + + def test_candidate_fraction_above_bound_fails(self): + n_loc = int(CANDIDATE_FRACTION_MAX * 10) + 1 # 4/10 > 0.33 + assert decide(_corpus(n_loc=n_loc))["generalization"] == "does_not_generalize" + + def test_healthy_abstention_below_bound_fails(self): + claims = MIN_HEALTHY_NEGATIVES - int(HEALTHY_ABSTENTION_MIN * MIN_HEALTHY_NEGATIVES) + 1 + assert decide(_corpus(healthy_claims=claims))["generalization"] == "does_not_generalize" + + def test_abstention_at_the_otel_fresh_value_passes(self): + # 7/12 = 0.583 is the otel-fresh value the bar was set from; it must pass (>= 0.58) + assert decide(_corpus(healthy_claims=5))["generalization"] == "generalizes" + + def test_no_healthy_windows_is_untestable_and_a_deviation(self): + v = decide(_corpus(n_healthy=0)) + assert v["generalization"] == "untestable" + assert any("healthy negatives" in d for d in v["protocol_deviations"]) + + def test_m2b_validated_by_a_resource_fault_retained_via_util(self): + extra = [_case("otel_adHighCpu", "ad", families=("util",), util_measured=True, + util_present=("util:ad:cpu",))] + assert decide(_corpus(extra=extra))["m2b"] == "validated" + + def test_m2b_untestable_when_no_resource_fault_is_measured(self): + extra = [_case("otel_adHighCpu", "ad", retained=False, families=(), util_measured=False)] + assert decide(_corpus(extra=extra))["m2b"] == "untestable" + + def test_m2b_not_validated_when_measured_but_not_retained_via_util(self): + extra = [_case("otel_chaos_recommendationCpuStress", "recommendation", families=("sig",), + util_measured=True)] + assert decide(_corpus(extra=extra))["m2b"] == "not_validated" + + def test_util_present_on_a_healthy_window_blocks_validation(self): + extra = [_case("otel_adHighCpu", "ad", families=("util",), util_measured=True), + _healthy(99, util_present=("util:cart:cpu",))] + v = decide(_corpus(extra=extra)) + assert v["m2b"] == "not_validated" and v["m2b_healthy_util_present"] == ["otel_healthy_99"] + + def test_missing_resource_scenarios_are_reported_as_deviations(self): + v = decide(_corpus()) + assert v["m2b"] == "untestable" + assert any("resource-fault scenarios absent" in d for d in v["protocol_deviations"]) + + def test_retained_without_a_present_family_is_topology(self): + assert _case("otel_a", "a", families=()).retained_via == ("topology",) + assert _case("otel_a", "a", retained=False).retained_via == () + + +def test_scenario_names(): + assert scenario_of("otel_adHighCpu") == "adHighCpu" + assert scenario_of("otel_chaos_recommendationCpuStress") == "recommendationCpuStress" + + +def test_report_renders_every_arm_and_case(): + report = build_m3_report(_corpus()) + md = render_markdown(report, provenance={"model_commit": "abc123"}) + assert "`abc123`" in md and "**Generalization:** `generalizes`" in md + assert all(arm in md for arm in ARMS) + assert all(c["id"] in md for c in report["cases"]) From 3af69ad3347148336bde3926ac5fc0fbc96381ca Mon Sep 17 00:00:00 2001 From: Leonardo Araujo Date: Sat, 26 Sep 2026 10:21:49 -0300 Subject: [PATCH 2/3] fix(eval): decide() is the frozen M3 protocol, not a paraphrase (#214 review) - M2b "retained via util" is now the arm delta: the frozen arm retains the cause, M2a (the same inputs without util) does not, and util on the cause is PRESENT. A co-present util on a fault that sig already retains no longer validates. - M2b "untestable" now reads series-identity coverage, as the protocol says: exactly no named utilization sample (otel-fresh was 0/37282). High coverage with unmeasured resource causes is not_validated. - The >=2x edge rate is counted on its own, including edges whose error branch also fired. The error branch and the combined PRESENT are reported separately. - The results post labels unobservable cases (no telemetry from the cause). - The gated statistics are now the otel-fresh reference definitions: - the universe is span OR metric services; - generated means any hypothesis; - selectivity is over generated positives, with an empty set counting as 0; - abstention means no hypothesis. "No localization claim" is reported as a secondary figure. Rerun on the spent otel-fresh corpus, the evaluator reproduces every reference number exactly. - The candidate-fraction bound stays the frozen literal 0.33, pinned by a test that exact 1/3 (otel-fresh M1's value) fails it. Co-Authored-By: Claude Opus 5.5 --- src/eval/structural_m3.py | 124 +++++++++++++-------- tests/unit/test_structural_m3.py | 181 ++++++++++++++++++++++++------- 2 files changed, 222 insertions(+), 83 deletions(-) diff --git a/src/eval/structural_m3.py b/src/eval/structural_m3.py index 259d708..5fa0681 100644 --- a/src/eval/structural_m3.py +++ b/src/eval/structural_m3.py @@ -12,6 +12,10 @@ (:func:`~src.core.rca.structural.load_structural_rows`), so series identity is exercised exactly as it is in production, not on a hand-built sample list. * The verdicts (:func:`decide`) apply the thresholds fixed in the protocol; none of them is a tunable. +* Each gated statistic is computed exactly as the otel-fresh reference numbers the thresholds came + from: the service universe is every service seen in span *or* metric rows; a positive is + *generated* when any hypothesis exists; selectivity is over generated positives (an empty + compatible set counts as 0); healthy abstention is *no hypothesis generated*. Pure: no DB. ``scripts/eval/m3_structural.py`` ingests a corpus and feeds rows in.""" from __future__ import annotations @@ -24,11 +28,10 @@ from src.core.rca.metric_series import instance_identity from src.core.rca.observable import State -from src.core.rca.partition import discretize_rate +from src.core.rca.partition import discretize_rate, discretize_ratio from src.core.rca.structural import resolve_structural from src.core.rca.structural_model import ( StructuralInputs, - service_universe, structural_signals, util_metric_class, ) @@ -63,10 +66,9 @@ class ArmResult: localization: tuple[str, ...] @property - def claims(self) -> bool: - """Did this arm localize anything? An empty localization (no hypotheses, or none compatible) - is an abstention.""" - return bool(self.localization) + def generated(self) -> bool: + """Did this arm generate any hypothesis? Not generating is the (reference) abstention.""" + return self.outcome != "no_candidates" @dataclass(frozen=True) @@ -81,9 +83,11 @@ class CaseEval: cause_families: tuple[str, ...] = () util_measured_for_cause: bool = False util_present: tuple[str, ...] = () - edges_measured: int = 0 - edges_present: int = 0 - edges_present_latency_only: int = 0 + edges_measured: int = 0 # edges whose sig is PRESENT or ABSENT + edges_present: int = 0 # ... PRESENT (error >= cutoff OR latency >= 2x) + edges_latency_measured: int = 0 # edges with a measured latency branch + edges_latency_high: int = 0 # ... at >= 2x, whether or not the error branch also fired + edges_error_present: int = 0 # edges whose explicit-status error rate crossed the cutoff util_samples: int = 0 # utilization-class metric samples in the rows util_samples_named: int = 0 # ... of which name their reporting instance notes: list[str] = field(default_factory=list) @@ -93,12 +97,21 @@ def retained(self, arm: str) -> bool: @property def retained_via(self) -> tuple[str, ...]: - """Families that witnessed the retained truth under the frozen model; ``topology`` when the - cause was retained with no PRESENT observable about it (a callee of an unexplained anomaly).""" + """PRESENT families about the retained truth under the frozen model — what *witnessed* it, not + what retained it (see :attr:`retained_by_util`); ``topology`` when no PRESENT observable concerns + the cause (a callee of an unexplained anomaly).""" if not self.retained(FULL_ARM): return () return self.cause_families or ("topology",) + @property + def retained_by_util(self) -> bool: + """Retained *via* ``util:`` — the arm delta: the frozen model retains the cause, the same inputs + without util (M2a) do not, and a ``util:`` observable on the cause is PRESENT. A measured util + observable is already a named-instance series (``summarize_utilization`` will not emit + PRESENT/ABSENT without one).""" + return self.retained(FULL_ARM) and not self.retained("M2a") and "util" in self.cause_families + def scenario_of(case_id: str) -> str: """``otel_adHighCpu`` / ``otel_chaos_recommendationCpuStress`` -> the scenario name.""" @@ -109,7 +122,9 @@ def evaluate_case(case_id: str, cause: Optional[str], positive: bool, spans: lis window_start: datetime) -> CaseEval: """Run every arm on one case's rows (spans + metrics over ``[baseline_start, window_end]``).""" inp = structural_signals(spans, metrics, window_start) - universe = service_universe(inp.signals, {sp.service for sp in spans if sp.service}) + # every service seen in span OR metric rows (the reference denominator) + universe = ({sp.service for sp in spans if sp.service} + | {m.service for m in metrics if getattr(m, "service", None)}) arms: dict[str, ArmResult] = {} for arm in ARMS: resolved = resolve_structural(arm_inputs(inp, arm)) @@ -130,9 +145,8 @@ def evaluate_case(case_id: str, cause: Optional[str], positive: bool, spans: lis families.append("util") edges = list(inp.edge_signals.values()) - present_edges = [e for e in edges if e.sig_state == State.PRESENT] - latency_only = [e for e in present_edges - if not (e.error_measured and discretize_rate(e.error_rate) == State.PRESENT)] + latency_high = [e for e in edges if e.latency_measured and discretize_ratio(e.latency_ratio) == State.HIGH] + error_present = [e for e in edges if e.error_measured and discretize_rate(e.error_rate) == State.PRESENT] util_rows = [m for m in metrics if util_metric_class(getattr(m, "metric", None)) is not None] named = [m for m in util_rows if instance_identity(getattr(m, "attributes", None)) is not None] return CaseEval( @@ -143,8 +157,10 @@ def evaluate_case(case_id: str, cause: Optional[str], positive: bool, spans: lis u.service == cause and u.sig_state is not None for u in inp.util_signals.values()), util_present=tuple(sorted(u.id for u in inp.util_signals.values() if u.sig_state == State.PRESENT)), edges_measured=sum(e.sig_state is not None for e in edges), - edges_present=len(present_edges), - edges_present_latency_only=len(latency_only), + edges_present=sum(e.sig_state == State.PRESENT for e in edges), + edges_latency_measured=sum(e.latency_measured for e in edges), + edges_latency_high=len(latency_high), + edges_error_present=len(error_present), util_samples=len(util_rows), util_samples_named=len(named), ) @@ -158,22 +174,24 @@ def summarize_arm(cases: list[CaseEval], arm: str) -> dict: pos = [c for c in cases if c.positive] neg = [c for c in cases if not c.positive] pos_tel = [c for c in pos if c.cause_has_telemetry] - claiming = [c for c in pos if c.arms[arm].claims] - fractions = [len(c.arms[arm].localization) / c.n_services for c in claiming if c.n_services] + generated = [c for c in pos if c.arms[arm].generated] + fractions = [len(c.arms[arm].localization) / max(c.n_services, 1) for c in generated] identified = [c for c in cases if c.arms[arm].outcome == "identified"] ident_correct = [c for c in identified if c.positive and c.arms[arm].localization == (c.cause,)] return { "positives": len(pos), - "generated": sum(c.arms[arm].outcome != "no_candidates" for c in pos), + "generated": len(generated), "truth_retained_all": [sum(c.retained(arm) for c in pos), len(pos)], "truth_retained_cause_has_telemetry": [sum(c.retained(arm) for c in pos_tel), len(pos_tel)], "truth_retained_cause_has_telemetry_rate": _rate(sum(c.retained(arm) for c in pos_tel), len(pos_tel)), - "median_candidates": statistics.median([len(c.arms[arm].localization) for c in claiming]) - if claiming else None, + "median_candidates": statistics.median([len(c.arms[arm].localization) for c in generated]) + if generated else None, "median_candidate_fraction": statistics.median(fractions) if fractions else None, "negatives": len(neg), - "healthy_abstained": [sum(not c.arms[arm].claims for c in neg), len(neg)], - "healthy_abstention_rate": _rate(sum(not c.arms[arm].claims for c in neg), len(neg)), + "healthy_abstained": [sum(not c.arms[arm].generated for c in neg), len(neg)], + "healthy_abstention_rate": _rate(sum(not c.arms[arm].generated for c in neg), len(neg)), + # secondary, not gated: healthy windows with no localization claim at all (incl. NO_COMPATIBLE) + "healthy_no_claim": [sum(not c.arms[arm].localization for c in neg), len(neg)], "outcomes": dict(sorted(Counter(c.arms[arm].outcome for c in cases).items())), "identified_precision": [len(ident_correct), len(identified)], "identified_on_healthy": sum(not c.positive for c in identified), @@ -181,14 +199,18 @@ def summarize_arm(cases: list[CaseEval], arm: str) -> dict: def edge_specificity(cases: list[CaseEval]) -> dict: - """Healthy-window flagged-edge rate at the unchanged M2a cutoff (reported, never recalibrated here).""" + """Healthy-window edge flags at the unchanged M2a cutoffs (reported, never recalibrated here). The + >=2x latency rate is counted on its own — including edges whose error branch also fired — so the + rate the protocol asks for is exact; the error branch and the combined PRESENT are reported too.""" neg = [c for c in cases if not c.positive] - measured = sum(c.edges_measured for c in neg) - present = sum(c.edges_present for c in neg) return { + "healthy_windows_with_latency_high_edge": [sum(c.edges_latency_high > 0 for c in neg), len(neg)], + "healthy_latency_high_over_latency_measured_edges": [sum(c.edges_latency_high for c in neg), + sum(c.edges_latency_measured for c in neg)], + "healthy_windows_with_error_present_edge": [sum(c.edges_error_present > 0 for c in neg), len(neg)], "healthy_windows_with_present_edge": [sum(c.edges_present > 0 for c in neg), len(neg)], - "healthy_present_over_measured_edges": [present, measured], - "healthy_present_edges_latency_only": [sum(c.edges_present_latency_only for c in neg), present], + "healthy_present_over_measured_edges": [sum(c.edges_present for c in neg), + sum(c.edges_measured for c in neg)], } @@ -218,13 +240,17 @@ def decide(cases: list[CaseEval]) -> dict: else: generalizes = "does_not_generalize" + # M2b, as frozen: not validated if util fires on a healthy window; untestable only when series- + # identity coverage is ~0 (exactly: no named utilization sample — otel-fresh was 0/37282); + # validated iff a resource-fault positive is retained via util (the arm delta); else not validated. resource = [c for c in cases if c.positive and scenario_of(c.case_id) in RESOURCE_FAULT_SCENARIOS] healthy_util = sorted(c.case_id for c in cases if not c.positive and c.util_present) + coverage = identity_coverage(cases)["rate"] if healthy_util: - m2b = "not_validated" # util PRESENT on a healthy window - elif not any(c.util_measured_for_cause for c in resource): - m2b = "untestable" # no resource fault's cause has a measured util observable - elif any("util" in c.retained_via for c in resource): + m2b = "not_validated" + elif coverage is None or coverage == 0.0: + m2b = "untestable" + elif any(c.retained_by_util for c in resource): m2b = "validated" else: m2b = "not_validated" @@ -240,7 +266,11 @@ def decide(cases: list[CaseEval]) -> dict: "generalization": generalizes, "generalization_criteria": criteria, "m2b": m2b, + "m2b_identity_coverage": coverage, "m2b_resource_cases": {c.case_id: {"util_measured_for_cause": c.util_measured_for_cause, + "retained_M2a": c.retained("M2a"), + "retained_full": c.retained(FULL_ARM), + "retained_by_util": c.retained_by_util, "retained_via": list(c.retained_via)} for c in resource}, "m2b_healthy_util_present": healthy_util, "protocol_deviations": deviations, @@ -258,7 +288,8 @@ def build_m3_report(cases: list[CaseEval]) -> dict: "cause_has_telemetry": c.cause_has_telemetry, "arms": {a: {"outcome": r.outcome, "localization": list(r.localization), "retained": c.retained(a)} for a, r in c.arms.items()}, - "retained_via": list(c.retained_via), "util_present": list(c.util_present)} + "retained_via": list(c.retained_via), "retained_by_util": c.retained_by_util, + "util_present": list(c.util_present)} for c in cases ], } @@ -290,7 +321,8 @@ def render_markdown(report: dict, *, provenance: dict) -> str: ("truth retained (cause has telemetry)", lambda a: _frac(a["truth_retained_cause_has_telemetry"])), ("median candidates", lambda a: str(a["median_candidates"])), ("median candidate fraction", lambda a: _num(a["median_candidate_fraction"])), - ("healthy abstention", lambda a: _frac(a["healthy_abstained"])), + ("healthy abstention (no hypothesis)", lambda a: _frac(a["healthy_abstained"])), + ("healthy no localization claim", lambda a: _frac(a["healthy_no_claim"])), ("IDENTIFIED precision", lambda a: f"{a['identified_precision'][0]}/{a['identified_precision'][1]}"), ("IDENTIFIED on healthy", lambda a: str(a["identified_on_healthy"])), ("outcomes", lambda a: ", ".join(f"{k}={n}" for k, n in a["outcomes"].items())), @@ -300,15 +332,19 @@ def render_markdown(report: dict, *, provenance: dict) -> str: es, ic = report["edge_specificity"], report["identity_coverage"] lines += [ "", - f"- Edge specificity (healthy): windows with a PRESENT edge " - f"{_frac(es['healthy_windows_with_present_edge'])}; PRESENT/measured edges " - f"{_frac(es['healthy_present_over_measured_edges'])}; latency-only " - f"{_frac(es['healthy_present_edges_latency_only'])}", + f"- Edge specificity (healthy), latency >= 2x: windows " + f"{_frac(es['healthy_windows_with_latency_high_edge'])}, edges " + f"{_frac(es['healthy_latency_high_over_latency_measured_edges'])} of latency-measured; error " + f"branch: windows {_frac(es['healthy_windows_with_error_present_edge'])}; any PRESENT edge: windows " + f"{_frac(es['healthy_windows_with_present_edge'])}, edges " + f"{_frac(es['healthy_present_over_measured_edges'])} of measured", f"- Series-identity coverage: {ic['named_instance']}/{ic['util_samples']} utilization samples name " f"their instance ({_num(ic['rate'])})", "", - "| case | cause | " + " | ".join(ARMS) + " | retained via | util PRESENT |", - "|---|---|" + "---|" * len(ARMS) + "---|---|", + "Unobservable = no telemetry from the cause service at all; listed as such, not a model failure.", + "", + "| case | cause | cause telemetry | " + " | ".join(ARMS) + " | witnessed by | util PRESENT |", + "|---|---|---|" + "---|" * len(ARMS) + "---|---|", ] for c in report["cases"]: cells = [] @@ -316,6 +352,8 @@ def render_markdown(report: dict, *, provenance: dict) -> str: r = c["arms"][arm] mark = ("✓ " if r["retained"] else "✗ ") if c["positive"] else "" cells.append(f"{mark}{r['outcome']} ({len(r['localization'])})") - lines.append(f"| {c['id']} | {c['cause'] or '—'} | " + " | ".join(cells) - + f" | {', '.join(c['retained_via']) or '—'} | {', '.join(c['util_present']) or '—'} |") + tel = "—" if not c["positive"] else ("yes" if c["cause_has_telemetry"] else "**UNOBSERVABLE**") + via = ", ".join(c["retained_via"]) + (" (util retained it)" if c["retained_by_util"] else "") + lines.append(f"| {c['id']} | {c['cause'] or '—'} | {tel} | " + " | ".join(cells) + + f" | {via or '—'} | {', '.join(c['util_present']) or '—'} |") return "\n".join(lines) + "\n" diff --git a/tests/unit/test_structural_m3.py b/tests/unit/test_structural_m3.py index 223ffe1..d30ff3b 100644 --- a/tests/unit/test_structural_m3.py +++ b/tests/unit/test_structural_m3.py @@ -17,9 +17,11 @@ arm_inputs, build_m3_report, decide, + edge_specificity, evaluate_case, render_markdown, scenario_of, + summarize_arm, ) _W = datetime(2026, 1, 1, 12, 0, 0, tzinfo=timezone.utc) @@ -44,22 +46,25 @@ def __init__(self, service, metric, value, ts, *, attributes=_ONE, metric_type=" _n = 0 -def _call(ts, *, status=OK, child=True, callee="payment", op=CHARGE): +def _call(ts, *, status=OK, child=True, callee="payment", op=CHARGE, dur=10.0, child_status=OK): global _n _n += 1 - out = [_Span("frontend", op, ts, trace=f"t{_n}", sid=f"p{_n}", status=status)] + out = [_Span("frontend", op, ts, trace=f"t{_n}", sid=f"p{_n}", status=status, dur=dur)] if child: out.append(_Span(callee, "handle", ts + timedelta(milliseconds=1), trace=f"t{_n}", sid=f"c{_n}", - parent=f"p{_n}", dur=9.0, status=OK)) + parent=f"p{_n}", dur=dur - 1, status=child_status)) return out -def _traffic(n=10, *, incident_status=OK, incident_child=True, callee="payment", op=CHARGE): +def _traffic(n=10, *, incident_status=OK, incident_child=True, callee="payment", op=CHARGE, incident_dur=10.0, + incident_child_status=OK, incident_errors=None): + """``n`` baseline + ``n`` incident calls. ``incident_errors`` fails only the first k incident calls.""" spans = [] for i in range(n): spans += _call(_W - timedelta(seconds=i + 1), callee=callee, op=op) - spans += _call(_W + timedelta(seconds=i + 1), status=incident_status, child=incident_child, - callee=callee, op=op) + status = incident_status if incident_errors is None else (ERR if i < incident_errors else OK) + spans += _call(_W + timedelta(seconds=i + 1), status=status, child=incident_child, callee=callee, + op=op, dur=incident_dur, child_status=incident_child_status) return spans @@ -70,6 +75,9 @@ def _gauge(service, metric, base, inc, n=5, attributes=_ONE): for i in range(n)] +_GET_ADS = "oteldemo.AdService/GetAds" + + class TestArms: def test_unknown_arm_is_rejected(self): with pytest.raises(ValueError): @@ -83,49 +91,97 @@ def test_unreachable_callee_is_recovered_by_the_edge_family_only(self): assert c.retained("M2a") and c.retained(FULL_ARM) assert "edge" in c.retained_via - def test_locally_silent_cpu_fault_is_recovered_by_the_util_family_only(self): - spans = _traffic() + _traffic(callee="ad", op="oteldemo.AdService/GetAds") + def test_locally_silent_cpu_fault_is_retained_by_util(self): + spans = _traffic() + _traffic(callee="ad", op=_GET_ADS) metrics = _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9) c = evaluate_case("otel_adHighCpu", "ad", True, spans, metrics, _W) assert not c.retained("M1") and not c.retained("M2a") - assert c.retained(FULL_ARM) and c.retained_via == ("util",) + assert c.retained(FULL_ARM) and c.retained_via == ("util",) and c.retained_by_util assert c.util_measured_for_cause and c.util_present == ("util:ad:cpu",) assert (c.util_samples, c.util_samples_named) == (10, 10) + def test_cpu_fault_already_retained_by_sig_is_not_retained_by_util(self): + # ad's own spans fail (sig:ad PRESENT, so M1 already retains ad) AND its CPU is pegged: util + # witnessed the case but did not retain it — the arm delta says so. + spans = _traffic() + _traffic(callee="ad", op=_GET_ADS, incident_child_status=ERR) + metrics = _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9) + c = evaluate_case("otel_adHighCpu", "ad", True, spans, metrics, _W) + assert c.retained("M1") and c.retained("M2a") and c.retained(FULL_ARM) + assert set(c.retained_via) >= {"sig", "util"} + assert not c.retained_by_util + def test_unnamed_utilization_counts_against_identity_coverage(self): metrics = _gauge("ad", "jvm.cpu.recent_utilization", 0.01, 0.9, attributes={}) c = evaluate_case("otel_adHighCpu", "ad", True, _traffic(), metrics, _W) assert (c.util_samples, c.util_samples_named) == (10, 0) assert not c.util_measured_for_cause # unnamed -> UNKNOWN, never measured + def test_universe_counts_metric_only_services(self): + c = evaluate_case("otel_x", "ad", True, _traffic(), _gauge("ad", "jvm.cpu.recent_utilization", 0.3, 0.3), _W) + assert c.n_services == 3 and c.cause_has_telemetry # frontend, payment (spans) + ad (metrics only) + + def test_cause_absent_from_all_telemetry_is_unobservable(self): + c = evaluate_case("otel_kafkaQueueProblems", "kafka", True, _traffic(), [], _W) + assert not c.cause_has_telemetry + def test_healthy_window_abstains_in_every_arm(self): c = evaluate_case("otel_healthy_1", None, False, _traffic(), _gauge("ad", "jvm.cpu.recent_utilization", 0.3, 0.31), _W) - assert all(not c.arms[a].claims for a in ARMS) + assert all(not c.arms[a].generated for a in ARMS) assert c.edges_measured == 1 and c.edges_present == 0 and c.util_present == () -def _case(cid, cause, *, retained=True, n_loc=1, n_services=10, families=("sig",), util_measured=False, - util_present=(), outcome="uncertain"): - loc = ((cause,) + tuple(f"x{i}" for i in range(n_loc - 1))) if (cause and retained) else \ - tuple(f"x{i}" for i in range(n_loc)) - arm = ArmResult(outcome if loc else "no_candidates", loc) - return CaseEval(cid, cause, cause is not None, n_services, cause is not None, {a: arm for a in ARMS}, +class TestEdgeSpecificity: + def test_error_only_edge_is_not_a_latency_flag(self): + # 2/10 explicit errors (>= 5% cutoff), latency unchanged: PRESENT, but never crossed 2x + c = evaluate_case("otel_healthy_1", None, False, _traffic(incident_errors=2), [], _W) + es = edge_specificity([c]) + assert es["healthy_windows_with_present_edge"] == [1, 1] + assert es["healthy_windows_with_error_present_edge"] == [1, 1] + assert es["healthy_windows_with_latency_high_edge"] == [0, 1] + + def test_latency_high_is_counted_even_when_the_error_branch_also_fired(self): + c = evaluate_case("otel_healthy_1", None, False, _traffic(incident_errors=1, incident_dur=30.0), [], _W) + es = edge_specificity([c]) + assert es["healthy_windows_with_latency_high_edge"] == [1, 1] + assert es["healthy_latency_high_over_latency_measured_edges"] == [1, 1] + assert es["healthy_windows_with_error_present_edge"] == [1, 1] + + +def _case(cid, cause, *, retained=True, retained_m2a=None, n_loc=1, n_services=10, families=("sig",), + util_measured=False, util_present=(), outcome="uncertain", util_samples=(0, 0), has_tel=True): + """A hand-built CaseEval. ``retained_m2a`` (default: same as ``retained``) sets the M1/M2a arms + independently of the frozen arm, so the arm delta is under test, not assumed.""" + def arm(kept): + loc = ((cause,) + tuple(f"x{i}" for i in range(n_loc - 1))) if (cause and kept) else \ + tuple(f"x{i}" for i in range(n_loc)) + return ArmResult(outcome if (loc or outcome != "uncertain") else "no_candidates", loc) + kept_m2a = retained if retained_m2a is None else retained_m2a + arms = {"M1": arm(kept_m2a), "M2a": arm(kept_m2a), FULL_ARM: arm(retained)} + return CaseEval(cid, cause, cause is not None, n_services, cause is not None and has_tel, arms, cause_families=tuple(families) if cause else (), util_measured_for_cause=util_measured, - util_present=tuple(util_present)) + util_present=tuple(util_present), util_samples=util_samples[0], + util_samples_named=util_samples[1]) -def _healthy(i, *, claims=False, util_present=()): - return _case(f"otel_healthy_{i}", None, n_loc=1 if claims else 0, util_present=util_present) +def _healthy(i, *, generated=False, util_present=(), outcome="uncertain"): + return _case(f"otel_healthy_{i}", None, n_loc=1 if generated else 0, util_present=util_present, + outcome=outcome) -def _corpus(*, retained=(True, True), n_loc=2, healthy_claims=0, n_healthy=MIN_HEALTHY_NEGATIVES, extra=()): - pos = [_case("otel_a", "a", retained=retained[0], n_loc=n_loc), +def _corpus(*, retained=(True, True), n_loc=2, healthy_generated=0, n_healthy=MIN_HEALTHY_NEGATIVES, extra=(), + coverage=(100, 80)): + pos = [_case("otel_a", "a", retained=retained[0], n_loc=n_loc, util_samples=coverage), _case("otel_b", "b", retained=retained[1], n_loc=n_loc)] - neg = [_healthy(i, claims=i < healthy_claims) for i in range(n_healthy)] + neg = [_healthy(i, generated=i < healthy_generated) for i in range(n_healthy)] return pos + neg + list(extra) +def _util_fault(cid="otel_adHighCpu", cause="ad", **kw): + return _case(cid, cause, **{"families": ("util",), "util_measured": True, + "util_present": (f"util:{cause}:cpu",), **kw}) + + class TestDecisionRules: def test_passes_all_three_pre_registered_bars(self): assert decide(_corpus())["generalization"] == "generalizes" @@ -134,46 +190,83 @@ def test_truth_retained_below_half_fails(self): cases = _corpus(retained=(True, False)) + [_case("otel_c", "c", retained=False)] assert decide(cases)["generalization"] == "does_not_generalize" + def test_unobservable_cause_is_outside_the_retention_denominator(self): + cases = _corpus() + [_case("otel_kafka", "kafka", retained=False, has_tel=False)] + v = decide(cases) + assert v["generalization"] == "generalizes" + assert v["generalization_criteria"]["truth_retained_cause_has_telemetry"][0] == 1.0 + def test_candidate_fraction_above_bound_fails(self): n_loc = int(CANDIDATE_FRACTION_MAX * 10) + 1 # 4/10 > 0.33 assert decide(_corpus(n_loc=n_loc))["generalization"] == "does_not_generalize" + def test_bound_is_the_frozen_literal_so_exactly_one_third_fails(self): + # otel-fresh M1's median fraction is exactly 1/3; the frozen text is "<= 0.33", applied literally + assert decide([ + _case("otel_a", "a", n_loc=1, n_services=3), _case("otel_b", "b", n_loc=1, n_services=3), + *[_healthy(i) for i in range(MIN_HEALTHY_NEGATIVES)]])["generalization"] == "does_not_generalize" + def test_healthy_abstention_below_bound_fails(self): - claims = MIN_HEALTHY_NEGATIVES - int(HEALTHY_ABSTENTION_MIN * MIN_HEALTHY_NEGATIVES) + 1 - assert decide(_corpus(healthy_claims=claims))["generalization"] == "does_not_generalize" + generated = MIN_HEALTHY_NEGATIVES - int(HEALTHY_ABSTENTION_MIN * MIN_HEALTHY_NEGATIVES) + 1 + assert decide(_corpus(healthy_generated=generated))["generalization"] == "does_not_generalize" def test_abstention_at_the_otel_fresh_value_passes(self): # 7/12 = 0.583 is the otel-fresh value the bar was set from; it must pass (>= 0.58) - assert decide(_corpus(healthy_claims=5))["generalization"] == "generalizes" + assert decide(_corpus(healthy_generated=5))["generalization"] == "generalizes" + + def test_abstention_is_no_hypothesis_generated_as_in_the_reference(self): + # a healthy window that generated hypotheses but none compatible is NOT an abstention + cases = _corpus(n_healthy=0) + [_healthy(i, outcome="no_compatible_hypothesis") for i in range(12)] + arm = summarize_arm(cases, FULL_ARM) + assert arm["healthy_abstained"] == [0, 12] and arm["healthy_no_claim"] == [12, 12] + + def test_selectivity_counts_generated_positives_with_an_empty_set_as_zero(self): + cases = [_case("otel_a", "a", n_loc=4), _case("otel_b", "b", retained=False, n_loc=0, + outcome="no_compatible_hypothesis"), + _case("otel_c", "c", retained=False, n_loc=0, outcome="uncertain")] # not generated + arm = summarize_arm(cases, FULL_ARM) + assert arm["generated"] == 2 and arm["median_candidate_fraction"] == pytest.approx(0.2) def test_no_healthy_windows_is_untestable_and_a_deviation(self): v = decide(_corpus(n_healthy=0)) assert v["generalization"] == "untestable" assert any("healthy negatives" in d for d in v["protocol_deviations"]) - def test_m2b_validated_by_a_resource_fault_retained_via_util(self): - extra = [_case("otel_adHighCpu", "ad", families=("util",), util_measured=True, - util_present=("util:ad:cpu",))] - assert decide(_corpus(extra=extra))["m2b"] == "validated" - def test_m2b_untestable_when_no_resource_fault_is_measured(self): - extra = [_case("otel_adHighCpu", "ad", retained=False, families=(), util_measured=False)] - assert decide(_corpus(extra=extra))["m2b"] == "untestable" +class TestM2bRule: + def test_validated_by_the_arm_delta(self): + assert decide(_corpus(extra=[_util_fault(retained_m2a=False)]))["m2b"] == "validated" + + def test_util_present_on_a_resource_fault_sig_already_retained_is_not_validated(self): + # counterexample 1: M2a (no util) retains the cause too, so util did not retain it + fault = _util_fault(retained_m2a=True, families=("sig", "util")) + assert decide(_corpus(extra=[fault]))["m2b"] == "not_validated" + + def test_arm_delta_without_util_on_the_cause_is_not_validated(self): + fault = _case("otel_adHighCpu", "ad", retained_m2a=False, families=("sig",), util_measured=True) + assert decide(_corpus(extra=[fault]))["m2b"] == "not_validated" - def test_m2b_not_validated_when_measured_but_not_retained_via_util(self): - extra = [_case("otel_chaos_recommendationCpuStress", "recommendation", families=("sig",), - util_measured=True)] - assert decide(_corpus(extra=extra))["m2b"] == "not_validated" + def test_high_coverage_with_unmeasured_resource_causes_is_not_validated(self): + # counterexample 2: coverage 0.8, the faults' util never measured -> testable, and it failed + fault = _case("otel_adHighCpu", "ad", retained=False, families=(), util_measured=False) + v = decide(_corpus(extra=[fault], coverage=(100, 80))) + assert v["m2b"] == "not_validated" and v["m2b_identity_coverage"] == pytest.approx(0.8) + + def test_zero_identity_coverage_is_untestable(self): + fault = _case("otel_adHighCpu", "ad", retained=False, families=(), util_measured=False) + assert decide(_corpus(extra=[fault], coverage=(37282, 0)))["m2b"] == "untestable" + + def test_no_utilization_samples_at_all_is_untestable(self): + assert decide(_corpus(coverage=(0, 0)))["m2b"] == "untestable" def test_util_present_on_a_healthy_window_blocks_validation(self): - extra = [_case("otel_adHighCpu", "ad", families=("util",), util_measured=True), - _healthy(99, util_present=("util:cart:cpu",))] + extra = [_util_fault(retained_m2a=False), _healthy(99, util_present=("util:cart:cpu",))] v = decide(_corpus(extra=extra)) assert v["m2b"] == "not_validated" and v["m2b_healthy_util_present"] == ["otel_healthy_99"] - def test_missing_resource_scenarios_are_reported_as_deviations(self): + def test_missing_resource_scenarios_are_a_deviation_and_cannot_validate(self): v = decide(_corpus()) - assert v["m2b"] == "untestable" + assert v["m2b"] == "not_validated" assert any("resource-fault scenarios absent" in d for d in v["protocol_deviations"]) def test_retained_without_a_present_family_is_topology(self): @@ -192,3 +285,11 @@ def test_report_renders_every_arm_and_case(): assert "`abc123`" in md and "**Generalization:** `generalizes`" in md assert all(arm in md for arm in ARMS) assert all(c["id"] in md for c in report["cases"]) + + +def test_report_flags_unobservable_cases(): + cases = _corpus() + [_case("otel_kafka", "kafka", retained=False, has_tel=False)] + md = render_markdown(build_m3_report(cases), provenance={}) + row = next(line for line in md.splitlines() if line.startswith("| otel_kafka ")) + assert "UNOBSERVABLE" in row + assert "UNOBSERVABLE" not in next(line for line in md.splitlines() if line.startswith("| otel_a ")) From a70bd3b2ce613ebeb9afd17e0742f2d4622e6936 Mon Sep 17 00:00:00 2001 From: Leonardo Araujo Date: Sat, 26 Sep 2026 10:28:49 -0300 Subject: [PATCH 3/3] fix(eval): M3 default path refuses artifacts that don't load (#214 review r2) The product's load_ranker/load_calibrator fail open: None falls back to the volume selector or to ordinal confidence. The driver still labelled the output "learned ranker", which would have been a false result on a run that happens once. - It now resolves both paths and loads them with the product loaders before any DB work. It refuses (exit 2) if either load fails. - It records absolute paths, sha256s and *_loaded in provenance. - It refuses if settings did not take rare_event or the ranker path. - The heading names the ranker only when it loaded. Also corrects the constant comments: the abstention gate is "no hypothesis generated", not the ungated healthy_no_claim, and selectivity is taken over generated positives. Co-Authored-By: Claude Opus 5.5 --- scripts/eval/m3_structural.py | 43 +++++++++++++++++++++++++++----- src/eval/structural_m3.py | 5 ++-- tests/unit/test_m3_driver.py | 47 +++++++++++++++++++++++++++++++++++ 3 files changed, 87 insertions(+), 8 deletions(-) create mode 100644 tests/unit/test_m3_driver.py diff --git a/scripts/eval/m3_structural.py b/scripts/eval/m3_structural.py index 4fab832..9236ac2 100644 --- a/scripts/eval/m3_structural.py +++ b/scripts/eval/m3_structural.py @@ -47,6 +47,27 @@ def corpus_hash(cases_dir: Path) -> str: return h.hexdigest() +def _sha256(path: Path) -> str: + return hashlib.sha256(path.read_bytes()).hexdigest() + + +def load_artifacts(ranker: str, calibrator: str) -> dict | str: + """Resolve and load the ranker + calibrator with the product loaders. Returns provenance (absolute + paths, sha256, ``*_loaded``) or, if either fails to load, the reason to refuse the run.""" + from src.core.rca.calibration import load_calibrator + from src.core.rca.ranker import load_ranker + + out: dict = {} + for name, raw, loader in (("ranker", ranker, load_ranker), ("calibrator", calibrator, load_calibrator)): + path = Path(raw).resolve() + if loader(str(path)) is None: + return f"{name} artifact {path} did not load (missing or unparseable); the default path would fall back" + out[name] = str(path) + out[f"{name}_sha256"] = _sha256(path) + out[f"{name}_loaded"] = True + return out + + def main() -> int: ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) ap.add_argument("cases", type=Path, help="captured case directory (one sub-dir per case)") @@ -64,13 +85,22 @@ def main() -> int: "(--allow-dirty for a non-canonical run)", file=sys.stderr) return 2 - # The default explain path as pre-registered: learned ranker + rare_event trigger. + # The default explain path as pre-registered: learned ranker + rare_event trigger. The product + # loaders fail open (None -> volume selector / ordinal confidence), which on a one-shot run would + # be a false result, so both artifacts must load here or the run is refused before any DB work. + artifacts = load_artifacts(args.ranker, args.calibrator) + if isinstance(artifacts, str): + print(f"refusing: {artifacts}", file=sys.stderr) + return 2 os.environ["TRIGGER_MODE"] = "rare_event" - os.environ["RCA_RANKER_MODEL_PATH"] = args.ranker - os.environ["RCA_CALIBRATOR_MODEL_PATH"] = args.calibrator + os.environ["RCA_RANKER_MODEL_PATH"] = artifacts["ranker"] + os.environ["RCA_CALIBRATOR_MODEL_PATH"] = artifacts["calibrator"] from src.config import reload_settings settings = reload_settings() + if settings.trigger_mode != "rare_event" or settings.rca_ranker_model_path != artifacts["ranker"]: + print("refusing: settings did not take the pre-registered trigger mode / ranker path", file=sys.stderr) + return 2 from sqlalchemy import text @@ -92,8 +122,7 @@ def main() -> int: "corpus_sha256": corpus_hash(args.cases), "n_cases": len(cases), "trigger_mode": settings.trigger_mode, - "ranker": args.ranker, - "calibrator": args.calibrator, + **artifacts, } with get_db() as db: @@ -119,7 +148,9 @@ def main() -> int: md = render_markdown(report, provenance=provenance) r = default["raglogs"] - md += ("\n**Default path (learned ranker + rare_event):** " + heading = ("learned ranker + rare_event" if provenance.get("ranker_loaded") + else "RANKER NOT LOADED — volume selector") + md += (f"\n**Default path ({heading}):** " + ", ".join(f"{k}={r[k]}" for k in sorted(r) if not isinstance(r[k], (dict, list))) + f"\n\nLift over baseline: {default['lift_over_baseline']}\n") args.md.write_text(md) diff --git a/src/eval/structural_m3.py b/src/eval/structural_m3.py index 5fa0681..45eb9bb 100644 --- a/src/eval/structural_m3.py +++ b/src/eval/structural_m3.py @@ -41,8 +41,9 @@ # Pre-registered on #209 (M3 protocol), from the otel-fresh results; fixed before capture. TRUTH_RETAINED_MIN = 0.50 # cause-has-telemetry truth retained (M2a's trace-only ceiling) -CANDIDATE_FRACTION_MAX = 0.33 # median |candidates| / |services| over positives that returned any -HEALTHY_ABSTENTION_MIN = 0.58 # healthy windows with no localization claim (7/12 on otel-fresh) +CANDIDATE_FRACTION_MAX = 0.33 # median |localization| / |services| over GENERATED positives (empty = 0); literal +HEALTHY_ABSTENTION_MIN = 0.58 # healthy windows with NO HYPOTHESIS GENERATED (7/12 on otel-fresh); + # not the ungated healthy_no_claim, which also counts NO_COMPATIBLE MIN_HEALTHY_NEGATIVES = 12 # Positives whose cause is a resource fault — the cases M2b exists to recover. RESOURCE_FAULT_SCENARIOS = ("adHighCpu", "recommendationCpuStress") diff --git a/tests/unit/test_m3_driver.py b/tests/unit/test_m3_driver.py new file mode 100644 index 0000000..ad8fa1b --- /dev/null +++ b/tests/unit/test_m3_driver.py @@ -0,0 +1,47 @@ +"""#209 M3 driver: the default-path check runs only with a ranker and calibrator that actually load. +The product loaders fail open (None -> volume selector / ordinal confidence); on a one-shot run that +would be a false result, so the driver refuses — before any DB work.""" +import importlib.util +from pathlib import Path + +_REPO = Path(__file__).resolve().parents[2] +_SCRIPT = _REPO / "scripts" / "eval" / "m3_structural.py" +_RANKER = _REPO / "models" / "rca_ranker.json" +_CALIBRATOR = _REPO / "models" / "rca_calibrator.json" + + +def _mod(): + spec = importlib.util.spec_from_file_location("m3_structural", _SCRIPT) + mod = importlib.util.module_from_spec(spec) + spec.loader.exec_module(mod) + return mod + + +def test_committed_artifacts_load_with_provenance(): + out = _mod().load_artifacts(str(_RANKER), str(_CALIBRATOR)) + assert isinstance(out, dict) + assert out["ranker_loaded"] and out["calibrator_loaded"] + assert Path(out["ranker"]).is_absolute() and len(out["ranker_sha256"]) == 64 + + +def test_missing_ranker_is_refused(tmp_path): + reason = _mod().load_artifacts(str(tmp_path / "nope.json"), str(_CALIBRATOR)) + assert isinstance(reason, str) and "ranker" in reason + + +def test_corrupt_calibrator_is_refused(tmp_path): + bad = tmp_path / "cal.json" + bad.write_text("{not json") + reason = _mod().load_artifacts(str(_RANKER), str(bad)) + assert isinstance(reason, str) and "calibrator" in reason + + +def test_main_refuses_before_touching_the_db(tmp_path, monkeypatch, capsys): + mod = _mod() + monkeypatch.setattr("sys.argv", ["m3_structural.py", str(tmp_path), "--json", str(tmp_path / "o.json"), + "--md", str(tmp_path / "o.md"), "--allow-dirty", + "--ranker", str(tmp_path / "missing.json")]) + monkeypatch.setattr("src.db.session.get_db", lambda: (_ for _ in ()).throw(AssertionError("DB touched"))) + assert mod.main() == 2 + assert "ranker" in capsys.readouterr().err + assert not (tmp_path / "o.md").exists()