Files
IdeA/crates/infrastructure/src/session/mod.rs
Blomios aab4bcafb6 feat(session): adapter HTTP OpenAI-compatible pour profils locaux/LAN (#14)
Ajoute un adapter de session HTTP OpenAI-compatible, purement additif,
permettant d'intégrer des modèles locaux/LAN comme profils IdeA canoniques
avec parité tool-calling/MCP.

- domain: extension du profil et des ports pour l'adapter OpenAI-compatible
- infrastructure: adapter openai_compat + routage factory
- app-tauri: mapping des outils OpenAI (openai_tools) + wiring state/lib
- application: catalogue d'agents + tests de use-cases profils

Validé QA (backend GO): round-trip byte-identique Claude/Codex, mapping
erreurs, dégradation tools, conversation_id None, conformance un seul Final,
routage factory. Suites vertes domain 467 / application 530 /
infrastructure 523 / app-tauri 247, 0 échec.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-07 22:09:16 +02:00

2021 lines
79 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

//! 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 openai_compat;
pub mod process;
/// Tests bout-en-bout de l'enforcement Landlock sur le chemin structuré (lot LP4-4),
/// Linux uniquement (pair du module `pty::sandbox_e2e_tests` du lot LP4-3).
#[cfg(all(test, target_os = "linux"))]
mod sandbox_e2e;
pub use claude::ClaudeSdkSession;
pub use codex::CodexExecSession;
pub use conformance::FakeCli;
pub use factory::StructuredSessionFactory;
pub use openai_compat::OpenAiCompatibleSession;
#[cfg(test)]
mod tests {
use std::sync::Arc;
use std::time::Duration;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
use domain::ids::ProfileId;
use domain::ports::{
AgentSession, AgentSessionError, AgentSessionFactory, ContextInjectionPlan,
PreparedContext, ReplyEvent, SessionPlan,
};
use domain::profile::{AgentProfile, ContextInjection, HttpChatConfig, 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(),
project_root: "/project".to_owned(),
}
}
fn cwd() -> ProjectPath {
ProjectPath::new("/").expect("cwd valide")
}
fn temp_cwd(name: &str) -> ProjectPath {
let path = std::env::temp_dir().join(format!(
"idea-structured-session-{name}-{}",
uuid::Uuid::new_v4()
));
std::fs::create_dir_all(&path).expect("temp cwd");
ProjectPath::new(path.to_string_lossy().into_owned()).expect("temp cwd path")
}
fn structured_profile(adapter: StructuredAdapter, command: &str) -> AgentProfile {
let profile = 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);
if adapter == StructuredAdapter::OpenAiCompatible {
profile.with_chat_http(
HttpChatConfig::new(
"http://127.0.0.1:9/v1",
"local-model",
None,
Some(1_000),
Some(100),
Some(1),
)
.expect("valid http config"),
)
} else {
profile
}
}
async fn one_shot_chat_server(content: &'static str) -> (String, tokio::task::JoinHandle<()>) {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let handle = tokio::spawn(async move {
let Ok((mut socket, _)) = listener.accept().await else {
return;
};
let mut buffer = Vec::new();
let mut chunk = [0_u8; 1024];
loop {
let n = socket.read(&mut chunk).await.expect("read request");
if n == 0 {
return;
}
buffer.extend_from_slice(&chunk[..n]);
if buffer.windows(4).any(|window| window == b"\r\n\r\n") {
break;
}
}
let body = format!(
r#"{{"choices":[{{"message":{{"role":"assistant","content":"{content}"}}}}]}}"#
);
let response = format!(
"HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{}",
body.len(),
body
);
socket
.write_all(response.as_bytes())
.await
.expect("write response");
});
(format!("http://{addr}/v1"), handle)
}
// -- 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, None, 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,
sandbox: None,
};
let err = run_turn(&spec, None, None, None)
.await
.expect_err("doit échouer");
assert!(matches!(err, AgentSessionError::Start(_)), "vu: {err:?}");
}
// -- parse_event Claude (format RÉEL vérifié 2026-06-09) --------------
#[test]
fn claude_parse_init_captures_session_id_and_heartbeats() {
let parsed = claude::parse_event(
r#"{"type":"system","subtype":"init","session_id":"conv-123","cwd":"/tmp","tools":[],"model":"claude-opus-4-8"}"#,
)
.expect("parse ok");
assert_eq!(parsed.session_id.as_deref(), Some("conv-123"));
// L'init capte le session_id ET émet un heartbeat (vivacité non terminale, lot 1).
assert_eq!(parsed.events, vec![ReplyEvent::Heartbeat]);
}
/// §21 (LS2) : un `rate_limit_event` n'est PLUS un heartbeat — il porte désormais
/// un [`ReplyEvent::RateLimited`] (niveau 1 structuré). Sans heure de reset
/// exploitable dans `rate_limit_info` ⇒ `RateLimited{None}` (filet humain en aval).
/// Le `session_id` reste capté.
#[test]
fn claude_parse_rate_limit_event_without_reset_is_rate_limited_none() {
let parsed = claude::parse_event(
r#"{"type":"rate_limit_event","rate_limit_info":{"x":1},"session_id":"conv-123"}"#,
)
.expect("parse ok");
assert_eq!(
parsed.events,
vec![ReplyEvent::RateLimited { resets_at_ms: None }]
);
assert_eq!(parsed.session_id.as_deref(), Some("conv-123"));
}
#[test]
fn claude_parse_assistant_text_and_tool() {
let text = claude::parse_event(
r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"bonjour"}]},"session_id":"c","parent_tool_use_id":null}"#,
)
.expect("parse ok");
assert_eq!(
text.events,
vec![ReplyEvent::TextDelta {
text: "bonjour".to_owned()
}]
);
let tool = claude::parse_event(
r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"tool_use","name":"Read"}]},"session_id":"c","parent_tool_use_id":null}"#,
)
.expect("parse ok");
assert_eq!(
tool.events,
vec![ReplyEvent::ToolActivity {
label: "Read".to_owned()
}]
);
}
/// Bug multi-blocs CORRIGÉ : une ligne `assistant` portant `[text, tool_use, text]`
/// produit **3** ReplyEvent dans l'ordre (et non plus le seul premier bloc).
#[test]
fn claude_parse_assistant_multiblock_yields_all_events_in_order() {
let parsed = claude::parse_event(
r#"{"type":"assistant","message":{"role":"assistant","content":[
{"type":"text","text":"un"},
{"type":"tool_use","name":"Read"},
{"type":"text","text":"deux"}]},"session_id":"c","parent_tool_use_id":null}"#,
)
.expect("parse ok");
assert_eq!(
parsed.events,
vec![
ReplyEvent::TextDelta {
text: "un".to_owned()
},
ReplyEvent::ToolActivity {
label: "Read".to_owned()
},
ReplyEvent::TextDelta {
text: "deux".to_owned()
},
]
);
}
#[test]
fn claude_parse_result_is_final() {
let parsed = claude::parse_event(
r#"{"type":"result","subtype":"success","is_error":false,"result":"réponse finale","session_id":"conv-123","num_turns":1}"#,
)
.expect("parse ok");
assert_eq!(
parsed.events,
vec![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!(unknown.events.is_empty());
}
// -- parse_event Codex (format RÉEL vérifié 2026-06-09) ---------------
#[test]
fn codex_parse_thread_started_message_and_final() {
let sess =
codex::parse_event(r#"{"type":"thread.started","thread_id":"cx-9"}"#).expect("ok");
assert_eq!(sess.conversation_id.as_deref(), Some("cx-9"));
// Le handshake ne capte que le thread_id, sans événement (pas un heartbeat).
assert!(sess.events.is_empty());
// turn.started / turn.completed ⇒ heartbeat (vivacité non terminale, lot 1).
let started = codex::parse_event(r#"{"type":"turn.started"}"#).expect("ok");
assert_eq!(started.events, vec![ReplyEvent::Heartbeat]);
let completed = codex::parse_event(r#"{"type":"turn.completed","usage":{}}"#).expect("ok");
assert_eq!(completed.events, vec![ReplyEvent::Heartbeat]);
let msg = codex::parse_event(
r#"{"type":"item.completed","item":{"id":"item_0","type":"agent_message","text":"fini"}}"#,
)
.expect("ok");
assert_eq!(
msg.events,
vec![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 (format RÉEL) : init → texte → tool_use → texte → result.
fn claude_script() -> Vec<&'static str> {
vec![
r#"{"type":"system","subtype":"init","session_id":"claude-conv-1","cwd":"/tmp","tools":[]}"#,
r#"{"type":"rate_limit_event","rate_limit_info":{"x":1},"session_id":"claude-conv-1"}"#,
r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"un "}]},"session_id":"claude-conv-1","parent_tool_use_id":null}"#,
r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"tool_use","name":"Read"}]},"session_id":"claude-conv-1","parent_tool_use_id":null}"#,
r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"deux"}]},"session_id":"claude-conv-1","parent_tool_use_id":null}"#,
r#"{"type":"result","subtype":"success","is_error":false,"result":"réponse Claude","session_id":"claude-conv-1","num_turns":1}"#,
]
}
/// Script Codex (format RÉEL) : thread.started → turn.started → reasoning item →
/// agent_message (= Final) → turn.completed.
fn codex_script() -> Vec<&'static str> {
vec![
r#"{"type":"thread.started","thread_id":"codex-conv-1"}"#,
r#"{"type":"turn.started"}"#,
r#"{"type":"item.completed","item":{"id":"item_0","type":"reasoning","text":"…"}}"#,
r#"{"type":"item.completed","item":{"id":"item_1","type":"agent_message","text":"réponse Codex"}}"#,
r#"{"type":"turn.completed","usage":{"input_tokens":10848}}"#,
]
}
#[tokio::test]
async fn claude_session_respects_port_contract() {
let fake = FakeCli::printing(&claude_script());
let session: Arc<dyn AgentSession> = Arc::new(ClaudeSdkSession::new(
SessionId::new_random(),
fake.command(),
"/",
None,
None,
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<dyn AgentSession> = Arc::new(CodexExecSession::new(
SessionId::new_random(),
fake.command(),
"/",
None,
Vec::new(),
None,
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,
None,
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,
None,
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));
}
/// §17.3/D7 — **strict coherence**: `AgentProfile::is_selectable` (the menu's
/// selection gate) and `StructuredSessionFactory::supports` (the runtime's
/// routing gate) must agree on **every** reference profile. If they ever
/// diverged, the wizard could offer a profile the runtime cannot drive (or
/// hide one it can). Asserting profile-by-profile over the real catalogue —
/// including the expected `claude=codex=true`, `gemini=aider=false` truth
/// table — means this fails the moment either gate changes without the other.
#[tokio::test]
async fn supports_and_is_selectable_agree_on_every_reference_profile() {
use std::collections::HashMap;
let factory = StructuredSessionFactory::new();
let profiles = application::reference_profiles();
let mut expected: HashMap<&str, bool> = HashMap::new();
expected.insert("claude", true);
expected.insert("codex", true);
expected.insert("openai-compatible", true);
expected.insert("gemini", false);
expected.insert("aider", false);
assert_eq!(
profiles.len(),
5,
"catalogue has the five reference profiles"
);
for profile in &profiles {
let selectable = profile.is_selectable();
assert_eq!(
selectable,
factory.supports(profile),
"is_selectable and supports must agree for `{}`",
profile.command
);
assert_eq!(
Some(&selectable),
expected.get(profile.command.as_str()),
"unexpected selectability for `{}`",
profile.command
);
}
}
#[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, 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, 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_routes_openai_compatible_to_http_session() {
let factory = StructuredSessionFactory::new();
let (endpoint, handle) = one_shot_chat_server("réponse HTTP").await;
let openai = structured_profile(StructuredAdapter::OpenAiCompatible, "openai-compatible")
.with_chat_http(
HttpChatConfig::new(
endpoint,
"local-model",
None,
Some(10_000),
Some(1_000),
None,
)
.expect("valid http config"),
);
let session = factory
.start(
&openai,
&prepared_ctx(),
&temp_cwd("factory-openai"),
&SessionPlan::None,
None,
)
.await
.expect("start OpenAI-compatible ok");
assert_eq!(
session.conversation_id(),
None,
"OpenAI-compatible route must not expose a provider conversation id"
);
let content = drain_final(session.as_ref()).await;
assert_eq!(content, "réponse HTTP");
handle.abort();
}
#[tokio::test]
async fn factory_passes_project_root_to_codex_add_dir() {
let factory = StructuredSessionFactory::new();
let (cmd, argv) = make_recording_fake(&[
r#"{"type":"thread.started","thread_id":"cx-new"}"#,
r#"{"type":"item.completed","item":{"id":"i0","type":"agent_message","text":"ok"}}"#,
]);
let codex = structured_profile(StructuredAdapter::Codex, &cmd);
let ctx = PreparedContext {
content: MarkdownDoc::new("# ctx"),
relative_path: "AGENTS.md".to_owned(),
project_root: "/project/root".to_owned(),
};
let session = factory
.start(&codex, &ctx, &cwd(), &SessionPlan::None, None)
.await
.expect("start Codex ok");
let content = drain_final(session.as_ref()).await;
assert_eq!(content, "ok");
let recorded = std::fs::read_to_string(&argv).expect("argv");
let args: Vec<&str> = recorded.lines().collect();
assert!(
args.windows(2).any(|w| w == ["--add-dir", "/project/root"]),
"factory must relay PreparedContext.project_root to Codex --add-dir, got: {args:?}"
);
assert!(
!args.contains(&"--ask-for-approval"),
"codex exec must not receive unsupported approval flags, got: {args:?}"
);
let _ = std::fs::remove_file(&cmd);
let _ = std::fs::remove_file(&argv);
}
#[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(),
},
None,
)
.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,
sandbox: None,
};
let err = run_turn(&spec, Some(Duration::from_millis(50)), None, None)
.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 <id>`).
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 dir = std::env::current_dir()
.expect("cwd")
.join("target")
.join("test-fakes")
.join("session");
std::fs::create_dir_all(&dir).expect("create rec fake dir");
let bin = dir.join(format!("idea-rec-cli-{}-{n}", std::process::id()));
let argv = dir.join(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 ---------------
/// Messages successifs : chaque message produit ses événements ; ici un par
/// message (texte puis tool_use).
#[test]
fn claude_multiple_messages_each_yield_their_events() {
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.events, vec![ReplyEvent::TextDelta { text: "a".into() }]);
assert_eq!(
t2.events,
vec![ReplyEvent::ToolActivity {
label: "Bash".into()
}]
);
}
/// Bug multi-blocs CORRIGÉ : un **seul** message portant plusieurs blocs
/// `text`/`tool_use` rend TOUS ses blocs dans l'ordre (plus de perte).
#[test]
fn claude_multiblock_message_yields_every_block() {
let parsed = claude::parse_event(
r#"{"type":"assistant","message":{"content":[
{"type":"text","text":"un"},
{"type":"tool_use","name":"Read"},
{"type":"text","text":"deux"}]}}"#,
)
.unwrap();
assert_eq!(
parsed.events,
vec![
ReplyEvent::TextDelta { text: "un".into() },
ReplyEvent::ToolActivity {
label: "Read".into()
},
ReplyEvent::TextDelta {
text: "deux".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.events,
vec![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!(parsed.events.is_empty());
}
/// 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!(claude::parse_event("[1,2,3]").unwrap().events.is_empty());
assert!(claude::parse_event("42").unwrap().events.is_empty());
}
// ---- Codex parse_event (format RÉEL vérifié 2026-06-09) -------------
/// Un item non-`agent_message` (reasoning/command/…) ⇒ ToolActivity (label = type).
#[test]
fn codex_non_agent_message_item_is_tool_activity() {
let r = codex::parse_event(
r#"{"type":"item.completed","item":{"id":"i0","type":"reasoning","text":"…"}}"#,
)
.unwrap();
assert_eq!(
r.events,
vec![ReplyEvent::ToolActivity {
label: "reasoning".into()
}]
);
let c =
codex::parse_event(r#"{"type":"item.completed","item":{"id":"i1","type":"command"}}"#)
.unwrap();
assert_eq!(
c.events,
vec![ReplyEvent::ToolActivity {
label: "command".into()
}]
);
}
/// `agent_message` ⇒ Final (porte le texte de réponse).
#[test]
fn codex_agent_message_is_final() {
let m = codex::parse_event(
r#"{"type":"item.completed","item":{"id":"i0","type":"agent_message","text":"bonjour"}}"#,
)
.unwrap();
assert_eq!(
m.events,
vec![ReplyEvent::Final {
content: "bonjour".into()
}]
);
}
/// `agent_message` sans `text` ⇒ Final avec contenu vide (pas de panic).
#[test]
fn codex_agent_message_without_text_is_empty_final() {
let m = codex::parse_event(
r#"{"type":"item.completed","item":{"id":"i0","type":"agent_message"}}"#,
)
.unwrap();
assert_eq!(
m.events,
vec![ReplyEvent::Final {
content: String::new()
}]
);
}
/// `thread.started` capte le `thread_id` ; pas d'événement.
#[test]
fn codex_thread_started_captures_thread_id() {
let p = codex::parse_event(r#"{"type":"thread.started","thread_id":"cx-7"}"#).unwrap();
assert_eq!(p.conversation_id.as_deref(), Some("cx-7"));
assert!(p.events.is_empty());
}
/// Ligne vide / type inconnu ⇒ ignorés sans erreur. (turn.started/completed sont
/// désormais des heartbeats : couverts par `codex_parse_thread_started_message_and_final`.)
#[test]
fn codex_empty_and_unknown_ignored() {
assert_eq!(codex::parse_event("").unwrap(), Default::default());
// Un `type` inconnu reste ignoré (robustesse), pas un heartbeat.
assert!(codex::parse_event(r#"{"type":"telemetry"}"#)
.unwrap()
.events
.is_empty());
}
// ---- Machinerie process via FakeCli ---------------------------------
/// Deltas PUIS Final : **exactement un** `Final`, et aucun autre `Final` après lui
/// (substituabilité Liskov). Note (lot 1) : un `turn.completed` postérieur émet un
/// `Heartbeat` non terminal — légitimement après le `Final` —, donc on ne teste plus
/// « rien après le Final » mais « pas de second Final, et seul un heartbeat peut
/// suivre ». Le rendez-vous synchrone (`drain_to_final`) s'arrête de toute façon au
/// premier `Final`.
#[tokio::test]
async fn codex_stream_closed_after_final() {
let fake = FakeCli::printing(&[
r#"{"type":"thread.started","thread_id":"c"}"#,
r#"{"type":"item.completed","item":{"id":"i0","type":"reasoning","text":"a"}}"#,
r#"{"type":"item.completed","item":{"id":"i1","type":"agent_message","text":"fin"}}"#,
r#"{"type":"turn.completed","usage":{}}"#,
]);
let s = CodexExecSession::new(
SessionId::new_random(),
fake.command(),
"/",
None,
Vec::new(),
None,
None,
);
let events: Vec<_> = s.send("x").await.expect("send").collect();
let finals = events
.iter()
.filter(|e| matches!(e, ReplyEvent::Final { .. }))
.count();
assert_eq!(finals, 1, "exactement un Final");
// Après le Final, seuls des événements non terminaux (heartbeat) peuvent suivre.
let after_final_terminals = events
.iter()
.skip_while(|e| !matches!(e, ReplyEvent::Final { .. }))
.skip(1)
.filter(|e| matches!(e, ReplyEvent::Final { .. }))
.count();
assert_eq!(
after_final_terminals, 0,
"aucun second Final après le premier"
);
}
#[tokio::test]
async fn codex_two_agent_messages_preserve_first_as_announcement_and_last_as_final() {
let fake = FakeCli::printing(&[
r#"{"type":"thread.started","thread_id":"c"}"#,
r#"{"type":"item.completed","item":{"id":"i0","type":"agent_message","text":"je regarde"}}"#,
r#"{"type":"item.completed","item":{"id":"i1","type":"agent_message","text":"résultat"}}"#,
]);
let s = CodexExecSession::new(
SessionId::new_random(),
fake.command(),
"/",
None,
Vec::new(),
None,
None,
);
let events: Vec<_> = s.send("x").await.expect("send").collect();
assert_eq!(
events,
vec![
ReplyEvent::Announcement {
text: "je regarde".into()
},
ReplyEvent::Final {
content: "résultat".into()
}
]
);
}
#[tokio::test]
async fn codex_single_agent_message_is_final_without_announcement() {
let fake = FakeCli::printing(&[
r#"{"type":"thread.started","thread_id":"c"}"#,
r#"{"type":"item.completed","item":{"id":"i0","type":"agent_message","text":"résultat"}}"#,
]);
let s = CodexExecSession::new(
SessionId::new_random(),
fake.command(),
"/",
None,
Vec::new(),
None,
None,
);
let events: Vec<_> = s.send("x").await.expect("send").collect();
assert_eq!(
events,
vec![ReplyEvent::Final {
content: "résultat".into()
}]
);
}
#[tokio::test]
async fn codex_zero_agent_message_has_no_final() {
let fake = FakeCli::printing(&[
r#"{"type":"thread.started","thread_id":"c"}"#,
r#"{"type":"item.completed","item":{"id":"i0","type":"reasoning","text":"analyse"}}"#,
]);
let s = CodexExecSession::new(
SessionId::new_random(),
fake.command(),
"/",
None,
Vec::new(),
None,
None,
);
let events: Vec<_> = s.send("x").await.expect("send").collect();
assert!(
events
.iter()
.all(|e| !matches!(e, ReplyEvent::Final { .. })),
"aucun agent_message => aucun Final"
);
}
#[tokio::test]
async fn codex_send_with_tap_emits_each_agent_message_live_as_announcement() {
let fake = FakeCli::printing(&[
r#"{"type":"thread.started","thread_id":"c"}"#,
r#"{"type":"item.completed","item":{"id":"i0","type":"agent_message","text":"je regarde"}}"#,
r#"{"type":"item.completed","item":{"id":"i1","type":"agent_message","text":"résultat"}}"#,
]);
let s = CodexExecSession::new(
SessionId::new_random(),
fake.command(),
"/",
None,
Vec::new(),
None,
None,
);
let (tx, rx) = std::sync::mpsc::channel();
let events: Vec<_> = s.send_with_tap("x", tx).await.expect("send").collect();
let live: Vec<_> = rx.into_iter().collect();
assert_eq!(
live,
vec![
ReplyEvent::Announcement {
text: "je regarde".into()
},
ReplyEvent::Announcement {
text: "résultat".into()
},
],
"le tap live publie chaque agent_message, y compris celui qui deviendra Final"
);
assert_eq!(
events,
vec![ReplyEvent::Final {
content: "résultat".into()
}],
"le flux final du chemin tap garde seulement le dernier Final pour le demandeur"
);
}
/// 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,
None,
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, None, 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, None, 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,
sandbox: None,
};
let err = run_turn(&spec, Some(Duration::from_millis(50)), None, None)
.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 <id>` 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","is_error":false,"result":"ok","session_id":"resume-id","num_turns":1}"#,
]);
let session = ClaudeSdkSession::new(
SessionId::new_random(),
cmd.clone(),
"/",
Some("resume-id".to_owned()),
None,
None,
);
// 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:?}"
);
// Format RÉEL : --output-format stream-json --verbose sont requis.
assert!(
args.contains(&"--output-format") && args.contains(&"stream-json"),
"argv doit porter --output-format stream-json, vu: {args:?}"
);
assert!(
args.contains(&"--verbose"),
"argv doit porter --verbose, vu: {args:?}"
);
let _ = std::fs::remove_file(&cmd);
let _ = std::fs::remove_file(&argv);
}
/// Idem Codex : `codex exec --json --skip-git-repo-check resume <thread_id> <prompt>`.
/// Les options globales de `exec` doivent précéder la sous-commande `resume`.
#[tokio::test]
async fn codex_resume_command_carries_resume_subcommand() {
let (cmd, argv) = make_recording_fake(&[
r#"{"type":"item.completed","item":{"id":"i0","type":"agent_message","text":"ok"}}"#,
]);
let session = CodexExecSession::new(
SessionId::new_random(),
cmd.clone(),
"/",
Some("cx-id".to_owned()),
Vec::new(),
None,
None,
);
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(&"--json"), "vu: {args:?}");
assert!(args.contains(&"--skip-git-repo-check"), "vu: {args:?}");
assert!(args.contains(&"vas-y"), "vu: {args:?}");
let resume_pos = args.iter().position(|a| *a == "resume").expect("resume");
let id_pos = args.iter().position(|a| *a == "cx-id").expect("id");
let json_pos = args.iter().position(|a| *a == "--json").expect("json");
let skip_pos = args
.iter()
.position(|a| *a == "--skip-git-repo-check")
.expect("skip");
assert!(
json_pos < resume_pos && resume_pos < id_pos,
"exec options must precede resume, then the session id, vu: {args:?}"
);
assert!(
skip_pos < resume_pos && resume_pos < id_pos,
"exec options must precede resume, then the session id, 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, None, 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":"item.completed","item":{"id":"i0","type":"agent_message","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(),
},
None,
)
.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(),
},
None,
)
.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, 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":"thread.started","thread_id":"cx-0"}"#,
r#"{"type":"item.completed","item":{"id":"i0","type":"agent_message","text":"direct"}}"#,
]);
let session: Arc<dyn AgentSession> = Arc::new(CodexExecSession::new(
SessionId::new_random(),
fake.command(),
"/",
None,
Vec::new(),
None,
None,
));
assert_agent_session_contract(session, "cx-0", "direct").await;
}
// =====================================================================
// DURCISSEMENT QA (lot D2-bis) — bouche les axes RÉSIDUELS du périmètre :
// (a) flot COMPLET init→assistant(MULTI-blocs sur une SEULE ligne)→result
// drainé via `send()` (et pas seulement `parse_event`) : on prouve que
// l'aplatissement multi-blocs tient bout-en-bout et qu'il y a UN Final ;
// (b) commande NEUVE : Claude ne porte PAS `--resume` au 1er tour (isolé,
// pas un `contains` global) et porte la base réelle ;
// (c) commande NEUVE : Codex porte `exec --json --skip-git-repo-check` SANS
// sous-commande `resume`, et `resume` n'apparaît PAS avant capture.
// Tout passe par le FakeCli/sidecar argv — jamais le vrai claude/codex.
// =====================================================================
/// Flot COMPLET Claude où l'assistant émet ses blocs sur **une seule ligne**
/// `content:[text,tool_use,text]` : drainé via `send()`, on doit obtenir, dans
/// l'ordre, TextDelta("a"), ToolActivity("T"), TextDelta("b") PUIS exactement
/// **un** Final("..."). Couvre l'aplatissement multi-blocs bout-en-bout (le
/// harnais de conformité, lui, n'utilise que des lignes mono-bloc).
#[tokio::test]
async fn claude_full_flow_multiblock_line_flattens_then_single_final() {
let fake = FakeCli::printing(&[
r#"{"type":"system","subtype":"init","session_id":"flow-1","cwd":"/tmp","tools":[]}"#,
r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"a"},{"type":"tool_use","name":"T"},{"type":"text","text":"b"}]},"session_id":"flow-1","parent_tool_use_id":null}"#,
r#"{"type":"result","subtype":"success","is_error":false,"result":"final-ok","session_id":"flow-1","num_turns":1}"#,
]);
let session = ClaudeSdkSession::new(
SessionId::new_random(),
fake.command(),
"/",
None,
None,
None,
);
let events: Vec<ReplyEvent> = session.send("x").await.expect("send ok").collect();
assert_eq!(
events,
vec![
// L'init `system` émet un heartbeat (vivacité non terminale, lot 1).
ReplyEvent::Heartbeat,
ReplyEvent::TextDelta { text: "a".into() },
ReplyEvent::ToolActivity { label: "T".into() },
ReplyEvent::TextDelta { text: "b".into() },
ReplyEvent::Final {
content: "final-ok".into()
},
],
"heartbeat d'init, puis les 3 blocs aplatis, PUIS un seul Final"
);
// Un seul Final, en dernière position (redondant mais explicite).
assert_eq!(
events
.iter()
.filter(|e| matches!(e, ReplyEvent::Final { .. }))
.count(),
1
);
assert_eq!(session.conversation_id().as_deref(), Some("flow-1"));
}
/// Commande NEUVE Claude : au 1er tour (aucun seed), l'argv ne doit PAS porter
/// `--resume` ni d'id, mais DOIT porter `-p <prompt> --output-format stream-json
/// --verbose`. (Le test existant prouve l'apparition au 2e tour via un `contains`
/// global ; ici on isole le 1er tour pour prouver l'ABSENCE.)
#[tokio::test]
async fn claude_new_conversation_command_has_no_resume() {
let (cmd, argv) = make_recording_fake(&[
r#"{"type":"system","subtype":"init","session_id":"new-1"}"#,
r#"{"type":"result","subtype":"success","result":"r","session_id":"new-1"}"#,
]);
let session =
ClaudeSdkSession::new(SessionId::new_random(), cmd.clone(), "/", None, None, None);
let _ = session.send("bonjour").await.expect("send ok");
let recorded = std::fs::read_to_string(&argv).expect("argv");
let args: Vec<&str> = recorded.lines().collect();
assert!(
!args.contains(&"--resume"),
"1er tour NEUF ne doit PAS porter --resume, vu: {args:?}"
);
// La base RÉELLE est bien présente.
assert!(args.contains(&"-p"), "vu: {args:?}");
assert!(args.contains(&"bonjour"), "vu: {args:?}");
assert!(
args.contains(&"--output-format") && args.contains(&"stream-json"),
"vu: {args:?}"
);
assert!(args.contains(&"--verbose"), "vu: {args:?}");
// L'ordre RÉEL : -p précède son prompt, qui précède --output-format.
let p = args.iter().position(|a| *a == "-p").unwrap();
let of = args.iter().position(|a| *a == "--output-format").unwrap();
assert!(p < of, "-p doit précéder --output-format, vu: {args:?}");
let _ = std::fs::remove_file(&cmd);
let _ = std::fs::remove_file(&argv);
}
/// Commande NEUVE Codex : au 1er tour (aucun seed), l'argv doit porter
/// `exec --json --skip-git-repo-check <prompt>` SANS la sous-commande `resume`.
#[tokio::test]
async fn codex_new_conversation_command_has_no_resume_subcommand() {
let (cmd, argv) = make_recording_fake(&[
r#"{"type":"thread.started","thread_id":"cx-new"}"#,
r#"{"type":"item.completed","item":{"id":"i0","type":"agent_message","text":"ok"}}"#,
]);
let session = CodexExecSession::new(
SessionId::new_random(),
cmd.clone(),
"/",
None,
Vec::new(),
None,
None,
);
let _ = session.send("salut").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"),
"1er tour NEUF ne doit PAS porter la sous-commande resume, vu: {args:?}"
);
assert!(args.contains(&"--json"), "vu: {args:?}");
assert!(args.contains(&"--skip-git-repo-check"), "vu: {args:?}");
assert!(args.contains(&"salut"), "vu: {args:?}");
// `exec` est bien la 1re sous-commande (position 0 de l'argv).
assert_eq!(
args.first(),
Some(&"exec"),
"exec doit ouvrir l'argv, vu: {args:?}"
);
let _ = std::fs::remove_file(&cmd);
let _ = std::fs::remove_file(&argv);
}
// =====================================================================
// DURCISSEMENT QA (lot D3, §17.9 D3 — fix codex 0.137) — autonomie
// d'écriture Codex : la commande générée porte EXACTEMENT
// [exec, --json, --skip-git-repo-check, --sandbox, workspace-write,
// --add-dir, <project-root>, <prompt>]
// (`resume <id>` après les options `exec` pour une reprise). Le flag `--ask-for-approval never`
// a été RETIRÉ : `codex exec` 0.137 ne le connaît pas (`error: unexpected
// argument`) et est déjà non-interactif. Ce test verrouille l'argv exact pour
// qu'aucune régression ne réintroduise un flag inconnu de la sous-commande.
// Prouvé via le sidecar argv du fake enregistreur (jamais le vrai codex).
// =====================================================================
/// Conversation NEUVE : argv EXACT `[exec, --json, --skip-git-repo-check,
/// --sandbox, workspace-write, --add-dir, <project-root>, <prompt>]`.
/// Pas de sous-commande `resume`, pas de `--ask-for-approval`.
#[tokio::test]
async fn codex_new_conversation_command_carries_exact_args() {
let (cmd, argv) = make_recording_fake(&[
r#"{"type":"thread.started","thread_id":"cx-new"}"#,
r#"{"type":"item.completed","item":{"id":"i0","type":"agent_message","text":"ok"}}"#,
]);
let session = CodexExecSession::new(
SessionId::new_random(),
cmd.clone(),
"/",
None,
vec!["/project/root".to_owned()],
None,
None,
);
let _ = session.send("salut").await.expect("send ok");
let recorded = std::fs::read_to_string(&argv).expect("argv");
let args: Vec<&str> = recorded.lines().collect();
assert_eq!(
args,
vec![
"exec",
"--json",
"--skip-git-repo-check",
"--sandbox",
"workspace-write",
"--add-dir",
"/project/root",
"salut",
],
"argv neuf doit être exact (sans resume, sans --ask-for-approval), vu: {args:?}"
);
let _ = std::fs::remove_file(&cmd);
let _ = std::fs::remove_file(&argv);
}
/// REPRISE (seed d'id) : argv EXACT `[exec, --json, --skip-git-repo-check,
/// --sandbox, workspace-write, --add-dir, <project-root>, resume, <id>, <prompt>]`.
/// Toujours pas de `--ask-for-approval`.
#[tokio::test]
async fn codex_resume_command_carries_exact_args() {
let (cmd, argv) = make_recording_fake(&[
r#"{"type":"item.completed","item":{"id":"i0","type":"agent_message","text":"ok"}}"#,
]);
let session = CodexExecSession::new(
SessionId::new_random(),
cmd.clone(),
"/",
Some("cx-id".to_owned()),
vec!["/project/root".to_owned()],
None,
None,
);
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_eq!(
args,
vec![
"exec",
"--json",
"--skip-git-repo-check",
"--sandbox",
"workspace-write",
"--add-dir",
"/project/root",
"resume",
"cx-id",
"vas-y",
],
"argv reprise doit être exact (options exec avant resume, sans --ask-for-approval), vu: {args:?}"
);
let _ = std::fs::remove_file(&cmd);
let _ = std::fs::remove_file(&argv);
}
// =====================================================================
// LS2 — adapter Claude niveau 1 (§21) : `parse_reset_ms` (parseur ISO-8601
// À LA MAIN + heuristique secondes/ms + days_from_civil) et le mapping
// `parse_event` du `rate_limit_event` vers `ReplyEvent::RateLimited`, plus
// la NON-TERMINALITÉ (T4). Tout passe par les fonctions pures (jamais de
// process) sauf le test de séquence via `send()` (FakeCli).
// =====================================================================
use serde_json::json;
use super::claude::parse_reset_ms;
// ---- parse_reset_ms : noms de champ reconnus + priorité ----------------
/// Les cinq noms de champ plausibles sont chacun reconnus (valeur en secondes
/// ⇒ ×1000). Couvre `resetsAt`, `resets_at`, `reset_at`, `resetAt`, `reset`.
#[test]
fn parse_reset_ms_recognises_every_field_name() {
for key in ["resetsAt", "resets_at", "reset_at", "resetAt", "reset"] {
let info = json!({ key: 1_700_000_000_i64 });
assert_eq!(
parse_reset_ms(&info),
Some(1_700_000_000_000),
"le champ `{key}` doit être reconnu (epoch secondes ×1000)"
);
}
}
/// Priorité : si plusieurs clés sont présentes, la PREMIÈRE de l'ordre
/// (`resetsAt` avant `reset`) gagne.
#[test]
fn parse_reset_ms_first_known_key_wins() {
// resetsAt (priorité 1) = 1_700_000_000 s ; reset (priorité 5) = 5 s.
let info = json!({ "reset": 5, "resetsAt": 1_700_000_000_i64 });
assert_eq!(
parse_reset_ms(&info),
Some(1_700_000_000_000),
"resetsAt prime sur reset"
);
}
// ---- parse_reset_ms : heuristique secondes vs millisecondes ------------
#[test]
fn parse_reset_ms_integer_seconds_are_scaled_to_ms() {
// < 10^12 ⇒ secondes ⇒ ×1000.
assert_eq!(
parse_reset_ms(&json!({ "reset": 1_700_000_000_i64 })),
Some(1_700_000_000_000)
);
}
#[test]
fn parse_reset_ms_integer_millis_are_kept_as_is() {
// ≥ 10^12 ⇒ déjà des millisecondes ⇒ tel quel.
assert_eq!(
parse_reset_ms(&json!({ "reset": 1_700_000_000_000_i64 })),
Some(1_700_000_000_000)
);
}
/// Le SEUIL exact (10^12) : juste en-dessous ⇒ secondes (×1000) ; pile/au-dessus
/// ⇒ millisecondes (tel quel).
#[test]
fn parse_reset_ms_threshold_boundary() {
// 10^12 - 1 ⇒ secondes ⇒ ×1000.
assert_eq!(
parse_reset_ms(&json!({ "reset": 999_999_999_999_i64 })),
Some(999_999_999_999_000)
);
// 10^12 pile ⇒ millisecondes ⇒ tel quel (la borne est inclusive côté ms).
assert_eq!(
parse_reset_ms(&json!({ "reset": 1_000_000_000_000_i64 })),
Some(1_000_000_000_000)
);
}
// ---- parse_reset_ms : floats ------------------------------------------
#[test]
fn parse_reset_ms_float_seconds_preserve_fraction() {
// 1_700_000_000.5 s < 10^12 ⇒ ×1000 = 1_700_000_000_500 ms.
assert_eq!(
parse_reset_ms(&json!({ "reset": 1_700_000_000.5_f64 })),
Some(1_700_000_000_500)
);
}
#[test]
fn parse_reset_ms_float_millis_kept_as_is() {
// 1.7e12 ≥ 10^12 ⇒ déjà ms ⇒ tronqué tel quel.
assert_eq!(
parse_reset_ms(&json!({ "reset": 1_700_000_000_000.0_f64 })),
Some(1_700_000_000_000)
);
}
// ---- parse_reset_ms : chaînes numériques (même heuristique) ------------
#[test]
fn parse_reset_ms_string_integer_uses_seconds_heuristic() {
assert_eq!(
parse_reset_ms(&json!({ "reset": "1700000000" })),
Some(1_700_000_000_000)
);
}
#[test]
fn parse_reset_ms_string_float_uses_seconds_heuristic() {
assert_eq!(
parse_reset_ms(&json!({ "reset": "1700000000.5" })),
Some(1_700_000_000_500)
);
}
// ---- parse_reset_ms : ISO-8601 / RFC3339 (parseur maison) --------------
/// `...Z` (UTC) : un instant rond connu. `2023-11-14T22:13:20Z` correspond à
/// l'epoch 1_700_000_000 s ⇒ 1_700_000_000_000 ms. Recoupe le parseur ISO
/// maison (days_from_civil + math d'heure) contre l'heuristique secondes.
#[test]
fn parse_reset_ms_iso_utc_z() {
assert_eq!(
parse_reset_ms(&json!({ "reset": "2023-11-14T22:13:20Z" })),
Some(1_700_000_000_000)
);
}
/// Offset `+hh:mm` : `2023-11-14T23:13:20+01:00` est le MÊME instant que
/// `22:13:20Z` ⇒ doit donner exactement le même epoch-ms (offset soustrait).
#[test]
fn parse_reset_ms_iso_positive_offset_converts_to_utc() {
assert_eq!(
parse_reset_ms(&json!({ "reset": "2023-11-14T23:13:20+01:00" })),
Some(1_700_000_000_000),
"+01:00 ⇒ on soustrait 1h pour revenir à l'UTC"
);
}
/// Offset `-hh:mm` : `2023-11-14T21:13:20-01:00` est aussi `22:13:20Z`.
#[test]
fn parse_reset_ms_iso_negative_offset_converts_to_utc() {
assert_eq!(
parse_reset_ms(&json!({ "reset": "2023-11-14T21:13:20-01:00" })),
Some(1_700_000_000_000),
"-01:00 ⇒ on ajoute 1h pour revenir à l'UTC"
);
}
/// Offset compact `±hhmm` (sans `:`) supporté par `split_tz`.
#[test]
fn parse_reset_ms_iso_compact_offset() {
assert_eq!(
parse_reset_ms(&json!({ "reset": "2023-11-14T23:13:20+0100" })),
Some(1_700_000_000_000)
);
}
/// Fraction de seconde `.fff` : tronquée/complétée à 3 chiffres (précision ms).
#[test]
fn parse_reset_ms_iso_fraction_padded_and_truncated() {
// `.5` ⇒ "500" ms.
assert_eq!(
parse_reset_ms(&json!({ "reset": "1970-01-01T00:00:00.5Z" })),
Some(500)
);
// `.123456` ⇒ tronqué à "123" ms.
assert_eq!(
parse_reset_ms(&json!({ "reset": "1970-01-01T00:00:00.123456Z" })),
Some(123)
);
// `.7` ⇒ complété à "700" ms.
assert_eq!(
parse_reset_ms(&json!({ "reset": "1970-01-01T00:00:00.7Z" })),
Some(700)
);
}
// ---- parse_reset_ms : robustesse (jamais de panic, jamais d'erreur) ----
#[test]
fn parse_reset_ms_unknown_key_yields_none() {
// Aucune clé connue ⇒ None.
assert_eq!(parse_reset_ms(&json!({ "retryAfter": 60 })), None);
assert_eq!(parse_reset_ms(&json!({})), None);
}
#[test]
fn parse_reset_ms_non_numeric_garbage_yields_none() {
// Valeurs inexploitables (booléen, null, tableau, objet, chaîne pourrie) ⇒ None.
assert_eq!(parse_reset_ms(&json!({ "reset": true })), None);
assert_eq!(parse_reset_ms(&json!({ "reset": null })), None);
assert_eq!(parse_reset_ms(&json!({ "reset": [1, 2, 3] })), None);
assert_eq!(parse_reset_ms(&json!({ "reset": { "nested": 1 } })), None);
assert_eq!(parse_reset_ms(&json!({ "reset": "pas une date" })), None);
}
/// Formes ISO **structurellement** malformées ⇒ None (pas de panic). NB : le
/// parseur maison ne valide PAS les plages (un mois 13 / jour 99 calcule une
/// valeur sans erreur) ; ce qui produit `None`, c'est l'ABSENCE de séparateur
/// `T`, une composante non numérique, ou un nombre de composantes invalide.
#[test]
fn parse_reset_ms_invalid_iso_string_yields_none() {
// Pas de séparateur de date/heure.
assert_eq!(parse_reset_ms(&json!({ "reset": "2023-11-14" })), None);
// Année non numérique.
assert_eq!(
parse_reset_ms(&json!({ "reset": "abcd-11-14T00:00:00Z" })),
None
);
// Composante de date manquante (pas de jour).
assert_eq!(
parse_reset_ms(&json!({ "reset": "2023-11T00:00:00Z" })),
None
);
// Trop de composantes de date.
assert_eq!(
parse_reset_ms(&json!({ "reset": "2023-11-14-9T00:00:00Z" })),
None
);
// Minute manquante dans l'heure.
assert_eq!(parse_reset_ms(&json!({ "reset": "2023-11-14T22Z" })), None);
}
// ---- days_from_civil & bissextiles (via le parseur ISO) ----------------
/// Référence absolue : l'époque Unix elle-même. `1970-01-01T00:00:00Z` ⇒ 0 ms
/// (days_from_civil(1970,1,1) == 0).
#[test]
fn parse_reset_ms_unix_epoch_is_zero() {
assert_eq!(
parse_reset_ms(&json!({ "reset": "1970-01-01T00:00:00Z" })),
Some(0)
);
}
/// Année bissextile : le 29 février 2024 existe et donne l'epoch attendu.
/// `2024-02-29T00:00:00Z` = 1_709_164_800 s = 1_709_164_800_000 ms (calculé à la
/// main : 2024-01-01 = 1_704_067_200 ; +31j (janvier) ; +28j pour atteindre le 29).
#[test]
fn parse_reset_ms_leap_day_2024_02_29() {
assert_eq!(
parse_reset_ms(&json!({ "reset": "2024-02-29T00:00:00Z" })),
Some(1_709_164_800_000)
);
}
/// Date post-2001 connue, recoupée indépendamment : `2021-01-01T00:00:00Z`
/// = 1_609_459_200 s = 1_609_459_200_000 ms.
#[test]
fn parse_reset_ms_known_post_2001_date() {
assert_eq!(
parse_reset_ms(&json!({ "reset": "2021-01-01T00:00:00Z" })),
Some(1_609_459_200_000)
);
}
// ---- parse_event : mapping rate_limit_event ----------------------------
/// `rate_limit_event` avec `rate_limit_info.resetsAt` exploitable ⇒
/// `RateLimited{Some(...)}` (l'heure de reset est extraite et normalisée).
#[test]
fn parse_event_rate_limit_with_reset_yields_rate_limited_some() {
let parsed = claude::parse_event(
r#"{"type":"rate_limit_event","rate_limit_info":{"resetsAt":1700000000},"session_id":"c"}"#,
)
.expect("parse ok");
assert_eq!(
parsed.events,
vec![ReplyEvent::RateLimited {
resets_at_ms: Some(1_700_000_000_000)
}]
);
assert_eq!(parsed.session_id.as_deref(), Some("c"));
}
/// `rate_limit_event` SANS `rate_limit_info` exploitable ⇒ `RateLimited{None}`
/// (et surtout PAS un `Heartbeat` : c'est le changement §21/LS2).
#[test]
fn parse_event_rate_limit_without_info_is_rate_limited_none_not_heartbeat() {
// rate_limit_info absent.
let absent = claude::parse_event(r#"{"type":"rate_limit_event","session_id":"c"}"#)
.expect("parse ok");
assert_eq!(
absent.events,
vec![ReplyEvent::RateLimited { resets_at_ms: None }]
);
assert_ne!(absent.events, vec![ReplyEvent::Heartbeat]);
// rate_limit_info présent mais sans clé de reset connue.
let no_key = claude::parse_event(
r#"{"type":"rate_limit_event","rate_limit_info":{"x":1},"session_id":"c"}"#,
)
.expect("parse ok");
assert_eq!(
no_key.events,
vec![ReplyEvent::RateLimited { resets_at_ms: None }]
);
}
// ---- Non-terminalité (T4) : RateLimited n'interrompt PAS ----------------
/// Au niveau séquence de `parse_event` : `rate_limit_event` puis `result` ⇒ la
/// concaténation des events est `[RateLimited, Final]` — le RateLimited s'intercale
/// et seul le Final clôt.
#[test]
fn parse_event_sequence_rate_limited_then_final_is_not_interrupted() {
let rl = claude::parse_event(
r#"{"type":"rate_limit_event","rate_limit_info":{"resetsAt":1700000000},"session_id":"c"}"#,
)
.expect("parse ok");
let fin = claude::parse_event(
r#"{"type":"result","subtype":"success","result":"fini","session_id":"c"}"#,
)
.expect("parse ok");
let mut seq = rl.events;
seq.extend(fin.events);
assert_eq!(
seq,
vec![
ReplyEvent::RateLimited {
resets_at_ms: Some(1_700_000_000_000)
},
ReplyEvent::Final {
content: "fini".to_owned()
},
]
);
}
/// Au niveau `send()` (FakeCli) : init → rate_limit_event → assistant → result ⇒
/// le flux émis est `[Heartbeat, RateLimited, TextDelta, Final]`. Le RateLimited
/// NE rompt PAS la boucle d'émission (T4) ; seul le Final clôt — on le PROUVE
/// bout-en-bout, pas seulement au niveau parse.
#[tokio::test]
async fn send_emits_rate_limited_intercalated_only_final_closes() {
let fake = FakeCli::printing(&[
r#"{"type":"system","subtype":"init","session_id":"rl-1","cwd":"/tmp","tools":[]}"#,
r#"{"type":"rate_limit_event","rate_limit_info":{"resetsAt":1700000000},"session_id":"rl-1"}"#,
r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"ap"}]},"session_id":"rl-1","parent_tool_use_id":null}"#,
r#"{"type":"result","subtype":"success","is_error":false,"result":"ok","session_id":"rl-1","num_turns":1}"#,
]);
let session = ClaudeSdkSession::new(
SessionId::new_random(),
fake.command(),
"/",
None,
None,
None,
);
let events: Vec<ReplyEvent> = session.send("x").await.expect("send ok").collect();
assert_eq!(
events,
vec![
ReplyEvent::Heartbeat,
ReplyEvent::RateLimited {
resets_at_ms: Some(1_700_000_000_000)
},
ReplyEvent::TextDelta { text: "ap".into() },
ReplyEvent::Final {
content: "ok".into()
},
],
"RateLimited s'intercale (non terminal), seul Final clôt"
);
// Exactement un Final, en dernière position.
assert_eq!(
events
.iter()
.filter(|e| matches!(e, ReplyEvent::Final { .. }))
.count(),
1
);
}
}