1use std::sync::Arc;
35
36use bamboo_domain::session::types::Session;
37use bamboo_domain::storage::Storage;
38use bamboo_domain::RuntimeSessionPersistence;
39use dashmap::DashMap;
40use tokio::sync::{Mutex, OwnedMutexGuard};
41
42const AUTHORITATIVE_METADATA_KEYS: &[&str] = &["gold_config", "workflow.run_ids.v1"];
43
44pub struct LockedSessionStore {
52 storage: Arc<dyn Storage>,
53 locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
54}
55
56pub struct SessionLockGuard {
79 guard: Option<OwnedMutexGuard<()>>,
81 locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
82 session_id: String,
83}
84
85impl Drop for SessionLockGuard {
86 fn drop(&mut self) {
87 self.guard.take();
90 self.locks
91 .remove_if(&self.session_id, |_, arc| Arc::strong_count(arc) == 1);
92 }
93}
94
95impl LockedSessionStore {
96 pub fn new(storage: Arc<dyn Storage>) -> Self {
98 Self {
99 storage,
100 locks: Arc::new(DashMap::new()),
101 }
102 }
103
104 pub fn storage(&self) -> &Arc<dyn Storage> {
106 &self.storage
107 }
108
109 pub async fn acquire_lock(&self, session_id: &str) -> SessionLockGuard {
121 let lock = self
126 .locks
127 .entry(session_id.to_string())
128 .or_insert_with(|| Arc::new(Mutex::new(())))
129 .clone();
130 let guard = lock.lock_owned().await;
131 SessionLockGuard {
132 guard: Some(guard),
133 locks: self.locks.clone(),
134 session_id: session_id.to_string(),
135 }
136 }
137
138 pub async fn save_runtime_only(&self, session: &mut Session) -> std::io::Result<()> {
154 let _guard = self.acquire_lock(&session.id).await;
155 if let Ok(Some(latest)) = self.storage.load_runtime_control_plane(&session.id).await {
156 apply_authoritative_metadata(session, &latest);
157 adopt_disk_bypass_permissions(session, &latest);
160 }
161 self.storage.save_runtime_state(session).await
162 }
163
164 pub async fn commit_metadata(&self, session: &Session) -> std::io::Result<()> {
174 let _guard = self.acquire_lock(&session.id).await;
175 self.storage.save_session(session).await
176 }
177
178 pub async fn merge_save_runtime(&self, session: &mut Session) -> std::io::Result<()> {
195 self.merge_save_runtime_inner(session, true).await
196 }
197
198 pub async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
207 let _guard = self.acquire_lock(&session.id).await;
208 let latest = self.storage.load_session(&session.id).await?;
209
210 if let Some(latest) = latest.as_ref() {
211 let incoming_count = session.messages.len();
212 let durable_count = latest.messages.len();
213 let appended = bamboo_domain::append_missing_runtime_messages(session, latest);
214 bamboo_domain::merge_session_inbox_admission(session, latest);
215 tracing::debug!(
216 "[{}] append-safe runtime checkpoint: durable={}, incoming={}, appended={}, saved={}",
217 session.id,
218 durable_count,
219 incoming_count,
220 appended,
221 session.messages.len(),
222 );
223 apply_authoritative_metadata(session, latest);
224 adopt_disk_bypass_permissions(session, latest);
225 }
226
227 self.storage.save_session(session).await
228 }
229
230 pub async fn save_runtime_authoritative_flags(
239 &self,
240 session: &mut Session,
241 ) -> std::io::Result<()> {
242 self.merge_save_runtime_inner(session, false).await
243 }
244
245 async fn merge_save_runtime_inner(
246 &self,
247 session: &mut Session,
248 adopt_bypass: bool,
249 ) -> std::io::Result<()> {
250 let _guard = self.acquire_lock(&session.id).await;
251
252 let latest = self.storage.load_session(&session.id).await.ok().flatten();
259
260 let existing_message_count = latest.as_ref().map(|s| s.messages.len());
266 let incoming_message_count = session.messages.len();
267 if existing_message_count.is_some_and(|existing| existing > incoming_message_count) {
268 tracing::warn!(
269 "[{}] merge_save_runtime SHRINK: disk has {:?} messages, saving {} (last_role={:?}, updated_at={}); a stale writer is reverting a concurrent append",
270 session.id,
271 existing_message_count,
272 incoming_message_count,
273 session.messages.last().map(|m| format!("{:?}", m.role)),
274 session.updated_at,
275 );
276 } else {
277 tracing::debug!(
278 "[{}] merge_save_runtime: disk={:?} messages, saving {} (updated_at={})",
279 session.id,
280 existing_message_count,
281 incoming_message_count,
282 session.updated_at,
283 );
284 }
285
286 if let Some(latest) = latest.as_ref() {
287 apply_authoritative_metadata(session, latest);
288 let restored = bamboo_domain::restore_missing_admitted_inbox_messages(session, latest);
289 if restored > 0 {
290 tracing::warn!(
291 session_id = %session.id,
292 restored,
293 "restored durable SessionInbox transcript messages into stale runtime save"
294 );
295 }
296 bamboo_domain::merge_session_inbox_admission(session, latest);
297 if adopt_bypass {
301 adopt_disk_bypass_permissions(session, latest);
302 }
303 }
304 self.storage.save_session(session).await
305 }
306
307 pub async fn update_runtime_config<F>(
319 &self,
320 session_id: &str,
321 mutate: F,
322 ) -> std::io::Result<Option<Session>>
323 where
324 F: FnOnce(&mut Session),
325 {
326 let _guard = self.acquire_lock(session_id).await;
327 let Some(mut session) = self.storage.load_session(session_id).await? else {
328 return Ok(None);
329 };
330 mutate(&mut session);
331 self.storage.save_session(&session).await?;
332 Ok(Some(session))
333 }
334}
335
336#[async_trait::async_trait]
340impl RuntimeSessionPersistence for LockedSessionStore {
341 async fn save_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
342 self.merge_save_runtime(session).await
343 }
344
345 async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
346 LockedSessionStore::checkpoint_runtime_session(self, session).await
347 }
348
349 async fn load_runtime_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
350 self.storage.load_session(session_id).await
351 }
352
353 async fn clear_legacy_pending_messages(
354 &self,
355 session_id: &str,
356 expected: &[serde_json::Value],
357 ) -> std::io::Result<bool> {
358 let _guard = self.acquire_lock(session_id).await;
359 let Some(mut latest) = self.storage.load_session(session_id).await? else {
360 return Ok(false);
361 };
362 if latest.pending_injected_messages().as_deref() != Some(expected) {
363 return Ok(false);
364 }
365 latest.clear_pending_injected_messages();
366 self.storage.save_runtime_state(&latest).await?;
367 Ok(true)
368 }
369}
370
371async fn merge_authoritative_metadata_into_stale(
380 storage: &Arc<dyn Storage>,
381 session: &mut Session,
382) {
383 if let Ok(Some(latest)) = storage.load_session(&session.id).await {
384 apply_authoritative_metadata(session, &latest);
385 bamboo_domain::restore_missing_admitted_inbox_messages(session, &latest);
386 bamboo_domain::merge_session_inbox_admission(session, &latest);
387 adopt_disk_bypass_permissions(session, &latest);
388 }
389}
390
391fn adopt_disk_bypass_permissions(session: &mut Session, latest: &Session) {
401 let Some(disk_bypass) = latest
406 .agent_runtime_state
407 .as_ref()
408 .map(|state| state.bypass_permissions)
409 else {
410 return;
411 };
412 match session.agent_runtime_state.as_mut() {
413 Some(state) => state.bypass_permissions = disk_bypass,
414 None if disk_bypass => {
417 session
418 .agent_runtime_state
419 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
420 .bypass_permissions = true;
421 }
422 None => {}
423 }
424}
425
426fn apply_authoritative_metadata(session: &mut Session, latest: &Session) {
432 if latest.metadata_version >= session.metadata_version {
433 session.title = latest.title.clone();
434 session.title_version = latest.title_version;
435 session.pinned = latest.pinned;
436 for key in AUTHORITATIVE_METADATA_KEYS {
437 if let Some(value) = latest.metadata.get(*key) {
438 session.metadata.insert((*key).to_string(), value.clone());
439 } else {
440 session.metadata.remove(*key);
441 }
442 }
443 session.metadata_version = latest.metadata_version;
444 }
445}
446
447pub async fn merge_save_session(
460 storage: &Arc<dyn Storage>,
461 session: &mut Session,
462) -> std::io::Result<()> {
463 merge_authoritative_metadata_into_stale(storage, session).await;
464 storage.save_session(session).await
465}
466
467#[cfg(test)]
470mod tests {
471 use super::*;
472 use crate::v2::SessionStoreV2;
473 use bamboo_domain::session::types::Session;
474
475 async fn make_storage() -> (tempfile::TempDir, Arc<dyn Storage>) {
476 let temp = tempfile::tempdir().unwrap();
477 let storage = SessionStoreV2::new(temp.path().to_path_buf())
478 .await
479 .expect("storage init");
480 (temp, Arc::new(storage) as Arc<dyn Storage>)
481 }
482
483 fn fresh(id: &str) -> Session {
484 Session::new(id.to_string(), "test-model".to_string())
485 }
486
487 #[tokio::test]
490 async fn update_runtime_config_preserves_concurrently_appended_messages() {
491 use bamboo_domain::session::types::Message;
492 use bamboo_domain::ReasoningEffort;
493
494 let (_temp, storage) = make_storage().await;
495 let store = LockedSessionStore::new(storage.clone());
496 let session_id = "cfg-preserve";
497
498 let mut initial = fresh(session_id);
500 initial.add_message(Message::user("hello"));
501 initial.add_message(Message::assistant("hi", None));
502 storage.save_session(&initial).await.unwrap();
503
504 let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
506 after_chat.add_message(Message::user("second question"));
507 storage.save_session(&after_chat).await.unwrap();
508 assert_eq!(after_chat.messages.len(), 3);
509
510 let updated = store
514 .update_runtime_config(session_id, |s| {
515 s.reasoning_effort = Some(ReasoningEffort::Max);
516 })
517 .await
518 .unwrap()
519 .expect("session exists");
520
521 assert_eq!(updated.reasoning_effort, Some(ReasoningEffort::Max));
522 assert_eq!(
523 updated.messages.len(),
524 3,
525 "config patch must not revert a concurrently-appended message"
526 );
527
528 let on_disk = storage.load_session(session_id).await.unwrap().unwrap();
529 assert_eq!(on_disk.messages.len(), 3);
530 assert_eq!(on_disk.reasoning_effort, Some(ReasoningEffort::Max));
531 }
532
533 #[tokio::test]
534 async fn update_runtime_config_returns_none_for_missing_session() {
535 use bamboo_domain::ReasoningEffort;
536
537 let (_temp, storage) = make_storage().await;
538 let store = LockedSessionStore::new(storage);
539 let result = store
540 .update_runtime_config("does-not-exist", |s| {
541 s.reasoning_effort = Some(ReasoningEffort::Low);
542 })
543 .await
544 .unwrap();
545 assert!(result.is_none());
546 }
547
548 #[tokio::test]
549 async fn merge_save_runtime_overwrites_messages_from_stale_snapshot() {
550 use bamboo_domain::session::types::Message;
555
556 let (_temp, storage) = make_storage().await;
557 let store = LockedSessionStore::new(storage.clone());
558 let session_id = "stale-clobber";
559
560 let mut baseline = fresh(session_id);
562 baseline.add_message(Message::user("hello"));
563 storage.save_session(&baseline).await.unwrap();
564 let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
565
566 let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
568 after_chat.add_message(Message::user("second"));
569 storage.save_session(&after_chat).await.unwrap();
570 assert_eq!(
571 storage
572 .load_session(session_id)
573 .await
574 .unwrap()
575 .unwrap()
576 .messages
577 .len(),
578 2
579 );
580
581 store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
583 let after = storage.load_session(session_id).await.unwrap().unwrap();
584 assert_eq!(
585 after.messages.len(),
586 1,
587 "merge_save_runtime clobbers concurrent appends — this is why config patches must use update_runtime_config"
588 );
589 }
590
591 #[tokio::test]
592 async fn stale_runtime_save_cannot_remove_admitted_inbox_transcript() {
593 use bamboo_domain::session::types::Message;
594 use bamboo_domain::SessionMessageId;
595
596 let (_temp, storage) = make_storage().await;
597 let store = LockedSessionStore::new(storage.clone());
598 let session_id = "stale-inbox-preserve";
599
600 let mut baseline = fresh(session_id);
601 let mut base = Message::user("base");
602 base.id = "base".to_string();
603 baseline.add_message(base);
604 storage.save_session(&baseline).await.unwrap();
605 let mut stale = baseline.clone();
606 let mut later_assistant = Message::assistant("runner output", None);
607 later_assistant.id = "later-assistant".to_string();
608 stale.add_message(later_assistant);
609
610 let mut durable = baseline;
611 let inbox_id = SessionMessageId::parse("durable-inbox-id").unwrap();
612 let mut admitted = Message::user("durable inbox message");
613 admitted.id = inbox_id.as_str().to_string();
614 durable.add_message(admitted);
615 durable
616 .session_inbox_admission_mut()
617 .record(inbox_id.clone(), 7);
618 storage.save_session(&durable).await.unwrap();
619
620 store.merge_save_runtime(&mut stale).await.unwrap();
621 let saved = storage.load_session(session_id).await.unwrap().unwrap();
622 let ids = saved
623 .messages
624 .iter()
625 .map(|message| message.id.as_str())
626 .collect::<Vec<_>>();
627 assert_eq!(ids, vec!["base", "durable-inbox-id", "later-assistant"]);
628 assert_eq!(ids.iter().filter(|id| **id == inbox_id.as_str()).count(), 1);
629 assert!(saved
630 .session_inbox_admission()
631 .is_some_and(|state| state.contains(&inbox_id)));
632 }
633
634 #[tokio::test]
635 async fn stale_runtime_save_preserves_typed_inbox_message_after_cursor_eviction() {
636 use bamboo_domain::{
637 SessionMessageEnvelope, SessionMessageId, SESSION_INBOX_ADMITTED_CAPACITY,
638 };
639
640 let (_temp, storage) = make_storage().await;
641 let store = LockedSessionStore::new(storage.clone());
642 let session_id = "evicted-inbox-preserve";
643 let mut durable = fresh(session_id);
644 let mut envelope = SessionMessageEnvelope::user_input(session_id, "old durable inbox");
645 envelope.id = SessionMessageId::parse("old-inbox-id").unwrap();
646 durable.add_message(envelope.to_provider_message().unwrap());
647 durable
648 .session_inbox_admission_mut()
649 .record(envelope.id.clone(), 1);
650 for sequence in 2..=(SESSION_INBOX_ADMITTED_CAPACITY as u64 + 1) {
651 durable.session_inbox_admission_mut().record(
652 SessionMessageId::parse(format!("newer-{sequence}")).unwrap(),
653 sequence,
654 );
655 }
656 assert!(!durable
657 .session_inbox_admission()
658 .unwrap()
659 .contains(&envelope.id));
660 storage.save_session(&durable).await.unwrap();
661
662 let mut stale = fresh(session_id);
663 store.merge_save_runtime(&mut stale).await.unwrap();
664 let saved = storage.load_session(session_id).await.unwrap().unwrap();
665 assert_eq!(
666 saved
667 .messages
668 .iter()
669 .filter(|message| message.id == envelope.id.as_str())
670 .count(),
671 1
672 );
673 }
674
675 #[tokio::test]
676 async fn checkpoint_runtime_session_preserves_disk_suffix_and_appends_live_messages() {
677 use bamboo_domain::session::types::Message;
678
679 let (_temp, storage) = make_storage().await;
680 let store = LockedSessionStore::new(storage.clone());
681 let session_id = "checkpoint-no-shrink";
682
683 let mut baseline = fresh(session_id);
684 baseline.add_message(Message::user("base"));
685 storage.save_session(&baseline).await.unwrap();
686 let mut runner_snapshot = baseline.clone();
687
688 let mut durable = baseline;
689 let mut disk_only = Message::user("concurrent injected message");
690 disk_only.id = "disk-only".to_string();
691 durable.add_message(disk_only);
692 storage.save_session(&durable).await.unwrap();
693
694 let mut live_only = Message::assistant("partial runner output", None);
695 live_only.id = "live-only".to_string();
696 runner_snapshot.add_message(live_only);
697
698 store
699 .checkpoint_runtime_session(&mut runner_snapshot)
700 .await
701 .unwrap();
702
703 let saved = storage.load_session(session_id).await.unwrap().unwrap();
704 let ids = saved
705 .messages
706 .iter()
707 .map(|message| message.id.as_str())
708 .collect::<Vec<_>>();
709 assert_eq!(
710 ids,
711 vec![durable.messages[0].id.as_str(), "disk-only", "live-only"]
712 );
713 assert_eq!(runner_snapshot.messages.len(), saved.messages.len());
714 assert_eq!(runner_snapshot.messages[1].id, saved.messages[1].id);
715 assert_eq!(runner_snapshot.messages[2].id, saved.messages[2].id);
716 assert_eq!(saved.messages[1].content, "concurrent injected message");
717 assert_eq!(saved.messages[2].content, "partial runner output");
718 }
719
720 #[tokio::test]
721 async fn activation_checkpoint_clears_presentation_without_shrinking_concurrent_turn() {
722 use bamboo_domain::session::runtime_state::{
723 AgentRuntimeState, AgentStatusState, WaitingForChildrenState,
724 };
725 use bamboo_domain::session::types::Message;
726
727 let (_temp, storage) = make_storage().await;
728 let store = LockedSessionStore::new(storage.clone());
729 let session_id = "activation-no-shrink";
730 let mut baseline = fresh(session_id);
731 baseline.add_message(Message::user("base"));
732 let mut state = AgentRuntimeState::new("activation-run");
733 state.status = AgentStatusState::Suspended;
734 state.waiting_for_children = Some(WaitingForChildrenState::for_children(
735 vec!["child-1".to_string()],
736 bamboo_domain::session::runtime_state::ChildWaitPolicy::All,
737 chrono::Utc::now(),
738 ));
739 baseline.agent_runtime_state = Some(state);
740 baseline.metadata.insert(
741 "runtime.suspend_reason".to_string(),
742 "waiting_for_children".to_string(),
743 );
744 storage.save_session(&baseline).await.unwrap();
745 let mut activation_snapshot = baseline.clone();
746
747 let mut concurrent = baseline;
748 let mut normal = Message::assistant("normal concurrent answer", None);
749 normal.id = "normal-concurrent".to_string();
750 concurrent.add_message(normal);
751 storage.save_session(&concurrent).await.unwrap();
752
753 let state = activation_snapshot.agent_runtime_state.as_mut().unwrap();
754 state.status = AgentStatusState::Idle;
755 state.suspension = None;
756 activation_snapshot
757 .metadata
758 .remove("runtime.suspend_reason");
759 store
760 .checkpoint_runtime_session(&mut activation_snapshot)
761 .await
762 .unwrap();
763
764 let saved = storage.load_session(session_id).await.unwrap().unwrap();
765 assert!(saved
766 .messages
767 .iter()
768 .any(|message| message.id == "normal-concurrent"));
769 let state = saved.agent_runtime_state.unwrap();
770 assert_eq!(state.status, AgentStatusState::Idle);
771 assert!(state.waiting_for_children.is_some());
772 assert!(!saved.metadata.contains_key("runtime.suspend_reason"));
773 }
774
775 #[tokio::test]
776 async fn merge_save_runtime_preserves_disk_authoritative_metadata_with_single_load() {
777 let (_temp, storage) = make_storage().await;
783 let store = LockedSessionStore::new(storage.clone());
784 let session_id = "runtime-merge-meta";
785
786 let mut baseline = fresh(session_id);
788 baseline.title = "Auto Title".to_string();
789 baseline.metadata_version = 0;
790 storage.save_session(&baseline).await.unwrap();
791
792 let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
794
795 let mut renamed = storage.load_session(session_id).await.unwrap().unwrap();
797 renamed.title = "User Renamed".to_string();
798 renamed.title_version = 1;
799 renamed.pinned = true;
800 renamed.metadata_version = 1;
801 store.commit_metadata(&renamed).await.unwrap();
802
803 stale_snapshot.title = "Auto Title".to_string();
805 store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
806
807 let after = storage.load_session(session_id).await.unwrap().unwrap();
808 assert_eq!(after.title, "User Renamed");
809 assert!(after.pinned);
810 assert_eq!(after.metadata_version, 1);
811 assert_eq!(stale_snapshot.title, "User Renamed");
813 assert_eq!(stale_snapshot.metadata_version, 1);
814 }
815
816 #[tokio::test]
817 async fn merge_save_runtime_preserves_durable_workflow_run_index_from_stale_runner() {
818 let (_temp, storage) = make_storage().await;
819 let store = LockedSessionStore::new(storage.clone());
820 let session_id = "runtime-workflow-run-index";
821
822 let baseline = fresh(session_id);
823 storage.save_session(&baseline).await.unwrap();
824 let mut stale_runner = storage.load_session(session_id).await.unwrap().unwrap();
825
826 store
827 .update_runtime_config(session_id, |session| {
828 session.metadata.insert(
829 "workflow.run_ids.v1".to_string(),
830 r#"["http-started-run"]"#.to_string(),
831 );
832 })
833 .await
834 .unwrap()
835 .expect("session exists");
836
837 store.merge_save_runtime(&mut stale_runner).await.unwrap();
838
839 assert_eq!(
840 stale_runner
841 .metadata
842 .get("workflow.run_ids.v1")
843 .map(String::as_str),
844 Some(r#"["http-started-run"]"#)
845 );
846 let durable = storage.load_session(session_id).await.unwrap().unwrap();
847 assert_eq!(
848 durable
849 .metadata
850 .get("workflow.run_ids.v1")
851 .map(String::as_str),
852 Some(r#"["http-started-run"]"#)
853 );
854 }
855
856 #[tokio::test]
860 async fn merge_save_runtime_adopts_disk_bypass_permissions() {
861 use bamboo_domain::AgentRuntimeState;
862
863 let (_temp, storage) = make_storage().await;
864 let store = LockedSessionStore::new(storage.clone());
865 let session_id = "runtime-bypass";
866
867 let baseline = fresh(session_id);
869 storage.save_session(&baseline).await.unwrap();
870
871 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
873 loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
874
875 store
877 .update_runtime_config(session_id, |s| {
878 s.agent_runtime_state
879 .get_or_insert_with(AgentRuntimeState::default)
880 .bypass_permissions = true;
881 })
882 .await
883 .unwrap()
884 .expect("session exists");
885
886 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
889
890 let after = storage.load_session(session_id).await.unwrap().unwrap();
891 assert!(
892 after
893 .agent_runtime_state
894 .as_ref()
895 .is_some_and(|s| s.bypass_permissions),
896 "disk bypass=ON must survive a stale runtime save (#540)"
897 );
898 assert!(loop_snapshot
900 .agent_runtime_state
901 .as_ref()
902 .is_some_and(|s| s.bypass_permissions));
903 }
904
905 #[tokio::test]
908 async fn merge_save_runtime_adopts_disk_bypass_off() {
909 use bamboo_domain::AgentRuntimeState;
910
911 let (_temp, storage) = make_storage().await;
912 let store = LockedSessionStore::new(storage.clone());
913 let session_id = "runtime-bypass-off";
914
915 let mut baseline = fresh(session_id);
917 let mut on_state = AgentRuntimeState::default();
918 on_state.bypass_permissions = true;
919 baseline.agent_runtime_state = Some(on_state);
920 storage.save_session(&baseline).await.unwrap();
921
922 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
924
925 store
927 .update_runtime_config(session_id, |s| {
928 s.agent_runtime_state
929 .get_or_insert_with(AgentRuntimeState::default)
930 .bypass_permissions = false;
931 })
932 .await
933 .unwrap()
934 .expect("session exists");
935
936 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
937
938 let after = storage.load_session(session_id).await.unwrap().unwrap();
939 assert!(
940 !after
941 .agent_runtime_state
942 .as_ref()
943 .is_some_and(|s| s.bypass_permissions),
944 "disk bypass=OFF must survive a stale runtime save (#540)"
945 );
946 }
947
948 #[tokio::test]
951 async fn save_runtime_authoritative_flags_persists_in_memory_bypass() {
952 use bamboo_domain::AgentRuntimeState;
953
954 let (_temp, storage) = make_storage().await;
955 let store = LockedSessionStore::new(storage.clone());
956 let session_id = "child-reseed";
957
958 let mut baseline = fresh(session_id);
960 let mut on_state = AgentRuntimeState::default();
961 on_state.bypass_permissions = true;
962 baseline.agent_runtime_state = Some(on_state);
963 storage.save_session(&baseline).await.unwrap();
964
965 let mut child = storage.load_session(session_id).await.unwrap().unwrap();
968 child
969 .agent_runtime_state
970 .get_or_insert_with(AgentRuntimeState::default)
971 .bypass_permissions = false;
972
973 store
975 .save_runtime_authoritative_flags(&mut child)
976 .await
977 .unwrap();
978
979 let after = storage.load_session(session_id).await.unwrap().unwrap();
980 assert!(
981 !after
982 .agent_runtime_state
983 .as_ref()
984 .is_some_and(|s| s.bypass_permissions),
985 "authoritative re-seed of bypass=OFF must persist, not be reverted (#540/#74)"
986 );
987 }
988
989 #[tokio::test]
991 async fn merge_save_runtime_leaves_bypass_when_disk_has_no_runtime_state() {
992 use bamboo_domain::AgentRuntimeState;
993
994 let (_temp, storage) = make_storage().await;
995 let store = LockedSessionStore::new(storage.clone());
996 let session_id = "no-runtime-state";
997
998 let baseline = fresh(session_id);
1000 assert!(baseline.agent_runtime_state.is_none());
1001 storage.save_session(&baseline).await.unwrap();
1002
1003 let mut running = storage.load_session(session_id).await.unwrap().unwrap();
1005 let mut on_state = AgentRuntimeState::default();
1006 on_state.bypass_permissions = true;
1007 running.agent_runtime_state = Some(on_state);
1008
1009 store.merge_save_runtime(&mut running).await.unwrap();
1010
1011 assert!(
1012 running
1013 .agent_runtime_state
1014 .as_ref()
1015 .is_some_and(|s| s.bypass_permissions),
1016 "a runtime-state-less disk copy must not force bypass OFF (#540)"
1017 );
1018 }
1019
1020 #[tokio::test]
1023 async fn merge_preserves_disk_title_when_versions_equal() {
1024 let (_temp, storage) = make_storage().await;
1025 let session_id = "merge-equal";
1026
1027 let mut on_disk = fresh(session_id);
1028 on_disk.title = "User Set This".to_string();
1029 on_disk.title_version = 0;
1030 on_disk.metadata_version = 0;
1031 storage.save_session(&on_disk).await.unwrap();
1032
1033 let mut runtime_copy = fresh(session_id);
1034 runtime_copy.title = "Stale Default".to_string();
1035 runtime_copy.title_version = 0;
1036 runtime_copy.metadata_version = 0;
1037 runtime_copy.messages = vec![];
1038
1039 merge_save_session(&storage, &mut runtime_copy)
1040 .await
1041 .unwrap();
1042
1043 let after = storage.load_session(session_id).await.unwrap().unwrap();
1044 assert_eq!(after.title, "User Set This");
1045 assert_eq!(after.title_version, 0);
1046 assert_eq!(runtime_copy.title, "User Set This");
1047 }
1048
1049 #[tokio::test]
1050 async fn merge_preserves_disk_when_disk_version_higher() {
1051 let (_temp, storage) = make_storage().await;
1052 let session_id = "merge-higher";
1053
1054 let mut on_disk = fresh(session_id);
1055 on_disk.title = "User Title v3".to_string();
1056 on_disk.title_version = 3;
1057 on_disk.metadata_version = 5;
1058 storage.save_session(&on_disk).await.unwrap();
1059
1060 let mut runtime_copy = fresh(session_id);
1061 runtime_copy.title = "Stale".to_string();
1062 runtime_copy.title_version = 1;
1063 runtime_copy.metadata_version = 0;
1064
1065 merge_save_session(&storage, &mut runtime_copy)
1066 .await
1067 .unwrap();
1068
1069 let after = storage.load_session(session_id).await.unwrap().unwrap();
1070 assert_eq!(after.title, "User Title v3");
1071 assert_eq!(after.title_version, 3);
1072 assert_eq!(after.metadata_version, 5);
1073 }
1074
1075 #[tokio::test]
1076 async fn merge_now_preserves_disk_pinned_in_metadata_group() {
1077 let (_temp, storage) = make_storage().await;
1078 let session_id = "pinned-merge";
1079
1080 let mut on_disk = fresh(session_id);
1081 on_disk.pinned = true;
1082 on_disk.metadata_version = 2;
1083 storage.save_session(&on_disk).await.unwrap();
1084
1085 let mut runtime_copy = fresh(session_id);
1086 runtime_copy.pinned = false;
1087 runtime_copy.metadata_version = 0;
1088
1089 merge_save_session(&storage, &mut runtime_copy)
1090 .await
1091 .unwrap();
1092
1093 let after = storage.load_session(session_id).await.unwrap().unwrap();
1094 assert!(
1095 after.pinned,
1096 "disk pinned=true should win over runtime false"
1097 );
1098 assert_eq!(after.metadata_version, 2);
1099 }
1100
1101 #[tokio::test]
1102 async fn merge_keeps_in_memory_when_session_version_higher() {
1103 let (_temp, storage) = make_storage().await;
1104 let session_id = "merge-bumped";
1105
1106 let mut on_disk = fresh(session_id);
1107 on_disk.title = "Old".to_string();
1108 on_disk.title_version = 1;
1109 on_disk.metadata_version = 3;
1110 storage.save_session(&on_disk).await.unwrap();
1111
1112 let mut authoritative_copy = fresh(session_id);
1113 authoritative_copy.title = "New Authoritative".to_string();
1114 authoritative_copy.title_version = 2;
1115 authoritative_copy.metadata_version = 4;
1116 authoritative_copy.pinned = true;
1117
1118 merge_save_session(&storage, &mut authoritative_copy)
1119 .await
1120 .unwrap();
1121
1122 let after = storage.load_session(session_id).await.unwrap().unwrap();
1123 assert_eq!(after.title, "New Authoritative");
1124 assert_eq!(after.title_version, 2);
1125 assert_eq!(after.metadata_version, 4);
1126 assert!(after.pinned);
1127 }
1128
1129 #[tokio::test]
1130 async fn merge_keeps_runtime_messages_when_disk_only_changed_metadata() {
1131 let (_temp, storage) = make_storage().await;
1132 let session_id = "merge-messages";
1133
1134 let mut on_disk = fresh(session_id);
1135 on_disk.title = "Fresh Title".to_string();
1136 on_disk.title_version = 2;
1137 on_disk.metadata_version = 5;
1138 storage.save_session(&on_disk).await.unwrap();
1139
1140 let mut runtime_copy = fresh(session_id);
1141 runtime_copy.title = "Stale".to_string();
1142 runtime_copy.metadata_version = 0;
1143 runtime_copy.messages = vec![bamboo_domain::session::types::Message {
1144 role: bamboo_domain::session::types::Role::User,
1145 content: "keep me".to_string(),
1146 id: "msg-1".to_string(),
1147 created_at: chrono::Utc::now(),
1148 reasoning: None,
1149 reasoning_signature: None,
1150 content_parts: None,
1151 image_ocr: None,
1152 phase: None,
1153 tool_calls: None,
1154 tool_call_id: None,
1155 tool_success: None,
1156 compressed: false,
1157 compressed_by_event_id: None,
1158 never_compress: false,
1159 compression_level: 0,
1160 metadata: None,
1161 }];
1162
1163 merge_save_session(&storage, &mut runtime_copy)
1164 .await
1165 .unwrap();
1166
1167 let after = storage.load_session(session_id).await.unwrap().unwrap();
1168 assert_eq!(after.title, "Fresh Title");
1169 assert_eq!(after.metadata_version, 5);
1170 assert_eq!(after.messages.len(), 1);
1171 assert_eq!(after.messages[0].content, "keep me");
1172 }
1173
1174 #[tokio::test]
1177 async fn locked_merge_save_runtime_serialises_concurrent_writes() {
1178 let (_temp, storage) = make_storage().await;
1179 let store = Arc::new(LockedSessionStore::new(storage));
1180 let session_id = "lock-serial".to_string();
1181
1182 let base = fresh(&session_id);
1184 store.storage().save_session(&base).await.unwrap();
1185
1186 let store_a = store.clone();
1189 let store_b = store.clone();
1190 let sid_a = session_id.clone();
1191 let sid_b = session_id.clone();
1192
1193 let a = tokio::spawn(async move {
1194 let _guard = store_a.acquire_lock(&sid_a).await;
1195 let mut s = store_a
1196 .storage()
1197 .load_session(&sid_a)
1198 .await
1199 .unwrap()
1200 .unwrap();
1201 s.title = "Writer A".to_string();
1202 s.title_version = s.title_version.saturating_add(1);
1203 s.metadata_version = s.metadata_version.saturating_add(1);
1204 s.updated_at = chrono::Utc::now();
1205 store_a.storage().save_session(&s).await.unwrap();
1206 s.title_version
1207 });
1208
1209 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
1211
1212 let b = tokio::spawn(async move {
1213 let _guard = store_b.acquire_lock(&sid_b).await;
1214 let mut s = store_b
1215 .storage()
1216 .load_session(&sid_b)
1217 .await
1218 .unwrap()
1219 .unwrap();
1220 s.title = "Writer B".to_string();
1221 s.title_version = s.title_version.saturating_add(1);
1222 s.metadata_version = s.metadata_version.saturating_add(1);
1223 s.updated_at = chrono::Utc::now();
1224 store_b.storage().save_session(&s).await.unwrap();
1225 s.title_version
1226 });
1227
1228 let (ver_a, ver_b) = tokio::join!(a, b);
1229 let final_s = store
1230 .storage()
1231 .load_session(&session_id)
1232 .await
1233 .unwrap()
1234 .unwrap();
1235 assert!(
1236 ver_a.unwrap() != ver_b.unwrap(),
1237 "concurrent writers must produce distinct versions"
1238 );
1239 assert_eq!(final_s.metadata_version, 2);
1240 }
1241
1242 #[tokio::test]
1243 async fn commit_metadata_is_plain_save_inside_lock() {
1244 let (_temp, storage) = make_storage().await;
1245 let store = LockedSessionStore::new(storage);
1246 let session_id = "commit-plain";
1247
1248 let mut s = fresh(session_id);
1249 s.title = "Committed".to_string();
1250 s.metadata_version = 1;
1251 s.title_version = 2;
1252
1253 store.commit_metadata(&s).await.unwrap();
1254
1255 let after = store
1256 .storage()
1257 .load_session(session_id)
1258 .await
1259 .unwrap()
1260 .unwrap();
1261 assert_eq!(after.title, "Committed");
1262 assert_eq!(after.metadata_version, 1);
1263 assert_eq!(after.title_version, 2);
1264 }
1265
1266 #[tokio::test]
1269 async fn acquire_lock_self_evicts_when_no_other_holder() {
1270 let (_temp, storage) = make_storage().await;
1271 let store = LockedSessionStore::new(storage);
1272
1273 {
1274 let _guard = store.acquire_lock("solo").await;
1275 assert_eq!(store.locks.len(), 1, "entry present while the lock is held");
1276 }
1277 assert_eq!(
1280 store.locks.len(),
1281 0,
1282 "lock entry must be evicted once released with no other holder"
1283 );
1284 }
1285
1286 #[tokio::test]
1287 async fn acquire_lock_many_distinct_ids_do_not_accumulate() {
1288 let (_temp, storage) = make_storage().await;
1289 let store = LockedSessionStore::new(storage);
1290
1291 for i in 0..100 {
1293 let _guard = store.acquire_lock(&format!("sess-{i}")).await;
1294 }
1295 assert_eq!(
1296 store.locks.len(),
1297 0,
1298 "acquiring locks for many distinct ids must not grow the map"
1299 );
1300 }
1301
1302 #[tokio::test]
1303 async fn acquire_lock_concurrent_waiter_keeps_valid_lock_and_map_drains() {
1304 use std::sync::atomic::{AtomicUsize, Ordering};
1305
1306 let (_temp, storage) = make_storage().await;
1307 let store = Arc::new(LockedSessionStore::new(storage));
1308
1309 let active = Arc::new(AtomicUsize::new(0));
1311 let max_seen = Arc::new(AtomicUsize::new(0));
1312
1313 let mut handles = Vec::new();
1314 for _ in 0..8 {
1315 let store = store.clone();
1316 let active = active.clone();
1317 let max_seen = max_seen.clone();
1318 handles.push(tokio::spawn(async move {
1319 let _guard = store.acquire_lock("contended").await;
1320 let now = active.fetch_add(1, Ordering::SeqCst) + 1;
1321 max_seen.fetch_max(now, Ordering::SeqCst);
1322 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1324 active.fetch_sub(1, Ordering::SeqCst);
1325 }));
1326 }
1327 for h in handles {
1328 h.await.unwrap();
1329 }
1330
1331 assert_eq!(
1336 max_seen.load(Ordering::SeqCst),
1337 1,
1338 "at most one holder of a given session lock at a time"
1339 );
1340 assert_eq!(
1341 store.locks.len(),
1342 0,
1343 "after all holders release, the contended entry must be fully evicted"
1344 );
1345 }
1346}