Files
IdeA/crates/infrastructure/src/session/claude.rs
Blomios 98bfcf4f22 feat(session-limits): LS5 — détecteur niveau 2 déclaratif + parsing temps partagé
Ajoute un détecteur niveau 2 (RateLimitParser, regex confiné à l'infra) qui
repère les mentions de limite de session dans la sortie textuelle, et factorise
le parsing d'heures dans un module pur (timeparse) partagé entre les niveaux 1
et 2. Le détecteur Claude niveau 1 est refactoré vers timeparse (~-121 lignes).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-16 20:02:24 +02:00

344 lines
16 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.

//! [`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. Le **spike S1 est résolu** : le schéma réel est vérifié
//! (2026-06-09) ; **seule [`parse_event`] (et la composition de la commande) porte
//! le format**, pas la machinerie ni le reste de l'adapter.
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use serde_json::Value;
use domain::ports::{AgentSession, AgentSessionError, ReplyEvent, ReplyStream};
use domain::sandbox::{SandboxEnforcer, SandboxPlan};
use domain::SessionId;
use super::process::{run_turn, SpawnLine};
/// Résultat du parsing d'une ligne : **zéro ou plusieurs** événements à émettre et/ou
/// un `session_id` capté (init/result). Permet à [`parse_event`] de rester **pure**
/// (aucun effet de bord) tout en remontant les deux informations.
///
/// Le champ `events` est un **vecteur** : une seule ligne `assistant` peut porter
/// **plusieurs** blocs (`content[]`) et donc produire **plusieurs** [`ReplyEvent`].
#[derive(Debug, Default, PartialEq, Eq)]
pub struct ParsedLine {
/// Événements universels à émettre (dans l'ordre), vide pour une ligne de contrôle.
pub events: Vec<ReplyEvent>,
/// `session_id` Claude capté sur cette ligne (id de conversation pour la reprise).
pub session_id: Option<String>,
}
/// **Parse une ligne du flux `stream-json` de Claude** vers le contrat universel.
///
/// # Format RÉEL vérifié 2026-06-09 (spike S1 résolu)
///
/// Commande : `claude -p "<prompt>" --output-format stream-json --verbose` ;
/// reprise : `claude --resume <session_id> -p … --output-format stream-json --verbose`.
///
/// Le flux est du **JSONL** (un objet JSON par ligne). Types réels :
///
/// - `{"type":"system","subtype":"init","session_id":"<uuid>","cwd":…,"tools":…,…}`
/// ⇒ capture le `session_id` (= id de conversation pour la reprise) **et** émet un
/// [`ReplyEvent::Heartbeat`] (preuve de vivacité non terminale : la CLI a démarré).
/// - `{"type":"rate_limit_event","rate_limit_info":{…},"session_id":"…"}`
/// ⇒ [`ReplyEvent::RateLimited`] (ARCHITECTURE §21, niveau 1 structuré) : l'heure
/// de reset est extraite de `rate_limit_info` par [`parse_reset_ms`] et normalisée
/// en époche-ms. Sans `rate_limit_info` exploitable ⇒ `RateLimited{None}` (limite
/// détectée, heure inconnue ⇒ filet humain en aval). **Non terminal** (comme le
/// `Heartbeat`) : il s'intercale, le flux continue jusqu'au `Final` ou la clôture.
/// - `{"type":"assistant","message":{"role":"assistant","content":[
/// {"type":"text","text":"…"} | {"type":"tool_use","name":"…", …}
/// ], …},"session_id":"…","parent_tool_use_id":null}`
/// ⇒ **chaque** bloc `text` ⇒ [`ReplyEvent::TextDelta`] ; **chaque** bloc `tool_use`
/// ⇒ [`ReplyEvent::ToolActivity`] (`label` = `name`). Une ligne `assistant` peut
/// donc produire **plusieurs** événements (contenu multi-blocs).
/// - `{"type":"result","subtype":"success","is_error":false,"result":"<texte final>","session_id":"<uuid>","num_turns":…,…}`
/// ⇒ [`ReplyEvent::Final`] (`content` = `result`) et confirme le `session_id`.
///
/// 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<ParsedLine, AgentSessionError> {
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 events = match value.get("type").and_then(Value::as_str) {
// init/handshake : on capte le session_id ET on émet un battement de cœur
// (preuve de vivacité non terminale : la CLI a démarré et répond).
Some("system") => vec![ReplyEvent::Heartbeat],
// Limite de session/débit (ARCHITECTURE §21, niveau 1) : on lit l'heure de
// reset dans `rate_limit_info` (au lieu de la jeter) et on émet un
// `RateLimited{resets_at_ms}` **non terminal**. Robuste : absence/illisibilité
// de `rate_limit_info` ⇒ `RateLimited{None}` (jamais d'erreur), filet humain.
Some("rate_limit_event") => {
let resets_at_ms = value.get("rate_limit_info").and_then(parse_reset_ms);
vec![ReplyEvent::RateLimited { resets_at_ms }]
}
Some("assistant") => assistant_events(&value),
Some("result") => value
.get("result")
.and_then(Value::as_str)
.map(|content| {
vec![ReplyEvent::Final {
content: content.to_owned(),
}]
})
.unwrap_or_default(),
_ => Vec::new(), // type inconnu / non pertinent : ignoré (robustesse).
};
Ok(ParsedLine { events, session_id })
}
/// Itère **TOUS** les blocs de contenu d'un message `assistant`, dans l'ordre :
/// chaque `text` ⇒ `TextDelta`, chaque `tool_use` ⇒ `ToolActivity`. Le `content`
/// est un **tableau** : un message multi-blocs produit donc plusieurs événements.
fn assistant_events(value: &Value) -> Vec<ReplyEvent> {
let Some(content) = value
.get("message")
.and_then(|m| m.get("content"))
.and_then(Value::as_array)
else {
return Vec::new();
};
let mut events = Vec::new();
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) {
events.push(ReplyEvent::TextDelta {
text: text.to_owned(),
});
}
}
Some("tool_use") => {
let label = block
.get("name")
.and_then(Value::as_str)
.unwrap_or("outil")
.to_owned();
events.push(ReplyEvent::ToolActivity { label });
}
_ => {}
}
}
events
}
/// **Extrait l'heure de reset d'une limite de débit** depuis l'objet
/// `rate_limit_info` d'un `rate_limit_event` Claude, **normalisée en époche-ms**
/// (ARCHITECTURE §21, niveau 1 structuré). Fonction **pure** (aucune I/O, aucun
/// `now`), isolée du reste de l'adapter pour être testable sans process.
///
/// # Spike §21.10-1 — format réel non garanti, parsing défensif
///
/// Le schéma exact du champ de reset n'est **pas** documenté de façon stable. On
/// gère donc défensivement plusieurs **noms de champ** plausibles, dans l'ordre :
/// `resetsAt`, `resets_at`, `reset_at`, `resetAt`, `reset`. La première clé présente
/// gagne. Sa valeur est convertie en époche-ms selon son type ([`value_to_epoch_ms`]) :
///
/// - **entier/float epoch en secondes** (magnitude `< 10^12`) ⇒ `× 1000` ;
/// - **entier/float epoch en millisecondes** (magnitude `≥ 10^12`) ⇒ tel quel ;
/// - **chaîne** : d'abord tentée comme entier/float epoch (même heuristique), sinon
/// parsée comme **ISO-8601 / RFC3339** ([`parse_rfc3339_to_ms`]).
///
/// L'heuristique secondes-vs-ms repose sur le **seuil de magnitude** [`EPOCH_MS_THRESHOLD`]
/// (`10^12`) : `10^12 ms ≈ 2001-09`, `10^12 s ≈ an 33658` — toute date plausible
/// (1970…) tombe sans ambiguïté du bon côté.
///
/// # Hypothèses & limites assumées
///
/// - Un champ **relatif** (`retryAfter` / `retry_after`, en secondes) **n'est pas**
/// exploité ici : le convertir en instant absolu exigerait `now`, or cette fonction
/// est **pure et sans horloge** (cf. `Clock` confiné à l'application). Il est donc
/// **ignoré** (documenté) ⇒ `None` ⇒ `RateLimited{None}` ⇒ filet humain. Une
/// résolution `now + delta` pourra être ajoutée côté appelant/application (LS4) qui,
/// lui, détient `Clock`.
/// - Une chaîne ISO **sans fuseau** est traitée comme **UTC** (best-effort), le reset
/// Claude étant une heure absolue. La gestion d'heure murale locale relève du
/// niveau 2 (LS5, §21.10-2), pas d'ici.
///
/// Retourne `None` (jamais d'erreur) si aucune clé connue n'est présente ou si la
/// valeur est inexploitable — l'appelant émet alors `RateLimited{None}`.
#[must_use]
pub fn parse_reset_ms(rate_limit_info: &Value) -> Option<i64> {
const RESET_KEYS: [&str; 5] = ["resetsAt", "resets_at", "reset_at", "resetAt", "reset"];
let raw = RESET_KEYS.iter().find_map(|k| rate_limit_info.get(*k))?;
value_to_epoch_ms(raw)
}
/// Convertit une [`Value`] JSON (entier / float / chaîne) en époche-ms, défensivement.
///
/// La conversion **générique** (heuristique secondes-vs-ms, parseur ISO-8601, algo
/// jour-civil) est factorisée dans [`crate::timeparse`] et **partagée** avec le
/// détecteur de niveau 2 (`ratelimit`) — pas de duplication (§21.2-T2). Seule
/// l'extraction depuis un `serde_json::Value` (typage int/float/str) reste ici.
fn value_to_epoch_ms(v: &Value) -> Option<i64> {
if let Some(i) = v.as_i64() {
return Some(crate::timeparse::int_epoch_to_ms(i));
}
if let Some(u) = v.as_u64() {
return Some(crate::timeparse::int_epoch_to_ms(i64::try_from(u).ok()?));
}
if let Some(f) = v.as_f64() {
return Some(crate::timeparse::float_epoch_to_ms(f));
}
if let Some(s) = v.as_str() {
return crate::timeparse::parse_absolute_ms(s);
}
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<Option<String>>,
/// Plan de sandbox OS **par lancement** (lot LP4-4), porté dans chaque
/// [`SpawnLine`]. `None` ⇒ aucun sandboxing (drain async natif).
sandbox: Option<SandboxPlan>,
/// Enforcer OS **par instance** (lot LP4-4), passé à [`run_turn`]. `None` ⇒ pas
/// de sandboxing même si un plan est présent (cohérent avec le chemin PTY).
sandbox_enforcer: Option<Arc<dyn SandboxEnforcer>>,
}
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. `sandbox` /
/// `sandbox_enforcer` (lot LP4-4) pilotent le sandboxing OS du tour ; `None`/`None`
/// ⇒ chemin natif inchangé.
#[must_use]
pub fn new(
id: SessionId,
command: impl Into<String>,
cwd: impl Into<String>,
seed_conversation_id: Option<String>,
sandbox: Option<SandboxPlan>,
sandbox_enforcer: Option<Arc<dyn SandboxEnforcer>>,
) -> Self {
Self {
id,
command: command.into(),
cwd: cwd.into(),
conversation_id: Mutex::new(seed_conversation_id),
sandbox,
sandbox_enforcer,
}
}
/// Compose la ligne de commande d'un tour selon l'état de conversation.
///
/// Format RÉEL vérifié 2026-06-09 :
/// - Conversation neuve : `claude -p <prompt> --output-format stream-json --verbose`.
/// - Reprise (id connu) : `claude --resume <id> -p <prompt> --output-format
/// stream-json --verbose`.
///
/// Le flag `--verbose` est **requis** : sans lui, `--output-format stream-json`
/// n'émet pas le flux JSONL ligne-à-ligne attendu par le parser.
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());
args.push("--verbose".to_owned());
SpawnLine {
command: self.command.clone(),
args,
cwd: self.cwd.clone(),
env: Vec::new(),
stdin: None,
sandbox: self.sandbox.clone(),
}
}
}
#[async_trait]
impl AgentSession for ClaudeSdkSession {
fn id(&self) -> SessionId {
self.id
}
fn conversation_id(&self) -> Option<String> {
self.conversation_id.lock().expect("mutex sain").clone()
}
async fn send(&self, prompt: &str) -> Result<ReplyStream, AgentSessionError> {
let spec = self.build_spawn_line(prompt);
let raw_lines = run_turn(&spec, None, self.sandbox_enforcer.as_ref()).await?;
let mut events = Vec::new();
let mut captured_id = None;
'lines: for line in &raw_lines {
let parsed = parse_event(line)?;
if let Some(id) = parsed.session_id {
captured_id = Some(id);
}
// Aplatit : une ligne `assistant` multi-blocs rend plusieurs événements.
// Le `Final` est le **seul** événement terminal (contrat de port, §21-T4) :
// on arrête d'émettre dès qu'on l'a vu. Un `RateLimited` (comme un
// `Heartbeat`) est **non terminal** : il s'intercale et ne rompt JAMAIS la
// boucle — un tour limité clos sans `Final` reste une fin gracieuse (le
// traitement « limité » est en aval, application lot LS4).
for event in parsed.events {
let is_final = matches!(event, ReplyEvent::Final { .. });
events.push(event);
if is_final {
break 'lines;
}
}
}
// 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(())
}
}