//! [`FsBackgroundTaskStore`] — durable store for first-class background tasks. //! //! Persistence is segmented by project id under the project root: //! //! ```text //! /.ideai/background-tasks/.json //! ``` //! //! Each file is a small JSON document `{ version, projectId, tasks }`. Mutations //! are serialized per project and writes are atomic: serialize to a unique tmp //! path, then rename over the target. The in-memory registry indexes open tasks //! by owner agent and task id; it is rebuilt lazily from disk on first use and //! updated on every successful mutation. //! //! Boot reconcile rule for B2: because the runner registry is not wired yet, a //! task in `Running` or `Waiting` whose id is not present in the caller-provided //! live handle list is marked `Failed` with a synthetic restart-loss error and //! `completion_delivered = false`. Existing terminal tasks that are not delivered //! are only reported as delivery-pending. `Queued` tasks are left unchanged. use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet}; use std::path::{Path, PathBuf}; use std::sync::Arc; use async_trait::async_trait; use serde::{Deserialize, Serialize}; use tokio::sync::{Mutex, RwLock}; use uuid::Uuid; use domain::{ AgentId, BackgroundTask, BackgroundTaskPortError, BackgroundTaskResult, BackgroundTaskState, BackgroundTaskStore, ProjectId, ProjectPath, TaskId, }; const IDEAI_DIR: &str = ".ideai"; const BACKGROUND_TASKS_DIR: &str = "background-tasks"; const TASK_DOC_VERSION: u32 = 1; /// File-backed implementation of [`BackgroundTaskStore`]. pub struct FsBackgroundTaskStore { dir: PathBuf, registry: RwLock, project_locks: Mutex>>>, } #[derive(Debug, Clone, Default, PartialEq, Eq)] struct RegistryIndex { loaded: bool, task_to_project: HashMap, open_by_agent: HashMap>, } #[derive(Debug, Clone, Default, PartialEq, Eq)] struct StoreState { docs: BTreeMap, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] struct TaskDoc { version: u32, project_id: ProjectId, tasks: Vec, } /// Summary returned by [`FsBackgroundTaskStore::reconcile_boot`]. #[derive(Debug, Clone, Default, PartialEq, Eq)] pub struct BackgroundTaskReconcileReport { /// Open tasks converted to `Failed` because no live runtime handle exists. pub failed_task_ids: Vec, /// Terminal tasks whose completion remains to be delivered. pub delivery_pending_task_ids: Vec, } impl FsBackgroundTaskStore { /// Builds the store for a project root. #[must_use] pub fn new(root: &ProjectPath) -> Self { let dir = PathBuf::from(root.as_str()) .join(IDEAI_DIR) .join(BACKGROUND_TASKS_DIR); Self { dir, registry: RwLock::new(RegistryIndex::default()), project_locks: Mutex::new(HashMap::new()), } } /// `/.ideai/background-tasks`. #[must_use] pub fn dir(&self) -> &Path { &self.dir } /// Returns the open task ids known by the in-memory registry for `agent_id`. /// /// This is intentionally a registry view: it does not read the filesystem. pub async fn registry_open_task_ids_for_agent(&self, agent_id: AgentId) -> Vec { self.registry .read() .await .open_by_agent .get(&agent_id) .map(|tasks| tasks.keys().copied().collect()) .unwrap_or_default() } /// Reconciles persisted tasks against live runtime handles at boot. /// /// B2 rule: `Running`/`Waiting` tasks missing from `live_task_ids` are marked /// `Failed` with a synthetic restart-loss failure and become delivery-pending. /// Already-terminal undelivered tasks are reported as delivery-pending. `Queued` /// tasks are left open. /// /// # Errors /// [`BackgroundTaskPortError`] on I/O, serialization or invalid transition. pub async fn reconcile_boot( &self, live_task_ids: &[TaskId], now_ms: u64, ) -> Result { let live: HashSet = live_task_ids.iter().copied().collect(); let mut state = self.read_all_docs().await?; let mut report = BackgroundTaskReconcileReport::default(); let mut changed_projects = BTreeSet::new(); for (project_id, doc) in &mut state.docs { for task in &mut doc.tasks { match task.state { BackgroundTaskState::Running | BackgroundTaskState::Waiting if !live.contains(&task.id) => { let finished_at_ms = now_ms.max(task.updated_at_ms); let result = BackgroundTaskResult::Failure { finished_at_ms, exit_code: None, error: "background task lost its runtime handle during IdeA restart" .into(), stdout_tail: None, stderr_tail: None, }; *task = task.complete(result).map_err(invalid_task)?; report.failed_task_ids.push(task.id); report.delivery_pending_task_ids.push(task.id); changed_projects.insert(*project_id); } _ if task.has_pending_completion_delivery() => { report.delivery_pending_task_ids.push(task.id); } _ => {} } } } for project_id in changed_projects { let live = &live; self.mutate_project_doc(project_id, |doc| { for task in &mut doc.tasks { if matches!( task.state, BackgroundTaskState::Running | BackgroundTaskState::Waiting ) && !live.contains(&task.id) { let finished_at_ms = now_ms.max(task.updated_at_ms); let result = BackgroundTaskResult::Failure { finished_at_ms, exit_code: None, error: "background task lost its runtime handle during IdeA restart" .into(), stdout_tail: None, stderr_tail: None, }; *task = task.complete(result).map_err(invalid_task)?; } } Ok(()) }) .await?; } let state = self.read_all_docs().await?; self.replace_registry_from_state(&state).await; Ok(report) } async fn ensure_registry_loaded(&self) -> Result<(), BackgroundTaskPortError> { if self.registry.read().await.loaded { return Ok(()); } let state = self.read_all_docs().await?; let mut registry = self.registry.write().await; if !registry.loaded { *registry = RegistryIndex::from_state(&state); } Ok(()) } async fn replace_registry_from_state(&self, state: &StoreState) { let mut registry = self.registry.write().await; *registry = RegistryIndex::from_state(state); } fn path_for_project(&self, project_id: ProjectId) -> PathBuf { self.dir.join(format!("{project_id}.json")) } fn tmp_path_for_project(&self, project_id: ProjectId) -> PathBuf { self.dir .join(format!("{project_id}.json.{}.tmp", Uuid::new_v4())) } async fn project_lock(&self, project_id: ProjectId) -> Arc> { let mut locks = self.project_locks.lock().await; Arc::clone( locks .entry(project_id) .or_insert_with(|| Arc::new(Mutex::new(()))), ) } async fn mutate_project_doc( &self, project_id: ProjectId, mutate: impl FnOnce(&mut TaskDoc) -> Result<(), BackgroundTaskPortError>, ) -> Result { let lock = self.project_lock(project_id).await; let _guard = lock.lock().await; let mut doc = self.read_doc(project_id).await?; mutate(&mut doc)?; self.write_doc(&doc).await?; Ok(doc) } async fn read_all_docs(&self) -> Result { let mut state = StoreState::default(); let mut entries = match tokio::fs::read_dir(&self.dir).await { Ok(entries) => entries, Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(state), Err(e) => return Err(io_error(e)), }; while let Some(entry) = entries.next_entry().await.map_err(io_error)? { let path = entry.path(); if path.extension().and_then(|ext| ext.to_str()) != Some("json") { continue; } let doc = self.read_doc_path(&path).await?; state.docs.insert(doc.project_id, doc); } Ok(state) } async fn read_doc(&self, project_id: ProjectId) -> Result { let path = self.path_for_project(project_id); match self.read_doc_path(&path).await { Ok(doc) => Ok(doc), Err(BackgroundTaskPortError::NotFound) => Ok(TaskDoc { version: TASK_DOC_VERSION, project_id, tasks: Vec::new(), }), Err(e) => Err(e), } } async fn read_doc_path(&self, path: &Path) -> Result { match tokio::fs::read(path).await { Ok(bytes) => serde_json::from_slice(&bytes) .map_err(|e| BackgroundTaskPortError::Store(e.to_string())), Err(e) if e.kind() == std::io::ErrorKind::NotFound => { Err(BackgroundTaskPortError::NotFound) } Err(e) => Err(io_error(e)), } } async fn write_doc(&self, doc: &TaskDoc) -> Result<(), BackgroundTaskPortError> { tokio::fs::create_dir_all(&self.dir) .await .map_err(io_error)?; let bytes = serde_json::to_vec_pretty(doc) .map_err(|e| BackgroundTaskPortError::Store(e.to_string()))?; let tmp = self.tmp_path_for_project(doc.project_id); tokio::fs::write(&tmp, &bytes).await.map_err(io_error)?; tokio::fs::rename(&tmp, self.path_for_project(doc.project_id)) .await .map_err(io_error)?; Ok(()) } async fn find_task_project( &self, task_id: TaskId, ) -> Result, BackgroundTaskPortError> { self.ensure_registry_loaded().await?; Ok(self .registry .read() .await .task_to_project .get(&task_id) .copied()) } async fn update_registry_for_task(&self, task: &BackgroundTask) { let mut registry = self.registry.write().await; registry.loaded = true; registry.task_to_project.insert(task.id, task.project_id); for tasks in registry.open_by_agent.values_mut() { tasks.remove(&task.id); } if !task.is_terminal() { registry .open_by_agent .entry(task.owner_agent_id) .or_default() .insert(task.id, task.clone()); } } } impl RegistryIndex { fn from_state(state: &StoreState) -> Self { let mut registry = Self { loaded: true, task_to_project: HashMap::new(), open_by_agent: HashMap::new(), }; for doc in state.docs.values() { for task in &doc.tasks { registry.task_to_project.insert(task.id, task.project_id); if !task.is_terminal() { registry .open_by_agent .entry(task.owner_agent_id) .or_default() .insert(task.id, task.clone()); } } } registry } } #[async_trait] impl BackgroundTaskStore for FsBackgroundTaskStore { async fn create(&self, task: &BackgroundTask) -> Result<(), BackgroundTaskPortError> { self.ensure_registry_loaded().await?; if self .registry .read() .await .task_to_project .contains_key(&task.id) { return Err(BackgroundTaskPortError::AlreadyExists); } self.mutate_project_doc(task.project_id, |doc| { if doc.tasks.iter().any(|existing| existing.id == task.id) { return Err(BackgroundTaskPortError::AlreadyExists); } doc.tasks.push(task.clone()); Ok(()) }) .await?; self.update_registry_for_task(task).await; Ok(()) } async fn get(&self, id: TaskId) -> Result, BackgroundTaskPortError> { let Some(project_id) = self.find_task_project(id).await? else { return Ok(None); }; let doc = self.read_doc(project_id).await?; Ok(doc.tasks.into_iter().find(|task| task.id == id)) } async fn save(&self, task: &BackgroundTask) -> Result<(), BackgroundTaskPortError> { self.ensure_registry_loaded().await?; self.mutate_project_doc(task.project_id, |doc| { if let Some(slot) = doc.tasks.iter_mut().find(|existing| existing.id == task.id) { *slot = task.clone(); Ok(()) } else { Err(BackgroundTaskPortError::NotFound) } }) .await?; self.update_registry_for_task(task).await; Ok(()) } async fn list_open_for_agent( &self, agent_id: AgentId, ) -> Result, BackgroundTaskPortError> { self.ensure_registry_loaded().await?; let registry = self.registry.read().await; let mut tasks: Vec = registry .open_by_agent .get(&agent_id) .map(|tasks| tasks.values().cloned().collect()) .unwrap_or_default(); tasks.sort_by_key(|task| (task.created_at_ms, task.id)); Ok(tasks) } async fn list_undelivered_completions( &self, ) -> Result, BackgroundTaskPortError> { let mut tasks: Vec = self .read_all_docs() .await? .docs .into_values() .flat_map(|doc| doc.tasks) .filter(BackgroundTask::has_pending_completion_delivery) .collect(); tasks.sort_by_key(|task| (task.updated_at_ms, task.id)); Ok(tasks) } async fn mark_completion_delivered( &self, task_id: TaskId, ) -> Result<(), BackgroundTaskPortError> { let task = self .get(task_id) .await? .ok_or(BackgroundTaskPortError::NotFound)?; let delivered = match task.mark_completion_delivered() { Ok(delivered) => delivered, Err(domain::BackgroundTaskError::CompletionAlreadyDelivered) => return Ok(()), Err(e) => return Err(invalid_task(e)), }; self.save(&delivered).await } } fn io_error(error: std::io::Error) -> BackgroundTaskPortError { BackgroundTaskPortError::Store(error.to_string()) } fn invalid_task(error: domain::BackgroundTaskError) -> BackgroundTaskPortError { BackgroundTaskPortError::Invalid(error.to_string()) }