Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 25 additions & 8 deletions src/crates/assembly/core/src/agentic/coordination/coordinator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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!(
Expand Down
167 changes: 167 additions & 0 deletions src/crates/assembly/core/src/agentic/persistence/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -485,6 +488,8 @@ pub struct PersistenceManager {
#[cfg(test)]
fail_next_session_state_write: std::sync::Mutex<Option<String>>,
#[cfg(test)]
fail_next_evidence_ledger_write: std::sync::Mutex<Option<String>>,
#[cfg(test)]
fail_next_session_metadata_write: std::sync::Mutex<Option<String>>,
#[cfg(test)]
fail_next_session_metadata_rollback: std::sync::Mutex<Option<String>>,
Expand All @@ -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),
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -1559,6 +1580,152 @@ impl PersistenceManager {
.await
}

pub(crate) async fn load_evidence_ledger_events(
&self,
workspace_path: &Path,
session_id: &str,
) -> BitFunResult<Vec<EvidenceLedgerEvent>> {
Self::validate_session_id(session_id)?;
let path = self.evidence_ledger_path(workspace_path, session_id);
let file = JsonFileStore
.read_locked_optional::<PersistedEvidenceLedgerFile>(&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<Vec<EvidenceLedgerEvent>> {
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::<PersistedEvidenceLedgerFile>(&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(crate) async fn retain_evidence_ledger_events(
&self,
workspace_path: &Path,
session_id: &str,
surviving_turn_ids: &std::collections::HashSet<String>,
) -> BitFunResult<Option<Vec<EvidenceLedgerEvent>>> {
Self::validate_session_id(session_id)?;
let _session_write = self.lock_session_write_operation(workspace_path, session_id)?;
let persistence_lock = self
.get_session_persistence_lock(workspace_path, session_id)
.await;
let _persistence_guard = persistence_lock.lock().await;

let path = self.evidence_ledger_path(workspace_path, session_id);
let _file_lock = JsonFileStore
.acquire_cross_process_lock(&path)
.await
.map_err(Self::json_store_error)?;
if !path.exists() {
return Ok(None);
}
let Some(mut file) = JsonFileStore
.read_optional::<PersistedEvidenceLedgerFile>(&path)
.await
.map_err(Self::json_store_error)?
else {
return Err(BitFunError::io(format!(
"Evidence ledger disappeared while retaining: {}",
path.display()
)));
};
let retained = file
.retain_turn_ids(session_id, surviving_turn_ids)
.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)?;
Ok(Some(retained))
}

/// Write a complete evidence ledger sidecar for a session. Used by session
/// branching to copy inherited evidence into the fork target. The caller
/// must already hold the session write lock for `session_id`.
pub(crate) async fn save_evidence_ledger_events(
&self,
workspace_path: &Path,
session_id: &str,
events: Vec<EvidenceLedgerEvent>,
) -> BitFunResult<()> {
Self::validate_session_id(session_id)?;
let persistence_lock = self
.get_session_persistence_lock(workspace_path, session_id)
.await;
let _persistence_guard = persistence_lock.lock().await;
self.ensure_session_dir(workspace_path, session_id).await?;
let path = self.evidence_ledger_path(workspace_path, session_id);
let _file_lock = JsonFileStore
.acquire_cross_process_lock(&path)
.await
.map_err(Self::json_store_error)?;
let file = PersistedEvidenceLedgerFile {
schema_version: EVIDENCE_LEDGER_SCHEMA_VERSION,
session_id: session_id.to_string(),
events,
};
// Validate before writing so a bad session_id on an event is caught.
file.clone()
.validated_events(session_id)
.map_err(|error| BitFunError::parse(error.to_string()))?;
JsonFileStore
.write_atomic_strict(&path, &file)
.await
.map_err(Self::json_store_error)?;
Ok(())
}

pub async fn load_prompt_cache(
&self,
workspace_path: &Path,
Expand Down
29 changes: 29 additions & 0 deletions src/crates/assembly/core/src/agentic/persistence/session_branch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,35 @@ impl PersistenceManager {
.await?;
}

// Copy evidence ledger events for the branched turns, rewriting
// session_id to the target session so the fork inherits
// checkpoints, failed commands, and partial subagent results.
let source_evidence_events = self
.load_evidence_ledger_events(workspace_path, &request.source_session_id)
.await?;
if !source_evidence_events.is_empty() {
let copied_turn_ids: std::collections::HashSet<String> = branched_turns
.iter()
.map(|turn| turn.turn_id.clone())
.collect();
let branched_evidence_events = source_evidence_events
.into_iter()
.filter(|event| copied_turn_ids.contains(&event.turn_id))
.map(|mut event| {
event.session_id = target_session_id.clone();
event
})
.collect::<Vec<_>>();
if !branched_evidence_events.is_empty() {
self.save_evidence_ledger_events(
workspace_path,
&target_session_id,
branched_evidence_events,
)
.await?;
}
}

let now_ms = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
Expand Down
Loading