GitHub Actions
deploy d092bea3608b7a29952f16357fda39b7a29e399b
2e818da
Raw
History Blame Contribute Delete
10.4 kB
"""Aggregation of one benchmark run's `operational.jsonl` + `quality.jsonl` into `summary.json`.
Scope note: this benchmark drives only the retrieval critical path
(`_get_chunks()` -> `rag.retrieval` -> its embedding/chroma/result-prepare
children), not a full backend request spanning memory/context/LLM/citation
subsystems. So "per-request subsystem shares" here are computed relative to
`retrieval_ms` (this workload's own measured parent operation), not
`backend_operation_ms` -- that full-request denominator only applies to a
workload that drives a whole chat/lesson turn, which this one deliberately
does not (see plan preamble: "does not introduce independent application
timers" / drives only the retrieval call under test).
Per the workload protocol, cold and warm rows are never merged into one
figure -- every stratum lives under its own `cache_state` key. Within each
cache_state, rows are additionally stratified by `pipeline_version` and
`query_category`, per the task-4 brief's Step 4.
"""
from __future__ import annotations
import json
import math
from collections import defaultdict
from pathlib import Path
from typing import Any, Iterable
# Retrieval stages this workload actually measures, in critical-path order.
_RETRIEVAL_STAGES: tuple[str, ...] = (
"embedding.query",
"chroma.collection_lookup",
"chroma.collection_count",
"chroma.vector_search",
"rag.result_prepare",
)
def _percentiles(values: list[float]) -> dict[str, float | int | None]:
if not values:
return {"p50": None, "p90": None, "p95": None, "max": None, "count": 0}
ordered = sorted(values)
n = len(ordered)
def _pct(p: float) -> float:
# Nearest-rank percentile over the sorted sample -- simple, deterministic,
# and stable on the small sample sizes this benchmark produces (10-ish
# warm repetitions per query), unlike interpolation methods that can
# imply more precision than a 10-sample distribution actually supports.
rank = max(0, min(n - 1, math.ceil(p * n) - 1))
return ordered[rank]
return {
"p50": _pct(0.50),
"p90": _pct(0.90),
"p95": _pct(0.95),
"max": ordered[-1],
"count": n,
}
def _mean(values: list[float]) -> float | None:
return sum(values) / len(values) if values else None
def load_jsonl(path: Path) -> list[dict[str, Any]]:
path = Path(path)
if not path.exists():
return []
rows: list[dict[str, Any]] = []
with path.open("r", encoding="utf-8") as handle:
for line in handle:
if line.strip():
rows.append(json.loads(line))
return rows
def _experiment(row: dict[str, Any]) -> dict[str, Any]:
return row.get("experiment") or {}
def cache_state_of(row: dict[str, Any]) -> str:
return str(_experiment(row).get("cache_state") or "unknown")
def pipeline_version_of(row: dict[str, Any]) -> str:
return str(_experiment(row).get("pipeline_version") or "unknown")
def query_category_of(row: dict[str, Any]) -> str:
return str(_experiment(row).get("query_category") or "unknown")
def _retrieval_ms_of(row: dict[str, Any]) -> float | None:
return (row.get("stage_durations_ms") or {}).get("retrieval_ms")
def _stage_ms_of(row: dict[str, Any], stage: str) -> float | None:
return (row.get("stage_durations_ms") or {}).get(stage)
def _retrieval_status_of(row: dict[str, Any]) -> str | None:
return (row.get("retrieval") or {}).get("status")
def _retrieval_outcome_counts(rows: Iterable[dict[str, Any]]) -> dict[str, Any]:
counts = {"success": 0, "success_empty": 0, "error_fallback": 0, "unknown": 0}
for row in rows:
status = _retrieval_status_of(row) or "unknown"
counts[status] = counts.get(status, 0) + 1
total = sum(counts.values())
rates = {f"{k}_rate": (v / total if total else None) for k, v in counts.items()}
return {"counts": counts, "total": total, **rates}
def _operational_summary(rows: list[dict[str, Any]]) -> dict[str, Any]:
retrieval_ms = [ms for r in rows if (ms := _retrieval_ms_of(r)) is not None]
stage_summaries: dict[str, Any] = {}
for stage in _RETRIEVAL_STAGES:
values = [ms for r in rows if (ms := _stage_ms_of(r, stage)) is not None]
stage_summaries[stage] = _percentiles(values)
return {"retrieval_ms": _percentiles(retrieval_ms), "stage_ms": stage_summaries}
def _retrieval_stage_shares(rows: list[dict[str, Any]]) -> dict[str, Any]:
"""Per-request stage_ms / retrieval_ms, then aggregated -- never
p95(stage) / p95(retrieval_ms), per the plan preamble's explicit warning
against dividing already-aggregated percentiles.
Named ``retrieval_stage_shares`` (not ``subsystem_shares``) because this
workload only ever drives the retrieval critical path (see module
docstring): every share here is a *retrieval stage's* fraction of
`retrieval_ms`, not a share of the plan's canonical cross-subsystem
denominator (memory/retrieval/context/llm/citation over
`backend_operation_ms`). The two are easy to confuse by name alone, so
this field is named for exactly what it scopes over.
"""
shares: dict[str, list[float]] = defaultdict(list)
for row in rows:
retrieval_ms = _retrieval_ms_of(row)
if not retrieval_ms:
continue
for stage in _RETRIEVAL_STAGES:
value = _stage_ms_of(row, stage)
if value is None:
continue
shares[stage].append(value / retrieval_ms)
return {stage: {"mean": _mean(vals), **_percentiles(vals)} for stage, vals in shares.items()}
def _quality_summary(rows: list[dict[str, Any]]) -> dict[str, Any]:
n = len(rows)
if n == 0:
return {
"count": 0,
"document_hit_rate": None,
"section_hit_rate": None,
"section_hit_applicable_count": 0,
"mrr": None,
"citation_validity_rate": None,
"citation_validity_applicable_count": 0,
"empty_result_rate": None,
"retrieval_error_rate": None,
}
section_applicable = [r for r in rows if r.get("section_hit") is not None]
citation_applicable = [r for r in rows if r.get("citation_validity") is not None]
return {
"count": n,
"document_hit_rate": sum(1 for r in rows if r.get("document_hit")) / n,
"section_hit_rate": (
sum(1 for r in section_applicable if r.get("section_hit")) / len(section_applicable)
if section_applicable
else None
),
"section_hit_applicable_count": len(section_applicable),
"mrr": _mean([r.get("reciprocal_rank", 0.0) for r in rows]),
"citation_validity_rate": (
sum(1 for r in citation_applicable if r.get("citation_validity")) / len(citation_applicable)
if citation_applicable
else None
),
"citation_validity_applicable_count": len(citation_applicable),
"empty_result_rate": sum(1 for r in rows if r.get("empty_result")) / n,
"retrieval_error_rate": sum(1 for r in rows if r.get("retrieval_error")) / n,
}
def _stratum_summary(operational_rows: list[dict[str, Any]], quality_rows: list[dict[str, Any]]) -> dict[str, Any]:
return {
"request_count": len(operational_rows),
"operational": _operational_summary(operational_rows),
"retrieval_outcomes": _retrieval_outcome_counts(operational_rows),
"retrieval_stage_shares": _retrieval_stage_shares(operational_rows),
"quality": _quality_summary(quality_rows),
}
def summarize_run(operational_rows: list[dict[str, Any]], quality_rows: list[dict[str, Any]]) -> dict[str, Any]:
"""Build the nested `cache_state -> pipeline_version -> {overall, by_category}` summary.
Cold and warm strata are always separate top-level keys -- never merged --
per the workload protocol.
"""
quality_by_trace = {r["trace_id"]: r for r in quality_rows if r.get("trace_id")}
by_cache: dict[str, list[dict[str, Any]]] = defaultdict(list)
for row in operational_rows:
by_cache[cache_state_of(row)].append(row)
summary: dict[str, Any] = {"cache_states": {}}
for cache_state, op_rows in by_cache.items():
by_pipeline: dict[str, list[dict[str, Any]]] = defaultdict(list)
for row in op_rows:
by_pipeline[pipeline_version_of(row)].append(row)
pipeline_summary: dict[str, Any] = {}
for pipeline_version, p_rows in by_pipeline.items():
p_quality = [quality_by_trace[r["trace_id"]] for r in p_rows if r.get("trace_id") in quality_by_trace]
overall = _stratum_summary(p_rows, p_quality)
by_category: dict[str, Any] = {}
categories = {query_category_of(r) for r in p_rows}
for category in sorted(categories):
cat_op_rows = [r for r in p_rows if query_category_of(r) == category]
cat_quality = [
quality_by_trace[r["trace_id"]] for r in cat_op_rows if r.get("trace_id") in quality_by_trace
]
by_category[category] = _stratum_summary(cat_op_rows, cat_quality)
pipeline_summary[pipeline_version] = {"overall": overall, "by_category": by_category}
summary["cache_states"][cache_state] = pipeline_summary
return summary
def summarize_run_dir(output_dir: Path) -> dict[str, Any]:
output_dir = Path(output_dir)
operational_rows = load_jsonl(output_dir / "operational.jsonl")
quality_rows = load_jsonl(output_dir / "quality.jsonl")
summary = summarize_run(operational_rows, quality_rows)
manifest_path = output_dir / "manifest.json"
if manifest_path.exists():
manifest = json.loads(manifest_path.read_text(encoding="utf-8"))
summary["run_id"] = manifest.get("run_id")
summary["experiment_id"] = manifest.get("experiment_id")
summary["smoke"] = manifest.get("smoke")
# Carried over so a reader of summary.json alone -- not manifest.json,
# not the source -- still sees the cold-cache-state caveat.
summary["cache_state_semantics"] = manifest.get("cache_state_semantics")
return summary
def write_summary(summary: dict[str, Any], path: Path) -> None:
Path(path).write_text(json.dumps(summary, indent=2) + "\n", encoding="utf-8")