From f87d08206017eb50ac20619af8e890982debab8d Mon Sep 17 00:00:00 2001 From: weishao Date: Sat, 22 Aug 2026 14:36:03 +0800 Subject: [PATCH] fix(agent-runtime): persist evidence ledger --- .../src/agentic/coordination/coordinator.rs | 33 +- .../core/src/agentic/persistence/manager.rs | 89 +++++ .../src/agentic/session/session_manager.rs | 374 +++++++++++++++++- .../tools/implementations/bash_tool.rs | 2 +- .../tools/implementations/delete_file_tool.rs | 2 +- .../tools/implementations/file_edit_tool.rs | 2 +- .../tools/implementations/file_write_tool.rs | 2 +- .../agentic/tools/implementations/git_tool.rs | 2 +- .../src/agentic/tools/tool_context_runtime.rs | 12 +- .../agent-runtime/src/evidence_ledger.rs | 289 +++++++++++++- 10 files changed, 767 insertions(+), 40 deletions(-) diff --git a/src/crates/assembly/core/src/agentic/coordination/coordinator.rs b/src/crates/assembly/core/src/agentic/coordination/coordinator.rs index 0c84931852..9e5cda3d2c 100644 --- a/src/crates/assembly/core/src/agentic/coordination/coordinator.rs +++ b/src/crates/assembly/core/src/agentic/coordination/coordinator.rs @@ -10062,14 +10062,31 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet partial_result.text.len() ); if let Some(parent_info) = subagent_parent_info.as_ref() { - let event = self.session_manager.record_subagent_partial_timeout( - &parent_info.session_id, - &parent_info.dialog_turn_id, - &logical_agent_type, - &partial_result.text, - Some("timeout"), - ); - partial_result = partial_result.with_ledger_event_id(event.event_id); + match self + .session_manager + .record_subagent_partial_timeout( + &parent_info.session_id, + &parent_info.dialog_turn_id, + &logical_agent_type, + &partial_result.text, + Some("timeout"), + ) + .await + { + Ok(event) => { + partial_result = + partial_result.with_ledger_event_id(event.event_id); + } + Err(error) => { + warn!( + "Failed to persist partial subagent evidence: parent_session_id={}, parent_turn_id={}, agent_type={}, error={}", + parent_info.session_id, + parent_info.dialog_turn_id, + logical_agent_type, + error + ); + } + } } if let Err(cleanup_err) = self.cleanup_subagent_resources(&session_id).await { warn!( diff --git a/src/crates/assembly/core/src/agentic/persistence/manager.rs b/src/crates/assembly/core/src/agentic/persistence/manager.rs index 20fc63200d..89d118a727 100644 --- a/src/crates/assembly/core/src/agentic/persistence/manager.rs +++ b/src/crates/assembly/core/src/agentic/persistence/manager.rs @@ -17,6 +17,9 @@ use crate::agentic::session::transcript_render::{ use crate::agentic::session::{ CoreSessionStorePort, SessionPromptCache, TokenAnchor, PROMPT_CACHE_SCHEMA_VERSION, }; +use crate::agentic::session::{ + EvidenceLedgerEvent, PersistedEvidenceLedgerFile, EVIDENCE_LEDGER_SCHEMA_VERSION, +}; use crate::agentic::skill_agent_snapshot::TurnSkillAgentSnapshot; use crate::infrastructure::PathManager; use crate::service::config::get_global_config_service; @@ -485,6 +488,8 @@ pub struct PersistenceManager { #[cfg(test)] fail_next_session_state_write: std::sync::Mutex>, #[cfg(test)] + fail_next_evidence_ledger_write: std::sync::Mutex>, + #[cfg(test)] fail_next_session_metadata_write: std::sync::Mutex>, #[cfg(test)] fail_next_session_metadata_rollback: std::sync::Mutex>, @@ -500,6 +505,8 @@ impl PersistenceManager { #[cfg(test)] fail_next_session_state_write: std::sync::Mutex::new(None), #[cfg(test)] + fail_next_evidence_ledger_write: std::sync::Mutex::new(None), + #[cfg(test)] fail_next_session_metadata_write: std::sync::Mutex::new(None), #[cfg(test)] fail_next_session_metadata_rollback: std::sync::Mutex::new(None), @@ -529,6 +536,14 @@ impl PersistenceManager { .expect("session state fault lock") = Some(session_id.to_string()); } + #[cfg(test)] + pub(crate) fn fail_next_evidence_ledger_write_for_test(&self, session_id: &str) { + *self + .fail_next_evidence_ledger_write + .lock() + .expect("evidence ledger fault lock") = Some(session_id.to_string()); + } + #[cfg(test)] pub(crate) fn fail_next_session_metadata_write_for_test(&self, session_id: &str) { *self @@ -608,6 +623,12 @@ impl PersistenceManager { self.session_layout(workspace_path).state_path(session_id) } + fn evidence_ledger_path(&self, workspace_path: &Path, session_id: &str) -> PathBuf { + self.session_layout(workspace_path) + .session_dir(session_id) + .join("evidence-ledger.json") + } + fn prompt_cache_path(&self, workspace_path: &Path, session_id: &str) -> PathBuf { self.session_layout(workspace_path) .prompt_cache_path(session_id) @@ -1559,6 +1580,74 @@ impl PersistenceManager { .await } + pub(crate) async fn load_evidence_ledger_events( + &self, + workspace_path: &Path, + session_id: &str, + ) -> BitFunResult> { + Self::validate_session_id(session_id)?; + let path = self.evidence_ledger_path(workspace_path, session_id); + let file = JsonFileStore + .read_locked_optional::(&path) + .await + .map_err(Self::json_store_error)?; + file.map(|file| { + file.validated_events(session_id) + .map_err(|error| BitFunError::parse(error.to_string())) + }) + .transpose() + .map(Option::unwrap_or_default) + } + + pub(crate) async fn append_evidence_ledger_event( + &self, + workspace_path: &Path, + event: &EvidenceLedgerEvent, + ) -> BitFunResult> { + Self::validate_session_id(&event.session_id)?; + let _session_write = + self.lock_session_write_operation(workspace_path, &event.session_id)?; + self.ensure_runtime_for_write(workspace_path).await?; + let persistence_lock = self + .get_session_persistence_lock(workspace_path, &event.session_id) + .await; + let _persistence_guard = persistence_lock.lock().await; + self.ensure_session_dir(workspace_path, &event.session_id) + .await?; + + #[cfg(test)] + { + let mut fault = self + .fail_next_evidence_ledger_write + .lock() + .expect("evidence ledger fault lock"); + if fault.as_deref() == Some(event.session_id.as_str()) { + *fault = None; + return Err(BitFunError::io("Injected evidence ledger write failure")); + } + } + + let path = self.evidence_ledger_path(workspace_path, &event.session_id); + let _file_lock = JsonFileStore + .acquire_cross_process_lock(&path) + .await + .map_err(Self::json_store_error)?; + let mut file = JsonFileStore + .read_optional::(&path) + .await + .map_err(Self::json_store_error)? + .unwrap_or_else(|| PersistedEvidenceLedgerFile::new(event.session_id.clone())); + file.append(event.clone()) + .map_err(|error| BitFunError::parse(error.to_string()))?; + file.schema_version = EVIDENCE_LEDGER_SCHEMA_VERSION; + JsonFileStore + .write_atomic_strict(&path, &file) + .await + .map_err(Self::json_store_error)?; + file.validated_events(&event.session_id) + .map_err(|error| BitFunError::parse(error.to_string())) + } + pub async fn load_prompt_cache( &self, workspace_path: &Path, diff --git a/src/crates/assembly/core/src/agentic/session/session_manager.rs b/src/crates/assembly/core/src/agentic/session/session_manager.rs index 5bee32f8e8..a08f14acf5 100644 --- a/src/crates/assembly/core/src/agentic/session/session_manager.rs +++ b/src/crates/assembly/core/src/agentic/session/session_manager.rs @@ -348,6 +348,7 @@ pub struct SessionManager { Arc>, file_read_state_store: Arc, evidence_ledger: Arc, + evidence_ledger_operation_locks: Arc, persistence_manager: Arc, memory_database: Arc, @@ -2027,6 +2028,7 @@ impl SessionManager { edit_constraints_store: Arc::new(DashMap::new()), file_read_state_store: Arc::new(FileReadStateStore::new()), evidence_ledger: Arc::new(SessionEvidenceLedger::new()), + evidence_ledger_operation_locks: Arc::new(KeyedAsyncLock::default()), persistence_manager, memory_database, config, @@ -2046,21 +2048,60 @@ impl SessionManager { self.persistence_manager.clone() } - pub fn append_evidence_event(&self, event: EvidenceLedgerEvent) -> EvidenceLedgerEvent { - self.evidence_ledger.append(event) + pub async fn append_evidence_event( + &self, + event: EvidenceLedgerEvent, + ) -> BitFunResult { + let _mutation_guard = self.lock_session_mutation(&event.session_id).await; + let _operation_guard = self + .evidence_ledger_operation_locks + .lock(&event.session_id) + .await; + let should_persist = self.config.enable_persistence + && self + .sessions + .get(&event.session_id) + .is_some_and(|session| self.should_persist_session(&session)); + if !should_persist { + return Ok(self.evidence_ledger.append(event)); + } + + let storage_path = self + .effective_session_storage_path(&event.session_id) + .await + .or_else(|| { + self.session_storage_path_index + .get(&event.session_id) + .map(|entry| entry.value().path.clone()) + }) + .ok_or_else(|| { + BitFunError::session(format!( + "Session storage path unavailable while persisting evidence: {}", + event.session_id + )) + })?; + let persisted_events = self + .persistence_manager + .append_evidence_ledger_event(&storage_path, &event) + .await?; + self.evidence_ledger + .replace_session(&event.session_id, persisted_events) + .map_err(|error| BitFunError::parse(error.to_string()))?; + Ok(event) } - pub fn record_checkpoint_created( + pub async fn record_checkpoint_created( &self, session_id: &str, turn_id: &str, tool_name: &str, target: &str, checkpoint: EvidenceLedgerCheckpoint, - ) -> EvidenceLedgerEvent { + ) -> BitFunResult { self.append_evidence_event(EvidenceLedgerEvent::checkpoint_created( session_id, turn_id, tool_name, target, checkpoint, )) + .await } pub fn evidence_events_for_turn( @@ -2089,14 +2130,14 @@ impl SessionManager { (!contract.is_empty()).then_some(contract) } - pub fn record_subagent_partial_timeout( + pub async fn record_subagent_partial_timeout( &self, session_id: &str, turn_id: &str, subagent_type: &str, partial_output: &str, error_kind: Option<&str>, - ) -> EvidenceLedgerEvent { + ) -> BitFunResult { let summary = format!( "Subagent {} timed out after producing partial output.", subagent_type @@ -2113,7 +2154,7 @@ impl SessionManager { .with_error_kind(error_kind.unwrap_or("timeout")) .with_partial_output(partial_output); - self.append_evidence_event(event) + self.append_evidence_event(event).await } /// Decide whether the given session model id is still usable. @@ -2360,6 +2401,7 @@ impl SessionManager { let edit_constraints_store = self.edit_constraints_store.clone(); let file_read_state_store = self.file_read_state_store.clone(); let evidence_ledger = self.evidence_ledger.clone(); + let evidence_ledger_operation_locks = self.evidence_ledger_operation_locks.clone(); let persistence_manager = self.persistence_manager.clone(); let memory_database = self.memory_database.clone(); let manager_config = self.config.clone(); @@ -2393,6 +2435,7 @@ impl SessionManager { edit_constraints_store, file_read_state_store, evidence_ledger, + evidence_ledger_operation_locks, persistence_manager, memory_database, config: manager_config, @@ -5519,6 +5562,7 @@ impl SessionManager { include_internal: bool, ) -> BitFunResult<(Session, Vec)> { let _mutation_guard = self.lock_session_mutation(session_id).await; + let _evidence_ledger_guard = self.evidence_ledger_operation_locks.lock(session_id).await; if self.is_session_loaded_from_storage_path(session_storage_path, session_id)? { let session = self.get_session(session_id).ok_or_else(|| { @@ -5621,6 +5665,10 @@ impl SessionManager { if let Some(revert) = staged_revert.as_ref() { persisted_turns.retain(|turn| turn.turn_index < revert.boundary_turn); } + let restored_evidence_events = self + .persistence_manager + .load_evidence_ledger_events(session_storage_path, session_id) + .await?; debug!( "Session restore phase completed: session_id={}, phase=load_session_with_turns, turn_count={}, duration_ms={}", session_id, @@ -5992,6 +6040,9 @@ impl SessionManager { self.evidence_ledger.as_ref(), ); } + self.evidence_ledger + .replace_session(session_id, restored_evidence_events) + .map_err(|error| BitFunError::parse(error.to_string()))?; let context_replace_started_at = Instant::now(); self.context_store @@ -9378,9 +9429,11 @@ mod tests { use crate::agentic::persistence::PersistenceManager; use crate::agentic::session::{ revert::{SessionRevertPhase, SessionRevertState, SESSION_REVERT_SCHEMA_VERSION}, - PromptCachePolicy, PromptCacheScope, SessionContextStore, SystemPromptCacheIdentity, - UserContextCacheIdentity, + EvidenceLedgerCheckpoint, PersistedEvidenceLedgerFile, PromptCachePolicy, PromptCacheScope, + SessionContextStore, SystemPromptCacheIdentity, UserContextCacheIdentity, }; + #[cfg(feature = "remote-workspace")] + use crate::agentic::session::{EvidenceLedgerEventStatus, EvidenceLedgerTargetKind}; use crate::agentic::skill_agent_snapshot::{SkillSnapshotEntry, TurnSkillAgentSnapshot}; use crate::infrastructure::ai::reasoning_catalog::{ project_model_reasoning_catalog as project_test_model_reasoning_catalog, @@ -15630,13 +15683,16 @@ mod tests { Arc::new(PersistenceManager::new(test_path_manager()).expect("persistence manager")); let manager = test_manager(persistence_manager); - let event = manager.record_subagent_partial_timeout( - "session-a", - "turn-a", - "ReviewSecurity", - "Found token logging before timeout.", - Some("timeout"), - ); + let event = manager + .record_subagent_partial_timeout( + "session-a", + "turn-a", + "ReviewSecurity", + "Found token logging before timeout.", + Some("timeout"), + ) + .await + .expect("in-memory evidence should record"); assert!(!event.event_id.is_empty()); let events = manager.evidence_events_for_turn("session-a", "turn-a"); @@ -15646,6 +15702,292 @@ mod tests { assert_eq!(summary.partial_subagent_results[0].event_id, event.event_id); } + #[tokio::test] + async fn evidence_ledger_persists_across_session_unload_and_restore() { + let workspace = TestWorkspace::new(); + let persistence_manager = Arc::new( + PersistenceManager::new(workspace.path_manager()).expect("persistence manager"), + ); + let manager = test_manager(persistence_manager); + let session = manager + .create_session( + "Durable evidence".to_string(), + "agentic".to_string(), + SessionConfig { + workspace_path: Some(workspace.path().to_string_lossy().to_string()), + ..Default::default() + }, + ) + .await + .expect("session should create"); + let storage_path = manager + .effective_session_storage_path(&session.session_id) + .await + .expect("storage path"); + let event = manager + .record_checkpoint_created( + &session.session_id, + "turn-a", + "Edit", + "src/lib.rs", + EvidenceLedgerCheckpoint { + current_branch: Some("feature/evidence".to_string()), + dirty_state_summary: "staged=0, unstaged=1, untracked=0".to_string(), + touched_files: vec!["src/lib.rs".to_string()], + diff_hash: Some("abc123".to_string()), + }, + ) + .await + .expect("checkpoint should persist before mutation"); + let ledger_path = storage_path + .join(&session.session_id) + .join("evidence-ledger.json"); + let stored: PersistedEvidenceLedgerFile = serde_json::from_slice( + &std::fs::read(&ledger_path).expect("ledger sidecar should exist"), + ) + .expect("ledger sidecar should deserialize"); + assert_eq!(stored.session_id, session.session_id); + assert_eq!(stored.events, vec![event.clone()]); + + assert!(manager + .unload_session_from_memory(&session.session_id) + .await + .expect("session should unload")); + assert!(manager + .evidence_events_for_turn(&session.session_id, "turn-a") + .is_empty()); + + manager + .restore_session_from_storage_path(&storage_path, &session.session_id) + .await + .expect("session should restore with evidence"); + assert_eq!( + manager.evidence_events_for_turn(&session.session_id, "turn-a"), + vec![event] + ); + let summary = manager.evidence_summary_for_session(&session.session_id, 10); + assert_eq!(summary.latest_checkpoints.len(), 1); + assert_eq!(summary.latest_checkpoints[0].target, "src/lib.rs"); + } + + #[tokio::test] + async fn evidence_ledger_write_failure_does_not_publish_memory_only_evidence() { + let workspace = TestWorkspace::new(); + let persistence_manager = Arc::new( + PersistenceManager::new(workspace.path_manager()).expect("persistence manager"), + ); + let manager = test_manager(persistence_manager.clone()); + let session = manager + .create_session( + "Evidence write failure".to_string(), + "agentic".to_string(), + SessionConfig { + workspace_path: Some(workspace.path().to_string_lossy().to_string()), + ..Default::default() + }, + ) + .await + .expect("session should create"); + persistence_manager.fail_next_evidence_ledger_write_for_test(&session.session_id); + + manager + .record_subagent_partial_timeout( + &session.session_id, + "turn-a", + "ReviewSecurity", + "Partial result", + Some("timeout"), + ) + .await + .expect_err("durable append failure must be visible"); + + assert!(manager + .evidence_events_for_turn(&session.session_id, "turn-a") + .is_empty()); + assert!(manager + .evidence_summary_for_session(&session.session_id, 10) + .partial_subagent_results + .is_empty()); + } + + #[tokio::test] + async fn concurrent_evidence_appends_keep_disk_and_memory_complete() { + let workspace = TestWorkspace::new(); + let persistence_manager = Arc::new( + PersistenceManager::new(workspace.path_manager()).expect("persistence manager"), + ); + let manager = Arc::new(test_manager(persistence_manager)); + let session = manager + .create_session( + "Concurrent evidence".to_string(), + "agentic".to_string(), + SessionConfig { + workspace_path: Some(workspace.path().to_string_lossy().to_string()), + ..Default::default() + }, + ) + .await + .expect("session should create"); + let storage_path = manager + .effective_session_storage_path(&session.session_id) + .await + .expect("storage path"); + + let first_manager = manager.clone(); + let first_session_id = session.session_id.clone(); + let first = tokio::spawn(async move { + first_manager + .record_subagent_partial_timeout( + &first_session_id, + "turn-a", + "ReviewSecurity", + "First partial result", + Some("timeout"), + ) + .await + .expect("first evidence append") + }); + let second_manager = manager.clone(); + let second_session_id = session.session_id.clone(); + let second = tokio::spawn(async move { + second_manager + .record_subagent_partial_timeout( + &second_session_id, + "turn-b", + "ReviewTests", + "Second partial result", + Some("timeout"), + ) + .await + .expect("second evidence append") + }); + let first = first.await.expect("first append task"); + let second = second.await.expect("second append task"); + + let ledger_path = storage_path + .join(&session.session_id) + .join("evidence-ledger.json"); + let stored: PersistedEvidenceLedgerFile = serde_json::from_slice( + &std::fs::read(ledger_path).expect("ledger sidecar should exist"), + ) + .expect("ledger sidecar should deserialize"); + let mut stored_ids = stored + .events + .into_iter() + .map(|event| event.event_id) + .collect::>(); + let mut memory_ids = manager + .evidence_ledger + .events_for_session(&session.session_id) + .into_iter() + .map(|event| event.event_id) + .collect::>(); + let mut expected_ids = vec![first.event_id, second.event_id]; + stored_ids.sort(); + memory_ids.sort(); + expected_ids.sort(); + + assert_eq!(stored_ids, expected_ids); + assert_eq!(memory_ids, expected_ids); + } + + #[tokio::test] + async fn corrupt_evidence_ledger_blocks_restore_without_overwriting_original_bytes() { + let workspace = TestWorkspace::new(); + let persistence_manager = Arc::new( + PersistenceManager::new(workspace.path_manager()).expect("persistence manager"), + ); + let manager = test_manager(persistence_manager); + let session = manager + .create_session( + "Corrupt evidence".to_string(), + "agentic".to_string(), + SessionConfig { + workspace_path: Some(workspace.path().to_string_lossy().to_string()), + ..Default::default() + }, + ) + .await + .expect("session should create"); + let storage_path = manager + .effective_session_storage_path(&session.session_id) + .await + .expect("storage path"); + let ledger_path = storage_path + .join(&session.session_id) + .join("evidence-ledger.json"); + let corrupt_bytes = b"{not valid evidence"; + std::fs::write(&ledger_path, corrupt_bytes).expect("corrupt fixture should write"); + assert!(manager + .unload_session_from_memory(&session.session_id) + .await + .expect("session should unload")); + + manager + .restore_session_from_storage_path(&storage_path, &session.session_id) + .await + .expect_err("corrupt evidence must not degrade to an empty ledger"); + + assert!(manager.get_session(&session.session_id).is_none()); + assert_eq!( + std::fs::read(&ledger_path).expect("corrupt sidecar should remain"), + corrupt_bytes + ); + } + + #[cfg(feature = "remote-workspace")] + #[tokio::test] + async fn remote_workspace_evidence_uses_the_resolved_session_mirror() { + let workspace = TestWorkspace::new(); + let path_manager = workspace.path_manager(); + let persistence_manager = + Arc::new(PersistenceManager::new(path_manager.clone()).expect("persistence manager")); + let manager = test_manager(persistence_manager); + let session = manager + .create_session( + "Remote evidence".to_string(), + "agentic".to_string(), + SessionConfig { + workspace_path: Some("/home/wsp/project".to_string()), + remote_connection_id: Some("ssh-1".to_string()), + remote_ssh_host: Some("dev-host".to_string()), + ..Default::default() + }, + ) + .await + .expect("remote session should create"); + let event = manager + .record_subagent_partial_timeout( + &session.session_id, + "turn-a", + "ReviewSecurity", + "Remote partial result", + Some("timeout"), + ) + .await + .expect("remote evidence should persist"); + let sessions_dir = crate::service::WorkspaceRuntimeService::new(path_manager) + .context_for_remote_workspace("dev-host", "/home/wsp/project") + .sessions_dir; + let ledger_path = sessions_dir + .join(&session.session_id) + .join("evidence-ledger.json"); + let stored: PersistedEvidenceLedgerFile = serde_json::from_slice( + &std::fs::read(ledger_path).expect("remote ledger sidecar should exist"), + ) + .expect("remote ledger sidecar should deserialize"); + + assert_eq!(stored.events, vec![event]); + assert_eq!( + stored.events[0].target_kind, + EvidenceLedgerTargetKind::Subagent + ); + assert_eq!( + stored.events[0].status, + EvidenceLedgerEventStatus::PartialTimeout + ); + } + #[tokio::test] async fn prompt_cache_persists_across_session_restore() { let workspace = TestWorkspace::new(); diff --git a/src/crates/assembly/core/src/agentic/tools/implementations/bash_tool.rs b/src/crates/assembly/core/src/agentic/tools/implementations/bash_tool.rs index 803ac5abd2..e0bb6e2588 100644 --- a/src/crates/assembly/core/src/agentic/tools/implementations/bash_tool.rs +++ b/src/crates/assembly/core/src/agentic/tools/implementations/bash_tool.rs @@ -605,7 +605,7 @@ Usage notes: if command_needs_light_checkpoint(command_str) { context .record_light_checkpoint("Bash", command_str, Vec::new()) - .await; + .await?; } // Remote workspace: execute via injected workspace shell diff --git a/src/crates/assembly/core/src/agentic/tools/implementations/delete_file_tool.rs b/src/crates/assembly/core/src/agentic/tools/implementations/delete_file_tool.rs index 56b27386e4..a15c93b152 100644 --- a/src/crates/assembly/core/src/agentic/tools/implementations/delete_file_tool.rs +++ b/src/crates/assembly/core/src/agentic/tools/implementations/delete_file_tool.rs @@ -338,7 +338,7 @@ Important notes: &resolved.logical_path, vec![resolved.logical_path.clone()], ) - .await; + .await?; // Remote workspace path: delete via shell command if resolved.uses_remote_workspace_backend() { diff --git a/src/crates/assembly/core/src/agentic/tools/implementations/file_edit_tool.rs b/src/crates/assembly/core/src/agentic/tools/implementations/file_edit_tool.rs index 42eab82815..18ff12be79 100644 --- a/src/crates/assembly/core/src/agentic/tools/implementations/file_edit_tool.rs +++ b/src/crates/assembly/core/src/agentic/tools/implementations/file_edit_tool.rs @@ -338,7 +338,7 @@ impl Tool for FileEditTool { &resolved.logical_path, vec![resolved.logical_path.clone()], ) - .await; + .await?; // For remote workspace paths, use the abstract FS to read → edit in memory → write back. if resolved.uses_remote_workspace_backend() { diff --git a/src/crates/assembly/core/src/agentic/tools/implementations/file_write_tool.rs b/src/crates/assembly/core/src/agentic/tools/implementations/file_write_tool.rs index 01e9452626..7763bdc184 100644 --- a/src/crates/assembly/core/src/agentic/tools/implementations/file_write_tool.rs +++ b/src/crates/assembly/core/src/agentic/tools/implementations/file_write_tool.rs @@ -529,7 +529,7 @@ impl Tool for FileWriteTool { &resolved.logical_path, vec![resolved.logical_path.clone()], ) - .await; + .await?; let file_already_exists = Self::file_exists(context, &resolved).await; if file_already_exists diff --git a/src/crates/assembly/core/src/agentic/tools/implementations/git_tool.rs b/src/crates/assembly/core/src/agentic/tools/implementations/git_tool.rs index 0fd24abec7..46ad76688a 100644 --- a/src/crates/assembly/core/src/agentic/tools/implementations/git_tool.rs +++ b/src/crates/assembly/core/src/agentic/tools/implementations/git_tool.rs @@ -1408,7 +1408,7 @@ When creating commits, use this format for the commit message: &format!("git {} {}", operation, args.unwrap_or("").trim()), Vec::new(), ) - .await; + .await?; } let start_time = std::time::Instant::now(); diff --git a/src/crates/assembly/core/src/agentic/tools/tool_context_runtime.rs b/src/crates/assembly/core/src/agentic/tools/tool_context_runtime.rs index 0d406800c9..ec377ee43b 100644 --- a/src/crates/assembly/core/src/agentic/tools/tool_context_runtime.rs +++ b/src/crates/assembly/core/src/agentic/tools/tool_context_runtime.rs @@ -396,21 +396,23 @@ impl ToolUseContext { tool_name: &str, target: &str, touched_files: Vec, - ) { + ) -> BitFunResult<()> { let Some(session_id) = self.session_id.as_deref() else { - return; + return Ok(()); }; let Some(turn_id) = self.dialog_turn_id.as_deref() else { - return; + return Ok(()); }; let Some(coordinator) = get_global_coordinator() else { - return; + return Ok(()); }; let checkpoint = self.build_light_checkpoint(touched_files).await; coordinator .get_session_manager() - .record_checkpoint_created(session_id, turn_id, tool_name, target, checkpoint); + .record_checkpoint_created(session_id, turn_id, tool_name, target, checkpoint) + .await?; + Ok(()) } async fn build_light_checkpoint(&self, touched_files: Vec) -> EvidenceLedgerCheckpoint { diff --git a/src/crates/execution/agent-runtime/src/evidence_ledger.rs b/src/crates/execution/agent-runtime/src/evidence_ledger.rs index 55035035b0..a5a6676f59 100644 --- a/src/crates/execution/agent-runtime/src/evidence_ledger.rs +++ b/src/crates/execution/agent-runtime/src/evidence_ledger.rs @@ -2,10 +2,12 @@ use crate::checkpoint::LightCheckpoint; use bitfun_runtime_ports::{CompressionContract, CompressionContractItem}; use dashmap::DashMap; use serde::{Deserialize, Serialize}; +use std::collections::HashMap; use std::sync::Arc; use std::time::{SystemTime, UNIX_EPOCH}; const MAX_PARTIAL_OUTPUT_BYTES: usize = 8_000; +pub const EVIDENCE_LEDGER_SCHEMA_VERSION: u32 = 1; #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub enum EvidenceLedgerTargetKind { @@ -19,7 +21,7 @@ pub enum EvidenceLedgerTargetKind { Artifact, #[serde(rename = "checkpoint")] Checkpoint, - #[serde(rename = "unknown")] + #[serde(other, rename = "unknown")] Unknown, } @@ -35,7 +37,7 @@ pub enum EvidenceLedgerEventStatus { PartialTimeout, #[serde(rename = "cancelled")] Cancelled, - #[serde(rename = "unknown")] + #[serde(other, rename = "unknown")] Unknown, } @@ -59,16 +61,54 @@ pub struct EvidenceLedgerEvent { pub target_kind: EvidenceLedgerTargetKind, pub target: String, pub status: EvidenceLedgerEventStatus, + #[serde(default)] pub exit_code_or_error_kind: Option, + #[serde(default)] pub touched_files: Vec, + #[serde(default)] pub artifact_path: Option, pub summary: String, + #[serde(default)] pub partial_output: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub checkpoint: Option, + #[serde(default)] pub created_at_ms: u64, } +/// Versioned sidecar owned by the runtime and written by product persistence. +/// +/// Keeping evidence separate from mutable Session state lets older builds +/// rewrite `state.json` without erasing evidence produced by a newer build. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PersistedEvidenceLedgerFile { + #[serde(default = "default_evidence_ledger_schema_version")] + pub schema_version: u32, + pub session_id: String, + #[serde(default)] + pub events: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)] +pub enum EvidenceLedgerPersistenceError { + #[error( + "unsupported evidence ledger schema version {actual}; maximum supported is {supported}" + )] + UnsupportedSchema { actual: u32, supported: u32 }, + #[error("evidence ledger session mismatch: expected {expected}, found {actual}")] + SessionMismatch { expected: String, actual: String }, + #[error("evidence ledger event has an empty event_id")] + EmptyEventId, + #[error("evidence ledger event {event_id} belongs to session {actual}, expected {expected}")] + EventSessionMismatch { + event_id: String, + expected: String, + actual: String, + }, + #[error("evidence ledger contains conflicting payloads for event {event_id}")] + ConflictingEvent { event_id: String }, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct EvidenceLedgerSummaryItem { pub event_id: String, @@ -188,6 +228,110 @@ impl From for EvidenceLedgerCheckpoint { } } +impl PersistedEvidenceLedgerFile { + pub fn new(session_id: impl Into) -> Self { + Self { + schema_version: EVIDENCE_LEDGER_SCHEMA_VERSION, + session_id: session_id.into(), + events: Vec::new(), + } + } + + pub fn append( + &mut self, + event: EvidenceLedgerEvent, + ) -> Result { + self.validate_header(&event.session_id)?; + self.validate_events()?; + validate_event(&event, &self.session_id)?; + + if let Some(existing) = self + .events + .iter() + .find(|existing| existing.event_id == event.event_id) + { + if existing == &event { + return Ok(false); + } + return Err(EvidenceLedgerPersistenceError::ConflictingEvent { + event_id: event.event_id, + }); + } + + self.events.push(event); + Ok(true) + } + + pub fn validated_events( + self, + expected_session_id: &str, + ) -> Result, EvidenceLedgerPersistenceError> { + self.validate_header(expected_session_id)?; + self.validate_events() + } + + fn validate_events(&self) -> Result, EvidenceLedgerPersistenceError> { + let mut unique_events = Vec::with_capacity(self.events.len()); + let mut event_indices = HashMap::::new(); + + for event in &self.events { + validate_event(event, &self.session_id)?; + if let Some(index) = event_indices.get(&event.event_id).copied() { + if unique_events[index] != *event { + return Err(EvidenceLedgerPersistenceError::ConflictingEvent { + event_id: event.event_id.clone(), + }); + } + continue; + } + event_indices.insert(event.event_id.clone(), unique_events.len()); + unique_events.push(event.clone()); + } + + Ok(unique_events) + } + + fn validate_header( + &self, + expected_session_id: &str, + ) -> Result<(), EvidenceLedgerPersistenceError> { + if self.schema_version > EVIDENCE_LEDGER_SCHEMA_VERSION { + return Err(EvidenceLedgerPersistenceError::UnsupportedSchema { + actual: self.schema_version, + supported: EVIDENCE_LEDGER_SCHEMA_VERSION, + }); + } + if self.session_id != expected_session_id { + return Err(EvidenceLedgerPersistenceError::SessionMismatch { + expected: expected_session_id.to_string(), + actual: self.session_id.clone(), + }); + } + Ok(()) + } +} + +fn default_evidence_ledger_schema_version() -> u32 { + EVIDENCE_LEDGER_SCHEMA_VERSION +} + +fn validate_event( + event: &EvidenceLedgerEvent, + expected_session_id: &str, +) -> Result<(), EvidenceLedgerPersistenceError> { + if event.event_id.is_empty() { + return Err(EvidenceLedgerPersistenceError::EmptyEventId); + } + if event.session_id != expected_session_id { + return Err(EvidenceLedgerPersistenceError::EventSessionMismatch { + event_id: event.event_id.clone(), + expected: expected_session_id.to_string(), + actual: event.session_id.clone(), + }); + } + Ok(()) +} + impl SessionEvidenceLedger { pub fn new() -> Self { Self::default() @@ -198,13 +342,48 @@ impl SessionEvidenceLedger { } pub fn append(&self, event: EvidenceLedgerEvent) -> EvidenceLedgerEvent { - self.events_by_session + let mut events = self + .events_by_session .entry(event.session_id.clone()) - .or_default() - .push(event.clone()); + .or_default(); + if let Some(existing) = events + .iter() + .find(|existing| existing.event_id == event.event_id) + { + return existing.clone(); + } + events.push(event.clone()); event } + pub fn replace_session( + &self, + session_id: &str, + events: Vec, + ) -> Result<(), EvidenceLedgerPersistenceError> { + let events = PersistedEvidenceLedgerFile { + schema_version: EVIDENCE_LEDGER_SCHEMA_VERSION, + session_id: session_id.to_string(), + events, + } + .validated_events(session_id)?; + + if events.is_empty() { + self.events_by_session.remove(session_id); + } else { + self.events_by_session + .insert(session_id.to_string(), events); + } + Ok(()) + } + + pub fn events_for_session(&self, session_id: &str) -> Vec { + self.events_by_session + .get(session_id) + .map(|events| events.clone()) + .unwrap_or_default() + } + pub fn events_for_turn(&self, session_id: &str, turn_id: &str) -> Vec { self.events_by_session .get(session_id) @@ -375,7 +554,8 @@ fn truncate_string_at_char_boundary(value: &str, max_bytes: usize) -> String { mod tests { use super::{ CompressionContract, EvidenceLedgerCheckpoint, EvidenceLedgerEvent, - EvidenceLedgerEventStatus, EvidenceLedgerTargetKind, SessionEvidenceLedger, + EvidenceLedgerEventStatus, EvidenceLedgerPersistenceError, EvidenceLedgerTargetKind, + PersistedEvidenceLedgerFile, SessionEvidenceLedger, EVIDENCE_LEDGER_SCHEMA_VERSION, }; use crate::checkpoint::LightCheckpoint; @@ -405,6 +585,103 @@ mod tests { assert!(ledger.events_for_turn("other-session", "turn-a").is_empty()); } + #[test] + fn persisted_ledger_append_is_idempotent_and_rejects_conflicting_replays() { + let event = EvidenceLedgerEvent::new( + "session-a", + "turn-a", + "Bash", + EvidenceLedgerTargetKind::Command, + "cargo test", + EvidenceLedgerEventStatus::Succeeded, + "Tests passed.", + ); + let mut file = PersistedEvidenceLedgerFile::new("session-a"); + + assert!(file.append(event.clone()).expect("first append")); + assert!(!file.append(event.clone()).expect("idempotent replay")); + + let mut conflicting = event.clone(); + conflicting.summary = "Different payload.".to_string(); + assert!(matches!( + file.append(conflicting), + Err(EvidenceLedgerPersistenceError::ConflictingEvent { .. }) + )); + assert_eq!(file.events, vec![event]); + } + + #[test] + fn persisted_ledger_rejects_newer_schema_and_cross_session_events() { + let newer = PersistedEvidenceLedgerFile { + schema_version: EVIDENCE_LEDGER_SCHEMA_VERSION + 1, + session_id: "session-a".to_string(), + events: Vec::new(), + }; + assert!(matches!( + newer.validated_events("session-a"), + Err(EvidenceLedgerPersistenceError::UnsupportedSchema { .. }) + )); + + let mut mismatched = PersistedEvidenceLedgerFile::new("session-a"); + mismatched.events.push(EvidenceLedgerEvent::new( + "session-b", + "turn-a", + "Bash", + EvidenceLedgerTargetKind::Command, + "cargo test", + EvidenceLedgerEventStatus::Succeeded, + "Tests passed.", + )); + assert!(matches!( + mismatched.validated_events("session-a"), + Err(EvidenceLedgerPersistenceError::EventSessionMismatch { .. }) + )); + } + + #[test] + fn legacy_evidence_event_payload_loads_with_additive_defaults() { + let file: PersistedEvidenceLedgerFile = serde_json::from_value(serde_json::json!({ + "session_id": "session-a", + "events": [{ + "event_id": "event-a", + "session_id": "session-a", + "turn_id": "turn-a", + "tool_name": "Bash", + "target_kind": "command", + "target": "cargo test", + "status": "succeeded", + "summary": "Tests passed." + }] + })) + .expect("legacy payload should deserialize"); + + assert_eq!(file.schema_version, EVIDENCE_LEDGER_SCHEMA_VERSION); + let events = file + .validated_events("session-a") + .expect("legacy payload should validate"); + assert_eq!(events.len(), 1); + assert_eq!(events[0].created_at_ms, 0); + assert!(events[0].touched_files.is_empty()); + assert!(events[0].partial_output.is_none()); + + let round_trip: PersistedEvidenceLedgerFile = serde_json::from_value( + serde_json::to_value(PersistedEvidenceLedgerFile { + schema_version: EVIDENCE_LEDGER_SCHEMA_VERSION, + session_id: "session-a".to_string(), + events, + }) + .expect("upgraded payload should serialize"), + ) + .expect("upgraded payload should deserialize"); + assert_eq!( + round_trip + .validated_events("session-a") + .expect("round trip should validate")[0] + .event_id, + "event-a" + ); + } + #[test] fn deleting_a_session_releases_its_evidence_events() { let ledger = SessionEvidenceLedger::new();