- Add 6 retrieval-augmented routing tests (3 live retrieval, 3 off-topic fake memories) to unblock Phase 5 - Defer Composio tool loading and graph construction to first use so expired or missing keys don't crash imports - Atomic cache write in retriever via temp file (open item #2) - Log rotation in weekly_summary.py, pruning JSONL >90 days (open item #3)
148 lines
5 KiB
Python
Executable file
148 lines
5 KiB
Python
Executable file
#!/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]:
|
|
# UTC date aligns with telemetry.py's log filename convention.
|
|
today = datetime.datetime.now(datetime.UTC).date()
|
|
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)
|
|
failed_cost = sum(estimate_cost(r) for r in records if not 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}** (incl. ${failed_cost:.2f} on failed runs)",
|
|
"",
|
|
"## 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)
|
|
|
|
|
|
RETENTION_DAYS = 90
|
|
|
|
|
|
def prune_old_logs(retention_days: int = RETENTION_DAYS) -> int:
|
|
"""Delete JSONL logs older than retention_days. Returns count deleted."""
|
|
if not LOG_DIR.is_dir():
|
|
return 0
|
|
cutoff = datetime.datetime.now(datetime.UTC).date() - datetime.timedelta(
|
|
days=retention_days
|
|
)
|
|
deleted = 0
|
|
for path in LOG_DIR.glob("*.jsonl"):
|
|
try:
|
|
file_date = datetime.date.fromisoformat(path.stem)
|
|
except ValueError:
|
|
continue
|
|
if file_date < cutoff:
|
|
path.unlink()
|
|
deleted += 1
|
|
return deleted
|
|
|
|
|
|
if __name__ == "__main__":
|
|
print(summarize(load_recent_records()))
|
|
pruned = prune_old_logs()
|
|
if pruned:
|
|
print(f"\n_Pruned {pruned} log file(s) older than {RETENTION_DAYS} days._")
|