Observabilité Chapitre 20 / 42

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_usd depuis 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 t3 porte 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_at et ended_at datent 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 Run gagne un compteur tokens que 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

MesureD’où elle vientLa question qu’elle règle
Latence (started_at/ended_at)les horodatages du chapitre 18où vont les minutes
Coût (cost_usd)le delta du compteur, sommé sur les tentativesoù va l’argent
Jetons (tokens)le même delta, quand l’action relève l’usagece que pèse la conversation
Verdict et tentatives (status, attempt)le runner, depuis le chapitre 18ce 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 build qui retente, l’autre dans un plan trop 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/json est une variante officielle du protocole) : un POST sur localhost:4318/v1/traces, le port par défaut d’un collector, avec Content-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 le phase_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éevous, votre repo, votre terminall’équipe, ses dashboards existants
Répond àquoi, où, quand, combien, au grain fintendances, comparaisons, alertes
Dépendancesaucuneun collector ou un Prometheus debout
Source de véritéoui, toujoursnon : 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.


Quiz — teste tes connaissances
Observabilité 7 questions Objectif : 5/7 minimum
0/7
bonnes reponses
Objectif non atteint (minimum 5/7 requis).
Remonte relire la fiche memo en pretant attention aux points manques, puis cliquer sur « Recommencer » pour retenter.