multi-agent-system / scripts /profile_pipeline.py
firepenguindisopanda
Refactor test documentation and improve consistency in comments
e382248
Raw
History Blame Contribute Delete
9.99 kB
#!/usr/bin/env python
"""Deterministic pipeline profiler.
Replaces every LLM call with a fake model of *constant* latency, so what is
measured is the shape of the graph rather than NVIDIA NIM's variance. The hosted
endpoint returned 7.5s and 146.4s for identical calls during this project, which
made single real runs useless for judging a 10-20% orchestration change.
With a constant per-call latency L the numbers become structural:
effective serial calls = wall_clock / L <- critical path depth
parallelism factor = (calls * L) / wall_clock
agent invocations = should equal the number of agents, exactly once
That last one is the cheap regression guard: a bug where every phase-3 agent ran
three times (once per judge completion, via a gate that re-fired) showed up here
in seconds, having previously cost three 20-minute real runs to find.
Usage:
python scripts/profile_pipeline.py # 0.5s per LLM call
python scripts/profile_pipeline.py --latency 0.1 # faster sweep
python scripts/profile_pipeline.py --json # machine readable
"""
from __future__ import annotations
import argparse
import asyncio
import json
import sys
import time
from collections import Counter
from pathlib import Path
from unittest.mock import patch
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from langchain_core.messages import AIMessage # noqa: E402
from langchain_core.runnables import Runnable # noqa: E402
FAKE_MARKDOWN = (
"# Section\n\n"
"This is generated specification prose used only for timing.\n\n"
"- **Point one:** detail\n- **Point two:** detail\n"
)
# (kind, start, end) for every simulated LLM call
CALLS: list[tuple[str, float, float]] = []
T0 = 0.0
def _minimal(schema):
"""Smallest valid instance of each structured-output schema."""
name = schema.__name__
if name == "JudgeVerdict":
return schema(
is_approved=True, score=8, issues=[],
recommended_action="accept", reasoning="timing harness",
)
if name == "CriticScores":
return schema(
role="product_owner", dimensions=[], overall_score=8.0,
passed=True, recommendation="proceed", summary="timing harness",
)
if name == "SkepticAttackVectors":
return schema(
role="product_owner", attack_vectors=[],
risk_level="low", summary="timing harness",
)
return schema()
def make_fake(kind: str, latency: float, role=None):
"""A Runnable that sleeps `latency` and returns canned content."""
label = f"{kind}:{getattr(role, 'value', role)}" if role else kind
async def _sleep_and_record():
start = time.time() - T0
await asyncio.sleep(latency)
CALLS.append((label, start, time.time() - T0))
class Fake(Runnable):
model = f"fake-{kind}"
async def ainvoke(self, input, config=None, **kw):
await _sleep_and_record()
return AIMessage(content=FAKE_MARKDOWN)
def invoke(self, input, config=None, **kw):
CALLS.append((label, time.time() - T0, time.time() - T0))
return AIMessage(content=FAKE_MARKDOWN)
async def astream(self, input, config=None, **kw):
# Stream in a few pieces so the token path is exercised too.
pieces = 4
for i in range(pieces):
await asyncio.sleep(latency / pieces)
yield AIMessage(content=FAKE_MARKDOWN[
i * len(FAKE_MARKDOWN) // pieces:
(i + 1) * len(FAKE_MARKDOWN) // pieces
])
CALLS.append((label, time.time() - T0 - latency, time.time() - T0))
def with_structured_output(self, schema, **kw):
outer = self
class Structured(Runnable):
async def ainvoke(self, input, config=None, **k):
await _sleep_and_record()
return _minimal(schema)
def invoke(self, input, config=None, **k):
return _minimal(schema)
return Structured()
def bind_tools(self, *a, **k):
return self
return Fake()
class StubRAG:
"""No network, no embeddings."""
def retrieve(self, *a, **k):
return []
def retrieve_with_books(self, *a, **k):
return []
def retrieve_book_rules(self, *a, **k):
return []
def format_docs(self, *a, **k):
return ""
def format_docs_with_books(self, *a, **k):
return ""
def ensure_initialized(self):
return None
async def run(latency: float) -> dict:
global T0
CALLS.clear()
from app.core.schemas import ProjectRequest
patches = [
patch("app.core.pipeline_orchestrator.RAGService", StubRAG),
patch("app.agents.specialists.generic_specialist.get_rag_service",
lambda *a, **k: StubRAG()),
patch("app.agents.specialists.generic_specialist.get_chat_model",
lambda *a, **k: make_fake("agent", latency, role=k.get("role"))),
patch("app.agents.auxiliary.critic_agent.get_critic_model",
lambda *a, **k: make_fake("critic", latency)),
patch("app.agents.auxiliary.skeptic_agent.get_skeptic_model",
lambda *a, **k: make_fake("skeptic", latency)),
patch("app.agents.auxiliary.verdict_judge.get_judge_model",
lambda *a, **k: make_fake("judge", latency)),
]
for p in patches:
p.start()
from app.core.pipeline_orchestrator import PipelineOrchestrator
orch = PipelineOrchestrator()
req = ProjectRequest(description="A tool-sharing app for neighbours.")
T0 = time.time()
events = Counter()
agent_starts = Counter()
async for ev in orch.run_pipeline_streaming(req):
events[ev.get("type")] += 1
if ev.get("type") == "agent_start":
agent_starts[ev["role"]] += 1
wall = time.time() - T0
for p in patches:
p.stop()
kinds = Counter(k.split(":", 1)[0] for k, _, _ in CALLS)
per_role = Counter(k.split(":", 1)[1] for k, _, _ in CALLS
if k.startswith("agent:"))
n_agents = len(orch.PHASE_1 + orch.PHASE_2 + orch.PHASE_3)
return {
"latency": latency,
"wall_clock": wall,
"llm_calls": dict(kinds),
"total_llm_calls": len(CALLS),
"effective_serial_calls": wall / latency if latency else 0,
"parallelism_factor": (len(CALLS) * latency / wall) if wall else 0,
"expected_agent_calls": n_agents,
"agent_starts": dict(agent_starts),
"agent_calls_per_role": dict(per_role),
"events": dict(events),
}
def report(r: dict) -> None:
print("\n" + "=" * 64)
print(f"DETERMINISTIC PIPELINE PROFILE (constant {r['latency']}s per LLM call)")
print("=" * 64)
print(f" wall clock {r['wall_clock']:8.2f}s")
if "critical_path_depth" in r:
print(f" critical path depth {r['critical_path_depth']:8.1f} <- overhead-corrected")
print(f" fixed overhead {r['fixed_overhead_s']:8.2f}s")
else:
print(f" effective serial calls {r['effective_serial_calls']:8.1f} <- inflated at low latency; use --calibrate")
print(f" parallelism factor {r['parallelism_factor']:8.2f}x")
print(f" total simulated LLM calls {r['total_llm_calls']:8d}")
print("\n calls by kind:")
for kind, n in sorted(r["llm_calls"].items()):
print(f" {kind:<10}{n:>4}")
per_role = r["agent_calls_per_role"]
repeats = {k: v for k, v in per_role.items() if v > 1}
print(f"\n agents invoked: {len(per_role)} distinct, "
f"{sum(per_role.values())} calls (1 per agent = "
f"{r['expected_agent_calls']})")
if repeats:
print("\n roles running more than once:")
for role, n in sorted(repeats.items(), key=lambda x: -x[1]):
print(f" {role:<22}{n}x")
print("\n NOTE: re-runs are legitimate when the ARI check re-prompts a")
print(" role. They are a BUG when a whole phase repeats together -")
print(" that means a gate re-fired and re-fanned-out behind it.")
phase3 = {"api_designer", "qa_strategist", "devops_architect",
"spec_coordinator"}
if phase3 <= set(repeats):
print("\n !! ALL of phase 3 repeated - this is the fan-out bug.")
else:
print(" OK - every agent ran exactly once")
print()
async def calibrate(lo: float = 0.2, hi: float = 0.5) -> dict:
"""Two-point measurement that cancels fixed interpreter overhead.
wall(L) = depth * L + overhead, so the slope between two latencies is the
critical path depth and the intercept is the overhead. A single division
(wall / L) conflates the two and inflates the depth badly at small L -
0.02s latency reported 39.9 serial calls against a true value near 8.
"""
a = await run(lo)
b = await run(hi)
depth = (b["wall_clock"] - a["wall_clock"]) / (hi - lo)
overhead = a["wall_clock"] - depth * lo
return {
**b,
"critical_path_depth": depth,
"fixed_overhead_s": overhead,
"calibration": {"lo": lo, "hi": hi,
"wall_lo": a["wall_clock"], "wall_hi": b["wall_clock"]},
}
if __name__ == "__main__":
ap = argparse.ArgumentParser()
ap.add_argument("--latency", type=float, default=0.5,
help="simulated seconds per LLM call (default 0.5)")
ap.add_argument("--json", action="store_true", help="emit JSON only")
ap.add_argument("--calibrate", action="store_true",
help="two-point measurement; cancels fixed overhead")
args = ap.parse_args()
result = asyncio.run(calibrate() if args.calibrate else run(args.latency))
if args.json:
print(json.dumps(result, indent=2))
else:
report(result)