From f05a98dba13faa787ae9b8fcfbb4350f13865b99 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Fri, 15 May 2026 11:49:05 -0400 Subject: [PATCH] Add JSONL telemetry and weekly summary (Phase 4) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- run.py | 27 ++++++++- scripts/weekly_summary.py | 121 ++++++++++++++++++++++++++++++++++++++ telemetry.py | 91 ++++++++++++++++++++++++++++ 3 files changed, 238 insertions(+), 1 deletion(-) create mode 100755 scripts/weekly_summary.py create mode 100644 telemetry.py diff --git a/run.py b/run.py index 9c49517..79eb636 100755 --- a/run.py +++ b/run.py @@ -8,6 +8,7 @@ Usage: import os import sys +import time from dotenv import load_dotenv @@ -15,6 +16,7 @@ os.chdir(os.path.dirname(os.path.abspath(__file__))) load_dotenv(".env") 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: @@ -41,7 +43,30 @@ def main(): print(out["route"]) 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(f"[{result['route']}]") print() diff --git a/scripts/weekly_summary.py b/scripts/weekly_summary.py new file mode 100755 index 0000000..89b91c7 --- /dev/null +++ b/scripts/weekly_summary.py @@ -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())) diff --git a/telemetry.py b/telemetry.py new file mode 100644 index 0000000..9408413 --- /dev/null +++ b/telemetry.py @@ -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