//! 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::sync::Arc; use std::time::Duration; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::process::Command; use domain::ports::AgentSessionError; use domain::sandbox::{SandboxEnforcer, SandboxPlan}; /// 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, /// Plan de sandbox OS à appliquer sur l'enfant (lot LP4-4). `None` ⇒ aucun /// sandboxing : `run_turn` emprunte le drain async tokio **inchangé** (invariant /// produit : rien posé ⇒ comportement natif). `Some` **et** un enforcer fourni à /// [`run_turn`] ⇒ le plan est appliqué sur l'enfant via [`drain_sandboxed`]. pub sandbox: 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. /// /// Quand `spec.sandbox` porte un plan **et** qu'un `enforcer` est fourni (chemin /// Linux uniquement, lot LP4-4), le tour passe par [`run_turn_sandboxed`] : l'enfant /// est lancé sous le domaine Landlock. Sinon — et **partout** hors Linux — c'est le /// drain async tokio historique, strictement inchangé (zéro régression). /// /// # Errors /// - [`AgentSessionError::Start`] si le process ne démarre pas (binaire introuvable, /// ou enforcement de sandbox impossible : fail-closed, **aucun** enfant ne tourne) ; /// - [`AgentSessionError::Io`] sur échec de lecture/écriture des pipes ; /// - [`AgentSessionError::Timeout`] si `timeout` expire. pub async fn run_turn( spec: &SpawnLine, timeout: Option, enforcer: Option<&Arc>, ) -> Result, AgentSessionError> { // Chemin SANDBOXÉ (Linux + plan posé + enforcer câblé) : transpose la technique // du PTY (`spawn_command_sandboxed`) — enforce sur un thread jetable AVANT le fork, // l'enfant hérite le domaine via fork+exec. #[cfg(target_os = "linux")] if let (Some(plan), Some(enforcer)) = (spec.sandbox.as_ref(), enforcer) { return run_turn_sandboxed(spec, plan.clone(), Arc::clone(enforcer), timeout).await; } // Hors Linux : aucun sandboxing OS ⇒ on ignore l'enforcer (Noop de toute façon) et // on garde le drain async historique. `let _` évite l'avertissement « unused ». #[cfg(not(target_os = "linux"))] let _ = enforcer; match timeout { Some(dur) => match tokio::time::timeout(dur, drain(spec)).await { Ok(result) => result, Err(_elapsed) => Err(AgentSessionError::Timeout), }, None => drain(spec).await, } } /// Variante **sandboxée** du tour (lot LP4-4, Linux seulement). Le drain bloquant /// `std::process` est exécuté sur un **thread jetable** : on y restreint le thread /// (`enforcer.enforce(plan)`) AVANT le `spawn`, puis on lance l'enfant. La technique /// reproduit celle du PTY ([`crate::pty`]) : `pre_exec(enforce)` est INTERDIT (Landlock /// alloue ⇒ risque de deadlock `malloc` post-`fork` en process multithreadé), on /// s'appuie donc sur l'**héritage du domaine Landlock**. /// /// ## Pourquoi pas de `pre_exec` (divergence assumée du cadrage) /// /// Le cadrage proposait un `pre_exec` **vide** pour forcer `std` sur le chemin /// déterministe `fork`+`exec` (jamais `posix_spawn`). Or `pre_exec` est `unsafe`, et /// cette crate est `#![forbid(unsafe_code)]` (invariant non contournable localement). /// On s'en passe sans perte de garantie : `landlock_restrict_self` restreint le /// **thread courant et toute sa descendance**, héritage assuré par le noyau à travers /// `fork`/`clone`/`vfork` **et** préservé par `execve` — donc aussi via `posix_spawn` /// (qui est `clone`+`execve` sous le capot), car l'enforcement vit au niveau des /// *credentials* de la tâche, hors d'atteinte de l'espace utilisateur. Le PTY n'obtenait /// le `fork`+`exec` que comme **effet de bord** du `pre_exec` interne de `portable-pty` ; /// la garantie de sécurité, elle, ne repose que sur cet héritage. Le thread meurt avec /// sa restriction, les autres threads d'IdeA ne sont jamais touchés. /// /// `enforce` fail-closed : un `Err` ⇒ [`AgentSessionError::Start`] et **aucun** enfant. /// /// ## Timeout sous sandbox /// /// Le thread bloquant n'est pas annulable « de l'extérieur ». On le réconcilie avec /// l'async par deux canaux oneshot : le thread renvoie un **killer** /// (`Arc>`) juste après le spawn, puis son résultat à la fin. On pose /// `tokio::time::timeout` sur la réception du résultat ; à expiration on **tue /// l'enfant** via le killer ⇒ EOF côté stdout ⇒ le thread sort de sa boucle de drain, /// `wait()` (reap, pas de zombie) et se termine. On renvoie alors [`AgentSessionError::Timeout`]. #[cfg(target_os = "linux")] async fn run_turn_sandboxed( spec: &SpawnLine, plan: SandboxPlan, enforcer: Arc, timeout: Option, ) -> Result, AgentSessionError> { use std::sync::Mutex as StdMutex; // Données possédées : rien n'emprunte le thread jetable. let command = spec.command.clone(); let args = spec.args.clone(); let cwd = spec.cwd.clone(); let env = spec.env.clone(); let stdin = spec.stdin.clone(); let (killer_tx, killer_rx) = tokio::sync::oneshot::channel::>>(); let (done_tx, done_rx) = tokio::sync::oneshot::channel::, AgentSessionError>>(); // Thread JETABLE : sa restriction Landlock meurt avec lui. std::thread::spawn(move || { let result = drain_sandboxed(command, args, cwd, env, stdin, &enforcer, &plan, killer_tx); // Le récepteur peut avoir abandonné (timeout) : on ignore l'erreur d'envoi. let _ = done_tx.send(result); }); match timeout { Some(dur) => match tokio::time::timeout(dur, done_rx).await { // Le thread a fini dans les temps (succès ou erreur métier). Ok(Ok(result)) => result, // Sender lâché sans valeur (panique du thread) ⇒ I/O. Ok(Err(_canceled)) => Err(AgentSessionError::Io( "thread sandbox terminé sans résultat".to_owned(), )), // Expiration : tue l'enfant (⇒ EOF ⇒ le thread se termine et reap). Err(_elapsed) => { if let Ok(child) = killer_rx.await { if let Ok(mut c) = child.lock() { let _ = c.kill(); } } Err(AgentSessionError::Timeout) } }, None => match done_rx.await { Ok(result) => result, Err(_canceled) => Err(AgentSessionError::Io( "thread sandbox terminé sans résultat".to_owned(), )), }, } } /// Drain **bloquant et sandboxé** exécuté sur le thread jetable (lot LP4-4, Linux). /// /// 1. `enforce(plan)` restreint CE thread (fail-closed) — la descendance hérite le /// domaine Landlock (cf. doc de [`run_turn_sandboxed`]) ; /// 2. `std::process::Command::spawn` (aucun `pre_exec` : `unsafe` interdit ici) ; /// 3. pousse le prompt sur stdin puis EOF ; /// 4. sort `stdout` du child **avant** de le partager : le drain lit sans tenir le /// `Mutex`, donc le killer (timeout) peut verrouiller et tuer à tout moment ; /// 5. draine ligne-à-ligne jusqu'à EOF ; /// 6. `wait()` (reap) — pas de zombie. #[cfg(target_os = "linux")] #[allow(clippy::too_many_arguments)] fn drain_sandboxed( command: String, args: Vec, cwd: String, env: Vec<(String, String)>, stdin_content: Option, enforcer: &Arc, plan: &SandboxPlan, killer_tx: tokio::sync::oneshot::Sender>>, ) -> Result, AgentSessionError> { use std::io::{BufRead, BufReader as StdBufReader, Write as _}; use std::process::{Command as StdCommand, Stdio}; use std::sync::Mutex as StdMutex; // 1. Restreint CE thread AVANT tout fork (fail-closed : Err ⇒ aucun child ne tourne). enforcer .enforce(plan) .map_err(|e| AgentSessionError::Start(format!("sandbox enforcement failed: {e}")))?; // 2. Commande std (≡ `drain` async, mais synchrone). let mut cmd = StdCommand::new(&command); cmd.args(&args) .stdin(Stdio::piped()) .stdout(Stdio::piped()) .stderr(Stdio::piped()); if cwd != "/" && !cwd.is_empty() { cmd.current_dir(&cwd); } for (key, value) in &env { cmd.env(key, value); } // Pas de `pre_exec` : il serait `unsafe` (interdit dans cette crate). L'enfant // hérite de toute façon le domaine Landlock posé sur ce thread (cf. doc de // `run_turn_sandboxed`), que `std` emprunte `posix_spawn` ou `fork`+`exec`. let mut child = cmd .spawn() .map_err(|e| AgentSessionError::Start(format!("{command}: {e}")))?; // 3. Pousse le prompt sur stdin puis le ferme (EOF). Erreur d'I/O ⇒ `Io`. if let Some(input) = &stdin_content { let mut si = child .stdin .take() .ok_or_else(|| AgentSessionError::Io("stdin pipe indisponible".to_owned()))?; si.write_all(input.as_bytes()) .map_err(|e| AgentSessionError::Io(e.to_string()))?; drop(si); } else { drop(child.stdin.take()); } // 4. Sort stdout AVANT de partager le child (drain sans lock ⇒ killer libre). let stdout = child .stdout .take() .ok_or_else(|| AgentSessionError::Io("stdout pipe indisponible".to_owned()))?; let child = Arc::new(StdMutex::new(child)); // Donne au côté async de quoi tuer l'enfant en cas de timeout. let _ = killer_tx.send(Arc::clone(&child)); // 5. Drain ligne-à-ligne jusqu'à EOF. Un kill côté async ferme stdout ⇒ EOF. let mut collected = Vec::new(); for line in StdBufReader::new(stdout).lines() { match line { Ok(l) => collected.push(l), Err(e) => return Err(AgentSessionError::Io(e.to_string())), } } // 6. Reap (pas de zombie). Lock tenu brièvement. if let Ok(mut c) = child.lock() { let _ = c.wait(); } Ok(collected) } /// 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()) }