Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
162 changes: 162 additions & 0 deletions scripts/eval/m3_structural.py
Original file line number Diff line number Diff line change
@@ -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())
64 changes: 40 additions & 24 deletions src/core/rca/structural.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,30 +21,23 @@
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,
)
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,
Expand All @@ -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)
Loading
Loading