From 5e10b5eb42b88c64ba7006a30867d480ad20b57e Mon Sep 17 00:00:00 2001 From: Blomios Date: Tue, 9 Jun 2026 17:54:48 +0200 Subject: [PATCH] =?UTF-8?q?feat(agent):=20fondation=20ex=C3=A9cution=20str?= =?UTF-8?q?uctur=C3=A9e=20des=20agents=20IA=20(D0+D1)=20=E2=80=94=20=C2=A7?= =?UTF-8?q?17?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Pivot orchestration : agents IA pilotés via leur mode programmatique/JSON (capture déterministe), au lieu du TUI brut + self-report. §16 (idea/MCP) marquée remplacée comme voie principale. - D0 (domaine) : port AgentSession + AgentSessionFactory, types ReplyEvent /ReplyStream/AgentSessionError, champ AgentProfile.structured_adapter (Option, skip si None ⇒ zéro régression), catalogue Claude/Codex annotés. - D1 (application) : registre StructuredSessions (jumeau de TerminalSessions), agrégateur LiveSessions{pty,structured} derrière LiveAgentRegistry (vivant si PTY OU structuré, surface du trait inchangée), helper send_blocking (draine le ReplyStream jusqu'au Final, Timeout sans tuer la session). Tests : domaine 16+2 ; application registre 11 + send_blocking 9 ; workspace 0 échec. A/B intacts. Aucun adapter concret (D2), pas de Tauri/front. Co-Authored-By: Claude Opus 4.8 --- ARCHITECTURE.md | 622 ++++++++++++++++++ Cargo.lock | 1 + Cargo.toml | 2 +- crates/application/Cargo.toml | 3 + crates/application/src/agent/catalogue.rs | 8 +- crates/application/src/agent/mod.rs | 3 + crates/application/src/agent/structured.rs | 73 ++ crates/application/src/lib.rs | 8 +- crates/application/src/terminal/mod.rs | 2 +- crates/application/src/terminal/registry.rs | 255 ++++++- crates/application/tests/profile_usecases.rs | 36 + crates/application/tests/send_blocking_d1.rs | 222 +++++++ .../tests/structured_registry_d1.rs | 305 +++++++++ crates/domain/Cargo.toml | 1 + crates/domain/src/ports.rs | 123 ++++ crates/domain/src/profile.rs | 35 + crates/domain/tests/structured_session_d0.rs | 396 +++++++++++ 17 files changed, 2084 insertions(+), 11 deletions(-) create mode 100644 crates/application/src/agent/structured.rs create mode 100644 crates/application/tests/send_blocking_d1.rs create mode 100644 crates/application/tests/structured_registry_d1.rs create mode 100644 crates/domain/tests/structured_session_d0.rs diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 0beba57..d62de68 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -633,6 +633,7 @@ IdeA/ | L13 | **OrchestratorApi** | File-watcher `.ideai/requests/`, port `OrchestratorApi`, adapter `FsOrchestratorAdapter`, protocole `agent.run`/`agent.stop`/`agent.attach`/`agent.detach`/`agent.message`, sessions agent visibles ou arrière-plan. | `domain/agent`, `application/agent`, `infrastructure/orchestrator`, `app-tauri`, `frontend/features/agents` | | L14 | **Mémoire** | Base de connaissance projet model-agnostic. Modèle 2 étages : `.md` source de vérité (port `MemoryStore`), rappel adaptatif (port `MemoryRecall`), embeddings déclaratifs (port `Embedder`), bascule auto sur seuil. Découpé en sous-lots **A/B/C** (voir §14.5). | `domain/memory`, `application/memory`, `infrastructure/store`, `frontend/features/memory` | | L15 | **Agent = entité à session persistante** | Hot-swap de l'AI profile d'un agent existant (chantier A) + reprise des sessions au redémarrage d'IdeA / réouverture projet (chantier B). Fondation commune « l'agent porte un cycle de vie de session ». Découpé en sous-lots **A0/A1/A2 + B0/B1/B2** (voir §15). | `domain/agent`, `application/agent`, `application/layout`, `infrastructure/store`, `app-tauri`, `frontend/features/agents`, `frontend/features/layout` | +| L16 | **Orchestration v3 — invocation native (CLI `idea` universelle + MCP optionnel)** | **Voie principale (cross-model)** : binaire `idea` (`ask` synchrone, `launch`, `reply`, `list-agents`, `next`) + skill built-in « Orchestration IdeA » auto-assigné, branchés sur le **même** `OrchestratorService` ; messagerie inter-agents synchrone via port `AgentReplyChannel` + outbox `.ideai/outbox/` ; tâche à un agent vivant via **inbox** `.ideai/inbox/`. **Optionnel/postérieur** : surface MCP (`idea_*`) + capacité `mcp` déclarative sur le profil, par-dessus le même backend. Découpé en bloc **universel (C0, C1, C-univ-1/2)** puis **MCP (C-mcp-0/1, C4)** — voir §16. | `crates/idea-cli`, `domain/{orchestrator,skill,events,ports,profile}`, `application/orchestrator`, `infrastructure/orchestrator` (+`mcp` optionnel), `app-tauri`, `frontend/features/agents` | --- @@ -739,6 +740,8 @@ Exemple `skill.create` : **Impact UI/UX** : le layout n'est pas la source de vérité du travail agent. La grille affiche des vues attachées aux sessions. L'UI doit exposer un registre des sessions visibles et arrière-plan, avec actions ouvrir dans une cellule, détacher, déplacer, arrêter, et afficher la relation `requestedBy`/`targetAgent` quand elle existe. +> **Évolution v3 (§16)** : ce protocole fichier reste le **contrat partagé**. L'**orchestration v3** en fait la **voie principale et universelle** via un **binaire `idea`** (client mince sur le PATH de l'agent) + un **skill built-in** auto-assigné, qui écrivent ces fichiers à la place de l'agent — et comble le manque de §14.3 : la **messagerie inter-agents synchrone** (`agent.ask` : transmet une tâche et **renvoie la réponse de contenu** inline via l'outbox, alors qu'ici `task` est ignoré pour un agent déjà vivant et la réponse n'est qu'un ACK). Une tâche à un agent **déjà vivant** passe par une **inbox** qu'il relit (pas un write dans son TUI). Une **surface MCP** (outils `idea_*`) est ajoutée **optionnellement** par-dessus le même `OrchestratorService`. Voir §16. + --- ### 14.4 Git = intégration optionnelle, zéro dépendance fonctionnelle @@ -1195,6 +1198,625 @@ La **reprise effective** d'un agent choisi dans le panneau **ne crée aucun nouv --- +## 16. Orchestration v3 — invocation native d'agents (CLI `idea` universelle + surface MCP optionnelle) — révisé 2026-06-09 + +> **⚠️ REMPLACÉE COMME VOIE PRINCIPALE par §17 (pivot 2026-06-09).** Le chef d'orchestre a tranché : on abandonne, **comme voie principale**, l'orchestration via TUI brut + binaire `idea`/skill auto-rapporté (fiabilité insuffisante : dépend du bon vouloir du modèle d'appeler `idea reply`/`idea next`). La nouvelle voie principale est **§17 — exécution structurée des agents IA via le port `AgentSession`** (mode programmatique par modèle, capture déterministe de la réponse, rendez-vous synchrone intrinsèque à `send()`). En conséquence : +> - **Abandonné/déprécié (voie principale)** : le binaire `idea` (`idea ask/reply/next`), le skill built-in « Orchestration IdeA » comme *mécanisme de délégation auto-rapporté*, l'**inbox** `.ideai/inbox/`, l'**outbox** `.ideai/outbox/`, le port `AgentReplyChannel`/`OutboxReplyChannel`, et le rendez-vous outbox des lots **C0/C1/C-univ-***. Le rendez-vous synchrone est désormais **intrinsèque** à `AgentSession::send() -> Reply` (§17.1) : plus besoin d'outbox ni de corrélation fichier. +> - **Conservé (repli/compat)** : le protocole fichier `.ideai/requests/` + `FsOrchestratorWatcher` (§14.3) reste un **adapter entrant de repli** (un agent ou un script qui écrit une requête à la main). `OrchestratorService` route désormais la délégation inter-agents via le port `AgentSession` (§17.4), pas via l'outbox. +> - **Non démarré ⇒ supprimé du périmètre** : les lots C0/C1/C-univ-1/C-univ-2 et le bloc MCP (C-mcp-*) **ne sont plus à livrer** tels quels. La « version B » (UI chat / sortie structurée) évoquée en §16.9 comme épic futur **devient la voie principale §17**. +> Le reste de §16 est laissé **pour mémoire/historique** (raisonnement, état du terrain) ; ne pas l'implémenter sans relire §17. + +> **Fondation** : v3 ne réécrit pas §14.3. Elle comble sa lacune fonctionnelle — la **messagerie inter-agents synchrone** (`ask_agent`) — et ajoute des **portes d'entrée** au-dessus du **même** `OrchestratorService`. Une seule logique applicative, plusieurs **adapters entrants** qui se ramènent tous au même `OrchestratorCommand` enrichi. +> +> **Décision produit verrouillée (révision 2026-06-09, non rediscutée — actée ici)** : la **voie principale et la garantie cross-model** est un **plancher universel** = un **skill built-in « Orchestration IdeA » auto-assigné à TOUT agent** + un **petit binaire CLI `idea`** posé par IdeA sur le `PATH` du sandbox de l'agent. L'agent délègue par une simple commande shell (`idea ask ""` bloque et imprime la réponse inline ; `idea launch`, `idea reply`, `idea list-agents`) — **aucun JSON manipulé par l'agent, aucun parsing de TUI**. Sous le capot, `idea` est un **client mince** qui écrit dans `.ideai/requests/` et attend `.ideai/outbox/` : il s'appuie EXACTEMENT sur `OrchestratorService::dispatch` + le port `AgentReplyChannel` + l'outbox (C0/C1 ci-dessous). Tout agent sachant lancer une commande shell sait déléguer ⇒ **zéro support modèle spécial requis** ; valider avec Claude + Codex garantit le cross-model. +> +> **MCP est rétrogradé en confort OPTIONNEL** par-dessus le **même** backend : un adapter entrant supplémentaire (outils typés `idea_*`) pour les CLIs qui le supportent, postérieur et non bloquant. Les spikes MCP ne conditionnent plus la garantie cross-model. +> +> **Hors périmètre C (épic futur séparé, noté pour cohérence)** : une « version B » (UI chat / agent headless à sortie structurée) reste compatible avec ce backend mais n'est pas requise pour le cross-model. Voir §16.9. +> +> Sémantique `ask` (les **deux** voies, identique) : lance/réveille la cible, transmet la tâche, **attend et renvoie le contenu** de sa réponse inline, corrélé via l'outbox. + +### 16.0 État du terrain (lu dans le code, pas présumé) + +| Pièce | Existe ? | Référence code | +|---|---|---| +| `OrchestratorRequest`/`OrchestratorCommand` (modèle pur, validé) | ✅ — actions `agent.run/stop/update_context`, `skill.create` | `domain/src/orchestrator.rs` | +| `OrchestratorService::dispatch` (un seul chemin applicatif, réutilise les use cases UI) | ✅ | `application/src/orchestrator/service.rs` | +| `FsOrchestratorWatcher` (adapter entrant fichier + `*.response.json` ACK) | ✅ | `infrastructure/src/orchestrator/mod.rs` | +| Profil déclaratif `AgentProfile` (+ `SessionStrategy` optionnel) | ✅ | `domain/src/profile.rs` | +| `SessionInspector` (lecture best-effort d'un transcript CLI, optionnel) | ✅ | `domain/src/ports.rs`, `application/src/agent/inspect.rs` | +| Invariant « 1 session vivante par agent » + `rebind_agent_node`/`session_for_agent` | ✅ | `application/src/terminal/registry.rs` | +| Prose « # Orchestration IdeA » injectée dans le convention file | ✅ — mais **prose libre**, à transformer en **skill built-in** documentant `idea` | `application/src/agent/lifecycle.rs` (`compose_convention_file`) | +| `Skill` (entité, scope `Global`/`Project`, injection convention file) | ✅ — pas de notion `Builtin` ni d'auto-assignation universelle | `domain/src/skill.rs`, `application/src/skill/*`, `compose_convention_file` (param `skills`) | +| `SpawnSpec.env: Vec<(String,String)>` (env injectable au lancement CLI) | ✅ — vecteur d'overrides d'env passé au `ProcessSpawner` | `domain/src/ports.rs` (`SpawnSpec`), `infrastructure/src/runtime/mod.rs` | +| **Binaire CLI `idea`** (client mince requests→outbox sur le PATH du run dir) | ❌ totalement absent | — | +| **Skill built-in « Orchestration IdeA » auto-assigné à tout agent** | ❌ — aujourd'hui prose libre, non modélisée comme skill | — | +| **`task` transmis à un agent déjà vivant** | ❌ — ignoré (replié en `context`, utilisé seulement à la création) | `service.rs` `spawn_agent` | +| **Réponse de contenu** (réveil du demandeur, corrélation requête↔réponse) | ❌ — la réponse n'est qu'un ACK de cycle de vie (`detail`) | `service.rs` / watcher | +| **Capacité MCP sur le profil** | ❌ totalement absent | — | +| **Serveur MCP / config MCP par CLI** | ❌ totalement absent | — | + +**Conclusion** : v3 = **(1)** ajouter une variante de commande qui *transmet une tâche et attend une réponse de contenu* (le vrai trou, port `AgentReplyChannel` + outbox) ; **(2)** le **plancher universel** — un binaire `idea` (client mince requests→outbox, posé sur le PATH du run dir) + le **skill built-in « Orchestration IdeA »** auto-assigné qui le documente — qui devient **la voie principale et la garantie cross-model** ; **(3) optionnel/postérieur** : un adapter entrant MCP qui se branche sur le `OrchestratorService` *exactement* comme le watcher et la CLI, + capacité déclarative sur le profil + injection de la config MCP par CLI. Aucun use case agent/terminal n'est réécrit ; `idea` et MCP partagent le **même** `OrchestratorService::dispatch` et le **même** outbox. + +### 16.1 Décisions tranchées (avec justification) + +0. **Plancher universel = binaire `idea` + skill built-in, voie PRINCIPALE et garantie cross-model** — la conscience d'orchestration d'un agent ne repose plus sur une prose libre « rappelle-toi d'écrire un JSON » mais sur **deux artefacts concrets** : (a) un **petit binaire CLI `idea`** posé par IdeA sur le `PATH` de l'agent (via son run dir isolé `.ideai/run//bin`, §14.1), et (b) un **skill built-in « Orchestration IdeA »** auto-assigné à **tout** agent, dont le `.md` documente les commandes `idea ask/launch/reply/list-agents`. *Justification* : universalité (principe fondateur) — toute CLI sait lancer une commande shell, donc tout modèle (Claude/Codex/Gemini/custom) sait déléguer **sans support spécial** ; le mécanisme testé (« l'agent exécute `idea` ») est **identique pour tous les modèles**, donc valider sur 2 CLIs garantit le cross-model. `idea` est un **client mince sans logique métier** : il (dé)sérialise vers `.ideai/requests/` et attend `.ideai/outbox/` — la logique vit dans `OrchestratorService` (DRY). + +1. **MCP rétrogradé en adapter entrant OPTIONNEL et postérieur** — le serveur MCP reste un **driving adapter d'infrastructure** (`infrastructure/src/orchestrator/mcp/`) qui appelle le **même** `OrchestratorService::dispatch` et lit/écrit le **même** outbox, mais il n'est **plus** la voie principale : c'est un **confort** (outils typés natifs) pour les CLIs qui le déclarent, ajouté **après** le plancher universel. Trois portes d'entrée substituables : **CLI `idea`** (universelle, principale), `FsOrchestratorWatcher` (fichier brut, repli historique §14.3), serveur MCP (optionnel). *Justification* : DRY + hexagonal — cible, identité, mémoire, observabilité UI passent par le seul chemin applicatif ; les spikes MCP (transport par CLI) ne conditionnent plus la garantie cross-model. + +2. **`OrchestratorCommand` gagne une variante `AskAgent` (transmission de tâche + attente de réponse)** — distincte de `SpawnAgent` (fire-and-forget). C'est la brique manquante de §14.3. `SpawnAgent` reste l'équivalent de `idea_launch_agent` (fire-and-forget) ; `AskAgent` porte `target`, `task`, et une **corrélation** (`request_id`). *Justification* : `parse, don't validate` — le modèle pur rend explicite « j'attends une réponse » vs « je lance et j'oublie », au lieu de surcharger `task` silencieusement comme aujourd'hui. + +3. **Le retour synchrone passe par un nouveau port `AgentReplyChannel` (corrélation requête↔réponse), PAS par `SessionInspector`** — `SessionInspector` lit best-effort un transcript *propre à chaque CLI* (fragile, non universel, déjà « best-effort par construction »). Pour un retour **fiable et model-agnostic**, on ne *devine* pas la fin de tour : on demande à la cible d'**écrire sa réponse dans un outbox** `.ideai/outbox/.json` (instruction injectée + outil MCP `idea_reply`), et l'appelant **attend cette corrélation** (await/poll + timeout). *Justification* : universalité (principe fondateur — marche pour Claude/Codex/Gemini/custom sans parser leur format) et frontière nette (le domaine ne connaît qu'un id de corrélation + un contenu, jamais un transcript). + +4. **Capacité MCP = champ optionnel `mcp` sur `AgentProfile`** (descripteur déclaratif), `None` par défaut ⇒ comportement actuel (repli fichier). Ajouter une CLI MCP = **donnée, pas code** (Open/Closed), comme `session`/`contextInjection`. *Justification* : cohérence avec §9 ; zéro régression pour les profils existants (sérialisation `skip_serializing_if = None`). + +5. **Repli homogène** — un agent dont le profil n'a **pas** de bloc `mcp` continue d'utiliser le protocole fichier `.ideai/requests` (prose injectée inchangée). Un agent MCP voit les outils typés. Les **deux** routes produisent le même `OrchestratorCommand` et, pour `ask`, écrivent/lisent le **même outbox**. *Justification* : « rien d'imposé, tout fonctionnel » — un runtime sans MCP n'est jamais bloqué ; un runtime MCP gagne la conscience native + arguments validés. + +6. **Timeout borné + sémantique d'erreur explicite** — `ask_agent` a un `timeout` (défaut configurable). À l'expiration : la cible **reste vivante** (on ne tue rien), l'outil renvoie une **erreur typée** `Timeout` (l'appelant décide). *Justification* : pas de blocage indéfini d'une conversation appelante ; cohérent avec l'invariant « stop est une action explicite » (§14.3). + +7. **Interaction avec « 1 session vivante par agent »** — `ask_agent` sur une cible **déjà vivante** ne relance pas : il **transmet la tâche à la session existante** (write PTY de la consigne + corrélation) et attend l'outbox. Sur une cible **éteinte** : `LaunchAgent` d'abord (même chemin que `SpawnAgent`), puis transmission. *Justification* : respecte l'invariant déjà enforce ; réutilise `session_for_agent`/`rebind_agent_node`. + +8. **`idea reply` et l'outbox sont model-agnostic** — la cible rend sa réponse soit par la **commande `idea reply ""`** (voie universelle : `idea` écrit l'outbox corrélé pour elle), soit par l'**outil MCP** `idea_reply(requestId, content)` (si profil MCP). Un seul format d'outbox `.ideai/outbox/.json` lu par l'adapter. *Justification* : symétrie parfaite des routes ⇒ l'appelant attend la même chose quelle que soit la CLI cible ; l'agent cible ne manipule jamais de JSON. + +9. **Le skill built-in « Orchestration IdeA » remplace la prose libre** — on introduit un **scope `Builtin`** sur l'entité `Skill` (à côté de `Global`/`Project`, §14.2). Un skill built-in est **fourni par IdeA** (contenu `.md` embarqué, non éditable par l'utilisateur), **auto-assigné à tout agent** à l'activation (pré-pendu à la liste de skills déjà injectée par `compose_convention_file`), et documente la CLI `idea`. *Justification* : cohérence stricte avec le système de skills §14.2 (« abstraction universelle de workflows réutilisables ») et la philosophie « skills intégrés / principe universel IdeA » — la conscience d'orchestration devient un workflow **versionné et testable**, pas un littéral en dur dans `compose_convention_file`. Le bloc prose « # Orchestration IdeA » actuel est **retiré** de `compose_convention_file` et **migré** dans le `.md` du skill built-in (réutilise le canal d'injection existant, zéro mécanisme neuf). + +10. **Délivrer une tâche à un agent DÉJÀ VIVANT = inbox relue par la cible, PAS write stdin** — décision tranchée du point dur §16.7. On **n'injecte pas** la consigne par `write` PTY dans le TUI en cours de rendu : on dépose la tâche dans une **inbox** `.ideai/inbox//.json` que la cible **relit elle-même** (le skill built-in lui apprend : « à chaque tour, traite ta prochaine tâche via `idea next` / lis ton inbox »). *Justification* : (a) écrire dans le PTY d'un TUI en train de rendre **corrompt l'affichage et entrelace les frappes** (bug connu « accents / ordre d'écriture » — writes non sérialisés par handle, cf. mémoire `terminal-input-accents-ordering`) ; un TUI plein écran (Claude Code, etc.) **n'a pas de prompt shell** où coller du texte. (b) L'inbox est **durable et corrélée** (survit au redémarrage, sérialise naturellement N `ask` concurrents par cible en une **file FIFO par agent**, §16.7-3). (c) Symétrie avec l'outbox : requête et réponse transitent par le **même médium fichier**, model-agnostic, déjà éprouvé (§14.3 notify+poll). Le `write` PTY reste réservé au **premier lancement** d'une cible éteinte (consigne initiale passée comme argument/contexte au spawn, pas dans un TUI vivant). + +### 16.2 Modèle de domaine (ajouts purs, I/O-free) + +Tout vit dans `domain/src/orchestrator.rs` (modèle) + `domain/src/profile.rs` (capacité) + `domain/src/events.rs` (event). **Aucun accès I/O** : la corrélation est un VO ; l'attente/poll/écriture outbox sont infra. + +```rust +// domain/src/orchestrator.rs — VO de corrélation (newtype validé, non vide) +pub struct CorrelationId(String); // ex. un Uuid stringifié, généré par l'adapter entrant + +// Nouvelle variante de commande : transmettre une tâche ET attendre une réponse de contenu. +pub enum OrchestratorCommand { + SpawnAgent { /* … inchangé … */ }, // = idea_launch_agent (fire-and-forget) + StopAgent { name: String }, + UpdateAgentContext { name: String, context: String }, + CreateSkill { /* … inchangé … */ }, + /// NOUVEAU : `idea_ask_agent` — lance/réveille `target`, lui transmet `task`, + /// et l'appelant attend la réponse corrélée par `correlation`. + AskAgent { + target: String, + task: String, + correlation: CorrelationId, + visibility: OrchestratorVisibility, // background par défaut + }, +} + +// Réponse de CONTENU (distincte de l'ACK de cycle de vie OrchestratorResponse infra). +// Pure : ce que la cible a produit, corrélé. L'infra la (dé)sérialise depuis l'outbox. +pub struct AgentReply { + pub correlation: CorrelationId, + pub from_agent: String, // nom de la cible qui répond + pub content: String, // sortie inline rendue à l'appelant +} +``` + +`OrchestratorRequest::validate` apprend l'action **`agent.ask`** (et l'alias outil MCP `idea_ask_agent`) ⇒ `AskAgent` (champs requis : `targetAgent`, `task` ; `correlation` injectée par l'adapter si absente du fichier). Les invariants existants (champs requis, scopes, visibility) sont **inchangés** ; on ajoute une branche + ses tests, façon `parse, don't validate`. + +> **Note (révision)** : la capacité MCP ci-dessous appartient désormais au **bloc MCP optionnel** (lot `C-mcp-0`), **pas** au cœur C0. Elle reste cadrée ici par cohérence, mais n'est plus un prérequis de la voie universelle. + +```rust +// domain/src/profile.rs — capacité MCP déclarative (Open/Closed, comme SessionStrategy) +pub struct McpCapability { + /// Comment IdeA déclare son serveur MCP à CETTE CLI. Chaque CLI a sa propre + /// conf MCP : on décrit le « où/comment écrire » de façon déclarative. + pub config_strategy: McpConfigStrategy, + /// Nom logique sous lequel les outils idea_* sont exposés (ex. "idea"). + pub server_name: String, +} +pub enum McpConfigStrategy { + /// Écrire un fichier de conf MCP au chemin attendu par la CLI (relatif au cwd + /// agent), au format JSON propre à la CLI (ex. .mcp.json pour Claude Code). + ConfigFile { target: String }, + /// Passer le serveur via un flag de lancement (ex. --mcp-config {path}). + Flag { flag: String }, + /// Variable d'environnement pointant la conf. + Env { var: String }, +} + +pub struct AgentProfile { + // … champs existants inchangés … + /// Capacité MCP optionnelle. `None` (défaut) ⇒ repli protocole fichier §14.3. + pub mcp: Option, +} +``` + +```rust +// domain/src/skill.rs — nouveau scope pour le skill built-in d'orchestration (Open/Closed). +pub enum SkillScope { + Global, // store global IDE (existant) + Project, // .ideai/skills/ (existant) + Builtin, // NOUVEAU : fourni par IdeA, non éditable, auto-assigné à tout agent +} +// Le skill built-in « Orchestration IdeA » (contenu .md embarqué) est exposé par +// le SkillStore (scope Builtin) et pré-pendu aux skills d'un agent à l'activation. +``` + +```rust +// domain/src/events.rs — event de contenu (calqué sur AgentLaunched) +DomainEvent::AgentReplied { + from_agent: AgentId, // la cible qui a répondu + correlation: String, // pour relier la réponse à la demande dans l'UI +} +``` +> `AgentReplied` est **observabilité** (l'UI montre « Architect a répondu à Main »). Le **retour de valeur** à l'appelant MCP ne passe **pas** par l'EventBus (qui est fire-and-forget) mais par l'attente de l'outbox côté adapter (§16.4) — l'event ne fait que *notifier* l'UI. + +### 16.3 Port(s) — frontière domaine + +Un **seul** nouveau port, fin (ISP), pour le rendez-vous requête↔réponse. Tout le reste réutilise l'existant. + +```rust +// domain/src/ports.rs +#[async_trait] +pub trait AgentReplyChannel: Send + Sync { + /// Publie la réponse d'une cible (appelé quand l'outbox `.json` + /// apparaît, ou par la commande `idea_reply`). Idempotent par corrélation. + async fn publish_reply(&self, reply: AgentReply) -> Result<(), ReplyError>; + + /// Attend (await, borné par `timeout`) la réponse corrélée. C'est ce que + /// `idea_ask_agent` bloque dessus. Universel : ne connaît qu'un id + un contenu. + async fn await_reply( + &self, + correlation: &CorrelationId, + timeout: Duration, + ) -> Result; // ReplyError::Timeout à l'expiration +} +``` + +- **Consommé par** : `OrchestratorService` (côté `AskAgent` : `await_reply`) et l'adapter qui détecte l'outbox / l'outil `idea_reply` (`publish_reply`). +- **Implémenté par** : `OutboxReplyChannel` (`infrastructure/src/orchestrator/`) — un registre de `oneshot`/`Notify` en mémoire **adossé** au répertoire `.ideai/outbox/` : l'écriture d'un `.json` (par une cible repli-fichier) **ou** un appel MCP `idea_reply` résolvent la même attente. Pour les cibles distantes/redémarrage, l'outbox fichier est la source durable ; l'in-memory `Notify` est l'optimisation latence (même philosophie que notify+poll du watcher §14.3). + +> **Pourquoi pas `SessionInspector`** : il est **best-effort** et **par-CLI** ; en faire la brique d'un retour *fiable* violerait l'universalité. `AgentReplyChannel` est *explicite* : la cible *déclare* sa réponse, on n'infère rien. + +### 16.3bis Plancher universel — binaire `idea` (adapter entrant principal) + skill built-in + +**Nouvel artefact : un binaire `idea`.** C'est un **driving adapter entrant**, pair universel du `FsOrchestratorWatcher` et du serveur MCP, mais qui vit dans un **processus séparé** (lancé par l'agent depuis son shell) et qui parle au backend IdeA **par les mêmes fichiers** que §14.3. + +- **Crate** : nouveau binaire `crates/idea-cli/` (binaire autonome, dépendances minimales). Il **ne** lie **pas** `application`/`infrastructure` ; c'est un **client mince** qui ne connaît que le **protocole fichier** `.ideai/{requests,inbox,outbox}/` (le contrat partagé). Il découvre le project root via une variable d'env injectée (`IDEA_PROJECT_ROOT`) et son identité d'agent appelant via `IDEA_AGENT` (toutes deux posées dans `SpawnSpec.env` au lancement, comme le run dir). *Justification hexagonale* : `idea` est un **adapter entrant out-of-process** ; la frontière entre lui et le cœur est le **protocole fichier**, pas un appel de fonction. Le cœur (`OrchestratorService` + `FsOrchestratorWatcher`) ne sait pas si le fichier de requête vient de `idea`, d'un agent qui l'a écrit à la main, ou d'un test. + +| Commande `idea` | Écrit | Attend | Effet rendu à l'agent | +|---|---|---|---| +| `idea ask ""` | `.ideai/requests//.json` (`type: agent.ask`, `correlation`) | `.ideai/outbox/.json` (poll + timeout) | **bloque**, imprime `reply.content` sur stdout (feeling natif type outil `Task`) | +| `idea launch ` | `.ideai/requests//.json` (`type: agent.run`) | rien (fire-and-forget) | retourne immédiatement (ACK) | +| `idea reply ""` | `.ideai/outbox/.json` (corrélation lue depuis `IDEA_CORRELATION`/inbox courante) | — | la cible rend sa réponse à l'appelant | +| `idea list-agents` | requête de découverte | la liste | imprime les agents du projet | +| `idea next` *(cible vivante)* | — | lit `.ideai/inbox//` (FIFO) | imprime la prochaine tâche + sa `correlation` (cf. décision 10) | + +- **Mise sur le PATH (§14.1)** : à l'activation d'un agent, IdeA matérialise `idea` dans le run dir (`/bin/idea`, par symlink/copie du binaire embarqué dans le bundle Tauri) et **préfixe `PATH`** via `SpawnSpec.env` (`PATH=/bin:`). Ainsi la commande `idea` est résolue **sans installation système**, par agent, exactement où vit déjà le convention file. Aucun nouveau port : on réutilise le `env` déjà transporté par `SpawnSpec`. +- **Skill built-in** : le `.md` du skill « Orchestration IdeA » (scope `Builtin`, décision 9) documente ces commandes et la consigne « ne jamais utiliser les subagents natifs du fournisseur ; pour traiter une tâche entrante, lis ton inbox via `idea next` ». Il est **auto-assigné à tout agent** et injecté par le canal skills existant de `compose_convention_file`. +- **Réutilisation DRY** : la requête `agent.ask` produite par `idea` est **le même** `OrchestratorRequest` que celui du watcher → **même** `validate` → **même** `AskAgent` → **même** `OrchestratorService::dispatch` → **même** `await_reply` sur le **même** `OutboxReplyChannel`. `idea` n'ajoute **aucune** logique métier ; il ne fait que **traduire une ligne de commande en fichier** et **attendre l'outbox**. + +### 16.4 Adapter MCP (OPTIONNEL, postérieur) — `infrastructure/src/orchestrator/mcp/` + +Nouvel adapter **entrant** (driving), strict pair du `FsOrchestratorWatcher` : + +- **Serveur MCP** (un par projet ouvert, comme un watcher par projet) exposant les outils : + + | Outil MCP | Mappe vers | Effet | + |---|---|---| + | `idea_ask_agent(target, task) → reply` | `OrchestratorCommand::AskAgent` | génère `CorrelationId`, `dispatch`, **await_reply** (timeout), renvoie `content` inline | + | `idea_launch_agent(target, visibility)` | `OrchestratorCommand::SpawnAgent` | fire-and-forget (équiv. `agent.run`) | + | `idea_list_agents() → […]` | `ListAgents` (via service) | découverte | + | `idea_reply(requestId, content)` | `AgentReplyChannel::publish_reply` | la **cible** rend sa réponse corrélée | + | (déjà couverts) `idea_create_skill`, `idea_update_context`, `idea_stop_agent` | commandes existantes | parité avec §14.3 | + +- **Transport** : le serveur MCP est lancé par IdeA et **branché à chaque CLI MCP** via la `McpConfigStrategy` du profil cible, au moment du `LaunchAgent` (IdeA matérialise la conf MCP — `ConfigFile`/`Flag`/`Env` — dans le cwd isolé `.ideai/run//`, comme le convention file §14.1). Choix stdio vs socket = détail d'implémentation de l'adapter (point ouvert §16.7), **invisible au domaine/application**. +- **Réutilisation** : l'adapter ne contient **aucune** logique de cycle de vie — il (dé)sérialise les appels d'outils → `OrchestratorCommand` → `OrchestratorService::dispatch`, exactement comme `dispatch_file`. La seule logique neuve est l'`await_reply` pour `idea_ask_agent`. + +**Injection / composition** : `app-tauri/src/state.rs` instancie `OutboxReplyChannel` (port `AgentReplyChannel`), le passe à `OrchestratorService` (nouvelle dépendance) **et** démarre, par projet ouvert, le serveur MCP **à côté** du `FsOrchestratorWatcher` (même hook `ensure_orchestrator_watch`). Aucun autre crate ne connaît MCP. + +### 16.5 Évolution de `OrchestratorService` (réutilisation maximale) + +`OrchestratorService` gagne **un** champ (`reply_channel: Arc`) et **une** branche `AskAgent` : + +``` +dispatch(AskAgent { target, task, correlation, visibility }): + 1. Résoudre l'agent cible (find_agent_id_by_name) — NotFound sinon. + 2. Vivant ? (sessions.session_for_agent) + a. Oui → déposer la tâche dans l'INBOX de la cible (décision 10) : + écrire `.ideai/inbox/{target}/{correlation}.json` { task, correlation }. + La cible la relit via `idea next` à son tour suivant — PAS de write PTY + dans le TUI vivant (évite la corruption d'affichage / l'entrelacement). + N `ask` concurrents ⇒ FIFO naturelle par cible (§16.7-3). + b. Non → LaunchAgent (comme SpawnAgent), avec la tâche en consigne initiale + (argument/contexte de spawn, pas un write dans un TUI) + la même + instruction de réponse corrélée. + 3. reply = reply_channel.await_reply(&correlation, timeout).await? // borné, Timeout typé + 4. publish(AgentReplied { from_agent, correlation }). + 5. Retourner reply.content (l'adapter MCP le renvoie inline ; le watcher l'écrit + dans `*.response.json` pour la route fichier). +``` + +> `SpawnAgent`, `StopAgent`, `UpdateAgentContext`, `CreateSkill` **inchangés**. La transmission de tâche (étape 2a, via inbox) corrige enfin le bug §14.3 « task ignoré pour un agent existant » — et le fait pour **les trois** routes entrantes (`idea`, watcher fichier, MCP optionnel), qui portent toutes `agent.ask`. + +### 16.6 Conformité hexagonale & SOLID + +- **Règle de dépendance** : ajouts domaine **purs** (`CorrelationId`, `AskAgent`, `AgentReply`, `SkillScope::Builtin`, `McpCapability`, event `AgentReplied`) ; le seul port neuf (`AgentReplyChannel`) est un **trait du domaine**, implémenté en infra. Le binaire `idea`, le serveur MCP, l'inbox et l'outbox sont **exclusivement** hors-domaine. Le domaine ignore la CLI `idea`, MCP, stdio, JSON-RPC, l'inbox/outbox FS. +- **Le binaire `idea` est un adapter entrant out-of-process** : sa frontière avec le cœur est le **protocole fichier** `.ideai/{requests,inbox,outbox}/` (le contrat), pas un appel de fonction. `crates/idea-cli/` ne lie ni `application` ni `infrastructure` ⇒ pas de fuite de couche. Il est **substituable** au watcher (mêmes fichiers) et au serveur MCP (même `dispatch`/outbox). +- **S** : `AgentReplyChannel` = un seul rôle (rendez-vous corrélé). `idea` = une seule techno d'entrée (CLI→fichier). L'adapter MCP = une seule techno d'entrée. `OrchestratorService` garde sa responsabilité (traduire commande → use cases). +- **O** : ajouter le plancher universel = données + un binaire client, **sans toucher** au domaine ni aux use cases. Ajouter une CLI MCP = un bloc `mcp` sur le profil (donnée). Le skill built-in = un scope `Builtin` + un `.md` embarqué, injecté par le canal skills existant. +- **L** : `idea`, fichier et MCP sont **substituables** comme adapters entrants — `OrchestratorService` se comporte identiquement. Un profil sans `mcp` n'a aucun manque : la CLI `idea` (universelle) reste sa voie de délégation ; MCP n'est qu'un confort en plus. +- **I** : `OrchestratorService` ne reçoit que `AgentReplyChannel` en plus ; `idea` ne dépend que du protocole fichier ; l'adapter MCP ne dépend que du service + du port reply. +- **D** : tout est injecté au composition root (`state.rs`) ; aucun `new ConcreteAdapter` ailleurs. Le PATH/env de `idea` est posé via `SpawnSpec.env` au `LaunchAgent` (réutilise l'existant). + +### 16.7 Points ouverts (spikes v3) + +1. **Détection « la cible a répondu » sans outil** : la cible rend sa réponse via la commande **`idea reply`** (voie universelle). Risque résiduel : la cible **oublie** d'appeler `idea reply` ou de traiter son inbox. Mitigation : timeout borné + relance déposée dans l'inbox ; l'outbox reste la seule source fiable (on n'infère jamais la fin de tour). Lot **C-univ-2**. +2. **Boucles d'`ask` (A demande à B qui redemande à A)** : détection de cycle / profondeur max de délégation pour éviter l'interblocage (A bloqué sur B bloqué sur A). Garde-fou applicatif (chaîne de corrélation transportée dans la requête). Lot **C-univ-2**. +3. *(OPTIONNEL, MCP)* **Transport MCP par CLI** : stdio (process enfant) vs serveur local (socket/HTTP) ; chaque CLI déclare sa conf différemment. Spike : matérialiser `.mcp.json` Claude Code, conf Codex, conf Gemini depuis `McpConfigStrategy`. Lot **C-mcp-1**. **Ne conditionne plus le cross-model.** + +> **Tranché (n'est plus ouvert)** : +> - **Tâche à un agent vivant** : **inbox relue par la cible**, pas write stdin (décision 10). Un TUI plein écran n'a pas de prompt shell ; écrire dans le PTY corrompt le rendu et entrelace les frappes (bug « accents/ordre d'écriture »). +> - **Concurrence des corrélations** : résolue par l'**inbox FIFO par cible** — N `ask` simultanés s'empilent comme N fichiers ordonnés que la cible draine un par un via `idea next`. Plus besoin d'un mécanisme de sérialisation de writes PTY. + +### 16.8 Découpage en LOTS testables (cycle dev↔QA) — réordonné (universel d'abord, MCP optionnel ensuite) + +> **Principe d'ordonnancement révisé** : la **voie universelle** (plancher : domaine + rendez-vous + binaire `idea` + skill built-in + wiring) est livrée **en premier et en entier** — elle suffit à elle seule à garantir le cross-model et à fermer le chantier C sur le plan produit. La **voie MCP** est un **bloc optionnel postérieur**, livrable plus tard sans rien bloquer. Chaque lot = binôme dev+test, vert avant le suivant. +> +> **Cœur inchangé** : C0 (fondation domaine) et C1 (port `AgentReplyChannel` + outbox + rendez-vous applicatif) restent le socle commun aux deux voies — **non réécrits**, seulement étendus (C1 délivre désormais par **inbox**, pas par write PTY). + +**Bloc UNIVERSEL (principal — ferme le cross-model)** + +| Lot | Périmètre | Crates/dossiers | Contrats (ports/DTO) | Tests attendus | +|---|---|---|---|---| +| **C0 (fondation domaine)** | `CorrelationId` (VO validé), `OrchestratorCommand::AskAgent`, `AgentReply`, action `agent.ask` dans `validate`, **`SkillScope::Builtin`**, event `DomainEvent::AgentReplied`. *(`McpCapability`/`McpConfigStrategy` repoussés au bloc MCP — plus dans C0.)* | `domain/src/{orchestrator,skill,events}.rs` | variante + VO + scope + event, tous purs/sérialisés | unit purs : `agent.ask` valide (target+task requis) ; `AskAgent` rejette task vide ; `CorrelationId` non vide ; `SkillScope::Builtin` round-trip ; event sérialisé. **Réutilise** le harnais `orchestrator.rs`. | +| **C1 (port + rendez-vous, inbox)** | Port `AgentReplyChannel` (domaine) + adapter `OutboxReplyChannel` (infra : Notify in-memory + outbox `.ideai/outbox/`) ; branche `AskAgent` dans `OrchestratorService` : cible vivante ⇒ **dépôt inbox** `.ideai/inbox//.json` (PAS de write PTY), cible morte ⇒ launch + consigne initiale ; `await_reply` + timeout ; publish `AgentReplied`. | `domain/src/ports.rs`, `infrastructure/src/orchestrator/`, `application/src/orchestrator/service.rs` | `AgentReplyChannel`, `OutboxReplyChannel`, `OrchestratorService(+reply_channel)` | unit (fakes) : `await_reply` débloqué par `publish_reply` corrélé ; **timeout** → `ReplyError::Timeout`, cible **non tuée** ; cible vivante ⇒ **inbox écrite, aucun write PTY** ; cible morte ⇒ launch ; 2 `ask` ⇒ 2 fichiers inbox ordonnés (FIFO) ; corrélation croisée ignorée. infra : outbox écrit ⇒ `await_reply` résout. | +| **C-univ-1 (binaire `idea` + skill built-in + wiring PATH)** | Nouveau `crates/idea-cli/` (client mince : `ask`/`launch`/`reply`/`list-agents`/`next` ↔ protocole fichier) ; `.md` du skill built-in « Orchestration IdeA » (embarqué, scope `Builtin`) ; **retirer** le bloc prose de `compose_convention_file` et **auto-assigner** le skill built-in à tout agent ; poser `idea` sur le PATH du run dir via `SpawnSpec.env` au `LaunchAgent` (`PATH=/bin:…`, `IDEA_PROJECT_ROOT`, `IDEA_AGENT`). Démarrer le watcher fichier par projet (déjà existant) suffit côté backend. | `crates/idea-cli/`, `application/src/agent/lifecycle.rs`, `domain|application/src/skill/*`, `infrastructure/src/runtime/`, `app-tauri/src/state.rs` | binaire `idea` ↔ `.ideai/{requests,inbox,outbox}/` ; skill `Builtin` injecté ; env PATH/IDEA_* | `idea-cli` : `ask` écrit la requête + bloque jusqu'à l'outbox + imprime `content` ; `reply` écrit l'outbox corrélé ; `next` draine l'inbox FIFO ; timeout → code d'erreur. application : `compose_convention_file` n'a **plus** de prose libre mais le skill built-in est présent en tête des skills ; tout agent reçoit `Builtin`. lifecycle : `SpawnSpec.env` porte `PATH`/`IDEA_*`. e2e (Claude **et** Codex) : `idea ask` rend la réponse inline. | +| **C-univ-2 (garde-fous délégation)** | Sérialisation FIFO par cible (déjà naturelle via inbox — tests de non-entrelacement) + chaîne de corrélation transportée ⇒ **profondeur max / détection de cycle** (A→B→A) ; relance inbox sur timeout. | `application/src/orchestrator/service.rs`, `domain/src/orchestrator.rs` (chaîne corr.) | profondeur/anti-cycle dans la commande | unit : 2 `ask` concurrents ⇒ réponses non entrelacées (corrélation correcte) ; boucle A→B→A bloquée à la profondeur max → erreur typée ; timeout ⇒ relance inbox. | + +**Bloc MCP (OPTIONNEL — confort natif, postérieur, ne bloque rien)** + +| Lot | Périmètre | Crates/dossiers | Contrats (outils) | Tests attendus | +|---|---|---|---|---| +| **C-mcp-0 (capacité profil)** | `McpCapability`/`McpConfigStrategy` sur `AgentProfile` (défaut `None` ⇒ comportement universel inchangé). | `domain/src/profile.rs` | capacité déclarative pure | unit : `mcp=None` round-trip inchangé ; `mcp=Some(...)` sérialisé. | +| **C-mcp-1 (adapter MCP entrant + conf par CLI)** | Serveur MCP `infrastructure/src/orchestrator/mcp/` exposant `idea_ask_agent`/`idea_launch_agent`/`idea_list_agents`/`idea_reply` (+ parité `stop`/`update_context`/`create_skill`) → **même** `OrchestratorService`/outbox ; matérialiser la conf MCP (`McpConfigStrategy`) au `LaunchAgent` ssi `profile.mcp.is_some()` ; démarrer le serveur MCP par projet dans `state.rs` à côté du watcher ; mention des outils dans le skill built-in pour profils MCP. | `infrastructure/src/orchestrator/mcp/`, `application/src/agent/lifecycle.rs`, `app-tauri/src/state.rs` | outils MCP ↔ `OrchestratorCommand` ; `idea_reply` ↔ `publish_reply` | intégration : `idea_ask_agent` → dispatch → réponse inline corrélée (même outbox que la voie `idea`) ; `idea_launch_agent` fire-and-forget ; conf MCP matérialisée ssi profil MCP ; serveur MCP + watcher démarrés à l'open. | +| **C4 (route fichier `agent.ask` + UI observabilité)** | Le `FsOrchestratorWatcher` gère `agent.ask` brut (parité avec `idea`, pour un agent qui écrit le fichier à la main) ⇒ `*.response.json` porte `reply.content` ; relais event `agentReplied` → l'UI montre la relation requête↔réponse (extension du registre §14.3). | `infrastructure/src/orchestrator/mod.rs`, `app-tauri/src/{events,lib}.rs`, `frontend/src/features/agents` | `*.response.json` porte `reply.content` ; event front `agentReplied` | infra : fichier `agent.ask` → réponse contient le contenu corrélé. app-tauri : event relayé. Vitest : l'UI affiche « X a répondu à Y ». | +| **C4 (UI observabilité + route fichier `agent.ask`)** | Le `FsOrchestratorWatcher` gère `agent.ask` (écrit la réponse de contenu dans `*.response.json`) ⇒ parité repli ; relais event `agentReplied` → l'UI montre la relation requête↔réponse (extension du registre §14.3). | `infrastructure/src/orchestrator/mod.rs`, `app-tauri/src/{events,lib}.rs`, `frontend/src/features/agents` | `*.response.json` porte `reply.content` ; event front `agentReplied` | infra : fichier `agent.ask` → réponse contient le contenu corrélé. app-tauri : event relayé. Vitest : l'UI affiche « X a répondu à Y ». | + +**Ordre conseillé** : **C0 → C1 → C-univ-1 → C-univ-2** *(le chantier C est produit-complet et cross-model garanti ici)*, **puis, optionnellement et plus tard** : C-mcp-0 → C-mcp-1, et C4 (route fichier brute + observabilité UI, utile aux deux voies — peut être avancé après C-univ-1 si l'on veut l'UI tôt). C0 débloque tout ; C1 livre le rendez-vous (cœur, inbox) ; C-univ-1 livre la voie principale `idea`+skill ; C-univ-2 durcit la délégation ; le bloc MCP n'ajoute qu'un confort natif sur le même backend. + +### 16.9 Situer les chantiers adjacents (hors v3, cohérence) + +- **Hot-swap de profil** (§15.1, chantier A) : orthogonal — changer le profil d'un agent peut *changer* sa capacité MCP (`mcp`) ; `ChangeAgentProfile` re-matérialisera/retirera la conf MCP à la relance à chaud (réutilise C-mcp-1 au lieu de dupliquer). La voie `idea` (universelle) ne dépend pas du profil et reste disponible quel que soit le hot-swap. +- **Reprise des sessions** (§15.2, chantier B) : une session reprise ré-injecte son convention file (donc le skill built-in + le PATH `idea`) et, si profil MCP, sa conf MCP ; aucune interaction nouvelle (C-univ-1 et C-mcp-1 couvrent l'injection au `LaunchAgent`, que la reprise réutilise). +- **« Version B » (UI chat / agent headless à sortie structurée)** — **épic futur séparé, hors chantier C.** Une interface où l'utilisateur (ou un agent) dialogue avec un agent via une UI dédiée et reçoit une **sortie structurée** d'un mode headless. Elle se brancherait sur le **même** `OrchestratorService` + `AgentReplyChannel` + outbox (donc compatible, zéro réécriture du cœur). Notée ici pour cohérence ; **non requise** pour la garantie cross-model, qui est entièrement assurée par le bloc universel. + +--- + +## 17. Exécution structurée des agents IA — port `AgentSession` (PIVOT 2026-06-09, voie principale) + +> **Pivot verrouillé par le chef d'orchestre (acté, non rediscuté).** IdeA ne lit plus le terminal d'un agent IA et ne lui demande plus de se rapporter. Pour un agent **IA**, IdeA le **pilote via son mode programmatique/structuré** (ex. `claude -p --output-format stream-json` ou l'Agent SDK ; `codex exec` à sortie structurée) et **lit la réponse comme du JSON déterministe** (un message `result` final bien défini). La **plomberie devient 100 % fiable** ; seul reste irréductible le *contenu* de la réponse (propre à tout LLM). Cette section **remplace §16** comme voie principale et **réconcilie** avec §15 (chantiers A « hot-swap profil » et B « reprise session », tous deux LIVRÉS). +> +> Cette section **complète** §6 (use cases agent), §7 (layout), §9 (profils déclaratifs), §14.1 (run dir isolé), §14.3 (registre visible/arrière-plan) et §15 (agent = entité reprenable). Elle **ajoute un port domaine** (`AgentSession` + sa factory), **deux adapters infra** (Claude/Codex), **un type de cellule** (cellule IA vs terminal brut), **un registre de sessions structurées**, et le câblage frontend (UI chat). Elle **ne casse pas** les terminaux non-IA (PTY + xterm inchangés) ni A/B. + +### 17.0 État du terrain (lu dans le code, pas présumé) + +| Pièce | Existe ? | Référence code | +|---|---|---| +| Port `AgentRuntime` (prépare un `SpawnSpec`, pur) + `PtyPort` (PTY interactif) | ✅ | `domain/src/ports.rs` | +| Adapter unique `CliAgentRuntime` (profil déclaratif → `SpawnSpec`) | ✅ | `infrastructure/src/runtime/mod.rs` | +| `LaunchAgent` (résout profil+contexte, injecte, **spawn PTY**, registre) — invariant « 1 session vivante/agent » | ✅ — aujourd'hui **toujours** un PTY, même pour un agent IA | `application/src/agent/lifecycle.rs`, `terminal/registry.rs` | +| `TerminalSessions` (registre `SessionId → PtyHandle + TerminalSession`) | ✅ — `session_for_agent`, `rebind_agent_node`, `node_for_agent`, `remove` | `application/src/terminal/registry.rs` | +| `SessionKind::{Plain, Agent{agent_id}}` (ce qui tourne dans une cellule) | ✅ — pas de distinction « IA structurée » vs « TUI brut » | `domain/src/terminal.rs` | +| `LeafCell { session?, agent?, conversation_id?, agent_was_running }` + ops pures | ✅ — `conversation_id` = id opaque de conversation CLI, persistant | `domain/src/layout.rs` | +| `SessionStrategy { assign_flag?, resume_flag }` + `SessionPlan{None,Assign,Resume}` + `resolve_session_plan` | ✅ — pleinement câblé (A/B) | `domain/src/profile.rs`, `domain/src/ports.rs`, `agent/lifecycle.rs` | +| `ChangeAgentProfile` (chantier A, hot-swap) | ✅ — compose `LaunchAgent` ; jette `conversation_id` | `application/src/agent/lifecycle.rs` | +| `ListResumableAgents` (chantier B, inventaire reprise) | ✅ | `application/src/agent/resume.rs` | +| Bridge `PtyBridge` (PTY ↔ `tauri::ipc::Channel>` par session, generation-tracked) | ✅ — **transport incrémental réutilisable** | `app-tauri/src/pty.rs`, `commands.rs` | +| Catalogue de profils de référence (Claude, Codex, Gemini, Aider) + profil **custom** | ✅ — first-run wizard propose **tous** | `application/src/agent/catalogue.rs`, `frontend` first-run | +| **Port `AgentSession`** (session programmatique persistante par agent IA) | ❌ **totalement absent** | — | +| **Adapters `ClaudeSdkSession` / `CodexExecSession`** | ❌ totalement absent | — | +| **Notion « adapter structuré supporté » sur le profil** (registre Claude/Codex) | ❌ — n'importe quel `command` est accepté | — | +| **Distinction cellule « agent IA » (chat) vs « terminal brut »** côté layout + frontend | ❌ — toute cellule d'agent = TUI dans xterm | — | + +**Conclusion** : §17 = **(1)** ajouter le port domaine `AgentSession` + sa factory `AgentSessionFactory` (sélection par profil) ; **(2)** deux adapters infra structurés (Claude SDK / Codex exec) qui maintiennent la session et **parsent leur JSON documenté** (zéro scraping de TUI) ; **(3)** router `LaunchAgent` (et donc A/B + l'orchestrateur) vers `AgentSession` **pour les agents IA**, en gardant le PTY pour les terminaux bruts ; **(4)** modéliser **deux types de cellules** (IA/chat vs terminal/xterm) ; **(5)** UI chat alimentée par le **même mécanisme `Channel`** que les PTY ; **(6)** restreindre le menu de profils aux modèles **ayant un adapter** (Claude/Codex) et **retirer le custom** pour l'instant. + +--- + +### 17.1 Le port `AgentSession` (frontière domaine) — signatures tranchées + +**Décision : `AgentSession` est un port domaine (trait `async`), une instance vivante = une conversation persistante avec UN agent IA.** Il incarne l'invariant « 1 session vivante par agent » au niveau *type* (un agent IA possède au plus une `AgentSession` vivante, dans le registre §17.5). Il ne fuit **aucun** détail Claude/Codex (pas de `stream-json`, pas de `--output-format`, pas de chemin de transcript) : seulement « envoie un prompt, reçois un flux d'événements incrémentaux puis un contenu final déterministe ». + +#### Décision streaming **vs** bloquant : **les deux, via un flux + un terminal déterministe** +`send()` retourne un **flux d'événements de réponse** (`ReplyStream`), exactement comme `PtyPort::subscribe_output` retourne un `OutputStream` — mais **typé** (deltas de texte, événements d'outil, puis **un** événement terminal `Final`), pas des octets bruts. Le **rendez-vous synchrone** dont l'orchestrateur a besoin (§17.4) s'obtient en **drainant le flux jusqu'au `Final`** : c'est un helper applicatif `send_blocking()` au-dessus du même flux (DRY — pas deux chemins). Justification : +- **L'UI chat** veut le **rendu incrémental** (deltas live) ⇒ le flux. +- **L'orchestrateur** (Main demande à Architect) veut **la réponse complète** ⇒ draine jusqu'à `Final` ⇒ déterministe, sans deviner la fin de tour (le `Final` est **émis par l'adapter** quand il a lu le message `result` documenté de la CLI). +- **Une seule primitive** (`send → stream`) sert les deux besoins : zéro duplication, frontière minimale. + +```rust +// domain/src/ports.rs — nouveau port (frontière domaine, infra l'implémente) + +use std::time::Duration; + +/// Un événement incrémental d'un tour de réponse d'un agent IA. Universel : +/// l'adapter (Claude/Codex) traduit SON format structuré documenté vers ces +/// variantes ; aucun détail propre à une CLI ne franchit cette frontière. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum ReplyEvent { + /// Un fragment de texte assistant (rendu incrémental côté UI chat). + TextDelta { text: String }, + /// Une activité d'outil de l'agent (best-effort, pour l'observabilité chat : + /// « lit un fichier », « lance une commande »). Le `label` est déjà + /// humain-lisible ; le détail brut reste dans l'adapter. + ToolActivity { label: String }, + /// **Événement terminal déterministe** d'un tour : l'adapter l'émet quand il a + /// lu le message `result` documenté de la CLI. Porte le contenu final agrégé. + /// Après `Final`, le flux se termine (plus aucun événement). + Final { content: String }, +} + +/// Flux borné d'événements de réponse d'UN tour. Se termine après le `Final` +/// (ou sur erreur). Calqué sur `OutputStream`, mais typé. +pub type ReplyStream = Box + Send>; + +/// Erreurs d'une `AgentSession`. +#[derive(Debug, Clone, PartialEq, Eq, Error)] +pub enum AgentSessionError { + /// La session programmatique n'a pas pu démarrer (CLI introuvable, mode + /// structuré indisponible, handshake invalide). + #[error("agent session start failed: {0}")] + Start(String), + /// Échec d'envoi/de communication avec la session vivante. + #[error("agent session io failed: {0}")] + Io(String), + /// La sortie structurée de la CLI n'a pas pu être décodée (JSON cassé, + /// schéma inattendu). Frontière nette : on ne propage jamais le JSON brut. + #[error("agent session decode failed: {0}")] + Decode(String), + /// `send_blocking` n'a pas observé de `Final` dans le temps imparti. La session + /// **reste vivante** (on ne tue rien) ; l'appelant décide (cf. §16 décision 6). + #[error("agent session reply timed out")] + Timeout, +} + +/// Une **session programmatique persistante** avec un agent IA : une conversation +/// vivante que l'on pilote en mode structuré et dont on lit la réponse de façon +/// déterministe. Une instance ⇔ un agent IA (invariant « 1 session vivante/agent »). +/// +/// Hexagonal : ce trait est **domaine** ; les adapters Claude/Codex (infra) ne +/// fuient aucun détail de CLI à travers lui. Substituable (Liskov) : Claude et +/// Codex offrent les mêmes garanties (flux d'événements → `Final` déterministe), +/// seul le moteur diffère. +#[async_trait] +pub trait AgentSession: Send + Sync { + /// L'id de session IdeA (mappe la cellule/agent, comme un `PtyHandle.session_id`). + fn id(&self) -> SessionId; + + /// L'id de conversation **du moteur** (opaque), persisté sur la cellule pour la + /// reprise (§15.2/B). `None` tant que le moteur n'en a pas attribué. Permet à + /// `LeafCell.conversation_id` de rester le pivot de reprise, model-agnostic. + fn conversation_id(&self) -> Option; + + /// Transmet `prompt` à la session vivante et retourne le **flux** d'événements + /// du tour (deltas → `Final`). Rendu incrémental (UI chat) ET base du rendez-vous + /// synchrone (cf. `send_blocking`, helper applicatif). + /// + /// # Errors + /// [`AgentSessionError::Io`]/[`Decode`] sur échec de communication/décodage. + async fn send(&self, prompt: &str) -> Result; + + /// Termine proprement la session (tue le process/SDK sous-jacent). Idempotent. + /// + /// # Errors + /// [`AgentSessionError::Io`] si l'arrêt échoue. + async fn shutdown(&self) -> Result<(), AgentSessionError>; +} + +/// **Factory** sélectionnée par le profil : crée/reprend une `AgentSession` pour un +/// agent IA. C'est elle qui sait *quel adapter* instancier (Claude/Codex) selon +/// `profile.structured_adapter` (§17.3). Open/Closed : ajouter un moteur structuré +/// = ajouter un adapter + une variante de registre, sans toucher au cœur. +#[async_trait] +pub trait AgentSessionFactory: Send + Sync { + /// Vrai si cette factory sait piloter `profile` en mode structuré (sert au + /// menu de sélection §17.6 : ne proposer que les profils supportés). + fn supports(&self, profile: &AgentProfile) -> bool; + + /// Démarre une session structurée pour `profile` dans `cwd` (run dir isolé + /// §14.1), avec le contexte déjà préparé (`PreparedContext`) et l'intention de + /// session (`SessionPlan` : neuf / assign / resume — réutilise §15). + /// + /// # Errors + /// [`AgentSessionError::Start`] si la CLI/SDK est indisponible ou le mode + /// structuré ne peut s'initialiser. + async fn start( + &self, + profile: &AgentProfile, + ctx: &PreparedContext, + cwd: &ProjectPath, + session: &SessionPlan, + ) -> Result, AgentSessionError>; +} +``` + +> **Helper applicatif (pas un nouveau port)** — le rendez-vous synchrone de l'orchestrateur : +> ```rust +> // application/src/agent/structured.rs (ou un util) +> /// Draine un ReplyStream jusqu'au `Final` (borné par timeout), agrège les +> /// TextDelta de secours si l'adapter n'a pas pré-agrégé. Renvoie le contenu final. +> pub async fn send_blocking( +> session: &dyn AgentSession, prompt: &str, timeout: Duration, +> ) -> Result; +> ``` +> Ainsi `OrchestratorService` (§17.4) obtient un **rendez-vous synchrone intrinsèque** sans outbox, sans corrélation fichier, sans deviner la fin de tour : le `Final` *est* la fin de tour, **émise par l'adapter** depuis le message `result` documenté de la CLI. + +--- + +### 17.2 Les deux adapters structurés (infra) — Claude & Codex + +**Emplacement** : `infrastructure/src/session/` (nouveau module), pair de `infrastructure/src/runtime/` et `infrastructure/src/pty/`. Chaque adapter implémente `AgentSession` ; une `AgentSessionFactory` infra (`StructuredSessionFactory`) route un profil vers le bon adapter selon `profile.structured_adapter` (§17.3) et **agrège** les deux (registre interne `{ ClaudeSdkSession, CodexExecSession }`). Aucun de ces types ne franchit la frontière domaine : seuls `Arc` et `ReplyEvent` sortent. + +#### `ClaudeSdkSession` (Claude) +- **Mode** : `claude` en mode **non-interactif structuré** — `claude -p --output-format stream-json --input-format stream-json` (flux JSONL bidirectionnel), ou l'**Agent SDK** Claude si retenu au spike S1. La session **persiste** : le process reste vivant entre les `send()` (on écrit le prompt sur stdin au format documenté, on lit les lignes JSON sur stdout). +- **Persistance / reprise** : capte l'`session_id`/`conversation` exposé par le premier message structuré ⇒ `conversation_id()`. Réouverture = relancer avec le flag de reprise (réutilise la sémantique `SessionStrategy.resume_flag` / `SessionPlan::Resume` de §15 ; le profil Claude porte déjà un bloc `session`). +- **Détection du `Final`** : la CLI émet un message `{"type":"result", ...}` documenté en fin de tour ⇒ l'adapter émet `ReplyEvent::Final { content }`. Les messages `assistant`/`content_block_delta` ⇒ `TextDelta` ; les `tool_use` ⇒ `ToolActivity`. **Le format exact du `stream-json` est un spike (S1)** mais n'invalide pas l'ossature : le contrat de sortie reste `ReplyEvent`. + +#### `CodexExecSession` (Codex) +- **Mode** : `codex exec` (mode non-interactif/automation) avec sa **sortie structurée** documentée (JSON par tour). Selon que Codex maintient ou non un process long, deux incarnations possibles derrière le **même** trait : (a) process persistant piloté en flux, ou (b) **un `exec` par `send()`** réattaché via l'id de conversation Codex (`resume`). **L'incarnation est un détail d'adapter** ; le port ne change pas. +- **Persistance / reprise** : id de conversation Codex ⇒ `conversation_id()` ; reprise via le flag Codex (`SessionStrategy`/`SessionPlan::Resume`, §15). +- **Détection du `Final`** : message terminal documenté de `codex exec` ⇒ `ReplyEvent::Final`. **Format exact = spike (S2).** + +> **Liskov garanti** : un test de **conformité de port** partagé (un harnais `agent_session_contract`) vérifie pour chaque adapter : `send()` émet ≥0 deltas puis **exactement un** `Final` ; après `Final` le flux est clos ; `shutdown()` est idempotent ; `conversation_id()` devient `Some` après le premier tour assignant. Les adapters réels sont testés derrière un **fake CLI** scriptable (un binaire de test qui imprime des lignes JSON canned) pour rester déterministes et hors-réseau en CI. + +--- + +### 17.3 Évolution du modèle `AgentProfile` (adapter structuré supporté, retrait du custom) + +**Décision : un profil IA déclare quel adapter structuré le pilote.** On ajoute un champ **optionnel** `structured_adapter` sur `AgentProfile` (Open/Closed, comme `session`/`mcp` ; `skip_serializing_if = None` ⇒ zéro régression de sérialisation). Un profil **sans** `structured_adapter` reste un profil **TUI/PTY** (terminal brut) — c'est le cas des profils non encore couverts (Gemini, Aider) et de tout profil legacy. + +```rust +// domain/src/profile.rs — nouvel enum + champ optionnel sur AgentProfile +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub enum StructuredAdapter { + /// Piloté par `ClaudeSdkSession` (mode `-p --output-format stream-json` / SDK). + Claude, + /// Piloté par `CodexExecSession` (`codex exec` structuré). + Codex, +} + +pub struct AgentProfile { + // … champs existants inchangés … + /// Adapter d'exécution **structurée** (§17). `None` ⇒ agent **TUI/PTY** (cellule + /// terminal brut, comportement historique). `Some(_)` ⇒ agent **IA structuré** + /// (cellule chat + port `AgentSession`). Open/Closed : ajouter un moteur = une + /// variante + un adapter. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub structured_adapter: Option, +} +``` + +**Catalogue (`application/src/agent/catalogue.rs`)** — les profils de référence **Claude** et **Codex** portent désormais `structured_adapter: Some(Claude|Codex)`. **Décision produit (verrouillée) : le menu de sélection ne propose QUE les profils ayant un adapter** (donc Claude + Codex) et **le profil custom est retiré pour l'instant**. Concrètement : +- Le **catalogue** ne liste plus Gemini/Aider/custom dans le wizard et la création d'agent (ils restent dans le code comme références mais ne sont **pas proposés** tant qu'ils n'ont pas d'adapter structuré). *Alternative pragmatique* : un prédicat `is_selectable(profile) = factory.supports(profile)` filtre la liste exposée à l'UI — **un seul point de vérité** (la factory), pas une liste en dur. +- **First-run wizard (§9)** : ne propose que Claude + Codex ; pas de saisie de commande custom. +- **UI création/édition d'agent (§17.6)** : sélecteur restreint aux profils sélectionnables ; le bouton « profil custom » est masqué (flag produit, réactivable plus tard sans changer les contrats). + +> **Réconciliation A (hot-swap, §15.1)** : `ChangeAgentProfile` continue de composer `LaunchAgent` ; comme `LaunchAgent` route maintenant selon `structured_adapter` (§17.4), un swap Claude→Codex **change d'adapter** `AgentSession` (la conversation repart à neuf — décision déjà verrouillée : on jette `conversation_id`, on garde `.md`/mémoire). Aucun nouveau use case. Un swap d'un profil structuré vers un profil PTY (ou l'inverse) **change aussi le type de cellule** (§17.4) — `ChangeAgentProfile` republie l'event de relance et la cellule se reconstruit dans le bon mode. + +--- + +### 17.4 Lancement / reprise / hot-swap via le port — évolution de `LaunchAgent` + +**Décision : `LaunchAgent` devient le point de routage `IA structuré` vs `terminal brut`.** Il garde **toute** sa logique amont (résolution agent+profil+contexte, run dir isolé §14.1, seed permissions, recall mémoire §14.5.4, composition du convention file, `resolve_session_plan` §15) — **inchangée** — puis **branche** selon le profil : + +``` +LaunchAgent::execute(input): + … (étapes 1→5 actuelles : résoudre agent/profil/contexte, run dir, seed, + prepare_invocation OU prepared context, apply_injection) … // INCHANGÉ + if profile.structured_adapter.is_some(): // ── AGENT IA STRUCTURÉ ── + session = agent_session_factory.start(profile, prepared, run_dir, session_plan) + structured_sessions.insert(agent_id, session) // registre §17.5 + publish(AgentLaunched { agent_id, session_id: session.id() }) + // PAS de pty.spawn ; la cellule est de type « chat » (§17.7) + return LaunchAgentOutput { session: TerminalSession{kind: Agent, …}, + assigned_conversation_id: session.conversation_id() } + else: // ── TERMINAL BRUT (PTY) ── + … pty.spawn + TerminalSessions.insert (chemin ACTUEL, INCHANGÉ) … +``` + +- **Invariant « 1 session vivante/agent »** : la garde existante (`session_for_agent` au début de `execute`) est **généralisée** pour interroger **les deux** registres (PTY + structuré) — un agent est vivant s'il a une session vivante dans l'un OU l'autre. Le `rebind`/idempotence se duplique trivialement côté structuré (même sémantique : une cellule est une **vue**). +- **Reprise (B, §15.2)** : `ListResumableAgents` est **inchangé** (lecture pure du layout + manifeste + `resume_supported` via le profil). La reprise effective appelle `LaunchAgent` ⇒ pour un profil structuré, la factory démarre la session avec `SessionPlan::Resume { conversation_id }` (l'adapter passe le flag de reprise du moteur). Le `conversation_id` reste **persisté sur la cellule** (`LeafCell.conversation_id`) : pivot model-agnostic déjà en place. +- **Hot-swap (A, §15.1)** : `ChangeAgentProfile` **inchangé dans sa structure** ; le « kill PTY » devient « **shutdown de la session** » polymorphe : on résout la session vivante (PTY *ou* structurée), on l'arrête, puis on **rappelle `LaunchAgent`** dans la même cellule avec le nouveau profil ⇒ le bon adapter/type de cellule est re-sélectionné. *Détail d'implémentation* : `relaunch_if_live` interroge les deux registres ; une petite abstraction `LiveSession::shutdown()` (énum interne PTY/structuré) évite de dupliquer la branche. + +#### Messagerie inter-agents via le port — `OrchestratorService` +**Décision : « Main demande à Architect » = `architectSession.send_blocking(task, timeout)` ; routage model-agnostic au-dessus du port.** Le rendez-vous synchrone est **intrinsèque** (§17.1) ⇒ **on supprime, pour la voie principale, l'outbox / `idea reply` / l'inbox** (§16 déprécié). `OrchestratorService` gagne le registre `StructuredSessions` (ou un petit port `AgentMessenger` qui l'enveloppe) et **une** branche `AskAgent` : + +``` +OrchestratorService::dispatch(AskAgent { target, task, … }): + 1. Résoudre l'agent cible (find_agent_id_by_name) — NotFound sinon. + 2. Session structurée vivante ? (structured_sessions.session_for_agent) + a. Oui → session.send_blocking(task, timeout) // rendez-vous direct + b. Non → LaunchAgent(target, structuré) puis send_blocking(task, timeout) + 3. publish(AgentReplied { from_agent, … }) // observabilité UI (inchangé esprit §16) + 4. return reply.content +``` +> **Réconciliation §16** : `AskAgent`/`AgentReply` (modèle pur) **restent utiles** (la commande exprime « j'attends une réponse »), mais leur *réalisation* n'est plus l'outbox : c'est `AgentSession::send_blocking`. Le port `AgentReplyChannel`/`OutboxReplyChannel`, l'inbox et l'outbox **disparaissent de la voie principale**. La cible **doit** être pilotable en mode structuré (profil Claude/Codex) — cohérent avec le retrait du custom et la restriction du menu (§17.3/§17.6). Un agent cible **PTY** (profil sans adapter) n'est **pas** adressable par `ask` synchrone (erreur typée `Start`/`NotFound` explicite) : c'est acceptable car le menu ne crée plus que des agents structurés. + +--- + +### 17.5 Registre des sessions structurées — où il vit + +**Décision : un registre applicatif `StructuredSessions`, jumeau de `TerminalSessions`, dans `application/src/terminal/` (ou `agent/`).** Il mappe `SessionId → Arc` et expose la **même** surface que `TerminalSessions` côté liveness/agent (`session_for_agent`, `node_for_agent`, `rebind_agent_node`, `remove`, `live_agents`), de sorte que les use cases (garde d'unicité, orchestrateur, snapshot) traitent les deux registres derrière un **trait commun de liveness**. + +- **Réutilisation maximale** : on **généralise le trait existant `LiveAgentRegistry`** (`is_agent_live`/`is_node_live`) pour qu'il couvre les deux registres. Un agrégateur `LiveSessions { pty: TerminalSessions, structured: StructuredSessions }` répond « cet agent est-il vivant ? » en interrogeant les deux. `LaunchAgent` et `OrchestratorService` dépendent de l'agrégateur (ISP : ils ne voient que la capacité « liveness + résolution »). +- **Pourquoi pas dans le domaine** : comme `TerminalSessions` (cf. ses docs), un registre d'instances **vivantes** (avec `Arc` = ressource process/SDK) est un **état d'exécution applicatif**, pas du modèle métier. Le domaine ne tient que des **ids** et des **snapshots** (`TerminalSession`), jamais une poignée de process. +- **Pourquoi pas dans l'adapter** : le registre arbitre l'invariant produit « 1 session/agent » (règle applicative) et sert plusieurs use cases ⇒ il vit au-dessus des adapters, injecté au composition root (`state.rs`), exactement comme `TerminalSessions`. + +--- + +### 17.6 Deux types de cellules — modèle de layout & frontend + +**Décision : la distinction « cellule IA (chat) » vs « cellule terminal brut » est DÉRIVÉE, pas un nouveau champ de layout.** Le modèle `LeafCell` (§7) reste **inchangé** (`session?`, `agent?`, `conversation_id?`, `agent_was_running`). Le **type de rendu** d'une cellule se déduit à l'attache : +- cellule **sans agent** ⇒ terminal brut (PTY + xterm), inchangé ; +- cellule **avec agent** ⇒ on lit le `structured_adapter` du profil de l'agent : `Some` ⇒ **cellule chat** ; `None` ⇒ **cellule terminal brut** (un agent TUI legacy). + +Justification : aucune migration de layout persisté, aucune duplication d'invariants, et le type suit **toujours** le profil courant de l'agent (donc un hot-swap A change le rendu automatiquement). Le backend expose le type dérivé dans le DTO de session/cellule (`cellKind: "chat" | "terminal"`), calculé depuis le profil — **un seul point de vérité**. + +#### Frontend (TypeScript/React) +- **Nouveau composant `AgentChatView`** (`features/agents/` ou un nouveau `features/chat/`), pair de `TerminalView` : rend la **conversation** (bulles user/assistant, deltas incrémentaux, badges d'activité d'outil), avec une zone de saisie qui appelle `agentSession.send`. +- **`LayoutGrid`** choisit le composant par `cellKind` : `terminal` ⇒ `TerminalView` (xterm, inchangé) ; `chat` ⇒ `AgentChatView`. Les terminaux **non-IA gardent strictement** le chemin actuel. +- **Transport incrémental = réutilisation du `Channel`** : le flux `ReplyEvent` est poussé au front par le **même mécanisme** que les octets PTY — un `tauri::ipc::Channel` par session, enregistré dans un **`ChatBridge`** jumeau du `PtyBridge` (generation-tracked, ré-attachable après navigation/layout, **exactement** le pattern §17.5/§terminal-lifecycle). `ReplyChunk` est le DTO sérialisé d'un `ReplyEvent` (`{kind:"textDelta"|"toolActivity"|"final", …}`). Une cellule chat ré-attachée **repaint** depuis un **scrollback de conversation** (les tours déjà rendus), miroir du scrollback PTY. +- **Reprise de vue (bug lifecycle, mémoire)** : changer de layout/onglet **ne tue pas** la session structurée (elle vit dans le registre backend, comme un PTY) ; la vue se ré-attache via un `reattach_agent_chat` (jumeau de `reattach`), repeint le scrollback de conversation, re-subscribe au `Channel`. **Même garantie** que les PTY. + +--- + +### 17.7 Commandes Tauri & DTO (app-tauri) + +Nouvelles commandes (jumelles des commandes PTY existantes ; réutilisent `resolve_project`, le pattern `Channel`, le bridge) : + +``` +| Commande Tauri | Request (camelCase) | Réponse / Channel | +|-------------------------|------------------------------------------------------|-----------------------------------| +| agent_send | { sessionId, prompt, onReply: Channel } | () + flux ReplyChunk sur le Channel| +| reattach_agent_chat | { sessionId, onReply: Channel } | ReattachChatDto { scrollback } | +| close_agent_session | { sessionId } | () (shutdown + unregister) | +``` +- `launch_agent` (existant) renvoie déjà `assignedConversationId` et la `TerminalSession` ; on ajoute `cellKind` au DTO de session (dérivé §17.6) pour que le front choisisse le composant. +- `ChatBridge` (`app-tauri/src/chat.rs`, jumeau de `pty.rs`) route les `ReplyEvent` du registre `StructuredSessions` vers le bon `Channel`, generation-tracked. La boucle de pompe (drainer le `ReplyStream` d'un `send` et `send_output` chaque event) **calque** la pompe PTY de `commands.rs`. +- Events : `AgentReplied` (observabilité) relayé comme aujourd'hui ; `AgentLaunched`/`AgentProfileChanged` inchangés. + +--- + +### 17.8 Conformité hexagonale & SOLID + +- **Règle de dépendance** : le port `AgentSession`/`AgentSessionFactory` + `ReplyEvent`/`AgentSessionError` sont **domaine** (purs, I/O via trait `async`). Les adapters Claude/Codex, le décodage `stream-json`/`codex exec`, le process/SDK, le `ChatBridge` et le `Channel` sont **exclusivement** hors-domaine. Le domaine ignore `claude`/`codex`, JSON, stdin/stdout, transcripts. +- **S** : `AgentSession` = piloter UNE conversation ; `AgentSessionFactory` = créer/reprendre selon profil ; `StructuredSessions` = registre de liveness ; `ChatBridge` = transport. Aucune fonction fourre-tout. +- **O** : ajouter un moteur structuré = un adapter + une variante `StructuredAdapter` (donnée) ; **zéro** modification du cœur. Le menu se filtre via `factory.supports` (un point de vérité). +- **L** : Claude et Codex sont **substituables** derrière `AgentSession` (contrat partagé `agent_session_contract`). Un profil sans `structured_adapter` reste un terminal brut substituable au comportement historique. +- **I** : `LaunchAgent`/`OrchestratorService` ne reçoivent que la capacité « liveness + résolution + factory », pas le détail des adapters. `OrchestratorService` ne gagne qu'un registre/messager, pas l'outbox. +- **D** : tout est injecté au composition root (`state.rs`) ; aucun `new ClaudeSdkSession` ailleurs. `LaunchAgent` dépend de `Arc` et de l'agrégateur de registres. + +--- + +### 17.9 Découpage en LOTS testables (cycle dev↔QA) — ordonné + +> **Principe d'ordonnancement** : fondation domaine d'abord (port + profil), puis adapters derrière un fake CLI (déterministes, hors-réseau), puis le routage `LaunchAgent`, puis le frontend chat, puis la messagerie inter-agents, enfin le durcissement A/B + retrait custom. **Backend et frontend séparés par lot.** Chaque lot = binôme dev+test, vert avant le suivant. Les **spikes S1 (format Claude stream-json) / S2 (format Codex exec)** sont isolés dans les adapters (lot D2) et n'invalident pas l'ossature (le contrat de sortie est `ReplyEvent`). + +| Lot | Côté | Périmètre | Crates/dossiers | Contrats (port/DTO/gateway) | Tests attendus | +|---|---|---|---|---|---| +| **D0 (fondation domaine)** | back | Port `AgentSession` + `AgentSessionFactory`, types `ReplyEvent`/`ReplyStream`/`AgentSessionError` ; champ `AgentProfile.structured_adapter: Option` (+ enum). Catalogue : Claude/Codex portent `Some(...)`. | `domain/src/{ports,profile}.rs`, `application/src/agent/catalogue.rs` | trait `AgentSession`, factory, enum, champ optionnel sérialisé | unit purs : `structured_adapter=None` round-trip **inchangé** (zéro régression sérialisation) ; `Some(Claude/Codex)` round-trip ; catalogue Claude/Codex annotés ; signatures compilent (stub). | +| **D1 (registre + agrégateur liveness)** | back | `StructuredSessions` (jumeau `TerminalSessions`) ; généraliser `LiveAgentRegistry` ; agrégateur `LiveSessions{pty,structured}` ; helper `send_blocking`. | `application/src/terminal/registry.rs` (ou `agent/`), `application/src/agent/structured.rs` | `StructuredSessions`, `LiveSessions`, `send_blocking` | unit (fakes) : `session_for_agent`/`rebind`/`remove` côté structuré ; agrégateur = vivant si PTY **ou** structuré ; `send_blocking` draine jusqu'au `Final` ; timeout → `Timeout`, session non tuée. | +| **D2 (adapters Claude/Codex + contrat)** | back | `ClaudeSdkSession`, `CodexExecSession`, `StructuredSessionFactory` (route par `structured_adapter`). **Spikes S1/S2** isolés ici. Harnais de **conformité de port** + **fake CLI** scriptable. | `infrastructure/src/session/` | impls `AgentSession`/`AgentSessionFactory` | contrat partagé : ≥0 deltas puis **un** `Final` ; flux clos après `Final` ; `shutdown` idempotent ; `conversation_id` `Some` après assign ; `factory.supports` vrai pour Claude/Codex, faux sinon. Décodage JSON cassé → `Decode` (jamais de panic, jamais de JSON brut propagé). **Hors-réseau** (fake CLI). | +| **D3 (routage `LaunchAgent` + use cases)** | back | `LaunchAgent` route structuré vs PTY ; garde d'unicité sur l'agrégateur ; `ChangeAgentProfile` (A) et reprise (B) passent par `AgentSession` (shutdown polymorphe, `SessionPlan::Resume`). | `application/src/agent/lifecycle.rs`, `…/resume.rs` | `LaunchAgentOutput(+cellKind dérivé)` | unit (fakes) : profil structuré ⇒ **pas** de `pty.spawn`, session via factory, registre peuplé ; profil PTY ⇒ chemin actuel inchangé ; A : swap Claude→Codex ⇒ shutdown session + relance nouvel adapter, `conversation_id` jeté ; B : `Resume` passe l'id moteur ; invariant 1 session/agent across registres. | +| **D4 (commandes + bridge chat)** | back | Commandes `agent_send`/`reattach_agent_chat`/`close_agent_session` ; `ChatBridge` (jumeau `PtyBridge`, generation-tracked) ; `cellKind` au DTO de session ; pompe `ReplyStream → Channel`. | `app-tauri/src/{chat,commands,dto,events,state,lib}.rs` | `ReplyChunk` DTO, `ReattachChatDto`, `cellKind` | app-tauri : `agent_send` pompe les events sur le `Channel` ; ré-attache → scrollback conversation repeint ; `close_agent_session` → `shutdown`+unregister ; generation supersede (pas de double pompe). | +| **D5 (frontend chat)** | front | `AgentChatView` (deltas live, activité d'outil, saisie) ; `LayoutGrid` choisit par `cellKind` ; `AgentGateway.sendPrompt/reattachChat/closeAgentSession` + adapter Tauri + mock ; scrollback conversation + ré-attache. | `frontend/src/features/{chat,agents,layout}`, `frontend/src/ports`, `frontend/src/adapters` | `AgentGateway.sendPrompt(sessionId, prompt, onReply)`, `reattachChat`, `closeAgentSession` | Vitest : cellule `chat` rend `AgentChatView`, `terminal` rend `TerminalView` ; deltas s'accumulent → `final` fige le tour ; ré-attache repeint sans re-spawn ; mock gateway streame des `ReplyChunk`. | +| **D6 (messagerie inter-agents)** | back | `OrchestratorService` route `AskAgent` via `send_blocking` (registre structuré) ; **retrait voie principale** de l'outbox/inbox/`AgentReplyChannel` ; `AgentReplied` (observabilité) conservé ; cible PTY ⇒ erreur typée. | `application/src/orchestrator/service.rs` | `OrchestratorService(+structured registry/messager)` | unit (fakes) : cible vivante ⇒ `send_blocking` ; cible morte ⇒ launch + send ; timeout → typé, cible vivante ; `AgentReplied` émis ; cible PTY non adressable ⇒ erreur explicite ; **plus aucun accès outbox**. | +| **D7 (retrait custom + menu restreint)** | back+front | `is_selectable = factory.supports` filtre la liste exposée (first-run wizard, création/édition agent) ; retrait du profil **custom** ; Gemini/Aider non proposés (pas d'adapter). | `application/src/agent/catalogue.rs`, `app-tauri` (commande de liste de profils sélectionnables), `frontend` first-run + `features/agents` | prédicat `is_selectable` / liste filtrée exposée | unit : seuls Claude/Codex sélectionnables ; custom absent. Vitest : wizard et sélecteur n'affichent que Claude/Codex, bouton custom masqué. | + +**Ordre conseillé** : **D0 → D1 → D2 → D3 → D4 → D5 → D6 → D7**. D0 débloque tout ; D1/D2 sont parallélisables après D0 (D1 sur fakes, D2 sur fake CLI) ; D3 branche le routage ; D4/D5 livrent le chat (back puis front) ; D6 bascule la messagerie inter-agents sur le port ; D7 ferme le produit (menu restreint). Les terminaux **non-IA** restent verts à chaque lot (chemin PTY jamais modifié). + +### 17.10 Spikes restants (n'invalident pas l'ossature) + +1. **S1 — format exact du `stream-json` Claude** (`-p --output-format stream-json --input-format stream-json` vs Agent SDK) : schéma des messages `assistant`/`result`/`tool_use`, flag de reprise, propagation de l'`session_id`. Isolé dans `ClaudeSdkSession` (lot D2) ; le contrat `ReplyEvent` ne bouge pas. **À valider en priorité (réf. doc API Claude / Agent SDK).** +2. **S2 — mode structuré exact de Codex** (`codex exec` : process persistant vs `exec` par tour + `resume`, format JSON de fin de tour) : isolé dans `CodexExecSession` (lot D2). +3. **Backpressure du flux chat** (gros tours, throttling/coalescing côté front) : réutilise la mitigation PTY existante (§13.5) ; pas un bloqueur d'ossature. + +--- + ## 13. Risques techniques & points ouverts (spikes) 1. **PTY cross-platform** : portable-pty + xterm.js OK sur les 3 OS, mais signaux/resize/exit codes diffèrent (Windows ConPTY). **Spike** L3. diff --git a/Cargo.lock b/Cargo.lock index 3e94304..5dc05a3 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -879,6 +879,7 @@ dependencies = [ "serde", "serde_json", "thiserror 2.0.18", + "tokio", "uuid", ] diff --git a/Cargo.toml b/Cargo.toml index 2b4bf4b..ad9005c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -18,7 +18,7 @@ serde = { version = "1", features = ["derive"] } serde_json = "1" thiserror = "2" async-trait = "0.1" -tokio = { version = "1", features = ["rt-multi-thread", "macros", "sync", "fs", "io-util"] } +tokio = { version = "1", features = ["rt-multi-thread", "macros", "sync", "fs", "io-util", "time"] } # Local git via libgit2. Network features (https/ssh → openssl) are off for L8: # only local operations (status/commit/branch/checkout/log) are in scope; remote # push/pull and static vendoring for the AppImage are deferred to L9/L11. diff --git a/crates/application/Cargo.toml b/crates/application/Cargo.toml index 94555af..da1a968 100644 --- a/crates/application/Cargo.toml +++ b/crates/application/Cargo.toml @@ -14,6 +14,9 @@ serde = { workspace = true } serde_json = { workspace = true } # `v5` derives stable reference-profile ids from a fixed namespace (catalogue). uuid = { workspace = true } +# `time` feature only : borne le rendez-vous synchrone `send_blocking` (§17.4). +# Déjà le runtime async du workspace ; pas une nouvelle dépendance externe. +tokio = { workspace = true } [dev-dependencies] tokio = { workspace = true } diff --git a/crates/application/src/agent/catalogue.rs b/crates/application/src/agent/catalogue.rs index 5734307..5126bb4 100644 --- a/crates/application/src/agent/catalogue.rs +++ b/crates/application/src/agent/catalogue.rs @@ -18,7 +18,7 @@ //! making the reference profiles addressable across runs without a registry. use domain::ids::ProfileId; -use domain::profile::{AgentProfile, ContextInjection}; +use domain::profile::{AgentProfile, ContextInjection, StructuredAdapter}; /// A fixed UUID namespace used to derive stable ids for reference profiles. /// (Random-looking but constant; only its stability matters.) @@ -56,7 +56,8 @@ pub fn reference_profiles() -> Vec { "{agentRunDir}", None, ) - .expect("claude reference profile is valid"), + .expect("claude reference profile is valid") + .with_structured_adapter(StructuredAdapter::Claude), AgentProfile::new( reference_id("codex"), "OpenAI Codex CLI", @@ -68,7 +69,8 @@ pub fn reference_profiles() -> Vec { "{agentRunDir}", None, ) - .expect("codex reference profile is valid"), + .expect("codex reference profile is valid") + .with_structured_adapter(StructuredAdapter::Codex), AgentProfile::new( reference_id("gemini"), "Gemini CLI", diff --git a/crates/application/src/agent/mod.rs b/crates/application/src/agent/mod.rs index 00cd81f..445e201 100644 --- a/crates/application/src/agent/mod.rs +++ b/crates/application/src/agent/mod.rs @@ -10,10 +10,13 @@ mod catalogue; mod inspect; mod lifecycle; mod resume; +mod structured; mod usecases; pub(crate) use lifecycle::unique_md_path; +pub use structured::send_blocking; + pub use catalogue::{reference_profile_id, reference_profiles}; pub use inspect::{InspectConversation, InspectConversationInput, InspectConversationOutput}; pub use resume::{ diff --git a/crates/application/src/agent/structured.rs b/crates/application/src/agent/structured.rs new file mode 100644 index 0000000..eeaa5ef --- /dev/null +++ b/crates/application/src/agent/structured.rs @@ -0,0 +1,73 @@ +//! Helper applicatif `send_blocking` — le rendez-vous synchrone inter-agents +//! au-dessus du port [`AgentSession`] (ARCHITECTURE §17.1 / §17.4). +//! +//! `AgentSession::send` retourne un **flux** d'événements (`ReplyStream`), à la +//! manière de `PtyPort::subscribe_output`, mais **typé** : deltas de texte → +//! activités d'outil → **un** événement terminal déterministe +//! [`ReplyEvent::Final`]. Le rendez-vous synchrone dont l'orchestrateur a besoin +//! (§17.4) s'obtient en **drainant ce flux jusqu'au `Final`** : c'est *la* primitive +//! de la messagerie inter-agents, **sans outbox, sans corrélation fichier** (le +//! `Final` *est* la fin de tour). +//! +//! DRY : un seul chemin de lecture (le flux). `send_blocking` n'est qu'un *consom- +//! mateur* du même flux que la cellule chat utilise pour le rendu incrémental. + +use std::time::Duration; + +use domain::ports::{AgentSession, AgentSessionError, ReplyEvent}; + +/// Envoie `prompt` à la session vivante puis **draine le flux de réponse jusqu'au +/// [`ReplyEvent::Final`]**, et retourne son contenu agrégé. +/// +/// C'est le rendez-vous synchrone (§17.4) : on attend que le tour soit +/// déterministiquement terminé (`Final`) avant de rendre la main. Les deltas de +/// texte et les activités d'outil traversés en chemin sont **ignorés** ici (ils +/// servent le rendu incrémental côté UI, pas l'appelant synchrone). +/// +/// `timeout`, lorsqu'il est fourni, borne l'attente : si aucun `Final` n'est observé +/// dans le délai, on retourne [`AgentSessionError::Timeout`] **sans tuer la +/// session** (elle reste vivante dans le registre ; l'appelant décide de la suite). +/// `None` ⇒ pas de borne temporelle (on attend la fin du tour). +/// +/// # Errors +/// - [`AgentSessionError::Io`]/[`AgentSessionError::Decode`] remontées par `send` +/// (échec de communication / décodage de la sortie structurée) ; +/// - [`AgentSessionError::Io`] si le flux se termine **sans** `Final` (tour +/// interrompu) ; +/// - [`AgentSessionError::Timeout`] si `timeout` expire avant le `Final`. +pub async fn send_blocking( + session: &dyn AgentSession, + prompt: &str, + timeout: Option, +) -> Result { + match timeout { + Some(dur) => match tokio::time::timeout(dur, drain_to_final(session, prompt)).await { + Ok(result) => result, + // La session **reste vivante** : on ne `shutdown` rien ici (§17.1). + Err(_elapsed) => Err(AgentSessionError::Timeout), + }, + None => drain_to_final(session, prompt).await, + } +} + +/// Ouvre le flux du tour (`send`) et le **draine jusqu'au `Final`**. +/// +/// Le flux ([`domain::ports::ReplyStream`]) est un itérateur synchrone et borné : +/// après le `Final` il ne produit plus rien. On le parcourt donc simplement +/// jusqu'à rencontrer le `Final` (et on retourne son contenu) ; si le flux +/// s'épuise avant, c'est un tour interrompu → erreur [`AgentSessionError::Io`]. +async fn drain_to_final( + session: &dyn AgentSession, + prompt: &str, +) -> Result { + let stream = session.send(prompt).await?; + for event in stream { + if let ReplyEvent::Final { content } = event { + return Ok(content); + } + // TextDelta / ToolActivity : ignorés par le rendez-vous synchrone. + } + Err(AgentSessionError::Io( + "le flux de réponse s'est terminé sans événement Final".to_string(), + )) +} diff --git a/crates/application/src/lib.rs b/crates/application/src/lib.rs index 6bb1f72..9e0a2a3 100644 --- a/crates/application/src/lib.rs +++ b/crates/application/src/lib.rs @@ -38,7 +38,7 @@ pub use agent::{ ListResumableAgentsOutput, ProfileAvailability, ReadAgentContext, ReadAgentContextInput, ReadAgentContextOutput, ReferenceProfiles, ReferenceProfilesOutput, ResumableAgent, SaveProfile, SaveProfileInput, SaveProfileOutput, UpdateAgentContext, UpdateAgentContextInput, - AGENT_MEMORY_RECALL_BUDGET, + send_blocking, AGENT_MEMORY_RECALL_BUDGET, }; pub use embedder::{ CheckEmbedderSuggestion, CheckEmbedderSuggestionInput, CheckEmbedderSuggestionOutput, @@ -94,8 +94,8 @@ pub use template::{ UpdateTemplateOutput, }; pub use terminal::{ - CloseTerminal, CloseTerminalInput, CloseTerminalOutput, LiveAgentRegistry, OpenTerminal, - OpenTerminalInput, OpenTerminalOutput, ResizeTerminal, ResizeTerminalInput, TerminalSessions, - WriteToTerminal, WriteToTerminalInput, + CloseTerminal, CloseTerminalInput, CloseTerminalOutput, LiveAgentRegistry, LiveSessions, + OpenTerminal, OpenTerminalInput, OpenTerminalOutput, ResizeTerminal, ResizeTerminalInput, + StructuredSessions, TerminalSessions, WriteToTerminal, WriteToTerminalInput, }; pub use window::{MoveTabToNewWindow, MoveTabToNewWindowInput, MoveTabToNewWindowOutput}; diff --git a/crates/application/src/terminal/mod.rs b/crates/application/src/terminal/mod.rs index b4536c3..9f250fc 100644 --- a/crates/application/src/terminal/mod.rs +++ b/crates/application/src/terminal/mod.rs @@ -28,7 +28,7 @@ mod registry; mod usecases; -pub use registry::{LiveAgentRegistry, TerminalSessions}; +pub use registry::{LiveAgentRegistry, LiveSessions, StructuredSessions, TerminalSessions}; pub use usecases::{ CloseTerminal, CloseTerminalInput, CloseTerminalOutput, OpenTerminal, OpenTerminalInput, OpenTerminalOutput, ResizeTerminal, ResizeTerminalInput, WriteToTerminal, WriteToTerminalInput, diff --git a/crates/application/src/terminal/registry.rs b/crates/application/src/terminal/registry.rs index e4d245e..79c03c9 100644 --- a/crates/application/src/terminal/registry.rs +++ b/crates/application/src/terminal/registry.rs @@ -7,9 +7,9 @@ //! application layer rather than the domain or the adapter. use std::collections::HashMap; -use std::sync::Mutex; +use std::sync::{Arc, Mutex}; -use domain::ports::PtyHandle; +use domain::ports::{AgentSession, PtyHandle}; use domain::{AgentId, NodeId, SessionId, SessionKind, TerminalSession}; /// A registered, live terminal: its PTY handle plus the domain snapshot. @@ -204,3 +204,254 @@ impl TerminalSessions { self.len() == 0 } } + +// --------------------------------------------------------------------------- +// StructuredSessions — le jumeau de TerminalSessions pour les sessions IA (§17.5) +// --------------------------------------------------------------------------- + +/// Une session structurée enregistrée : la session vivante plus les coordonnées +/// (agent + cellule hôte) que [`AgentSession`] ne porte pas lui-même. +/// +/// `AgentSession` n'expose que `id()`/`conversation_id()` ; comme le snapshot +/// [`TerminalSession`] côté PTY, on associe ici l'`agent_id` (clé de liveness) et +/// le `node_id` (cellule-vue, rebindable) pour offrir la **même** surface que +/// [`TerminalSessions`]. +struct StructuredEntry { + /// La session vivante (ressource process/SDK), derrière le port domaine. + session: Arc, + /// L'agent IA pilotant cette session (invariant « 1 session vivante/agent »). + agent_id: AgentId, + /// La cellule (feuille de layout) qui héberge actuellement la vue. + node_id: NodeId, +} + +/// Registre en mémoire des sessions IA structurées vivantes (ARCHITECTURE §17.5). +/// +/// **Jumeau de [`TerminalSessions`]** : même rôle (état d'exécution applicatif, pas +/// du modèle métier — cf. les docs de [`TerminalSessions`]), même surface côté +/// liveness/agent (`session_for_agent`, `node_for_agent`, `live_agents`, +/// `rebind_agent_node`, `insert`/`remove`/`session`). La seule différence : il +/// stocke des `Arc` (sessions programmatiques) au lieu de +/// [`PtyHandle`]/[`TerminalSession`]. +/// +/// Respecte l'invariant produit **« 1 session vivante par agent »** : la garde +/// d'unicité (généralisée sur les deux registres via [`LiveSessions`]) interroge +/// `session_for_agent` avant tout lancement. +#[derive(Default)] +pub struct StructuredSessions { + entries: Mutex>, +} + +impl LiveAgentRegistry for StructuredSessions { + fn is_agent_live(&self, agent_id: &AgentId) -> bool { + self.session_for_agent(agent_id).is_some() + } + + fn is_node_live(&self, node_id: &NodeId) -> bool { + self.entries + .lock() + .map(|m| m.values().any(|e| e.node_id == *node_id)) + .unwrap_or(false) + } +} + +impl StructuredSessions { + /// Crée un registre vide. + #[must_use] + pub fn new() -> Self { + Self { + entries: Mutex::new(HashMap::new()), + } + } + + /// Enregistre une session fraîchement démarrée pour `agent_id`, hébergée par + /// la cellule `node_id`. Clé par l'id de session ([`AgentSession::id`]). + pub fn insert(&self, session: Arc, agent_id: AgentId, node_id: NodeId) { + if let Ok(mut map) = self.entries.lock() { + let id = session.id(); + map.insert( + id, + StructuredEntry { + session, + agent_id, + node_id, + }, + ); + } + } + + /// Retourne la session enregistrée pour un id, si présente. + #[must_use] + pub fn session(&self, id: &SessionId) -> Option> { + self.entries + .lock() + .ok() + .and_then(|m| m.get(id).map(|e| Arc::clone(&e.session))) + } + + /// Retourne la session vivante hébergeant `agent_id`, si elle existe. + /// + /// Jumeau de [`TerminalSessions::session_for_agent`] : **non ambigu** par + /// l'invariant « 1 session vivante/agent » — le premier match est *le* match. + #[must_use] + pub fn session_for_agent(&self, agent_id: &AgentId) -> Option> { + self.entries.lock().ok().and_then(|m| { + m.values() + .find(|e| &e.agent_id == agent_id) + .map(|e| Arc::clone(&e.session)) + }) + } + + /// Retourne l'[`SessionId`] de la session vivante hébergeant `agent_id`, si any. + #[must_use] + pub fn session_id_for_agent(&self, agent_id: &AgentId) -> Option { + self.entries.lock().ok().and_then(|m| { + m.values() + .find(|e| &e.agent_id == agent_id) + .map(|e| e.session.id()) + }) + } + + /// Retourne le [`NodeId`] de la cellule vivante hébergeant `agent_id`, si any. + /// + /// Jumeau de [`TerminalSessions::node_for_agent`]. + #[must_use] + pub fn node_for_agent(&self, agent_id: &AgentId) -> Option { + self.entries.lock().ok().and_then(|m| { + m.values() + .find(|e| &e.agent_id == agent_id) + .map(|e| e.node_id) + }) + } + + /// Liste chaque agent IA vivant, sa cellule hôte et son id de session. + /// + /// Jumeau de [`TerminalSessions::live_agents`] : un tuple + /// `(AgentId, NodeId, SessionId)` par session structurée vivante. + #[must_use] + pub fn live_agents(&self) -> Vec<(AgentId, NodeId, SessionId)> { + self.entries + .lock() + .map(|m| { + m.values() + .map(|e| (e.agent_id, e.node_id, e.session.id())) + .collect() + }) + .unwrap_or_default() + } + + /// Rebinde la session vivante d'un agent vers une nouvelle cellule-vue sans + /// redémarrer la conversation (« la cellule est une vue », §17.6). + /// + /// Jumeau de [`TerminalSessions::rebind_agent_node`] : seul le `node_id` + /// change ; la session, son id et sa conversation restent intacts. Retourne la + /// session rebindée, ou `None` si l'agent n'a pas de session vivante. + #[must_use] + pub fn rebind_agent_node( + &self, + agent_id: &AgentId, + node_id: NodeId, + ) -> Option> { + self.entries.lock().ok().and_then(|mut m| { + let entry = m.values_mut().find(|e| &e.agent_id == agent_id)?; + entry.node_id = node_id; + Some(Arc::clone(&entry.session)) + }) + } + + /// Retire une session du registre, retournant la session si présente (pour que + /// l'appelant la `shutdown` hors du verrou). + pub fn remove(&self, id: &SessionId) -> Option> { + self.entries + .lock() + .ok() + .and_then(|mut m| m.remove(id).map(|e| e.session)) + } + + /// Retourne toutes les sessions vivantes (pour un arrêt global propre au + /// shutdown applicatif, jumeau de [`TerminalSessions::handles`]). + #[must_use] + pub fn sessions(&self) -> Vec> { + self.entries + .lock() + .map(|m| m.values().map(|e| Arc::clone(&e.session)).collect()) + .unwrap_or_default() + } + + /// Nombre de sessions structurées vivantes. + #[must_use] + pub fn len(&self) -> usize { + self.entries.lock().map(|m| m.len()).unwrap_or(0) + } + + /// Si le registre est vide. + #[must_use] + pub fn is_empty(&self) -> bool { + self.len() == 0 + } +} + +// --------------------------------------------------------------------------- +// LiveSessions — agrégateur des deux registres derrière LiveAgentRegistry (§17.5) +// --------------------------------------------------------------------------- + +/// Agrégateur de liveness sur **les deux** registres (PTY + structuré). +/// +/// Un agent terminal brut vit dans [`TerminalSessions`] ; un agent IA structuré vit +/// dans [`StructuredSessions`]. La garde d'unicité et l'orchestrateur (§17.4) +/// dépendent de cet agrégateur (ISP : ils ne voient que la capacité « liveness + +/// résolution »), et une requête de liveness/agent voit donc les deux registres : +/// un agent est vivant s'il a une session vivante dans **l'un OU l'autre**. +/// +/// Implémente [`LiveAgentRegistry`] : le trait existant n'est **pas modifié** (ses +/// implémenteurs et appelants actuels — `TerminalSessions`, le snapshot — restent +/// inchangés), il est simplement **réalisé par un troisième implémenteur** qui +/// agrège, ce qui *généralise* sa portée à l'ensemble PTY+structuré sans régression. +pub struct LiveSessions { + /// Registre des sessions terminal brut (PTY). + pub pty: Arc, + /// Registre des sessions IA structurées. + pub structured: Arc, +} + +impl LiveSessions { + /// Construit l'agrégateur à partir des deux registres partagés. + #[must_use] + pub fn new(pty: Arc, structured: Arc) -> Self { + Self { pty, structured } + } + + /// L'[`SessionId`] de la session vivante d'un agent, PTY **ou** structurée. + #[must_use] + pub fn session_id_for_agent(&self, agent_id: &AgentId) -> Option { + self.pty + .session_for_agent(agent_id) + .or_else(|| self.structured.session_id_for_agent(agent_id)) + } + + /// La cellule hôte de la session vivante d'un agent, PTY **ou** structurée. + #[must_use] + pub fn node_for_agent(&self, agent_id: &AgentId) -> Option { + self.pty + .node_for_agent(agent_id) + .or_else(|| self.structured.node_for_agent(agent_id)) + } + + /// Tous les agents vivants des deux registres (PTY puis structurés). + #[must_use] + pub fn live_agents(&self) -> Vec<(AgentId, NodeId, SessionId)> { + let mut all = self.pty.live_agents(); + all.extend(self.structured.live_agents()); + all + } +} + +impl LiveAgentRegistry for LiveSessions { + fn is_agent_live(&self, agent_id: &AgentId) -> bool { + self.pty.is_agent_live(agent_id) || self.structured.is_agent_live(agent_id) + } + + fn is_node_live(&self, node_id: &NodeId) -> bool { + self.pty.is_node_live(node_id) || self.structured.is_node_live(node_id) + } +} diff --git a/crates/application/tests/profile_usecases.rs b/crates/application/tests/profile_usecases.rs index edeb2a5..711cffb 100644 --- a/crates/application/tests/profile_usecases.rs +++ b/crates/application/tests/profile_usecases.rs @@ -324,3 +324,39 @@ fn catalogue_ids_are_stable_across_calls() { // And match the slug-derived id helper. assert_eq!(first[0].id, reference_profile_id("claude")); } + +// --------------------------------------------------------------------------- +// LOT D0 (§17.3) — structured_adapter sur les profils de référence +// --------------------------------------------------------------------------- + +#[test] +fn catalogue_claude_and_codex_carry_their_structured_adapter() { + use domain::profile::StructuredAdapter; + + let profiles = reference_profiles(); + let by_command: HashMap<&str, &AgentProfile> = + profiles.iter().map(|p| (p.command.as_str(), p)).collect(); + + // Claude / Codex sont pilotés en mode structuré (cellule chat + AgentSession). + assert_eq!( + by_command["claude"].structured_adapter, + Some(StructuredAdapter::Claude), + "Claude reference profile must declare the Claude structured adapter" + ); + assert_eq!( + by_command["codex"].structured_adapter, + Some(StructuredAdapter::Codex), + "Codex reference profile must declare the Codex structured adapter" + ); +} + +#[test] +fn catalogue_gemini_and_aider_stay_pty_without_adapter() { + // §17.3 : les profils non encore couverts restent TUI/PTY (pas d'adapter). + let profiles = reference_profiles(); + let by_command: HashMap<&str, &AgentProfile> = + profiles.iter().map(|p| (p.command.as_str(), p)).collect(); + + assert_eq!(by_command["gemini"].structured_adapter, None); + assert_eq!(by_command["aider"].structured_adapter, None); +} diff --git a/crates/application/tests/send_blocking_d1.rs b/crates/application/tests/send_blocking_d1.rs new file mode 100644 index 0000000..e46cd36 --- /dev/null +++ b/crates/application/tests/send_blocking_d1.rs @@ -0,0 +1,222 @@ +//! LOT D1 (ARCHITECTURE §17.1 / §17.4 / §17.9 ligne D1) — tests unitaires du +//! helper applicatif `send_blocking`, **100 % fakes**. +//! +//! `send_blocking` draine le flux d'un tour jusqu'au [`ReplyEvent::Final`] : +//! - cas nominal : `TextDelta*` puis `Final{content}` ⇒ renvoie `content` (les +//! deltas/activités traversés sont ignorés par le rendez-vous synchrone) ; +//! - flux **sans** `Final` (épuisé) ⇒ [`AgentSessionError::Io`] ; +//! - `send` renvoyant `Err(Decode/Io)` ⇒ propagé tel quel ; +//! - **timeout** ⇒ [`AgentSessionError::Timeout`] **sans tuer la session** +//! (aucun `shutdown` appelé) ; le fake retarde son `send` de façon +//! déterministe (au-delà du `timeout`) pour forcer l'expiration. +//! +//! Fake `AgentSession` local et **scriptable** (inspiré du fake D0 du domaine) : +//! il produit le `ReplyStream` qu'on lui a scripté, ou une erreur, ou un délai ; +//! il compte ses `shutdown` pour prouver la non-mise-à-mort sur timeout. + +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::time::Duration; + +use async_trait::async_trait; + +use application::send_blocking; +use domain::ports::{AgentSession, AgentSessionError, ReplyEvent, ReplyStream}; +use domain::SessionId; +use uuid::Uuid; + +fn sid(n: u128) -> SessionId { + SessionId::from_uuid(Uuid::from_u128(n)) +} + +/// Ce que le fake doit faire au prochain `send`. +enum Script { + /// Renvoyer ce flux d'événements (consommé tel quel). + Stream(Vec), + /// `send` échoue avec cette erreur. + Err(AgentSessionError), + /// `send` dort `delay` avant de renvoyer le flux (force le timeout). + Delayed { + delay: Duration, + events: Vec, + }, +} + +/// Fake scriptable d'`AgentSession`. Mono-usage (un `send` scripté). +struct ScriptedSession { + id: SessionId, + script: std::sync::Mutex>, + shutdowns: AtomicUsize, +} + +impl ScriptedSession { + fn new(script: Script) -> Self { + Self { + id: sid(1), + script: std::sync::Mutex::new(Some(script)), + shutdowns: AtomicUsize::new(0), + } + } + fn shutdown_count(&self) -> usize { + self.shutdowns.load(Ordering::SeqCst) + } +} + +#[async_trait] +impl AgentSession for ScriptedSession { + fn id(&self) -> SessionId { + self.id + } + fn conversation_id(&self) -> Option { + None + } + async fn send(&self, _prompt: &str) -> Result { + let script = self + .script + .lock() + .unwrap() + .take() + .expect("send scripté une seule fois"); + match script { + Script::Stream(events) => { + let stream: ReplyStream = Box::new(events.into_iter()); + Ok(stream) + } + Script::Err(e) => Err(e), + Script::Delayed { delay, events } => { + tokio::time::sleep(delay).await; + let stream: ReplyStream = Box::new(events.into_iter()); + Ok(stream) + } + } + } + async fn shutdown(&self) -> Result<(), AgentSessionError> { + self.shutdowns.fetch_add(1, Ordering::SeqCst); + Ok(()) + } +} + +fn delta(t: &str) -> ReplyEvent { + ReplyEvent::TextDelta { text: t.to_owned() } +} +fn tool(l: &str) -> ReplyEvent { + ReplyEvent::ToolActivity { label: l.to_owned() } +} +fn final_(c: &str) -> ReplyEvent { + ReplyEvent::Final { content: c.to_owned() } +} + +// --------------------------------------------------------------------------- +// Cas nominal +// --------------------------------------------------------------------------- + +#[tokio::test] +async fn returns_final_content_ignoring_deltas_and_tools() { + let session = ScriptedSession::new(Script::Stream(vec![ + delta("hel"), + tool("reads a file"), + delta("lo"), + final_("hello world"), + ])); + + let out = send_blocking(&session, "ping", None).await; + assert_eq!(out, Ok("hello world".to_owned())); + // Rendez-vous synchrone : on ne tue pas la session sur succès. + assert_eq!(session.shutdown_count(), 0); +} + +#[tokio::test] +async fn returns_final_content_with_no_intermediate_events() { + // Flux = juste le Final déterministe. + let session = ScriptedSession::new(Script::Stream(vec![final_("done")])); + let out = send_blocking(&session, "x", Some(Duration::from_secs(5))).await; + assert_eq!(out, Ok("done".to_owned())); +} + +#[tokio::test] +async fn stops_at_first_final_even_if_more_events_follow() { + // Robustesse : le drain s'arrête au PREMIER Final et renvoie son contenu. + let session = ScriptedSession::new(Script::Stream(vec![ + delta("a"), + final_("first"), + final_("second-should-be-ignored"), + ])); + let out = send_blocking(&session, "x", None).await; + assert_eq!(out, Ok("first".to_owned())); +} + +// --------------------------------------------------------------------------- +// Bord : flux sans Final +// --------------------------------------------------------------------------- + +#[tokio::test] +async fn stream_without_final_is_io_error() { + let session = ScriptedSession::new(Script::Stream(vec![delta("a"), tool("b")])); + let out = send_blocking(&session, "x", None).await; + assert!( + matches!(out, Err(AgentSessionError::Io(_))), + "flux épuisé sans Final ⇒ Io, obtenu {out:?}" + ); +} + +#[tokio::test] +async fn empty_stream_is_io_error() { + let session = ScriptedSession::new(Script::Stream(vec![])); + let out = send_blocking(&session, "x", None).await; + assert!( + matches!(out, Err(AgentSessionError::Io(_))), + "flux vide ⇒ Io, obtenu {out:?}" + ); +} + +// --------------------------------------------------------------------------- +// Bord : erreur de `send` propagée +// --------------------------------------------------------------------------- + +#[tokio::test] +async fn send_decode_error_is_propagated() { + let session = + ScriptedSession::new(Script::Err(AgentSessionError::Decode("bad json".to_owned()))); + let out = send_blocking(&session, "x", None).await; + assert_eq!(out, Err(AgentSessionError::Decode("bad json".to_owned()))); +} + +#[tokio::test] +async fn send_io_error_is_propagated() { + let session = + ScriptedSession::new(Script::Err(AgentSessionError::Io("broken pipe".to_owned()))); + let out = send_blocking(&session, "x", Some(Duration::from_secs(5))).await; + assert_eq!(out, Err(AgentSessionError::Io("broken pipe".to_owned()))); +} + +// --------------------------------------------------------------------------- +// Bord : timeout — Timeout renvoyé ET session NON tuée +// --------------------------------------------------------------------------- + +#[tokio::test] +async fn timeout_returns_timeout_error_and_does_not_kill_session() { + // `send` dort 1 s ; timeout fixé à 20 ms ⇒ l'attente expire avant le Final. + let session = ScriptedSession::new(Script::Delayed { + delay: Duration::from_secs(1), + events: vec![final_("too late")], + }); + + let out = send_blocking(&session, "x", Some(Duration::from_millis(20))).await; + assert_eq!(out, Err(AgentSessionError::Timeout)); + // §17.1 : on ne `shutdown` rien sur timeout — la session reste vivante. + assert_eq!( + session.shutdown_count(), + 0, + "timeout ne doit PAS tuer la session" + ); +} + +#[tokio::test] +async fn no_timeout_bound_waits_for_final() { + // `timeout = None` ⇒ pas de borne : on attend le Final même après un délai. + let session = ScriptedSession::new(Script::Delayed { + delay: Duration::from_millis(10), + events: vec![final_("eventually")], + }); + let out = send_blocking(&session, "x", None).await; + assert_eq!(out, Ok("eventually".to_owned())); +} diff --git a/crates/application/tests/structured_registry_d1.rs b/crates/application/tests/structured_registry_d1.rs new file mode 100644 index 0000000..68226d2 --- /dev/null +++ b/crates/application/tests/structured_registry_d1.rs @@ -0,0 +1,305 @@ +//! LOT D1 (ARCHITECTURE §17.5 / §17.9 ligne D1) — tests unitaires des registres +//! structurés et de l'agrégateur de liveness, **100 % fakes**. +//! +//! Couvre : +//! - [`StructuredSessions`] : `insert` + résolution `session_for_agent` / +//! `session_id_for_agent` / `node_for_agent` / `session` (`None` si inconnu) ; +//! invariant « 1 session vivante par agent » ; `rebind_agent_node` (change la +//! cellule, pas la session) ; `remove` (retire + renvoie la session) ; +//! `live_agents` (triplets `(AgentId, NodeId, SessionId)`) ; impl +//! [`LiveAgentRegistry`] (`is_agent_live` / `is_node_live`). +//! - [`LiveSessions`] (agrégateur) : OR des deux registres pour `is_agent_live` / +//! `is_node_live` ; `live_agents` agrège (PTY puis structuré) ; +//! `node_for_agent` / `session_id_for_agent` trouvent dans l'un ou l'autre. +//! +//! Style calqué sur `terminal_usecases.rs` (constructeurs validants du domaine, +//! pas de littéraux fragiles). Le fake `AgentSession` est local et minimal : +//! ces tests n'exercent QUE le registre, pas `send`. + +use std::sync::Arc; + +use async_trait::async_trait; + +use application::{LiveAgentRegistry, LiveSessions, StructuredSessions, TerminalSessions}; +use domain::ports::{ + AgentSession, AgentSessionError, PtyHandle, ReplyStream, +}; +use domain::{ + AgentId, NodeId, ProjectPath, PtySize, SessionId, SessionKind, TerminalSession, +}; +use uuid::Uuid; + +// --- petits constructeurs déterministes ------------------------------------ + +fn sid(n: u128) -> SessionId { + SessionId::from_uuid(Uuid::from_u128(n)) +} +fn aid(n: u128) -> AgentId { + AgentId::from_uuid(Uuid::from_u128(n)) +} +fn nid(n: u128) -> NodeId { + NodeId::from_uuid(Uuid::from_u128(n)) +} + +/// Fake minimal d'`AgentSession` : porte juste l'id (ce que le registre clé). +/// `send`/`shutdown` ne sont pas exercés ici. +struct FakeSession { + id: SessionId, +} + +#[async_trait] +impl AgentSession for FakeSession { + fn id(&self) -> SessionId { + self.id + } + fn conversation_id(&self) -> Option { + None + } + async fn send(&self, _prompt: &str) -> Result { + let stream: ReplyStream = Box::new(std::iter::empty()); + Ok(stream) + } + async fn shutdown(&self) -> Result<(), AgentSessionError> { + Ok(()) + } +} + +fn fake(id: SessionId) -> Arc { + Arc::new(FakeSession { id }) +} + +// =========================================================================== +// StructuredSessions +// =========================================================================== + +#[test] +fn structured_insert_resolve_and_remove() { + let reg = StructuredSessions::new(); + assert!(reg.is_empty()); + assert_eq!(reg.len(), 0); + + let s = sid(1); + let a = aid(10); + let n = nid(100); + reg.insert(fake(s), a, n); + + assert_eq!(reg.len(), 1); + assert!(!reg.is_empty()); + + // Résolution par id de session. + assert!(reg.session(&s).is_some()); + assert_eq!(reg.session(&s).unwrap().id(), s); + + // Résolution par agent. + assert_eq!(reg.session_id_for_agent(&a), Some(s)); + assert_eq!(reg.node_for_agent(&a), Some(n)); + assert!(reg.session_for_agent(&a).is_some()); + assert_eq!(reg.session_for_agent(&a).unwrap().id(), s); + + // Inconnu ⇒ None partout. + assert!(reg.session(&sid(999)).is_none()); + assert!(reg.session_for_agent(&aid(999)).is_none()); + assert!(reg.session_id_for_agent(&aid(999)).is_none()); + assert!(reg.node_for_agent(&aid(999)).is_none()); + + // remove retire ET renvoie la session. + let removed = reg.remove(&s).expect("remove returns the session"); + assert_eq!(removed.id(), s); + assert!(reg.is_empty()); + assert!(reg.remove(&s).is_none(), "second remove is a no-op"); + assert!(reg.session_for_agent(&a).is_none()); +} + +#[test] +fn structured_live_agents_lists_triples() { + let reg = StructuredSessions::new(); + reg.insert(fake(sid(1)), aid(10), nid(100)); + reg.insert(fake(sid(2)), aid(20), nid(200)); + + let mut live = reg.live_agents(); + live.sort_by_key(|(a, _, _)| a.as_uuid()); + + assert_eq!( + live, + vec![ + (aid(10), nid(100), sid(1)), + (aid(20), nid(200), sid(2)), + ] + ); +} + +#[test] +fn structured_one_live_session_per_agent_invariant() { + // L'agent n'a pas de session vivante avant insertion. + let reg = StructuredSessions::new(); + let a = aid(10); + assert!(!reg.is_agent_live(&a)); + assert!(reg.session_for_agent(&a).is_none()); + + reg.insert(fake(sid(1)), a, nid(100)); + assert!(reg.is_agent_live(&a)); + + // `session_for_agent` est non ambigu : il rend LA session de l'agent. + let resolved = reg.session_for_agent(&a).unwrap().id(); + assert_eq!(resolved, sid(1)); + // Et un seul triplet figure pour cet agent dans live_agents. + let count = reg.live_agents().iter().filter(|(x, _, _)| *x == a).count(); + assert_eq!(count, 1, "one live session per agent"); +} + +#[test] +fn structured_rebind_changes_cell_not_session() { + let reg = StructuredSessions::new(); + let a = aid(10); + reg.insert(fake(sid(1)), a, nid(100)); + + let new_cell = nid(200); + let rebound = reg + .rebind_agent_node(&a, new_cell) + .expect("rebind returns the session"); + + // Session inchangée (même id), cellule mise à jour. + assert_eq!(rebound.id(), sid(1)); + assert_eq!(reg.node_for_agent(&a), Some(new_cell)); + assert_eq!(reg.session_id_for_agent(&a), Some(sid(1))); + // Toujours une seule entrée (pas de duplication). + assert_eq!(reg.len(), 1); + + // Rebind d'un agent inconnu ⇒ None. + assert!(reg.rebind_agent_node(&aid(999), nid(300)).is_none()); +} + +#[test] +fn structured_live_agent_registry_impl() { + let reg = StructuredSessions::new(); + let a = aid(10); + let n = nid(100); + + assert!(!reg.is_agent_live(&a)); + assert!(!reg.is_node_live(&n)); + + reg.insert(fake(sid(1)), a, n); + + assert!(reg.is_agent_live(&a)); + assert!(reg.is_node_live(&n)); + assert!(!reg.is_agent_live(&aid(999))); + assert!(!reg.is_node_live(&nid(999))); + + // is_node_live suit le rebind (la cellule vivante change). + let _ = reg.rebind_agent_node(&a, nid(200)); + assert!(!reg.is_node_live(&n), "old cell no longer live"); + assert!(reg.is_node_live(&nid(200)), "new cell is live"); +} + +#[test] +fn structured_sessions_snapshot_for_global_shutdown() { + let reg = StructuredSessions::new(); + reg.insert(fake(sid(1)), aid(10), nid(100)); + reg.insert(fake(sid(2)), aid(20), nid(200)); + + let mut ids: Vec = reg.sessions().iter().map(|s| s.id()).collect(); + ids.sort_by_key(|s| s.as_uuid()); + assert_eq!(ids, vec![sid(1), sid(2)]); +} + +// =========================================================================== +// LiveSessions — agrégateur (PTY OR structuré) +// =========================================================================== + +/// Insère un agent PTY dans `TerminalSessions`. +fn insert_pty(pty: &TerminalSessions, s: SessionId, a: AgentId, n: NodeId) { + let session = TerminalSession::starting( + s, + n, + ProjectPath::new("/p").unwrap(), + SessionKind::Agent { agent_id: a }, + PtySize::new(24, 80).unwrap(), + ); + pty.insert(PtyHandle { session_id: s }, session); +} + +#[test] +fn aggregator_agent_live_via_structured_only() { + let pty = Arc::new(TerminalSessions::new()); + let structured = Arc::new(StructuredSessions::new()); + let agg = LiveSessions::new(Arc::clone(&pty), Arc::clone(&structured)); + + let a = aid(10); + structured.insert(fake(sid(1)), a, nid(100)); + + assert!(agg.is_agent_live(&a), "live in structured ⇒ true"); + assert!(agg.is_node_live(&nid(100))); + assert_eq!(agg.session_id_for_agent(&a), Some(sid(1))); + assert_eq!(agg.node_for_agent(&a), Some(nid(100))); +} + +#[test] +fn aggregator_agent_live_via_pty_only() { + let pty = Arc::new(TerminalSessions::new()); + let structured = Arc::new(StructuredSessions::new()); + let agg = LiveSessions::new(Arc::clone(&pty), Arc::clone(&structured)); + + let a = aid(20); + insert_pty(&pty, sid(2), a, nid(200)); + + assert!(agg.is_agent_live(&a), "live in PTY ⇒ true"); + assert!(agg.is_node_live(&nid(200))); + assert_eq!(agg.session_id_for_agent(&a), Some(sid(2))); + assert_eq!(agg.node_for_agent(&a), Some(nid(200))); +} + +#[test] +fn aggregator_agent_absent_from_both_is_not_live() { + let pty = Arc::new(TerminalSessions::new()); + let structured = Arc::new(StructuredSessions::new()); + let agg = LiveSessions::new(pty, structured); + + let a = aid(30); + assert!(!agg.is_agent_live(&a), "absent from both ⇒ false"); + assert!(!agg.is_node_live(&nid(300))); + assert!(agg.session_id_for_agent(&a).is_none()); + assert!(agg.node_for_agent(&a).is_none()); + assert!(agg.live_agents().is_empty()); +} + +#[test] +fn aggregator_live_agents_concatenates_pty_then_structured() { + let pty = Arc::new(TerminalSessions::new()); + let structured = Arc::new(StructuredSessions::new()); + let agg = LiveSessions::new(Arc::clone(&pty), Arc::clone(&structured)); + + insert_pty(&pty, sid(1), aid(10), nid(100)); + structured.insert(fake(sid(2)), aid(20), nid(200)); + + let all = agg.live_agents(); + assert_eq!(all.len(), 2, "both registries contribute"); + // PTY d'abord, structuré ensuite (ordre documenté de l'agrégateur). + assert_eq!(all[0], (aid(10), nid(100), sid(1))); + assert_eq!(all[1], (aid(20), nid(200), sid(2))); +} + +#[test] +fn aggregator_resolution_prefers_pty_then_falls_back_to_structured() { + // Deux agents distincts, un par registre : la résolution trouve chacun + // dans le bon registre (OR / fallback). + let pty = Arc::new(TerminalSessions::new()); + let structured = Arc::new(StructuredSessions::new()); + let agg = LiveSessions::new(Arc::clone(&pty), Arc::clone(&structured)); + + let pty_agent = aid(10); + let struct_agent = aid(20); + insert_pty(&pty, sid(1), pty_agent, nid(100)); + structured.insert(fake(sid(2)), struct_agent, nid(200)); + + // Agent PTY : résolu par le registre PTY. + assert_eq!(agg.session_id_for_agent(&pty_agent), Some(sid(1))); + assert_eq!(agg.node_for_agent(&pty_agent), Some(nid(100))); + // Agent structuré : fallback sur le registre structuré. + assert_eq!(agg.session_id_for_agent(&struct_agent), Some(sid(2))); + assert_eq!(agg.node_for_agent(&struct_agent), Some(nid(200))); + + // is_node_live : OR sur les deux. + assert!(agg.is_node_live(&nid(100))); + assert!(agg.is_node_live(&nid(200))); + assert!(!agg.is_node_live(&nid(999))); +} diff --git a/crates/domain/Cargo.toml b/crates/domain/Cargo.toml index b9c2df0..99997c9 100644 --- a/crates/domain/Cargo.toml +++ b/crates/domain/Cargo.toml @@ -14,3 +14,4 @@ async-trait = { workspace = true } [dev-dependencies] serde_json = { workspace = true } +tokio = { workspace = true } diff --git a/crates/domain/src/ports.rs b/crates/domain/src/ports.rs index d87511b..af49b1c 100644 --- a/crates/domain/src/ports.rs +++ b/crates/domain/src/ports.rs @@ -189,6 +189,41 @@ pub type OutputStream = Box> + Send>; /// A boxed stream of domain events, returned by [`EventBus::subscribe`]. pub type EventStream = Box + Send>; +/// Un événement incrémental d'un tour de réponse d'un agent IA (ARCHITECTURE §17.1). +/// +/// Universel : l'adapter (Claude/Codex) traduit SON format structuré documenté +/// vers ces variantes ; **aucun** détail propre à une CLI (pas de `stream-json`, +/// pas de `--output-format`, pas de chemin de transcript) ne franchit cette +/// frontière domaine. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum ReplyEvent { + /// Un fragment de texte assistant (rendu incrémental côté UI chat). + TextDelta { + /// Le fragment de texte. + text: String, + }, + /// Une activité d'outil de l'agent (best-effort, pour l'observabilité chat : + /// « lit un fichier », « lance une commande »). Le `label` est déjà + /// humain-lisible ; le détail brut reste dans l'adapter. + ToolActivity { + /// Libellé humain-lisible de l'activité. + label: String, + }, + /// **Événement terminal déterministe** d'un tour : l'adapter l'émet quand il a + /// lu le message `result` documenté de la CLI. Porte le contenu final agrégé. + /// Après `Final`, le flux se termine (plus aucun événement). + Final { + /// Le contenu final agrégé du tour. + content: String, + }, +} + +/// Flux borné d'événements de réponse d'UN tour (ARCHITECTURE §17.1). Se termine +/// après le [`ReplyEvent::Final`] (ou sur erreur). Calqué sur [`OutputStream`], +/// mais **typé** : deltas de texte → activités d'outil → un `Final` déterministe, +/// plutôt que des octets bruts. +pub type ReplyStream = Box + Send>; + // --------------------------------------------------------------------------- // Per-port error types // --------------------------------------------------------------------------- @@ -218,6 +253,29 @@ pub enum PtyError { NotFound, } +/// Errors from an [`AgentSession`] / [`AgentSessionFactory`] (ARCHITECTURE §17.1). +/// +/// Frontière nette : on ne propage **jamais** le JSON brut d'une CLI à travers +/// ces erreurs (cf. [`AgentSessionError::Decode`]). +#[derive(Debug, Clone, PartialEq, Eq, Error)] +pub enum AgentSessionError { + /// La session programmatique n'a pas pu démarrer (CLI introuvable, mode + /// structuré indisponible, handshake invalide). + #[error("agent session start failed: {0}")] + Start(String), + /// Échec d'envoi/de communication avec la session vivante. + #[error("agent session io failed: {0}")] + Io(String), + /// La sortie structurée de la CLI n'a pas pu être décodée (JSON cassé, schéma + /// inattendu). On ne propage jamais le JSON brut. + #[error("agent session decode failed: {0}")] + Decode(String), + /// `send_blocking` n'a pas observé de [`ReplyEvent::Final`] dans le temps + /// imparti. La session **reste vivante** (on ne tue rien) ; l'appelant décide. + #[error("agent session reply timed out")] + Timeout, +} + /// Errors from [`ProcessSpawner`]. #[derive(Debug, Clone, PartialEq, Eq, Error)] pub enum ProcessError { @@ -402,6 +460,71 @@ pub trait AgentRuntime: Send + Sync { ) -> Result; } +/// Une **session programmatique persistante** avec un agent IA (ARCHITECTURE §17.1) : +/// une conversation vivante que l'on pilote en mode structuré et dont on lit la +/// réponse de façon déterministe. Une instance ⇔ un agent IA (invariant « 1 session +/// vivante/agent », porté au niveau *type*). +/// +/// Hexagonal : ce trait est **domaine** ; les adapters Claude/Codex (infra) ne +/// fuient aucun détail de CLI à travers lui. Substituable (Liskov) : Claude et +/// Codex offrent les mêmes garanties (flux d'événements → [`ReplyEvent::Final`] +/// déterministe), seul le moteur diffère. +/// +/// Calqué sur [`PtyPort`] (consommé comme trait-objet `Arc`, +/// d'où `#[async_trait]` pour rester object-safe — cf. note d'en-tête du module). +#[async_trait] +pub trait AgentSession: Send + Sync { + /// L'id de session IdeA (mappe la cellule/agent, comme un [`PtyHandle::session_id`]). + fn id(&self) -> SessionId; + + /// L'id de conversation **du moteur** (opaque), persisté sur la cellule pour la + /// reprise (§15.2). `None` tant que le moteur n'en a pas attribué. Permet à + /// `LeafCell.conversation_id` de rester le pivot de reprise, model-agnostic. + fn conversation_id(&self) -> Option; + + /// Transmet `prompt` à la session vivante et retourne le **flux** d'événements + /// du tour (deltas → [`ReplyEvent::Final`]). Rendu incrémental (UI chat) ET + /// base du rendez-vous synchrone (helper applicatif `send_blocking`). + /// + /// # Errors + /// [`AgentSessionError::Io`]/[`AgentSessionError::Decode`] sur échec de + /// communication/décodage. + async fn send(&self, prompt: &str) -> Result; + + /// Termine proprement la session (tue le process/SDK sous-jacent). Idempotent. + /// + /// # Errors + /// [`AgentSessionError::Io`] si l'arrêt échoue. + async fn shutdown(&self) -> Result<(), AgentSessionError>; +} + +/// **Factory** sélectionnée par le profil (ARCHITECTURE §17.1) : crée/reprend une +/// [`AgentSession`] pour un agent IA. C'est elle qui sait *quel adapter* instancier +/// (Claude/Codex) selon `profile.structured_adapter` (§17.3). Open/Closed : ajouter +/// un moteur structuré = ajouter un adapter + une variante de registre, sans +/// toucher au cœur. +#[async_trait] +pub trait AgentSessionFactory: Send + Sync { + /// Vrai si cette factory sait piloter `profile` en mode structuré (sert au menu + /// de sélection §17.6 : ne proposer que les profils supportés). + fn supports(&self, profile: &AgentProfile) -> bool; + + /// Démarre une session structurée pour `profile` dans `cwd` (run dir isolé + /// §14.1), avec le contexte déjà préparé ([`PreparedContext`]) et l'intention de + /// session ([`SessionPlan`] : neuf / assign / resume — réutilise §15). + /// + /// # Errors + /// [`AgentSessionError::Start`] si la CLI/SDK est indisponible ou le mode + /// structuré ne peut s'initialiser. + async fn start( + &self, + profile: &AgentProfile, + ctx: &PreparedContext, + cwd: &ProjectPath, + session: &SessionPlan, + ) -> Result, AgentSessionError>; +} + /// Open and drive pseudo-terminals. #[async_trait] pub trait PtyPort: Send + Sync { diff --git a/crates/domain/src/profile.rs b/crates/domain/src/profile.rs index 7811236..64e751f 100644 --- a/crates/domain/src/profile.rs +++ b/crates/domain/src/profile.rs @@ -118,6 +118,22 @@ impl SessionStrategy { } } +/// Adapter d'**exécution structurée** qui pilote un profil IA (ARCHITECTURE §17). +/// +/// Déclaratif, Open/Closed (comme [`EmbedderStrategy`]) : un profil déclare quel +/// adapter le pilote en mode programmatique. Ajouter un moteur structuré = une +/// variante ici + un adapter infra (`infrastructure/src/session/`), sans toucher +/// au cœur. Un profil **sans** `structured_adapter` reste un profil **TUI/PTY** +/// (terminal brut, comportement historique). +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub enum StructuredAdapter { + /// Piloté par `ClaudeSdkSession` (mode `-p --output-format stream-json` / SDK). + Claude, + /// Piloté par `CodexExecSession` (`codex exec` structuré). + Codex, +} + /// Declarative runtime configuration for one AI CLI. /// /// Invariants: @@ -148,6 +164,15 @@ pub struct AgentProfile { /// keeps today's behaviour. #[serde(default, skip_serializing_if = "Option::is_none")] pub session: Option, + /// Adapter d'exécution **structurée** (ARCHITECTURE §17). `None` ⇒ agent + /// **TUI/PTY** (cellule terminal brut, comportement historique). `Some(_)` ⇒ + /// agent **IA structuré** (cellule chat + port [`crate::ports::AgentSession`]). + /// Open/Closed : ajouter un moteur = une variante + un adapter. + /// + /// `skip_serializing_if = Option::is_none` ⇒ **zéro régression** de + /// sérialisation : un profil sans adapter sérialise exactement comme avant. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub structured_adapter: Option, } /// Embedding strategy of an [`EmbedderProfile`] (LOT C, étage 2 vectoriel). @@ -281,6 +306,16 @@ impl AgentProfile { detect, cwd_template, session, + structured_adapter: None, }) } + + /// Builder : fixe l'[`StructuredAdapter`] d'exécution structurée (§17) et + /// renvoie le profil. Laisse [`AgentProfile::new`] stable (zéro régression + /// d'appel) : les profils TUI/PTY ne l'appellent simplement pas. + #[must_use] + pub fn with_structured_adapter(mut self, adapter: StructuredAdapter) -> Self { + self.structured_adapter = Some(adapter); + self + } } diff --git a/crates/domain/tests/structured_session_d0.rs b/crates/domain/tests/structured_session_d0.rs new file mode 100644 index 0000000..ab9a95f --- /dev/null +++ b/crates/domain/tests/structured_session_d0.rs @@ -0,0 +1,396 @@ +//! LOT D0 (fondation §17.1) — tests purs du domaine : +//! - `StructuredAdapter` : round-trip serde camelCase (`"claude"` / `"codex"`) ; +//! - `AgentProfile.structured_adapter` : zéro régression serde (omission via +//! `skip_serializing_if`, défaut `None` à la désérialisation), round-trip avec +//! adapter, et builder `with_structured_adapter` qui ne touche rien d'autre ; +//! - types du port `AgentSession` : `ReplyEvent` (3 variantes + égalité), +//! `AgentSessionError` (variantes + `Display`/`Error`), et un **fake in-module** +//! prouvant que `AgentSession`/`AgentSessionFactory` sont implémentables (D0 = pas +//! d'impl réelle, seulement la conformité de signature). +//! +//! Style calqué sur `serde_roundtrip.rs` / `agent_profile_a0.rs` (réutilise les +//! constructeurs validants du domaine, pas de littéraux fragiles). + +use std::sync::Arc; + +use async_trait::async_trait; + +use domain::ids::{ProfileId, SessionId}; +use domain::ports::{ + AgentSession, AgentSessionError, AgentSessionFactory, PreparedContext, ReplyEvent, ReplyStream, + SessionPlan, +}; +use domain::profile::{AgentProfile, ContextInjection, StructuredAdapter}; +use domain::project::ProjectPath; +use domain::MarkdownDoc; +use uuid::Uuid; + +fn profid(n: u128) -> ProfileId { + ProfileId::from_uuid(Uuid::from_u128(n)) +} + +fn roundtrip(value: &T) -> T +where + T: serde::Serialize + serde::de::DeserializeOwned + PartialEq + std::fmt::Debug, +{ + let json = serde_json::to_string(value).expect("serialize"); + serde_json::from_str(&json).expect("deserialize") +} + +// --------------------------------------------------------------------------- +// StructuredAdapter — round-trip serde camelCase ("claude" / "codex") +// --------------------------------------------------------------------------- + +#[test] +fn structured_adapter_serialises_camel_case() { + assert_eq!( + serde_json::to_string(&StructuredAdapter::Claude).unwrap(), + "\"claude\"" + ); + assert_eq!( + serde_json::to_string(&StructuredAdapter::Codex).unwrap(), + "\"codex\"" + ); +} + +#[test] +fn structured_adapter_deserialises_from_camel_case() { + let claude: StructuredAdapter = serde_json::from_str("\"claude\"").unwrap(); + let codex: StructuredAdapter = serde_json::from_str("\"codex\"").unwrap(); + assert_eq!(claude, StructuredAdapter::Claude); + assert_eq!(codex, StructuredAdapter::Codex); +} + +#[test] +fn structured_adapter_roundtrips_both_variants() { + for adapter in [StructuredAdapter::Claude, StructuredAdapter::Codex] { + assert_eq!(roundtrip(&adapter), adapter); + } +} + +#[test] +fn structured_adapter_rejects_unknown_variant() { + // Un JSON hors vocabulaire ne doit pas se faufiler (pas de défaut silencieux). + let parsed: Result = serde_json::from_str("\"gemini\""); + assert!(parsed.is_err(), "unknown adapter must fail to deserialise"); +} + +// --------------------------------------------------------------------------- +// AgentProfile.structured_adapter — zéro régression serde +// --------------------------------------------------------------------------- + +/// Profil minimal TUI/PTY (sans adapter), construit via le constructeur validant. +fn pty_profile() -> AgentProfile { + AgentProfile::new( + profid(1), + "Gemini CLI", + "gemini", + vec![], + ContextInjection::convention_file("GEMINI.md").unwrap(), + Some("gemini --version".to_owned()), + "{agentRunDir}", + None, + ) + .unwrap() +} + +#[test] +fn profile_without_adapter_omits_the_field() { + // skip_serializing_if = Option::is_none ⇒ la clé `structuredAdapter` est ABSENTE. + let p = pty_profile(); + assert_eq!(p.structured_adapter, None); + let json = serde_json::to_string(&p).unwrap(); + assert!( + !json.contains("structuredAdapter"), + "field must be omitted when None; json was {json}" + ); + // Et le round-trip reste fidèle. + assert_eq!(roundtrip(&p), p); +} + +#[test] +fn legacy_profile_without_adapter_deserialises_to_none() { + // Un `profiles.json` produit AVANT que le champ existe : pas de clé + // `structuredAdapter` ⇒ défaut `None` (zéro régression de désérialisation). + let json = r#"{ + "id": "00000000-0000-0000-0000-000000000001", + "name": "Gemini CLI", + "command": "gemini", + "args": [], + "contextInjection": { "strategy": "conventionFile", "target": "GEMINI.md" }, + "detect": null, + "cwdTemplate": "{agentRunDir}" + }"#; + let p: AgentProfile = serde_json::from_str(json).expect("legacy profile must deserialise"); + assert_eq!(p.structured_adapter, None); +} + +#[test] +fn profile_with_adapter_roundtrips_and_uses_camel_case() { + let p = pty_profile().with_structured_adapter(StructuredAdapter::Claude); + let json = serde_json::to_string(&p).unwrap(); + // Clé camelCase + valeur camelCase de la variante. + assert!( + json.contains("\"structuredAdapter\":\"claude\""), + "json was {json}" + ); + assert!(!json.contains("structured_adapter"), "json was {json}"); + assert_eq!(roundtrip(&p), p); + + // Codex aussi. + let c = pty_profile().with_structured_adapter(StructuredAdapter::Codex); + let json = serde_json::to_string(&c).unwrap(); + assert!( + json.contains("\"structuredAdapter\":\"codex\""), + "json was {json}" + ); + assert_eq!(roundtrip(&c), c); +} + +#[test] +fn with_structured_adapter_sets_only_that_field() { + let before = pty_profile(); + let after = before.clone().with_structured_adapter(StructuredAdapter::Codex); + + // Le seul champ muté : + assert_eq!(after.structured_adapter, Some(StructuredAdapter::Codex)); + assert_ne!(after.structured_adapter, before.structured_adapter); + + // Tous les autres champs strictement inchangés : + assert_eq!(after.id, before.id); + assert_eq!(after.name, before.name); + assert_eq!(after.command, before.command); + assert_eq!(after.args, before.args); + assert_eq!(after.context_injection, before.context_injection); + assert_eq!(after.detect, before.detect); + assert_eq!(after.cwd_template, before.cwd_template); + assert_eq!(after.session, before.session); + + // Preuve d'équivalence : ne diffère du `before` que par l'adapter posé. + let mut patched = before; + patched.structured_adapter = Some(StructuredAdapter::Codex); + assert_eq!(after, patched); +} + +#[test] +fn with_structured_adapter_is_last_write_wins() { + // Re-poser un adapter écrase le précédent (idempotent par valeur). + let p = pty_profile() + .with_structured_adapter(StructuredAdapter::Claude) + .with_structured_adapter(StructuredAdapter::Codex); + assert_eq!(p.structured_adapter, Some(StructuredAdapter::Codex)); +} + +#[test] +fn new_defaults_structured_adapter_to_none() { + // Le constructeur `new` reste stable : il ne pose jamais d'adapter. + assert_eq!(pty_profile().structured_adapter, None); +} + +// --------------------------------------------------------------------------- +// ReplyEvent — 3 variantes construites, égalité, clone +// --------------------------------------------------------------------------- + +#[test] +fn reply_event_three_variants_construct_and_carry_payload() { + let delta = ReplyEvent::TextDelta { + text: "hel".to_owned(), + }; + let tool = ReplyEvent::ToolActivity { + label: "reads a file".to_owned(), + }; + let final_ = ReplyEvent::Final { + content: "hello world".to_owned(), + }; + + match &delta { + ReplyEvent::TextDelta { text } => assert_eq!(text, "hel"), + other => panic!("expected TextDelta, got {other:?}"), + } + match &tool { + ReplyEvent::ToolActivity { label } => assert_eq!(label, "reads a file"), + other => panic!("expected ToolActivity, got {other:?}"), + } + match &final_ { + ReplyEvent::Final { content } => assert_eq!(content, "hello world"), + other => panic!("expected Final, got {other:?}"), + } +} + +#[test] +fn reply_event_equality_and_clone() { + let e = ReplyEvent::Final { + content: "done".to_owned(), + }; + assert_eq!(e.clone(), e); + + // Même variante, payload différent ⇒ inégal. + assert_ne!( + e, + ReplyEvent::Final { + content: "other".to_owned() + } + ); + // Variantes différentes ⇒ inégal. + assert_ne!( + ReplyEvent::TextDelta { text: "x".into() }, + ReplyEvent::ToolActivity { label: "x".into() } + ); +} + +// --------------------------------------------------------------------------- +// AgentSessionError — variantes + Display / std::error::Error +// --------------------------------------------------------------------------- + +#[test] +fn agent_session_error_display_messages() { + assert_eq!( + AgentSessionError::Start("cli missing".to_owned()).to_string(), + "agent session start failed: cli missing" + ); + assert_eq!( + AgentSessionError::Io("broken pipe".to_owned()).to_string(), + "agent session io failed: broken pipe" + ); + assert_eq!( + AgentSessionError::Decode("bad json".to_owned()).to_string(), + "agent session decode failed: bad json" + ); + assert_eq!( + AgentSessionError::Timeout.to_string(), + "agent session reply timed out" + ); +} + +#[test] +fn agent_session_error_is_std_error_and_equates() { + // Conformité au trait std::error::Error (thiserror). + fn assert_error(_e: &E) {} + let e = AgentSessionError::Decode("x".to_owned()); + assert_error(&e); + + // Clone + égalité par variante/payload. + assert_eq!(e.clone(), e); + assert_ne!( + AgentSessionError::Start("a".into()), + AgentSessionError::Start("b".into()) + ); + assert_ne!(AgentSessionError::Timeout, AgentSessionError::Io("t".into())); +} + +// --------------------------------------------------------------------------- +// AgentSession / AgentSessionFactory — fake in-module (conformité de signature) +// +// D0 ne fournit PAS d'impl réelle : ce fake prouve seulement que les traits sont +// implémentables (object-safe, signatures cohérentes). Aucune logique réelle. +// --------------------------------------------------------------------------- + +struct FakeSession { + id: SessionId, + conversation_id: Option, +} + +#[async_trait] +impl AgentSession for FakeSession { + fn id(&self) -> SessionId { + self.id + } + + fn conversation_id(&self) -> Option { + self.conversation_id.clone() + } + + async fn send(&self, prompt: &str) -> Result { + // Flux borné minimal : un delta puis le Final déterministe. + let stream: ReplyStream = Box::new( + vec![ + ReplyEvent::TextDelta { + text: prompt.to_owned(), + }, + ReplyEvent::Final { + content: prompt.to_owned(), + }, + ] + .into_iter(), + ); + Ok(stream) + } + + async fn shutdown(&self) -> Result<(), AgentSessionError> { + Ok(()) + } +} + +struct FakeFactory; + +#[async_trait] +impl AgentSessionFactory for FakeFactory { + fn supports(&self, profile: &AgentProfile) -> bool { + profile.structured_adapter.is_some() + } + + async fn start( + &self, + _profile: &AgentProfile, + _ctx: &PreparedContext, + _cwd: &ProjectPath, + _session: &SessionPlan, + ) -> Result, AgentSessionError> { + Ok(Arc::new(FakeSession { + id: SessionId::from_uuid(Uuid::from_u128(7)), + conversation_id: None, + })) + } +} + +#[tokio::test] +async fn fake_session_proves_trait_is_implementable() { + let sid = SessionId::from_uuid(Uuid::from_u128(42)); + let session = FakeSession { + id: sid, + conversation_id: Some("conv-1".to_owned()), + }; + + // Consommé comme trait-objet (object-safety du port via #[async_trait]). + let dyn_session: Arc = Arc::new(session); + assert_eq!(dyn_session.id(), sid); + assert_eq!(dyn_session.conversation_id(), Some("conv-1".to_owned())); + + // send -> flux d'événements borné se terminant par Final (contrat §17.1). + let events: Vec = dyn_session.send("ping").await.unwrap().collect(); + assert_eq!( + events, + vec![ + ReplyEvent::TextDelta { + text: "ping".to_owned() + }, + ReplyEvent::Final { + content: "ping".to_owned() + }, + ] + ); + assert!(matches!(events.last(), Some(ReplyEvent::Final { .. }))); + + dyn_session.shutdown().await.expect("shutdown ok"); +} + +#[tokio::test] +async fn fake_factory_supports_only_structured_profiles_and_starts() { + let factory: Arc = Arc::new(FakeFactory); + + let structured = pty_profile().with_structured_adapter(StructuredAdapter::Claude); + let pty = pty_profile(); + assert!(factory.supports(&structured)); + assert!(!factory.supports(&pty)); + + let ctx = PreparedContext { + content: MarkdownDoc::new("# ctx"), + relative_path: "CLAUDE.md".to_owned(), + }; + let cwd = ProjectPath::new("/srv/run").unwrap(); + let session = factory + .start(&structured, &ctx, &cwd, &SessionPlan::None) + .await + .expect("factory starts a session"); + assert_eq!(session.id(), SessionId::from_uuid(Uuid::from_u128(7))); +}