diff --git a/scripts/eval/m3_structural.py b/scripts/eval/m3_structural.py new file mode 100644 index 0000000..9236ac2 --- /dev/null +++ b/scripts/eval/m3_structural.py @@ -0,0 +1,162 @@ +#!/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 _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)") + 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. 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"] = 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 + + 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, + **artifacts, + } + + 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"] + 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) + 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..45eb9bb --- /dev/null +++ b/src/eval/structural_m3.py @@ -0,0 +1,360 @@ +"""#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. +* 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 + +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, discretize_ratio +from src.core.rca.structural import resolve_structural +from src.core.rca.structural_model import ( + StructuralInputs, + 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 |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") + + +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 generated(self) -> bool: + """Did this arm generate any hypothesis? Not generating is the (reference) abstention.""" + return self.outcome != "no_candidates" + + +@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 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) + + 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, ...]: + """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.""" + 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) + # 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)) + 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()) + 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( + 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=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), + ) + + +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] + 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": 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 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].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), + } + + +def edge_specificity(cases: list[CaseEval]) -> dict: + """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] + 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": [sum(c.edges_present for c in neg), + sum(c.edges_measured for c in neg)], + } + + +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" + + # 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" + 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" + + 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_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, + } + + +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), "retained_by_util": c.retained_by_util, + "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 (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())), + ] + 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), 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'])})", + "", + "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 = [] + 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'])})") + 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/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_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() diff --git a/tests/unit/test_structural_m3.py b/tests/unit/test_structural_m3.py new file mode 100644 index 0000000..d30ff3b --- /dev/null +++ b/tests/unit/test_structural_m3.py @@ -0,0 +1,295 @@ +"""#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, + edge_specificity, + evaluate_case, + render_markdown, + scenario_of, + summarize_arm, +) + +_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, 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, 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=dur - 1, status=child_status)) + return out + + +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) + 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 + + +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)] + + +_GET_ADS = "oteldemo.AdService/GetAds" + + +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_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",) 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].generated for a in ARMS) + assert c.edges_measured == 1 and c.edges_present == 0 and c.util_present == () + + +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_samples=util_samples[0], + util_samples_named=util_samples[1]) + + +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_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, 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" + + 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): + 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_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"]) + + +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_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 = [_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_a_deviation_and_cannot_validate(self): + v = decide(_corpus()) + 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): + 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"]) + + +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 "))