Add JSONL telemetry and weekly summary (Phase 4)
Observability slim layer. Every full run appends one JSON line to ~/.claude/logs/orchestrator/YYYY-MM-DD.jsonl with timestamp, sha256-prefix task hash (raw task is never logged), retrieved memory names, router choice, runtime, tokens in/out, success/error. risk_class and confidence fields are reserved nulls for Phase 5. - telemetry.py: log_run(), build_record(), task_hash(), token-usage extraction from AIMessage.usage_metadata. log_run swallows all exceptions — telemetry never kills a run. - run.py: wraps app.invoke in try/except with a monotonic-clock window; logs on both success and failure. --route-only path is left unlogged (no agent work, doesn't represent a "run"). - scripts/weekly_summary.py: scans the last 7 days of JSONL and prints a markdown digest (routes, unknown rate, cross-review rate, success rate, total spend, mean tokens/route). Schedule via /schedule and pipe stdout to Slack from the scheduler. Cost rates per route are rough Sonnet/Haiku/GPT/Gemini/DeepSeek defaults suitable for spotting runaway prompts, not finance. Router tokens for structured-output calls aren't captured (they don't surface through the message trail); agent tokens are the dominant component anyway. Validated: golden-set 21/21 still passing; full-run smoke writes expected fields; weekly_summary.py prints clean markdown from a 2-run log.
This commit is contained in:
parent
67180e124e
commit
f05a98dba1
3 changed files with 238 additions and 1 deletions
27
run.py
27
run.py
|
|
@ -8,6 +8,7 @@ Usage:
|
||||||
|
|
||||||
import os
|
import os
|
||||||
import sys
|
import sys
|
||||||
|
import time
|
||||||
|
|
||||||
from dotenv import load_dotenv
|
from dotenv import load_dotenv
|
||||||
|
|
||||||
|
|
@ -15,6 +16,7 @@ os.chdir(os.path.dirname(os.path.abspath(__file__)))
|
||||||
load_dotenv(".env")
|
load_dotenv(".env")
|
||||||
|
|
||||||
from graph import app, retriever_node, router_node # noqa: E402 (env must be loaded before graph imports composio)
|
from graph import app, retriever_node, router_node # noqa: E402 (env must be loaded before graph imports composio)
|
||||||
|
from telemetry import build_record, log_run # noqa: E402
|
||||||
|
|
||||||
|
|
||||||
def _format_retrieved(retrieved: list[dict] | None) -> str:
|
def _format_retrieved(retrieved: list[dict] | None) -> str:
|
||||||
|
|
@ -41,7 +43,30 @@ def main():
|
||||||
print(out["route"])
|
print(out["route"])
|
||||||
return
|
return
|
||||||
|
|
||||||
result = app.invoke({"task": task, "messages": []})
|
started = time.monotonic()
|
||||||
|
result: dict | None = None
|
||||||
|
error: str | None = None
|
||||||
|
try:
|
||||||
|
result = app.invoke({"task": task, "messages": []})
|
||||||
|
except Exception as exc:
|
||||||
|
error = f"{type(exc).__name__}: {exc}"
|
||||||
|
finished = time.monotonic()
|
||||||
|
|
||||||
|
log_run(
|
||||||
|
build_record(
|
||||||
|
task=task,
|
||||||
|
started=started,
|
||||||
|
finished=finished,
|
||||||
|
result=result,
|
||||||
|
success=error is None,
|
||||||
|
error=error,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
if error is not None:
|
||||||
|
print(f"[error] {error}", file=sys.stderr)
|
||||||
|
sys.exit(1)
|
||||||
|
|
||||||
print(_format_retrieved(result.get("retrieved")))
|
print(_format_retrieved(result.get("retrieved")))
|
||||||
print(f"[{result['route']}]")
|
print(f"[{result['route']}]")
|
||||||
print()
|
print()
|
||||||
|
|
|
||||||
121
scripts/weekly_summary.py
Executable file
121
scripts/weekly_summary.py
Executable file
|
|
@ -0,0 +1,121 @@
|
||||||
|
#!/usr/bin/env python3
|
||||||
|
"""Weekly orchestrator summary.
|
||||||
|
|
||||||
|
Scans the last 7 days of telemetry JSONL under ~/.claude/logs/orchestrator/
|
||||||
|
and prints a markdown digest to stdout. Schedule via cron or /schedule;
|
||||||
|
the scheduler pipes the output to Slack.
|
||||||
|
|
||||||
|
Cost estimates use a route -> model-rate map; runs that hit done/unknown
|
||||||
|
are billed at Sonnet (router-only) rates. Numbers are rough — for spotting
|
||||||
|
runaway prompts, not finance.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import datetime
|
||||||
|
import json
|
||||||
|
from collections import Counter, defaultdict
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
LOG_DIR = Path("~/.claude/logs/orchestrator").expanduser()
|
||||||
|
WINDOW_DAYS = 7
|
||||||
|
|
||||||
|
# Rough $/MTok (input, output) by route. Sonnet for router-only routes.
|
||||||
|
COST_RATES: dict[str | None, tuple[float, float]] = {
|
||||||
|
"implementer": (3.00, 15.00),
|
||||||
|
"reviewer": (3.00, 15.00),
|
||||||
|
"researcher": (1.00, 5.00),
|
||||||
|
"cross_reviewer": (2.00, 8.00),
|
||||||
|
"scanner": (1.25, 5.00),
|
||||||
|
"fast_coder": (0.14, 0.28),
|
||||||
|
"connector": (3.00, 15.00),
|
||||||
|
"done": (3.00, 15.00),
|
||||||
|
"unknown": (3.00, 15.00),
|
||||||
|
None: (3.00, 15.00),
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def load_recent_records(window_days: int = WINDOW_DAYS) -> list[dict]:
|
||||||
|
today = datetime.date.today()
|
||||||
|
out: list[dict] = []
|
||||||
|
for i in range(window_days):
|
||||||
|
day = today - datetime.timedelta(days=i)
|
||||||
|
path = LOG_DIR / f"{day.isoformat()}.jsonl"
|
||||||
|
if not path.exists():
|
||||||
|
continue
|
||||||
|
for line in path.read_text().splitlines():
|
||||||
|
line = line.strip()
|
||||||
|
if not line:
|
||||||
|
continue
|
||||||
|
try:
|
||||||
|
out.append(json.loads(line))
|
||||||
|
except json.JSONDecodeError:
|
||||||
|
continue
|
||||||
|
return out
|
||||||
|
|
||||||
|
|
||||||
|
def estimate_cost(record: dict) -> float:
|
||||||
|
rate_in, rate_out = COST_RATES.get(record.get("route"), COST_RATES[None])
|
||||||
|
tokens_in = record.get("tokens_in", 0) or 0
|
||||||
|
tokens_out = record.get("tokens_out", 0) or 0
|
||||||
|
return (tokens_in * rate_in + tokens_out * rate_out) / 1_000_000
|
||||||
|
|
||||||
|
|
||||||
|
def summarize(records: list[dict], window_days: int = WINDOW_DAYS) -> str:
|
||||||
|
if not records:
|
||||||
|
return (
|
||||||
|
"# Orchestrator weekly summary\n\n"
|
||||||
|
f"_No runs logged in the last {window_days} days._"
|
||||||
|
)
|
||||||
|
|
||||||
|
total = len(records)
|
||||||
|
successes = sum(1 for r in records if r.get("success"))
|
||||||
|
routes = Counter(r.get("route") for r in records)
|
||||||
|
|
||||||
|
by_route_tokens: dict[str, list[tuple[int, int]]] = defaultdict(list)
|
||||||
|
for r in records:
|
||||||
|
by_route_tokens[r.get("route") or "null"].append(
|
||||||
|
(r.get("tokens_in", 0) or 0, r.get("tokens_out", 0) or 0)
|
||||||
|
)
|
||||||
|
|
||||||
|
total_cost = sum(estimate_cost(r) for r in records if r.get("success"))
|
||||||
|
total_tokens_in = sum((r.get("tokens_in", 0) or 0) for r in records)
|
||||||
|
total_tokens_out = sum((r.get("tokens_out", 0) or 0) for r in records)
|
||||||
|
|
||||||
|
unknown_rate = routes.get("unknown", 0) / total
|
||||||
|
cross_rate = routes.get("cross_reviewer", 0) / total
|
||||||
|
success_rate = successes / total
|
||||||
|
|
||||||
|
lines = [
|
||||||
|
"# Orchestrator weekly summary",
|
||||||
|
f"_Last {window_days} days · {total} runs · {successes} succeeded_",
|
||||||
|
"",
|
||||||
|
"## Routes",
|
||||||
|
]
|
||||||
|
for route, count in routes.most_common():
|
||||||
|
label = route if route is not None else "null"
|
||||||
|
lines.append(f"- `{label}`: {count} ({count / total:.0%})")
|
||||||
|
lines += [
|
||||||
|
"",
|
||||||
|
"## Health",
|
||||||
|
f"- Unknown route rate: **{unknown_rate:.1%}** (target <2%)",
|
||||||
|
f"- Cross-review rate: **{cross_rate:.1%}**",
|
||||||
|
f"- Success rate: **{success_rate:.1%}**",
|
||||||
|
"",
|
||||||
|
"## Spend (rough)",
|
||||||
|
f"- Total tokens: {total_tokens_in:,} in / {total_tokens_out:,} out",
|
||||||
|
f"- Estimated cost: **${total_cost:.2f}**",
|
||||||
|
"",
|
||||||
|
"## Tokens per route (mean in/out per run)",
|
||||||
|
]
|
||||||
|
for route, tokens in sorted(by_route_tokens.items(), key=lambda x: -len(x[1])):
|
||||||
|
n = len(tokens)
|
||||||
|
mean_in = sum(t[0] for t in tokens) / n
|
||||||
|
mean_out = sum(t[1] for t in tokens) / n
|
||||||
|
lines.append(f"- `{route}`: {mean_in:,.0f} in / {mean_out:,.0f} out ({n} runs)")
|
||||||
|
|
||||||
|
return "\n".join(lines)
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
print(summarize(load_recent_records()))
|
||||||
91
telemetry.py
Normal file
91
telemetry.py
Normal file
|
|
@ -0,0 +1,91 @@
|
||||||
|
"""JSONL telemetry for orchestrator runs.
|
||||||
|
|
||||||
|
Appends one JSON line per run to ~/.claude/logs/orchestrator/YYYY-MM-DD.jsonl.
|
||||||
|
Logging failures are swallowed — telemetry must never kill the orchestrator.
|
||||||
|
|
||||||
|
Fields (Phase 4):
|
||||||
|
timestamp ISO-8601 UTC
|
||||||
|
task_hash sha256(task)[:16] — never log raw task content
|
||||||
|
retrieved list[str] — memory names surfaced by the retriever
|
||||||
|
route str | None — the router's choice (or None on early failure)
|
||||||
|
risk_class null — placeholder for Phase 5
|
||||||
|
confidence null — placeholder for Phase 5
|
||||||
|
runtime_seconds float
|
||||||
|
tokens_in int
|
||||||
|
tokens_out int
|
||||||
|
success bool
|
||||||
|
error str | None — type+message if the run raised
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import datetime
|
||||||
|
import hashlib
|
||||||
|
import json
|
||||||
|
import os
|
||||||
|
from pathlib import Path
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
LOG_DIR = Path(os.path.expanduser("~/.claude/logs/orchestrator"))
|
||||||
|
|
||||||
|
|
||||||
|
def task_hash(task: str) -> str:
|
||||||
|
return hashlib.sha256(task.encode("utf-8")).hexdigest()[:16]
|
||||||
|
|
||||||
|
|
||||||
|
def extract_token_usage(messages: list[Any]) -> tuple[int, int]:
|
||||||
|
"""Sum input/output token usage across all AIMessages in the run.
|
||||||
|
|
||||||
|
Each LangChain AIMessage from a provider with usage telemetry exposes a
|
||||||
|
`usage_metadata` dict with at least input_tokens / output_tokens. Messages
|
||||||
|
without usage are skipped.
|
||||||
|
"""
|
||||||
|
in_tok = 0
|
||||||
|
out_tok = 0
|
||||||
|
for m in messages or []:
|
||||||
|
usage = getattr(m, "usage_metadata", None)
|
||||||
|
if not usage:
|
||||||
|
continue
|
||||||
|
in_tok += int(usage.get("input_tokens", 0) or 0)
|
||||||
|
out_tok += int(usage.get("output_tokens", 0) or 0)
|
||||||
|
return in_tok, out_tok
|
||||||
|
|
||||||
|
|
||||||
|
def build_record(
|
||||||
|
task: str,
|
||||||
|
started: float,
|
||||||
|
finished: float,
|
||||||
|
result: dict | None,
|
||||||
|
success: bool,
|
||||||
|
error: str | None,
|
||||||
|
) -> dict:
|
||||||
|
retrieved_names = [
|
||||||
|
m.get("name", "?") for m in (result or {}).get("retrieved", []) or []
|
||||||
|
]
|
||||||
|
tokens_in, tokens_out = extract_token_usage((result or {}).get("messages", []))
|
||||||
|
return {
|
||||||
|
"timestamp": datetime.datetime.now(datetime.UTC).isoformat(),
|
||||||
|
"task_hash": task_hash(task),
|
||||||
|
"retrieved": retrieved_names,
|
||||||
|
"route": (result or {}).get("route"),
|
||||||
|
"risk_class": None,
|
||||||
|
"confidence": None,
|
||||||
|
"runtime_seconds": round(finished - started, 3),
|
||||||
|
"tokens_in": tokens_in,
|
||||||
|
"tokens_out": tokens_out,
|
||||||
|
"success": success,
|
||||||
|
"error": error,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def log_run(record: dict, log_dir: Path = LOG_DIR) -> None:
|
||||||
|
"""Append a single record as one JSON line. Never raises."""
|
||||||
|
try:
|
||||||
|
log_dir.mkdir(parents=True, exist_ok=True)
|
||||||
|
date = datetime.date.today().isoformat()
|
||||||
|
path = log_dir / f"{date}.jsonl"
|
||||||
|
with path.open("a") as f:
|
||||||
|
f.write(json.dumps(record) + "\n")
|
||||||
|
except Exception:
|
||||||
|
# Telemetry never kills the run.
|
||||||
|
pass
|
||||||
Reference in a new issue