use std::sync::Arc;
use bamboo_domain::session::types::Session;
use bamboo_domain::storage::Storage;
use bamboo_domain::{PermissionAuditSeed, PermissionAuditSnapshot, RuntimeSessionPersistence};
use dashmap::DashMap;
use tokio::sync::{Mutex, OwnedMutexGuard};
const AUTHORITATIVE_METADATA_KEYS: &[&str] = &["gold_config", "workflow.run_ids.v1"];
pub struct LockedSessionStore {
storage: Arc<dyn Storage>,
locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
}
pub struct SessionLockGuard {
guard: Option<OwnedMutexGuard<()>>,
locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
session_id: String,
}
impl Drop for SessionLockGuard {
fn drop(&mut self) {
self.guard.take();
self.locks
.remove_if(&self.session_id, |_, arc| Arc::strong_count(arc) == 1);
}
}
impl LockedSessionStore {
pub fn new(storage: Arc<dyn Storage>) -> Self {
Self {
storage,
locks: Arc::new(DashMap::new()),
}
}
pub fn storage(&self) -> &Arc<dyn Storage> {
&self.storage
}
pub async fn acquire_lock(&self, session_id: &str) -> SessionLockGuard {
let lock = self
.locks
.entry(session_id.to_string())
.or_insert_with(|| Arc::new(Mutex::new(())))
.clone();
let guard = lock.lock_owned().await;
SessionLockGuard {
guard: Some(guard),
locks: self.locks.clone(),
session_id: session_id.to_string(),
}
}
pub async fn save_runtime_only(&self, session: &mut Session) -> std::io::Result<()> {
self.save_runtime_only_and_publish(session, |_| {}).await
}
pub async fn save_runtime_only_and_publish<F>(
&self,
session: &mut Session,
publish: F,
) -> std::io::Result<()>
where
F: FnOnce(&Session) + Send,
{
let _guard = self.acquire_lock(&session.id).await;
if let Ok(Some(latest)) = self.storage.load_runtime_control_plane(&session.id).await {
apply_authoritative_metadata(session, &latest);
adopt_fresher_disk_permission_posture(session, &latest);
}
let result = self.storage.save_runtime_state(session).await;
publish(session);
result
}
pub async fn update_task_list_control_plane_and_publish<F>(
&self,
session_id: &str,
task_list: &bamboo_domain::TaskList,
version: &str,
publish: F,
) -> std::io::Result<bool>
where
F: FnOnce(&Session) + Send,
{
let _guard = self.acquire_lock(session_id).await;
let Some(mut latest) = self.storage.load_runtime_control_plane(session_id).await? else {
return Ok(false);
};
latest.set_task_list(task_list.clone());
latest.set_task_list_version_meta(version.to_string());
self.storage.save_runtime_state(&latest).await?;
publish(&latest);
Ok(true)
}
pub async fn commit_metadata(&self, session: &Session) -> std::io::Result<()> {
let _guard = self.acquire_lock(&session.id).await;
self.storage.save_session(session).await
}
pub async fn merge_save_runtime(&self, session: &mut Session) -> std::io::Result<()> {
self.merge_save_runtime_and_publish(session, |_, _| {})
.await
}
pub async fn merge_save_runtime_and_publish<F>(
&self,
session: &mut Session,
publish: F,
) -> std::io::Result<()>
where
F: FnOnce(&Session, bool) + Send,
{
self.merge_save_runtime_inner_and_publish(session, true, publish)
.await
}
pub async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
self.checkpoint_runtime_session_and_publish(session, |_, _| {})
.await
}
pub async fn checkpoint_runtime_session_and_publish<F>(
&self,
session: &mut Session,
publish: F,
) -> std::io::Result<()>
where
F: FnOnce(&Session, bool) + Send,
{
let _guard = self.acquire_lock(&session.id).await;
let latest = self.storage.load_session(&session.id).await?;
if let Some(latest) = latest.as_ref() {
let incoming_count = session.messages.len();
let durable_count = latest.messages.len();
let appended = bamboo_domain::append_missing_runtime_messages(session, latest);
bamboo_domain::merge_session_inbox_admission(session, latest);
tracing::debug!(
"[{}] append-safe runtime checkpoint: durable={}, incoming={}, appended={}, saved={}",
session.id,
durable_count,
incoming_count,
appended,
session.messages.len(),
);
apply_authoritative_metadata(session, latest);
adopt_fresher_disk_permission_posture(session, latest);
}
let result = self.storage.save_session(session).await;
publish(session, result.is_ok());
result
}
pub async fn save_runtime_authoritative_flags(
&self,
session: &mut Session,
) -> std::io::Result<()> {
self.merge_save_runtime_inner_and_publish(session, false, |_, _| {})
.await
}
async fn merge_save_runtime_inner_and_publish<F>(
&self,
session: &mut Session,
adopt_bypass: bool,
publish: F,
) -> std::io::Result<()>
where
F: FnOnce(&Session, bool) + Send,
{
let _guard = self.acquire_lock(&session.id).await;
let latest = self.storage.load_session(&session.id).await.ok().flatten();
let existing_message_count = latest.as_ref().map(|s| s.messages.len());
let incoming_message_count = session.messages.len();
if existing_message_count.is_some_and(|existing| existing > incoming_message_count) {
tracing::warn!(
"[{}] merge_save_runtime SHRINK: disk has {:?} messages, saving {} (last_role={:?}, updated_at={}); a stale writer is reverting a concurrent append",
session.id,
existing_message_count,
incoming_message_count,
session.messages.last().map(|m| format!("{:?}", m.role)),
session.updated_at,
);
} else {
tracing::debug!(
"[{}] merge_save_runtime: disk={:?} messages, saving {} (updated_at={})",
session.id,
existing_message_count,
incoming_message_count,
session.updated_at,
);
}
if let Some(latest) = latest.as_ref() {
apply_authoritative_metadata(session, latest);
let restored = bamboo_domain::restore_missing_admitted_inbox_messages(session, latest);
if restored > 0 {
tracing::warn!(
session_id = %session.id,
restored,
"restored durable SessionInbox transcript messages into stale runtime save"
);
}
bamboo_domain::merge_session_inbox_admission(session, latest);
if adopt_bypass {
adopt_fresher_disk_permission_posture(session, latest);
}
}
let result = self.storage.save_session(session).await;
publish(session, result.is_ok());
result
}
pub async fn seed_runtime_activation_and_publish<F>(
&self,
session: &mut Session,
publish: F,
) -> std::io::Result<()>
where
F: FnOnce(&Session, bool) + Send,
{
let _guard = self.acquire_lock(&session.id).await;
let mut incoming_audit = PermissionAuditSnapshot::from_metadata(&session.metadata)
.ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"activation seed requires a complete permission audit record",
)
})?;
if let Some(latest) = self.storage.load_session(&session.id).await? {
apply_authoritative_metadata(session, &latest);
bamboo_domain::restore_missing_admitted_inbox_messages(session, &latest);
bamboo_domain::merge_session_inbox_admission(session, &latest);
let durable_audit = PermissionAuditSnapshot::from_metadata(&latest.metadata);
let durable_floor = durable_audit
.as_ref()
.map(|snapshot| snapshot.audit_revision)
.unwrap_or_default();
if let Some(durable_audit) = durable_audit {
if durable_audit.resolution == incoming_audit.resolution {
incoming_audit.transitioned_at = durable_audit.transitioned_at;
}
}
incoming_audit.audit_revision = bamboo_domain::next_permission_audit_revision_after(
durable_floor.max(incoming_audit.audit_revision),
)
.map_err(|error| {
std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string())
})?;
}
session
.agent_runtime_state
.get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
.set_permission_mode(incoming_audit.resolution.requested);
incoming_audit.write_to(&mut session.metadata);
let result = self.storage.save_session(session).await;
publish(session, result.is_ok());
result
}
pub async fn update_authoritative_permission_posture_and_publish<M, P>(
&self,
session_id: &str,
seed: &PermissionAuditSeed,
mutate: M,
publish: P,
) -> std::io::Result<Option<Session>>
where
M: FnOnce(&mut Session),
P: FnOnce(&Session),
{
let _guard = self.acquire_lock(session_id).await;
let Some(mut latest) = self.storage.load_session(session_id).await? else {
return Ok(None);
};
let previous_mode = latest
.agent_runtime_state
.as_ref()
.map(|state| state.effective_permission_mode())
.unwrap_or_default();
let previous_resolution = PermissionAuditSnapshot::from_metadata(&latest.metadata)
.map(|snapshot| snapshot.resolution);
mutate(&mut latest);
latest
.agent_runtime_state
.get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
.set_permission_mode(seed.resolution.requested);
let mode_changed = previous_mode != seed.resolution.requested;
let posture_changed = previous_resolution != Some(seed.resolution);
let transitioned_at = posture_changed.then(|| chrono::Utc::now().to_rfc3339());
bamboo_domain::record_permission_audit(
&mut latest.metadata,
seed,
transitioned_at.as_deref(),
)
.map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string()))?;
if mode_changed {
latest.metadata_version = latest.metadata_version.saturating_add(1);
}
self.storage.save_session(&latest).await?;
publish(&latest);
Ok(Some(latest))
}
pub async fn record_permission_posture_activation_and_publish<P>(
&self,
session_id: &str,
expected_audit_revision: Option<u64>,
seed: &PermissionAuditSeed,
publish: P,
) -> std::io::Result<Option<Session>>
where
P: FnOnce(&Session),
{
let _guard = self.acquire_lock(session_id).await;
let Some(mut latest) = self.storage.load_session(session_id).await? else {
return Ok(None);
};
let durable_audit = PermissionAuditSnapshot::from_metadata(&latest.metadata);
let durable_revision = durable_audit
.as_ref()
.map(|snapshot| snapshot.audit_revision);
if durable_revision != expected_audit_revision {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"stale permission posture activation: durable audit changed after dispatch",
));
}
let durable_requested = latest
.agent_runtime_state
.as_ref()
.map(|state| state.effective_permission_mode())
.unwrap_or_default();
if durable_requested != seed.resolution.requested || !seed.resolution.is_consistent() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"stale or inconsistent permission posture activation",
));
}
bamboo_domain::record_permission_audit(&mut latest.metadata, seed, None).map_err(
|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string()),
)?;
self.storage.save_session(&latest).await?;
publish(&latest);
Ok(Some(latest))
}
pub async fn update_runtime_config<F>(
&self,
session_id: &str,
mutate: F,
) -> std::io::Result<Option<Session>>
where
F: FnOnce(&mut Session),
{
self.update_runtime_config_and_publish(session_id, mutate, |_| {})
.await
}
pub async fn update_runtime_config_and_publish<M, P>(
&self,
session_id: &str,
mutate: M,
publish: P,
) -> std::io::Result<Option<Session>>
where
M: FnOnce(&mut Session),
P: FnOnce(&Session),
{
let _guard = self.acquire_lock(session_id).await;
let Some(mut session) = self.storage.load_session(session_id).await? else {
return Ok(None);
};
mutate(&mut session);
self.storage.save_session(&session).await?;
publish(&session);
Ok(Some(session))
}
pub async fn clear_legacy_pending_messages_and_publish<F>(
&self,
session_id: &str,
expected: &[serde_json::Value],
publish: F,
) -> std::io::Result<bool>
where
F: FnOnce(&Session) + Send,
{
let _guard = self.acquire_lock(session_id).await;
let Some(mut latest) = self.storage.load_session(session_id).await? else {
return Ok(false);
};
if latest.pending_injected_messages().as_deref() != Some(expected) {
return Ok(false);
}
latest.clear_pending_injected_messages();
self.storage.save_runtime_state(&latest).await?;
publish(&latest);
Ok(true)
}
}
#[async_trait::async_trait]
impl RuntimeSessionPersistence for LockedSessionStore {
async fn save_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
self.merge_save_runtime(session).await
}
async fn seed_runtime_activation(&self, session: &mut Session) -> std::io::Result<()> {
self.seed_runtime_activation_and_publish(session, |_, _| {})
.await
}
async fn record_permission_posture_activation(
&self,
session_id: &str,
expected_audit_revision: Option<u64>,
seed: &PermissionAuditSeed,
) -> std::io::Result<Option<Session>> {
self.record_permission_posture_activation_and_publish(
session_id,
expected_audit_revision,
seed,
|_| {},
)
.await
}
async fn save_runtime_control_plane(&self, session: &mut Session) -> std::io::Result<()> {
self.save_runtime_only(session).await
}
async fn load_runtime_control_plane(
&self,
session_id: &str,
) -> std::io::Result<Option<Session>> {
self.storage.load_runtime_control_plane(session_id).await
}
async fn update_task_list_control_plane(
&self,
session_id: &str,
task_list: &bamboo_domain::TaskList,
version: &str,
) -> std::io::Result<bool> {
self.update_task_list_control_plane_and_publish(session_id, task_list, version, |_| {})
.await
}
async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
LockedSessionStore::checkpoint_runtime_session(self, session).await
}
async fn load_runtime_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
self.storage.load_session(session_id).await
}
async fn clear_legacy_pending_messages(
&self,
session_id: &str,
expected: &[serde_json::Value],
) -> std::io::Result<bool> {
self.clear_legacy_pending_messages_and_publish(session_id, expected, |_| {})
.await
}
}
async fn merge_authoritative_metadata_into_stale(
storage: &Arc<dyn Storage>,
session: &mut Session,
) {
if let Ok(Some(latest)) = storage.load_session(&session.id).await {
apply_authoritative_metadata(session, &latest);
bamboo_domain::restore_missing_admitted_inbox_messages(session, &latest);
bamboo_domain::merge_session_inbox_admission(session, &latest);
adopt_fresher_disk_permission_posture(session, &latest);
}
}
fn adopt_fresher_disk_permission_posture(session: &mut Session, latest: &Session) {
let Some(disk_mode) = latest
.agent_runtime_state
.as_ref()
.map(|state| state.effective_permission_mode())
else {
return;
};
let current_mode = session
.agent_runtime_state
.as_ref()
.map(|state| state.effective_permission_mode())
.unwrap_or_default();
let Some(disk_audit) = bamboo_domain::fresher_disk_permission_audit(
current_mode,
&session.metadata,
disk_mode,
&latest.metadata,
) else {
return;
};
match session.agent_runtime_state.as_mut() {
Some(state) => state.set_permission_mode(disk_mode),
None if disk_mode != bamboo_domain::SessionPermissionMode::Default => {
let state = session
.agent_runtime_state
.get_or_insert_with(bamboo_domain::AgentRuntimeState::default);
state.set_permission_mode(disk_mode);
}
None => {}
}
disk_audit.write_to(&mut session.metadata);
}
fn apply_authoritative_metadata(session: &mut Session, latest: &Session) {
if latest.metadata_version >= session.metadata_version {
session.title = latest.title.clone();
session.title_version = latest.title_version;
session.pinned = latest.pinned;
for key in AUTHORITATIVE_METADATA_KEYS {
if let Some(value) = latest.metadata.get(*key) {
session.metadata.insert((*key).to_string(), value.clone());
} else {
session.metadata.remove(*key);
}
}
session.metadata_version = latest.metadata_version;
}
}
pub async fn merge_save_session(
storage: &Arc<dyn Storage>,
session: &mut Session,
) -> std::io::Result<()> {
merge_authoritative_metadata_into_stale(storage, session).await;
storage.save_session(session).await
}
#[cfg(test)]
mod tests {
use super::*;
use crate::v2::SessionStoreV2;
use bamboo_domain::{session::types::Session, PermissionMode};
use std::sync::atomic::{AtomicUsize, Ordering};
struct CountingControlPlaneStorage {
inner: Arc<SessionStoreV2>,
control_plane_loads: AtomicUsize,
}
#[async_trait::async_trait]
impl Storage for CountingControlPlaneStorage {
async fn save_session(&self, session: &Session) -> std::io::Result<()> {
self.inner.save_session(session).await
}
async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
self.inner.load_session(session_id).await
}
async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
self.inner.delete_session(session_id).await
}
async fn save_runtime_state(&self, session: &Session) -> std::io::Result<()> {
self.inner.save_runtime_state(session).await
}
async fn load_runtime_control_plane(
&self,
session_id: &str,
) -> std::io::Result<Option<Session>> {
self.control_plane_loads.fetch_add(1, Ordering::SeqCst);
self.inner.load_runtime_control_plane(session_id).await
}
}
async fn make_storage() -> (tempfile::TempDir, Arc<dyn Storage>) {
let temp = tempfile::tempdir().unwrap();
let storage = SessionStoreV2::new(temp.path().to_path_buf())
.await
.expect("storage init");
(temp, Arc::new(storage) as Arc<dyn Storage>)
}
fn fresh(id: &str) -> Session {
Session::new(id.to_string(), "test-model".to_string())
}
fn set_permission_audit(
session: &mut Session,
requested: bamboo_domain::SessionPermissionMode,
policy_revision: u64,
mapping: &str,
transitioned_at: &str,
) -> u64 {
let resolution = bamboo_domain::resolve_permission_mode(requested, PermissionMode::Default);
session
.agent_runtime_state
.get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
.set_permission_mode(requested);
bamboo_domain::record_permission_audit(
&mut session.metadata,
&PermissionAuditSeed::new(policy_revision, resolution, mapping),
Some(transitioned_at),
)
.unwrap()
}
#[tokio::test]
async fn same_mode_newer_run_start_audit_survives_every_runtime_save_path() {
for path in ["merge", "checkpoint", "control-plane"] {
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = format!("same-mode-newer-{path}");
let mut durable = fresh(&session_id);
set_permission_audit(
&mut durable,
bamboo_domain::SessionPermissionMode::Default,
1,
"bamboo_runtime:old-policy",
"2026-07-31T12:00:00Z",
);
storage.save_session(&durable).await.unwrap();
let mut run_start = durable.clone();
let old_revision = PermissionAuditSnapshot::from_metadata(&durable.metadata)
.unwrap()
.audit_revision;
let new_revision = set_permission_audit(
&mut run_start,
bamboo_domain::SessionPermissionMode::Default,
2,
"bamboo_runtime:new-policy",
"2026-07-31T12:00:00Z",
);
assert!(new_revision > old_revision);
match path {
"merge" => store.merge_save_runtime(&mut run_start).await.unwrap(),
"checkpoint" => store
.checkpoint_runtime_session(&mut run_start)
.await
.unwrap(),
"control-plane" => store.save_runtime_only(&mut run_start).await.unwrap(),
_ => unreachable!(),
}
let saved = storage.load_session(&session_id).await.unwrap().unwrap();
let audit = PermissionAuditSnapshot::from_metadata(&saved.metadata).unwrap();
assert_eq!(audit.audit_revision, new_revision, "path={path}");
assert_eq!(audit.policy_revision, 2, "path={path}");
assert_eq!(audit.executor_mapping, "bamboo_runtime:new-policy");
}
}
#[tokio::test]
async fn newer_disk_transition_wins_after_mode_cycles_back() {
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "permission-cycle-back";
let mut baseline = fresh(session_id);
let stale_revision = set_permission_audit(
&mut baseline,
bamboo_domain::SessionPermissionMode::Default,
1,
"bamboo_runtime:initial",
"2026-07-31T12:00:00Z",
);
storage.save_session(&baseline).await.unwrap();
let mut stale_runtime = baseline.clone();
let mut durable = baseline;
set_permission_audit(
&mut durable,
bamboo_domain::SessionPermissionMode::Auto,
2,
"bamboo_runtime:auto",
"2026-07-31T12:01:00Z",
);
let durable_revision = set_permission_audit(
&mut durable,
bamboo_domain::SessionPermissionMode::Default,
3,
"bamboo_runtime:cycled-default",
"2026-07-31T12:02:00Z",
);
assert!(durable_revision > stale_revision);
storage.save_session(&durable).await.unwrap();
store.merge_save_runtime(&mut stale_runtime).await.unwrap();
let saved = storage.load_session(session_id).await.unwrap().unwrap();
let audit = PermissionAuditSnapshot::from_metadata(&saved.metadata).unwrap();
assert_eq!(audit.audit_revision, durable_revision);
assert_eq!(audit.policy_revision, 3);
assert_eq!(audit.executor_mapping, "bamboo_runtime:cycled-default");
}
#[tokio::test]
async fn authoritative_activation_seed_replaces_every_warm_worker_posture() {
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "warm-permission-matrix";
let cases = [
(
bamboo_domain::SessionPermissionMode::Auto,
PermissionMode::Default,
PermissionMode::Auto,
),
(
bamboo_domain::SessionPermissionMode::Default,
PermissionMode::Default,
PermissionMode::Default,
),
(
bamboo_domain::SessionPermissionMode::Auto,
PermissionMode::Default,
PermissionMode::Auto,
),
(
bamboo_domain::SessionPermissionMode::Bypass,
PermissionMode::Auto,
PermissionMode::BypassPermissions,
),
];
let mut previous_revision = 0;
for (index, (requested, configured, expected_effective)) in cases.into_iter().enumerate() {
let mut activation = fresh(session_id);
activation
.agent_runtime_state
.get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
.set_permission_mode(requested);
let resolution = bamboo_domain::resolve_permission_mode(requested, configured);
bamboo_domain::record_permission_audit(
&mut activation.metadata,
&PermissionAuditSeed::new(
index as u64 + 1,
resolution,
format!("bamboo_worker:{}", resolution.effective.as_str()),
),
Some("2026-07-31T12:00:00Z"),
)
.unwrap();
RuntimeSessionPersistence::seed_runtime_activation(&store, &mut activation)
.await
.unwrap();
let durable = storage.load_session(session_id).await.unwrap().unwrap();
assert_eq!(
durable
.agent_runtime_state
.as_ref()
.unwrap()
.effective_permission_mode(),
requested,
"activation {index}"
);
let audit = PermissionAuditSnapshot::from_metadata(&durable.metadata).unwrap();
assert_eq!(audit.resolution.requested, requested);
assert_eq!(audit.resolution.effective, expected_effective);
assert!(audit.audit_revision > previous_revision);
previous_revision = audit.audit_revision;
}
}
#[tokio::test]
async fn resident_reseed_bumps_etag_only_for_typed_transition() {
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "resident-atomic-permission";
let mut baseline = fresh(session_id);
baseline.metadata_version = 7;
set_permission_audit(
&mut baseline,
bamboo_domain::SessionPermissionMode::Auto,
1,
"bamboo_runtime:auto",
"2026-07-31T12:00:00Z",
);
storage.save_session(&baseline).await.unwrap();
let initial_audit = PermissionAuditSnapshot::from_metadata(&baseline.metadata).unwrap();
let same_mode_seed = PermissionAuditSeed::bamboo_runtime(
2,
bamboo_domain::resolve_permission_mode(
bamboo_domain::SessionPermissionMode::Auto,
PermissionMode::Default,
),
);
let refreshed = store
.update_authoritative_permission_posture_and_publish(
session_id,
&same_mode_seed,
|session| {
session
.metadata
.insert("resident.marker".to_string(), "same-mode".to_string());
},
|_| {},
)
.await
.unwrap()
.unwrap();
let refreshed_audit = PermissionAuditSnapshot::from_metadata(&refreshed.metadata).unwrap();
assert_eq!(refreshed.metadata_version, 7);
assert!(refreshed_audit.audit_revision > initial_audit.audit_revision);
assert_eq!(refreshed_audit.policy_revision, 2);
let transition_seed = PermissionAuditSeed::bamboo_runtime(
3,
bamboo_domain::resolve_permission_mode(
bamboo_domain::SessionPermissionMode::Default,
PermissionMode::Default,
),
);
let transitioned = store
.update_authoritative_permission_posture_and_publish(
session_id,
&transition_seed,
|session| {
session
.metadata
.insert("resident.marker".to_string(), "transition".to_string());
},
|_| {},
)
.await
.unwrap()
.unwrap();
let transitioned_audit =
PermissionAuditSnapshot::from_metadata(&transitioned.metadata).unwrap();
assert_eq!(transitioned.metadata_version, 8, "old ETag must be invalid");
assert_eq!(
transitioned
.agent_runtime_state
.as_ref()
.unwrap()
.effective_permission_mode(),
bamboo_domain::SessionPermissionMode::Default
);
assert_eq!(
transitioned_audit.resolution.requested,
bamboo_domain::SessionPermissionMode::Default
);
assert!(transitioned_audit.audit_revision > refreshed_audit.audit_revision);
assert_eq!(
transitioned
.metadata
.get("resident.marker")
.map(String::as_str),
Some("transition")
);
}
#[tokio::test]
async fn worker_activation_cas_cannot_overwrite_concurrent_permission_patch() {
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "permission-activation-cas";
let mut baseline = fresh(session_id);
set_permission_audit(
&mut baseline,
bamboo_domain::SessionPermissionMode::Default,
1,
"bamboo_runtime:default",
"2026-07-31T12:00:00Z",
);
storage.save_session(&baseline).await.unwrap();
let dispatched_revision = PermissionAuditSnapshot::from_metadata(&baseline.metadata)
.unwrap()
.audit_revision;
let patched_resolution = bamboo_domain::resolve_permission_mode(
bamboo_domain::SessionPermissionMode::Auto,
PermissionMode::Default,
);
let patched = store
.update_authoritative_permission_posture_and_publish(
session_id,
&PermissionAuditSeed::new(2, patched_resolution, "patch:auto"),
|_| {},
|_| {},
)
.await
.unwrap()
.unwrap();
let patched_audit = PermissionAuditSnapshot::from_metadata(&patched.metadata).unwrap();
assert!(patched_audit.audit_revision > dispatched_revision);
let stale_worker_seed = PermissionAuditSeed::new(
1,
bamboo_domain::resolve_permission_mode(
bamboo_domain::SessionPermissionMode::Default,
PermissionMode::Default,
),
"worker:stale-default",
);
let error = store
.record_permission_posture_activation_and_publish(
session_id,
Some(dispatched_revision),
&stale_worker_seed,
|_| {},
)
.await
.unwrap_err();
assert!(error.to_string().contains("durable audit changed"));
let durable = storage.load_session(session_id).await.unwrap().unwrap();
let durable_audit = PermissionAuditSnapshot::from_metadata(&durable.metadata).unwrap();
assert_eq!(durable_audit, patched_audit);
assert_eq!(durable_audit.executor_mapping, "patch:auto");
}
#[tokio::test]
async fn update_runtime_config_preserves_concurrently_appended_messages() {
use bamboo_domain::session::types::Message;
use bamboo_domain::ReasoningEffort;
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "cfg-preserve";
let mut initial = fresh(session_id);
initial.add_message(Message::user("hello"));
initial.add_message(Message::assistant("hi", None));
storage.save_session(&initial).await.unwrap();
let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
after_chat.add_message(Message::user("second question"));
storage.save_session(&after_chat).await.unwrap();
assert_eq!(after_chat.messages.len(), 3);
let updated = store
.update_runtime_config(session_id, |s| {
s.reasoning_effort = Some(ReasoningEffort::Max);
})
.await
.unwrap()
.expect("session exists");
assert_eq!(updated.reasoning_effort, Some(ReasoningEffort::Max));
assert_eq!(
updated.messages.len(),
3,
"config patch must not revert a concurrently-appended message"
);
let on_disk = storage.load_session(session_id).await.unwrap().unwrap();
assert_eq!(on_disk.messages.len(), 3);
assert_eq!(on_disk.reasoning_effort, Some(ReasoningEffort::Max));
}
#[tokio::test]
async fn update_runtime_config_returns_none_for_missing_session() {
use bamboo_domain::ReasoningEffort;
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage);
let result = store
.update_runtime_config("does-not-exist", |s| {
s.reasoning_effort = Some(ReasoningEffort::Low);
})
.await
.unwrap();
assert!(result.is_none());
}
#[tokio::test]
async fn merge_save_runtime_overwrites_messages_from_stale_snapshot() {
use bamboo_domain::session::types::Message;
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "stale-clobber";
let mut baseline = fresh(session_id);
baseline.add_message(Message::user("hello"));
storage.save_session(&baseline).await.unwrap();
let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
after_chat.add_message(Message::user("second"));
storage.save_session(&after_chat).await.unwrap();
assert_eq!(
storage
.load_session(session_id)
.await
.unwrap()
.unwrap()
.messages
.len(),
2
);
store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
let after = storage.load_session(session_id).await.unwrap().unwrap();
assert_eq!(
after.messages.len(),
1,
"merge_save_runtime clobbers concurrent appends — this is why config patches must use update_runtime_config"
);
}
#[tokio::test]
async fn stale_runtime_save_cannot_remove_admitted_inbox_transcript() {
use bamboo_domain::session::types::Message;
use bamboo_domain::SessionMessageId;
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "stale-inbox-preserve";
let mut baseline = fresh(session_id);
let mut base = Message::user("base");
base.id = "base".to_string();
baseline.add_message(base);
storage.save_session(&baseline).await.unwrap();
let mut stale = baseline.clone();
let mut later_assistant = Message::assistant("runner output", None);
later_assistant.id = "later-assistant".to_string();
stale.add_message(later_assistant);
let mut durable = baseline;
let inbox_id = SessionMessageId::parse("durable-inbox-id").unwrap();
let mut admitted = Message::user("durable inbox message");
admitted.id = inbox_id.as_str().to_string();
durable.add_message(admitted);
durable
.session_inbox_admission_mut()
.record(inbox_id.clone(), 7);
storage.save_session(&durable).await.unwrap();
store.merge_save_runtime(&mut stale).await.unwrap();
let saved = storage.load_session(session_id).await.unwrap().unwrap();
let ids = saved
.messages
.iter()
.map(|message| message.id.as_str())
.collect::<Vec<_>>();
assert_eq!(ids, vec!["base", "durable-inbox-id", "later-assistant"]);
assert_eq!(ids.iter().filter(|id| **id == inbox_id.as_str()).count(), 1);
assert!(saved
.session_inbox_admission()
.is_some_and(|state| state.contains(&inbox_id)));
}
#[tokio::test]
async fn stale_runtime_save_preserves_typed_inbox_message_after_cursor_eviction() {
use bamboo_domain::{
SessionMessageEnvelope, SessionMessageId, SESSION_INBOX_ADMITTED_CAPACITY,
};
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "evicted-inbox-preserve";
let mut durable = fresh(session_id);
let mut envelope = SessionMessageEnvelope::user_input(session_id, "old durable inbox");
envelope.id = SessionMessageId::parse("old-inbox-id").unwrap();
durable.add_message(envelope.to_provider_message().unwrap());
durable
.session_inbox_admission_mut()
.record(envelope.id.clone(), 1);
for sequence in 2..=(SESSION_INBOX_ADMITTED_CAPACITY as u64 + 1) {
durable.session_inbox_admission_mut().record(
SessionMessageId::parse(format!("newer-{sequence}")).unwrap(),
sequence,
);
}
assert!(!durable
.session_inbox_admission()
.unwrap()
.contains(&envelope.id));
storage.save_session(&durable).await.unwrap();
let mut stale = fresh(session_id);
store.merge_save_runtime(&mut stale).await.unwrap();
let saved = storage.load_session(session_id).await.unwrap().unwrap();
assert_eq!(
saved
.messages
.iter()
.filter(|message| message.id == envelope.id.as_str())
.count(),
1
);
}
#[tokio::test]
async fn checkpoint_runtime_session_preserves_disk_suffix_and_appends_live_messages() {
use bamboo_domain::session::types::Message;
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "checkpoint-no-shrink";
let mut baseline = fresh(session_id);
baseline.add_message(Message::user("base"));
storage.save_session(&baseline).await.unwrap();
let mut runner_snapshot = baseline.clone();
let mut durable = baseline;
let mut disk_only = Message::user("concurrent injected message");
disk_only.id = "disk-only".to_string();
durable.add_message(disk_only);
storage.save_session(&durable).await.unwrap();
let mut live_only = Message::assistant("partial runner output", None);
live_only.id = "live-only".to_string();
runner_snapshot.add_message(live_only);
store
.checkpoint_runtime_session(&mut runner_snapshot)
.await
.unwrap();
let saved = storage.load_session(session_id).await.unwrap().unwrap();
let ids = saved
.messages
.iter()
.map(|message| message.id.as_str())
.collect::<Vec<_>>();
assert_eq!(
ids,
vec![durable.messages[0].id.as_str(), "disk-only", "live-only"]
);
assert_eq!(runner_snapshot.messages.len(), saved.messages.len());
assert_eq!(runner_snapshot.messages[1].id, saved.messages[1].id);
assert_eq!(runner_snapshot.messages[2].id, saved.messages[2].id);
assert_eq!(saved.messages[1].content, "concurrent injected message");
assert_eq!(saved.messages[2].content, "partial runner output");
}
#[tokio::test]
async fn activation_checkpoint_clears_presentation_without_shrinking_concurrent_turn() {
use bamboo_domain::session::runtime_state::{
AgentRuntimeState, AgentStatusState, WaitingForChildrenState,
};
use bamboo_domain::session::types::Message;
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "activation-no-shrink";
let mut baseline = fresh(session_id);
baseline.add_message(Message::user("base"));
let mut state = AgentRuntimeState::new("activation-run");
state.status = AgentStatusState::Suspended;
state.waiting_for_children = Some(WaitingForChildrenState::for_children(
vec!["child-1".to_string()],
bamboo_domain::session::runtime_state::ChildWaitPolicy::All,
chrono::Utc::now(),
));
baseline.agent_runtime_state = Some(state);
baseline.metadata.insert(
"runtime.suspend_reason".to_string(),
"waiting_for_children".to_string(),
);
storage.save_session(&baseline).await.unwrap();
let mut activation_snapshot = baseline.clone();
let mut concurrent = baseline;
let mut normal = Message::assistant("normal concurrent answer", None);
normal.id = "normal-concurrent".to_string();
concurrent.add_message(normal);
storage.save_session(&concurrent).await.unwrap();
let state = activation_snapshot.agent_runtime_state.as_mut().unwrap();
state.status = AgentStatusState::Idle;
state.suspension = None;
activation_snapshot
.metadata
.remove("runtime.suspend_reason");
store
.checkpoint_runtime_session(&mut activation_snapshot)
.await
.unwrap();
let saved = storage.load_session(session_id).await.unwrap().unwrap();
assert!(saved
.messages
.iter()
.any(|message| message.id == "normal-concurrent"));
let state = saved.agent_runtime_state.unwrap();
assert_eq!(state.status, AgentStatusState::Idle);
assert!(state.waiting_for_children.is_some());
assert!(!saved.metadata.contains_key("runtime.suspend_reason"));
}
#[tokio::test]
async fn merge_save_runtime_preserves_disk_authoritative_metadata_with_single_load() {
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "runtime-merge-meta";
let mut baseline = fresh(session_id);
baseline.title = "Auto Title".to_string();
baseline.metadata_version = 0;
storage.save_session(&baseline).await.unwrap();
let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
let mut renamed = storage.load_session(session_id).await.unwrap().unwrap();
renamed.title = "User Renamed".to_string();
renamed.title_version = 1;
renamed.pinned = true;
renamed.metadata_version = 1;
store.commit_metadata(&renamed).await.unwrap();
stale_snapshot.title = "Auto Title".to_string();
store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
let after = storage.load_session(session_id).await.unwrap().unwrap();
assert_eq!(after.title, "User Renamed");
assert!(after.pinned);
assert_eq!(after.metadata_version, 1);
assert_eq!(stale_snapshot.title, "User Renamed");
assert_eq!(stale_snapshot.metadata_version, 1);
}
#[tokio::test]
async fn merge_save_runtime_preserves_durable_workflow_run_index_from_stale_runner() {
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "runtime-workflow-run-index";
let baseline = fresh(session_id);
storage.save_session(&baseline).await.unwrap();
let mut stale_runner = storage.load_session(session_id).await.unwrap().unwrap();
store
.update_runtime_config(session_id, |session| {
session.metadata.insert(
"workflow.run_ids.v1".to_string(),
r#"["http-started-run"]"#.to_string(),
);
})
.await
.unwrap()
.expect("session exists");
store.merge_save_runtime(&mut stale_runner).await.unwrap();
assert_eq!(
stale_runner
.metadata
.get("workflow.run_ids.v1")
.map(String::as_str),
Some(r#"["http-started-run"]"#)
);
let durable = storage.load_session(session_id).await.unwrap().unwrap();
assert_eq!(
durable
.metadata
.get("workflow.run_ids.v1")
.map(String::as_str),
Some(r#"["http-started-run"]"#)
);
}
#[tokio::test]
async fn merge_save_runtime_adopts_disk_bypass_permissions() {
use bamboo_domain::AgentRuntimeState;
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "runtime-bypass";
let baseline = fresh(session_id);
storage.save_session(&baseline).await.unwrap();
let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
store
.update_runtime_config(session_id, |s| {
s.agent_runtime_state
.get_or_insert_with(AgentRuntimeState::default)
.bypass_permissions = true;
})
.await
.unwrap()
.expect("session exists");
store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
let after = storage.load_session(session_id).await.unwrap().unwrap();
assert!(
after
.agent_runtime_state
.as_ref()
.is_some_and(|s| s.bypass_permissions),
"disk bypass=ON must survive a stale runtime save (#540)"
);
assert!(loop_snapshot
.agent_runtime_state
.as_ref()
.is_some_and(|s| s.bypass_permissions));
}
#[tokio::test]
async fn merge_save_runtime_adopts_disk_auto_permission_mode() {
use bamboo_domain::{AgentRuntimeState, SessionPermissionMode};
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "runtime-auto";
storage.save_session(&fresh(session_id)).await.unwrap();
let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
loop_snapshot.metadata.insert(
"permission.requested_mode".to_string(),
"default".to_string(),
);
loop_snapshot.metadata.insert(
"permission.effective_mode".to_string(),
"default".to_string(),
);
loop_snapshot.metadata.insert(
"permission.executor_mapping".to_string(),
"bamboo_runtime:default".to_string(),
);
store
.update_runtime_config(session_id, |session| {
session
.agent_runtime_state
.get_or_insert_with(AgentRuntimeState::default)
.set_permission_mode(SessionPermissionMode::Auto);
session
.metadata
.insert("permission.policy_revision".to_string(), "12".to_string());
session
.metadata
.insert("permission.requested_mode".to_string(), "auto".to_string());
session
.metadata
.insert("permission.effective_mode".to_string(), "auto".to_string());
session.metadata.insert(
"permission.executor_mapping".to_string(),
"bamboo_runtime:auto".to_string(),
);
session.metadata.insert(
"permission.transitioned_at".to_string(),
"2026-07-31T12:00:00Z".to_string(),
);
session.metadata_version = session.metadata_version.saturating_add(1);
})
.await
.unwrap()
.expect("session exists");
store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
let durable = storage.load_session(session_id).await.unwrap().unwrap();
for state in [
durable.agent_runtime_state.as_ref(),
loop_snapshot.agent_runtime_state.as_ref(),
] {
assert_eq!(
state.map(AgentRuntimeState::effective_permission_mode),
Some(SessionPermissionMode::Auto)
);
}
for session in [&durable, &loop_snapshot] {
assert_eq!(
session.metadata.get("permission.policy_revision"),
Some(&"12".to_string())
);
assert_eq!(
session.metadata.get("permission.requested_mode"),
Some(&"auto".to_string())
);
assert_eq!(
session.metadata.get("permission.effective_mode"),
Some(&"auto".to_string())
);
assert_eq!(
session.metadata.get("permission.executor_mapping"),
Some(&"bamboo_runtime:auto".to_string())
);
assert_eq!(
session.metadata.get("permission.transitioned_at"),
Some(&"2026-07-31T12:00:00Z".to_string())
);
}
}
#[tokio::test]
async fn merge_save_runtime_adopts_disk_bypass_off() {
use bamboo_domain::AgentRuntimeState;
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "runtime-bypass-off";
let mut baseline = fresh(session_id);
let on_state = AgentRuntimeState {
bypass_permissions: true,
..AgentRuntimeState::default()
};
baseline.agent_runtime_state = Some(on_state);
storage.save_session(&baseline).await.unwrap();
let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
store
.update_runtime_config(session_id, |s| {
s.agent_runtime_state
.get_or_insert_with(AgentRuntimeState::default)
.bypass_permissions = false;
})
.await
.unwrap()
.expect("session exists");
store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
let after = storage.load_session(session_id).await.unwrap().unwrap();
assert!(
!after
.agent_runtime_state
.as_ref()
.is_some_and(|s| s.bypass_permissions),
"disk bypass=OFF must survive a stale runtime save (#540)"
);
}
#[tokio::test]
async fn save_runtime_authoritative_flags_persists_in_memory_posture_and_audit() {
use bamboo_domain::AgentRuntimeState;
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "child-reseed";
let mut baseline = fresh(session_id);
let on_state = AgentRuntimeState {
bypass_permissions: true,
..AgentRuntimeState::default()
};
baseline.agent_runtime_state = Some(on_state);
for (key, value) in [
("permission.policy_revision", "12"),
("permission.requested_mode", "bypass"),
("permission.effective_mode", "bypass"),
("permission.executor_mapping", "bamboo_runtime:bypass"),
("permission.transitioned_at", "2026-07-31T12:00:00Z"),
] {
baseline.metadata.insert(key.to_string(), value.to_string());
}
storage.save_session(&baseline).await.unwrap();
let mut child = storage.load_session(session_id).await.unwrap().unwrap();
child
.agent_runtime_state
.get_or_insert_with(AgentRuntimeState::default)
.bypass_permissions = false;
for (key, value) in [
("permission.policy_revision", "13"),
("permission.requested_mode", "default"),
("permission.effective_mode", "default"),
("permission.executor_mapping", "bamboo_runtime:default"),
("permission.transitioned_at", "2026-07-31T12:01:00Z"),
] {
child.metadata.insert(key.to_string(), value.to_string());
}
store
.save_runtime_authoritative_flags(&mut child)
.await
.unwrap();
let after = storage.load_session(session_id).await.unwrap().unwrap();
assert!(
!after
.agent_runtime_state
.as_ref()
.is_some_and(|s| s.bypass_permissions),
"authoritative re-seed of bypass=OFF must persist, not be reverted (#540/#74)"
);
for (key, value) in [
("permission.policy_revision", "13"),
("permission.requested_mode", "default"),
("permission.effective_mode", "default"),
("permission.executor_mapping", "bamboo_runtime:default"),
("permission.transitioned_at", "2026-07-31T12:01:00Z"),
] {
assert_eq!(after.metadata.get(key).map(String::as_str), Some(value));
}
}
#[tokio::test]
async fn merge_save_runtime_leaves_bypass_when_disk_has_no_runtime_state() {
use bamboo_domain::AgentRuntimeState;
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "no-runtime-state";
let baseline = fresh(session_id);
assert!(baseline.agent_runtime_state.is_none());
storage.save_session(&baseline).await.unwrap();
let mut running = storage.load_session(session_id).await.unwrap().unwrap();
let on_state = AgentRuntimeState {
bypass_permissions: true,
..AgentRuntimeState::default()
};
running.agent_runtime_state = Some(on_state);
store.merge_save_runtime(&mut running).await.unwrap();
assert!(
running
.agent_runtime_state
.as_ref()
.is_some_and(|s| s.bypass_permissions),
"a runtime-state-less disk copy must not force bypass OFF (#540)"
);
}
#[tokio::test]
async fn merge_preserves_disk_title_when_versions_equal() {
let (_temp, storage) = make_storage().await;
let session_id = "merge-equal";
let mut on_disk = fresh(session_id);
on_disk.title = "User Set This".to_string();
on_disk.title_version = 0;
on_disk.metadata_version = 0;
storage.save_session(&on_disk).await.unwrap();
let mut runtime_copy = fresh(session_id);
runtime_copy.title = "Stale Default".to_string();
runtime_copy.title_version = 0;
runtime_copy.metadata_version = 0;
runtime_copy.messages = vec![];
merge_save_session(&storage, &mut runtime_copy)
.await
.unwrap();
let after = storage.load_session(session_id).await.unwrap().unwrap();
assert_eq!(after.title, "User Set This");
assert_eq!(after.title_version, 0);
assert_eq!(runtime_copy.title, "User Set This");
}
#[tokio::test]
async fn merge_preserves_disk_when_disk_version_higher() {
let (_temp, storage) = make_storage().await;
let session_id = "merge-higher";
let mut on_disk = fresh(session_id);
on_disk.title = "User Title v3".to_string();
on_disk.title_version = 3;
on_disk.metadata_version = 5;
storage.save_session(&on_disk).await.unwrap();
let mut runtime_copy = fresh(session_id);
runtime_copy.title = "Stale".to_string();
runtime_copy.title_version = 1;
runtime_copy.metadata_version = 0;
merge_save_session(&storage, &mut runtime_copy)
.await
.unwrap();
let after = storage.load_session(session_id).await.unwrap().unwrap();
assert_eq!(after.title, "User Title v3");
assert_eq!(after.title_version, 3);
assert_eq!(after.metadata_version, 5);
}
#[tokio::test]
async fn merge_now_preserves_disk_pinned_in_metadata_group() {
let (_temp, storage) = make_storage().await;
let session_id = "pinned-merge";
let mut on_disk = fresh(session_id);
on_disk.pinned = true;
on_disk.metadata_version = 2;
storage.save_session(&on_disk).await.unwrap();
let mut runtime_copy = fresh(session_id);
runtime_copy.pinned = false;
runtime_copy.metadata_version = 0;
merge_save_session(&storage, &mut runtime_copy)
.await
.unwrap();
let after = storage.load_session(session_id).await.unwrap().unwrap();
assert!(
after.pinned,
"disk pinned=true should win over runtime false"
);
assert_eq!(after.metadata_version, 2);
}
#[tokio::test]
async fn merge_keeps_in_memory_when_session_version_higher() {
let (_temp, storage) = make_storage().await;
let session_id = "merge-bumped";
let mut on_disk = fresh(session_id);
on_disk.title = "Old".to_string();
on_disk.title_version = 1;
on_disk.metadata_version = 3;
storage.save_session(&on_disk).await.unwrap();
let mut authoritative_copy = fresh(session_id);
authoritative_copy.title = "New Authoritative".to_string();
authoritative_copy.title_version = 2;
authoritative_copy.metadata_version = 4;
authoritative_copy.pinned = true;
merge_save_session(&storage, &mut authoritative_copy)
.await
.unwrap();
let after = storage.load_session(session_id).await.unwrap().unwrap();
assert_eq!(after.title, "New Authoritative");
assert_eq!(after.title_version, 2);
assert_eq!(after.metadata_version, 4);
assert!(after.pinned);
}
#[tokio::test]
async fn merge_keeps_runtime_messages_when_disk_only_changed_metadata() {
let (_temp, storage) = make_storage().await;
let session_id = "merge-messages";
let mut on_disk = fresh(session_id);
on_disk.title = "Fresh Title".to_string();
on_disk.title_version = 2;
on_disk.metadata_version = 5;
storage.save_session(&on_disk).await.unwrap();
let mut runtime_copy = fresh(session_id);
runtime_copy.title = "Stale".to_string();
runtime_copy.metadata_version = 0;
runtime_copy.messages = vec![bamboo_domain::session::types::Message {
role: bamboo_domain::session::types::Role::User,
content: "keep me".to_string(),
id: "msg-1".to_string(),
created_at: chrono::Utc::now(),
reasoning: None,
reasoning_signature: None,
content_parts: None,
image_ocr: None,
phase: None,
tool_calls: None,
tool_call_id: None,
tool_success: None,
compressed: false,
compressed_by_event_id: None,
never_compress: false,
compression_level: 0,
metadata: None,
}];
merge_save_session(&storage, &mut runtime_copy)
.await
.unwrap();
let after = storage.load_session(session_id).await.unwrap().unwrap();
assert_eq!(after.title, "Fresh Title");
assert_eq!(after.metadata_version, 5);
assert_eq!(after.messages.len(), 1);
assert_eq!(after.messages[0].content, "keep me");
}
#[tokio::test]
async fn runtime_control_plane_port_uses_sidecar_without_rewriting_messages() {
use bamboo_domain::session::types::Message;
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage.clone());
let session_id = "runtime-control-plane";
let mut durable = fresh(session_id);
durable.add_message(Message::user("durable transcript"));
storage.save_session(&durable).await.unwrap();
let mut runtime = durable.clone();
runtime.model = "updated-control-plane-model".to_string();
runtime.add_message(Message::assistant("uncheckpointed runtime message", None));
RuntimeSessionPersistence::save_runtime_control_plane(&store, &mut runtime)
.await
.unwrap();
let control_plane =
RuntimeSessionPersistence::load_runtime_control_plane(&store, session_id)
.await
.unwrap()
.expect("control-plane exists");
assert!(
control_plane.messages.is_empty(),
"LockedSessionStore must expose its message-free sidecar"
);
assert_eq!(control_plane.model, "updated-control-plane-model");
let reloaded = storage
.load_session(session_id)
.await
.unwrap()
.expect("session exists");
assert_eq!(reloaded.model, "updated-control-plane-model");
assert_eq!(
reloaded.messages.len(),
1,
"control-plane save must not write the uncheckpointed message"
);
assert_eq!(reloaded.messages[0].content, "durable transcript");
}
#[tokio::test]
async fn atomic_task_patch_loads_inside_lock_and_preserves_interleaved_runtime_state() {
let temp = tempfile::tempdir().unwrap();
let inner = Arc::new(
SessionStoreV2::new(temp.path().to_path_buf())
.await
.expect("storage init"),
);
let session_id = "atomic-task-patch";
inner
.save_session(&fresh(session_id))
.await
.expect("seed session");
let counted = Arc::new(CountingControlPlaneStorage {
inner: inner.clone(),
control_plane_loads: AtomicUsize::new(0),
});
let storage: Arc<dyn Storage> = counted.clone();
let store = Arc::new(LockedSessionStore::new(storage));
let guard = store.acquire_lock(session_id).await;
let now = chrono::Utc::now();
let task_list = bamboo_domain::TaskList {
session_id: session_id.to_string(),
title: "Atomic Task patch".to_string(),
items: Vec::new(),
created_at: now,
updated_at: now,
};
let (started_tx, started_rx) = tokio::sync::oneshot::channel();
let patch_store = store.clone();
let patch = tokio::spawn(async move {
let _ = started_tx.send(());
RuntimeSessionPersistence::update_task_list_control_plane(
patch_store.as_ref(),
session_id,
&task_list,
"9",
)
.await
});
started_rx.await.expect("patch task started");
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
assert_eq!(
counted.control_plane_loads.load(Ordering::SeqCst),
0,
"Task patch must acquire the session lock before loading its snapshot"
);
let mut latest = inner
.load_runtime_control_plane(session_id)
.await
.expect("load latest control-plane")
.expect("control-plane exists");
latest.agent_runtime_state = Some(bamboo_domain::AgentRuntimeState::new("latest-run"));
latest
.metadata
.insert("concurrent.runtime".to_string(), "preserve".to_string());
inner
.save_runtime_state(&latest)
.await
.expect("publish concurrent runtime transition");
drop(guard);
assert!(
patch.await.expect("patch join").expect("patch succeeds"),
"existing root must be patched"
);
assert_eq!(counted.control_plane_loads.load(Ordering::SeqCst), 1);
let reloaded = inner
.load_session(session_id)
.await
.expect("reload")
.expect("session exists");
assert_eq!(
reloaded
.agent_runtime_state
.as_ref()
.map(|state| state.run_id.as_str()),
Some("latest-run")
);
assert_eq!(
reloaded
.metadata
.get("concurrent.runtime")
.map(String::as_str),
Some("preserve")
);
assert_eq!(reloaded.task_list_version_meta().as_deref(), Some("9"));
assert_eq!(
reloaded.task_list.as_ref().map(|list| list.title.as_str()),
Some("Atomic Task patch")
);
}
#[tokio::test]
async fn locked_merge_save_runtime_serialises_concurrent_writes() {
let (_temp, storage) = make_storage().await;
let store = Arc::new(LockedSessionStore::new(storage));
let session_id = "lock-serial".to_string();
let base = fresh(&session_id);
store.storage().save_session(&base).await.unwrap();
let store_a = store.clone();
let store_b = store.clone();
let sid_a = session_id.clone();
let sid_b = session_id.clone();
let a = tokio::spawn(async move {
let _guard = store_a.acquire_lock(&sid_a).await;
let mut s = store_a
.storage()
.load_session(&sid_a)
.await
.unwrap()
.unwrap();
s.title = "Writer A".to_string();
s.title_version = s.title_version.saturating_add(1);
s.metadata_version = s.metadata_version.saturating_add(1);
s.updated_at = chrono::Utc::now();
store_a.storage().save_session(&s).await.unwrap();
s.title_version
});
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
let b = tokio::spawn(async move {
let _guard = store_b.acquire_lock(&sid_b).await;
let mut s = store_b
.storage()
.load_session(&sid_b)
.await
.unwrap()
.unwrap();
s.title = "Writer B".to_string();
s.title_version = s.title_version.saturating_add(1);
s.metadata_version = s.metadata_version.saturating_add(1);
s.updated_at = chrono::Utc::now();
store_b.storage().save_session(&s).await.unwrap();
s.title_version
});
let (ver_a, ver_b) = tokio::join!(a, b);
let final_s = store
.storage()
.load_session(&session_id)
.await
.unwrap()
.unwrap();
assert!(
ver_a.unwrap() != ver_b.unwrap(),
"concurrent writers must produce distinct versions"
);
assert_eq!(final_s.metadata_version, 2);
}
#[tokio::test]
async fn commit_metadata_is_plain_save_inside_lock() {
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage);
let session_id = "commit-plain";
let mut s = fresh(session_id);
s.title = "Committed".to_string();
s.metadata_version = 1;
s.title_version = 2;
store.commit_metadata(&s).await.unwrap();
let after = store
.storage()
.load_session(session_id)
.await
.unwrap()
.unwrap();
assert_eq!(after.title, "Committed");
assert_eq!(after.metadata_version, 1);
assert_eq!(after.title_version, 2);
}
#[tokio::test]
async fn acquire_lock_self_evicts_when_no_other_holder() {
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage);
{
let _guard = store.acquire_lock("solo").await;
assert_eq!(store.locks.len(), 1, "entry present while the lock is held");
}
assert_eq!(
store.locks.len(),
0,
"lock entry must be evicted once released with no other holder"
);
}
#[tokio::test]
async fn acquire_lock_many_distinct_ids_do_not_accumulate() {
let (_temp, storage) = make_storage().await;
let store = LockedSessionStore::new(storage);
for i in 0..100 {
let _guard = store.acquire_lock(&format!("sess-{i}")).await;
}
assert_eq!(
store.locks.len(),
0,
"acquiring locks for many distinct ids must not grow the map"
);
}
#[tokio::test]
async fn acquire_lock_concurrent_waiter_keeps_valid_lock_and_map_drains() {
use std::sync::atomic::{AtomicUsize, Ordering};
let (_temp, storage) = make_storage().await;
let store = Arc::new(LockedSessionStore::new(storage));
let active = Arc::new(AtomicUsize::new(0));
let max_seen = Arc::new(AtomicUsize::new(0));
let mut handles = Vec::new();
for _ in 0..8 {
let store = store.clone();
let active = active.clone();
let max_seen = max_seen.clone();
handles.push(tokio::spawn(async move {
let _guard = store.acquire_lock("contended").await;
let now = active.fetch_add(1, Ordering::SeqCst) + 1;
max_seen.fetch_max(now, Ordering::SeqCst);
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
active.fetch_sub(1, Ordering::SeqCst);
}));
}
for h in handles {
h.await.unwrap();
}
assert_eq!(
max_seen.load(Ordering::SeqCst),
1,
"at most one holder of a given session lock at a time"
);
assert_eq!(
store.locks.len(),
0,
"after all holders release, the contended entry must be fully evicted"
);
}
}