Coût, tokens et latence par phase & Exporter vers OpenTelemetry, Prometheus et Grafana
La salle de contrôle apprend à compter : coût, jetons et latence phase par phase dans la trace — et un pont d'une commande qui projette vos runs dans OpenTelemetry, Prometheus et Grafana.
Hier, les couloirs de nage vous ont montré où passent les minutes d’un run. Mais posez-vous la
question qui fâche : où passe l’argent ? Votre journal sait qu’un SDLC a coûté 47 centimes,
un seul chiffre, au niveau du run. La phase build en a-t-elle mangé les trois quarts ? Le
plan vaut-il son prix ? Aujourd’hui, vous ne pouvez pas répondre sans deviner.
À la fin de ce chapitre, chaque phase de votre trace portera ses trois mesures (coût, jetons, durée) sans qu’une ligne de vos ADW ne change, et un pont d’une commande projettera vos runs dans la langue commune de l’observabilité, OpenTelemetry, pour qui veut les voir dans Grafana à côté du reste de son infrastructure. Deux pièces closent le jalon du module : le journal du chapitre 18 gagne son grain comptable, et la salle de contrôle gagne sa porte vers l’extérieur.
Coût, tokens et latence par phase
L’idée en une phrase
Le runner mesure ce que chaque tentative ajoute aux compteurs du run, en delta avant/après l’action, et l’écrit sur la ligne de phase du journal : la phase devient l’unité comptable de l’usine, et toute la mécanique vit côté code déterministe, sans toucher ni à vos ADW ni au port harnais.
Points clés
- Le delta ne demande rien à personne. Vos actions créditent déjà
run.cost_usddepuis le chapitre 8. Le runner lit le compteur avant et après chaque tentative, et la différence est le prix de la tentative. Zéro changement dans vos scripts, zéro changement dans le port : la couture ne bouge pas d’un millimètre. - Les reprises s’additionnent. La ligne de phase accumule tentative après tentative : un
ok t3porte le prix de ses trois essais, pas celui du dernier. C’est le vrai coût de la phase, celui que la vue d’hier vous laissait seulement soupçonner. - La latence était déjà là.
started_atetended_atdatent chaque phase depuis le chapitre 18. Le coût et les jetons les rejoignent, et les trois questions (où vont les minutes, où va l’argent, qui a retenté) se répondent d’une seule requête. - Les jetons voyagent par le même canal, honnêtement. Le
Rungagne un compteurtokensque toute action peut créditer comme elle crédite le coût, le jour où elle relève l’usage que son harnais rapporte (entrée, sortie, lectures et écritures de cache, le raisonnement comptant dans la sortie). Tant qu’une action ne rapporte rien, la colonne reste à zéro : « rien relevé », jamais un chiffre inventé.
Exemple concret
Reprenez le SDLC du chapitre 13 : dix-huit phases, ~40 à 60 centimes, une dizaine de
minutes. Une requête sur les phases du run, et la répartition tombe : build, verte à la
deuxième tentative, porte environ les deux tiers du coût, plan un gros quart, scout
quelques centimes, et les quatorze phases code affichent 0,00 $ tout rond. La décision
suit dans la seconde : si ce run est trop cher, c’est le prompt du builder qu’on retravaille
(chapitre 9) ou son étage qu’on descend dans le roster (chapitre 15), pas les gates, qui ne
coûtent rien, ni le scout, qui ne pèse rien.
Les mesures de la ligne de phase
| Mesure | D’où elle vient | La question qu’elle règle |
|---|---|---|
Latence (started_at/ended_at) | les horodatages du chapitre 18 | où vont les minutes |
Coût (cost_usd) | le delta du compteur, sommé sur les tentatives | où va l’argent |
Jetons (tokens) | le même delta, quand l’action relève l’usage | ce que pèse la conversation |
Verdict et tentatives (status, attempt) | le runner, depuis le chapitre 18 | ce que le prix a acheté |
Config — le cœur du delta, dans le runner
Deux lignes encadrent chaque tentative, c’est toute la mécanique. La pièce complète est dans le TP, l’extrait montre où elle vit :
# Le grain comptable : ce que CETTE tentative ajoute aux compteurs du run.
# Mesure en delta — vos actions creditent, le runner constate.
cost_before, tokens_before = self.cost_usd, self.tokens
self.results[spec.name] = spec.action(self, attempt)
# ... et le verdict part au tracer AVEC son prix :
# cost_usd=self.cost_usd - cost_before, tokens=self.tokens - tokens_before
Et la question « où va l’argent » devient une commande, en une ligne portable :
sqlite3 adws/adw_data/factory.db "SELECT name, ROUND(SUM(cost_usd),2), COUNT(*) FROM phases GROUP BY name ORDER BY 2 DESC LIMIT 5;"
Piège courant : « le coût total du run suffit pour piloter » est inexact. Deux runs à 50 centimes peuvent raconter deux histoires opposées : l’un brûle tout dans un
buildqui retente, l’autre dans unplantrop bavard. Le total dit combien, seul le grain par phase dit quoi améliorer, et c’est la seule information qui rende une décision.
Exporter vers OpenTelemetry, Prometheus et Grafana
L’idée en une phrase
Le pont vers l’observabilité d’équipe est un export après coup : un script déterministe
relit factory.db et le traduit, les runs en traces OpenTelemetry, les agrégats en
métriques Prometheus, sans qu’aucune dépendance ne remonte vers l’usine. Le chemin reste
usine → SQLite → lecture, et l’export n’est qu’une lecture de plus.
Points clés
- OpenTelemetry est la langue commune des traces. Un run devient une trace, chaque phase un span enfant du run, avec les mêmes horodatages, les mêmes verdicts, le coût et les jetons en attributs. Tout backend qui parle OTLP les accepte : un collector local, Jaeger, Grafana et son monde.
- La stdlib suffit. OTLP/HTTP accepte le JSON (
http/jsonest une variante officielle du protocole) : un POST surlocalhost:4318/v1/traces, le port par défaut d’un collector, avecContent-Type: application/json, et vos spans partent. Aucun SDK, aucune dépendance. - Les identifiants sont dérivés, l’export est rejouable. L’id de trace se calcule depuis
l’
adw_id, l’id de span depuis lephase_id: exporter deux fois le même run produit exactement les mêmes spans, sans doublons chez le backend, et un export raté se relance sans arrière-pensée. - Prometheus prend les agrégats, Grafana fait la vitrine. L’exposition Prometheus est un format texte lisible à l’œil : runs par verdict, dollars cumulés par ADW et par phase, secondes par phase. Prometheus vient les chercher, Grafana les dessine, et vos coûts d’usine s’affichent à côté de vos métriques d’infrastructure.
Exemple concret
Votre équipe a déjà un Grafana. Vous lancez un collector local, puis
uv run adws/obs_export.py --otlp : vos dix derniers runs partent en quelques centaines de
millisecondes. Le SDLC d’hier arrive en trace de dix-neuf spans, et son flame graph
reproduit exactement les couloirs du chapitre 19, build en bloc massif, les gates en traits
fins. Coût de l’opération : zéro token, puisque ces runs sont déjà payés et que l’export ne fait que
les raconter ailleurs. Et si vous n’avez ni collector ni Grafana : --prom imprime les
agrégats sur la sortie standard, et le journal local reste, comme toujours, la vérité.
Le journal local et le pont : deux portées
| SQLite local (ch. 18-19) | OTel / Prometheus / Grafana | |
|---|---|---|
| Portée | vous, votre repo, votre terminal | l’équipe, ses dashboards existants |
| Répond à | quoi, où, quand, combien, au grain fin | tendances, comparaisons, alertes |
| Dépendances | aucune | un collector ou un Prometheus debout |
| Source de vérité | oui, toujours | non : une projection, reconstructible |
Commande — le pont, dans les deux sens
Une seule version suffit ici : aucun harnais en vue, le pont est du code déterministe qui lit du SQLite. Une commande par ligne, exécutables telles quelles sous bash comme sous PowerShell :
uv run adws/obs_export.py --prom
uv run adws/obs_export.py --otlp
uv run adws/obs_export.py --otlp http://mon-collector:4318
Piège courant : « pour voir mes runs dans Grafana, il faut instrumenter le runner avec un SDK OpenTelemetry » est inexact. L’instrumentation vivante inverserait le sens de la dépendance : chaque run se mettrait à parler à un collector, et un collector éteint deviendrait votre problème pendant le run. La trace du chapitre 18 contient déjà tout. L’export après coup rejoue des runs déjà payés, même ceux de la semaine dernière, et un collector éteint ne coûte qu’un export raté, jamais un run.
Fil rouge — la pièce posée aujourd’hui
La zone « salle de contrôle » du plan atteint son jalon : chaque run laisse une trace SQLite
exploitable : coût, jetons, durée, par phase. Trois pièces s’emboîtent : le tracer du
chapitre 18 gagne ses colonnes comptables (et migre en place le journal que vous avez déjà),
le runner du chapitre 8 mesure le delta de chaque tentative, et adws/obs_export.py ouvre la
porte vers OpenTelemetry et Prometheus. La couture ne bouge pas : l’agent propose, le code
dispose, et le code compte. Les agents ignorent toujours qu’ils sont observés, aucune
enveloppe ne traverse, la dépendance va de l’export vers le journal, jamais du run vers un
dashboard. À l’usage : zéro token, quelques millisecondes par verdict tracé, quelques
centaines de millisecondes par export. À l’économie : la question « quel prompt, quel étage
de roster améliorer » se répondait à l’aveugle, elle se répond désormais par une requête
gratuite sur des runs déjà payés.
Travaux pratiques — la pièce du jour
Une pièce complète à poser dans le repo compagnon plume-factory, qui devient, chapitre après
chapitre, votre usine logicielle agentique. Aujourd’hui, trois fichiers : le tracer et le
runner qui apprennent à compter, et le pont vers l’extérieur.
Pièce — adws/adw_modules/tracer.py
Cette version remplace celle du chapitre 18. Les phases portent désormais leur coût et
leurs jetons, accumulés tentative après tentative, et les runs leur total de jetons. Votre
journal existant est migré en place au premier lancement, sans perdre une ligne. Tout le
reste est inchangé : JSONL d’abord, WAL, --last, le vert qui se mérite.
# /// script
# requires-python = ">=3.11"
# ///
"""tracer — le journal de l'usine : chaque evenement, au moment ou il arrive.
Deux supports, une verite. Le JSONL est l'enregistrement brut — un fichier
par run, un evenement par ligne, ecrit au fil de l'eau ; SQLite
(factory.db, mode WAL) est le miroir requetable, mis a jour dans le meme
geste. Pas de serveur, pas de push : le chemin est toujours
usine -> sqlite -> lecture.
Version chapitre 20 : le grain comptable. Les phases portent leur cout et
leurs jetons — accumules tentative apres tentative — et les runs leur
total de jetons. Un journal deja rempli est migre en place, sans perdre
une ligne.
Lance directement, le module porte sa propre gate :
uv run adws/adw_modules/tracer.py # auto-test, zero token
uv run adws/adw_modules/tracer.py --last # relire le dernier run
"""
from __future__ import annotations
import json
import sqlite3
import sys
import tempfile
import uuid
from datetime import datetime, timezone
from pathlib import Path
DB_PATH = Path("adws/adw_data/factory.db")
JSONL_DIR = Path("adws/adw_data/traces")
# Le schema du journal : trois tables, une par question.
# runs — qui a tourne, quand, verdict, cout, jetons ;
# phases — le deroule d'un run — et, depuis le chapitre 20, ce que
# chaque phase a coute en dollars et en jetons ;
# events — le grain fin : tout ce qui s'est passe, date, type.
# Regle de la maison, gravee en SQL : status vaut 'fail' par defaut —
# le vert se merite, un crash ne fabrique jamais un faux succes.
SCHEMA = """
CREATE TABLE IF NOT EXISTS runs (
adw_id TEXT PRIMARY KEY,
adw_name TEXT,
request TEXT,
status TEXT DEFAULT 'fail',
started_at TEXT,
ended_at TEXT,
cost_usd REAL DEFAULT 0,
total_tokens INTEGER DEFAULT 0
);
CREATE TABLE IF NOT EXISTS phases (
phase_id TEXT PRIMARY KEY,
adw_id TEXT REFERENCES runs,
seq INTEGER,
name TEXT,
kind TEXT,
status TEXT DEFAULT 'fail',
attempt INTEGER DEFAULT 0,
retries INTEGER DEFAULT 0,
error TEXT,
started_at TEXT,
ended_at TEXT,
cost_usd REAL DEFAULT 0,
tokens INTEGER DEFAULT 0
);
CREATE TABLE IF NOT EXISTS events (
event_id TEXT PRIMARY KEY,
adw_id TEXT REFERENCES runs,
phase_id TEXT REFERENCES phases,
type TEXT,
name TEXT,
payload_json TEXT,
ts TEXT
);
"""
# La migration en place d'un journal anterieur au chapitre 20 :
# CREATE TABLE IF NOT EXISTS n'ajoute pas de colonne a une table existante,
# alors chaque colonne nouvelle s'ajoute une a une — et « colonne deja
# presente » n'est pas une erreur, c'est un journal deja au niveau.
MIGRATIONS = (
"ALTER TABLE runs ADD COLUMN total_tokens INTEGER DEFAULT 0",
"ALTER TABLE phases ADD COLUMN cost_usd REAL DEFAULT 0",
"ALTER TABLE phases ADD COLUMN tokens INTEGER DEFAULT 0",
)
def now_iso() -> str:
"""L'heure UTC, milliseconde comprise : triable en SQL comme du texte."""
return datetime.now(timezone.utc).isoformat(timespec="milliseconds")
class Tracer:
"""Ecrit chaque evenement DEUX fois : la ligne JSONL, puis le miroir SQLite."""
def __init__(self, db_path: Path = DB_PATH, jsonl_dir: Path = JSONL_DIR):
db_path = Path(db_path)
db_path.parent.mkdir(parents=True, exist_ok=True)
self.jsonl_dir = Path(jsonl_dir)
self.jsonl_dir.mkdir(parents=True, exist_ok=True)
# isolation_level=None : autocommit — rien ne reste en l'air si le
# process meurt, chaque evenement ecrit est un evenement garde.
self.conn = sqlite3.connect(db_path, isolation_level=None)
# WAL : lire la base pendant que l'usine ecrit.
self.conn.execute("PRAGMA journal_mode=WAL;")
# Deux runs concurrents patientent au lieu d'echouer.
self.conn.execute("PRAGMA busy_timeout=5000;")
self.conn.executescript(SCHEMA)
for migration in MIGRATIONS:
try:
self.conn.execute(migration)
except sqlite3.OperationalError:
pass # colonne deja presente : rien a migrer
# ── evenements : le grain fin ────────────────────────────────────────
def event(self, adw_id: str, type_: str, name: str,
phase_id: str | None = None, payload: dict | None = None) -> str:
event_id = f"evt_{uuid.uuid4().hex[:12]}"
line = {"event_id": event_id, "ts": now_iso(), "adw_id": adw_id,
"phase_id": phase_id, "type": type_, "name": name,
"payload": payload or {}}
# 1. Le brut d'abord : la ligne JSONL est l'enregistrement de reference.
with (self.jsonl_dir / f"{adw_id}.jsonl").open("a", encoding="utf-8") as f:
f.write(json.dumps(line, ensure_ascii=False) + "\n")
# 2. Le miroir ensuite : la meme information, requetable en SQL.
self.conn.execute(
"INSERT INTO events (event_id, adw_id, phase_id, type, name,"
" payload_json, ts) VALUES (?,?,?,?,?,?,?)",
(event_id, adw_id, phase_id, type_, name,
json.dumps(payload or {}, ensure_ascii=False), line["ts"]))
return event_id
# ── runs ─────────────────────────────────────────────────────────────
def run_start(self, adw_id: str, adw_name: str, request: str) -> None:
self.conn.execute(
"INSERT INTO runs (adw_id, adw_name, request, status, started_at)"
" VALUES (?,?,?,?,?) ON CONFLICT(adw_id) DO NOTHING",
(adw_id, adw_name, request, "running", now_iso()))
self.event(adw_id, "run_start", adw_name, payload={"request": request})
def run_finish(self, adw_id: str, ok: bool, cost_usd: float,
tokens: int = 0) -> None:
# Un run tue net ne passe jamais ici : il reste 'running' sans
# ended_at — c'est la signature requetable d'un run interrompu.
status = "success" if ok else "fail"
self.conn.execute(
"UPDATE runs SET status=?, ended_at=?, cost_usd=?, total_tokens=?"
" WHERE adw_id=?",
(status, now_iso(), round(cost_usd, 4), tokens, adw_id))
self.event(adw_id, "run_end", status,
payload={"cost_usd": round(cost_usd, 4), "tokens": tokens})
# ── phases ───────────────────────────────────────────────────────────
def phase_start(self, adw_id: str, seq: int, name: str, kind: str,
retries: int = 0) -> str:
"""Declare une phase AVANT sa premiere tentative ; rend son phase_id."""
phase_id = f"{adw_id}-{seq:02d}"
self.conn.execute(
"INSERT INTO phases (phase_id, adw_id, seq, name, kind, retries,"
" started_at) VALUES (?,?,?,?,?,?,?) ON CONFLICT(phase_id) DO NOTHING",
(phase_id, adw_id, seq, name, kind, retries, now_iso()))
self.event(adw_id, "phase_start", name, phase_id=phase_id,
payload={"seq": seq, "kind": kind})
return phase_id
def phase_attempt(self, phase_id: str, adw_id: str, name: str,
attempt: int, ok: bool, error: str = "",
cost_usd: float = 0.0, tokens: int = 0) -> None:
"""Le verdict d'une tentative tombe TOUT DE SUITE, jamais en fin de run.
La ligne de phase garde le DERNIER etat connu — mais le cout et les
jetons S'ACCUMULENT : une phase verte a la tentative 3 porte le prix
de ses trois essais, pas celui du dernier.
"""
status = "success" if ok else "fail"
self.conn.execute(
"UPDATE phases SET status=?, attempt=?, error=?, ended_at=?,"
" cost_usd=cost_usd+?, tokens=tokens+? WHERE phase_id=?",
(status, attempt, error, now_iso(), round(cost_usd, 4), tokens,
phase_id))
payload = {"attempt": attempt}
if error:
payload["error"] = error
if cost_usd:
payload["cost_usd"] = round(cost_usd, 4)
if tokens:
payload["tokens"] = tokens
self.event(adw_id, "phase_ok" if ok else "phase_fail", name,
phase_id=phase_id, payload=payload)
# ── lecture : le dernier run, sans client sqlite3 ────────────────────────
def last_run(db_path: Path = DB_PATH) -> int:
"""Relit le dernier run trace : verdict, phases, couts, volume d'evenements."""
if not Path(db_path).is_file():
print("aucun journal — lancez un ADW, puis revenez", file=sys.stderr)
return 1
conn = sqlite3.connect(db_path)
row = conn.execute(
"SELECT adw_id, adw_name, request, status, cost_usd, started_at"
" FROM runs ORDER BY started_at DESC LIMIT 1").fetchone()
if row is None:
print("journal vide — lancez un ADW, puis revenez", file=sys.stderr)
return 1
adw_id, adw_name, request, status, cost, started = row
print(f"run {adw_id} ({adw_name}) — {status} — ~{cost:.4f} $ — {started}")
print(f" demande : {request or '—'}")
for seq, name, kind, pstatus, attempt, pcost in conn.execute(
"SELECT seq, name, kind, status, attempt, cost_usd FROM phases"
" WHERE adw_id=? ORDER BY seq", (adw_id,)):
print(f" {seq:02d} {name:<22} {kind:<6} {pstatus:<8}"
f" tentative {attempt} ~{pcost or 0:.4f} $")
total = conn.execute("SELECT COUNT(*) FROM events WHERE adw_id=?",
(adw_id,)).fetchone()[0]
print(f" {total} evenements dans le journal")
return 0
def selftest() -> int:
"""La gate du module : l'accumulation verifiee sur deux tentatives, zero token."""
with tempfile.TemporaryDirectory() as tmp:
tracer = Tracer(db_path=Path(tmp) / "test.db",
jsonl_dir=Path(tmp) / "traces")
tracer.run_start("selftest", "tracer_selftest", "auto-test du tracer")
phase_id = tracer.phase_start("selftest", 1, "demo", "agent", retries=1)
# Deux tentatives payees : la ligne de phase doit porter la SOMME.
tracer.phase_attempt(phase_id, "selftest", "demo", 1, ok=False,
error="enveloppe invalide", cost_usd=0.02, tokens=800)
tracer.phase_attempt(phase_id, "selftest", "demo", 2, ok=True,
cost_usd=0.03, tokens=1400)
tracer.run_finish("selftest", ok=True, cost_usd=0.05, tokens=2200)
cost, tokens, status = tracer.conn.execute(
"SELECT cost_usd, tokens, status FROM phases WHERE phase_id=?",
(phase_id,)).fetchone()
events = tracer.conn.execute("SELECT COUNT(*) FROM events").fetchone()[0]
lines = (Path(tmp) / "traces" / "selftest.jsonl").read_text(
encoding="utf-8").strip().splitlines()
ok = (abs(cost - 0.05) < 1e-9 and tokens == 2200
and status == "success" and events == 5 and len(lines) == events)
verdict = "OK" if ok else "KO"
print(f"tracer {verdict} — phase a {cost:.2f} $ / {tokens} jetons"
f" sur 2 tentatives, {events} evenements, statut {status}")
return 0 if ok else 1
if __name__ == "__main__":
if "--last" in sys.argv:
raise SystemExit(last_run())
raise SystemExit(selftest())
Pièce — adws/adw_modules/runner.py
Cette version remplace celle du chapitre 18. Deux lignes encadrent désormais chaque
tentative pour en mesurer le delta de coût et de jetons, et le Run gagne son compteur
tokens. Tout le reste est inchangé : PhaseFailure, PhaseSpec, l’API de Run, les
messages stderr et la ligne de bilan du banc. Vos ADW des chapitres 8 à 19 tournent tels quels.
"""runner — le squelette de l'usine : phases, sequencement, retries, traces.
Un ADW declare ses phases ; le runner les execute dans l'ordre, mesure,
retente les phases agent en session vivante, et rend un code retour.
Le succes se merite : toute PhaseFailure marque la tentative en echec,
et un run n'est vert que si toutes ses phases le sont.
Version chapitre 20 : le grain comptable. Le runner mesure en delta ce
que chaque tentative ajoute aux compteurs du run (cost_usd, tokens) et
l'ecrit au tracer avec le verdict. L'API ne bouge pas : vos ADW des
chapitres 8 a 19 tournent tels quels.
"""
from __future__ import annotations
import sys
import time
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any, Callable
from .tracer import Tracer
class PhaseFailure(Exception):
"""L'echec motive d'une phase — le runner decide s'il retente."""
@dataclass(frozen=True)
class PhaseSpec:
"""Ce qu'un ADW declare : un nom, un cote de la couture, une action."""
name: str
kind: str # "agent" ou "code"
action: Callable[["Run", int], Any] # (run, tentative) -> resultat
retries: int = 0 # phases agent : reprises en session vivante
@dataclass
class Run:
"""L'etat partage d'un run : resultats des phases, sessions, cout, jetons."""
adw_id: str
results: dict[str, Any] = field(default_factory=dict)
sessions: dict[str, str] = field(default_factory=dict) # phase -> session_id
cost_usd: float = 0.0
# Le compteur de jetons voyage comme le cout : une action qui releve
# l'usage de son harnais le credite (run.tokens += ...) ; une action
# qui ne releve rien laisse zero — jamais un chiffre invente.
tokens: int = 0
tracer: Tracer | None = None # injectable pour les tests ; None = journal standard
def __post_init__(self) -> None:
if self.tracer is None:
self.tracer = Tracer()
def execute(self, phases: list[PhaseSpec]) -> int:
"""Sequence les phases declarees. Arret a la premiere phase en echec definitif."""
if not phases:
print("aucune phase declaree — un ADW vide n'est pas un ADW", file=sys.stderr)
return 1
# Le nom de l'ADW et la demande sont releves sur la ligne de commande :
# zero changement dans vos scripts, et la trace sait deja qui tourne.
adw_name = Path(sys.argv[0]).stem if sys.argv and sys.argv[0] else ""
request = " ".join(sys.argv[1:])[:500]
self.tracer.run_start(self.adw_id, adw_name, request)
for seq, spec in enumerate(phases, start=1):
if spec.kind not in ("agent", "code"):
print(f"[{self.adw_id}] {spec.name} : kind inconnu {spec.kind!r}",
file=sys.stderr)
self.tracer.run_finish(self.adw_id, ok=False,
cost_usd=self.cost_usd, tokens=self.tokens)
return 1
if not self._run_phase(seq, spec):
print(f"[{self.adw_id}] ECHEC en phase {spec.name} — arret du run",
file=sys.stderr)
self.tracer.run_finish(self.adw_id, ok=False,
cost_usd=self.cost_usd, tokens=self.tokens)
return 1
# La ligne de bilan reste l'API de l'oeil et du banc (ch. 17) —
# la trace s'ajoute, elle ne retire rien.
print(f"[{self.adw_id}] run vert — cout total ~{self.cost_usd:.4f} $",
file=sys.stderr)
self.tracer.run_finish(self.adw_id, ok=True,
cost_usd=self.cost_usd, tokens=self.tokens)
return 0
def _run_phase(self, seq: int, spec: PhaseSpec) -> bool:
# Retenter une phase code n'a pas de sens : meme entree, meme sortie.
# Seules les phases agent ont droit aux reprises — en session vivante.
attempts = 1 + (spec.retries if spec.kind == "agent" else 0)
phase_id = self.tracer.phase_start(self.adw_id, seq, spec.name,
spec.kind, retries=attempts - 1)
for attempt in range(attempts):
clock = time.monotonic()
# Le grain comptable (ch. 20) : ce que CETTE tentative ajoute aux
# compteurs du run — mesure en delta, vos actions ne changent pas.
cost_before, tokens_before = self.cost_usd, self.tokens
label = f"{spec.name} ({spec.kind}, tentative {attempt + 1}/{attempts})"
try:
# L'action recoit le run (etat partage) et le numero de tentative :
# a la tentative 1, une phase agent envoie la demande ; ensuite,
# elle envoie la correction dans la MEME session.
self.results[spec.name] = spec.action(self, attempt)
except PhaseFailure as error:
print(f"[{self.adw_id}] {label} : echec — {error}", file=sys.stderr)
self.tracer.phase_attempt(phase_id, self.adw_id, spec.name,
attempt + 1, ok=False, error=str(error),
cost_usd=self.cost_usd - cost_before,
tokens=self.tokens - tokens_before)
continue
duration = time.monotonic() - clock
print(f"[{self.adw_id}] {label} : OK en {duration:.1f} s", file=sys.stderr)
self.tracer.phase_attempt(phase_id, self.adw_id, spec.name,
attempt + 1, ok=True,
cost_usd=self.cost_usd - cost_before,
tokens=self.tokens - tokens_before)
return True
return False
Pièce — adws/obs_export.py
Le pont de la salle de contrôle. Bibliothèque standard uniquement : il lit le journal et le traduit en spans OTLP JSON vers un collector, ou en exposition Prometheus sur la sortie standard, sans jamais écrire dedans. Son auto-test vérifie la traduction sur une base jetable, sans réseau. Entièrement côté déterministe.
# /// script
# requires-python = ">=3.11"
# ///
"""obs_export — le pont : la trace de l'usine traduite pour OTel et Prometheus.
Aucune instrumentation : le pont est une LECTURE de plus sur le journal
du chapitre 18. Le chemin reste usine -> sqlite -> lecture ; ce module
traduit apres coup des runs deja payes — un collector eteint ne coute
qu'un export rate, jamais un run.
uv run adws/obs_export.py --prom # metriques Prometheus (stdout)
uv run adws/obs_export.py --otlp # spans OTLP -> localhost:4318
uv run adws/obs_export.py --otlp http://collector:4318 # autre collector
uv run adws/obs_export.py --selftest # la gate : zero reseau, zero token
"""
from __future__ import annotations
import argparse
import hashlib
import io
import json
import sqlite3
import sys
import tempfile
import urllib.error
import urllib.request
from datetime import datetime
from pathlib import Path
DB_PATH = Path("adws/adw_data/factory.db")
DEFAULT_ENDPOINT = "http://localhost:4318" # OTLP/HTTP : le port standard d'un collector
RUNS_LIMIT = 10 # les derniers runs exportes
# OTLP : 0 = UNSET, 1 = OK, 2 = ERROR — un run encore 'running' reste UNSET.
STATUS = {"success": 1, "fail": 2}
def connect(db_path: Path = DB_PATH) -> sqlite3.Connection:
if not Path(db_path).is_file():
sys.exit("aucun journal — lancez un ADW, puis revenez")
conn = sqlite3.connect(db_path)
conn.execute("PRAGMA busy_timeout=5000;")
return conn
def _ns(ts: str | None) -> str:
"""ISO -> nanosecondes Unix, en chaine : le format des horodatages OTLP."""
if not ts:
return "0"
return str(int(datetime.fromisoformat(ts).timestamp() * 1_000_000_000))
def _trace_id(adw_id: str) -> str:
"""32 hexas DERIVES du run : exporter deux fois le meme run produit la
meme trace — pas de doublons chez le backend, l'export est rejouable."""
return hashlib.sha256(f"plume-factory:{adw_id}".encode()).hexdigest()[:32]
def _span_id(key: str) -> str:
"""16 hexas derives de la phase — meme logique, meme garantie."""
return hashlib.sha256(key.encode()).hexdigest()[:16]
def _attrs(pairs: dict) -> list[dict]:
"""Les attributs OTLP : chaque valeur declare son type dans le JSON —
et les entiers 64 bits voyagent en chaine, regle du format."""
out = []
for key, value in pairs.items():
if isinstance(value, float):
out.append({"key": key, "value": {"doubleValue": value}})
elif isinstance(value, int):
out.append({"key": key, "value": {"intValue": str(value)}})
else:
out.append({"key": key, "value": {"stringValue": str(value)}})
return out
def build_payload(conn: sqlite3.Connection,
limit: int = RUNS_LIMIT) -> tuple[dict, int]:
"""Les derniers runs traduits en traces OTLP : run = span racine,
phase = span enfant. Tout vient de runs et phases — rien d'invente."""
spans = []
runs = conn.execute(
"SELECT adw_id, adw_name, status, cost_usd, total_tokens,"
" started_at, ended_at FROM runs ORDER BY started_at DESC LIMIT ?",
(limit,)).fetchall()
for adw_id, name, status, cost, tokens, started, ended in runs:
trace_id = _trace_id(adw_id)
root_id = _span_id(f"{adw_id}:run")
spans.append({
"traceId": trace_id, "spanId": root_id,
"name": name or "run", "kind": 1, # SPAN_KIND_INTERNAL
"startTimeUnixNano": _ns(started),
"endTimeUnixNano": _ns(ended or started),
"attributes": _attrs({"factory.adw_id": adw_id,
"factory.cost_usd": float(cost or 0),
"factory.tokens": int(tokens or 0)}),
"status": {"code": STATUS.get(status, 0)},
})
for (phase_id, seq, pname, kind, pstatus, attempt,
pcost, ptokens, p_start, p_end) in conn.execute(
"SELECT phase_id, seq, name, kind, status, attempt, cost_usd,"
" tokens, started_at, ended_at FROM phases WHERE adw_id=?"
" ORDER BY seq", (adw_id,)):
spans.append({
"traceId": trace_id, "spanId": _span_id(phase_id),
"parentSpanId": root_id,
"name": f"{seq:02d} {pname}", "kind": 1,
"startTimeUnixNano": _ns(p_start),
"endTimeUnixNano": _ns(p_end or p_start),
"attributes": _attrs({"factory.kind": kind,
"factory.attempt": int(attempt or 0),
"factory.cost_usd": float(pcost or 0),
"factory.tokens": int(ptokens or 0)}),
"status": {"code": STATUS.get(pstatus, 0)},
})
payload = {"resourceSpans": [{
"resource": {"attributes": _attrs({"service.name": "plume-factory"})},
"scopeSpans": [{"scope": {"name": "obs_export"}, "spans": spans}],
}]}
return payload, len(runs)
def otlp(conn: sqlite3.Connection, endpoint: str) -> int:
"""POSTe les spans au collector — OTLP/HTTP, variante JSON officielle."""
payload, nruns = build_payload(conn)
if not nruns:
sys.exit("journal vide — lancez un ADW, puis revenez")
url = endpoint.rstrip("/") + "/v1/traces"
body = json.dumps(payload).encode("utf-8")
request = urllib.request.Request(
url, data=body, headers={"Content-Type": "application/json"})
try:
with urllib.request.urlopen(request, timeout=10) as response:
print(f"{nruns} runs exportes vers {url} — HTTP {response.status}")
except (urllib.error.URLError, OSError) as error:
# Le sens de la dependance, verifie ici : un collector eteint ne
# coute qu'un export rate — le journal, lui, n'a rien perdu.
sys.exit(f"collector injoignable ({error}) — relancez l'export"
" quand il sera debout : le journal n'a rien perdu")
return 0
def prom(conn: sqlite3.Connection, out=sys.stdout) -> int:
"""L'exposition Prometheus : les agregats du journal, en texte lisible.
A rediriger vers un fichier ou a servir — Prometheus vient les chercher."""
w = out.write
w("# HELP factory_runs_total runs traces, par ADW et verdict\n")
w("# TYPE factory_runs_total counter\n")
for name, status, count in conn.execute(
"SELECT adw_name, status, COUNT(*) FROM runs GROUP BY 1, 2"):
w(f'factory_runs_total{{adw="{name}",status="{status}"}} {count}\n')
w("# HELP factory_cost_usd_total dollars depenses, par ADW\n")
w("# TYPE factory_cost_usd_total counter\n")
for name, cost in conn.execute(
"SELECT adw_name, ROUND(SUM(cost_usd), 4) FROM runs GROUP BY 1"):
w(f'factory_cost_usd_total{{adw="{name}"}} {cost or 0}\n')
w("# HELP factory_phase_cost_usd_total dollars par phase — ou va l'argent\n")
w("# TYPE factory_phase_cost_usd_total counter\n")
for name, cost in conn.execute(
"SELECT name, ROUND(SUM(cost_usd), 4) FROM phases GROUP BY 1"
" HAVING SUM(cost_usd) > 0 ORDER BY 2 DESC"):
w(f'factory_phase_cost_usd_total{{phase="{name}"}} {cost}\n')
w("# HELP factory_phase_seconds_total duree cumulee, par phase\n")
w("# TYPE factory_phase_seconds_total counter\n")
for name, seconds in conn.execute(
"SELECT name, ROUND(SUM((julianday(ended_at) - julianday(started_at))"
" * 86400), 1) FROM phases WHERE ended_at IS NOT NULL GROUP BY 1"):
w(f'factory_phase_seconds_total{{phase="{name}"}} {seconds}\n')
return 0
def selftest() -> int:
"""La gate du module : un journal jetable, la traduction verifiee — zero reseau."""
from adw_modules.tracer import Tracer
with tempfile.TemporaryDirectory() as tmp:
db = Path(tmp) / "test.db"
tracer = Tracer(db_path=db, jsonl_dir=Path(tmp) / "traces")
tracer.run_start("demo", "export_selftest", "run synthetique")
p1 = tracer.phase_start("demo", 1, "plan", "agent")
tracer.phase_attempt(p1, "demo", "plan", 1, ok=True,
cost_usd=0.05, tokens=2000)
p2 = tracer.phase_start("demo", 2, "gate_tests", "code")
tracer.phase_attempt(p2, "demo", "gate_tests", 1, ok=True)
tracer.run_finish("demo", ok=True, cost_usd=0.05, tokens=2000)
conn = sqlite3.connect(db)
payload, nruns = build_payload(conn)
spans = payload["resourceSpans"][0]["scopeSpans"][0]["spans"]
buffer = io.StringIO()
prom(conn, out=buffer)
text = buffer.getvalue()
ok = (nruns == 1 and len(spans) == 3
and all(len(s["traceId"]) == 32 and len(s["spanId"]) == 16
for s in spans)
and spans[0].get("parentSpanId") is None
and any(a["key"] == "factory.cost_usd"
for a in spans[1]["attributes"])
and 'factory_cost_usd_total{adw="export_selftest"} 0.05' in text)
verdict = "OK" if ok else "KO"
print(f"obs_export {verdict} — 1 trace, {len(spans)} spans,"
f" {len(text.splitlines())} lignes d'exposition prometheus")
return 0 if ok else 1
if __name__ == "__main__":
parser = argparse.ArgumentParser(
description="le pont OTel / Prometheus de l'usine")
parser.add_argument("--prom", action="store_true",
help="exposition Prometheus sur stdout")
parser.add_argument("--otlp", nargs="?", const=DEFAULT_ENDPOINT,
metavar="URL", help="POSTer les spans vers ce collector")
parser.add_argument("--selftest", action="store_true",
help="la gate du module")
args = parser.parse_args()
if args.selftest:
raise SystemExit(selftest())
conn = connect()
if args.prom:
raise SystemExit(prom(conn))
if args.otlp:
raise SystemExit(otlp(conn, args.otlp))
parser.print_help()
La gate du TP
Quatre commandes, une par ligne, depuis la racine de plume-factory :
uv run adws/adw_modules/tracer.py
uv run adws/obs_export.py --selftest
uv run adws/adw_prompt.py "Combien de fichiers dans apps/plume ? Reponds en une phrase."
uv run adws/obs_export.py --prom
Attendu : la première imprime tracer OK — phase a 0.05 $ / 2200 jetons sur 2 tentatives, 5 evenements, statut success, la deuxième obs_export OK — 1 trace, 3 spans, …, toutes deux
à zéro token. La troisième trace un vrai run, coût par phase compris, la quatrième imprime
l’exposition Prometheus, où factory_phase_cost_usd_total montre où va votre argent. En tout :
~quelques centimes, ~1 minute. Si un collector tourne chez vous,
uv run adws/obs_export.py --otlp en une ligne de plus. Sinon, la lecture à sec ci-dessus
suffit, rien d’autre n’est requis.