Skip to content
Open
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
89 changes: 89 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,74 @@ 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 async fn load_prompt_cache(
&self,
workspace_path: &Path,
Expand Down
Loading