feat(agent): adapters structurés Claude/Codex + fake CLI + conformité (D2) — §17
infrastructure/src/session/ : machinerie de process générique (paramétrable par la commande = seam d'injection du fake CLI), adapters ClaudeSdkSession/ CodexExecSession avec parsing ISOLÉ par adapter (parse_event), factory StructuredSessionFactory (routage par structured_adapter), FakeCli scriptable + harnais de conformité Liskov assert_agent_session_contract. Incarnation « un run par tour » (send relance claude -p / --resume <id>, continuité via conversation_id — colle au pivot reprise B). Tests : 41 contre le FAKE CLI (jamais le vrai claude/codex), workspace vert. Points en attente des spikes S1/S2 (format réel) — n'impactent que parse_event : - mapping JSON→ReplyEvent Claude (S1) et Codex (S2) sur schémas SUPPOSÉS ; - Claude multi-blocs : parse_event ne garde que le 1er bloc (à corriger si Claude émet plusieurs blocs/message — confirmer S1) ; - flux sans Final : permissif côté adapter, l'erreur est gérée par send_blocking (consommateur). À reconfirmer côté UI streaming (D4). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
141
crates/infrastructure/src/session/process.rs
Normal file
141
crates/infrastructure/src/session/process.rs
Normal file
@ -0,0 +1,141 @@
|
||||
//! Machinerie de process **générique et paramétrable par la commande** pour les
|
||||
//! adapters structurés (ARCHITECTURE §17.2). Strictement **séparée** du parsing du
|
||||
//! format de chaque CLI (cf. `claude::parse_event` / `codex::parse_event`).
|
||||
//!
|
||||
//! # Rôle
|
||||
//!
|
||||
//! Lancer un process CLI dont la sortie est **orientée lignes** (un objet/événement
|
||||
//! par ligne, typiquement du JSONL), lire ces lignes jusqu'à épuisement (ou jusqu'à
|
||||
//! un timeout), et les remettre au parser propre à l'adapter. Aucune connaissance
|
||||
//! de Claude/Codex/JSON ici : *seulement* spawn, pipes, drain ligne-à-ligne, kill.
|
||||
//!
|
||||
//! # Pourquoi paramétrable par la commande
|
||||
//!
|
||||
//! Le binaire et ses arguments sont fournis par l'appelant ([`SpawnLine`]). En prod
|
||||
//! l'adapter passe `claude … --output-format stream-json` ; en test, un **fake CLI
|
||||
//! scriptable** (cf. `conformance`) est substitué pour rejouer un script de lignes
|
||||
//! canned — déterministe, hors-réseau. La machinerie est identique dans les deux cas.
|
||||
//!
|
||||
//! # Incarnation « un run par tour »
|
||||
//!
|
||||
//! `run_turn` lance le process, lui passe (optionnellement) le `prompt` sur stdin,
|
||||
//! draine **toutes** les lignes de stdout jusqu'à EOF, attend la fin du process, et
|
||||
//! retourne les lignes brutes. C'est l'incarnation (b) de §17.2 (« un `exec` par
|
||||
//! `send()` », réattaché via l'id de conversation) — la plus simple et la plus
|
||||
//! déterministe ; elle suffit au contrat de port et au fake CLI. Une éventuelle
|
||||
//! incarnation « process persistant piloté en flux » resterait derrière le **même**
|
||||
//! port sans changer le parsing.
|
||||
|
||||
use std::io;
|
||||
use std::process::Stdio;
|
||||
use std::time::Duration;
|
||||
|
||||
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
||||
use tokio::process::Command;
|
||||
|
||||
use domain::ports::AgentSessionError;
|
||||
|
||||
/// Une invocation orientée lignes : binaire + arguments + cwd + env + prompt à
|
||||
/// pousser sur stdin (le cas échéant). **Paramétrable par la commande** : c'est ce
|
||||
/// que le test substitue pour brancher le fake CLI.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct SpawnLine {
|
||||
/// Exécutable à lancer (le vrai `claude`/`codex`, ou le fake CLI en test).
|
||||
pub command: String,
|
||||
/// Arguments (déjà composés par l'adapter : `--output-format stream-json`, …).
|
||||
pub args: Vec<String>,
|
||||
/// Répertoire de travail (run dir isolé §14.1).
|
||||
pub cwd: String,
|
||||
/// Variables d'environnement additionnelles.
|
||||
pub env: Vec<(String, String)>,
|
||||
/// Contenu à écrire sur stdin du process (`None` ⇒ stdin fermé immédiatement).
|
||||
/// Sert au mode `--input-format stream-json` / au passage du prompt.
|
||||
pub stdin: Option<String>,
|
||||
}
|
||||
|
||||
/// Lance `spec`, pousse `spec.stdin` sur l'entrée standard, **draine toutes les
|
||||
/// lignes de stdout jusqu'à EOF**, puis attend la fin du process. Retourne les
|
||||
/// lignes brutes (sans `\n`) dans l'ordre d'émission.
|
||||
///
|
||||
/// `timeout` borne l'ensemble run+drain : à expiration le process est tué et
|
||||
/// [`AgentSessionError::Timeout`] est retourné (la machinerie ne suppose jamais que
|
||||
/// l'appelant veut attendre indéfiniment). `None` ⇒ pas de borne.
|
||||
///
|
||||
/// # Errors
|
||||
/// - [`AgentSessionError::Start`] si le process ne démarre pas (binaire introuvable) ;
|
||||
/// - [`AgentSessionError::Io`] sur échec de lecture/écriture des pipes ;
|
||||
/// - [`AgentSessionError::Timeout`] si `timeout` expire.
|
||||
pub async fn run_turn(
|
||||
spec: &SpawnLine,
|
||||
timeout: Option<Duration>,
|
||||
) -> Result<Vec<String>, AgentSessionError> {
|
||||
match timeout {
|
||||
Some(dur) => match tokio::time::timeout(dur, drain(spec)).await {
|
||||
Ok(result) => result,
|
||||
Err(_elapsed) => Err(AgentSessionError::Timeout),
|
||||
},
|
||||
None => drain(spec).await,
|
||||
}
|
||||
}
|
||||
|
||||
/// Cœur du drain : spawn → écriture stdin → lecture ligne-à-ligne → wait.
|
||||
async fn drain(spec: &SpawnLine) -> Result<Vec<String>, AgentSessionError> {
|
||||
let mut cmd = Command::new(&spec.command);
|
||||
cmd.args(&spec.args)
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped());
|
||||
if spec.cwd != "/" && !spec.cwd.is_empty() {
|
||||
cmd.current_dir(&spec.cwd);
|
||||
}
|
||||
for (key, value) in &spec.env {
|
||||
cmd.env(key, value);
|
||||
}
|
||||
|
||||
let mut child = cmd
|
||||
.spawn()
|
||||
.map_err(|e| AgentSessionError::Start(format!("{}: {e}", spec.command)))?;
|
||||
|
||||
// Pousse le prompt sur stdin puis le ferme (EOF) — sans quoi un process en mode
|
||||
// `stream-json` attend indéfiniment. Erreur d'I/O ⇒ `Io` (pas `Start`).
|
||||
if let Some(input) = &spec.stdin {
|
||||
let mut stdin = child
|
||||
.stdin
|
||||
.take()
|
||||
.ok_or_else(|| AgentSessionError::Io("stdin pipe indisponible".to_owned()))?;
|
||||
stdin
|
||||
.write_all(input.as_bytes())
|
||||
.await
|
||||
.map_err(|e| AgentSessionError::Io(e.to_string()))?;
|
||||
// Drop explicite ⇒ EOF côté process.
|
||||
drop(stdin);
|
||||
} else {
|
||||
drop(child.stdin.take());
|
||||
}
|
||||
|
||||
let stdout = child
|
||||
.stdout
|
||||
.take()
|
||||
.ok_or_else(|| AgentSessionError::Io("stdout pipe indisponible".to_owned()))?;
|
||||
let mut lines = BufReader::new(stdout).lines();
|
||||
|
||||
let mut collected = Vec::new();
|
||||
loop {
|
||||
match lines.next_line().await {
|
||||
Ok(Some(line)) => collected.push(line),
|
||||
Ok(None) => break,
|
||||
Err(e) => return Err(map_io(e)),
|
||||
}
|
||||
}
|
||||
|
||||
// Réifie le statut du process pour ne pas laisser de zombie ; un statut non nul
|
||||
// n'est pas en soi une erreur de session (la CLI peut signaler une erreur
|
||||
// *dans* sa sortie structurée, captée par le parser).
|
||||
let _ = child.wait().await.map_err(map_io)?;
|
||||
Ok(collected)
|
||||
}
|
||||
|
||||
/// Mappe une erreur d'I/O système sur l'erreur de port appropriée.
|
||||
fn map_io(e: io::Error) -> AgentSessionError {
|
||||
AgentSessionError::Io(e.to_string())
|
||||
}
|
||||
Reference in New Issue
Block a user