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 save_runtime_authoritative_flags(
207 &self,
208 session: &mut Session,
209 ) -> std::io::Result<()> {
210 self.merge_save_runtime_inner(session, false).await
211 }
212
213 async fn merge_save_runtime_inner(
214 &self,
215 session: &mut Session,
216 adopt_bypass: bool,
217 ) -> std::io::Result<()> {
218 let _guard = self.acquire_lock(&session.id).await;
219
220 let latest = self.storage.load_session(&session.id).await.ok().flatten();
227
228 let existing_message_count = latest.as_ref().map(|s| s.messages.len());
234 let incoming_message_count = session.messages.len();
235 if existing_message_count.is_some_and(|existing| existing > incoming_message_count) {
236 tracing::warn!(
237 "[{}] merge_save_runtime SHRINK: disk has {:?} messages, saving {} (last_role={:?}, updated_at={}); a stale writer is reverting a concurrent append",
238 session.id,
239 existing_message_count,
240 incoming_message_count,
241 session.messages.last().map(|m| format!("{:?}", m.role)),
242 session.updated_at,
243 );
244 } else {
245 tracing::debug!(
246 "[{}] merge_save_runtime: disk={:?} messages, saving {} (updated_at={})",
247 session.id,
248 existing_message_count,
249 incoming_message_count,
250 session.updated_at,
251 );
252 }
253
254 if let Some(latest) = latest.as_ref() {
255 apply_authoritative_metadata(session, latest);
256 if adopt_bypass {
260 adopt_disk_bypass_permissions(session, latest);
261 }
262 }
263 self.storage.save_session(session).await
264 }
265
266 pub async fn update_runtime_config<F>(
278 &self,
279 session_id: &str,
280 mutate: F,
281 ) -> std::io::Result<Option<Session>>
282 where
283 F: FnOnce(&mut Session),
284 {
285 let _guard = self.acquire_lock(session_id).await;
286 let Some(mut session) = self.storage.load_session(session_id).await? else {
287 return Ok(None);
288 };
289 mutate(&mut session);
290 self.storage.save_session(&session).await?;
291 Ok(Some(session))
292 }
293}
294
295#[async_trait::async_trait]
299impl RuntimeSessionPersistence for LockedSessionStore {
300 async fn save_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
301 self.merge_save_runtime(session).await
302 }
303
304 async fn load_runtime_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
305 self.storage.load_session(session_id).await
306 }
307}
308
309async fn merge_authoritative_metadata_into_stale(
318 storage: &Arc<dyn Storage>,
319 session: &mut Session,
320) {
321 if let Ok(Some(latest)) = storage.load_session(&session.id).await {
322 apply_authoritative_metadata(session, &latest);
323 adopt_disk_bypass_permissions(session, &latest);
324 }
325}
326
327fn adopt_disk_bypass_permissions(session: &mut Session, latest: &Session) {
337 let Some(disk_bypass) = latest
342 .agent_runtime_state
343 .as_ref()
344 .map(|state| state.bypass_permissions)
345 else {
346 return;
347 };
348 match session.agent_runtime_state.as_mut() {
349 Some(state) => state.bypass_permissions = disk_bypass,
350 None if disk_bypass => {
353 session
354 .agent_runtime_state
355 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
356 .bypass_permissions = true;
357 }
358 None => {}
359 }
360}
361
362fn apply_authoritative_metadata(session: &mut Session, latest: &Session) {
368 if latest.metadata_version >= session.metadata_version {
369 session.title = latest.title.clone();
370 session.title_version = latest.title_version;
371 session.pinned = latest.pinned;
372 for key in AUTHORITATIVE_METADATA_KEYS {
373 if let Some(value) = latest.metadata.get(*key) {
374 session.metadata.insert((*key).to_string(), value.clone());
375 } else {
376 session.metadata.remove(*key);
377 }
378 }
379 session.metadata_version = latest.metadata_version;
380 }
381}
382
383pub async fn merge_save_session(
396 storage: &Arc<dyn Storage>,
397 session: &mut Session,
398) -> std::io::Result<()> {
399 merge_authoritative_metadata_into_stale(storage, session).await;
400 storage.save_session(session).await
401}
402
403#[cfg(test)]
406mod tests {
407 use super::*;
408 use crate::v2::SessionStoreV2;
409 use bamboo_domain::session::types::Session;
410
411 async fn make_storage() -> (tempfile::TempDir, Arc<dyn Storage>) {
412 let temp = tempfile::tempdir().unwrap();
413 let storage = SessionStoreV2::new(temp.path().to_path_buf())
414 .await
415 .expect("storage init");
416 (temp, Arc::new(storage) as Arc<dyn Storage>)
417 }
418
419 fn fresh(id: &str) -> Session {
420 Session::new(id.to_string(), "test-model".to_string())
421 }
422
423 #[tokio::test]
426 async fn update_runtime_config_preserves_concurrently_appended_messages() {
427 use bamboo_domain::session::types::Message;
428 use bamboo_domain::ReasoningEffort;
429
430 let (_temp, storage) = make_storage().await;
431 let store = LockedSessionStore::new(storage.clone());
432 let session_id = "cfg-preserve";
433
434 let mut initial = fresh(session_id);
436 initial.add_message(Message::user("hello"));
437 initial.add_message(Message::assistant("hi", None));
438 storage.save_session(&initial).await.unwrap();
439
440 let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
442 after_chat.add_message(Message::user("second question"));
443 storage.save_session(&after_chat).await.unwrap();
444 assert_eq!(after_chat.messages.len(), 3);
445
446 let updated = store
450 .update_runtime_config(session_id, |s| {
451 s.reasoning_effort = Some(ReasoningEffort::Max);
452 })
453 .await
454 .unwrap()
455 .expect("session exists");
456
457 assert_eq!(updated.reasoning_effort, Some(ReasoningEffort::Max));
458 assert_eq!(
459 updated.messages.len(),
460 3,
461 "config patch must not revert a concurrently-appended message"
462 );
463
464 let on_disk = storage.load_session(session_id).await.unwrap().unwrap();
465 assert_eq!(on_disk.messages.len(), 3);
466 assert_eq!(on_disk.reasoning_effort, Some(ReasoningEffort::Max));
467 }
468
469 #[tokio::test]
470 async fn update_runtime_config_returns_none_for_missing_session() {
471 use bamboo_domain::ReasoningEffort;
472
473 let (_temp, storage) = make_storage().await;
474 let store = LockedSessionStore::new(storage);
475 let result = store
476 .update_runtime_config("does-not-exist", |s| {
477 s.reasoning_effort = Some(ReasoningEffort::Low);
478 })
479 .await
480 .unwrap();
481 assert!(result.is_none());
482 }
483
484 #[tokio::test]
485 async fn merge_save_runtime_overwrites_messages_from_stale_snapshot() {
486 use bamboo_domain::session::types::Message;
491
492 let (_temp, storage) = make_storage().await;
493 let store = LockedSessionStore::new(storage.clone());
494 let session_id = "stale-clobber";
495
496 let mut baseline = fresh(session_id);
498 baseline.add_message(Message::user("hello"));
499 storage.save_session(&baseline).await.unwrap();
500 let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
501
502 let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
504 after_chat.add_message(Message::user("second"));
505 storage.save_session(&after_chat).await.unwrap();
506 assert_eq!(
507 storage
508 .load_session(session_id)
509 .await
510 .unwrap()
511 .unwrap()
512 .messages
513 .len(),
514 2
515 );
516
517 store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
519 let after = storage.load_session(session_id).await.unwrap().unwrap();
520 assert_eq!(
521 after.messages.len(),
522 1,
523 "merge_save_runtime clobbers concurrent appends — this is why config patches must use update_runtime_config"
524 );
525 }
526
527 #[tokio::test]
528 async fn merge_save_runtime_preserves_disk_authoritative_metadata_with_single_load() {
529 let (_temp, storage) = make_storage().await;
535 let store = LockedSessionStore::new(storage.clone());
536 let session_id = "runtime-merge-meta";
537
538 let mut baseline = fresh(session_id);
540 baseline.title = "Auto Title".to_string();
541 baseline.metadata_version = 0;
542 storage.save_session(&baseline).await.unwrap();
543
544 let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
546
547 let mut renamed = storage.load_session(session_id).await.unwrap().unwrap();
549 renamed.title = "User Renamed".to_string();
550 renamed.title_version = 1;
551 renamed.pinned = true;
552 renamed.metadata_version = 1;
553 store.commit_metadata(&renamed).await.unwrap();
554
555 stale_snapshot.title = "Auto Title".to_string();
557 store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
558
559 let after = storage.load_session(session_id).await.unwrap().unwrap();
560 assert_eq!(after.title, "User Renamed");
561 assert!(after.pinned);
562 assert_eq!(after.metadata_version, 1);
563 assert_eq!(stale_snapshot.title, "User Renamed");
565 assert_eq!(stale_snapshot.metadata_version, 1);
566 }
567
568 #[tokio::test]
569 async fn merge_save_runtime_preserves_durable_workflow_run_index_from_stale_runner() {
570 let (_temp, storage) = make_storage().await;
571 let store = LockedSessionStore::new(storage.clone());
572 let session_id = "runtime-workflow-run-index";
573
574 let baseline = fresh(session_id);
575 storage.save_session(&baseline).await.unwrap();
576 let mut stale_runner = storage.load_session(session_id).await.unwrap().unwrap();
577
578 store
579 .update_runtime_config(session_id, |session| {
580 session.metadata.insert(
581 "workflow.run_ids.v1".to_string(),
582 r#"["http-started-run"]"#.to_string(),
583 );
584 })
585 .await
586 .unwrap()
587 .expect("session exists");
588
589 store.merge_save_runtime(&mut stale_runner).await.unwrap();
590
591 assert_eq!(
592 stale_runner
593 .metadata
594 .get("workflow.run_ids.v1")
595 .map(String::as_str),
596 Some(r#"["http-started-run"]"#)
597 );
598 let durable = storage.load_session(session_id).await.unwrap().unwrap();
599 assert_eq!(
600 durable
601 .metadata
602 .get("workflow.run_ids.v1")
603 .map(String::as_str),
604 Some(r#"["http-started-run"]"#)
605 );
606 }
607
608 #[tokio::test]
612 async fn merge_save_runtime_adopts_disk_bypass_permissions() {
613 use bamboo_domain::AgentRuntimeState;
614
615 let (_temp, storage) = make_storage().await;
616 let store = LockedSessionStore::new(storage.clone());
617 let session_id = "runtime-bypass";
618
619 let baseline = fresh(session_id);
621 storage.save_session(&baseline).await.unwrap();
622
623 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
625 loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
626
627 store
629 .update_runtime_config(session_id, |s| {
630 s.agent_runtime_state
631 .get_or_insert_with(AgentRuntimeState::default)
632 .bypass_permissions = true;
633 })
634 .await
635 .unwrap()
636 .expect("session exists");
637
638 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
641
642 let after = storage.load_session(session_id).await.unwrap().unwrap();
643 assert!(
644 after
645 .agent_runtime_state
646 .as_ref()
647 .is_some_and(|s| s.bypass_permissions),
648 "disk bypass=ON must survive a stale runtime save (#540)"
649 );
650 assert!(loop_snapshot
652 .agent_runtime_state
653 .as_ref()
654 .is_some_and(|s| s.bypass_permissions));
655 }
656
657 #[tokio::test]
660 async fn merge_save_runtime_adopts_disk_bypass_off() {
661 use bamboo_domain::AgentRuntimeState;
662
663 let (_temp, storage) = make_storage().await;
664 let store = LockedSessionStore::new(storage.clone());
665 let session_id = "runtime-bypass-off";
666
667 let mut baseline = fresh(session_id);
669 let mut on_state = AgentRuntimeState::default();
670 on_state.bypass_permissions = true;
671 baseline.agent_runtime_state = Some(on_state);
672 storage.save_session(&baseline).await.unwrap();
673
674 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
676
677 store
679 .update_runtime_config(session_id, |s| {
680 s.agent_runtime_state
681 .get_or_insert_with(AgentRuntimeState::default)
682 .bypass_permissions = false;
683 })
684 .await
685 .unwrap()
686 .expect("session exists");
687
688 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
689
690 let after = storage.load_session(session_id).await.unwrap().unwrap();
691 assert!(
692 !after
693 .agent_runtime_state
694 .as_ref()
695 .is_some_and(|s| s.bypass_permissions),
696 "disk bypass=OFF must survive a stale runtime save (#540)"
697 );
698 }
699
700 #[tokio::test]
703 async fn save_runtime_authoritative_flags_persists_in_memory_bypass() {
704 use bamboo_domain::AgentRuntimeState;
705
706 let (_temp, storage) = make_storage().await;
707 let store = LockedSessionStore::new(storage.clone());
708 let session_id = "child-reseed";
709
710 let mut baseline = fresh(session_id);
712 let mut on_state = AgentRuntimeState::default();
713 on_state.bypass_permissions = true;
714 baseline.agent_runtime_state = Some(on_state);
715 storage.save_session(&baseline).await.unwrap();
716
717 let mut child = storage.load_session(session_id).await.unwrap().unwrap();
720 child
721 .agent_runtime_state
722 .get_or_insert_with(AgentRuntimeState::default)
723 .bypass_permissions = false;
724
725 store
727 .save_runtime_authoritative_flags(&mut child)
728 .await
729 .unwrap();
730
731 let after = storage.load_session(session_id).await.unwrap().unwrap();
732 assert!(
733 !after
734 .agent_runtime_state
735 .as_ref()
736 .is_some_and(|s| s.bypass_permissions),
737 "authoritative re-seed of bypass=OFF must persist, not be reverted (#540/#74)"
738 );
739 }
740
741 #[tokio::test]
743 async fn merge_save_runtime_leaves_bypass_when_disk_has_no_runtime_state() {
744 use bamboo_domain::AgentRuntimeState;
745
746 let (_temp, storage) = make_storage().await;
747 let store = LockedSessionStore::new(storage.clone());
748 let session_id = "no-runtime-state";
749
750 let baseline = fresh(session_id);
752 assert!(baseline.agent_runtime_state.is_none());
753 storage.save_session(&baseline).await.unwrap();
754
755 let mut running = storage.load_session(session_id).await.unwrap().unwrap();
757 let mut on_state = AgentRuntimeState::default();
758 on_state.bypass_permissions = true;
759 running.agent_runtime_state = Some(on_state);
760
761 store.merge_save_runtime(&mut running).await.unwrap();
762
763 assert!(
764 running
765 .agent_runtime_state
766 .as_ref()
767 .is_some_and(|s| s.bypass_permissions),
768 "a runtime-state-less disk copy must not force bypass OFF (#540)"
769 );
770 }
771
772 #[tokio::test]
775 async fn merge_preserves_disk_title_when_versions_equal() {
776 let (_temp, storage) = make_storage().await;
777 let session_id = "merge-equal";
778
779 let mut on_disk = fresh(session_id);
780 on_disk.title = "User Set This".to_string();
781 on_disk.title_version = 0;
782 on_disk.metadata_version = 0;
783 storage.save_session(&on_disk).await.unwrap();
784
785 let mut runtime_copy = fresh(session_id);
786 runtime_copy.title = "Stale Default".to_string();
787 runtime_copy.title_version = 0;
788 runtime_copy.metadata_version = 0;
789 runtime_copy.messages = vec![];
790
791 merge_save_session(&storage, &mut runtime_copy)
792 .await
793 .unwrap();
794
795 let after = storage.load_session(session_id).await.unwrap().unwrap();
796 assert_eq!(after.title, "User Set This");
797 assert_eq!(after.title_version, 0);
798 assert_eq!(runtime_copy.title, "User Set This");
799 }
800
801 #[tokio::test]
802 async fn merge_preserves_disk_when_disk_version_higher() {
803 let (_temp, storage) = make_storage().await;
804 let session_id = "merge-higher";
805
806 let mut on_disk = fresh(session_id);
807 on_disk.title = "User Title v3".to_string();
808 on_disk.title_version = 3;
809 on_disk.metadata_version = 5;
810 storage.save_session(&on_disk).await.unwrap();
811
812 let mut runtime_copy = fresh(session_id);
813 runtime_copy.title = "Stale".to_string();
814 runtime_copy.title_version = 1;
815 runtime_copy.metadata_version = 0;
816
817 merge_save_session(&storage, &mut runtime_copy)
818 .await
819 .unwrap();
820
821 let after = storage.load_session(session_id).await.unwrap().unwrap();
822 assert_eq!(after.title, "User Title v3");
823 assert_eq!(after.title_version, 3);
824 assert_eq!(after.metadata_version, 5);
825 }
826
827 #[tokio::test]
828 async fn merge_now_preserves_disk_pinned_in_metadata_group() {
829 let (_temp, storage) = make_storage().await;
830 let session_id = "pinned-merge";
831
832 let mut on_disk = fresh(session_id);
833 on_disk.pinned = true;
834 on_disk.metadata_version = 2;
835 storage.save_session(&on_disk).await.unwrap();
836
837 let mut runtime_copy = fresh(session_id);
838 runtime_copy.pinned = false;
839 runtime_copy.metadata_version = 0;
840
841 merge_save_session(&storage, &mut runtime_copy)
842 .await
843 .unwrap();
844
845 let after = storage.load_session(session_id).await.unwrap().unwrap();
846 assert!(
847 after.pinned,
848 "disk pinned=true should win over runtime false"
849 );
850 assert_eq!(after.metadata_version, 2);
851 }
852
853 #[tokio::test]
854 async fn merge_keeps_in_memory_when_session_version_higher() {
855 let (_temp, storage) = make_storage().await;
856 let session_id = "merge-bumped";
857
858 let mut on_disk = fresh(session_id);
859 on_disk.title = "Old".to_string();
860 on_disk.title_version = 1;
861 on_disk.metadata_version = 3;
862 storage.save_session(&on_disk).await.unwrap();
863
864 let mut authoritative_copy = fresh(session_id);
865 authoritative_copy.title = "New Authoritative".to_string();
866 authoritative_copy.title_version = 2;
867 authoritative_copy.metadata_version = 4;
868 authoritative_copy.pinned = true;
869
870 merge_save_session(&storage, &mut authoritative_copy)
871 .await
872 .unwrap();
873
874 let after = storage.load_session(session_id).await.unwrap().unwrap();
875 assert_eq!(after.title, "New Authoritative");
876 assert_eq!(after.title_version, 2);
877 assert_eq!(after.metadata_version, 4);
878 assert!(after.pinned);
879 }
880
881 #[tokio::test]
882 async fn merge_keeps_runtime_messages_when_disk_only_changed_metadata() {
883 let (_temp, storage) = make_storage().await;
884 let session_id = "merge-messages";
885
886 let mut on_disk = fresh(session_id);
887 on_disk.title = "Fresh Title".to_string();
888 on_disk.title_version = 2;
889 on_disk.metadata_version = 5;
890 storage.save_session(&on_disk).await.unwrap();
891
892 let mut runtime_copy = fresh(session_id);
893 runtime_copy.title = "Stale".to_string();
894 runtime_copy.metadata_version = 0;
895 runtime_copy.messages = vec![bamboo_domain::session::types::Message {
896 role: bamboo_domain::session::types::Role::User,
897 content: "keep me".to_string(),
898 id: "msg-1".to_string(),
899 created_at: chrono::Utc::now(),
900 reasoning: None,
901 reasoning_signature: None,
902 content_parts: None,
903 image_ocr: None,
904 phase: None,
905 tool_calls: None,
906 tool_call_id: None,
907 tool_success: None,
908 compressed: false,
909 compressed_by_event_id: None,
910 never_compress: false,
911 compression_level: 0,
912 metadata: None,
913 }];
914
915 merge_save_session(&storage, &mut runtime_copy)
916 .await
917 .unwrap();
918
919 let after = storage.load_session(session_id).await.unwrap().unwrap();
920 assert_eq!(after.title, "Fresh Title");
921 assert_eq!(after.metadata_version, 5);
922 assert_eq!(after.messages.len(), 1);
923 assert_eq!(after.messages[0].content, "keep me");
924 }
925
926 #[tokio::test]
929 async fn locked_merge_save_runtime_serialises_concurrent_writes() {
930 let (_temp, storage) = make_storage().await;
931 let store = Arc::new(LockedSessionStore::new(storage));
932 let session_id = "lock-serial".to_string();
933
934 let base = fresh(&session_id);
936 store.storage().save_session(&base).await.unwrap();
937
938 let store_a = store.clone();
941 let store_b = store.clone();
942 let sid_a = session_id.clone();
943 let sid_b = session_id.clone();
944
945 let a = tokio::spawn(async move {
946 let _guard = store_a.acquire_lock(&sid_a).await;
947 let mut s = store_a
948 .storage()
949 .load_session(&sid_a)
950 .await
951 .unwrap()
952 .unwrap();
953 s.title = "Writer A".to_string();
954 s.title_version = s.title_version.saturating_add(1);
955 s.metadata_version = s.metadata_version.saturating_add(1);
956 s.updated_at = chrono::Utc::now();
957 store_a.storage().save_session(&s).await.unwrap();
958 s.title_version
959 });
960
961 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
963
964 let b = tokio::spawn(async move {
965 let _guard = store_b.acquire_lock(&sid_b).await;
966 let mut s = store_b
967 .storage()
968 .load_session(&sid_b)
969 .await
970 .unwrap()
971 .unwrap();
972 s.title = "Writer B".to_string();
973 s.title_version = s.title_version.saturating_add(1);
974 s.metadata_version = s.metadata_version.saturating_add(1);
975 s.updated_at = chrono::Utc::now();
976 store_b.storage().save_session(&s).await.unwrap();
977 s.title_version
978 });
979
980 let (ver_a, ver_b) = tokio::join!(a, b);
981 let final_s = store
982 .storage()
983 .load_session(&session_id)
984 .await
985 .unwrap()
986 .unwrap();
987 assert!(
988 ver_a.unwrap() != ver_b.unwrap(),
989 "concurrent writers must produce distinct versions"
990 );
991 assert_eq!(final_s.metadata_version, 2);
992 }
993
994 #[tokio::test]
995 async fn commit_metadata_is_plain_save_inside_lock() {
996 let (_temp, storage) = make_storage().await;
997 let store = LockedSessionStore::new(storage);
998 let session_id = "commit-plain";
999
1000 let mut s = fresh(session_id);
1001 s.title = "Committed".to_string();
1002 s.metadata_version = 1;
1003 s.title_version = 2;
1004
1005 store.commit_metadata(&s).await.unwrap();
1006
1007 let after = store
1008 .storage()
1009 .load_session(session_id)
1010 .await
1011 .unwrap()
1012 .unwrap();
1013 assert_eq!(after.title, "Committed");
1014 assert_eq!(after.metadata_version, 1);
1015 assert_eq!(after.title_version, 2);
1016 }
1017
1018 #[tokio::test]
1021 async fn acquire_lock_self_evicts_when_no_other_holder() {
1022 let (_temp, storage) = make_storage().await;
1023 let store = LockedSessionStore::new(storage);
1024
1025 {
1026 let _guard = store.acquire_lock("solo").await;
1027 assert_eq!(store.locks.len(), 1, "entry present while the lock is held");
1028 }
1029 assert_eq!(
1032 store.locks.len(),
1033 0,
1034 "lock entry must be evicted once released with no other holder"
1035 );
1036 }
1037
1038 #[tokio::test]
1039 async fn acquire_lock_many_distinct_ids_do_not_accumulate() {
1040 let (_temp, storage) = make_storage().await;
1041 let store = LockedSessionStore::new(storage);
1042
1043 for i in 0..100 {
1045 let _guard = store.acquire_lock(&format!("sess-{i}")).await;
1046 }
1047 assert_eq!(
1048 store.locks.len(),
1049 0,
1050 "acquiring locks for many distinct ids must not grow the map"
1051 );
1052 }
1053
1054 #[tokio::test]
1055 async fn acquire_lock_concurrent_waiter_keeps_valid_lock_and_map_drains() {
1056 use std::sync::atomic::{AtomicUsize, Ordering};
1057
1058 let (_temp, storage) = make_storage().await;
1059 let store = Arc::new(LockedSessionStore::new(storage));
1060
1061 let active = Arc::new(AtomicUsize::new(0));
1063 let max_seen = Arc::new(AtomicUsize::new(0));
1064
1065 let mut handles = Vec::new();
1066 for _ in 0..8 {
1067 let store = store.clone();
1068 let active = active.clone();
1069 let max_seen = max_seen.clone();
1070 handles.push(tokio::spawn(async move {
1071 let _guard = store.acquire_lock("contended").await;
1072 let now = active.fetch_add(1, Ordering::SeqCst) + 1;
1073 max_seen.fetch_max(now, Ordering::SeqCst);
1074 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1076 active.fetch_sub(1, Ordering::SeqCst);
1077 }));
1078 }
1079 for h in handles {
1080 h.await.unwrap();
1081 }
1082
1083 assert_eq!(
1088 max_seen.load(Ordering::SeqCst),
1089 1,
1090 "at most one holder of a given session lock at a time"
1091 );
1092 assert_eq!(
1093 store.locks.len(),
1094 0,
1095 "after all holders release, the contended entry must be fully evicted"
1096 );
1097 }
1098}