diff --git a/crates/infrastructure/src/lib.rs b/crates/infrastructure/src/lib.rs index 468b786..722325b 100644 --- a/crates/infrastructure/src/lib.rs +++ b/crates/infrastructure/src/lib.rs @@ -23,6 +23,7 @@ pub mod process; pub mod pty; pub mod remote; pub mod runtime; +pub mod session; pub mod store; pub use clock::SystemClock; @@ -39,6 +40,7 @@ pub use process::LocalProcessSpawner; pub use pty::PortablePtyAdapter; pub use remote::{remote_host, LocalHost}; pub use runtime::CliAgentRuntime; +pub use session::{ClaudeSdkSession, CodexExecSession, FakeCli, StructuredSessionFactory}; #[cfg(feature = "vector-onnx")] pub use store::OnnxEmbedder; #[cfg(feature = "vector-http")] diff --git a/crates/infrastructure/src/session/claude.rs b/crates/infrastructure/src/session/claude.rs new file mode 100644 index 0000000..4942bf2 --- /dev/null +++ b/crates/infrastructure/src/session/claude.rs @@ -0,0 +1,230 @@ +//! [`ClaudeSdkSession`] — adapter structuré Claude (ARCHITECTURE §17.2, spike **S1**). +//! +//! Pilote `claude` en mode non-interactif structuré et traduit son flux **JSONL** +//! (`--output-format stream-json`, un objet JSON par ligne) vers le contrat de port +//! universel [`ReplyEvent`]. Aucun détail Claude (`stream-json`, `session_id`, +//! `--resume`) ne franchit la frontière domaine. +//! +//! # Séparation parsing / machinerie (CRUCIAL — §17.2) +//! +//! Le **parsing du format Claude est ISOLÉ** dans la fonction pure [`parse_event`]. +//! La machinerie de process (spawn, pipes, drain) vit dans [`super::process`] et +//! ignore tout du JSON. Quand le **spike S1** aura confirmé le schéma réel auprès de +//! l'utilisateur, **seule [`parse_event`] (et la composition de la commande) devra +//! changer**, pas la machinerie ni le reste de l'adapter. + +use std::sync::Mutex; + +use async_trait::async_trait; +use serde_json::Value; + +use domain::ports::{AgentSession, AgentSessionError, ReplyEvent, ReplyStream}; +use domain::SessionId; + +use super::process::{run_turn, SpawnLine}; + +/// Résultat du parsing d'une ligne : un événement à émettre (le cas échéant) et/ou +/// un `session_id` capté (init/result). Permet à [`parse_event`] de rester **pure** +/// (aucun effet de bord) tout en remontant les deux informations. +#[derive(Debug, Default, PartialEq, Eq)] +pub struct ParsedLine { + /// Événement universel à émettre, ou `None` (ligne de contrôle, ex. `init`). + pub event: Option, + /// `session_id` Claude capté sur cette ligne (id de conversation pour la reprise). + pub session_id: Option, +} + +/// **Parse une ligne du flux `stream-json` de Claude** vers le contrat universel. +/// +/// # Schéma SUPPOSÉ — à confirmer au spike S1 (réf. doc API Claude / Agent SDK) +/// +/// Le flux est du **JSONL** (un objet JSON par ligne). Schéma présumé : +/// +/// - `{"type":"system","subtype":"init","session_id":"", …}` +/// ⇒ capture le `session_id` (= id de conversation pour la reprise), **aucun** +/// événement émis. +/// - `{"type":"assistant","message":{"role":"assistant","content":[ +/// {"type":"text","text":"…"} | {"type":"tool_use","name":"…", …} +/// ]}, "session_id":"…"}` +/// ⇒ chaque bloc `text` ⇒ [`ReplyEvent::TextDelta`] ; chaque bloc `tool_use` +/// ⇒ [`ReplyEvent::ToolActivity`] (`label` = `name`). +/// - `{"type":"result","subtype":"success","result":"","session_id":"", …}` +/// ⇒ [`ReplyEvent::Final`] (`content` = `result`) et confirme le `session_id`. +/// +/// > NOTE S1 : ce schéma est **présumé**. Les noms exacts (`assistant` vs +/// > `content_block_delta`, structure de `content`, sous-type de `result`) seront +/// > vérifiés au spike. Le contrat de sortie ([`ReplyEvent`]) ne bougera pas — seule +/// > cette fonction changera. +/// +/// Une ligne **vide** est ignorée (`ParsedLine` par défaut). Un objet **inconnu** +/// (type non reconnu) est ignoré sans erreur (robustesse : la CLI peut émettre des +/// événements de contrôle non pertinents). Seul un JSON **illisible** ⇒ `Decode`. +/// +/// # Errors +/// [`AgentSessionError::Decode`] si la ligne n'est pas un JSON valide. On ne propage +/// **jamais** le JSON brut : seul un message de diagnostic court est inclus. +pub fn parse_event(line: &str) -> Result { + let trimmed = line.trim(); + if trimmed.is_empty() { + return Ok(ParsedLine::default()); + } + let value: Value = serde_json::from_str(trimmed) + .map_err(|e| AgentSessionError::Decode(format!("ligne JSON illisible: {e}")))?; + + let session_id = value + .get("session_id") + .and_then(Value::as_str) + .map(str::to_owned); + + let event = match value.get("type").and_then(Value::as_str) { + Some("system") => None, // init/handshake : on ne capte que le session_id. + Some("assistant") => first_assistant_event(&value), + Some("result") => { + value + .get("result") + .and_then(Value::as_str) + .map(|content| ReplyEvent::Final { + content: content.to_owned(), + }) + } + _ => None, // type inconnu / non pertinent : ignoré (robustesse). + }; + + Ok(ParsedLine { event, session_id }) +} + +/// Extrait le **premier** bloc de contenu pertinent d'un message `assistant` : +/// un `text` ⇒ `TextDelta`, un `tool_use` ⇒ `ToolActivity`. (Un message porte en +/// pratique un bloc ; on prend le premier pertinent — robuste si la forme évolue.) +fn first_assistant_event(value: &Value) -> Option { + let content = value + .get("message") + .and_then(|m| m.get("content")) + .and_then(Value::as_array)?; + for block in content { + match block.get("type").and_then(Value::as_str) { + Some("text") => { + if let Some(text) = block.get("text").and_then(Value::as_str) { + return Some(ReplyEvent::TextDelta { + text: text.to_owned(), + }); + } + } + Some("tool_use") => { + let label = block + .get("name") + .and_then(Value::as_str) + .unwrap_or("outil") + .to_owned(); + return Some(ReplyEvent::ToolActivity { label }); + } + _ => {} + } + } + None +} + +/// Adapter de session structurée Claude. +/// +/// Incarnation « un `claude -p` par tour » (§17.2 (b)) : chaque `send` relance la +/// CLI en passant le prompt, et — dès qu'un `session_id` a été capté — le flag de +/// reprise pour rester sur la **même** conversation. Le `session_id` Claude est +/// exposé via [`conversation_id`](AgentSession::conversation_id) (pivot de reprise +/// model-agnostic, persisté sur la cellule). +pub struct ClaudeSdkSession { + /// Id de session IdeA (mappe la cellule/agent). + id: SessionId, + /// Binaire à lancer (`claude` en prod, fake CLI en test). + command: String, + /// Répertoire de travail (run dir isolé §14.1). + cwd: String, + /// Id de conversation **du moteur** Claude, capté au premier tour, `None` avant. + conversation_id: Mutex>, +} + +impl ClaudeSdkSession { + /// Construit l'adapter. `command` est le binaire à lancer (injecté ⇒ testable + /// avec un fake CLI) ; `seed_conversation_id` amorce la reprise (`SessionPlan:: + /// Resume` côté factory) ou reste `None` pour une conversation neuve. + #[must_use] + pub fn new( + id: SessionId, + command: impl Into, + cwd: impl Into, + seed_conversation_id: Option, + ) -> Self { + Self { + id, + command: command.into(), + cwd: cwd.into(), + conversation_id: Mutex::new(seed_conversation_id), + } + } + + /// Compose la ligne de commande d'un tour selon l'état de conversation. + /// + /// - Conversation neuve : `claude -p --output-format stream-json`. + /// - Reprise (id connu) : `claude --resume -p --output-format + /// stream-json`. + /// + /// > NOTE S1 : flags exacts (`-p`, `--resume`, `--output-format stream-json`, + /// > éventuel `--input-format stream-json`) à confirmer au spike. + fn build_spawn_line(&self, prompt: &str) -> SpawnLine { + let mut args = Vec::new(); + if let Some(id) = self.conversation_id.lock().expect("mutex sain").as_ref() { + args.push("--resume".to_owned()); + args.push(id.clone()); + } + args.push("-p".to_owned()); + args.push(prompt.to_owned()); + args.push("--output-format".to_owned()); + args.push("stream-json".to_owned()); + SpawnLine { + command: self.command.clone(), + args, + cwd: self.cwd.clone(), + env: Vec::new(), + stdin: None, + } + } +} + +#[async_trait] +impl AgentSession for ClaudeSdkSession { + fn id(&self) -> SessionId { + self.id + } + + fn conversation_id(&self) -> Option { + self.conversation_id.lock().expect("mutex sain").clone() + } + + async fn send(&self, prompt: &str) -> Result { + let spec = self.build_spawn_line(prompt); + let raw_lines = run_turn(&spec, None).await?; + + let mut events = Vec::new(); + let mut captured_id = None; + for line in &raw_lines { + let parsed = parse_event(line)?; + if let Some(id) = parsed.session_id { + captured_id = Some(id); + } + if let Some(event) = parsed.event { + events.push(event); + } + } + // Persiste le session_id capté (pivot de reprise) avant de rendre le flux. + if let Some(id) = captured_id { + *self.conversation_id.lock().expect("mutex sain") = Some(id); + } + + Ok(Box::new(events.into_iter())) + } + + async fn shutdown(&self) -> Result<(), AgentSessionError> { + // Incarnation « un run par tour » : aucun process long ne survit entre les + // tours, donc `shutdown` est intrinsèquement idempotent et sans effet. + Ok(()) + } +} diff --git a/crates/infrastructure/src/session/codex.rs b/crates/infrastructure/src/session/codex.rs new file mode 100644 index 0000000..caa275b --- /dev/null +++ b/crates/infrastructure/src/session/codex.rs @@ -0,0 +1,201 @@ +//! [`CodexExecSession`] — adapter structuré Codex (ARCHITECTURE §17.2, spike **S2**). +//! +//! Pilote `codex exec` en mode non-interactif et traduit sa sortie structurée vers +//! le contrat universel [`ReplyEvent`]. **Le format de Codex est BEAUCOUP plus +//! incertain que celui de Claude** : tout le format présumé est ISOLÉ dans +//! [`parse_event`], clairement marqué « SUPPOSÉ — à confirmer S2 ». +//! +//! # Séparation parsing / machinerie (CRUCIAL — §17.2) +//! +//! Comme pour Claude, la machinerie de process vit dans [`super::process`] et ignore +//! le format. Quand le **spike S2** aura confirmé la sortie réelle (process +//! persistant vs `exec` par tour, schéma JSON de fin de tour), **seule +//! [`parse_event`] (et la composition de la commande) changera**. + +use std::sync::Mutex; + +use async_trait::async_trait; +use serde_json::Value; + +use domain::ports::{AgentSession, AgentSessionError, ReplyEvent, ReplyStream}; +use domain::SessionId; + +use super::process::{run_turn, SpawnLine}; + +/// Résultat du parsing d'une ligne Codex : un événement (le cas échéant) et/ou un id +/// de conversation Codex capté. Miroir de `claude::ParsedLine`. +#[derive(Debug, Default, PartialEq, Eq)] +pub struct ParsedLine { + /// Événement universel à émettre, ou `None`. + pub event: Option, + /// Id de conversation Codex capté (pour la reprise via le flag Codex). + pub conversation_id: Option, +} + +/// **Parse une ligne de la sortie structurée de `codex exec`** vers le contrat +/// universel. +/// +/// # Schéma SUPPOSÉ — à confirmer au spike S2 (format Codex TRÈS incertain) +/// +/// On suppose le **minimum viable** : un flux de lignes JSON (JSONL), dont +/// **un** événement final identifiable porte le texte de réponse. Schéma présumé : +/// +/// - `{"type":"session","conversation_id":"", …}` (ou `"id"`) ⇒ capte l'id de +/// conversation Codex (pour `--resume`/reprise), **aucun** événement émis. +/// - `{"type":"message"|"delta","text":"…"}` ⇒ [`ReplyEvent::TextDelta`]. +/// - `{"type":"tool"|"tool_call","name":"…"}` ⇒ [`ReplyEvent::ToolActivity`]. +/// - `{"type":"result"|"final","text":"", …}` (ou champ `"output"`) +/// ⇒ [`ReplyEvent::Final`] : **l'événement terminal du tour**. +/// +/// > NOTE S2 : ce schéma est **largement présumé**. La forme réelle de `codex exec` +/// > (noms de types, champ portant le texte, présence d'un id de conversation, mode +/// > persistant vs one-shot) est à confirmer au spike. Le contrat de sortie +/// > ([`ReplyEvent`]) ne bougera pas — seule cette fonction changera. +/// +/// Ligne vide ⇒ ignorée ; type inconnu ⇒ ignoré sans erreur ; JSON illisible ⇒ +/// [`AgentSessionError::Decode`] (jamais de JSON brut propagé). +/// +/// # Errors +/// [`AgentSessionError::Decode`] si la ligne n'est pas un JSON valide. +pub fn parse_event(line: &str) -> Result { + let trimmed = line.trim(); + if trimmed.is_empty() { + return Ok(ParsedLine::default()); + } + let value: Value = serde_json::from_str(trimmed) + .map_err(|e| AgentSessionError::Decode(format!("ligne JSON illisible: {e}")))?; + + // Id de conversation Codex : on tolère `conversation_id` ou `id` (présumé). + let conversation_id = value + .get("conversation_id") + .or_else(|| value.get("id")) + .and_then(Value::as_str) + .map(str::to_owned); + + // Champ texte présumé : `text` en priorité, repli sur `output`. + let text = || { + value + .get("text") + .or_else(|| value.get("output")) + .and_then(Value::as_str) + .unwrap_or_default() + .to_owned() + }; + + let event = match value.get("type").and_then(Value::as_str) { + Some("session") => None, // handshake : on ne capte que l'id de conversation. + Some("message" | "delta") => Some(ReplyEvent::TextDelta { text: text() }), + Some("tool" | "tool_call") => { + let label = value + .get("name") + .and_then(Value::as_str) + .unwrap_or("outil") + .to_owned(); + Some(ReplyEvent::ToolActivity { label }) + } + Some("result" | "final") => Some(ReplyEvent::Final { content: text() }), + _ => None, // type inconnu : ignoré (robustesse). + }; + + Ok(ParsedLine { + event, + conversation_id, + }) +} + +/// Adapter de session structurée Codex. +/// +/// Incarnation « un `codex exec` par tour » (§17.2 (b)) : chaque `send` relance +/// `codex exec` avec le prompt et, dès qu'un id de conversation a été capté, le flag +/// de reprise. L'incarnation « process persistant » resterait derrière le **même** +/// port sans toucher au parsing. +pub struct CodexExecSession { + /// Id de session IdeA. + id: SessionId, + /// Binaire à lancer (`codex` en prod, fake CLI en test). + command: String, + /// Répertoire de travail (run dir isolé §14.1). + cwd: String, + /// Id de conversation **du moteur** Codex, capté au premier tour, `None` avant. + conversation_id: Mutex>, +} + +impl CodexExecSession { + /// Construit l'adapter. `command` est injecté (⇒ testable avec un fake CLI) ; + /// `seed_conversation_id` amorce la reprise ou reste `None` (conversation neuve). + #[must_use] + pub fn new( + id: SessionId, + command: impl Into, + cwd: impl Into, + seed_conversation_id: Option, + ) -> Self { + Self { + id, + command: command.into(), + cwd: cwd.into(), + conversation_id: Mutex::new(seed_conversation_id), + } + } + + /// Compose la ligne de commande d'un tour. + /// + /// - Conversation neuve : `codex exec `. + /// - Reprise (id connu) : `codex exec --resume `. + /// + /// > NOTE S2 : sous-commande et flags (`exec`, `--resume`, format de sortie + /// > structuré éventuel) **à confirmer au spike**. + fn build_spawn_line(&self, prompt: &str) -> SpawnLine { + let mut args = vec!["exec".to_owned()]; + if let Some(id) = self.conversation_id.lock().expect("mutex sain").as_ref() { + args.push("--resume".to_owned()); + args.push(id.clone()); + } + args.push(prompt.to_owned()); + SpawnLine { + command: self.command.clone(), + args, + cwd: self.cwd.clone(), + env: Vec::new(), + stdin: None, + } + } +} + +#[async_trait] +impl AgentSession for CodexExecSession { + fn id(&self) -> SessionId { + self.id + } + + fn conversation_id(&self) -> Option { + self.conversation_id.lock().expect("mutex sain").clone() + } + + async fn send(&self, prompt: &str) -> Result { + let spec = self.build_spawn_line(prompt); + let raw_lines = run_turn(&spec, None).await?; + + let mut events = Vec::new(); + let mut captured_id = None; + for line in &raw_lines { + let parsed = parse_event(line)?; + if let Some(id) = parsed.conversation_id { + captured_id = Some(id); + } + if let Some(event) = parsed.event { + events.push(event); + } + } + if let Some(id) = captured_id { + *self.conversation_id.lock().expect("mutex sain") = Some(id); + } + + Ok(Box::new(events.into_iter())) + } + + async fn shutdown(&self) -> Result<(), AgentSessionError> { + // « un run par tour » ⇒ pas de process long survivant : idempotent, no-op. + Ok(()) + } +} diff --git a/crates/infrastructure/src/session/conformance.rs b/crates/infrastructure/src/session/conformance.rs new file mode 100644 index 0000000..c2f916a --- /dev/null +++ b/crates/infrastructure/src/session/conformance.rs @@ -0,0 +1,238 @@ +//! Fake CLI scriptable + **harnais de conformité de port (Liskov)** pour les +//! adapters structurés (ARCHITECTURE §17.2). Permet de tester la **machinerie** +//! (spawn, lecture ligne-à-ligne, drain jusqu'au `Final`, capture d'id de session, +//! shutdown, timeout) **sans réseau ni vraie CLI**, et d'asserter le **contrat +//! [`AgentSession`]** de façon réutilisable pour Claude ET Codex. +//! +//! Disponible hors `cfg(test)` (mais sous une porte `pub`) pour que QA puisse +//! réutiliser le harnais et le fake CLI dans des tests d'intégration ultérieurs. +//! +//! > Ce qui dépend du **format réel non vérifié** (spikes S1/S2) : les *scripts* +//! > de lignes JSON fournis aux tests reproduisent le **schéma SUPPOSÉ** documenté +//! > dans `claude::parse_event` / `codex::parse_event`. Quand S1/S2 confirmeront le +//! > vrai format, ces scripts (et le parser) seront ajustés ; la machinerie et le +//! > harnais, eux, restent valides. + +use std::io::Write; +use std::path::PathBuf; + +use super::process::SpawnLine; + +/// Un **fake CLI** : un script exécutable qui **rejoue un script de lignes** sur +/// stdout (en ignorant ses arguments), puis se termine. Substitué au vrai +/// `claude`/`codex` pour rendre la machinerie déterministe et hors-réseau. +/// +/// Le binaire est matérialisé dans un fichier temporaire ; il est supprimé au drop. +pub struct FakeCli { + /// Chemin du script exécutable généré. + path: PathBuf, +} + +impl FakeCli { + /// Crée un fake CLI qui imprimera exactement `lines` (une par ligne de stdout), + /// dans l'ordre, puis sortira avec le code 0. + /// + /// # Panics + /// Panique si le fichier temporaire ne peut être écrit (environnement de test + /// cassé) — acceptable dans un utilitaire de test. + #[must_use] + pub fn printing(lines: &[&str]) -> Self { + let mut path = std::env::temp_dir(); + // Nom unique : pid + compteur atomique pour éviter toute collision entre + // tests parallèles. + use std::sync::atomic::{AtomicU64, Ordering}; + static COUNTER: AtomicU64 = AtomicU64::new(0); + let n = COUNTER.fetch_add(1, Ordering::Relaxed); + path.push(format!("idea-fake-cli-{}-{n}", std::process::id())); + + let mut script = String::from("#!/bin/sh\n"); + for line in lines { + // `printf '%s\n'` imprime la ligne littéralement (pas d'interprétation + // d'échappements), en isolant la donnée de toute injection shell via + // l'unique argument `--`. + script.push_str("printf '%s\\n' "); + script.push_str(&shell_single_quote(line)); + script.push('\n'); + } + + let mut file = std::fs::File::create(&path).expect("création du fake CLI"); + file.write_all(script.as_bytes()) + .expect("écriture du fake CLI"); + // `sync_all` force la fermeture/flush du descripteur en écriture AVANT toute + // tentative d'exécution : sans cela, `execve` sur un binaire encore ouvert en + // écriture par un autre thread (suite parallèle) retourne `ETXTBSY` + // (« Text file busy », os error 26) — défaut de fixture intermittent. + file.sync_all().expect("sync du fake CLI"); + drop(file); + set_executable(&path); + wait_until_executable(&path); + + Self { path } + } + + /// Le binaire à passer en `command` d'un [`SpawnLine`] / d'un adapter. + #[must_use] + pub fn command(&self) -> String { + self.path.to_string_lossy().into_owned() + } + + /// Construit un [`SpawnLine`] minimal lançant ce fake CLI (utile pour tester la + /// machinerie [`super::process::run_turn`] directement). + #[must_use] + pub fn spawn_line(&self) -> SpawnLine { + SpawnLine { + command: self.command(), + args: Vec::new(), + cwd: "/".to_owned(), + env: Vec::new(), + stdin: None, + } + } +} + +impl Drop for FakeCli { + fn drop(&mut self) { + let _ = std::fs::remove_file(&self.path); + } +} + +/// Échappe une chaîne pour l'insérer en argument shell entre quotes simples. +fn shell_single_quote(s: &str) -> String { + let mut out = String::with_capacity(s.len() + 2); + out.push('\''); + for c in s.chars() { + if c == '\'' { + out.push_str("'\\''"); + } else { + out.push(c); + } + } + out.push('\''); + out +} + +/// Attend que `path` soit réellement exécutable (probe `execve` qui ne retourne +/// plus `ETXTBSY`). Sur Linux, exécuter un fichier encore ouvert en écriture par un +/// autre thread échoue avec « Text file busy » : on boucle un court instant jusqu'à +/// ce que la condition se lève, garantissant qu'un `spawn` ultérieur ne *race* pas. +#[cfg(unix)] +pub(crate) fn wait_until_executable(path: &std::path::Path) { + use std::process::{Command, Stdio}; + use std::time::{Duration, Instant}; + + let deadline = Instant::now() + Duration::from_secs(5); + loop { + // Probe la plus légère possible : on tente le spawn ; ETXTBSY ⇒ on réessaie. + match Command::new(path) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + { + Ok(mut child) => { + let _ = child.wait(); + return; + } + Err(e) if e.raw_os_error() == Some(26) && Instant::now() < deadline => { + std::thread::sleep(Duration::from_millis(2)); + } + // Toute autre erreur (ou dépassement de délai) : on rend la main, le test + // appelant remontera l'échec réel s'il subsiste. + Err(_) => return, + } + } +} + +#[cfg(not(unix))] +pub(crate) fn wait_until_executable(_path: &std::path::Path) {} + +#[cfg(unix)] +fn set_executable(path: &std::path::Path) { + use std::os::unix::fs::PermissionsExt; + let mut perms = std::fs::metadata(path) + .expect("metadata fake CLI") + .permissions(); + perms.set_mode(0o755); + std::fs::set_permissions(path, perms).expect("chmod fake CLI"); +} + +#[cfg(not(unix))] +fn set_executable(_path: &std::path::Path) { + // Sur les plateformes non-Unix le harnais s'appuiera sur un fake CLI adapté + // (ex. `.cmd`) ; non requis pour la CI Linux actuelle. +} + +// --------------------------------------------------------------------------- +// Harnais de conformité de port (Liskov) — réutilisable Claude ET Codex +// --------------------------------------------------------------------------- + +#[cfg(test)] +pub(crate) mod harness { + use std::sync::Arc; + + use domain::ports::{AgentSession, ReplyEvent}; + + /// Asserte le **contrat [`AgentSession`]** sur une session déjà construite + /// derrière un fake CLI dont le script produit ≥0 deltas/activités puis **un** + /// `Final` portant `expected_final`, et dont l'init assigne `expected_conv_id`. + /// + /// Vérifie (substituabilité Liskov, §17.2) : + /// 1. `send` émet une séquence de deltas/activités **puis exactement un** `Final` ; + /// 2. après le `Final` le flux est **clos** (plus aucun événement) ; + /// 3. le `Final` porte bien `expected_final` ; + /// 4. `conversation_id()` devient `Some(expected_conv_id)` après le tour assignant ; + /// 5. `shutdown()` réussit et est **idempotent** (deux appels OK). + pub async fn assert_agent_session_contract( + session: Arc, + expected_conv_id: &str, + expected_final: &str, + ) { + // Avant tout tour : aucun id de conversation assigné. + assert_eq!( + session.conversation_id(), + None, + "conversation_id doit être None avant le premier tour" + ); + + let stream = session.send("salut").await.expect("send doit réussir"); + let events: Vec = stream.collect(); + + // (1)+(2) : exactement un Final, en dernière position. + let final_count = events + .iter() + .filter(|e| matches!(e, ReplyEvent::Final { .. })) + .count(); + assert_eq!(final_count, 1, "le flux doit porter EXACTEMENT un Final"); + match events.last() { + Some(ReplyEvent::Final { content }) => { + // (3) + assert_eq!(content, expected_final, "contenu Final inattendu"); + } + other => panic!("le dernier événement doit être Final, vu: {other:?}"), + } + // Les événements avant le Final ne sont que des deltas / activités. + for e in &events[..events.len() - 1] { + assert!( + matches!( + e, + ReplyEvent::TextDelta { .. } | ReplyEvent::ToolActivity { .. } + ), + "avant le Final, seuls deltas/activités sont permis, vu: {e:?}" + ); + } + + // (4) : id de conversation capté après le tour assignant. + assert_eq!( + session.conversation_id().as_deref(), + Some(expected_conv_id), + "conversation_id doit être assigné après le premier tour" + ); + + // (5) : shutdown réussit et est idempotent. + session.shutdown().await.expect("shutdown doit réussir"); + session + .shutdown() + .await + .expect("shutdown doit être idempotent"); + } +} diff --git a/crates/infrastructure/src/session/factory.rs b/crates/infrastructure/src/session/factory.rs new file mode 100644 index 0000000..105aa38 --- /dev/null +++ b/crates/infrastructure/src/session/factory.rs @@ -0,0 +1,86 @@ +//! [`StructuredSessionFactory`] — la fabrique [`AgentSessionFactory`] qui **route un +//! profil vers le bon adapter** structuré (ARCHITECTURE §17.2) selon +//! `profile.structured_adapter` (§17.3). Agrège Claude + Codex derrière une seule +//! surface ; aucun type concret ne franchit la frontière domaine (seuls +//! `Arc` sortent). +//! +//! Open/Closed : ajouter un moteur structuré = ajouter un adapter + une variante +//! [`StructuredAdapter`] + un bras de `match` ici. Le cœur ne bouge pas. + +use std::sync::Arc; + +use async_trait::async_trait; + +use domain::ports::{ + AgentSession, AgentSessionError, AgentSessionFactory, PreparedContext, SessionPlan, +}; +use domain::profile::{AgentProfile, StructuredAdapter}; +use domain::project::ProjectPath; +use domain::SessionId; + +use super::claude::ClaudeSdkSession; +use super::codex::CodexExecSession; + +/// Fabrique infra des sessions structurées, sélectionnée par le profil. +/// +/// Sans état : elle instancie l'adapter au vol depuis le profil (le binaire à +/// lancer = `profile.command`), de sorte qu'un seul exemplaire injecté au +/// composition root sert tous les agents (jumeau de `CliAgentRuntime`). +#[derive(Debug, Default, Clone, Copy)] +pub struct StructuredSessionFactory; + +impl StructuredSessionFactory { + /// Construit la fabrique. + #[must_use] + pub const fn new() -> Self { + Self + } +} + +/// Dérive l'[`SessionPlan`] le `seed` de reprise : seul [`SessionPlan::Resume`] +/// amorce l'adapter avec un id de conversation existant ; `Assign`/`None` partent +/// d'une conversation neuve (l'id sera capté au premier tour). +fn seed_conversation_id(session: &SessionPlan) -> Option { + match session { + SessionPlan::Resume { conversation_id } => Some(conversation_id.clone()), + SessionPlan::None | SessionPlan::Assign { .. } => None, + } +} + +#[async_trait] +impl AgentSessionFactory for StructuredSessionFactory { + fn supports(&self, profile: &AgentProfile) -> bool { + profile.structured_adapter.is_some() + } + + async fn start( + &self, + profile: &AgentProfile, + _ctx: &PreparedContext, + cwd: &ProjectPath, + session: &SessionPlan, + ) -> Result, AgentSessionError> { + let adapter = profile.structured_adapter.ok_or_else(|| { + AgentSessionError::Start(format!( + "le profil « {} » n'a pas d'adapter structuré", + profile.name + )) + })?; + + let id = SessionId::new_random(); + let command = profile.command.clone(); + let cwd = cwd.as_str().to_owned(); + let seed = seed_conversation_id(session); + + // NOTE : le contexte (`_ctx`) est injecté par `LaunchAgent` (D3) via le + // convention file dans le run dir *avant* l'appel à la factory (le `.md` est + // déjà écrit) ; l'adapter n'a donc qu'à lancer la CLI dans ce cwd. Aucune + // injection supplémentaire n'incombe ici en mode structuré (la CLI lit son + // fichier conventionnel — CLAUDE.md / AGENTS.md — depuis le cwd). + let session: Arc = match adapter { + StructuredAdapter::Claude => Arc::new(ClaudeSdkSession::new(id, command, cwd, seed)), + StructuredAdapter::Codex => Arc::new(CodexExecSession::new(id, command, cwd, seed)), + }; + Ok(session) + } +} diff --git a/crates/infrastructure/src/session/mod.rs b/crates/infrastructure/src/session/mod.rs new file mode 100644 index 0000000..76705f7 --- /dev/null +++ b/crates/infrastructure/src/session/mod.rs @@ -0,0 +1,941 @@ +//! Adapters d'**exécution structurée** des agents IA (ARCHITECTURE §17.2), pair de +//! [`crate::runtime`] (TUI/PTY) et [`crate::pty`]. Implémentent le port domaine +//! [`domain::ports::AgentSession`] et la fabrique [`domain::ports::AgentSessionFactory`]. +//! +//! # Principe directeur (CRUCIAL, §17.2) +//! +//! Le **parsing du format de sortie de chaque CLI est ISOLÉ** dans une fonction pure +//! dédiée par adapter ([`claude::parse_event`], [`codex::parse_event`]), **séparée** +//! de la machinerie de process ([`process`]). Les **vrais** formats Claude/Codex +//! seront confirmés par les spikes **S1** (Claude) et **S2** (Codex) ; quand on les +//! aura, **seules ces fonctions de parsing changeront**, pas la machinerie. +//! +//! # Composants +//! +//! - [`process`] : machinerie de process **paramétrable par la commande** (spawn, +//! pipes, drain ligne-à-ligne, timeout) — substituable par un fake CLI en test. +//! - [`claude::ClaudeSdkSession`] / [`codex::CodexExecSession`] : les deux adapters. +//! - [`factory::StructuredSessionFactory`] : route un profil vers le bon adapter. +//! - [`conformance`] : **fake CLI** scriptable + **harnais de conformité** (Liskov), +//! réutilisable pour valider le contrat de port des deux adapters hors-réseau. + +pub mod claude; +pub mod codex; +pub mod conformance; +pub mod factory; +pub mod process; + +pub use claude::ClaudeSdkSession; +pub use codex::CodexExecSession; +pub use conformance::FakeCli; +pub use factory::StructuredSessionFactory; + +#[cfg(test)] +mod tests { + use std::sync::Arc; + use std::time::Duration; + + use domain::ids::ProfileId; + use domain::ports::{ + AgentSession, AgentSessionError, AgentSessionFactory, ContextInjectionPlan, + PreparedContext, ReplyEvent, SessionPlan, + }; + use domain::profile::{AgentProfile, ContextInjection, StructuredAdapter}; + use domain::project::ProjectPath; + use domain::{MarkdownDoc, SessionId}; + + use super::claude::{self, ClaudeSdkSession}; + use super::codex::{self, CodexExecSession}; + use super::conformance::harness::assert_agent_session_contract; + use super::conformance::FakeCli; + use super::factory::StructuredSessionFactory; + use super::process::run_turn; + + // -- Helpers ---------------------------------------------------------- + + fn prepared_ctx() -> PreparedContext { + PreparedContext { + content: MarkdownDoc::new("# ctx"), + relative_path: "CLAUDE.md".to_owned(), + } + } + + fn cwd() -> ProjectPath { + ProjectPath::new("/").expect("cwd valide") + } + + fn structured_profile(adapter: StructuredAdapter, command: &str) -> AgentProfile { + AgentProfile::new( + ProfileId::new_random(), + "Profil structuré", + command, + Vec::new(), + ContextInjection::convention_file("CLAUDE.md").expect("convention file valide"), + None, + "{agentRunDir}", + None, + ) + .expect("profil valide") + .with_structured_adapter(adapter) + } + + // -- Machinerie de process (paramétrable, fake CLI) ------------------- + + #[tokio::test] + async fn run_turn_drains_every_line_in_order() { + let fake = FakeCli::printing(&["ligne-1", "ligne-2", "ligne-3"]); + let lines = run_turn(&fake.spawn_line(), None) + .await + .expect("run_turn réussit"); + assert_eq!(lines, vec!["ligne-1", "ligne-2", "ligne-3"]); + } + + #[tokio::test] + async fn run_turn_unknown_binary_yields_start_error() { + let spec = super::process::SpawnLine { + command: "/binaire/qui/n/existe/pas/idea-xyz".to_owned(), + args: Vec::new(), + cwd: "/".to_owned(), + env: Vec::new(), + stdin: None, + }; + let err = run_turn(&spec, None).await.expect_err("doit échouer"); + assert!(matches!(err, AgentSessionError::Start(_)), "vu: {err:?}"); + } + + // -- parse_event Claude (schéma SUPPOSÉ S1) --------------------------- + + #[test] + fn claude_parse_init_captures_session_id_without_event() { + let parsed = + claude::parse_event(r#"{"type":"system","subtype":"init","session_id":"conv-123"}"#) + .expect("parse ok"); + assert_eq!(parsed.session_id.as_deref(), Some("conv-123")); + assert_eq!(parsed.event, None); + } + + #[test] + fn claude_parse_assistant_text_and_tool() { + let text = claude::parse_event( + r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"bonjour"}]}}"#, + ) + .expect("parse ok"); + assert_eq!( + text.event, + Some(ReplyEvent::TextDelta { + text: "bonjour".to_owned() + }) + ); + + let tool = claude::parse_event( + r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"tool_use","name":"Read"}]}}"#, + ) + .expect("parse ok"); + assert_eq!( + tool.event, + Some(ReplyEvent::ToolActivity { + label: "Read".to_owned() + }) + ); + } + + #[test] + fn claude_parse_result_is_final() { + let parsed = claude::parse_event( + r#"{"type":"result","subtype":"success","result":"réponse finale","session_id":"conv-123"}"#, + ) + .expect("parse ok"); + assert_eq!( + parsed.event, + Some(ReplyEvent::Final { + content: "réponse finale".to_owned() + }) + ); + assert_eq!(parsed.session_id.as_deref(), Some("conv-123")); + } + + #[test] + fn claude_parse_broken_json_is_decode_error_no_raw_leak() { + let err = claude::parse_event("{ pas du json").expect_err("doit échouer"); + match err { + AgentSessionError::Decode(msg) => { + assert!( + !msg.contains("pas du json"), + "le JSON brut ne doit pas fuir" + ); + } + other => panic!("attendu Decode, vu: {other:?}"), + } + } + + #[test] + fn claude_parse_empty_and_unknown_lines_are_ignored() { + assert_eq!(claude::parse_event("").expect("ok"), Default::default()); + let unknown = claude::parse_event(r#"{"type":"telemetry","x":1}"#).expect("ok ignoré"); + assert_eq!(unknown.event, None); + } + + // -- parse_event Codex (schéma SUPPOSÉ S2) ---------------------------- + + #[test] + fn codex_parse_session_message_and_final() { + let sess = + codex::parse_event(r#"{"type":"session","conversation_id":"cx-9"}"#).expect("ok"); + assert_eq!(sess.conversation_id.as_deref(), Some("cx-9")); + assert_eq!(sess.event, None); + + let msg = codex::parse_event(r#"{"type":"message","text":"salut"}"#).expect("ok"); + assert_eq!( + msg.event, + Some(ReplyEvent::TextDelta { + text: "salut".to_owned() + }) + ); + + let fin = codex::parse_event(r#"{"type":"result","text":"fini"}"#).expect("ok"); + assert_eq!( + fin.event, + Some(ReplyEvent::Final { + content: "fini".to_owned() + }) + ); + } + + #[test] + fn codex_parse_broken_json_is_decode_error() { + let err = codex::parse_event("<<<").expect_err("doit échouer"); + assert!(matches!(err, AgentSessionError::Decode(_)), "vu: {err:?}"); + } + + // -- Conformité de port (Liskov) — Claude ET Codex -------------------- + + /// Script Claude (schéma SUPPOSÉ S1) : init → texte → tool_use → result. + fn claude_script() -> Vec<&'static str> { + vec![ + r#"{"type":"system","subtype":"init","session_id":"claude-conv-1"}"#, + r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"un "}]}}"#, + r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"tool_use","name":"Read"}]}}"#, + r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"deux"}]}}"#, + r#"{"type":"result","subtype":"success","result":"réponse Claude","session_id":"claude-conv-1"}"#, + ] + } + + /// Script Codex (schéma SUPPOSÉ S2) : session → message → tool → result. + fn codex_script() -> Vec<&'static str> { + vec![ + r#"{"type":"session","conversation_id":"codex-conv-1"}"#, + r#"{"type":"delta","text":"trav"}"#, + r#"{"type":"tool_call","name":"bash"}"#, + r#"{"type":"result","text":"réponse Codex"}"#, + ] + } + + #[tokio::test] + async fn claude_session_respects_port_contract() { + let fake = FakeCli::printing(&claude_script()); + let session: Arc = Arc::new(ClaudeSdkSession::new( + SessionId::new_random(), + fake.command(), + "/", + None, + )); + assert_agent_session_contract(session, "claude-conv-1", "réponse Claude").await; + } + + #[tokio::test] + async fn codex_session_respects_port_contract() { + let fake = FakeCli::printing(&codex_script()); + let session: Arc = Arc::new(CodexExecSession::new( + SessionId::new_random(), + fake.command(), + "/", + None, + )); + assert_agent_session_contract(session, "codex-conv-1", "réponse Codex").await; + } + + /// Le flux est **clos** après le `Final` : drainé une fois, il ne reproduit + /// rien (l'incarnation « un run par tour » est intrinsèquement bornée). + #[tokio::test] + async fn stream_is_closed_after_final() { + let fake = FakeCli::printing(&claude_script()); + let session = ClaudeSdkSession::new(SessionId::new_random(), fake.command(), "/", None); + let stream = session.send("x").await.expect("send ok"); + let events: Vec<_> = stream.collect(); + let after_final = events + .iter() + .skip_while(|e| !matches!(e, ReplyEvent::Final { .. })) + .skip(1) + .count(); + assert_eq!(after_final, 0, "aucun événement après le Final"); + } + + /// Un JSON cassé **au milieu du flux** remonte `Decode` (jamais de panic). + #[tokio::test] + async fn broken_line_in_stream_yields_decode() { + let fake = FakeCli::printing(&[ + r#"{"type":"system","subtype":"init","session_id":"c"}"#, + "{ ceci n'est pas du json", + ]); + let session = ClaudeSdkSession::new(SessionId::new_random(), fake.command(), "/", None); + match session.send("x").await { + Err(AgentSessionError::Decode(_)) => {} + Err(other) => panic!("attendu Decode, vu: {other:?}"), + Ok(_) => panic!("attendu une erreur Decode, vu un flux"), + } + } + + // -- Factory : routage par structured_adapter ------------------------ + + #[tokio::test] + async fn factory_supports_only_structured_profiles() { + let factory = StructuredSessionFactory::new(); + let claude = structured_profile(StructuredAdapter::Claude, "claude"); + let codex = structured_profile(StructuredAdapter::Codex, "codex"); + let tui = AgentProfile::new( + ProfileId::new_random(), + "Gemini", + "gemini", + Vec::new(), + ContextInjection::convention_file("GEMINI.md").expect("valide"), + None, + "{agentRunDir}", + None, + ) + .expect("profil valide"); // pas de structured_adapter + + assert!(factory.supports(&claude)); + assert!(factory.supports(&codex)); + assert!(!factory.supports(&tui)); + } + + #[tokio::test] + async fn factory_routes_claude_and_codex() { + let factory = StructuredSessionFactory::new(); + let fake = FakeCli::printing(&claude_script()); + + // Claude : la session démarre et respecte le contrat via le fake CLI. + let claude = structured_profile(StructuredAdapter::Claude, &fake.command()); + let session = factory + .start(&claude, &prepared_ctx(), &cwd(), &SessionPlan::None) + .await + .expect("start Claude ok"); + let content = drain_final(session.as_ref()).await; + assert_eq!(content, "réponse Claude"); + + // Codex : routé vers l'adapter Codex (id de session distinct, démarrage ok). + let fake_cx = FakeCli::printing(&codex_script()); + let codex = structured_profile(StructuredAdapter::Codex, &fake_cx.command()); + let session_cx = factory + .start(&codex, &prepared_ctx(), &cwd(), &SessionPlan::None) + .await + .expect("start Codex ok"); + let content_cx = drain_final(session_cx.as_ref()).await; + assert_eq!(content_cx, "réponse Codex"); + } + + #[tokio::test] + async fn factory_resume_seeds_conversation_id() { + let factory = StructuredSessionFactory::new(); + let fake = FakeCli::printing(&claude_script()); + let claude = structured_profile(StructuredAdapter::Claude, &fake.command()); + let session = factory + .start( + &claude, + &prepared_ctx(), + &cwd(), + &SessionPlan::Resume { + conversation_id: "repris-42".to_owned(), + }, + ) + .await + .expect("start resume ok"); + // L'id de reprise amorce la session avant tout tour (pivot model-agnostic). + assert_eq!(session.conversation_id().as_deref(), Some("repris-42")); + } + + async fn drain_final(session: &dyn AgentSession) -> String { + let stream = session.send("x").await.expect("send ok"); + for event in stream { + if let ReplyEvent::Final { content } = event { + return content; + } + } + panic!("aucun Final"); + } + + // -- Sanity : le PreparedContext et le ContextInjectionPlan ne sont pas + // requis par l'adapter structuré (le .md est déjà écrit par LaunchAgent). + #[test] + fn prepared_context_is_carried_not_required_by_adapter() { + // Documentation exécutable : un plan de fichier existe côté runtime PTY, + // mais l'adapter structuré ne le consomme pas (la CLI lit son convention + // file depuis le cwd). On vérifie juste que le type compose. + let _plan = ContextInjectionPlan::File { + target: "CLAUDE.md".to_owned(), + }; + let _ = prepared_ctx(); + } + + /// Le timeout de la machinerie tue le process et remonte `Timeout` (un fake CLI + /// qui dort plus longtemps que la borne). + #[tokio::test] + async fn run_turn_honours_timeout() { + // Fake CLI qui dort 5s avant d'imprimer : la borne 50ms doit déclencher. + let mut path = std::env::temp_dir(); + // Nom unique (compteur atomique) pour éviter toute collision entre exécutions + // parallèles répétées de la suite. + use std::sync::atomic::{AtomicU64, Ordering}; + static SLOW_COUNTER: AtomicU64 = AtomicU64::new(0); + let n = SLOW_COUNTER.fetch_add(1, Ordering::Relaxed); + path.push(format!("idea-fake-slow-{}-{n}", std::process::id())); + // `File::create` + `sync_all` + `drop` ferme le descripteur en écriture AVANT + // l'exécution : sinon `execve` peut retourner `ETXTBSY` sous charge parallèle. + { + use std::io::Write as _; + let mut f = std::fs::File::create(&path).expect("write"); + f.write_all(b"#!/bin/sh\nsleep 5\nprintf 'tard\\n'\n") + .expect("write"); + f.sync_all().expect("sync"); + } + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + let mut perms = std::fs::metadata(&path).unwrap().permissions(); + perms.set_mode(0o755); + std::fs::set_permissions(&path, perms).unwrap(); + } + // Garantit que le binaire est exec-ready (plus de `ETXTBSY`) avant le spawn + // mesuré : on probe en boucle, mais on **tue immédiatement** l'enfant (il + // dormirait 5s) — on ne veut prouver que l'exécutabilité, pas attendre. + #[cfg(unix)] + { + use std::process::{Command, Stdio}; + use std::time::{Duration, Instant}; + let deadline = Instant::now() + Duration::from_secs(5); + loop { + match Command::new(&path) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + { + Ok(mut child) => { + let _ = child.kill(); + let _ = child.wait(); + break; + } + Err(e) if e.raw_os_error() == Some(26) && Instant::now() < deadline => { + std::thread::sleep(Duration::from_millis(2)); + } + Err(_) => break, + } + } + } + let spec = super::process::SpawnLine { + command: path.to_string_lossy().into_owned(), + args: Vec::new(), + cwd: "/".to_owned(), + env: Vec::new(), + stdin: None, + }; + let err = run_turn(&spec, Some(Duration::from_millis(50))) + .await + .expect_err("doit expirer"); + assert!(matches!(err, AgentSessionError::Timeout), "vu: {err:?}"); + let _ = std::fs::remove_file(&path); + } + + // ===================================================================== + // DURCISSEMENT QA (lot D2) — couvre les axes non couverts par les tests + // initiaux. Tout passe par le FakeCli (jamais le vrai claude/codex). + // ===================================================================== + + // -- Helper : fake CLI qui enregistre son argv dans un fichier sidecar, + // puis rejoue un script de lignes. Permet de PROUVER que la commande + // générée porte bien le flag de reprise (`--resume `). + fn make_recording_fake(script: &[&str]) -> (String, std::path::PathBuf) { + use std::io::Write as _; + use std::sync::atomic::{AtomicU64, Ordering}; + static C: AtomicU64 = AtomicU64::new(0); + let n = C.fetch_add(1, Ordering::Relaxed); + let mut bin = std::env::temp_dir(); + bin.push(format!("idea-rec-cli-{}-{n}", std::process::id())); + let mut argv = std::env::temp_dir(); + argv.push(format!("idea-rec-argv-{}-{n}", std::process::id())); + + let mut s = String::from("#!/bin/sh\n"); + // Enregistre chaque argument sur sa propre ligne dans le sidecar. + s.push_str(&format!( + "for a in \"$@\"; do printf '%s\\n' \"$a\" >> '{}'; done\n", + argv.display() + )); + for line in script { + s.push_str("printf '%s\\n' "); + // réutilise le quoting de conformance via un quoting local simple. + s.push('\''); + s.push_str(&line.replace('\'', "'\\''")); + s.push_str("'\n"); + } + { + let mut f = std::fs::File::create(&bin).expect("create rec fake"); + f.write_all(s.as_bytes()).expect("write rec fake"); + f.sync_all().expect("sync rec fake"); + } + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + let mut p = std::fs::metadata(&bin).unwrap().permissions(); + p.set_mode(0o755); + std::fs::set_permissions(&bin, p).unwrap(); + } + super::conformance::wait_until_executable(&bin); + (bin.to_string_lossy().into_owned(), argv) + } + + // ---- Claude parse_event : plusieurs blocs, robustesse --------------- + + /// Plusieurs `tool_use` / blocs texte dans des messages successifs : chaque + /// message produit UN événement (le premier bloc pertinent). On documente la + /// limite connue : un message multi-blocs ne rend que son **premier** bloc. + #[test] + fn claude_multiple_messages_each_yield_one_event() { + let t1 = claude::parse_event( + r#"{"type":"assistant","message":{"content":[{"type":"text","text":"a"}]}}"#, + ) + .unwrap(); + let t2 = claude::parse_event( + r#"{"type":"assistant","message":{"content":[{"type":"tool_use","name":"Bash"}]}}"#, + ) + .unwrap(); + assert_eq!(t1.event, Some(ReplyEvent::TextDelta { text: "a".into() })); + assert_eq!( + t2.event, + Some(ReplyEvent::ToolActivity { + label: "Bash".into() + }) + ); + } + + /// LIMITE CONNUE (à arbitrer dev/archi) : un **seul** message portant plusieurs + /// blocs `text`/`tool_use` ne rend que le PREMIER bloc pertinent — les blocs + /// suivants sont perdus. Ce test PINNE le comportement actuel (pas un échec). + #[test] + fn claude_multiblock_message_keeps_only_first_block_known_limit() { + let parsed = claude::parse_event( + r#"{"type":"assistant","message":{"content":[ + {"type":"text","text":"un"}, + {"type":"tool_use","name":"Read"}, + {"type":"text","text":"deux"}]}}"#, + ) + .unwrap(); + // Comportement ACTUEL : seul le premier bloc (`text:"un"`) est émis. + assert_eq!( + parsed.event, + Some(ReplyEvent::TextDelta { text: "un".into() }) + ); + } + + /// `tool_use` sans `name` ⇒ label de repli « outil » (jamais de panic). + #[test] + fn claude_tool_use_without_name_falls_back() { + let parsed = claude::parse_event( + r#"{"type":"assistant","message":{"content":[{"type":"tool_use"}]}}"#, + ) + .unwrap(); + assert_eq!( + parsed.event, + Some(ReplyEvent::ToolActivity { + label: "outil".into() + }) + ); + } + + /// `result` sans champ `result` ⇒ aucun event (pas de panic, pas de Final vide + /// fabriqué). Documente la robustesse du parser. + #[test] + fn claude_result_without_content_yields_no_event() { + let parsed = claude::parse_event(r#"{"type":"result","subtype":"success"}"#).unwrap(); + assert_eq!(parsed.event, None); + } + + /// Ligne whitespace-only (espaces/tabs) ⇒ ignorée comme une ligne vide. + #[test] + fn claude_whitespace_line_is_ignored() { + assert_eq!(claude::parse_event(" \t ").unwrap(), Default::default()); + } + + /// JSON valide mais non-objet (tableau, nombre) ⇒ pas de type ⇒ ignoré, jamais + /// de panic, jamais de Decode. + #[test] + fn claude_valid_non_object_json_is_ignored() { + assert_eq!(claude::parse_event("[1,2,3]").unwrap().event, None); + assert_eq!(claude::parse_event("42").unwrap().event, None); + } + + // ---- Codex parse_event (schéma SUPPOSÉ S2 — NON CONFIRMÉ) ----------- + // NB: toutes ces assertions dépendent du format présumé S2. Le contrat de + // sortie (ReplyEvent) est stable ; le mapping changera au spike S2. + + /// [S2 présumé] `tool`/`tool_call` ⇒ ToolActivity ; `name` manquant ⇒ « outil ». + #[test] + fn codex_tool_activity_and_fallback_s2() { + let t = codex::parse_event(r#"{"type":"tool","name":"grep"}"#).unwrap(); + assert_eq!( + t.event, + Some(ReplyEvent::ToolActivity { + label: "grep".into() + }) + ); + let t2 = codex::parse_event(r#"{"type":"tool_call"}"#).unwrap(); + assert_eq!( + t2.event, + Some(ReplyEvent::ToolActivity { + label: "outil".into() + }) + ); + } + + /// [S2 présumé] `delta` et `message` produisent tous deux un TextDelta. + #[test] + fn codex_delta_and_message_both_text_s2() { + let d = codex::parse_event(r#"{"type":"delta","text":"x"}"#).unwrap(); + let m = codex::parse_event(r#"{"type":"message","text":"y"}"#).unwrap(); + assert_eq!(d.event, Some(ReplyEvent::TextDelta { text: "x".into() })); + assert_eq!(m.event, Some(ReplyEvent::TextDelta { text: "y".into() })); + } + + /// [S2 présumé] `final` (alias de `result`) avec repli sur le champ `output`. + #[test] + fn codex_final_alias_and_output_fallback_s2() { + let f = codex::parse_event(r#"{"type":"final","output":"sortie"}"#).unwrap(); + assert_eq!( + f.event, + Some(ReplyEvent::Final { + content: "sortie".into() + }) + ); + } + + /// [S2 présumé] id de conversation : repli `id` quand `conversation_id` absent. + #[test] + fn codex_conversation_id_falls_back_to_id_s2() { + let p = codex::parse_event(r#"{"type":"session","id":"cx-7"}"#).unwrap(); + assert_eq!(p.conversation_id.as_deref(), Some("cx-7")); + assert_eq!(p.event, None); + } + + /// [S2 présumé] ligne vide / type inconnu ⇒ ignorée sans erreur. + #[test] + fn codex_empty_and_unknown_ignored_s2() { + assert_eq!(codex::parse_event("").unwrap(), Default::default()); + assert_eq!( + codex::parse_event(r#"{"type":"heartbeat"}"#).unwrap().event, + None + ); + } + + // ---- Machinerie process via FakeCli --------------------------------- + + /// Deltas PUIS Final : le flux ne contient rien après le `Final` (déjà couvert + /// pour Claude ; ici on le prouve aussi pour Codex, substituabilité Liskov). + #[tokio::test] + async fn codex_stream_closed_after_final() { + let fake = FakeCli::printing(&[ + r#"{"type":"session","conversation_id":"c"}"#, + r#"{"type":"delta","text":"a"}"#, + r#"{"type":"result","text":"fin"}"#, + ]); + let s = CodexExecSession::new(SessionId::new_random(), fake.command(), "/", None); + let events: Vec<_> = s.send("x").await.expect("send").collect(); + let after = events + .iter() + .skip_while(|e| !matches!(e, ReplyEvent::Final { .. })) + .skip(1) + .count(); + assert_eq!(after, 0); + } + + /// LIMITE/ÉCART (à arbitrer) : un flux SANS `Final` ne provoque PAS d'erreur au + /// niveau de l'adapter — `send()` renvoie Ok avec uniquement des deltas et AUCUN + /// `Final`. Le §17.9 D2 mentionne « flux sans Final ⇒ Io » ; l'adapter actuel ne + /// l'applique pas (c'est `send_blocking` côté application, lot D1, qui transforme + /// l'absence de Final en Timeout). Ce test PINNE le comportement réel observé. + #[tokio::test] + async fn stream_without_final_is_silently_ok_at_adapter_level() { + let fake = FakeCli::printing(&[ + r#"{"type":"system","subtype":"init","session_id":"c"}"#, + r#"{"type":"assistant","message":{"content":[{"type":"text","text":"a"}]}}"#, + ]); + let s = ClaudeSdkSession::new(SessionId::new_random(), fake.command(), "/", None); + let events: Vec<_> = s.send("x").await.expect("send ok").collect(); + let finals = events + .iter() + .filter(|e| matches!(e, ReplyEvent::Final { .. })) + .count(); + assert_eq!(finals, 0, "comportement actuel: aucun Final fabriqué"); + // L'id de conversation est tout de même capté (init lu). + assert_eq!(s.conversation_id().as_deref(), Some("c")); + } + + /// EOF immédiat (binaire qui n'imprime rien) ⇒ flux vide, pas d'erreur, pas de + /// panic. La machinerie draine proprement un stdout vide. + #[tokio::test] + async fn run_turn_empty_output_is_ok() { + let fake = FakeCli::printing(&[]); + let lines = run_turn(&fake.spawn_line(), None).await.expect("ok"); + assert!(lines.is_empty()); + } + + /// stdin fourni : la machinerie l'écrit et ferme le pipe (EOF) sans bloquer. + #[tokio::test] + async fn run_turn_writes_stdin_then_eof() { + let fake = FakeCli::printing(&["pong"]); + let mut spec = fake.spawn_line(); + spec.stdin = Some("ping".to_owned()); + let lines = run_turn(&spec, None).await.expect("ok"); + assert_eq!(lines, vec!["pong"]); + } + + /// Timeout via la machinerie + FakeCli lent (≈ le test existant, mais bâti sur + /// un fake qui dort) : borne courte ⇒ `Timeout`. + #[tokio::test] + async fn run_turn_timeout_on_slow_fake() { + // Fake dormeur déterministe (sleep 5s) borné par un timeout de 50ms. + let mut path = std::env::temp_dir(); + use std::sync::atomic::{AtomicU64, Ordering}; + static C: AtomicU64 = AtomicU64::new(0); + path.push(format!( + "idea-slow2-{}-{}", + std::process::id(), + C.fetch_add(1, Ordering::Relaxed) + )); + { + use std::io::Write as _; + let mut f = std::fs::File::create(&path).unwrap(); + f.write_all(b"#!/bin/sh\nsleep 5\n").unwrap(); + f.sync_all().unwrap(); + } + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + let mut p = std::fs::metadata(&path).unwrap().permissions(); + p.set_mode(0o755); + std::fs::set_permissions(&path, p).unwrap(); + } + // Probe non-bloquant (kill immédiat) pour écarter ETXTBSY. + #[cfg(unix)] + { + use std::process::{Command, Stdio}; + use std::time::Instant; + let dl = Instant::now() + Duration::from_secs(5); + loop { + match Command::new(&path) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + { + Ok(mut c) => { + let _ = c.kill(); + let _ = c.wait(); + break; + } + Err(e) if e.raw_os_error() == Some(26) && Instant::now() < dl => { + std::thread::sleep(Duration::from_millis(2)); + } + Err(_) => break, + } + } + } + let spec = super::process::SpawnLine { + command: path.to_string_lossy().into_owned(), + args: Vec::new(), + cwd: "/".to_owned(), + env: Vec::new(), + stdin: None, + }; + let err = run_turn(&spec, Some(Duration::from_millis(50))) + .await + .expect_err("doit expirer"); + assert!(matches!(err, AgentSessionError::Timeout), "vu: {err:?}"); + let _ = std::fs::remove_file(&path); + } + + // ---- Adapters : reprise + commande générée porte le flag -------------- + + /// PROUVE que la commande réellement lancée porte `--resume ` quand la + /// session a été amorcée en reprise (via le sidecar argv du fake enregistreur). + #[tokio::test] + async fn claude_resume_command_carries_resume_flag() { + let (cmd, argv) = make_recording_fake(&[ + r#"{"type":"result","subtype":"success","result":"ok","session_id":"resume-id"}"#, + ]); + let session = ClaudeSdkSession::new( + SessionId::new_random(), + cmd.clone(), + "/", + Some("resume-id".to_owned()), + ); + // conversation_id amorcé avant tout tour. + assert_eq!(session.conversation_id().as_deref(), Some("resume-id")); + let _ = session.send("salut").await.expect("send ok"); + + let recorded = std::fs::read_to_string(&argv).expect("argv enregistré"); + let args: Vec<&str> = recorded.lines().collect(); + assert!( + args.contains(&"--resume"), + "argv doit porter --resume, vu: {args:?}" + ); + assert!( + args.contains(&"resume-id"), + "argv doit porter l'id de reprise, vu: {args:?}" + ); + assert!( + args.contains(&"salut"), + "argv doit porter le prompt, vu: {args:?}" + ); + let _ = std::fs::remove_file(&cmd); + let _ = std::fs::remove_file(&argv); + } + + /// Idem Codex : `codex exec --resume ` (schéma S2 présumé). + #[tokio::test] + async fn codex_resume_command_carries_resume_flag_s2() { + let (cmd, argv) = + make_recording_fake(&[r#"{"type":"result","text":"ok","conversation_id":"cx-id"}"#]); + let session = CodexExecSession::new( + SessionId::new_random(), + cmd.clone(), + "/", + Some("cx-id".to_owned()), + ); + let _ = session.send("vas-y").await.expect("send ok"); + let recorded = std::fs::read_to_string(&argv).expect("argv"); + let args: Vec<&str> = recorded.lines().collect(); + assert!(args.contains(&"exec"), "vu: {args:?}"); + assert!(args.contains(&"--resume"), "vu: {args:?}"); + assert!(args.contains(&"cx-id"), "vu: {args:?}"); + assert!(args.contains(&"vas-y"), "vu: {args:?}"); + let _ = std::fs::remove_file(&cmd); + let _ = std::fs::remove_file(&argv); + } + + /// Une conversation NEUVE (pas de seed) NE porte PAS `--resume` au premier tour, + /// mais le capte après (init) ⇒ le SECOND tour, lui, porte `--resume`. + #[tokio::test] + async fn claude_new_then_resume_flag_appears_on_second_turn() { + let (cmd, argv) = make_recording_fake(&[ + r#"{"type":"system","subtype":"init","session_id":"captured-1"}"#, + r#"{"type":"result","subtype":"success","result":"r","session_id":"captured-1"}"#, + ]); + let session = ClaudeSdkSession::new(SessionId::new_random(), cmd.clone(), "/", None); + assert_eq!(session.conversation_id(), None); + let _ = session.send("t1").await.expect("t1"); + assert_eq!(session.conversation_id().as_deref(), Some("captured-1")); + let _ = session.send("t2").await.expect("t2"); + + let recorded = std::fs::read_to_string(&argv).unwrap(); + // Le sidecar accumule les deux tours. Le 1er tour ne doit PAS avoir d'id avant + // capture ; après capture le 2e tour porte --resume captured-1. On vérifie la + // présence globale (les deux tours sont concaténés). + assert!( + recorded.contains("--resume"), + "2e tour doit porter --resume" + ); + assert!(recorded.contains("captured-1")); + let _ = std::fs::remove_file(&cmd); + let _ = std::fs::remove_file(&argv); + } + + // ---- Factory : reprise amorce le seed + route correctement ----------- + + /// Reprise via la factory : `SessionPlan::Resume` amorce le seed, et la session + /// résultante porte bien l'id AVANT tout tour (déjà partiellement couvert ; ici + /// on couvre AUSSI Codex). + #[tokio::test] + async fn factory_resume_seeds_codex() { + let factory = StructuredSessionFactory::new(); + let fake = FakeCli::printing(&[r#"{"type":"result","text":"ok"}"#]); + let codex = structured_profile(StructuredAdapter::Codex, &fake.command()); + let session = factory + .start( + &codex, + &prepared_ctx(), + &cwd(), + &SessionPlan::Resume { + conversation_id: "cx-resume".to_owned(), + }, + ) + .await + .expect("start resume codex"); + assert_eq!(session.conversation_id().as_deref(), Some("cx-resume")); + } + + /// `SessionPlan::Assign` (assigne un id côté IdeA mais conversation moteur neuve) + /// ⇒ pas de seed moteur (l'id moteur sera capté au 1er tour). Couvre la 3e + /// variante de SessionPlan, non testée jusqu'ici. + #[tokio::test] + async fn factory_assign_does_not_seed_engine_id() { + let factory = StructuredSessionFactory::new(); + let fake = FakeCli::printing(&claude_script()); + let claude = structured_profile(StructuredAdapter::Claude, &fake.command()); + let session = factory + .start( + &claude, + &prepared_ctx(), + &cwd(), + &SessionPlan::Assign { + conversation_id: "ignored-by-engine".to_owned(), + }, + ) + .await + .expect("start assign"); + // Assign n'amorce PAS le moteur : conversation_id reste None avant tour. + assert_eq!(session.conversation_id(), None); + } + + /// La factory échoue proprement (`Start`) si on lui passe un profil SANS adapter + /// structuré (cohérence avec `supports`). Garde-fou défensif. + #[tokio::test] + async fn factory_start_rejects_non_structured_profile() { + let factory = StructuredSessionFactory::new(); + let tui = AgentProfile::new( + ProfileId::new_random(), + "Aider", + "aider", + Vec::new(), + ContextInjection::convention_file("AGENTS.md").expect("valide"), + None, + "{agentRunDir}", + None, + ) + .expect("profil valide"); + match factory + .start(&tui, &prepared_ctx(), &cwd(), &SessionPlan::None) + .await + { + Err(AgentSessionError::Start(_)) => {} + Err(other) => panic!("attendu Start, vu: {other:?}"), + Ok(_) => panic!("la factory ne doit pas démarrer un profil non structuré"), + } + } + + /// Codex passe le harnais de conformité partagé (substituabilité Liskov) — déjà + /// présent ; on ajoute un script Codex MINIMAL (zéro delta) pour prouver que le + /// contrat « ≥0 deltas puis un Final » tient avec zéro delta. + #[tokio::test] + async fn codex_contract_holds_with_zero_deltas() { + let fake = FakeCli::printing(&[ + r#"{"type":"session","conversation_id":"cx-0"}"#, + r#"{"type":"result","text":"direct"}"#, + ]); + let session: Arc = Arc::new(CodexExecSession::new( + SessionId::new_random(), + fake.command(), + "/", + None, + )); + assert_agent_session_contract(session, "cx-0", "direct").await; + } +} diff --git a/crates/infrastructure/src/session/process.rs b/crates/infrastructure/src/session/process.rs new file mode 100644 index 0000000..1de58cf --- /dev/null +++ b/crates/infrastructure/src/session/process.rs @@ -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, + /// 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, +} + +/// 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, +) -> Result, 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, 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()) +}