#!/usr/bin/env python3 """Token ledger + budget gate (권고 #1). Claude Code에서 서브에이전트가 끝나면 정확한 토큰 수를 Orchestrator에 반환한다. 그 수치를 append-only 원장에 적고(log), 대표용 대시보드로 렌더(dashboard)하며, tier 예산 초과를 게이트(check)한다. 강제(초과 시 collapse 강등)의 주체는 Orchestrator이고, 이 도구는 계측·게이트다. Usage: token_ledger.py log --workflow WF --role ROLE --tokens N [--wave V] [--tier T] token_ledger.py dashboard # -> reports/TOKENS.md token_ledger.py check --workflow WF --tier T [--wave V] [--add N] # 예산 초과면 exit 2 예산(tier)은 per-wave 다. check/dashboard 는 워크플로 전체가 아니라 해당 wave 만 대조한다 (finding #19). --wave 미지정이면 '-' wave 로 묶여 단일-wave 워크플로 동작이 보존된다. """ import json import os import sys from datetime import datetime, timezone import yaml ROOT = os.environ.get("CLAUDE_PROJECT_DIR") or os.path.dirname( os.path.dirname(os.path.dirname(os.path.abspath(__file__))) ) sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) import _workspace as W # noqa: E402 _ORGWORK = os.path.join(ROOT, "org-os", "06-agent-work") # KPI 예산 = SSOT, org-os 유지 LEDGER = os.path.join(W.state_dir(), "token-ledger.jsonl") KPI = os.path.join(_ORGWORK, "agent-operating-kpi.yaml") DASH = os.path.join(W.reports_dir(), "TOKENS.md") def budgets(): try: d = yaml.safe_load(open(KPI))["agent-operating-kpi"]["token-budgets"] return d.get("per-wave", {}), float(d.get("cost-per-1k-tokens-usd", 0.015)) except Exception: return {"light": 150000, "standard": 500000, "heavy": 2000000}, 0.015 def rows(): if not os.path.exists(LEDGER): return [] out = [] for line in open(LEDGER): line = line.strip() if line: try: out.append(json.loads(line)) except json.JSONDecodeError: pass return out def log(workflow, role, tokens, wave=None, tier=None, usage_source_id=None): per_wave, _ = budgets() if tier not in per_wave: raise ValueError(f"미등록 tier: {tier!r}") value = int(tokens) if value < 0: raise ValueError("tokens는 0 이상이어야 한다") if not str(workflow or "").strip() or workflow == "-": raise ValueError("canonical workflow id 필수") if usage_source_id and any(r.get("usage-source-id") == usage_source_id for r in rows()): raise ValueError(f"중복 usage-source-id: {usage_source_id}") os.makedirs(os.path.dirname(LEDGER), exist_ok=True) rec = { "at": datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"), "workflow": workflow, "wave": wave or "-", "role": role, "tokens": value, "tier": tier, } if usage_source_id: rec["usage-source-id"] = usage_source_id with open(LEDGER, "a") as f: try: import fcntl fcntl.flock(f.fileno(), fcntl.LOCK_EX) except Exception: pass f.write(json.dumps(rec, ensure_ascii=False) + "\n") f.flush() os.fsync(f.fileno()) print(f"[token_ledger] +{tokens} tok · {workflow}/{role}") def sum_workflow(workflow): return sum(r["tokens"] for r in rows() if r.get("workflow") == workflow) def sum_wave(workflow, wave): """단일 wave의 토큰 합. 예산(per-wave)은 wave 단위로 대조해야 하므로 이걸 쓴다. finding #19: 예전 check()/dashboard()는 sum_workflow(워크플로 전체 합)를 per-wave 예산과 비교해, 정상 wave가 여러 번 쌓이면 뒤 wave에서 허위 초과가 났다. wave 미지정 로그는 '-' 버킷으로 묶이므로 단일-wave 워크플로의 기존 동작은 그대로 보존된다.""" w = wave or "-" return sum(r["tokens"] for r in rows() if r.get("workflow") == workflow and (r.get("wave") or "-") == w) def status_icon(used, budget): if not budget: return "—" r = used / budget return "🚨" if r > 1 else ("⚠️" if r > 0.8 else "✅") def dashboard(): # finding #19: 예산 대비는 per-wave 이므로 (워크플로, wave) 단위로 그룹핑해 각 wave를 대조한다. per_wave, cost1k = budgets() data = rows() groups = {} for r in data: groups.setdefault((r.get("workflow", "-"), r.get("wave") or "-"), []).append(r) workflows = {k[0] for k in groups} ts = datetime.now().strftime("%Y-%m-%d %H:%M") total = sum(r["tokens"] for r in data) L = ["# 💰 토큰 대시보드 (대표용)", "", f"생성: {ts} · 총 {total:,} tok · 추정 ${total/1000*cost1k:,.2f} · 워크플로 {len(workflows)}개 · wave {len(groups)}개", f"> tier 예산(per-wave): light {per_wave.get('light',0):,} · standard {per_wave.get('standard',0):,} · heavy {per_wave.get('heavy',0):,} · 초과 시 Orchestrator가 collapse로 강등. (예산 대비는 wave 단위)", "", "| 워크플로 | wave | 워커수 | 토큰 | 추정$ | tier | wave예산대비 | 상태 |", "|---|---|--:|--:|--:|---|---|:--:|"] for wf, wv in sorted(groups): rs = groups[(wf, wv)] tok = sum(r["tokens"] for r in rs) tier = next((r.get("tier") for r in rs if r.get("tier") and r.get("tier") != "-"), "-") budget = per_wave.get(tier) pct = f"{tok/budget*100:.0f}% of {tier}" if budget else "-" L.append(f"| {wf} | {wv} | {len(rs)} | {tok:,} | ${tok/1000*cost1k:,.2f} | {tier} | {pct} | {status_icon(tok, budget)} |") L += ["", "## 워커별 상세", "", "| at | 워크플로 | 역할 | 토큰 |", "|---|---|---|--:|"] for r in sorted(data, key=lambda x: x.get("at", ""), reverse=True): L.append(f"| {r.get('at','-')} | {r.get('workflow','-')} | {r.get('role','-')} | {r['tokens']:,} |") os.makedirs(os.path.dirname(DASH), exist_ok=True) with open(DASH, "w") as f: f.write("\n".join(L) + "\n") print(f"[token_ledger] dashboard -> {os.path.relpath(DASH, ROOT)} ({total:,} tok)") def check(workflow, tier, add=0, wave=None): # finding #19: 예산은 per-wave 이므로 현재 wave의 토큰만 대조한다(워크플로 전체 합 아님). per_wave, _ = budgets() budget = per_wave.get(tier) if budget is None: sys.stderr.write(f"[token_ledger] 미등록 tier: {tier!r}\n") sys.exit(2) addition = int(add or 0) if addition < 0: sys.stderr.write("[token_ledger] --add는 0 이상이어야 한다\n") sys.exit(2) w = wave or "-" used = sum_wave(workflow, w) + addition if budget and used > budget: sys.stderr.write( f"[token_ledger] BUDGET EXCEEDED {workflow}/wave {w}: {used:,} > {tier} per-wave 예산 {budget:,}. " f"fan-out을 collapse(단일 종합)로 강등하거나 tier를 올려라.\n") sys.exit(2) print(f"[token_ledger] OK {workflow}/wave {w}: {used:,}/{budget or '∞'} ({tier})") sys.exit(0) def main(): a = sys.argv[1:] if not a: sys.stderr.write(__doc__) sys.exit(1) cmd, opt = a[0], {} i = 1 while i < len(a): if a[i].startswith("--"): opt[a[i][2:]] = a[i + 1] if i + 1 < len(a) and not a[i + 1].startswith("--") else True i += 2 else: i += 1 if cmd == "log": try: log(opt.get("workflow", "-"), opt.get("role", "-"), opt.get("tokens", 0), opt.get("wave"), opt.get("tier"), opt.get("usage-source-id")) except (TypeError, ValueError) as exc: sys.stderr.write(f"[token_ledger] log 거부: {exc}\n") sys.exit(2) elif cmd == "dashboard": dashboard() elif cmd == "check": check(opt.get("workflow", "-"), opt.get("tier", "standard"), opt.get("add", 0), opt.get("wave")) else: sys.stderr.write(f"unknown command: {cmd}\n") sys.exit(1) if __name__ == "__main__": main()