Session vivante : reprise, RPC et fork
Une session pi qu'on reprend sans la perdre — registre écrit avant l'appel, reprise refusée hors de son dossier — puis l'autre porte de pi, le mode RPC, réservée au seul usage qui la justifie dans l'usine : un best-of-N par clone, N builds sur une amorce payée une fois.
Depuis le chapitre A6, c’est votre runner qui nomme la session pi avant de l’ouvrir, et depuis le
chapitre A7 une tentative coupée par le coupe-circuit repart dans la même conversation. Mais que
reste-t-il de cette session si le runner lui-même est tué net, par un Ctrl-C pendant un build ou une
boîte du module 6 éteinte ? Un identifiant en mémoire, disparu avec le processus. Et que se
passe-t-il si vous reprenez cette session depuis un autre dossier ? pi en crée une seconde, au même
identifiant, sans rien dire. À la fin de ce chapitre, chaque session laissera une trace avant
même que pi démarre, une reprise hors de son dossier sera refusée avec son motif, et vous saurez
tenir une session pi vivante (un seul processus, plusieurs tours, des clones) par l’autre
porte de pi, le mode RPC. Vous saurez surtout quand la prendre : pas pour les phases, qui
restent un processus par tentative, mais pour le best-of-N du chapitre 25, où N builds partent
d’une même amorce payée une fois. Trois pièces : harness.py (qui remplace la version du chapitre
A8), un client RPC minimal, harness_rpc.py, et l’ADW qui l’utilise, adw_bon.py.
Reprendre une session sans la perdre
L’idée en une phrase
Une session pi est un fichier JSONL nommé par le runner (--session-id, chapitre A6) que pi
crée ou reprend. Pour qu’une reprise ne perde jamais rien, le code déterministe déclare
la session dans un registre (ledger.jsonl) avant de lancer pi, et refuse de la reprendre
depuis un autre dossier que celui de sa naissance : deux gestes côté runner, aucun côté agent.
Points clés
--session-idcrée ou reprend, vérifié sur pi 0.85. Si un fichier de session porte cet identifiant dans le dossier--session-dir, pi le rouvre et charge son historique, sinon il en crée un (avec un avertissement sur la sortie d’erreur, que le runner ignore). Le rechargement ne ré-émet pas l’historique sur stdout : seuls les messages nouveaux produisent desmessage_end, donc aucun double comptage de jetons ni de coût à la reprise.- La reprise exige le même dossier. Avec un
--session-dirautre que celui par défaut, pi filtre les sessions sur le champcwdde leur en-tête. Relancé depuis un autre répertoire (un worktree, une boîte du module 6), il ne trouve rien et crée une seconde session au même identifiant. Deux fichiers, un seul id, une trace harnais coupée en deux. Le port refuse désormais cette reprise avec son motif (resume_requires_cwd) : c’est un échec de phase, session conservée, comme tout refus du port. - Déclarer avant d’appeler. L’uuid est choisi par le runner depuis A6, mais il manquait
un endroit où il survive au runner. Le registre
adws/adw_data/sessions/ledger.jsonlreçoit une ligne (identifiant, dossier de naissance, modèle, phase et ADW quand le runner les a posés dans l’environnementFACTORY_PHASE,FACTORY_ADW_ID, chapitre A7) avantsubprocess.run. Un runner tué net laisse cette ligne, et la cléphase_sessiondu chapitre A4, elle, reste écrite par le runner après la tentative. - Le registre sert aussi de mémoire de dossier.
session_originlit le registre, et à défaut l’en-tête du fichier de session lui-même (*_<id>.jsonl, première ligne) : une session née avant ce chapitre est protégée de la même façon. - Un processus par tentative reste la règle. La reprise n’est pas un processus qui attend :
c’est un nouveau
pi -pqui rouvre le même fichier. C’est ce qui rend le kill propre, la trace bornée et l’isolation gratuite, et c’est ce qu’il faudra peser au second sous-thème.
Exemple concret
Un build SDLC (chapitre 13) sur le workhorse z-ai/glm-5.3, relevé ce jour sur
openrouter.ai/models : 1,40 $ le million de jetons en entrée, 4,40 $ en sortie. Première
tentative : douze tours, une quarantaine de centimes, puis les gates tombent. La correction repart
dans la même session : pi recharge les douze tours depuis le fichier, n’en ré-émet aucun, et la
tentative 2 ne facture que ses propres tours. Maintenant, tuez le runner pendant la tentative 2.
Avant ce chapitre, run.sessions mourait avec lui : la tentative suivante repartait à froid, avec une
quarantaine de centimes de contexte à reconstruire. Avec le registre, une seule requête
(grep sur ledger.jsonl ou la commande de la gate) vous rend l’identifiant et son dossier de
naissance, et la reprise coûte ce qu’elle coûtait : quelques centimes de correction. Relancez
enfin la même reprise depuis un worktree : le port refuse en zéro jeton, avant tout appel,
au lieu de laisser pi ouvrir une session jumelle que personne n’aurait vue.
Quatre façons de continuer un travail
| Geste | Ce que pi fait | Coût d’entrée | Quand l’usine l’emploie |
|---|---|---|---|
Nouveau pi -p, id neuf | session vide, contexte à reconstruire | tout le contexte, à nouveau | première tentative d’une phase |
Nouveau pi -p, même id, même dossier | recharge le fichier, n’émet rien d’ancien | zéro (la relecture du préfixe se paie au prochain tour, en cache quand le fournisseur en a) | reprise en session vivante (ch. 8, A7) |
Nouveau pi -p, même id, autre dossier | seconde session, même id | un contexte vide, et une trace scindée | jamais — refusé par resume_requires_cwd |
Processus --mode rpc gardé ouvert | plusieurs tours, clones, sans relancer | zéro entre les tours | best-of-N par clone (sous-thème 2) |
Config — une ligne de registre, et le refus qu’elle permet
Voici ce que le port écrit avant chaque appel, et le verdict qu’il rend quand on lui demande de reprendre cette session depuis un autre dossier. Rien de tout cela ne traverse la couture : le modèle n’en sait rien.
{"session_id": "5b1c…", "cwd": "/home/vous/plume-factory", "model": "z-ai/glm-5.3",
"phase": "build", "adw_id": "9f2e4a10", "declared_at": "2026-09-05T08:41:12+00:00"}
# adws/adw_modules/harness.py — extrait : refuser AVANT d'appeler, zero jeton.
def resume_requires_cwd(session_id, cwd, ledger=LEDGER, session_dir=SESSION_DIR):
origin = session_origin(session_id, ledger, session_dir)
if origin is None:
return # inconnue : pi la creera — c'est une premiere tentative
if Path(origin).resolve() != Path(cwd).resolve():
raise HarnessError(f"reprise refusee : la session {session_id} est nee dans {origin}, "
f"demandee depuis {Path(cwd).resolve()} — pi en creerait une "
"seconde au meme id", session_id)
Piège courant : « une session pi, c’est un processus qui attend » est inexact. Une session est un fichier : le processus
pi -pmeurt à chaque tentative et le suivant rouvre le fichier. C’est précisément pour cela que la reprise est robuste : rien ne vit en mémoire qu’il faudrait garder en vie. Le processus long existe, c’est le mode RPC, mais il se mérite par un usage précis, pas par confort.
RPC et fork pour le best-of-N
L’idée en une phrase
Le mode RPC de pi (pi --mode rpc) tient un seul processus ouvert et lui parle en JSONL :
une commande par ligne sur stdin, une réponse par commande, les événements de session sur stdout.
Il permet de cloner une session après une amorce et d’y lancer N builds sur un préfixe
partagé. Le code garde le graphe, le protocole, le verdict (_judge, le même qu’au
chapitre A8) et le refus de toute demande d’interface, et l’agent reste un nœud borné dans chaque
clone.
Points clés
- Le protocole, vérifié dans la doc
rpc.mdde pi 0.85. Chaque commande est un objet JSON sur une ligne ({"id": "c1", "type": "prompt", "message": "…"}), pi répond par un objettype: responseportant le mêmeid, puis diffuse les événements du tour (agent_start,message_end,tool_execution_end,entry_appended, jusqu’àagent_settled, le signal que rien ne suivra automatiquement. Le séparateur est\nseul, et le client découpe là-dessus, jamais sur un lecteur de lignes générique. - Les commandes que l’usine emploie, et pas d’autres.
prompt,clone,switch_session,get_state(pour liresessionIdetsessionFile),set_model,abort,get_session_stats.cloneduplique la branche courante à sa position dans une session neuve et bascule dessus.forkremonte à un message utilisateur antérieur, dont l’usine n’a pas besoin, l’amorce est le point de départ voulu. - Un clone rejoue le cycle de vie, pas l’historique. À chaque
cloneouswitch_session, pi relie les extensions à la nouvelle session :session_startest ré-émis, le coupe-circuit du chapitre A7 remet ses compteurs à zéro, l’inventaire du chapitre A8 se redépose au tour suivant, et chaque bras est donc jugé avec les mêmes preuves qu’une phase ordinaire. Ce qui est fixé pour tout le processus : l’environnement (le schéma de la porte, les limites), donc une enveloppe par processus et une garde par processus. ask= refus, ici aussi. En RPC,ctx.hasUIvauttrue: un damage control réglé surask(chapitre A3) émet une demandeextension_ui_requestet bloque jusqu’à la réponse. Le client répondcancelledà toute demande de dialogue et la journalise : la décision du chapitre A3 tient sans qu’un humain soit dans la boucle.- Décision du livre :
-ppour les phases, RPC pour le best-of-N. Un processus par tentative donne l’isolation, le kill propre, une trace bornée et un environnement par phase. Un processus long donne des tours sans relance, des clones et unabortexterne, au prix d’un état à piloter. Le seul endroit de l’usine où le second gagne nettement est celui où l’on répète un même départ : N builds sur une amorce.
Exemple concret
Un best-of-3 sur une spec du planner, workhorse GLM 5.3 (prix ci-dessus). L’amorce (lire la
spec en entier, ouvrir les quatre ou cinq fichiers qu’elle nomme, lire les tests) prend six à
huit tours. Comme chaque tour renvoie tout le contexte, elle cumule quelques dizaines de milliers
de jetons : de cinq à dix centimes. Puis trois clones : chaque bras implémente en une dizaine
de tours, une vingtaine de centimes chacun, et rend son enveloppe par la porte typée. Les gates
tranchent, le patch est moissonné, l’arbre remis à zéro. À froid, avec trois adw_build séparés,
l’amorce aurait été payée trois fois, et surtout refaite trois fois, avec trois lectures
différentes de la même spec. Par clone, elle est payée une fois et les trois bras partent du
même état lu. L’économie brute est de l’ordre de deux amorces, dix à vingt centimes sur un
lot d’un peu moins d’un dollar, soit dix à vingt pour cent, davantage quand l’amorce est
longue. Sur les fournisseurs qui facturent la relecture de cache à une fraction du prix
d’entrée, le préfixe partagé pèse moins encore. Ce que le lot vous rend en plus, et qui n’a pas de
prix au million de jetons : trois implémentations comparables, parce qu’une seule variable a
bougé.
Un processus par tentative, ou un processus qui dure
| Critère | pi -p par phase (ch. A6) | pi --mode rpc (ce chapitre) |
|---|---|---|
| Isolation d’une tentative | totale : le processus meurt | partielle : l’état survit, le clone isole la session |
| Environnement (schéma, limites) | par phase | par processus — une enveloppe, une garde |
| Reprise après une coupure | rouvrir le fichier | le processus est déjà là ; abort externe possible |
| Départs répétés (best-of-N) | N contextes à reconstruire | une amorce, N clones |
| Ce que le runner doit tenir | rien entre deux appels | un tube stdin/stdout, des id, un close() |
Script — un tour RPC, lu et jugé comme une phase
Le cœur du client : une commande, sa réponse, puis le flux jusqu’à agent_settled, et ce flux
passe dans le même _read_pi_stream et le même _judge que le port. Le RPC change la façon
de parler à pi, pas la façon de le juger.
# adws/adw_modules/harness_rpc.py — extrait : le tour, puis le verdict du port.
def prompt(self, message: str) -> Turn:
clock = time.monotonic()
_, seen = self._command({"type": "prompt", "message": message}) # accepte, ou refuse
lines, deadline = list(seen), clock + self.request.timeout
while True:
event, line = self._next(deadline) # les demandes d'interface y sont refusees
lines.append(line)
if event.get("type") == SETTLED: # agent_settled : plus rien ne suivra
break
return Turn(reading=harness._read_pi_stream(lines), lines=lines,
seconds=time.monotonic() - clock)
def judge(self, turn: Turn, request=None) -> harness.HarnessResult:
return harness._judge(turn.reading, request or self.request, self.session_id, 0, "")
Côté Claude Code, il n’y a pas d’équivalent à poser : pas de mode RPC documenté, pas de clone de session, ce qui est un argument de plus pour pi comme nœud d’usine dès qu’on répète un départ.
Piège courant : « le RPC est plus rapide et moins cher, autant tout y passer » est inexact. Entre deux tours d’un même processus vous ne gagnez que le démarrage de pi, moins d’une seconde, et le contexte, lui, est renvoyé à chaque tour dans les deux modes. Ce que le RPC achète, c’est le clone : un point de départ partagé. Hors de ce cas, un processus long est un état de plus à tenir (un tube à ne pas laisser se remplir, un
close()à garantir, un environnement figé pour toute sa vie) sans rien rendre en échange.
Fil rouge — la pièce posée aujourd’hui
Sur le plan de l’usine, la pièce du jour touche le port harnais (chapitre 7, adaptateur v6)
et ouvre, à côté de lui, une seconde porte vers pi, harness_rpc.py, réservée à un seul
client, adw_bon.py, le best-of-N local qui complète celui du module 6 (chapitre 25 : N
boîtes, N rosters). La couture ne bouge pas : le runner Python possède le graphe, il choisit
l’identifiant, écrit le registre, ouvre et ferme le processus, clone, juge, moissonne, remet
l’arbre à zéro. L’agent n’est qu’un tour borné dans chaque clone, et son enveloppe passe par la
même porte typée qu’au chapitre A6. Déterministe : le registre écrit avant l’appel, le refus de
reprise hors dossier, le protocole JSONL et ses id, le refus des dialogues, le clonage, les gates
par bras, le classement. Délégué : l’amorce (lire) et chaque bras (implémenter). Coût d’usage :
zéro jeton pour le registre et le refus, et un best-of-3 sur le workhorse coûte un peu moins d’un
dollar, dont une amorce payée une fois au lieu de trois, soit dix à vingt pour cent de moins qu’à
froid, et trois bras comparables. Deux notes de continuité : le miroir bloquant de
factory-obs.ts à session_shutdown (chapitre A4) reste en place, et en RPC il ne tourne qu’à la
fermeture du processus, une fois pour tout le lot, et adw_bon.py relance l’ingestion côté runner
après avoir fermé pi (elle est idempotente, chapitre A4). L’extension d’injection de contexte
annoncée au chapitre A6 attend toujours son tour, le budget du jour étant pris par les trois pièces.
Travaux pratiques — la pièce du jour
Trois fichiers à poser dans plume-factory, qui devient, chapitre après chapitre, votre usine
logicielle agentique : le port qui déclare avant d’appeler et refuse une reprise hors dossier, le
client RPC minimal, et le best-of-N par clone. Prérequis : pi (chapitre A1), Bun (chapitre 2), un
dépôt git propre (l’ADW du jour remet l’arbre à zéro entre deux bras : il l’exige et le
vérifie, zéro jeton).
Pièce — adws/adw_modules/harness.py
Cette version remplace celle du chapitre A8. Tout ce qui existait reste : mêmes dataclasses,
même run(), même verdict _judge, même adaptateur Claude Code. S’ajoutent le registre des
sessions (LEDGER, declare_session, session_origin), le refus resume_requires_cwd, et la
factorisation des drapeaux pi (_pi_flags) que le client RPC réutilise tels quels : un seul
endroit pour la ligne de commande, deux modes. Le profil d’authentification du chapitre 17bis (auth:, route, coffre) est conservé.
"""harness — le port de l'usine vers ses agents.
Une frontiere, deux adaptateurs. Aucun script de l'usine n'invoque `pi` ou
`claude` directement : tout passe par run(). Une HarnessRequest entre, un
HarnessResult sort — quel que soit le harnais derriere la porte.
Version chapitre 17bis : le profil d'authentification par agent. Une
requete peut porter `auth`, les variables de credentials que le roster a
resolues pour cet agent (noms ET valeurs, jamais dans le prompt), et
`direct`, la route hors passerelle d'un profil de forfait ou d'API native.
Quand elle porte un profil, .env devient un COFFRE : le noeud ne recoit
plus tout ce que .env avait charge, seulement ce que son profil nomme — a
la place du .env global du ch. 15, meme preseance (.env puis environnement
reel). Sans `auth`, l'heritage du ch. 15 s'applique tel quel.
Version annexe A6 (v4 de l'adaptateur pi) : le runner ne fait plus confiance
au noeud pi, il le VERIFIE. Le socle .pi/ est charge a coup sur (--approve),
le prompt voyage par stdin, l'enveloppe arrive par l'outil terminal
report_phase (tool_execution_end) plutot que devinee dans la prose, les
jetons et le stopReason sont lus, le modele observe est compare au modele
demande, et l'identifiant de session est choisi AVANT l'appel — une erreur
le porte, la reprise ne perd plus sa session.
Version annexe A8 (v5) : le runner exige la PREUVE que le socle est charge
(inventaire depose dans le flux, refus fail-closed s'il manque ou s'il
declare un manquant). Le verdict est une fonction pure (_judge).
Version annexe A9 (v6) : la session survit au runner. Chaque session est
DECLAREE dans un registre (adws/adw_data/sessions/ledger.jsonl) avant que pi
demarre — identifiant, dossier de naissance, modele, phase, ADW — et une
reprise depuis un autre dossier est REFUSEE avant tout appel : pi, avec un
--session-dir non standard, filtre ses sessions sur le cwd de leur en-tete
et en creerait une seconde au meme id. Les drapeaux pi sont factorises
(_pi_flags) pour que harness_rpc.py, la session vivante du best-of-N,
parle exactement la meme langue. L'API de run() ne bouge pas : vos ADW des
chapitres 8 a 27 tournent tels quels.
"""
from __future__ import annotations
import json
import os
import shutil
import subprocess
import uuid
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
from . import envelopes
# Les sessions pi et les schemas de porte vivent dans adw_data/ — couvert par
# le .gitignore du ch. 1. Chemins ABSOLUS a l'appel : pi filtre ses sessions
# sur le cwd de leur en-tete, un chemin relatif devient ambigu (module 6).
SESSION_DIR = Path("adws/adw_data/sessions")
SCHEMA_DIR = Path("adws/adw_data/schemas")
# Le registre des sessions (A9) : une ligne JSON par session ouverte, ecrite
# AVANT l'appel — l'uuid, le dossier de naissance, le modele, la phase et
# l'ADW. C'est ce qui survit a un runner tue net, et ce que lit
# resume_requires_cwd pour refuser une reprise hors de son dossier.
LEDGER = SESSION_DIR / "ledger.jsonl"
PHASE_ENV = "FACTORY_PHASE" # poses par la fabrique agent_action (A7)
ADW_ENV = "FACTORY_ADW_ID"
# La route par defaut de l'usine : le fournisseur passerelle integre de pi.
# Une seule cle (OPENROUTER_API_KEY) sert tous les moteurs du roster.
# Passer par les API directes des fournisseurs : GATEWAY = "" — et a vous
# de fournir une cle par fournisseur dans l'environnement.
GATEWAY = "openrouter"
# La variable que lit .pi/extensions/factory-report.ts pour construire
# l'outil report_phase : le chemin du schema JSON ecrit par ce module.
SCHEMA_ENV = "FACTORY_ENVELOPE_SCHEMA"
# Le socle : les extensions pi que TOUT noeud d'usine doit avoir chargees.
# Le runner les nomme (FACTORY_EXPECTED), factory-inventory.ts les compare a
# ce qu'il trouve et depose l'inventaire dans le flux (INVENTORY_ENTRY). Le
# poste (footer, chaine) n'en fait pas partie : il ne se charge qu'en TUI.
SOCLE = ("damage-control", "factory-obs", "factory-report", "factory-guard",
"factory-inventory")
EXPECTED_ENV = "FACTORY_EXPECTED"
INVENTORY_ENTRY = "factory-inventory"
# Le dialecte Claude Code : outils avec majuscules, et pas d'outil ls ni find
# dedies — Bash et Glob les couvrent. L'adaptateur absorbe l'asymetrie.
CLAUDE_TOOLS = {"read": "Read", "bash": "Bash", "edit": "Edit", "write": "Write",
"grep": "Grep", "find": "Glob", "ls": "Bash"}
# L'echelle de reflexion de pi, traduite en budget de tokens pour Claude Code
# (variable d'environnement MAX_THINKING_TOKENS).
THINKING_TOKENS = {"off": 0, "minimal": 1024, "low": 4096, "medium": 8192,
"high": 16384, "xhigh": 24576, "max": 32000}
# Les cles que .env a fournies — et que l'environnement reel ne portait pas.
# C'est le COFFRE (17bis) : ce que le port retire du noeud quand la requete
# porte un profil, pour n'y remettre que ce que le profil nomme.
ENV_FILE_KEYS: set[str] = set()
def _load_env(path: str | Path = ".env") -> None:
"""Charge .env dans l'environnement du process — une fois, au chargement.
Ni pi ni claude ne lisent .env d'eux-memes : sans ce chargement, la cle
de la passerelle n'atteindrait jamais les agents. Une variable deja
presente dans l'environnement reel gagne toujours — un export de session
ou un secret de CI ne sont jamais ecrases — et n'entre pas dans le coffre :
elle est a vous, pas a .env.
"""
env_file = Path(path)
if not env_file.is_file():
return
for line in env_file.read_text(encoding="utf-8").splitlines():
line = line.strip()
if not line or line.startswith("#") or "=" not in line:
continue
key, _, value = line.partition("=")
key = key.strip()
if key not in os.environ:
os.environ[key] = value.strip().strip("'\"")
ENV_FILE_KEYS.add(key)
_load_env()
def node_environment(base: dict[str, str], extra_env: dict[str, str] | None,
auth: dict[str, str] | None, vault: set[str]) -> dict[str, str]:
"""L'environnement d'un noeud — pure, donc testable a sec.
Sans profil (auth=None) : l'heritage du ch. 15, tout .env compris. Avec
profil : .env est un coffre — chaque cle qu'il a fournie est retiree,
puis le profil depose les siennes, avec leurs valeurs. Ce que le port ou
l'ADW pose lui-meme pour la phase passe toujours.
"""
environment = dict(base)
if auth is not None:
for key in vault:
environment.pop(key, None)
environment.update(auth)
if extra_env:
environment.update(extra_env)
return environment
class HarnessError(RuntimeError):
"""Le harnais n'a pas rendu de reponse exploitable.
Depuis A6, l'erreur porte la session qu'elle a interrompue : l'uuid est
choisi avant l'appel, donc une tentative tuee par le mur de temps a quand
meme une session — la reprise la poursuit au lieu de repartir a froid.
"""
def __init__(self, message: str, session_id: str | None = None) -> None:
super().__init__(message)
self.session_id = session_id
@dataclass(frozen=True)
class HarnessRequest:
"""Ce que l'usine a le droit de demander a un agent — rien de plus."""
prompt: str
session_id: str | None = None # None = nouvelle session
cwd: str = "."
timeout: int = 600 # un agent muet ne bloque pas l'usine
model: str | None = None # l'id du REGISTRE, tel quel dans le roster
thinking: str | None = None # off..max — None = defaut du harnais
tools: tuple[str, ...] = () # allowlist du roster — () = outils par defaut
schema: dict | None = None # schema JSON de la porte — None = Envelope de base
socle: tuple[str, ...] = SOCLE # extensions exigees du noeud pi — () = aucune exigence
auth: dict[str, str] | None = None # le profil resolu (17bis) — None = l'heritage du ch. 15
direct: bool = False # route directe (17bis) : le modele est deja une route pi
@dataclass(frozen=True)
class HarnessResult:
"""Ce qu'un agent rend a l'usine — quel que soit le harnais."""
text: str # la derniere reponse de l'agent (ou l'enveloppe en JSON)
session_id: str # de quoi poursuivre la MEME session
cost_usd: float # 0.0 si le harnais ne rapporte pas le cout
returncode: int
tokens: int = 0 # jetons factures sur cet appel, 0 si non rapportes
stop_reason: str = "stop" # stop | length | toolUse | error | aborted
envelope: dict | None = None # l'enveloppe rendue par la porte typee, sinon None
def run(harness: str, request: HarnessRequest) -> HarnessResult:
"""L'unique porte d'entree vers les agents : choisit l'adaptateur, normalise."""
try:
adapter = ADAPTERS[harness]
except KeyError:
raise HarnessError(
f"harnais inconnu {harness!r} — disponibles : {sorted(ADAPTERS)}"
) from None
return adapter(request)
def _route(model: str, direct: bool = False) -> str:
"""L'id du registre devient une route pi : prefixe du fournisseur passerelle.
Le roster parle le langage du registre (z-ai/glm-5.3) — la meme chaine
que verifie la jauge du ch. 14. Le prefixe est un detail de dialecte :
il vit ici, jamais dans le YAML ni dans vos scripts. Sous un profil a
route directe (17bis), l'identifiant est deja une route pi (zai/glm-5.3,
kimi-coding/k3) et part tel quel — la route est dans le profil, pas dans
le nom : deepseek/... est un auteur OpenRouter ET un fournisseur natif.
"""
if direct or not GATEWAY or model.startswith(GATEWAY + "/"):
return model
return f"{GATEWAY}/{model}"
def _spawn(cmd: list[str], request: HarnessRequest,
extra_env: dict[str, str] | None = None,
stdin_text: str | None = None) -> subprocess.CompletedProcess[str]:
# Resoudre l'executable via le PATH : sous Windows, les harnais sont des
# shims (pi.cmd, claude.cmd) que CreateProcess ne trouve pas par leur nom
# court — shutil.which respecte PATHEXT et regle les deux mondes d'un coup.
executable = shutil.which(cmd[0])
if executable is None:
raise HarnessError(f"{cmd[0]!r} introuvable dans le PATH — "
"le harnais est-il installe ?")
environment = node_environment(dict(os.environ), extra_env, request.auth, ENV_FILE_KEYS)
# Deux modes d'entree, jamais d'entre-deux :
# - stdin_text=None : le prompt voyage dans argv, et stdin est ferme
# (DEVNULL) — un enfant qui herite de notre stdin peut attendre
# indefiniment une entree qui ne viendra jamais : echec silencieux,
# 0 % CPU, aucune sortie.
# - stdin_text : le prompt voyage par stdin, puis le tube est referme.
# Indispensable quand le harnais est un shim .cmd Windows : cmd.exe
# tronque un argument a la premiere nouvelle ligne, et les asks de
# l'usine (brief + mission + contrat) sont multi-lignes.
io = ({"input": stdin_text} if stdin_text is not None
else {"stdin": subprocess.DEVNULL})
# Encodage explicite : les harnais emettent de l'UTF-8, mais text=True
# seul decode avec la locale — cp1252 sous Windows, qui mutile tirets
# et accents. Vaut pour la sortie ET pour le prompt ecrit sur stdin.
try:
return subprocess.run([executable, *cmd[1:]], **io,
capture_output=True, text=True,
encoding="utf-8", errors="replace",
env=environment,
timeout=request.timeout, cwd=request.cwd)
except subprocess.TimeoutExpired:
# Le timeout aussi sort par la porte normalisee : une seule exception.
raise HarnessError(f"harnais muet apres {request.timeout} s : {cmd[0]}") from None
def _text_of(message: dict) -> str:
"""Concatene les blocs de texte d'un message pi."""
return "".join(part.get("text", "") for part in message.get("content", []) or []
if isinstance(part, dict) and part.get("type") == "text")
# ---------- le registre des sessions (A9) : declarer avant d'appeler ----------
def declare_session(session_id: str, cwd: str | Path, model: str | None = None,
ledger: Path = LEDGER, env: dict[str, str] | None = None) -> dict:
"""Ecrit la session dans le registre — AVANT que pi demarre.
Ce que l'ADW a pose dans l'environnement (phase, adw_id — fabrique
agent_action, A7) est releve au passage : un runner tue net laisse ainsi
de quoi retrouver la session de chaque phase, sans passer par la trace.
env : l'environnement du noeud s'il differe de celui du runner (RPC).
"""
source = os.environ if env is None else env
entry = {"session_id": session_id, "cwd": str(Path(cwd).resolve()), "model": model,
"phase": source.get(PHASE_ENV), "adw_id": source.get(ADW_ENV),
"declared_at": datetime.now(timezone.utc).isoformat(timespec="seconds")}
ledger.parent.mkdir(parents=True, exist_ok=True)
with ledger.open("a", encoding="utf-8") as f:
f.write(json.dumps(entry, ensure_ascii=False) + "\n")
return entry
def session_origin(session_id: str, ledger: Path = LEDGER,
session_dir: Path = SESSION_DIR) -> str | None:
"""Le dossier de naissance d'une session : le registre d'abord, l'en-tete pi sinon.
Une session nee avant ce chapitre n'est pas dans le registre, mais son
fichier (<horodatage>_<id>.jsonl) commence par un en-tete qui porte le
cwd : elle est protegee de la meme facon. None = session inconnue.
"""
origin = None
if ledger.is_file():
for line in ledger.read_text(encoding="utf-8").splitlines():
try:
entry = json.loads(line)
except json.JSONDecodeError:
continue
if entry.get("session_id") == session_id and entry.get("cwd"):
origin = entry["cwd"] # la derniere declaration gagne
if origin:
return origin
for path in sorted(session_dir.glob(f"*_{session_id}.jsonl")):
try:
with path.open(encoding="utf-8") as f:
header = json.loads(f.readline())
except (OSError, json.JSONDecodeError):
continue
if header.get("type") == "session" and header.get("cwd"):
return str(header["cwd"])
return None
def resume_requires_cwd(session_id: str, cwd: str | Path, ledger: Path = LEDGER,
session_dir: Path = SESSION_DIR) -> None:
"""Refuse — avant tout appel, zero jeton — une reprise hors de son dossier.
Verifie sur pi 0.85 : avec un --session-dir non standard, pi filtre les
sessions sur le cwd de leur en-tete ; depuis un autre dossier il ne
retrouve rien et CREE une seconde session au meme id, sans erreur.
"""
origin = session_origin(session_id, ledger, session_dir)
if origin is None:
return # inconnue : pi la creera — c'est une premiere tentative
if Path(origin).resolve() != Path(cwd).resolve():
raise HarnessError(f"reprise refusee : la session {session_id} est nee dans {origin}, "
f"demandee depuis {Path(cwd).resolve()} — pi en creerait une "
"seconde au meme id", session_id)
# ---------- l'adaptateur pi : des fonctions pures autour d'un seul subprocess ----------
def _write_schema(schema: dict) -> Path:
"""La porte typee se construit depuis le schema ecrit ICI, avant la phase.
Une seule source de verite, cote Python : envelopes.schema(). L'extension
ne connait aucun champ d'avance — elle lit ce fichier au chargement.
"""
SCHEMA_DIR.mkdir(parents=True, exist_ok=True)
target = SCHEMA_DIR / f"{schema.get('title', 'Envelope')}.json"
target.write_text(json.dumps(schema, indent=2, ensure_ascii=False), encoding="utf-8")
return target.resolve()
def _pi_flags(request: HarnessRequest, session_id: str) -> list[str]:
"""Les drapeaux communs aux deux modes (json et rpc) — purs, testables a sec."""
flags = [
# --approve : le socle .pi/ (damage control, trace, porte typee) est
# charge a coup sur, sans dependre d'un trust.json de poste. Cela
# vaut confiance au depot COURANT — jamais sur un depot inconnu.
"--approve",
"--session-id", session_id, "--session-dir", str(SESSION_DIR.resolve()),
# Un noeud d'usine ne lit ni les skills ni les templates du poste :
# ils s'adressent a un humain qui pilote. AGENTS.md reste charge —
# c'est du contexte gouverne par le depot.
"--no-skills", "--no-prompt-templates"]
if request.model:
flags += ["--model", _route(request.model, request.direct)]
if request.thinking:
flags += ["--thinking", request.thinking]
if request.tools:
# --tools est une allowlist qui retire AUSSI les outils d'extension
# absents de la liste : la porte doit y figurer, sinon elle disparait.
flags += ["--tools", ",".join((*request.tools, envelopes.REPORT_TOOL))]
return flags
def _pi_argv(request: HarnessRequest, session_id: str) -> list[str]:
"""La ligne de commande pi d'une phase : un tour headless, flux JSON, prompt sur stdin."""
# `--` ferme les options : plus rien de positionnel — le prompt arrive par
# stdin, quel que soit son premier caractere ou son nombre de lignes.
return ["pi", "-p", "--mode", "json", *_pi_flags(request, session_id), "--"]
@dataclass
class PiReading:
"""Ce que le runner retient du flux JSON de pi : quatre chiffres, une enveloppe."""
text: str = ""
cost_usd: float = 0.0
tokens: int = 0
stop_reason: str = "stop"
observed_model: str | None = None
envelope: dict | None = None
inventory: dict | None = None # ce que le socle declare de lui-meme (A8)
def _read_pi_stream(lines: list[str]) -> PiReading:
"""Lit le flux ligne a ligne — pure, donc testable a sec.
message_end (assistant) : texte, cout, jetons, stopReason, modele observe.
tool_execution_end de report_phase : l'enveloppe, deja validee par schema
cote pi. entry_appended de factory-inventory : la preuve du socle. Le
dernier texte gagne ; les couts et jetons s'additionnent. Le flux d'un
tour RPC (harness_rpc.py) se lit avec la meme fonction.
"""
reading = PiReading()
for line in lines:
try:
event = json.loads(line)
except json.JSONDecodeError:
continue
kind = event.get("type")
if kind == "message_end":
message = event.get("message") or {}
if message.get("role") != "assistant":
continue
reading.text = _text_of(message) or reading.text
usage = message.get("usage") or {}
reading.cost_usd += (usage.get("cost") or {}).get("total") or 0.0
reading.tokens += int(usage.get("totalTokens") or 0)
reading.stop_reason = message.get("stopReason") or reading.stop_reason
if reading.observed_model is None and message.get("model"):
reading.observed_model = f"{message.get('provider')}/{message.get('model')}"
elif (kind == "tool_execution_end"
and event.get("toolName") == envelopes.REPORT_TOOL
and not event.get("isError")):
details = (event.get("result") or {}).get("details")
if isinstance(details, dict):
reading.envelope = details
elif kind == "entry_appended":
entry = event.get("entry") or {}
if entry.get("customType") == INVENTORY_ENTRY and isinstance(entry.get("data"), dict):
reading.inventory = entry["data"]
return reading
def _judge(reading: PiReading, request: HarnessRequest, session_id: str,
returncode: int, stderr: str) -> HarnessResult:
"""Le verdict du runner sur un flux pi — pure, donc testable a sec.
Dans l'ordre : le socle d'abord (sans preuve, pas de phase — meme avec
une enveloppe) ; puis la porte typee ; puis les arrets qui ne sont pas
une reponse (error, aborted — pi rend 0 en --mode json, le code retour
ne dit rien) ; enfin le texte, pour le chemin de secours ; et le modele
observe, qui doit etre celui du roster.
"""
evidence = stderr.strip()[-400:]
if request.socle:
if reading.inventory is None:
raise HarnessError(
"socle pi non charge : aucun inventaire dans le flux — .pi/ non approuve, "
"factory-inventory.ts absent, ou pi n'a pas demarre de tour"
+ (f" ({evidence})" if evidence else ""), session_id)
missing = [name for name in request.socle
if name not in (reading.inventory.get("present") or [])]
if missing:
raise HarnessError(f"socle pi incomplet : manquants {missing} — "
f"pi {reading.inventory.get('pi_version', '?')} n'a charge que "
f"{reading.inventory.get('present')}", session_id)
if reading.envelope is not None:
text = json.dumps(reading.envelope, ensure_ascii=False)
elif reading.stop_reason in ("error", "aborted"):
raise HarnessError(f"pi s'est arrete sur {reading.stop_reason} : "
f"{evidence or reading.text[-400:]}", session_id)
elif returncode != 0 and not reading.text:
raise HarnessError(f"pi a rendu {returncode} : {evidence}", session_id)
else:
text = reading.text
# Le modele observe doit etre celui du roster : pi resout un motif par
# sous-chaine, et une facture sur le mauvais moteur n'est pas un run vert.
if request.model and reading.observed_model \
and reading.observed_model != _route(request.model, request.direct):
raise HarnessError(f"modele observe {reading.observed_model!r} "
f"≠ modele demande {_route(request.model, request.direct)!r}", session_id)
return HarnessResult(text=text, session_id=session_id, cost_usd=reading.cost_usd,
returncode=returncode, tokens=reading.tokens,
stop_reason=reading.stop_reason, envelope=reading.envelope)
def _run_pi(request: HarnessRequest) -> HarnessResult:
# pi : c'est VOUS qui nommez la session — et vous la nommez AVANT l'appel.
# Meme id + meme dossier = meme contexte ; une erreur porte cet id.
session_id = request.session_id or str(uuid.uuid4())
SESSION_DIR.mkdir(parents=True, exist_ok=True)
if request.session_id:
# Une reprise hors de son dossier de naissance ouvrirait une session
# jumelle : refus avant tout appel, zero jeton (A9).
resume_requires_cwd(session_id, request.cwd)
schema_path = _write_schema(request.schema or envelopes.schema(envelopes.Envelope))
# Ce que le noeud doit savoir de l'usine passe par l'environnement : le
# schema de la porte, et la liste du socle qu'il devra prouver.
extra_env = {SCHEMA_ENV: str(schema_path), EXPECTED_ENV: ",".join(request.socle)}
# Le registre AVANT le processus : un runner tue pendant la phase laisse
# cette ligne, et la reprise sait quelle session poursuivre (A9).
declare_session(session_id, request.cwd, request.model)
try:
proc = _spawn(_pi_argv(request, session_id), request,
extra_env=extra_env, stdin_text=request.prompt)
except HarnessError as error:
error.session_id = session_id
raise
reading = _read_pi_stream(proc.stdout.splitlines())
return _judge(reading, request, session_id, proc.returncode, proc.stderr)
def _claude_evidence(stdout: str, stderr: str) -> str:
"""Le motif d'un echec claude — pure. Le JSON de stdout d'abord, stderr ensuite.
En --output-format json, claude ecrit son erreur dans l'objet de stdout
(result, error) et n'envoie sur stderr que des avertissements (« Ignoring
N permissions.allow entries… ») : lire stderr d'abord masquerait la vraie
cause — une cle absente, un depot non approuve, un modele inconnu.
"""
try:
payload = json.loads(stdout)
for key in ("result", "error", "message"):
if isinstance(payload, dict) and payload.get(key):
return str(payload[key]).strip()[-400:]
except (json.JSONDecodeError, TypeError):
pass
lines = [line for line in stderr.strip().splitlines() if not line.startswith("Ignoring ")]
return ("\n".join(lines).strip() or stderr.strip() or stdout.strip())[-400:]
def _run_claude(request: HarnessRequest) -> HarnessResult:
# Claude Code : c'est LUI qui nomme la session. On la poursuit en rendant
# son session_id via --resume. Pas d'outil terminal type de ce cote :
# l'enveloppe reste une convention de texte, parse() la lit en secours.
#
# Regime par defaut : natif Anthropic — sa propre authentification, des
# modeles Anthropic. Le pointer sur la passerelle est possible (trois
# variables : ANTHROPIC_BASE_URL, ANTHROPIC_AUTH_TOKEN, et
# ANTHROPIC_API_KEY explicitement vide), mais la compatibilite n'est
# garantie que sur les modeles Anthropic : l'adaptateur polyglotte de
# l'usine reste pi, et ce choix-la appartient a votre environnement,
# pas a cet adaptateur.
cmd = ["claude", "-p", "--output-format", "json"]
if request.model:
# Le dialecte claude ignore le fournisseur : provider/id -> id.
cmd += ["--model", request.model.split("/", 1)[-1]]
if request.tools:
# Traduire puis dedoublonner en gardant l'ordre : bash et ls donnent
# tous deux Bash, inutile de le declarer deux fois.
allowed = list(dict.fromkeys(
CLAUDE_TOOLS[tool] for tool in request.tools if tool in CLAUDE_TOOLS))
cmd += ["--allowedTools", ",".join(allowed)]
if request.session_id:
cmd += ["--resume", request.session_id]
# L'echelle de reflexion devient un budget de tokens — l'asymetrie reste
# dans l'adaptateur, le roster n'en sait rien.
extra_env = ({"MAX_THINKING_TOKENS": str(THINKING_TOKENS[request.thinking])}
if request.thinking in THINKING_TOKENS else None)
# Le prompt part par stdin, PAS dans argv : c'est un mode documente de
# claude -p, et le seul qui survive aux shims .cmd de Windows.
proc = _spawn(cmd, request, extra_env, stdin_text=request.prompt)
if proc.returncode != 0:
# Le JSON de stdout d'abord, stderr ensuite : les avertissements de
# claude (« Ignoring … ») ne doivent pas masquer la vraie cause.
evidence = _claude_evidence(proc.stdout, proc.stderr)
raise HarnessError(f"claude a rendu {proc.returncode} : {evidence}",
request.session_id)
try:
payload = json.loads(proc.stdout)
except json.JSONDecodeError:
raise HarnessError("claude n'a pas rendu l'objet JSON attendu "
"(--output-format json)", request.session_id) from None
usage = payload.get("usage") or {}
return HarnessResult(text=str(payload.get("result", "")),
session_id=str(payload.get("session_id", "")),
cost_usd=float(payload.get("total_cost_usd") or 0.0),
returncode=proc.returncode,
tokens=int(usage.get("input_tokens") or 0)
+ int(usage.get("output_tokens") or 0),
stop_reason="error" if payload.get("is_error") else "stop")
# Le registre des adaptateurs. Un harnais de plus = une fonction + une ligne.
ADAPTERS = {"pi": _run_pi, "claude": _run_claude}
if __name__ == "__main__":
# La gate du module — zero token, sans pi : les fonctions pures et le registre.
# Lancer depuis la racine : uv run python -m adws.adw_modules.harness
import tempfile
request = HarnessRequest(prompt="- une ligne qui commence par un tiret",
model="z-ai/glm-5.3", tools=("read", "bash"))
argv = _pi_argv(request, "sess-1")
assert "--approve" in argv and argv[-1] == "--", argv
assert Path(argv[argv.index("--session-dir") + 1]).is_absolute()
assert argv[argv.index("--tools") + 1] == "read,bash,report_phase"
assert argv[argv.index("--model") + 1] == "openrouter/z-ai/glm-5.3"
assert request.prompt not in argv # le prompt part par stdin, jamais par argv
assert argv[4:-1] == _pi_flags(request, "sess-1") # un seul jeu de drapeaux, deux modes
def inventory_line(present: list[str]) -> str:
return json.dumps({"type": "entry_appended", "entry": {
"type": "custom", "customType": INVENTORY_ENTRY,
"data": {"pi_version": "0.85.0", "present": present,
"missing": [n for n in SOCLE if n not in present]}}})
assistant = json.dumps({"type": "message_end", "message": {
"role": "assistant", "provider": "openrouter", "model": "z-ai/glm-5.3",
"content": [{"type": "text", "text": "je lis"}], "stopReason": "toolUse",
"usage": {"totalTokens": 1200, "cost": {"total": 0.002}}}})
envelope = json.dumps({"type": "tool_execution_end", "toolName": "report_phase",
"isError": False,
"result": {"details": {"status": "success", "summary": "fini"}}})
# 1. Le socle au complet : l'enveloppe passe, les chiffres sont lus.
reading = _read_pi_stream([inventory_line(list(SOCLE)), assistant, envelope, "pas du JSON"])
assert reading.inventory and reading.inventory["missing"] == []
result = _judge(reading, request, "sess-1", 0, "")
assert result.envelope == {"status": "success", "summary": "fini"}, result
assert result.tokens == 1200 and result.stop_reason == "toolUse"
assert reading.observed_model == _route(request.model, request.direct)
# 2. Pas d'inventaire : refus, MEME avec une enveloppe valide — fail-closed.
try:
_judge(_read_pi_stream([assistant, envelope]), request, "sess-2", 0, "")
raise AssertionError("un flux sans inventaire doit etre refuse")
except HarnessError as error:
assert "socle pi non charge" in str(error) and error.session_id == "sess-2"
# 3. Un manquant (factory-guard renomme) : refus motive, session conservee.
partial = [n for n in SOCLE if n != "factory-guard"]
try:
_judge(_read_pi_stream([inventory_line(partial), assistant, envelope]), request, "sess-3", 0, "")
raise AssertionError("un socle incomplet doit etre refuse")
except HarnessError as error:
assert "manquants ['factory-guard']" in str(error) and error.session_id == "sess-3"
# 4. socle=() : aucune exigence — le chemin d'un run « a sec » ou d'un harnais sans socle.
free = HarnessRequest(prompt="ping", socle=())
assert _judge(_read_pi_stream([assistant]), free, "sess-4", 0, "").text == "je lis"
# 5. Un arret aborted (coupe-circuit A7) reste une HarnessError porteuse de session.
aborted = _read_pi_stream([inventory_line(list(SOCLE)), json.dumps({"type": "message_end", "message": {
"role": "assistant", "content": [], "stopReason": "aborted", "usage": {}}})])
try:
_judge(aborted, request, "sess-5", 0, "")
raise AssertionError("aborted doit lever")
except HarnessError as error:
assert "aborted" in str(error) and error.session_id == "sess-5"
# 6. Le registre (A9) : declare avant l'appel, retrouve, et une reprise
# hors dossier refusee — dans un dossier temporaire, jamais le vrai registre.
with tempfile.TemporaryDirectory() as tmp:
ledger, sessions = Path(tmp) / "ledger.jsonl", Path(tmp) / "sessions"
home, elsewhere = Path(tmp) / "usine", Path(tmp) / "worktree"
home.mkdir(), elsewhere.mkdir(), sessions.mkdir()
os.environ[PHASE_ENV], os.environ[ADW_ENV] = "build", "9f2e4a10"
entry = declare_session("sess-6", home, "z-ai/glm-5.3", ledger=ledger)
del os.environ[PHASE_ENV], os.environ[ADW_ENV]
assert entry["phase"] == "build" and entry["adw_id"] == "9f2e4a10"
assert session_origin("sess-6", ledger, sessions) == str(home.resolve())
resume_requires_cwd("sess-6", home, ledger, sessions) # meme dossier : passe
resume_requires_cwd("sess-inconnue", elsewhere, ledger, sessions) # inconnue : passe
try:
resume_requires_cwd("sess-6", elsewhere, ledger, sessions)
raise AssertionError("une reprise hors dossier doit etre refusee")
except HarnessError as error:
assert "reprise refusee" in str(error) and error.session_id == "sess-6"
# Une session d'avant A9 : pas dans le registre, mais son en-tete pi suffit.
(sessions / "2026-09-05T08-00-00_sess-7.jsonl").write_text(
json.dumps({"type": "session", "version": 3, "id": "sess-7", "cwd": str(home)}) + "\n",
encoding="utf-8")
assert session_origin("sess-7", ledger, sessions) == str(home)
# La route directe (17bis) : un profil hors passerelle envoie l'identifiant tel quel.
assert _route("deepseek/deepseek-v4-flash-0731") == "openrouter/deepseek/deepseek-v4-flash-0731"
assert _route("deepseek/deepseek-v4-flash-0731", direct=True) == "deepseek/deepseek-v4-flash-0731"
assert _route("zai/glm-5.3", direct=True) == "zai/glm-5.3"
# Le coffre (17bis) : sans profil, tout .env passe ; avec profil, seul le profil passe —
# et ce que le port ou l'ADW pose lui-meme pour la phase passe toujours.
base = {"PATH": "/usr/bin", "OPENROUTER_API_KEY": "sk-or-env", "ZAI_API_KEY": "zai-env",
"TERM": "xterm"}
vault = {"OPENROUTER_API_KEY", "ZAI_API_KEY"}
legacy = node_environment(base, {"MAX_THINKING_TOKENS": "8192"}, None, vault)
assert legacy["OPENROUTER_API_KEY"] == "sk-or-env" and legacy["ZAI_API_KEY"] == "zai-env"
node = node_environment(base, {"MAX_THINKING_TOKENS": "8192"}, {"OPENROUTER_API_KEY": "sk-or-env"}, vault)
assert node["OPENROUTER_API_KEY"] == "sk-or-env" and "ZAI_API_KEY" not in node
assert node["PATH"] == "/usr/bin" and node["TERM"] == "xterm" and node["MAX_THINKING_TOKENS"] == "8192"
session = node_environment(base, None, {}, vault) # regime session : rien a injecter, coffre ferme
assert "OPENROUTER_API_KEY" not in session and "ZAI_API_KEY" not in session
exported = node_environment(base, None, {}, set()) # une variable de l'environnement reel reste a vous
assert exported["OPENROUTER_API_KEY"] == "sk-or-env"
print("harness OK — argv v4 (approve, stdin, --, report_phase) ; verdict v5 : socle "
f"prouve ({len(SOCLE)} extensions), refus sans inventaire, refus sur manquant, "
f"enveloppe par la porte ({result.tokens} jetons), arret 'aborted' detecte ; "
"registre v6 : declare avant l'appel, reprise hors dossier refusee"
" ; profil 17bis : route directe hors passerelle, coffre .env ferme sous profil, ouvert sans")
Pièce — adws/adw_modules/harness_rpc.py
Le client RPC minimal, nouveau. Une classe PiRpc qui ouvre pi --mode rpc avec les drapeaux
du port, parle le protocole (commandes corrélées par id, réponses, flux jusqu’à
agent_settled), refuse toute demande de dialogue, et rend chaque tour au même _judge que
les phases. Sept commandes de pi, pas une de plus : prompt, clone, switch_session,
get_state, set_model, abort, get_session_stats. Sa gate n’a pas besoin de pi : le module se
lance lui-même en faux pi RPC (--fake), même framing, même protocole, et se pilote de bout en
bout : une amorce, une demande de dialogue refusée, deux clones. Aucune dépendance.
"""harness_rpc — la session vivante : un processus pi, plusieurs tours, des clones (A9).
Le port harness.run() garde un processus par tentative : c'est l'isolation,
le kill propre, la trace bornee. Ce module ouvre l'autre porte de pi,
--mode rpc, pour le seul usage qui la justifie dans l'usine : le best-of-N
par clone (adw_bon.py) — une amorce payee une fois, N bras qui en partent.
Meme lecture du flux (_read_pi_stream), meme verdict (_judge) : le RPC
change la maniere de parler a pi, pas la maniere de le juger.
Protocole (pi 0.85, doc rpc.md) : une commande JSON par ligne sur stdin ;
une reponse {"type": "response", "command": ..., "success": ...} par
commande, correlee par "id" ; les evenements de session en JSONL sur
stdout, jusqu'a agent_settled. Le separateur est "\\n" seul. Toute demande
d'interface (extension_ui_request : select, confirm, input, editor) recoit
un refus (cancelled) et est journalisee : ask = refus en headless, la
decision du chapitre A3 tient sans humain dans la boucle.
Gate (zero token, sans pi) : uv run python -m adws.adw_modules.harness_rpc
— le module se lance lui-meme comme un faux pi RPC (--fake) et se pilote.
"""
from __future__ import annotations
import json
import os
import queue
import shutil
import subprocess
import sys
import threading
import time
import uuid
from dataclasses import dataclass, field, replace
from pathlib import Path
from . import envelopes, harness
from .harness import HarnessError, HarnessRequest, HarnessResult, PiReading
# Les methodes de dialogue du sous-protocole d'interface : elles BLOQUENT pi
# jusqu'a la reponse. notify, setStatus, setWidget... n'attendent rien.
DIALOGS = ("select", "confirm", "input", "editor")
SETTLED = "agent_settled" # plus rien ne suivra automatiquement : fin du tour
def rpc_argv(request: HarnessRequest, session_id: str) -> list[str]:
"""La ligne de commande d'une session vivante : memes drapeaux qu'une phase, autre mode."""
return ["pi", "--mode", "rpc", *harness._pi_flags(request, session_id)]
@dataclass
class Turn:
"""Ce qu'un tour RPC rend au runner : la lecture du port, le brut, la duree."""
reading: PiReading
lines: list[str] = field(default_factory=list)
seconds: float = 0.0
class PiRpc:
"""Un processus pi --mode rpc, pilote en JSONL. A fermer — `with` s'en charge.
Le processus herite de l'environnement du runner (schema de la porte,
socle attendu, limites du roster) UNE fois, a l'ouverture : une
enveloppe par processus, une garde par processus. Chaque clone rejoue
session_start — compteurs remis a zero, inventaire redepose — et se
juge avec les memes preuves qu'une phase ordinaire.
"""
def __init__(self, request: HarnessRequest, session_id: str | None = None,
extra_env: dict[str, str] | None = None,
executable: list[str] | None = None,
ledger: Path = harness.LEDGER) -> None:
self.request = request
self.session_id = session_id or request.session_id or str(uuid.uuid4())
self._argv = executable or rpc_argv(request, self.session_id)
self._extra_env = dict(extra_env or {})
self._ledger = ledger
self.proc: subprocess.Popen[str] | None = None
self._lines: queue.Queue[str | None] = queue.Queue()
self._env: dict[str, str] = {}
self._seq = 0
self.refused: list[dict] = [] # les demandes d'interface refusees (ask = refus)
def __enter__(self) -> "PiRpc":
self.start()
return self
def __exit__(self, *_exc) -> None:
self.close()
# ---------- ouverture / fermeture ----------
def start(self) -> None:
program = shutil.which(self._argv[0]) or self._argv[0]
if not Path(program).exists() and shutil.which(program) is None:
raise HarnessError(f"{self._argv[0]!r} introuvable dans le PATH — pi est-il installe ?",
self.session_id)
harness.SESSION_DIR.mkdir(parents=True, exist_ok=True)
# Le meme contexte d'environnement qu'une phase : la porte, le socle,
# puis ce que l'ADW ajoute (limites, identite de phase).
schema_path = harness._write_schema(self.request.schema
or envelopes.schema(envelopes.Envelope))
env = {**os.environ, harness.SCHEMA_ENV: str(schema_path),
harness.EXPECTED_ENV: ",".join(self.request.socle), **self._extra_env}
# Le registre AVANT le processus, comme pour une phase (A9).
self._env = env
harness.declare_session(self.session_id, self.request.cwd, self.request.model,
ledger=self._ledger, env=env)
# stderr va dans un fichier, jamais dans un tube que personne ne lit :
# un processus long qui remplit un tube ferme finit par se bloquer.
self._stderr_path = harness.SESSION_DIR / f"{self.session_id}.stderr.log"
self._stderr = self._stderr_path.open("w", encoding="utf-8")
self.proc = subprocess.Popen(
[program, *self._argv[1:]], stdin=subprocess.PIPE, stdout=subprocess.PIPE,
stderr=self._stderr, text=True, encoding="utf-8", errors="replace",
env=env, cwd=self.request.cwd, bufsize=1)
threading.Thread(target=self._pump, daemon=True).start()
def close(self) -> None:
"""Fin d'entree = arret propre : pi emet session_shutdown (le miroir A4 tourne)."""
if self.proc is None:
return
try:
self.proc.stdin.close()
except OSError:
pass
try:
self.proc.wait(timeout=30)
except subprocess.TimeoutExpired:
self.proc.kill()
self.proc.wait()
self._stderr.close()
self.proc = None
def _pump(self) -> None:
# Le lecteur du tube, dans son fil : stdout est vide en continu, pi
# ne se bloque jamais sur une sortie pleine. "\n" seul separe.
assert self.proc is not None and self.proc.stdout is not None
for line in self.proc.stdout:
self._lines.put(line.rstrip("\r\n"))
self._lines.put(None)
def _stderr_tail(self) -> str:
try:
self._stderr.flush()
return self._stderr_path.read_text(encoding="utf-8", errors="replace")[-400:].strip()
except (OSError, ValueError, AttributeError):
return ""
# ---------- le protocole ----------
def _write(self, obj: dict) -> None:
assert self.proc is not None and self.proc.stdin is not None
self.proc.stdin.write(json.dumps(obj, ensure_ascii=False) + "\n")
self.proc.stdin.flush()
def _next(self, deadline: float) -> tuple[dict, str]:
"""La prochaine ligne JSON du flux — les demandes de dialogue sont refusees en passant."""
while True:
remaining = deadline - time.monotonic()
if remaining <= 0:
raise HarnessError(f"pi (rpc) muet apres {self.request.timeout} s", self.session_id)
try:
line = self._lines.get(timeout=min(remaining, 1.0))
except queue.Empty:
if self.proc is not None and self.proc.poll() is not None:
raise HarnessError(f"pi (rpc) s'est arrete ({self.proc.returncode}) : "
f"{self._stderr_tail()}", self.session_id) from None
continue
if line is None:
raise HarnessError(f"pi (rpc) a ferme sa sortie : {self._stderr_tail()}",
self.session_id)
try:
event = json.loads(line)
except json.JSONDecodeError:
continue
if event.get("type") == "extension_ui_request" and event.get("method") in DIALOGS:
# ask = refus : personne n'est la pour repondre, et pi attendrait.
self.refused.append(event)
self._write({"type": "extension_ui_response", "id": event.get("id"),
"cancelled": True})
continue
return event, line
def _command(self, payload: dict, timeout: float = 60) -> tuple[dict, list[str]]:
"""Envoie une commande, attend SA reponse ; rend aussi les evenements croises."""
self._seq += 1
rid = f"c{self._seq}"
self._write({"id": rid, **payload})
deadline, seen = time.monotonic() + timeout, []
while True:
event, line = self._next(deadline)
if event.get("type") == "response" and event.get("id") == rid:
if not event.get("success"):
raise HarnessError(f"pi (rpc) refuse {payload['type']} : {event.get('error')}",
self.session_id)
return (event.get("data") or {}), seen
seen.append(line) # un evenement du tour, deja en route : il compte
def command(self, payload: dict, timeout: float = 60) -> dict:
return self._command(payload, timeout)[0]
# ---------- ce que l'usine emploie ----------
def prompt(self, message: str) -> Turn:
"""Un tour complet : la demande, puis tout le flux jusqu'a agent_settled — lu comme le port."""
clock = time.monotonic()
_, seen = self._command({"type": "prompt", "message": message}) # accepte, ou refuse
lines, deadline = list(seen), clock + self.request.timeout
while True:
event, line = self._next(deadline) # les demandes d'interface y sont refusees
lines.append(line)
if event.get("type") == SETTLED: # agent_settled : plus rien ne suivra
break
return Turn(reading=harness._read_pi_stream(lines), lines=lines,
seconds=time.monotonic() - clock)
def judge(self, turn: Turn, request: HarnessRequest | None = None) -> HarnessResult:
"""Le verdict du port, a l'identique : socle prouve, porte typee, modele observe."""
return harness._judge(turn.reading, request or self.request, self.session_id, 0, "")
def state(self) -> dict:
return self.command({"type": "get_state"})
def stats(self) -> dict:
"""Jetons et cout cumules de la session courante, tels que pi les tient."""
return self.command({"type": "get_session_stats"})
def clone(self) -> dict:
"""Une session neuve, copie de la branche courante a sa position : le prefixe partage.
pi bascule dessus et relie les extensions (session_start rejoue). Le
nouvel identifiant est lu dans get_state et declare au registre.
"""
data = self.command({"type": "clone"})
if data.get("cancelled"):
raise HarnessError("clone annule par une extension", self.session_id)
state = self.state()
self.session_id = str(state.get("sessionId") or self.session_id)
harness.declare_session(self.session_id, self.request.cwd, self.request.model,
ledger=self._ledger, env=self._env)
return state
def switch_session(self, session_file: str) -> dict:
"""Revenir a une session par son fichier (celui de l'amorce, entre deux clones)."""
data = self.command({"type": "switch_session", "sessionPath": session_file})
if data.get("cancelled"):
raise HarnessError("switch_session annule par une extension", self.session_id)
state = self.state()
self.session_id = str(state.get("sessionId") or self.session_id)
return state
def set_model(self, model: str) -> HarnessRequest:
"""Change le moteur de la session courante ; rend la requete a juger avec ce modele."""
provider, _, model_id = harness._route(model).partition("/")
self.command({"type": "set_model", "provider": provider, "modelId": model_id})
return replace(self.request, model=model)
def set_thinking(self, level: str) -> None:
self.command({"type": "set_thinking_level", "level": level})
def abort(self) -> None:
"""L'arret externe : pi interrompt le tour et repond quand il est inactif."""
self.command({"type": "abort"}, timeout=30)
# ---------- un faux pi RPC pour la gate : meme framing, meme protocole, zero jeton ----------
def _fake_pi() -> int:
"""Joue pi --mode rpc : repond aux commandes, diffuse un tour scripte par prompt."""
out = sys.stdout
turns, made, current = 0, 0, 0
# Le modele demande sur la ligne de commande est celui qu'on observe —
# comme le vrai pi quand l'id est exact (A6 verifie l'inverse).
route = sys.argv[sys.argv.index("--model") + 1] if "--model" in sys.argv else "openrouter/fake/model"
provider, _, model = route.partition("/")
def emit(obj: dict) -> None:
out.write(json.dumps(obj) + "\n")
out.flush()
def ok(rid, command, data=None):
response = {"id": rid, "type": "response", "command": command, "success": True}
if data is not None:
response["data"] = data
emit(response)
while True:
raw = sys.stdin.readline()
if not raw:
return 0 # fin d'entree : arret propre
cmd = json.loads(raw)
rid, kind = cmd.get("id"), cmd.get("type")
if kind == "prompt":
emit({"type": "agent_start"}) # deja en route avant la reponse : le client le garde
ok(rid, "prompt")
emit({"type": "entry_appended", "entry": {"type": "custom", "customType":
harness.INVENTORY_ENTRY, "data": {"pi_version": "fake",
"present": list(harness.SOCLE), "missing": []}}})
if turns == 0:
# Un damage control en mode ask : pi BLOQUE jusqu'a la reponse.
emit({"type": "extension_ui_request", "id": "ui-1", "method": "confirm",
"title": "rm -rf adws ?"})
answer = json.loads(sys.stdin.readline())
assert answer.get("cancelled") is True, answer
emit({"type": "message_end", "message": {
"role": "assistant", "provider": provider, "model": model,
"content": [{"type": "text", "text": cmd["message"][:24]}],
"stopReason": "toolUse", "usage": {"totalTokens": 100, "cost": {"total": 0.001}}}})
emit({"type": "tool_execution_end", "toolName": envelopes.REPORT_TOOL,
"isError": False, "result": {"details": {"status": "success",
"summary": f"tour {turns}"}}})
emit({"type": SETTLED})
turns += 1
elif kind == "clone":
made += 1
current = made
ok(rid, "clone", {"cancelled": False})
elif kind == "switch_session":
current = int(Path(cmd["sessionPath"]).stem.split("-")[-1])
ok(rid, "switch_session", {"cancelled": False})
elif kind == "get_state":
ok(rid, "get_state", {"sessionId": f"fake-{current}",
"sessionFile": f"/fake/session-{current}.jsonl"})
elif kind == "get_session_stats":
ok(rid, "get_session_stats", {"cost": round(0.001 * turns, 3),
"tokens": {"total": 100 * turns}})
elif kind in ("abort", "set_model", "set_thinking_level"):
ok(rid, kind)
else:
emit({"id": rid, "type": "response", "command": kind, "success": False,
"error": f"commande inconnue {kind}"})
if __name__ == "__main__":
if "--fake" in sys.argv:
raise SystemExit(_fake_pi())
# La gate du module — zero token, sans pi : argv pur, puis un lot complet
# pilote sur le faux pi. Lancer depuis la racine :
# uv run python -m adws.adw_modules.harness_rpc
import tempfile
request = HarnessRequest(prompt="", model="fake/model", tools=("read", "bash"), timeout=20)
argv = rpc_argv(request, "sess-rpc")
assert argv[:3] == ["pi", "--mode", "rpc"] and "--approve" in argv, argv
assert "-p" not in argv and "--" not in argv # pas de prompt positionnel en RPC
assert argv[argv.index("--tools") + 1] == "read,bash,report_phase"
fake = [sys.executable, "-m", "adws.adw_modules.harness_rpc", "--fake"]
with tempfile.TemporaryDirectory() as tmp:
with PiRpc(request, session_id="sess-rpc", executable=fake,
ledger=Path(tmp) / "ledger.jsonl") as pi:
amorce = pi.judge(pi.prompt("lis la spec")) # tour 0 : l'amorce
assert amorce.envelope == {"status": "success", "summary": "tour 0"}, amorce
assert amorce.tokens == 100 and amorce.stop_reason == "toolUse"
assert len(pi.refused) == 1 and pi.refused[0]["method"] == "confirm" # ask = refus
base = pi.state()["sessionFile"]
arms = []
for k in range(2):
if k:
pi.switch_session(base) # retour a l'amorce
pi.clone() # un bras = un clone
arms.append(pi.judge(pi.prompt(f"bras {k}")).envelope["summary"])
assert arms == ["tour 1", "tour 2"] and pi.session_id == "fake-2", (arms, pi.session_id)
assert pi.stats()["tokens"]["total"] == 300
declared = (Path(tmp) / "ledger.jsonl").read_text(encoding="utf-8").splitlines()
assert len(declared) == 3 # l'amorce et les deux clones
print("harness_rpc OK — argv rpc (approve, sans -p ni --), un tour lu et juge comme une "
f"phase ({amorce.tokens} jetons), une demande de dialogue refusee, deux clones "
"depuis la meme amorce, trois sessions au registre")
Pièce — adws/adw_bon.py
Le best-of-N par clone, nouveau : le pendant local de sandbox_bestof.py (chapitre 25). Il
importe BUILDER_BRIEF et BuildEnvelope d’adw_build.py (chapitre 27), phase_env de
adw_scout.py (chapitre A7), les gates du chapitre 12 et le roster, et il n’en modifie aucun. Un
constat qui exige un arbre propre (les specs non suivies sont tolérées : c’est là que le planner
les dépose), une amorce en session vivante, un bras par clone (chacun une phase du runner, donc
une ligne dans phases avec ses jetons et son coût), une moisson qui classe. Deux boutons sans
jeton : --demo écrit une spec d’essai, --dry éprouve le classement. Le processus pi est fermé
dans tous les cas, puis l’ingestion de la trace harnais est relancée côté runner.
#!/usr/bin/env -S uv run --script
# /// script
# requires-python = ">=3.11"
# dependencies = ["pyyaml"]
# ///
"""adw_bon — best-of-N par clone, dans une session vivante (annexe A9).
Usage :
uv run adws/adw_bon.py specs/ma-spec.md [--arms 3] [--models a/b,c/d] [--config ...]
uv run adws/adw_bon.py --demo # ecrit specs/bon-demo.md, la spec d'essai, puis sort
uv run adws/adw_bon.py --dry # la gate a sec : classement et lot fictifs, zero token
Un seul processus pi (--mode rpc, harness_rpc.py). D'abord l'AMORCE : le
builder lit la spec et les fichiers qu'elle nomme, n'ecrit rien, rend une
enveloppe par la porte typee — c'est le prefixe partage, paye une fois.
Puis N BRAS : chacun est un clone de l'amorce (meme contexte lu, zero
relecture), implemente, est juge par les gates du roster ; son patch est
moissonne, l'arbre est remis a zero. La MOISSON classe : vert d'abord, le
moins cher, puis le plus rapide. Le code choisit ; vous tranchez.
Le pendant local du best-of-N hors-site du chapitre 25 (N boites, N
rosters) : ici une seule variable bouge — le tirage, ou le modele par bras
avec --models — et l'amorce est commune.
"""
import argparse
import json
import subprocess
import sys
import time
import uuid
from dataclasses import dataclass, field, replace
from pathlib import Path
from adw_build import BUILDER_BRIEF, BuildEnvelope
from adw_modules import envelopes, gates, harness, harness_rpc, harness_trace, roster
from adw_modules.envelopes import EnvelopeError
from adw_modules.runner import PhaseFailure, PhaseSpec, Run
from adw_scout import phase_env
BON_DIR = Path("adws/adw_data/bon") # un dossier par lot : patches et moisson — jamais commite
DEMO_SPEC = Path("specs/bon-demo.md")
DEMO_TEXT = """# Compteur de mots
Dans le dossier du payload, ajouter un module `src/word-count.ts` qui exporte une
fonction `wordCount(text: string): number` : le nombre de mots separes par des blancs
(espaces, tabulations, retours a la ligne), 0 pour une chaine vide ou faite de blancs.
Ajouter son test `src/word-count.test.ts` (bun test) : chaine vide, un mot, plusieurs
blancs consecutifs, retours a la ligne.
Ne toucher a aucun autre fichier. Verifier avec la suite de tests du payload.
"""
AMORCE_BRIEF = """Tu es le builder de l'usine. Cette premiere etape est une LECTURE.
- Lis la spec EN ENTIER : {spec}
- Lis chaque fichier qu'elle nomme, et les tests existants du payload ({payload}).
- N'ecris RIEN et ne lance aucune commande qui modifie le depot.
- Quand tu as tout lu, resume en cinq lignes au plus ce que tu changeras
(dans 'summary'), avec changed_files: [] — tu n'as encore rien modifie."""
BRAS_ASK = """Tu as deja lu la spec et les fichiers concernes : implemente maintenant
la spec {spec}, sans la relire, exactement, rien de plus.
"""
# ---------- git : l'arbre propre, le patch de chaque bras, la remise a zero ----------
def git(*args: str) -> str:
proc = subprocess.run(["git", *args], capture_output=True, text=True, encoding="utf-8")
if proc.returncode != 0:
raise PhaseFailure(f"git {' '.join(args)} : {proc.stderr.strip()[-300:]}")
return proc.stdout
def dirty_lines() -> list[str]:
"""Ce qui sale l'arbre — sauf les specs non suivies, que le planner depose la."""
return [line for line in git("status", "--porcelain").splitlines()
if not line.startswith("?? specs/")]
def harvest_patch(target: Path) -> int:
"""Le patch du bras — tout sauf specs/ — puis l'index remis tel qu'il etait."""
git("add", "-A", "--", ".", ":(exclude)specs")
patch = git("diff", "--cached", "--binary")
git("reset", "-q")
target.write_text(patch, encoding="utf-8")
return len(patch.encode("utf-8"))
def reset_tree() -> None:
"""L'arbre revient a HEAD : fichiers suivis restaures, nouveaux supprimes — specs/ preservee."""
git("checkout", "--", ".")
git("clean", "-qfd", "-e", "specs")
# ---------- le lot : ce que les phases partagent ----------
@dataclass
class Lot:
spec: str
builder: roster.AgentSpec
factory: roster.Roster
models: list[str]
run_dir: Path
pi: harness_rpc.PiRpc | None = None
request: harness.HarnessRequest | None = None
base_file: str = "" # le fichier de session de l'amorce
amorce: dict = field(default_factory=dict)
results: list[dict] = field(default_factory=list)
def close(self) -> None:
if self.pi is not None:
self.pi.close() # session_shutdown : le miroir A4 tourne
self.pi = None
try:
# L'ingestion cote runner, apres la fermeture : idempotente (A4),
# elle rattrape ce que le miroir n'aurait pas eu le temps d'ecrire.
harness_trace.ingest(harness_trace.connect())
except Exception as error: # best effort : le brut est deja sur le disque
print(f"[bon] ingestion differee : {error}", file=sys.stderr)
# ---------- les phases ----------
def constat(lot: Lot):
"""Phase code : l'arbre est propre, la spec existe, le dossier du lot est cree — zero token."""
def action(run: Run, attempt: int):
git("rev-parse", "--is-inside-work-tree")
dirty = dirty_lines()
if dirty:
raise PhaseFailure("arbre de travail sale — commitez ou remisez d'abord, le lot "
f"remet l'arbre a zero entre deux bras :\n" + "\n".join(dirty[:10]))
spec = Path(lot.spec)
if not spec.is_file() or spec.stat().st_size == 0:
raise PhaseFailure(f"spec introuvable ou vide : {lot.spec}")
lot.run_dir.mkdir(parents=True, exist_ok=True)
return {"spec": lot.spec, "arms": len(lot.models)}
return action
def amorce(lot: Lot):
"""Phase agent : ouvrir la session vivante, lire — le prefixe que tous les bras partagent."""
def action(run: Run, attempt: int):
agent = lot.builder
lot.request = harness.HarnessRequest(
prompt="", model=agent.model, thinking=agent.thinking, tools=agent.tools,
timeout=agent.timeout, schema=envelopes.schema(BuildEnvelope))
# Les limites du roster et l'identite de la phase, par l'environnement
# (A7) — fixes pour tout le processus : une garde par processus.
lot.pi = harness_rpc.PiRpc(lot.request, extra_env=phase_env("amorce", agent, run))
lot.pi.start()
run.sessions["amorce"] = lot.pi.session_id
ask = (AMORCE_BRIEF.format(spec=lot.spec, payload=lot.factory.payload.dir)
+ "\n\n" + envelopes.contract(BuildEnvelope))
try:
turn = lot.pi.prompt(ask)
result = lot.pi.judge(turn)
except harness.HarnessError as error:
raise PhaseFailure(str(error)) from None
run.cost_usd += result.cost_usd
run.tokens += result.tokens
lot.amorce = {"tokens": result.tokens, "cost_usd": round(result.cost_usd, 4),
"seconds": round(turn.seconds, 1), "session_id": lot.pi.session_id}
lot.base_file = str(lot.pi.state().get("sessionFile") or "")
if not lot.base_file:
raise PhaseFailure("pi n'a pas rendu le fichier de session de l'amorce")
payload = result.envelope if result.envelope is not None else result.text
try:
return envelopes.parse(payload, envelopes.Envelope) # lecture : l'enveloppe de base
except EnvelopeError as error:
raise PhaseFailure(f"enveloppe d'amorce invalide : {error}") from None
return action
def bras(lot: Lot, k: int, model: str):
"""Phase agent : un clone de l'amorce, une implementation, les gates, le patch, la remise a zero.
Un bras rouge n'arrete pas le lot : son verdict et son motif sont dans le
resultat et dans la trace (evenement bon_bras) ; la moisson classera.
"""
def action(run: Run, attempt: int):
pi, request = lot.pi, lot.request
assert pi is not None and request is not None
try:
if k > 1:
pi.switch_session(lot.base_file) # retour a l'amorce
pi.clone() # session_start rejoue : garde et inventaire a neuf
if model != lot.builder.model:
request = pi.set_model(model) # une variable par bras, tenue par le code
except harness.HarnessError as error:
raise PhaseFailure(str(error)) from None # sans clone, pas de bras : le lot s'arrete
run.sessions[f"bras_{k}"] = pi.session_id
arm = {"bras": k, "model": model, "session_id": pi.session_id, "verdict": "rouge",
"cost_usd": 0.0, "tokens": 0, "seconds": 0.0, "patch_bytes": 0, "motif": ""}
clock = time.monotonic()
try:
turn = pi.prompt(BRAS_ASK.format(spec=lot.spec) + BUILDER_BRIEF + "\n\n"
+ envelopes.contract(BuildEnvelope))
result = pi.judge(turn, request)
arm["cost_usd"], arm["tokens"] = round(result.cost_usd, 4), result.tokens
payload = result.envelope if result.envelope is not None else result.text
envelope: BuildEnvelope = envelopes.parse(payload, BuildEnvelope)
if envelope.status != "success":
arm["motif"] = f"le builder declare un echec : {envelope.summary}"
else:
# Les gates du roster, sur l'arbre tel que le bras l'a laisse — zero token.
reports = [gates.changed_files_exist(envelope)]
reports += [gates.command(cmd, cwd)(envelope)
for cmd, cwd in lot.factory.payload.commands()]
arm["motif"] = gates.motif(reports)
arm["verdict"] = "vert" if not arm["motif"] else "rouge"
except (harness.HarnessError, EnvelopeError) as error:
arm["motif"] = str(error)
finally:
arm["seconds"] = round(time.monotonic() - clock, 1)
arm["patch_bytes"] = harvest_patch(lot.run_dir / f"bras-{k}.patch")
reset_tree()
run.cost_usd += arm["cost_usd"]
run.tokens += arm["tokens"]
run.tracer.event(run.adw_id, "bon_bras", f"bras_{k}", payload=arm)
print(f"[bon] bras {k} ({model}) : {arm['verdict']} — {arm['tokens']} jetons, "
f"{arm['cost_usd']:.4f} $, {arm['seconds']} s, patch {arm['patch_bytes']} o"
+ (f" — {arm['motif'].splitlines()[0][:80]}" if arm["motif"] else ""),
file=sys.stderr)
lot.results.append(arm)
return arm
return action
# ---------- la moisson : classer, chiffrer, proposer — zero token ----------
def rank(results: list[dict]) -> list[dict]:
"""Vert d'abord, puis le moins cher, puis le plus rapide. Pas de score compose."""
return sorted(results, key=lambda a: (a["verdict"] != "vert", a["cost_usd"], a["seconds"]))
def print_table(ranked: list[dict]) -> None:
print(f"{'bras':<5}{'modele':<32}{'verdict':<8}{'jetons':>9}{'$':>9}{'s':>7}{'patch':>8}")
for a in ranked:
print(f"{a['bras']:<5}{a['model'][:31]:<32}{a['verdict']:<8}{a['tokens']:>9}"
f"{a['cost_usd']:>9.4f}{a['seconds']:>7.1f}{a['patch_bytes']:>8}")
def moisson(lot: Lot):
def action(run: Run, attempt: int):
ranked = rank(lot.results)
print_table(ranked)
n = len(lot.results)
saved = lot.amorce.get("tokens", 0) * (n - 1)
print(f"amorce : {lot.amorce.get('tokens', 0)} jetons, {lot.amorce.get('cost_usd', 0):.4f} $, "
f"payee une fois — a froid elle l'aurait ete {n} fois : ~{saved} jetons de plus "
"(estimation, hors cache du fournisseur)")
out = lot.run_dir / "harvest.json"
out.write_text(json.dumps({"adw_id": run.adw_id, "spec": lot.spec, "amorce": lot.amorce,
"results": ranked}, indent=2, ensure_ascii=False),
encoding="utf-8")
winner = ranked[0]
if winner["verdict"] != "vert":
raise PhaseFailure(f"aucun bras vert sur {n} — motifs dans {out}")
patch = lot.run_dir / f"bras-{winner['bras']}.patch"
print(f"propose : bras {winner['bras']} ({winner['model']}) — vert, le moins cher des verts")
print(f"dispose : git apply --index {patch}")
return ranked
return action
# ---------- la gate a sec : le classement sur un lot fictif ----------
def selftest() -> int:
arms = [
{"bras": 1, "model": "m", "verdict": "rouge", "cost_usd": 0.05, "tokens": 40_000, "seconds": 60.0, "patch_bytes": 900, "motif": "bun test : 1 fail"},
{"bras": 2, "model": "m", "verdict": "vert", "cost_usd": 0.22, "tokens": 150_000, "seconds": 140.0, "patch_bytes": 1800, "motif": ""},
{"bras": 3, "model": "m", "verdict": "vert", "cost_usd": 0.19, "tokens": 130_000, "seconds": 170.0, "patch_bytes": 1700, "motif": ""},
]
ranked = rank(arms)
assert [a["bras"] for a in ranked] == [3, 2, 1], ranked # vert, puis le moins cher
assert rank(arms[:1])[0]["verdict"] == "rouge" # un lot tout rouge se classe quand meme
print_table(ranked)
print("adw_bon OK — classement : vert d'abord, le moins cher, puis le plus rapide ; "
"lot fictif, zero token")
return 0
def main() -> int:
parser = argparse.ArgumentParser(description="best-of-N par clone, en session vivante")
parser.add_argument("spec", nargs="?", help="la spec a implementer (sous specs/)")
parser.add_argument("--arms", type=int, default=3, help="nombre de bras, 2 a 5")
parser.add_argument("--models", default="",
help="un modele par bras, separes par des virgules (ids du registre)")
parser.add_argument("--config", default=str(roster.DEFAULT_PATH))
parser.add_argument("--demo", action="store_true", help=f"ecrire {DEMO_SPEC} puis sortir")
parser.add_argument("--dry", action="store_true", help="la gate a sec, zero token")
args = parser.parse_args()
if args.dry:
return selftest()
if args.demo:
DEMO_SPEC.parent.mkdir(exist_ok=True)
DEMO_SPEC.write_text(DEMO_TEXT, encoding="utf-8")
print(f"spec d'essai ecrite : {DEMO_SPEC}")
return 0
if not args.spec:
parser.error("spec manquante — uv run adws/adw_bon.py specs/ma-spec.md")
factory = roster.load(args.config) # zero token : tout echec est gratuit
builder = factory.agents["builder"]
models = ([m.strip() for m in args.models.split(",") if m.strip()]
if args.models else [builder.model] * args.arms)
if not 2 <= len(models) <= 5:
print("un best-of-N va de 2 a 5 bras — au-dela, la moisson n'est plus lisible",
file=sys.stderr)
return 1
run = Run(adw_id=uuid.uuid4().hex[:8])
lot = Lot(spec=args.spec, builder=builder, factory=factory, models=models,
run_dir=BON_DIR / run.adw_id)
phases = [PhaseSpec(name="constat_bon", kind="code", action=constat(lot)),
PhaseSpec(name="amorce", kind="agent", action=amorce(lot))]
phases += [PhaseSpec(name=f"bras_{k}", kind="agent", action=bras(lot, k, model))
for k, model in enumerate(models, start=1)]
phases.append(PhaseSpec(name="moisson", kind="code", action=moisson(lot)))
try:
return run.execute(phases)
finally:
lot.close() # toujours : le processus pi ne survit pas au lot
if __name__ == "__main__":
sys.exit(main())
La gate du TP
Depuis la racine de plume-factory, sur un arbre propre (git status vide, hors specs).
Une commande par ligne, identiques dans bash et PowerShell. Les trois premières ne coûtent rien,
la quatrième écrit la spec d’essai, la cinquième lance le lot avec trois bras sur le builder du
roster, et les deux dernières lisent la trace.
uv run python -m adws.adw_modules.harness
uv run python -m adws.adw_modules.harness_rpc
uv run adws/adw_bon.py --dry
uv run adws/adw_bon.py --demo
uv run adws/adw_bon.py specs/bon-demo.md --arms 3
sqlite3 adws/adw_data/factory.db "SELECT name, status, tokens, cost_usd FROM phases WHERE adw_id = (SELECT adw_id FROM runs WHERE adw_name = 'adw_bon' ORDER BY started_at DESC LIMIT 1) ORDER BY seq;"
sqlite3 adws/adw_data/factory.db "SELECT COUNT(DISTINCT session_id) FROM harness_by_phase WHERE adw_id = (SELECT adw_id FROM runs WHERE adw_name = 'adw_bon' ORDER BY started_at DESC LIMIT 1);"
Résultat attendu : harness OK — … registre v6 : declare avant l'appel, reprise hors dossier refusee, puis harness_rpc OK — argv rpc (approve, sans -p ni --), un tour lu et juge comme une phase (100 jetons), une demande de dialogue refusee, deux clones depuis la meme amorce, trois sessions au registre, puis adw_bon OK — classement : vert d'abord, le moins cher, puis le plus rapide. Vient ensuite le
lot : une ligne [bon] bras k (…) : vert — … jetons, … $, … s, patch … o par bras, le tableau de la
moisson, la ligne amorce : … jetons … payee une fois — a froid elle l'aurait ete 3 fois, et
propose : bras k … — vert, le moins cher des verts suivi de la commande git apply --index qui
pose le patch gagnant, et l’arbre est revenu propre. La première requête montre amorce, bras_1,
bras_2, bras_3 avec leurs jetons, l’amorce une seule fois, qui est le chiffre du jalon. La
seconde compte quatre sessions jointes au run (l’amorce et trois clones), chacune avec sa trace
harnais. Coût : les quatre gates à sec ne dépensent rien, et le lot coûte de cinquante
centimes à un dollar sur le workhorse, en cinq à dix minutes (les bras sont séquentiels : un
seul arbre de travail). Variante éco : --models avec trois fois le modèle léger du roster, ou
--arms 2. Une variante d’essai du refus, sans jeton : relancez n’importe quel ADW avec une
session du registre depuis un autre dossier : c’est le port qui refuse, pi ne démarre pas.