//! [`FsProviderSessionStore`] — l'adapter `tokio::fs` du port [`ProviderSessionStore`] //! (cadrage « persistance conversationnelle », lot P5). //! //! Le `resumable_id` **propre au moteur** (le `--resume`/`--continue` d'une CLI) est //! rangé **par (conversation, provider)** (ARCHITECTURE §19.2/§19.3) dans un seul fichier //! JSON par conversation : une **map `providerId → resumableId`** où plusieurs providers //! coexistent (après un swap Claude→Codex, on garde l'id de chacun) : //! //! ```text //! /.ideai/conversations/ //! └── / //! ├── log.jsonl # le log append-only (lot P2) //! ├── handoff.md # le dernier point de reprise (lot P3) //! └── providers.json # { "": "", ... } (ce module) //! ``` //! //! ## Format `providers.json` //! //! Un objet JSON plat `{ "claude": "abc-123", "codex": "def-456" }`. C'est l'unité //! écrite/relue : pas de wrapper, la clé est l'identifiant de profil/moteur. //! //! ## Robustesse //! //! - **Écriture atomique** : on écrit dans `providers.json.tmp` puis on `rename` — //! jamais de fichier à moitié écrit visible (même convention que [`super::FsHandoffStore`]). //! - **Lecture-modification-écriture** : `set` charge la map existante, insère/écrase la //! seule clé `provider_id`, puis réécrit — les autres providers ne sont jamais perdus. //! L'opération est **sérialisée par conversation** (un `Mutex` async par fichier, comme //! [`super::FsConversationLog`]) pour qu'un `set` concurrent ne perde pas l'écriture de //! l'autre (read-modify-write atomique du point de vue applicatif). //! - **Fichier ou clé absent** ⇒ `get` renvoie `Ok(None)` (jamais une erreur) : un //! provider sans session rangée est l'état normal au premier tour. //! - **Fichier présent mais illisible** (JSON corrompu) ⇒ [`StoreError::Serialization`]. //! Contrairement au log JSONL (où une ligne corrompue est sautée), il n'y a qu'un //! enregistrement : on ne peut pas « sauter », c'est une vraie erreur. use std::collections::HashMap; use std::path::PathBuf; use std::sync::{Arc, Mutex}; use async_trait::async_trait; use domain::conversation::ConversationId; use domain::conversation_log::ProviderSessionStore; use domain::ports::StoreError; use domain::project::ProjectPath; use super::{CONVERSATIONS_DIR, IDEAI_DIR}; /// Nom du fichier des sessions par provider, par conversation. const PROVIDERS_FILE: &str = "providers.json"; /// Nom du fichier temporaire d'écriture atomique (renommé sur `providers.json`). const PROVIDERS_TMP_FILE: &str = "providers.json.tmp"; /// La map persistée : `providerId → resumableId`. type ProviderMap = HashMap; /// Adapter `tokio::fs` du store des `resumable_id`, un `providers.json` par conversation. /// /// Même convention de construction que [`super::FsConversationLog`] / [`super::FsHandoffStore`] : /// le **project root** est fourni au constructeur, la base `/.ideai/conversations` /// en dérive, et chaque conversation a son sous-dossier `/`. pub struct FsProviderSessionStore { /// Racine projet canonique. `None` si le root fourni n'est pas résoluble : /// le store devient alors fail-closed (lecture `None`, écriture `Err`). root: Option, /// Racine `/.ideai/conversations`. base: Option, /// Verrous d'écriture, un par fichier de conversation (sérialise le read-modify-write). write_locks: Mutex>>>, } impl FsProviderSessionStore { /// Construit l'adapter à partir du **project root**. /// /// La base `/.ideai/conversations` en est dérivée ; le dossier de conversation /// est créé paresseusement au premier `set`. #[must_use] pub fn new(root: &ProjectPath) -> Self { let root = std::fs::canonicalize(root.as_str()).ok(); let base = root .as_ref() .map(|root| root.join(IDEAI_DIR).join(CONVERSATIONS_DIR)); Self { root, base, write_locks: Mutex::new(HashMap::new()), } } /// `/` — le dossier d'une conversation. fn conversation_dir(&self, conversation: ConversationId) -> Result { let Some(root) = self.root.as_ref() else { return Err(StoreError::Io( "project root canonique indisponible pour providers.json".to_owned(), )); }; let Some(base) = self.base.as_ref() else { return Err(StoreError::Io( "base providers.json indisponible".to_owned(), )); }; let dir = base.join(conversation.to_string()); if !dir.starts_with(root) { return Err(StoreError::Io( "chemin providers.json hors du project root canonique".to_owned(), )); } Ok(dir) } /// `//providers.json` — le fichier des sessions par provider. fn providers_path(&self, conversation: ConversationId) -> Result { Ok(self.conversation_dir(conversation)?.join(PROVIDERS_FILE)) } /// `//providers.json.tmp` — le fichier temporaire d'écriture. fn providers_tmp_path(&self, conversation: ConversationId) -> Result { Ok(self .conversation_dir(conversation)? .join(PROVIDERS_TMP_FILE)) } /// Renvoie (en le créant au besoin) le verrou d'écriture de `conversation`. fn write_lock(&self, conversation: ConversationId) -> Arc> { self.write_locks .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) .entry(conversation) .or_default() .clone() } /// Charge la map `providerId → resumableId` de `conversation`. /// /// Fichier absent ⇒ map vide (jamais une erreur). JSON illisible ⇒ /// [`StoreError::Serialization`]. async fn load_map(&self, conversation: ConversationId) -> Result { let path = match self.providers_path(conversation) { Ok(path) => path, Err(_) => return Ok(ProviderMap::new()), }; let bytes = match tokio::fs::read(path).await { Ok(bytes) => bytes, Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(ProviderMap::new()), Err(e) => return Err(StoreError::Io(e.to_string())), }; serde_json::from_slice(&bytes) .map_err(|e| StoreError::Serialization(format!("providers.json: {e}"))) } } #[async_trait] impl ProviderSessionStore for FsProviderSessionStore { async fn get( &self, conversation: ConversationId, provider_id: &str, ) -> Result, StoreError> { let map = self.load_map(conversation).await?; Ok(map.get(provider_id).cloned()) } async fn set( &self, conversation: ConversationId, provider_id: &str, resumable_id: &str, ) -> Result<(), StoreError> { // Read-modify-write sérialisé par conversation : le verrou est tenu le temps du // load + insert + write atomique, donc deux `set` concurrents (providers // différents ou non) s'ordonnent sans s'écraser. let lock = self.write_lock(conversation); let _guard = lock.lock().await; let mut map = self.load_map(conversation).await?; map.insert(provider_id.to_string(), resumable_id.to_string()); // Sérialiser **avant** toute I/O d'écriture : une erreur de sérialisation ne doit // pas laisser de fichier tmp partiel. let body = serde_json::to_vec(&map).map_err(|e| StoreError::Serialization(e.to_string()))?; let dir = self.conversation_dir(conversation)?; tokio::fs::create_dir_all(&dir) .await .map_err(|e| StoreError::Io(e.to_string()))?; // Écriture atomique : écrire le tmp puis `rename` sur la cible. Un lecteur ne voit // jamais de fichier à moitié écrit (le rename est atomique sur le FS). let tmp = self.providers_tmp_path(conversation)?; tokio::fs::write(&tmp, &body) .await .map_err(|e| StoreError::Io(e.to_string()))?; tokio::fs::rename(&tmp, self.providers_path(conversation)?) .await .map_err(|e| StoreError::Io(e.to_string()))?; Ok(()) } }