1use std::sync::Arc;
37
38use bamboo_domain::session::types::Session;
39use bamboo_domain::storage::Storage;
40use bamboo_domain::{
41 latest_response_occurrence, PermissionAuditSeed, PermissionAuditSnapshot, ResponseOccurrence,
42 RuntimeSessionPersistence, CONSUMED_CLARIFICATION_IDS_KEY, CONSUMED_RESPONSE_OCCURRENCES_KEY,
43};
44use dashmap::DashMap;
45use tokio::sync::{Mutex, OwnedMutexGuard};
46
47const AUTHORITATIVE_METADATA_KEYS: &[&str] = &["gold_config", "workflow.run_ids.v1"];
48const RESPONSE_CONTROL_METADATA_KEYS: &[&str] = &[
49 CONSUMED_CLARIFICATION_IDS_KEY,
50 CONSUMED_RESPONSE_OCCURRENCES_KEY,
51 "runtime.suspend_reason",
52 "clarification_resume_pending",
53 "conclusion_with_options_resume_pending",
54 "execute.startup_handoff_at",
55 "permission.reexecute_tool_call_id",
56 "permission.reexecute_request_generation",
57 "retry_resume_pending",
58 "retry_resume_reason",
59 "provider_name",
60];
61const TASK_CONTROL_PLANE_CONFLICT_PREFIX: &str = "Task control-plane changed while saving session ";
62const MAX_TASK_CONTROL_PLANE_REBASE_RETRIES: usize = 3;
63
64fn adopt_durable_consumed_clarification(session: &mut Session, durable: &Session) -> bool {
73 let Some(incoming_tool_call_id) = session
74 .pending_question
75 .as_ref()
76 .map(|pending| pending.tool_call_id.clone())
77 else {
78 return false;
79 };
80 let occurrence_ledger = durable
81 .metadata
82 .get(CONSUMED_RESPONSE_OCCURRENCES_KEY)
83 .map(|value| serde_json::from_str::<Vec<ResponseOccurrence>>(value));
84 let was_consumed = match occurrence_ledger {
85 Some(Ok(consumed)) => latest_response_occurrence(session, &incoming_tool_call_id)
86 .is_some_and(|incoming| consumed.iter().any(|entry| entry == &incoming)),
87 Some(Err(_)) => false,
90 None => {
91 let legacy_consumed = durable
92 .metadata
93 .get(CONSUMED_CLARIFICATION_IDS_KEY)
94 .and_then(|value| serde_json::from_str::<Vec<String>>(value).ok())
95 .unwrap_or_default()
96 .iter()
97 .any(|tool_call_id| tool_call_id == &incoming_tool_call_id);
98 legacy_consumed
102 && latest_response_occurrence(durable, &incoming_tool_call_id).is_some_and(
103 |durable_occurrence| {
104 latest_response_occurrence(session, &incoming_tool_call_id)
105 .is_some_and(|incoming| incoming == durable_occurrence)
106 },
107 )
108 }
109 };
110 if !was_consumed {
111 return false;
112 }
113
114 bamboo_domain::append_missing_runtime_messages(session, durable);
117 session
118 .pending_question
119 .clone_from(&durable.pending_question);
120 for key in RESPONSE_CONTROL_METADATA_KEYS {
121 if let Some(value) = durable.metadata.get(*key) {
122 session.metadata.insert((*key).to_string(), value.clone());
123 } else {
124 session.metadata.remove(*key);
125 }
126 }
127 session.model.clone_from(&durable.model);
128 session.model_ref.clone_from(&durable.model_ref);
129 session.reasoning_effort = durable.reasoning_effort;
130 session
131 .agent_runtime_state
132 .clone_from(&durable.agent_runtime_state);
133 if let Some(runtime_metadata) = session.runtime_metadata.as_mut() {
134 runtime_metadata.provider_name = durable
135 .runtime_metadata
136 .as_ref()
137 .and_then(|metadata| metadata.provider_name.clone());
138 } else if durable
139 .runtime_metadata
140 .as_ref()
141 .is_some_and(|metadata| metadata.provider_name.is_some())
142 {
143 session.runtime_metadata = durable.runtime_metadata.as_ref().map(|metadata| {
144 let mut response_metadata = bamboo_domain::SessionRuntimeMetadata::default();
145 response_metadata
146 .provider_name
147 .clone_from(&metadata.provider_name);
148 response_metadata
149 });
150 }
151 true
152}
153
154fn is_task_control_plane_save_conflict(error: &std::io::Error) -> bool {
155 error.kind() == std::io::ErrorKind::WouldBlock
156 && error
157 .to_string()
158 .starts_with(TASK_CONTROL_PLANE_CONFLICT_PREFIX)
159}
160
161fn adopt_durable_task_control_plane(session: &mut Session, durable: &Session) {
162 session.task_list = durable.task_list.clone();
163 session
164 .metadata
165 .remove(bamboo_domain::session::runtime_metadata::keys::TASK_LIST_VERSION);
166 if let Some(runtime_metadata) = session.runtime_metadata.as_mut() {
167 runtime_metadata.task_list_version = None;
168 }
169 if session
170 .runtime_metadata
171 .as_ref()
172 .is_some_and(bamboo_domain::session::SessionRuntimeMetadata::is_empty)
173 {
174 session.runtime_metadata = None;
175 }
176 if let Some(version) = durable.task_list_version_meta() {
177 session.set_task_list_version_meta(version);
178 }
179}
180
181fn task_list_snapshot_matches(
182 session: &Session,
183 expected_task_list: &bamboo_domain::TaskList,
184) -> std::io::Result<bool> {
185 Ok(serde_json::to_value(&session.task_list)
186 .map_err(|error| std::io::Error::other(error.to_string()))?
187 == serde_json::to_value(Some(expected_task_list))
188 .map_err(|error| std::io::Error::other(error.to_string()))?)
189}
190
191fn unconditional_task_patch_would_regress(
192 durable: &Session,
193 incoming_task_list: &bamboo_domain::TaskList,
194 incoming_version: &str,
195) -> std::io::Result<bool> {
196 let same_list = task_list_snapshot_matches(durable, incoming_task_list)?;
197 let Some(durable_version) = durable.task_list_version_meta() else {
198 return Ok(false);
199 };
200 match (
201 incoming_version.parse::<u64>(),
202 durable_version.parse::<u64>(),
203 ) {
204 (Ok(incoming), Ok(durable)) => {
205 Ok(incoming < durable || (incoming == durable && !same_list))
206 }
207 _ => Ok(incoming_version != durable_version || !same_list),
208 }
209}
210
211pub struct LockedSessionStore {
219 storage: Arc<dyn Storage>,
220 locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
221 task_pair_transaction_lock: Arc<Mutex<()>>,
225}
226
227pub struct SessionLockGuard {
250 guard: Option<OwnedMutexGuard<()>>,
252 locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
253 session_id: String,
254}
255
256impl Drop for SessionLockGuard {
257 fn drop(&mut self) {
258 self.guard.take();
261 self.locks
262 .remove_if(&self.session_id, |_, arc| Arc::strong_count(arc) == 1);
263 }
264}
265
266impl LockedSessionStore {
267 pub fn new(storage: Arc<dyn Storage>) -> Self {
269 Self {
270 storage,
271 locks: Arc::new(DashMap::new()),
272 task_pair_transaction_lock: Arc::new(Mutex::new(())),
273 }
274 }
275
276 pub fn storage(&self) -> &Arc<dyn Storage> {
278 &self.storage
279 }
280
281 pub async fn acquire_lock(&self, session_id: &str) -> SessionLockGuard {
293 let lock = self
298 .locks
299 .entry(session_id.to_string())
300 .or_insert_with(|| Arc::new(Mutex::new(())))
301 .clone();
302 let guard = lock.lock_owned().await;
303 SessionLockGuard {
304 guard: Some(guard),
305 locks: self.locks.clone(),
306 session_id: session_id.to_string(),
307 }
308 }
309
310 async fn save_session_rebasing_task_conflicts(
315 &self,
316 session: &mut Session,
317 ) -> std::io::Result<()> {
318 for attempt in 0..=MAX_TASK_CONTROL_PLANE_REBASE_RETRIES {
319 match self.storage.save_session(session).await {
320 Ok(()) => return Ok(()),
321 Err(error)
322 if is_task_control_plane_save_conflict(&error)
323 && attempt < MAX_TASK_CONTROL_PLANE_REBASE_RETRIES =>
324 {
325 let Some(durable) =
326 self.storage.load_runtime_control_plane(&session.id).await?
327 else {
328 return Err(error);
329 };
330 adopt_durable_task_control_plane(session, &durable);
331 }
332 Err(error) => {
333 if is_task_control_plane_save_conflict(&error) {
334 if let Some(durable) =
335 self.storage.load_runtime_control_plane(&session.id).await?
336 {
337 adopt_durable_task_control_plane(session, &durable);
338 }
339 }
340 return Err(error);
341 }
342 }
343 }
344 unreachable!("bounded Task conflict retry loop always returns")
345 }
346
347 async fn save_runtime_state_rebasing_task_conflicts(
350 &self,
351 session: &mut Session,
352 ) -> std::io::Result<()> {
353 for attempt in 0..=MAX_TASK_CONTROL_PLANE_REBASE_RETRIES {
354 match self.storage.save_runtime_state(session).await {
355 Ok(()) => return Ok(()),
356 Err(error)
357 if is_task_control_plane_save_conflict(&error)
358 && attempt < MAX_TASK_CONTROL_PLANE_REBASE_RETRIES =>
359 {
360 let Some(durable) =
361 self.storage.load_runtime_control_plane(&session.id).await?
362 else {
363 return Err(error);
364 };
365 adopt_durable_task_control_plane(session, &durable);
366 }
367 Err(error) => {
368 if is_task_control_plane_save_conflict(&error) {
369 if let Some(durable) =
370 self.storage.load_runtime_control_plane(&session.id).await?
371 {
372 adopt_durable_task_control_plane(session, &durable);
373 }
374 }
375 return Err(error);
376 }
377 }
378 }
379 unreachable!("bounded Task conflict retry loop always returns")
380 }
381
382 pub async fn save_runtime_only(&self, session: &mut Session) -> std::io::Result<()> {
400 self.save_runtime_only_and_publish(session, |_| {}).await
401 }
402
403 pub async fn save_runtime_only_and_publish<F>(
413 &self,
414 session: &mut Session,
415 publish: F,
416 ) -> std::io::Result<()>
417 where
418 F: FnOnce(&Session) + Send,
419 {
420 let _guard = self.acquire_lock(&session.id).await;
421 if let Some(latest) = self.storage.load_runtime_control_plane(&session.id).await? {
422 apply_authoritative_metadata(session, &latest);
423 adopt_fresher_disk_permission_posture(session, &latest);
426 adopt_durable_model_context_state(session, &latest);
431 }
432 let result = self
433 .save_runtime_state_rebasing_task_conflicts(session)
434 .await;
435 publish(session);
436 result
437 }
438
439 pub async fn update_task_list_control_plane_and_publish<F>(
446 &self,
447 session_id: &str,
448 task_list: &bamboo_domain::TaskList,
449 version: &str,
450 publish: F,
451 ) -> std::io::Result<bool>
452 where
453 F: FnOnce(&Session) + Send,
454 {
455 let _guard = self.acquire_lock(session_id).await;
456 let Some(mut latest) = self.storage.load_runtime_control_plane(session_id).await? else {
457 return Ok(false);
458 };
459 if unconditional_task_patch_would_regress(&latest, task_list, version)? {
460 return Err(std::io::Error::new(
461 std::io::ErrorKind::WouldBlock,
462 format!("Task control-plane changed while patching session {session_id}"),
463 ));
464 }
465 let original = latest.clone();
466 latest.task_list = Some(task_list.clone());
467 latest.set_task_list_version_meta(version.to_string());
468 if !self
469 .storage
470 .save_task_control_plane_if_matches(&original, &latest)
471 .await?
472 {
473 return Err(std::io::Error::new(
474 std::io::ErrorKind::WouldBlock,
475 format!("Task control-plane changed while patching session {session_id}"),
476 ));
477 }
478 publish(&latest);
479 Ok(true)
480 }
481
482 pub async fn update_task_list_control_plane_if_version_and_publish<F>(
486 &self,
487 session_id: &str,
488 expected_version: &str,
489 expected_task_list: &bamboo_domain::TaskList,
490 task_list: &bamboo_domain::TaskList,
491 version: &str,
492 publish: F,
493 ) -> std::io::Result<bool>
494 where
495 F: FnOnce(&Session) + Send,
496 {
497 let _guard = self.acquire_lock(session_id).await;
498 let Some(mut latest) = self.storage.load_runtime_control_plane(session_id).await? else {
499 return Ok(false);
500 };
501 if latest.task_list_version_meta().as_deref() != Some(expected_version)
502 || !task_list_snapshot_matches(&latest, expected_task_list)?
503 {
504 return Ok(false);
505 }
506 let original = latest.clone();
507 latest.task_list = Some(task_list.clone());
508 latest.set_task_list_version_meta(version.to_string());
509 if !self
510 .storage
511 .save_task_control_plane_if_matches(&original, &latest)
512 .await?
513 {
514 return Ok(false);
515 }
516 publish(&latest);
517 Ok(true)
518 }
519
520 pub async fn update_task_list_control_planes_if_version_and_publish<F>(
526 &self,
527 session_id: &str,
528 shared_session_id: &str,
529 expected_version: &str,
530 expected_task_list: &bamboo_domain::TaskList,
531 task_list: &bamboo_domain::TaskList,
532 version: &str,
533 publish: F,
534 ) -> std::io::Result<bool>
535 where
536 F: FnOnce(&Session, &Session) + Send,
537 {
538 if session_id == shared_session_id {
539 return self
540 .update_task_list_control_plane_if_version_and_publish(
541 session_id,
542 expected_version,
543 expected_task_list,
544 task_list,
545 version,
546 |session| publish(session, session),
547 )
548 .await;
549 }
550
551 let (first_id, second_id) = if session_id < shared_session_id {
552 (session_id, shared_session_id)
553 } else {
554 (shared_session_id, session_id)
555 };
556 let _transaction_guard = self.task_pair_transaction_lock.lock().await;
557 let _first_guard = self.acquire_lock(first_id).await;
558 let _second_guard = self.acquire_lock(second_id).await;
559
560 self.storage
564 .recover_task_control_plane_transaction(first_id, second_id)
565 .await?;
566
567 let Some(mut local) = self.storage.load_runtime_control_plane(session_id).await? else {
568 return Ok(false);
569 };
570 let Some(mut shared) = self
571 .storage
572 .load_runtime_control_plane(shared_session_id)
573 .await?
574 else {
575 return Ok(false);
576 };
577 if local.task_list_version_meta().as_deref() != Some(expected_version)
578 || shared.task_list_version_meta().as_deref() != Some(expected_version)
579 || !task_list_snapshot_matches(&local, expected_task_list)?
580 || !task_list_snapshot_matches(&shared, expected_task_list)?
581 {
582 return Ok(false);
583 }
584
585 let local_original = local.clone();
586 let shared_original = shared.clone();
587 local.task_list = Some(task_list.clone());
591 local.set_task_list_version_meta(version.to_string());
592 shared.task_list = Some(task_list.clone());
593 shared.set_task_list_version_meta(version.to_string());
594 let (first_original, first_updated, second_original, second_updated) =
595 if session_id < shared_session_id {
596 (&local_original, &local, &shared_original, &shared)
597 } else {
598 (&shared_original, &shared, &local_original, &local)
599 };
600 let committed = self
601 .storage
602 .save_task_control_planes_atomically(
603 first_original,
604 first_updated,
605 second_original,
606 second_updated,
607 )
608 .await?;
609 if !committed {
610 return Ok(false);
611 }
612 publish(&local, &shared);
613 Ok(true)
614 }
615
616 pub async fn commit_metadata(&self, session: &Session) -> std::io::Result<()> {
626 let _guard = self.acquire_lock(&session.id).await;
627 let mut committed = session.clone();
628 if let Some(latest) = self.storage.load_runtime_control_plane(&session.id).await? {
629 adopt_durable_model_context_state(&mut committed, &latest);
630 }
631 self.save_session_rebasing_task_conflicts(&mut committed)
632 .await
633 }
634
635 pub async fn merge_save_runtime(&self, session: &mut Session) -> std::io::Result<()> {
652 self.merge_save_runtime_and_publish(session, |_, _| {})
653 .await
654 }
655
656 pub async fn merge_save_runtime_and_publish<F>(
663 &self,
664 session: &mut Session,
665 publish: F,
666 ) -> std::io::Result<()>
667 where
668 F: FnOnce(&Session, bool) + Send,
669 {
670 self.merge_save_runtime_inner_and_publish(session, true, publish)
671 .await
672 }
673
674 pub async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
683 self.checkpoint_runtime_session_and_publish(session, |_, _| {})
684 .await
685 }
686
687 pub async fn checkpoint_runtime_session_and_publish<F>(
694 &self,
695 session: &mut Session,
696 publish: F,
697 ) -> std::io::Result<()>
698 where
699 F: FnOnce(&Session, bool) + Send,
700 {
701 let _guard = self.acquire_lock(&session.id).await;
702 let latest = self.storage.load_session(&session.id).await?;
703
704 if let Some(latest) = latest.as_ref() {
705 ensure_model_context_checkpoint_is_current(session, latest)?;
706 let incoming_count = session.messages.len();
707 let durable_count = latest.messages.len();
708 let appended = bamboo_domain::append_missing_runtime_messages(session, latest);
709 bamboo_domain::merge_session_inbox_admission(session, latest);
710 let adopted_response = adopt_durable_consumed_clarification(session, latest);
711 tracing::debug!(
712 "[{}] append-safe runtime checkpoint: durable={}, incoming={}, appended={}, adopted_response={}, saved={}",
713 session.id,
714 durable_count,
715 incoming_count,
716 appended,
717 adopted_response,
718 session.messages.len(),
719 );
720 apply_authoritative_metadata(session, latest);
721 adopt_fresher_disk_permission_posture(session, latest);
722 }
723
724 let result = self.save_session_rebasing_task_conflicts(session).await;
725 publish(session, result.is_ok());
726 result
727 }
728
729 pub async fn save_runtime_authoritative_flags(
738 &self,
739 session: &mut Session,
740 ) -> std::io::Result<()> {
741 self.merge_save_runtime_inner_and_publish(session, false, |_, _| {})
742 .await
743 }
744
745 async fn merge_save_runtime_inner_and_publish<F>(
746 &self,
747 session: &mut Session,
748 adopt_bypass: bool,
749 publish: F,
750 ) -> std::io::Result<()>
751 where
752 F: FnOnce(&Session, bool) + Send,
753 {
754 let _guard = self.acquire_lock(&session.id).await;
755
756 let latest = self.storage.load_session(&session.id).await?;
763
764 let existing_message_count = latest.as_ref().map(|s| s.messages.len());
770 let incoming_message_count = session.messages.len();
771 if existing_message_count.is_some_and(|existing| existing > incoming_message_count) {
772 tracing::warn!(
773 "[{}] merge_save_runtime SHRINK: disk has {:?} messages, saving {} (last_role={:?}, updated_at={}); a stale writer is reverting a concurrent append",
774 session.id,
775 existing_message_count,
776 incoming_message_count,
777 session.messages.last().map(|m| format!("{:?}", m.role)),
778 session.updated_at,
779 );
780 } else {
781 tracing::debug!(
782 "[{}] merge_save_runtime: disk={:?} messages, saving {} (updated_at={})",
783 session.id,
784 existing_message_count,
785 incoming_message_count,
786 session.updated_at,
787 );
788 }
789
790 if let Some(latest) = latest.as_ref() {
791 adopt_durable_consumed_clarification(session, latest);
792 apply_authoritative_metadata(session, latest);
793 let restored = bamboo_domain::restore_missing_admitted_inbox_messages(session, latest);
794 if restored > 0 {
795 tracing::warn!(
796 session_id = %session.id,
797 restored,
798 "restored durable SessionInbox transcript messages into stale runtime save"
799 );
800 }
801 bamboo_domain::merge_session_inbox_admission(session, latest);
802 if adopt_bypass {
807 adopt_fresher_disk_permission_posture(session, latest);
808 }
809 adopt_fresher_durable_model_context_state(session, latest);
810 }
811 let result = self.save_session_rebasing_task_conflicts(session).await;
812 publish(session, result.is_ok());
813 result
814 }
815
816 pub async fn seed_runtime_activation_and_publish<F>(
826 &self,
827 session: &mut Session,
828 publish: F,
829 ) -> std::io::Result<()>
830 where
831 F: FnOnce(&Session, bool) + Send,
832 {
833 let _guard = self.acquire_lock(&session.id).await;
834 let mut incoming_audit = PermissionAuditSnapshot::from_metadata(&session.metadata)
835 .ok_or_else(|| {
836 std::io::Error::new(
837 std::io::ErrorKind::InvalidInput,
838 "activation seed requires a complete permission audit record",
839 )
840 })?;
841
842 if let Some(latest) = self.storage.load_session(&session.id).await? {
843 apply_authoritative_metadata(session, &latest);
844 bamboo_domain::restore_missing_admitted_inbox_messages(session, &latest);
845 bamboo_domain::merge_session_inbox_admission(session, &latest);
846 adopt_fresher_durable_model_context_state(session, &latest);
847
848 let durable_audit = PermissionAuditSnapshot::from_metadata(&latest.metadata);
849 let durable_floor = durable_audit
850 .as_ref()
851 .map(|snapshot| snapshot.audit_revision)
852 .unwrap_or_default();
853 if let Some(durable_audit) = durable_audit {
854 if durable_audit.resolution == incoming_audit.resolution {
855 incoming_audit.transitioned_at = durable_audit.transitioned_at;
856 }
857 }
858 incoming_audit.audit_revision = bamboo_domain::next_permission_audit_revision_after(
859 durable_floor.max(incoming_audit.audit_revision),
860 )
861 .map_err(|error| {
862 std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string())
863 })?;
864 }
865
866 session
867 .agent_runtime_state
868 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
869 .set_permission_mode(incoming_audit.resolution.requested);
870 incoming_audit.write_to(&mut session.metadata);
871
872 let result = self.save_session_rebasing_task_conflicts(session).await;
873 publish(session, result.is_ok());
874 result
875 }
876
877 pub async fn update_authoritative_permission_posture_and_publish<M, P>(
883 &self,
884 session_id: &str,
885 seed: &PermissionAuditSeed,
886 mutate: M,
887 publish: P,
888 ) -> std::io::Result<Option<Session>>
889 where
890 M: FnOnce(&mut Session),
891 P: FnOnce(&Session),
892 {
893 let _guard = self.acquire_lock(session_id).await;
894 let Some(mut latest) = self.storage.load_session(session_id).await? else {
895 return Ok(None);
896 };
897 let previous_mode = latest
898 .agent_runtime_state
899 .as_ref()
900 .map(|state| state.effective_permission_mode())
901 .unwrap_or_default();
902 let previous_resolution = PermissionAuditSnapshot::from_metadata(&latest.metadata)
903 .map(|snapshot| snapshot.resolution);
904 mutate(&mut latest);
905 latest
906 .agent_runtime_state
907 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
908 .set_permission_mode(seed.resolution.requested);
909 let mode_changed = previous_mode != seed.resolution.requested;
910 let posture_changed = previous_resolution != Some(seed.resolution);
911 let transitioned_at = posture_changed.then(|| chrono::Utc::now().to_rfc3339());
912 bamboo_domain::record_permission_audit(
913 &mut latest.metadata,
914 seed,
915 transitioned_at.as_deref(),
916 )
917 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string()))?;
918 if mode_changed {
919 latest.metadata_version = latest.metadata_version.saturating_add(1);
920 }
921 self.save_session_rebasing_task_conflicts(&mut latest)
922 .await?;
923 publish(&latest);
924 Ok(Some(latest))
925 }
926
927 pub async fn record_permission_posture_activation_and_publish<P>(
932 &self,
933 session_id: &str,
934 expected_audit_revision: Option<u64>,
935 seed: &PermissionAuditSeed,
936 publish: P,
937 ) -> std::io::Result<Option<Session>>
938 where
939 P: FnOnce(&Session),
940 {
941 let _guard = self.acquire_lock(session_id).await;
942 let Some(mut latest) = self.storage.load_session(session_id).await? else {
943 return Ok(None);
944 };
945 let durable_audit = PermissionAuditSnapshot::from_metadata(&latest.metadata);
946 let durable_revision = durable_audit
947 .as_ref()
948 .map(|snapshot| snapshot.audit_revision);
949 if durable_revision != expected_audit_revision {
950 return Err(std::io::Error::new(
951 std::io::ErrorKind::InvalidData,
952 "stale permission posture activation: durable audit changed after dispatch",
953 ));
954 }
955 let durable_requested = latest
956 .agent_runtime_state
957 .as_ref()
958 .map(|state| state.effective_permission_mode())
959 .unwrap_or_default();
960 if durable_requested != seed.resolution.requested || !seed.resolution.is_consistent() {
961 return Err(std::io::Error::new(
962 std::io::ErrorKind::InvalidData,
963 "stale or inconsistent permission posture activation",
964 ));
965 }
966 bamboo_domain::record_permission_audit(&mut latest.metadata, seed, None).map_err(
967 |error| std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string()),
968 )?;
969 self.save_session_rebasing_task_conflicts(&mut latest)
970 .await?;
971 publish(&latest);
972 Ok(Some(latest))
973 }
974
975 pub async fn update_runtime_config<F>(
987 &self,
988 session_id: &str,
989 mutate: F,
990 ) -> std::io::Result<Option<Session>>
991 where
992 F: FnOnce(&mut Session),
993 {
994 self.update_runtime_config_and_publish(session_id, mutate, |_| {})
995 .await
996 }
997
998 pub async fn update_runtime_config_and_publish<M, P>(
1001 &self,
1002 session_id: &str,
1003 mutate: M,
1004 publish: P,
1005 ) -> std::io::Result<Option<Session>>
1006 where
1007 M: FnOnce(&mut Session),
1008 P: FnOnce(&Session),
1009 {
1010 let _guard = self.acquire_lock(session_id).await;
1011 let Some(mut session) = self.storage.load_session(session_id).await? else {
1012 return Ok(None);
1013 };
1014 mutate(&mut session);
1015 self.save_session_rebasing_task_conflicts(&mut session)
1016 .await?;
1017 publish(&session);
1018 Ok(Some(session))
1019 }
1020
1021 async fn load_response_candidate<C>(
1026 &self,
1027 session_id: &str,
1028 load_cached: C,
1029 ) -> std::io::Result<Option<Session>>
1030 where
1031 C: FnOnce() -> Option<Session> + Send,
1032 {
1033 let cached_candidate = load_cached();
1034 let durable = self.storage.load_session(session_id).await?;
1035 let Some(mut session) = (match (cached_candidate, durable.as_ref()) {
1036 (Some(cached), Some(durable)) => {
1037 let prefer_durable = durable.updated_at > cached.updated_at
1038 || (durable.updated_at == cached.updated_at
1039 && cached.pending_question.is_none()
1040 && durable.pending_question.is_some());
1041 Some(if prefer_durable {
1042 durable.clone()
1043 } else {
1044 cached
1045 })
1046 }
1047 (Some(cached), None) => Some(cached),
1048 (None, durable) => durable.cloned(),
1049 }) else {
1050 return Ok(None);
1051 };
1052 if let Some(latest) = durable.as_ref() {
1053 adopt_durable_consumed_clarification(&mut session, latest);
1058 bamboo_domain::append_missing_runtime_messages(&mut session, latest);
1059 apply_authoritative_metadata(&mut session, latest);
1060 let restored =
1061 bamboo_domain::restore_missing_admitted_inbox_messages(&mut session, latest);
1062 if restored > 0 {
1063 tracing::warn!(
1064 session_id,
1065 restored,
1066 "restored durable SessionInbox transcript messages into response transaction"
1067 );
1068 }
1069 bamboo_domain::merge_session_inbox_admission(&mut session, latest);
1070 adopt_fresher_disk_permission_posture(&mut session, latest);
1071 adopt_fresher_durable_model_context_state(&mut session, latest);
1072 }
1073 Ok(Some(session))
1074 }
1075
1076 pub async fn inspect_runtime_session_for_response<C>(
1081 &self,
1082 session_id: &str,
1083 load_cached: C,
1084 ) -> std::io::Result<Option<Session>>
1085 where
1086 C: FnOnce() -> Option<Session> + Send,
1087 {
1088 let _guard = self.acquire_lock(session_id).await;
1089 self.load_response_candidate(session_id, load_cached).await
1090 }
1091
1092 pub async fn mutate_runtime_session_and_publish<C, M, P, E>(
1102 &self,
1103 session_id: &str,
1104 load_cached: C,
1105 mutate: M,
1106 publish: P,
1107 ) -> std::io::Result<Result<Option<Session>, E>>
1108 where
1109 C: FnOnce() -> Option<Session> + Send,
1110 M: FnOnce(&mut Session) -> Result<(), E> + Send,
1111 P: FnOnce(&Session) + Send,
1112 E: Send,
1113 {
1114 let _guard = self.acquire_lock(session_id).await;
1115 let Some(mut session) = self
1116 .load_response_candidate(session_id, load_cached)
1117 .await?
1118 else {
1119 return Ok(Ok(None));
1120 };
1121 if let Err(error) = mutate(&mut session) {
1122 return Ok(Err(error));
1123 }
1124 self.save_session_rebasing_task_conflicts(&mut session)
1125 .await?;
1126 publish(&session);
1127 Ok(Ok(Some(session)))
1128 }
1129
1130 pub async fn clear_legacy_pending_messages_and_publish<F>(
1133 &self,
1134 session_id: &str,
1135 expected: &[serde_json::Value],
1136 publish: F,
1137 ) -> std::io::Result<bool>
1138 where
1139 F: FnOnce(&Session) + Send,
1140 {
1141 let _guard = self.acquire_lock(session_id).await;
1142 let Some(mut latest) = self.storage.load_session(session_id).await? else {
1143 return Ok(false);
1144 };
1145 if latest.pending_injected_messages().as_deref() != Some(expected) {
1146 return Ok(false);
1147 }
1148 latest.clear_pending_injected_messages();
1149 self.save_runtime_state_rebasing_task_conflicts(&mut latest)
1150 .await?;
1151 publish(&latest);
1152 Ok(true)
1153 }
1154}
1155
1156#[async_trait::async_trait]
1160impl RuntimeSessionPersistence for LockedSessionStore {
1161 async fn save_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
1162 self.merge_save_runtime(session).await
1163 }
1164
1165 async fn seed_runtime_activation(&self, session: &mut Session) -> std::io::Result<()> {
1166 self.seed_runtime_activation_and_publish(session, |_, _| {})
1167 .await
1168 }
1169
1170 async fn record_permission_posture_activation(
1171 &self,
1172 session_id: &str,
1173 expected_audit_revision: Option<u64>,
1174 seed: &PermissionAuditSeed,
1175 ) -> std::io::Result<Option<Session>> {
1176 self.record_permission_posture_activation_and_publish(
1177 session_id,
1178 expected_audit_revision,
1179 seed,
1180 |_| {},
1181 )
1182 .await
1183 }
1184
1185 async fn save_runtime_control_plane(&self, session: &mut Session) -> std::io::Result<()> {
1186 self.save_runtime_only(session).await
1187 }
1188
1189 async fn load_runtime_control_plane(
1190 &self,
1191 session_id: &str,
1192 ) -> std::io::Result<Option<Session>> {
1193 self.storage.load_runtime_control_plane(session_id).await
1194 }
1195
1196 async fn update_task_list_control_plane(
1197 &self,
1198 session_id: &str,
1199 task_list: &bamboo_domain::TaskList,
1200 version: &str,
1201 ) -> std::io::Result<bool> {
1202 self.update_task_list_control_plane_and_publish(session_id, task_list, version, |_| {})
1203 .await
1204 }
1205
1206 async fn update_task_list_control_plane_if_version(
1207 &self,
1208 session_id: &str,
1209 expected_version: &str,
1210 expected_task_list: &bamboo_domain::TaskList,
1211 task_list: &bamboo_domain::TaskList,
1212 version: &str,
1213 ) -> std::io::Result<bool> {
1214 self.update_task_list_control_plane_if_version_and_publish(
1215 session_id,
1216 expected_version,
1217 expected_task_list,
1218 task_list,
1219 version,
1220 |_| {},
1221 )
1222 .await
1223 }
1224
1225 async fn update_task_list_control_planes_if_version(
1226 &self,
1227 session_id: &str,
1228 shared_session_id: &str,
1229 expected_version: &str,
1230 expected_task_list: &bamboo_domain::TaskList,
1231 task_list: &bamboo_domain::TaskList,
1232 version: &str,
1233 ) -> std::io::Result<bool> {
1234 self.update_task_list_control_planes_if_version_and_publish(
1235 session_id,
1236 shared_session_id,
1237 expected_version,
1238 expected_task_list,
1239 task_list,
1240 version,
1241 |_, _| {},
1242 )
1243 .await
1244 }
1245
1246 async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
1247 LockedSessionStore::checkpoint_runtime_session(self, session).await
1248 }
1249
1250 async fn load_runtime_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
1251 self.storage.load_session(session_id).await
1252 }
1253
1254 async fn clear_legacy_pending_messages(
1255 &self,
1256 session_id: &str,
1257 expected: &[serde_json::Value],
1258 ) -> std::io::Result<bool> {
1259 self.clear_legacy_pending_messages_and_publish(session_id, expected, |_| {})
1260 .await
1261 }
1262}
1263
1264async fn merge_authoritative_metadata_into_stale(
1273 storage: &Arc<dyn Storage>,
1274 session: &mut Session,
1275) -> std::io::Result<()> {
1276 if let Some(latest) = storage.load_session(&session.id).await? {
1277 adopt_durable_consumed_clarification(session, &latest);
1278 apply_authoritative_metadata(session, &latest);
1279 bamboo_domain::restore_missing_admitted_inbox_messages(session, &latest);
1280 bamboo_domain::merge_session_inbox_admission(session, &latest);
1281 adopt_fresher_disk_permission_posture(session, &latest);
1282 adopt_fresher_durable_model_context_state(session, &latest);
1283 }
1284 Ok(())
1285}
1286
1287fn adopt_durable_model_context_state(session: &mut Session, latest: &Session) {
1292 session
1293 .model_context_state
1294 .clone_from(&latest.model_context_state);
1295}
1296
1297fn adopt_fresher_durable_model_context_state(session: &mut Session, latest: &Session) {
1303 let adopt = match (
1304 session.model_context_state.as_ref(),
1305 latest.model_context_state.as_ref(),
1306 ) {
1307 (None, Some(_)) => true,
1308 (Some(incoming), Some(durable)) => {
1309 durable.state_revision > incoming.state_revision
1310 || (durable.state_revision == incoming.state_revision && durable != incoming)
1311 }
1312 _ => false,
1313 };
1314 if adopt {
1315 adopt_durable_model_context_state(session, latest);
1316 }
1317}
1318
1319fn ensure_model_context_checkpoint_is_current(
1324 session: &Session,
1325 latest: &Session,
1326) -> std::io::Result<()> {
1327 let stale_or_conflicting = match (
1328 session.model_context_state.as_ref(),
1329 latest.model_context_state.as_ref(),
1330 ) {
1331 (None, Some(_)) => true,
1332 (Some(incoming), Some(durable)) => {
1333 durable.state_revision > incoming.state_revision
1334 || (durable.state_revision == incoming.state_revision && durable != incoming)
1335 }
1336 _ => false,
1337 };
1338 if stale_or_conflicting {
1339 return Err(std::io::Error::new(
1340 std::io::ErrorKind::WouldBlock,
1341 "stale or conflicting model-context ledger checkpoint",
1342 ));
1343 }
1344 Ok(())
1345}
1346
1347fn adopt_fresher_disk_permission_posture(session: &mut Session, latest: &Session) {
1359 let Some(disk_mode) = latest
1364 .agent_runtime_state
1365 .as_ref()
1366 .map(|state| state.effective_permission_mode())
1367 else {
1368 return;
1369 };
1370 let current_mode = session
1371 .agent_runtime_state
1372 .as_ref()
1373 .map(|state| state.effective_permission_mode())
1374 .unwrap_or_default();
1375 let Some(disk_audit) = bamboo_domain::fresher_disk_permission_audit(
1376 current_mode,
1377 &session.metadata,
1378 disk_mode,
1379 &latest.metadata,
1380 ) else {
1381 return;
1382 };
1383
1384 match session.agent_runtime_state.as_mut() {
1385 Some(state) => state.set_permission_mode(disk_mode),
1386 None if disk_mode != bamboo_domain::SessionPermissionMode::Default => {
1389 let state = session
1390 .agent_runtime_state
1391 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default);
1392 state.set_permission_mode(disk_mode);
1393 }
1394 None => {}
1395 }
1396
1397 disk_audit.write_to(&mut session.metadata);
1399}
1400
1401fn apply_authoritative_metadata(session: &mut Session, latest: &Session) {
1407 if latest.metadata_version >= session.metadata_version {
1408 session.title = latest.title.clone();
1409 session.title_version = latest.title_version;
1410 session.title_generated = latest.title_generated;
1411 session.pinned = latest.pinned;
1412 for key in AUTHORITATIVE_METADATA_KEYS {
1413 if let Some(value) = latest.metadata.get(*key) {
1414 session.metadata.insert((*key).to_string(), value.clone());
1415 } else {
1416 session.metadata.remove(*key);
1417 }
1418 }
1419 session.metadata_version = latest.metadata_version;
1420 }
1421}
1422
1423pub async fn merge_save_session(
1436 storage: &Arc<dyn Storage>,
1437 session: &mut Session,
1438) -> std::io::Result<()> {
1439 merge_authoritative_metadata_into_stale(storage, session).await?;
1440 storage.save_session(session).await
1441}
1442
1443#[cfg(test)]
1446mod tests {
1447 use super::*;
1448 use crate::v2::{RuntimeTaskTransactionFault, SessionStoreV2};
1449 use bamboo_domain::{session::types::Session, PermissionMode};
1450 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
1451
1452 struct CountingControlPlaneStorage {
1453 inner: Arc<SessionStoreV2>,
1454 control_plane_loads: AtomicUsize,
1455 full_saves: AtomicUsize,
1456 runtime_state_saves: AtomicUsize,
1457 }
1458
1459 struct PairCommitBarrierStorage {
1460 inner: Arc<SessionStoreV2>,
1461 before_commit: Arc<tokio::sync::Barrier>,
1462 }
1463
1464 struct SingleCommitBarrierStorage {
1465 inner: Arc<SessionStoreV2>,
1466 before_commit: Arc<tokio::sync::Barrier>,
1467 }
1468
1469 struct SingleCommitPauseStorage {
1470 inner: Arc<SessionStoreV2>,
1471 commit_reached: Arc<tokio::sync::Barrier>,
1472 release_commit: Arc<tokio::sync::Barrier>,
1473 }
1474
1475 #[async_trait::async_trait]
1476 impl Storage for SingleCommitPauseStorage {
1477 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
1478 self.inner.save_session(session).await
1479 }
1480
1481 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
1482 self.inner.load_session(session_id).await
1483 }
1484
1485 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
1486 self.inner.delete_session(session_id).await
1487 }
1488
1489 async fn load_runtime_control_plane(
1490 &self,
1491 session_id: &str,
1492 ) -> std::io::Result<Option<Session>> {
1493 self.inner.load_runtime_control_plane(session_id).await
1494 }
1495
1496 async fn save_task_control_plane_if_matches(
1497 &self,
1498 original: &Session,
1499 updated: &Session,
1500 ) -> std::io::Result<bool> {
1501 self.commit_reached.wait().await;
1502 self.release_commit.wait().await;
1503 self.inner
1504 .save_task_control_plane_if_matches(original, updated)
1505 .await
1506 }
1507 }
1508
1509 #[async_trait::async_trait]
1510 impl Storage for SingleCommitBarrierStorage {
1511 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
1512 self.inner.save_session(session).await
1513 }
1514
1515 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
1516 self.inner.load_session(session_id).await
1517 }
1518
1519 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
1520 self.inner.delete_session(session_id).await
1521 }
1522
1523 async fn load_runtime_control_plane(
1524 &self,
1525 session_id: &str,
1526 ) -> std::io::Result<Option<Session>> {
1527 self.inner.load_runtime_control_plane(session_id).await
1528 }
1529
1530 async fn save_task_control_plane_if_matches(
1531 &self,
1532 original: &Session,
1533 updated: &Session,
1534 ) -> std::io::Result<bool> {
1535 self.before_commit.wait().await;
1536 self.inner
1537 .save_task_control_plane_if_matches(original, updated)
1538 .await
1539 }
1540 }
1541
1542 #[async_trait::async_trait]
1543 impl Storage for PairCommitBarrierStorage {
1544 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
1545 self.inner.save_session(session).await
1546 }
1547
1548 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
1549 self.inner.load_session(session_id).await
1550 }
1551
1552 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
1553 self.inner.delete_session(session_id).await
1554 }
1555
1556 async fn save_runtime_state(&self, session: &Session) -> std::io::Result<()> {
1557 self.inner.save_runtime_state(session).await
1558 }
1559
1560 async fn load_runtime_control_plane(
1561 &self,
1562 session_id: &str,
1563 ) -> std::io::Result<Option<Session>> {
1564 self.inner.load_runtime_control_plane(session_id).await
1565 }
1566
1567 async fn recover_task_control_plane_transaction(
1568 &self,
1569 first_session_id: &str,
1570 second_session_id: &str,
1571 ) -> std::io::Result<()> {
1572 self.inner
1573 .recover_task_control_plane_transaction(first_session_id, second_session_id)
1574 .await
1575 }
1576
1577 async fn save_task_control_planes_atomically(
1578 &self,
1579 first_original: &Session,
1580 first_updated: &Session,
1581 second_original: &Session,
1582 second_updated: &Session,
1583 ) -> std::io::Result<bool> {
1584 self.before_commit.wait().await;
1589 self.inner
1590 .save_task_control_planes_atomically(
1591 first_original,
1592 first_updated,
1593 second_original,
1594 second_updated,
1595 )
1596 .await
1597 }
1598 }
1599
1600 #[async_trait::async_trait]
1601 impl Storage for CountingControlPlaneStorage {
1602 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
1603 self.full_saves.fetch_add(1, Ordering::SeqCst);
1604 self.inner.save_session(session).await
1605 }
1606
1607 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
1608 self.inner.load_session(session_id).await
1609 }
1610
1611 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
1612 self.inner.delete_session(session_id).await
1613 }
1614
1615 async fn save_runtime_state(&self, session: &Session) -> std::io::Result<()> {
1616 self.runtime_state_saves.fetch_add(1, Ordering::SeqCst);
1617 self.inner.save_runtime_state(session).await
1618 }
1619
1620 async fn load_runtime_control_plane(
1621 &self,
1622 session_id: &str,
1623 ) -> std::io::Result<Option<Session>> {
1624 self.control_plane_loads.fetch_add(1, Ordering::SeqCst);
1625 self.inner.load_runtime_control_plane(session_id).await
1626 }
1627
1628 async fn recover_task_control_plane_transaction(
1629 &self,
1630 first_session_id: &str,
1631 second_session_id: &str,
1632 ) -> std::io::Result<()> {
1633 self.inner
1634 .recover_task_control_plane_transaction(first_session_id, second_session_id)
1635 .await
1636 }
1637
1638 async fn save_task_control_plane_if_matches(
1639 &self,
1640 original: &Session,
1641 updated: &Session,
1642 ) -> std::io::Result<bool> {
1643 let committed = self
1644 .inner
1645 .save_task_control_plane_if_matches(original, updated)
1646 .await?;
1647 if committed {
1648 self.runtime_state_saves.fetch_add(1, Ordering::SeqCst);
1649 }
1650 Ok(committed)
1651 }
1652
1653 async fn save_task_control_planes_atomically(
1654 &self,
1655 first_original: &Session,
1656 first_updated: &Session,
1657 second_original: &Session,
1658 second_updated: &Session,
1659 ) -> std::io::Result<bool> {
1660 let committed = self
1661 .inner
1662 .save_task_control_planes_atomically(
1663 first_original,
1664 first_updated,
1665 second_original,
1666 second_updated,
1667 )
1668 .await?;
1669 if committed {
1670 self.runtime_state_saves.fetch_add(2, Ordering::SeqCst);
1671 }
1672 Ok(committed)
1673 }
1674 }
1675
1676 async fn make_storage() -> (tempfile::TempDir, Arc<dyn Storage>) {
1677 let temp = tempfile::tempdir().unwrap();
1678 let storage = SessionStoreV2::new(temp.path().to_path_buf())
1679 .await
1680 .expect("storage init");
1681 (temp, Arc::new(storage) as Arc<dyn Storage>)
1682 }
1683
1684 fn fresh(id: &str) -> Session {
1685 Session::new(id.to_string(), "test-model".to_string())
1686 }
1687
1688 fn typed_permission_result(
1689 tool_call_id: &str,
1690 message_id: &str,
1691 generation: &str,
1692 content: &str,
1693 ) -> bamboo_domain::session::types::Message {
1694 let mut message =
1695 bamboo_domain::session::types::Message::tool_result(tool_call_id, content);
1696 message.id = message_id.to_string();
1697 message.metadata = Some(serde_json::json!({
1698 "permission_request": {
1699 "request_generation": generation,
1700 }
1701 }));
1702 message
1703 }
1704
1705 fn ledger_state(state_revision: u64, marker: &str) -> bamboo_domain::ModelContextState {
1706 bamboo_domain::ModelContextState {
1707 state_revision,
1708 prefix_epoch: state_revision,
1709 cache_scope_sha256: Some("scope".to_string()),
1710 transcript_item_sha256: vec![marker.to_string()],
1711 ..bamboo_domain::ModelContextState::default()
1712 }
1713 }
1714
1715 fn set_permission_audit(
1716 session: &mut Session,
1717 requested: bamboo_domain::SessionPermissionMode,
1718 policy_revision: u64,
1719 mapping: &str,
1720 transitioned_at: &str,
1721 ) -> u64 {
1722 let resolution = bamboo_domain::resolve_permission_mode(requested, PermissionMode::Default);
1723 session
1724 .agent_runtime_state
1725 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
1726 .set_permission_mode(requested);
1727 bamboo_domain::record_permission_audit(
1728 &mut session.metadata,
1729 &PermissionAuditSeed::new(policy_revision, resolution, mapping),
1730 Some(transitioned_at),
1731 )
1732 .unwrap()
1733 }
1734
1735 #[tokio::test]
1736 async fn checked_runtime_mutation_persists_a_cache_only_pending_session() {
1737 let (_temp, storage) = make_storage().await;
1738 let store = LockedSessionStore::new(storage.clone());
1739 let mut cached = fresh("cache-only-response");
1740 cached.set_pending_question(
1741 "tool-1".to_string(),
1742 "ConclusionWithOptions".to_string(),
1743 "Choose".to_string(),
1744 vec!["A".to_string()],
1745 false,
1746 );
1747
1748 let saved = store
1749 .mutate_runtime_session_and_publish(
1750 &cached.id.clone(),
1751 move || Some(cached),
1752 |session| {
1753 assert!(session.pending_question.is_some());
1754 session.clear_pending_question();
1755 Ok::<_, ()>(())
1756 },
1757 |_| {},
1758 )
1759 .await
1760 .unwrap()
1761 .unwrap()
1762 .expect("cache-only session should be created durably");
1763
1764 assert!(saved.pending_question.is_none());
1765 assert!(storage
1766 .load_session("cache-only-response")
1767 .await
1768 .unwrap()
1769 .unwrap()
1770 .pending_question
1771 .is_none());
1772 }
1773
1774 #[tokio::test]
1775 async fn checked_runtime_mutation_preserves_durable_authorities_for_newer_cache() {
1776 let (_temp, storage) = make_storage().await;
1777 let store = LockedSessionStore::new(storage.clone());
1778 let session_id = "cached-response-authorities";
1779 let mut durable = fresh(session_id);
1780 durable.title = "Durable title".to_string();
1781 durable.title_version = 4;
1782 durable.title_generated = true;
1783 durable.metadata_version = 9;
1784 set_permission_audit(
1785 &mut durable,
1786 bamboo_domain::SessionPermissionMode::Auto,
1787 7,
1788 "bamboo_runtime:durable-auto",
1789 "2026-08-10T09:00:00Z",
1790 );
1791 storage.save_session(&durable).await.unwrap();
1792
1793 let mut cached = fresh(session_id);
1794 cached.title = "Stale cached title".to_string();
1795 cached.updated_at = durable.updated_at + chrono::Duration::seconds(1);
1796 cached.set_pending_question(
1797 "tool-1".to_string(),
1798 "ConclusionWithOptions".to_string(),
1799 "Choose".to_string(),
1800 vec!["A".to_string()],
1801 false,
1802 );
1803
1804 store
1805 .mutate_runtime_session_and_publish(
1806 session_id,
1807 move || Some(cached),
1808 |session| {
1809 session.clear_pending_question();
1810 Ok::<_, ()>(())
1811 },
1812 |_| {},
1813 )
1814 .await
1815 .unwrap()
1816 .unwrap()
1817 .expect("session should exist");
1818
1819 let saved = storage.load_session(session_id).await.unwrap().unwrap();
1820 assert_eq!(saved.title, "Durable title");
1821 assert_eq!(saved.title_version, 4);
1822 assert_eq!(saved.metadata_version, 9);
1823 assert_eq!(
1824 saved
1825 .agent_runtime_state
1826 .as_ref()
1827 .unwrap()
1828 .effective_permission_mode(),
1829 bamboo_domain::SessionPermissionMode::Auto
1830 );
1831 let audit = bamboo_domain::PermissionAuditSnapshot::from_metadata(&saved.metadata).unwrap();
1832 assert_eq!(audit.policy_revision, 7);
1833 assert_eq!(audit.executor_mapping, "bamboo_runtime:durable-auto");
1834 }
1835
1836 #[tokio::test]
1837 async fn checked_runtime_mutation_cannot_resurrect_consumed_ask_from_newer_cache() {
1838 use bamboo_domain::session::types::Message;
1839
1840 let (_temp, storage) = make_storage().await;
1841 let store = LockedSessionStore::new(storage.clone());
1842 let session_id = "cached-consumed-response";
1843 let mut stale_cached = fresh(session_id);
1844 stale_cached.add_message(Message::tool_result("call-1", "waiting"));
1845 stale_cached.set_pending_question(
1846 "call-1".to_string(),
1847 "ConclusionWithOptions".to_string(),
1848 "Choose".to_string(),
1849 vec!["A".to_string()],
1850 false,
1851 );
1852
1853 let mut durable = stale_cached.clone();
1854 durable.clear_pending_question();
1855 durable.metadata.insert(
1856 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
1857 r#"["call-1"]"#.to_string(),
1858 );
1859 durable.messages[0].content = "Selected response: A".to_string();
1860 durable.add_message(Message::user("durable concurrent message"));
1861 storage.save_session(&durable).await.unwrap();
1862
1863 stale_cached.updated_at = durable.updated_at + chrono::Duration::seconds(1);
1866 let saved = store
1867 .mutate_runtime_session_and_publish(
1868 session_id,
1869 move || Some(stale_cached),
1870 |session| {
1871 assert!(session.pending_question.is_none());
1872 Ok::<_, ()>(())
1873 },
1874 |_| {},
1875 )
1876 .await
1877 .unwrap()
1878 .unwrap()
1879 .unwrap();
1880
1881 assert!(saved.pending_question.is_none());
1882 assert_eq!(saved.messages.len(), 2);
1883 assert_eq!(saved.messages[0].content, "Selected response: A");
1884 assert_eq!(saved.messages[1].content, "durable concurrent message");
1885 }
1886
1887 #[tokio::test]
1888 async fn response_inspection_adopts_durable_consumption_without_writing() {
1889 let (_temp, storage) = make_storage().await;
1890 let store = LockedSessionStore::new(storage.clone());
1891 let session_id = "inspect-consumed-response";
1892 let mut stale_cached = fresh(session_id);
1893 stale_cached.add_message(bamboo_domain::session::types::Message::tool_result(
1898 "call-1", "waiting",
1899 ));
1900 stale_cached.set_pending_question(
1901 "call-1".to_string(),
1902 "ConclusionWithOptions".to_string(),
1903 "Choose".to_string(),
1904 vec!["A".to_string()],
1905 false,
1906 );
1907
1908 let mut durable = stale_cached.clone();
1909 durable.clear_pending_question();
1910 durable.metadata.insert(
1911 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
1912 r#"["call-1"]"#.to_string(),
1913 );
1914 storage.save_session(&durable).await.unwrap();
1915 stale_cached.updated_at = durable.updated_at + chrono::Duration::seconds(1);
1916
1917 let inspected = store
1918 .inspect_runtime_session_for_response(session_id, move || Some(stale_cached))
1919 .await
1920 .unwrap()
1921 .expect("session should be inspectable");
1922 assert!(inspected.pending_question.is_none());
1923
1924 let unchanged = storage.load_session(session_id).await.unwrap().unwrap();
1925 assert_eq!(unchanged.updated_at, durable.updated_at);
1926 assert!(unchanged.pending_question.is_none());
1927 }
1928
1929 #[tokio::test]
1930 async fn same_mode_newer_run_start_audit_survives_every_runtime_save_path() {
1931 for path in ["merge", "checkpoint", "control-plane"] {
1932 let (_temp, storage) = make_storage().await;
1933 let store = LockedSessionStore::new(storage.clone());
1934 let session_id = format!("same-mode-newer-{path}");
1935 let mut durable = fresh(&session_id);
1936 set_permission_audit(
1937 &mut durable,
1938 bamboo_domain::SessionPermissionMode::Default,
1939 1,
1940 "bamboo_runtime:old-policy",
1941 "2026-07-31T12:00:00Z",
1942 );
1943 storage.save_session(&durable).await.unwrap();
1944
1945 let mut run_start = durable.clone();
1946 let old_revision = PermissionAuditSnapshot::from_metadata(&durable.metadata)
1947 .unwrap()
1948 .audit_revision;
1949 let new_revision = set_permission_audit(
1950 &mut run_start,
1951 bamboo_domain::SessionPermissionMode::Default,
1952 2,
1953 "bamboo_runtime:new-policy",
1954 "2026-07-31T12:00:00Z",
1955 );
1956 assert!(new_revision > old_revision);
1957
1958 match path {
1959 "merge" => store.merge_save_runtime(&mut run_start).await.unwrap(),
1960 "checkpoint" => store
1961 .checkpoint_runtime_session(&mut run_start)
1962 .await
1963 .unwrap(),
1964 "control-plane" => store.save_runtime_only(&mut run_start).await.unwrap(),
1965 _ => unreachable!(),
1966 }
1967
1968 let saved = storage.load_session(&session_id).await.unwrap().unwrap();
1969 let audit = PermissionAuditSnapshot::from_metadata(&saved.metadata).unwrap();
1970 assert_eq!(audit.audit_revision, new_revision, "path={path}");
1971 assert_eq!(audit.policy_revision, 2, "path={path}");
1972 assert_eq!(audit.executor_mapping, "bamboo_runtime:new-policy");
1973 }
1974 }
1975
1976 #[tokio::test]
1977 async fn newer_disk_transition_wins_after_mode_cycles_back() {
1978 let (_temp, storage) = make_storage().await;
1979 let store = LockedSessionStore::new(storage.clone());
1980 let session_id = "permission-cycle-back";
1981 let mut baseline = fresh(session_id);
1982 let stale_revision = set_permission_audit(
1983 &mut baseline,
1984 bamboo_domain::SessionPermissionMode::Default,
1985 1,
1986 "bamboo_runtime:initial",
1987 "2026-07-31T12:00:00Z",
1988 );
1989 storage.save_session(&baseline).await.unwrap();
1990 let mut stale_runtime = baseline.clone();
1991
1992 let mut durable = baseline;
1993 set_permission_audit(
1994 &mut durable,
1995 bamboo_domain::SessionPermissionMode::Auto,
1996 2,
1997 "bamboo_runtime:auto",
1998 "2026-07-31T12:01:00Z",
1999 );
2000 let durable_revision = set_permission_audit(
2001 &mut durable,
2002 bamboo_domain::SessionPermissionMode::Default,
2003 3,
2004 "bamboo_runtime:cycled-default",
2005 "2026-07-31T12:02:00Z",
2006 );
2007 assert!(durable_revision > stale_revision);
2008 storage.save_session(&durable).await.unwrap();
2009
2010 store.merge_save_runtime(&mut stale_runtime).await.unwrap();
2011 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2012 let audit = PermissionAuditSnapshot::from_metadata(&saved.metadata).unwrap();
2013 assert_eq!(audit.audit_revision, durable_revision);
2014 assert_eq!(audit.policy_revision, 3);
2015 assert_eq!(audit.executor_mapping, "bamboo_runtime:cycled-default");
2016 }
2017
2018 #[tokio::test]
2019 async fn authoritative_activation_seed_replaces_every_warm_worker_posture() {
2020 let (_temp, storage) = make_storage().await;
2021 let store = LockedSessionStore::new(storage.clone());
2022 let session_id = "warm-permission-matrix";
2023 let cases = [
2024 (
2025 bamboo_domain::SessionPermissionMode::Auto,
2026 PermissionMode::Default,
2027 PermissionMode::Auto,
2028 ),
2029 (
2030 bamboo_domain::SessionPermissionMode::Default,
2031 PermissionMode::Default,
2032 PermissionMode::Default,
2033 ),
2034 (
2035 bamboo_domain::SessionPermissionMode::Auto,
2036 PermissionMode::Default,
2037 PermissionMode::Auto,
2038 ),
2039 (
2040 bamboo_domain::SessionPermissionMode::Bypass,
2041 PermissionMode::Auto,
2042 PermissionMode::BypassPermissions,
2043 ),
2044 ];
2045 let mut previous_revision = 0;
2046
2047 for (index, (requested, configured, expected_effective)) in cases.into_iter().enumerate() {
2048 let mut activation = fresh(session_id);
2049 activation
2050 .agent_runtime_state
2051 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
2052 .set_permission_mode(requested);
2053 let resolution = bamboo_domain::resolve_permission_mode(requested, configured);
2054 bamboo_domain::record_permission_audit(
2055 &mut activation.metadata,
2056 &PermissionAuditSeed::new(
2057 index as u64 + 1,
2058 resolution,
2059 format!("bamboo_worker:{}", resolution.effective.as_str()),
2060 ),
2061 Some("2026-07-31T12:00:00Z"),
2062 )
2063 .unwrap();
2064
2065 RuntimeSessionPersistence::seed_runtime_activation(&store, &mut activation)
2066 .await
2067 .unwrap();
2068 let durable = storage.load_session(session_id).await.unwrap().unwrap();
2069 assert_eq!(
2070 durable
2071 .agent_runtime_state
2072 .as_ref()
2073 .unwrap()
2074 .effective_permission_mode(),
2075 requested,
2076 "activation {index}"
2077 );
2078 let audit = PermissionAuditSnapshot::from_metadata(&durable.metadata).unwrap();
2079 assert_eq!(audit.resolution.requested, requested);
2080 assert_eq!(audit.resolution.effective, expected_effective);
2081 assert!(audit.audit_revision > previous_revision);
2082 previous_revision = audit.audit_revision;
2083 }
2084 }
2085
2086 #[tokio::test]
2087 async fn resident_reseed_bumps_etag_only_for_typed_transition() {
2088 let (_temp, storage) = make_storage().await;
2089 let store = LockedSessionStore::new(storage.clone());
2090 let session_id = "resident-atomic-permission";
2091 let mut baseline = fresh(session_id);
2092 baseline.metadata_version = 7;
2093 set_permission_audit(
2094 &mut baseline,
2095 bamboo_domain::SessionPermissionMode::Auto,
2096 1,
2097 "bamboo_runtime:auto",
2098 "2026-07-31T12:00:00Z",
2099 );
2100 storage.save_session(&baseline).await.unwrap();
2101 let initial_audit = PermissionAuditSnapshot::from_metadata(&baseline.metadata).unwrap();
2102
2103 let same_mode_seed = PermissionAuditSeed::bamboo_runtime(
2104 2,
2105 bamboo_domain::resolve_permission_mode(
2106 bamboo_domain::SessionPermissionMode::Auto,
2107 PermissionMode::Default,
2108 ),
2109 );
2110 let refreshed = store
2111 .update_authoritative_permission_posture_and_publish(
2112 session_id,
2113 &same_mode_seed,
2114 |session| {
2115 session
2116 .metadata
2117 .insert("resident.marker".to_string(), "same-mode".to_string());
2118 },
2119 |_| {},
2120 )
2121 .await
2122 .unwrap()
2123 .unwrap();
2124 let refreshed_audit = PermissionAuditSnapshot::from_metadata(&refreshed.metadata).unwrap();
2125 assert_eq!(refreshed.metadata_version, 7);
2126 assert!(refreshed_audit.audit_revision > initial_audit.audit_revision);
2127 assert_eq!(refreshed_audit.policy_revision, 2);
2128
2129 let transition_seed = PermissionAuditSeed::bamboo_runtime(
2130 3,
2131 bamboo_domain::resolve_permission_mode(
2132 bamboo_domain::SessionPermissionMode::Default,
2133 PermissionMode::Default,
2134 ),
2135 );
2136 let transitioned = store
2137 .update_authoritative_permission_posture_and_publish(
2138 session_id,
2139 &transition_seed,
2140 |session| {
2141 session
2142 .metadata
2143 .insert("resident.marker".to_string(), "transition".to_string());
2144 },
2145 |_| {},
2146 )
2147 .await
2148 .unwrap()
2149 .unwrap();
2150 let transitioned_audit =
2151 PermissionAuditSnapshot::from_metadata(&transitioned.metadata).unwrap();
2152 assert_eq!(transitioned.metadata_version, 8, "old ETag must be invalid");
2153 assert_eq!(
2154 transitioned
2155 .agent_runtime_state
2156 .as_ref()
2157 .unwrap()
2158 .effective_permission_mode(),
2159 bamboo_domain::SessionPermissionMode::Default
2160 );
2161 assert_eq!(
2162 transitioned_audit.resolution.requested,
2163 bamboo_domain::SessionPermissionMode::Default
2164 );
2165 assert!(transitioned_audit.audit_revision > refreshed_audit.audit_revision);
2166 assert_eq!(
2167 transitioned
2168 .metadata
2169 .get("resident.marker")
2170 .map(String::as_str),
2171 Some("transition")
2172 );
2173 }
2174
2175 #[tokio::test]
2176 async fn worker_activation_cas_cannot_overwrite_concurrent_permission_patch() {
2177 let (_temp, storage) = make_storage().await;
2178 let store = LockedSessionStore::new(storage.clone());
2179 let session_id = "permission-activation-cas";
2180 let mut baseline = fresh(session_id);
2181 set_permission_audit(
2182 &mut baseline,
2183 bamboo_domain::SessionPermissionMode::Default,
2184 1,
2185 "bamboo_runtime:default",
2186 "2026-07-31T12:00:00Z",
2187 );
2188 storage.save_session(&baseline).await.unwrap();
2189 let dispatched_revision = PermissionAuditSnapshot::from_metadata(&baseline.metadata)
2190 .unwrap()
2191 .audit_revision;
2192
2193 let patched_resolution = bamboo_domain::resolve_permission_mode(
2194 bamboo_domain::SessionPermissionMode::Auto,
2195 PermissionMode::Default,
2196 );
2197 let patched = store
2198 .update_authoritative_permission_posture_and_publish(
2199 session_id,
2200 &PermissionAuditSeed::new(2, patched_resolution, "patch:auto"),
2201 |_| {},
2202 |_| {},
2203 )
2204 .await
2205 .unwrap()
2206 .unwrap();
2207 let patched_audit = PermissionAuditSnapshot::from_metadata(&patched.metadata).unwrap();
2208 assert!(patched_audit.audit_revision > dispatched_revision);
2209
2210 let stale_worker_seed = PermissionAuditSeed::new(
2211 1,
2212 bamboo_domain::resolve_permission_mode(
2213 bamboo_domain::SessionPermissionMode::Default,
2214 PermissionMode::Default,
2215 ),
2216 "worker:stale-default",
2217 );
2218 let error = store
2219 .record_permission_posture_activation_and_publish(
2220 session_id,
2221 Some(dispatched_revision),
2222 &stale_worker_seed,
2223 |_| {},
2224 )
2225 .await
2226 .unwrap_err();
2227 assert!(error.to_string().contains("durable audit changed"));
2228
2229 let durable = storage.load_session(session_id).await.unwrap().unwrap();
2230 let durable_audit = PermissionAuditSnapshot::from_metadata(&durable.metadata).unwrap();
2231 assert_eq!(durable_audit, patched_audit);
2232 assert_eq!(durable_audit.executor_mapping, "patch:auto");
2233 }
2234
2235 #[tokio::test]
2238 async fn update_runtime_config_preserves_concurrently_appended_messages() {
2239 use bamboo_domain::session::types::Message;
2240 use bamboo_domain::ReasoningEffort;
2241
2242 let (_temp, storage) = make_storage().await;
2243 let store = LockedSessionStore::new(storage.clone());
2244 let session_id = "cfg-preserve";
2245
2246 let mut initial = fresh(session_id);
2248 initial.add_message(Message::user("hello"));
2249 initial.add_message(Message::assistant("hi", None));
2250 storage.save_session(&initial).await.unwrap();
2251
2252 let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
2254 after_chat.add_message(Message::user("second question"));
2255 storage.save_session(&after_chat).await.unwrap();
2256 assert_eq!(after_chat.messages.len(), 3);
2257
2258 let updated = store
2262 .update_runtime_config(session_id, |s| {
2263 s.reasoning_effort = Some(ReasoningEffort::Max);
2264 })
2265 .await
2266 .unwrap()
2267 .expect("session exists");
2268
2269 assert_eq!(updated.reasoning_effort, Some(ReasoningEffort::Max));
2270 assert_eq!(
2271 updated.messages.len(),
2272 3,
2273 "config patch must not revert a concurrently-appended message"
2274 );
2275
2276 let on_disk = storage.load_session(session_id).await.unwrap().unwrap();
2277 assert_eq!(on_disk.messages.len(), 3);
2278 assert_eq!(on_disk.reasoning_effort, Some(ReasoningEffort::Max));
2279 }
2280
2281 #[tokio::test]
2282 async fn update_runtime_config_returns_none_for_missing_session() {
2283 use bamboo_domain::ReasoningEffort;
2284
2285 let (_temp, storage) = make_storage().await;
2286 let store = LockedSessionStore::new(storage);
2287 let result = store
2288 .update_runtime_config("does-not-exist", |s| {
2289 s.reasoning_effort = Some(ReasoningEffort::Low);
2290 })
2291 .await
2292 .unwrap();
2293 assert!(result.is_none());
2294 }
2295
2296 #[tokio::test]
2297 async fn merge_save_runtime_overwrites_messages_from_stale_snapshot() {
2298 use bamboo_domain::session::types::Message;
2303
2304 let (_temp, storage) = make_storage().await;
2305 let store = LockedSessionStore::new(storage.clone());
2306 let session_id = "stale-clobber";
2307
2308 let mut baseline = fresh(session_id);
2310 baseline.add_message(Message::user("hello"));
2311 storage.save_session(&baseline).await.unwrap();
2312 let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
2313
2314 let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
2316 after_chat.add_message(Message::user("second"));
2317 storage.save_session(&after_chat).await.unwrap();
2318 assert_eq!(
2319 storage
2320 .load_session(session_id)
2321 .await
2322 .unwrap()
2323 .unwrap()
2324 .messages
2325 .len(),
2326 2
2327 );
2328
2329 store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
2331 let after = storage.load_session(session_id).await.unwrap().unwrap();
2332 assert_eq!(
2333 after.messages.len(),
2334 1,
2335 "merge_save_runtime clobbers concurrent appends — this is why config patches must use update_runtime_config"
2336 );
2337 }
2338
2339 #[tokio::test]
2340 async fn stale_runtime_save_cannot_remove_admitted_inbox_transcript() {
2341 use bamboo_domain::session::types::Message;
2342 use bamboo_domain::SessionMessageId;
2343
2344 let (_temp, storage) = make_storage().await;
2345 let store = LockedSessionStore::new(storage.clone());
2346 let session_id = "stale-inbox-preserve";
2347
2348 let mut baseline = fresh(session_id);
2349 let mut base = Message::user("base");
2350 base.id = "base".to_string();
2351 baseline.add_message(base);
2352 storage.save_session(&baseline).await.unwrap();
2353 let mut stale = baseline.clone();
2354 let mut later_assistant = Message::assistant("runner output", None);
2355 later_assistant.id = "later-assistant".to_string();
2356 stale.add_message(later_assistant);
2357
2358 let mut durable = baseline;
2359 let inbox_id = SessionMessageId::parse("durable-inbox-id").unwrap();
2360 let mut admitted = Message::user("durable inbox message");
2361 admitted.id = inbox_id.as_str().to_string();
2362 durable.add_message(admitted);
2363 durable
2364 .session_inbox_admission_mut()
2365 .record(inbox_id.clone(), 7);
2366 storage.save_session(&durable).await.unwrap();
2367
2368 store.merge_save_runtime(&mut stale).await.unwrap();
2369 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2370 let ids = saved
2371 .messages
2372 .iter()
2373 .map(|message| message.id.as_str())
2374 .collect::<Vec<_>>();
2375 assert_eq!(ids, vec!["base", "durable-inbox-id", "later-assistant"]);
2376 assert_eq!(ids.iter().filter(|id| **id == inbox_id.as_str()).count(), 1);
2377 assert!(saved
2378 .session_inbox_admission()
2379 .is_some_and(|state| state.contains(&inbox_id)));
2380 }
2381
2382 #[tokio::test]
2383 async fn stale_runtime_save_preserves_typed_inbox_message_after_cursor_eviction() {
2384 use bamboo_domain::{
2385 SessionMessageEnvelope, SessionMessageId, SESSION_INBOX_ADMITTED_CAPACITY,
2386 };
2387
2388 let (_temp, storage) = make_storage().await;
2389 let store = LockedSessionStore::new(storage.clone());
2390 let session_id = "evicted-inbox-preserve";
2391 let mut durable = fresh(session_id);
2392 let mut envelope = SessionMessageEnvelope::user_input(session_id, "old durable inbox");
2393 envelope.id = SessionMessageId::parse("old-inbox-id").unwrap();
2394 durable.add_message(envelope.to_provider_message().unwrap());
2395 durable
2396 .session_inbox_admission_mut()
2397 .record(envelope.id.clone(), 1);
2398 for sequence in 2..=(SESSION_INBOX_ADMITTED_CAPACITY as u64 + 1) {
2399 durable.session_inbox_admission_mut().record(
2400 SessionMessageId::parse(format!("newer-{sequence}")).unwrap(),
2401 sequence,
2402 );
2403 }
2404 assert!(!durable
2405 .session_inbox_admission()
2406 .unwrap()
2407 .contains(&envelope.id));
2408 storage.save_session(&durable).await.unwrap();
2409
2410 let mut stale = fresh(session_id);
2411 store.merge_save_runtime(&mut stale).await.unwrap();
2412 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2413 assert_eq!(
2414 saved
2415 .messages
2416 .iter()
2417 .filter(|message| message.id == envelope.id.as_str())
2418 .count(),
2419 1
2420 );
2421 }
2422
2423 #[tokio::test]
2424 async fn runtime_final_save_cannot_resurrect_a_consumed_clarification() {
2425 use bamboo_domain::session::types::Message;
2426
2427 let (_temp, storage) = make_storage().await;
2428 let store = LockedSessionStore::new(storage.clone());
2429 let session_id = "consumed-clarification-final-save";
2430
2431 let mut suspended = fresh(session_id);
2432 suspended.add_message(Message::tool_result(
2433 "call-1",
2434 r#"{"status":"awaiting_clarification"}"#,
2435 ));
2436 suspended.set_pending_question(
2437 "call-1".to_string(),
2438 "ConclusionWithOptions".to_string(),
2439 "Choose".to_string(),
2440 vec!["A".to_string()],
2441 false,
2442 );
2443 suspended.metadata.insert(
2444 "runtime.suspend_reason".to_string(),
2445 "awaiting_clarification".to_string(),
2446 );
2447 storage.save_session(&suspended).await.unwrap();
2448 let mut stale_runner = suspended.clone();
2449
2450 let mut answered = suspended;
2451 answered.clear_pending_question();
2452 answered.metadata.remove("runtime.suspend_reason");
2453 answered.metadata.insert(
2454 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
2455 r#"["call-1"]"#.to_string(),
2456 );
2457 answered.metadata.insert(
2458 "clarification_resume_pending".to_string(),
2459 "true".to_string(),
2460 );
2461 answered.metadata.insert(
2462 "conclusion_with_options_resume_pending".to_string(),
2463 "true".to_string(),
2464 );
2465 answered.metadata.insert(
2466 "execute.startup_handoff_at".to_string(),
2467 "2026-08-10T09:00:00.000Z".to_string(),
2468 );
2469 let answer = answered
2470 .messages
2471 .iter_mut()
2472 .find(|message| message.tool_call_id.as_deref() == Some("call-1"))
2473 .unwrap();
2474 answer.content = "Selected response: A".to_string();
2475 storage.save_session(&answered).await.unwrap();
2476
2477 store.merge_save_runtime(&mut stale_runner).await.unwrap();
2478
2479 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2480 assert!(saved.pending_question.is_none());
2481 assert!(!saved.metadata.contains_key("runtime.suspend_reason"));
2482 assert_eq!(
2483 saved
2484 .metadata
2485 .get("clarification_resume_pending")
2486 .map(String::as_str),
2487 Some("true")
2488 );
2489 assert_eq!(
2490 saved
2491 .metadata
2492 .get("execute.startup_handoff_at")
2493 .map(String::as_str),
2494 Some("2026-08-10T09:00:00.000Z")
2495 );
2496 let answers = saved
2497 .messages
2498 .iter()
2499 .filter(|message| message.tool_call_id.as_deref() == Some("call-1"))
2500 .collect::<Vec<_>>();
2501 assert_eq!(answers.len(), 1);
2502 assert_eq!(answers[0].content, "Selected response: A");
2503 }
2504
2505 #[tokio::test]
2506 async fn runtime_checkpoint_cannot_resurrect_a_consumed_clarification() {
2507 use bamboo_domain::session::types::Message;
2508
2509 let (_temp, storage) = make_storage().await;
2510 let store = LockedSessionStore::new(storage.clone());
2511 let session_id = "consumed-clarification-checkpoint";
2512 let mut stale_runner = fresh(session_id);
2513 stale_runner.add_message(Message::tool_result("call-1", "waiting"));
2514 stale_runner.set_pending_question(
2515 "call-1".to_string(),
2516 "ConclusionWithOptions".to_string(),
2517 "Choose".to_string(),
2518 vec!["A".to_string()],
2519 false,
2520 );
2521 stale_runner.metadata.insert(
2522 "runtime.suspend_reason".to_string(),
2523 "awaiting_clarification".to_string(),
2524 );
2525
2526 let mut answered = stale_runner.clone();
2527 answered.clear_pending_question();
2528 answered.metadata.remove("runtime.suspend_reason");
2529 answered.metadata.insert(
2530 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
2531 r#"["call-1"]"#.to_string(),
2532 );
2533 answered.metadata.insert(
2534 "clarification_resume_pending".to_string(),
2535 "true".to_string(),
2536 );
2537 answered.messages[0].content = "Selected response: A".to_string();
2538 storage.save_session(&answered).await.unwrap();
2539
2540 store
2541 .checkpoint_runtime_session(&mut stale_runner)
2542 .await
2543 .unwrap();
2544
2545 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2546 assert!(saved.pending_question.is_none());
2547 assert!(!saved.metadata.contains_key("runtime.suspend_reason"));
2548 assert_eq!(saved.messages.len(), 1);
2549 assert_eq!(saved.messages[0].content, "Selected response: A");
2550 }
2551
2552 #[tokio::test]
2553 async fn runtime_final_save_does_not_consume_a_new_reused_permission_occurrence() {
2554 let (_temp, storage) = make_storage().await;
2555 let store = LockedSessionStore::new(storage.clone());
2556 let session_id = "reused-permission-final-save";
2557
2558 let mut durable = fresh(session_id);
2559 durable.add_message(typed_permission_result(
2560 "reused-call",
2561 "old-result",
2562 "generation-old",
2563 "Selected response: Approve",
2564 ));
2565 durable.metadata.insert(
2566 CONSUMED_RESPONSE_OCCURRENCES_KEY.to_string(),
2567 serde_json::to_string(&vec![ResponseOccurrence {
2568 tool_call_id: "reused-call".to_string(),
2569 tool_result_message_id: "old-result".to_string(),
2570 permission_generation: Some("generation-old".to_string()),
2571 }])
2572 .unwrap(),
2573 );
2574 storage.save_session(&durable).await.unwrap();
2575
2576 let mut new_runner = durable;
2577 new_runner.add_message(typed_permission_result(
2578 "reused-call",
2579 "new-result",
2580 "generation-new",
2581 "waiting for the new decision",
2582 ));
2583 new_runner.set_pending_question(
2584 "reused-call".to_string(),
2585 "Permission".to_string(),
2586 "Approve new operation?".to_string(),
2587 vec!["Approve".to_string(), "Deny".to_string()],
2588 false,
2589 );
2590 new_runner.metadata.insert(
2591 "runtime.suspend_reason".to_string(),
2592 "awaiting_permission_approval".to_string(),
2593 );
2594
2595 store.merge_save_runtime(&mut new_runner).await.unwrap();
2596
2597 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2598 assert_eq!(
2599 saved
2600 .pending_question
2601 .as_ref()
2602 .map(|pending| pending.tool_call_id.as_str()),
2603 Some("reused-call")
2604 );
2605 assert_eq!(
2606 saved.messages.last().map(|message| message.id.as_str()),
2607 Some("new-result")
2608 );
2609 assert_eq!(
2610 saved.messages.last().unwrap().content,
2611 "waiting for the new decision"
2612 );
2613 assert_eq!(
2614 latest_response_occurrence(&saved, "reused-call")
2615 .and_then(|occurrence| occurrence.permission_generation),
2616 Some("generation-new".to_string())
2617 );
2618 }
2619
2620 #[tokio::test]
2621 async fn legacy_consumed_id_does_not_consume_a_new_reused_occurrence_after_upgrade() {
2622 let (_temp, storage) = make_storage().await;
2623 let store = LockedSessionStore::new(storage.clone());
2624 let session_id = "legacy-reused-permission-final-save";
2625
2626 let mut durable = fresh(session_id);
2627 durable.add_message(typed_permission_result(
2628 "reused-call",
2629 "old-result",
2630 "generation-old",
2631 "Selected response: Approve",
2632 ));
2633 durable.metadata.insert(
2634 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
2635 r#"["reused-call"]"#.to_string(),
2636 );
2637 storage.save_session(&durable).await.unwrap();
2638
2639 let mut new_runner = durable;
2640 new_runner.add_message(typed_permission_result(
2641 "reused-call",
2642 "new-result",
2643 "generation-new",
2644 "waiting for the new decision",
2645 ));
2646 new_runner.set_pending_question(
2647 "reused-call".to_string(),
2648 "Permission".to_string(),
2649 "Approve new operation?".to_string(),
2650 vec!["Approve".to_string(), "Deny".to_string()],
2651 false,
2652 );
2653
2654 store.merge_save_runtime(&mut new_runner).await.unwrap();
2655
2656 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2657 assert_eq!(
2658 saved
2659 .pending_question
2660 .as_ref()
2661 .map(|pending| pending.tool_call_id.as_str()),
2662 Some("reused-call")
2663 );
2664 assert_eq!(
2665 saved.messages.last().map(|message| message.id.as_str()),
2666 Some("new-result")
2667 );
2668 assert_eq!(
2669 latest_response_occurrence(&saved, "reused-call")
2670 .and_then(|occurrence| occurrence.permission_generation),
2671 Some("generation-new".to_string())
2672 );
2673 }
2674
2675 #[tokio::test]
2676 async fn runtime_checkpoint_does_not_consume_a_new_reused_permission_occurrence() {
2677 let (_temp, storage) = make_storage().await;
2678 let store = LockedSessionStore::new(storage.clone());
2679 let session_id = "reused-permission-checkpoint";
2680
2681 let mut durable = fresh(session_id);
2682 durable.add_message(typed_permission_result(
2683 "reused-call",
2684 "old-result",
2685 "generation-old",
2686 "Selected response: Deny",
2687 ));
2688 durable.metadata.insert(
2689 CONSUMED_RESPONSE_OCCURRENCES_KEY.to_string(),
2690 serde_json::to_string(&vec![ResponseOccurrence {
2691 tool_call_id: "reused-call".to_string(),
2692 tool_result_message_id: "old-result".to_string(),
2693 permission_generation: Some("generation-old".to_string()),
2694 }])
2695 .unwrap(),
2696 );
2697 storage.save_session(&durable).await.unwrap();
2698
2699 let mut new_runner = durable;
2700 new_runner.add_message(typed_permission_result(
2701 "reused-call",
2702 "new-result",
2703 "generation-new",
2704 "waiting for the new decision",
2705 ));
2706 new_runner.set_pending_question(
2707 "reused-call".to_string(),
2708 "Permission".to_string(),
2709 "Approve new operation?".to_string(),
2710 vec!["Approve".to_string(), "Deny".to_string()],
2711 false,
2712 );
2713
2714 store
2715 .checkpoint_runtime_session(&mut new_runner)
2716 .await
2717 .unwrap();
2718
2719 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2720 assert_eq!(
2721 saved
2722 .pending_question
2723 .as_ref()
2724 .map(|pending| pending.tool_call_id.as_str()),
2725 Some("reused-call")
2726 );
2727 assert_eq!(
2728 saved.messages.last().map(|message| message.id.as_str()),
2729 Some("new-result")
2730 );
2731 assert_eq!(
2732 latest_response_occurrence(&saved, "reused-call")
2733 .and_then(|occurrence| occurrence.permission_generation),
2734 Some("generation-new".to_string())
2735 );
2736 }
2737
2738 #[tokio::test]
2739 async fn consumed_permission_adoption_keeps_reexecute_id_and_generation_paired() {
2740 let (_temp, storage) = make_storage().await;
2741 let store = LockedSessionStore::new(storage.clone());
2742 let session_id = "consumed-permission-control-pair";
2743
2744 let mut stale_runner = fresh(session_id);
2745 stale_runner.add_message(typed_permission_result(
2746 "call-1",
2747 "result-1",
2748 "generation-1",
2749 "waiting",
2750 ));
2751 stale_runner.set_pending_question(
2752 "call-1".to_string(),
2753 "Permission".to_string(),
2754 "Approve?".to_string(),
2755 vec!["Approve".to_string(), "Deny".to_string()],
2756 false,
2757 );
2758
2759 let mut answered = stale_runner.clone();
2760 answered.clear_pending_question();
2761 answered.messages[0].content = "Selected response: Approve".to_string();
2762 answered.metadata.insert(
2763 CONSUMED_RESPONSE_OCCURRENCES_KEY.to_string(),
2764 serde_json::to_string(&vec![ResponseOccurrence {
2765 tool_call_id: "call-1".to_string(),
2766 tool_result_message_id: "result-1".to_string(),
2767 permission_generation: Some("generation-1".to_string()),
2768 }])
2769 .unwrap(),
2770 );
2771 answered.metadata.insert(
2772 "permission.reexecute_tool_call_id".to_string(),
2773 "call-1".to_string(),
2774 );
2775 answered.metadata.insert(
2776 "permission.reexecute_request_generation".to_string(),
2777 "generation-1".to_string(),
2778 );
2779 storage.save_session(&answered).await.unwrap();
2780
2781 store.merge_save_runtime(&mut stale_runner).await.unwrap();
2782
2783 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2784 assert!(saved.pending_question.is_none());
2785 assert_eq!(
2786 saved
2787 .metadata
2788 .get("permission.reexecute_tool_call_id")
2789 .map(String::as_str),
2790 Some("call-1")
2791 );
2792 assert_eq!(
2793 saved
2794 .metadata
2795 .get("permission.reexecute_request_generation")
2796 .map(String::as_str),
2797 Some("generation-1")
2798 );
2799 }
2800
2801 #[tokio::test]
2802 async fn checkpoint_runtime_session_preserves_disk_suffix_and_appends_live_messages() {
2803 use bamboo_domain::session::types::Message;
2804
2805 let (_temp, storage) = make_storage().await;
2806 let store = LockedSessionStore::new(storage.clone());
2807 let session_id = "checkpoint-no-shrink";
2808
2809 let mut baseline = fresh(session_id);
2810 baseline.add_message(Message::user("base"));
2811 storage.save_session(&baseline).await.unwrap();
2812 let mut runner_snapshot = baseline.clone();
2813
2814 let mut durable = baseline;
2815 let mut disk_only = Message::user("concurrent injected message");
2816 disk_only.id = "disk-only".to_string();
2817 durable.add_message(disk_only);
2818 storage.save_session(&durable).await.unwrap();
2819
2820 let mut live_only = Message::assistant("partial runner output", None);
2821 live_only.id = "live-only".to_string();
2822 runner_snapshot.add_message(live_only);
2823
2824 store
2825 .checkpoint_runtime_session(&mut runner_snapshot)
2826 .await
2827 .unwrap();
2828
2829 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2830 let ids = saved
2831 .messages
2832 .iter()
2833 .map(|message| message.id.as_str())
2834 .collect::<Vec<_>>();
2835 assert_eq!(
2836 ids,
2837 vec![durable.messages[0].id.as_str(), "disk-only", "live-only"]
2838 );
2839 assert_eq!(runner_snapshot.messages.len(), saved.messages.len());
2840 assert_eq!(runner_snapshot.messages[1].id, saved.messages[1].id);
2841 assert_eq!(runner_snapshot.messages[2].id, saved.messages[2].id);
2842 assert_eq!(saved.messages[1].content, "concurrent injected message");
2843 assert_eq!(saved.messages[2].content, "partial runner output");
2844 }
2845
2846 #[tokio::test]
2847 async fn runtime_only_save_preserves_checkpointed_ledger_and_publishes_merged_state() {
2848 let (_temp, storage) = make_storage().await;
2849 let store = LockedSessionStore::new(storage.clone());
2850 let session_id = "runtime-only-ledger-race";
2851 let baseline = fresh(session_id);
2852 storage.save_session(&baseline).await.unwrap();
2853 let mut stale_control = storage.load_session(session_id).await.unwrap().unwrap();
2854
2855 let mut runner = baseline;
2856 runner.model_context_state = Some(ledger_state(1, "runner-l1"));
2857 store.checkpoint_runtime_session(&mut runner).await.unwrap();
2858
2859 stale_control.metadata.insert(
2860 "runtime.suspend_reason".to_string(),
2861 "waiting_for_children".to_string(),
2862 );
2863 let published = std::sync::Arc::new(std::sync::Mutex::new(None));
2864 let published_clone = published.clone();
2865 store
2866 .save_runtime_only_and_publish(&mut stale_control, move |saved| {
2867 *published_clone.lock().unwrap() = Some(saved.clone());
2868 })
2869 .await
2870 .unwrap();
2871
2872 let expected = runner.model_context_state.clone();
2873 assert_eq!(stale_control.model_context_state, expected);
2874 assert_eq!(
2875 published
2876 .lock()
2877 .unwrap()
2878 .as_ref()
2879 .unwrap()
2880 .model_context_state,
2881 expected
2882 );
2883 let sidecar = storage
2884 .load_runtime_control_plane(session_id)
2885 .await
2886 .unwrap()
2887 .unwrap();
2888 assert_eq!(sidecar.model_context_state, expected);
2889 assert_eq!(
2890 sidecar
2891 .metadata
2892 .get("runtime.suspend_reason")
2893 .map(String::as_str),
2894 Some("waiting_for_children")
2895 );
2896 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
2897 assert_eq!(reloaded.model_context_state, expected);
2898 }
2899
2900 #[tokio::test]
2901 async fn full_runtime_save_preserves_newer_ledger_but_commits_control_mutation() {
2902 let (_temp, storage) = make_storage().await;
2903 let store = LockedSessionStore::new(storage.clone());
2904 let session_id = "full-save-ledger-race";
2905 let baseline = fresh(session_id);
2906 storage.save_session(&baseline).await.unwrap();
2907 let mut stale = storage.load_session(session_id).await.unwrap().unwrap();
2908
2909 let mut runner = baseline;
2910 runner.model_context_state = Some(ledger_state(1, "runner-l1"));
2911 store.checkpoint_runtime_session(&mut runner).await.unwrap();
2912
2913 stale
2914 .metadata
2915 .insert("activated_tools".to_string(), "[\"search\"]".to_string());
2916 let published = std::sync::Arc::new(std::sync::Mutex::new(None));
2917 let published_clone = published.clone();
2918 store
2919 .merge_save_runtime_and_publish(&mut stale, move |saved, committed| {
2920 assert!(committed);
2921 *published_clone.lock().unwrap() = Some(saved.clone());
2922 })
2923 .await
2924 .unwrap();
2925
2926 let expected = runner.model_context_state.clone();
2927 assert_eq!(stale.model_context_state, expected);
2928 assert_eq!(
2929 published
2930 .lock()
2931 .unwrap()
2932 .as_ref()
2933 .unwrap()
2934 .model_context_state,
2935 expected
2936 );
2937 let sidecar = storage
2938 .load_runtime_control_plane(session_id)
2939 .await
2940 .unwrap()
2941 .unwrap();
2942 assert_eq!(sidecar.model_context_state, expected);
2943 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
2944 assert_eq!(reloaded.model_context_state, expected);
2945 assert_eq!(
2946 reloaded.metadata.get("activated_tools").map(String::as_str),
2947 Some("[\"search\"]")
2948 );
2949 }
2950
2951 #[tokio::test]
2952 async fn newer_explicit_epoch_reset_wins_an_ordinary_full_runtime_save() {
2953 let (_temp, storage) = make_storage().await;
2954 let store = LockedSessionStore::new(storage.clone());
2955 let session_id = "full-save-ledger-reset";
2956 let mut durable = fresh(session_id);
2957 durable.model_context_state = Some(ledger_state(1, "runner-l1"));
2958 storage.save_session(&durable).await.unwrap();
2959
2960 let mut compression = storage.load_session(session_id).await.unwrap().unwrap();
2961 compression.reset_model_context_epoch(bamboo_domain::ModelContextResetReason::Compression);
2962 let reset = compression.model_context_state.clone();
2963 assert_eq!(reset.as_ref().unwrap().state_revision, 2);
2964 store.merge_save_runtime(&mut compression).await.unwrap();
2965
2966 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
2967 assert_eq!(reloaded.model_context_state, reset);
2968 assert_eq!(
2969 reloaded
2970 .model_context_state
2971 .as_ref()
2972 .and_then(|state| state.last_reset_reason),
2973 Some(bamboo_domain::ModelContextResetReason::Compression)
2974 );
2975 }
2976
2977 #[tokio::test]
2978 async fn checkpoint_rejects_equal_revision_divergence_without_overwriting_disk() {
2979 let (_temp, storage) = make_storage().await;
2980 let store = LockedSessionStore::new(storage.clone());
2981 let session_id = "ledger-checkpoint-cas";
2982 let mut baseline = fresh(session_id);
2983 baseline.model_context_state = Some(ledger_state(1, "runner-l1"));
2984 storage.save_session(&baseline).await.unwrap();
2985
2986 let mut first = baseline.clone();
2987 first.reset_model_context_epoch(bamboo_domain::ModelContextResetReason::Compression);
2988 let mut conflicting = baseline;
2989 conflicting.reset_model_context_epoch(bamboo_domain::ModelContextResetReason::Rollback);
2990 assert_eq!(
2991 first.model_context_state.as_ref().unwrap().state_revision,
2992 conflicting
2993 .model_context_state
2994 .as_ref()
2995 .unwrap()
2996 .state_revision
2997 );
2998
2999 store.checkpoint_runtime_session(&mut first).await.unwrap();
3000 let error = store
3001 .checkpoint_runtime_session(&mut conflicting)
3002 .await
3003 .unwrap_err();
3004 assert_eq!(error.kind(), std::io::ErrorKind::WouldBlock);
3005 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
3006 assert_eq!(reloaded.model_context_state, first.model_context_state);
3007 }
3008
3009 #[tokio::test]
3010 async fn activation_checkpoint_clears_presentation_without_shrinking_concurrent_turn() {
3011 use bamboo_domain::session::runtime_state::{
3012 AgentRuntimeState, AgentStatusState, WaitingForChildrenState,
3013 };
3014 use bamboo_domain::session::types::Message;
3015
3016 let (_temp, storage) = make_storage().await;
3017 let store = LockedSessionStore::new(storage.clone());
3018 let session_id = "activation-no-shrink";
3019 let mut baseline = fresh(session_id);
3020 baseline.add_message(Message::user("base"));
3021 let mut state = AgentRuntimeState::new("activation-run");
3022 state.status = AgentStatusState::Suspended;
3023 state.waiting_for_children = Some(WaitingForChildrenState::for_children(
3024 vec!["child-1".to_string()],
3025 bamboo_domain::session::runtime_state::ChildWaitPolicy::All,
3026 chrono::Utc::now(),
3027 ));
3028 baseline.agent_runtime_state = Some(state);
3029 baseline.metadata.insert(
3030 "runtime.suspend_reason".to_string(),
3031 "waiting_for_children".to_string(),
3032 );
3033 storage.save_session(&baseline).await.unwrap();
3034 let mut activation_snapshot = baseline.clone();
3035
3036 let mut concurrent = baseline;
3037 let mut normal = Message::assistant("normal concurrent answer", None);
3038 normal.id = "normal-concurrent".to_string();
3039 concurrent.add_message(normal);
3040 storage.save_session(&concurrent).await.unwrap();
3041
3042 let state = activation_snapshot.agent_runtime_state.as_mut().unwrap();
3043 state.status = AgentStatusState::Idle;
3044 state.suspension = None;
3045 activation_snapshot
3046 .metadata
3047 .remove("runtime.suspend_reason");
3048 store
3049 .checkpoint_runtime_session(&mut activation_snapshot)
3050 .await
3051 .unwrap();
3052
3053 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3054 assert!(saved
3055 .messages
3056 .iter()
3057 .any(|message| message.id == "normal-concurrent"));
3058 let state = saved.agent_runtime_state.unwrap();
3059 assert_eq!(state.status, AgentStatusState::Idle);
3060 assert!(state.waiting_for_children.is_some());
3061 assert!(!saved.metadata.contains_key("runtime.suspend_reason"));
3062 }
3063
3064 #[tokio::test]
3065 async fn merge_save_runtime_preserves_disk_authoritative_metadata_with_single_load() {
3066 let (_temp, storage) = make_storage().await;
3072 let store = LockedSessionStore::new(storage.clone());
3073 let session_id = "runtime-merge-meta";
3074
3075 let mut baseline = fresh(session_id);
3077 baseline.title = "Auto Title".to_string();
3078 baseline.metadata_version = 0;
3079 storage.save_session(&baseline).await.unwrap();
3080
3081 let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
3083
3084 let mut renamed = storage.load_session(session_id).await.unwrap().unwrap();
3086 renamed.title = "User Renamed".to_string();
3087 renamed.title_version = 1;
3088 renamed.pinned = true;
3089 renamed.metadata_version = 1;
3090 store.commit_metadata(&renamed).await.unwrap();
3091
3092 stale_snapshot.title = "Auto Title".to_string();
3094 store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
3095
3096 let after = storage.load_session(session_id).await.unwrap().unwrap();
3097 assert_eq!(after.title, "User Renamed");
3098 assert!(after.pinned);
3099 assert_eq!(after.metadata_version, 1);
3100 assert_eq!(stale_snapshot.title, "User Renamed");
3102 assert_eq!(stale_snapshot.metadata_version, 1);
3103 }
3104
3105 #[tokio::test]
3106 async fn merge_save_runtime_preserves_durable_workflow_run_index_from_stale_runner() {
3107 let (_temp, storage) = make_storage().await;
3108 let store = LockedSessionStore::new(storage.clone());
3109 let session_id = "runtime-workflow-run-index";
3110
3111 let baseline = fresh(session_id);
3112 storage.save_session(&baseline).await.unwrap();
3113 let mut stale_runner = storage.load_session(session_id).await.unwrap().unwrap();
3114
3115 store
3116 .update_runtime_config(session_id, |session| {
3117 session.metadata.insert(
3118 "workflow.run_ids.v1".to_string(),
3119 r#"["http-started-run"]"#.to_string(),
3120 );
3121 })
3122 .await
3123 .unwrap()
3124 .expect("session exists");
3125
3126 store.merge_save_runtime(&mut stale_runner).await.unwrap();
3127
3128 assert_eq!(
3129 stale_runner
3130 .metadata
3131 .get("workflow.run_ids.v1")
3132 .map(String::as_str),
3133 Some(r#"["http-started-run"]"#)
3134 );
3135 let durable = storage.load_session(session_id).await.unwrap().unwrap();
3136 assert_eq!(
3137 durable
3138 .metadata
3139 .get("workflow.run_ids.v1")
3140 .map(String::as_str),
3141 Some(r#"["http-started-run"]"#)
3142 );
3143 }
3144
3145 #[tokio::test]
3149 async fn merge_save_runtime_adopts_disk_bypass_permissions() {
3150 use bamboo_domain::AgentRuntimeState;
3151
3152 let (_temp, storage) = make_storage().await;
3153 let store = LockedSessionStore::new(storage.clone());
3154 let session_id = "runtime-bypass";
3155
3156 let baseline = fresh(session_id);
3158 storage.save_session(&baseline).await.unwrap();
3159
3160 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
3162 loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
3163
3164 store
3166 .update_runtime_config(session_id, |s| {
3167 s.agent_runtime_state
3168 .get_or_insert_with(AgentRuntimeState::default)
3169 .bypass_permissions = true;
3170 })
3171 .await
3172 .unwrap()
3173 .expect("session exists");
3174
3175 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
3178
3179 let after = storage.load_session(session_id).await.unwrap().unwrap();
3180 assert!(
3181 after
3182 .agent_runtime_state
3183 .as_ref()
3184 .is_some_and(|s| s.bypass_permissions),
3185 "disk bypass=ON must survive a stale runtime save (#540)"
3186 );
3187 assert!(loop_snapshot
3189 .agent_runtime_state
3190 .as_ref()
3191 .is_some_and(|s| s.bypass_permissions));
3192 }
3193
3194 #[tokio::test]
3197 async fn merge_save_runtime_adopts_disk_auto_permission_mode() {
3198 use bamboo_domain::{AgentRuntimeState, SessionPermissionMode};
3199
3200 let (_temp, storage) = make_storage().await;
3201 let store = LockedSessionStore::new(storage.clone());
3202 let session_id = "runtime-auto";
3203
3204 storage.save_session(&fresh(session_id)).await.unwrap();
3205 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
3206 loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
3207 loop_snapshot.metadata.insert(
3208 "permission.requested_mode".to_string(),
3209 "default".to_string(),
3210 );
3211 loop_snapshot.metadata.insert(
3212 "permission.effective_mode".to_string(),
3213 "default".to_string(),
3214 );
3215 loop_snapshot.metadata.insert(
3216 "permission.executor_mapping".to_string(),
3217 "bamboo_runtime:default".to_string(),
3218 );
3219
3220 store
3221 .update_runtime_config(session_id, |session| {
3222 session
3223 .agent_runtime_state
3224 .get_or_insert_with(AgentRuntimeState::default)
3225 .set_permission_mode(SessionPermissionMode::Auto);
3226 session
3227 .metadata
3228 .insert("permission.policy_revision".to_string(), "12".to_string());
3229 session
3230 .metadata
3231 .insert("permission.requested_mode".to_string(), "auto".to_string());
3232 session
3233 .metadata
3234 .insert("permission.effective_mode".to_string(), "auto".to_string());
3235 session.metadata.insert(
3236 "permission.executor_mapping".to_string(),
3237 "bamboo_runtime:auto".to_string(),
3238 );
3239 session.metadata.insert(
3240 "permission.transitioned_at".to_string(),
3241 "2026-07-31T12:00:00Z".to_string(),
3242 );
3243 session.metadata_version = session.metadata_version.saturating_add(1);
3244 })
3245 .await
3246 .unwrap()
3247 .expect("session exists");
3248
3249 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
3250
3251 let durable = storage.load_session(session_id).await.unwrap().unwrap();
3252 for state in [
3253 durable.agent_runtime_state.as_ref(),
3254 loop_snapshot.agent_runtime_state.as_ref(),
3255 ] {
3256 assert_eq!(
3257 state.map(AgentRuntimeState::effective_permission_mode),
3258 Some(SessionPermissionMode::Auto)
3259 );
3260 }
3261 for session in [&durable, &loop_snapshot] {
3262 assert_eq!(
3263 session.metadata.get("permission.policy_revision"),
3264 Some(&"12".to_string())
3265 );
3266 assert_eq!(
3267 session.metadata.get("permission.requested_mode"),
3268 Some(&"auto".to_string())
3269 );
3270 assert_eq!(
3271 session.metadata.get("permission.effective_mode"),
3272 Some(&"auto".to_string())
3273 );
3274 assert_eq!(
3275 session.metadata.get("permission.executor_mapping"),
3276 Some(&"bamboo_runtime:auto".to_string())
3277 );
3278 assert_eq!(
3279 session.metadata.get("permission.transitioned_at"),
3280 Some(&"2026-07-31T12:00:00Z".to_string())
3281 );
3282 }
3283 }
3284
3285 #[tokio::test]
3288 async fn merge_save_runtime_adopts_disk_bypass_off() {
3289 use bamboo_domain::AgentRuntimeState;
3290
3291 let (_temp, storage) = make_storage().await;
3292 let store = LockedSessionStore::new(storage.clone());
3293 let session_id = "runtime-bypass-off";
3294
3295 let mut baseline = fresh(session_id);
3297 let on_state = AgentRuntimeState {
3298 bypass_permissions: true,
3299 ..AgentRuntimeState::default()
3300 };
3301 baseline.agent_runtime_state = Some(on_state);
3302 storage.save_session(&baseline).await.unwrap();
3303
3304 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
3306
3307 store
3309 .update_runtime_config(session_id, |s| {
3310 s.agent_runtime_state
3311 .get_or_insert_with(AgentRuntimeState::default)
3312 .bypass_permissions = false;
3313 })
3314 .await
3315 .unwrap()
3316 .expect("session exists");
3317
3318 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
3319
3320 let after = storage.load_session(session_id).await.unwrap().unwrap();
3321 assert!(
3322 !after
3323 .agent_runtime_state
3324 .as_ref()
3325 .is_some_and(|s| s.bypass_permissions),
3326 "disk bypass=OFF must survive a stale runtime save (#540)"
3327 );
3328 }
3329
3330 #[tokio::test]
3333 async fn save_runtime_authoritative_flags_persists_in_memory_posture_and_audit() {
3334 use bamboo_domain::AgentRuntimeState;
3335
3336 let (_temp, storage) = make_storage().await;
3337 let store = LockedSessionStore::new(storage.clone());
3338 let session_id = "child-reseed";
3339
3340 let mut baseline = fresh(session_id);
3342 let on_state = AgentRuntimeState {
3343 bypass_permissions: true,
3344 ..AgentRuntimeState::default()
3345 };
3346 baseline.agent_runtime_state = Some(on_state);
3347 for (key, value) in [
3348 ("permission.policy_revision", "12"),
3349 ("permission.requested_mode", "bypass"),
3350 ("permission.effective_mode", "bypass"),
3351 ("permission.executor_mapping", "bamboo_runtime:bypass"),
3352 ("permission.transitioned_at", "2026-07-31T12:00:00Z"),
3353 ] {
3354 baseline.metadata.insert(key.to_string(), value.to_string());
3355 }
3356 storage.save_session(&baseline).await.unwrap();
3357
3358 let mut child = storage.load_session(session_id).await.unwrap().unwrap();
3361 child
3362 .agent_runtime_state
3363 .get_or_insert_with(AgentRuntimeState::default)
3364 .bypass_permissions = false;
3365 for (key, value) in [
3366 ("permission.policy_revision", "13"),
3367 ("permission.requested_mode", "default"),
3368 ("permission.effective_mode", "default"),
3369 ("permission.executor_mapping", "bamboo_runtime:default"),
3370 ("permission.transitioned_at", "2026-07-31T12:01:00Z"),
3371 ] {
3372 child.metadata.insert(key.to_string(), value.to_string());
3373 }
3374
3375 store
3377 .save_runtime_authoritative_flags(&mut child)
3378 .await
3379 .unwrap();
3380
3381 let after = storage.load_session(session_id).await.unwrap().unwrap();
3382 assert!(
3383 !after
3384 .agent_runtime_state
3385 .as_ref()
3386 .is_some_and(|s| s.bypass_permissions),
3387 "authoritative re-seed of bypass=OFF must persist, not be reverted (#540/#74)"
3388 );
3389 for (key, value) in [
3390 ("permission.policy_revision", "13"),
3391 ("permission.requested_mode", "default"),
3392 ("permission.effective_mode", "default"),
3393 ("permission.executor_mapping", "bamboo_runtime:default"),
3394 ("permission.transitioned_at", "2026-07-31T12:01:00Z"),
3395 ] {
3396 assert_eq!(after.metadata.get(key).map(String::as_str), Some(value));
3397 }
3398 }
3399
3400 #[tokio::test]
3402 async fn merge_save_runtime_leaves_bypass_when_disk_has_no_runtime_state() {
3403 use bamboo_domain::AgentRuntimeState;
3404
3405 let (_temp, storage) = make_storage().await;
3406 let store = LockedSessionStore::new(storage.clone());
3407 let session_id = "no-runtime-state";
3408
3409 let baseline = fresh(session_id);
3411 assert!(baseline.agent_runtime_state.is_none());
3412 storage.save_session(&baseline).await.unwrap();
3413
3414 let mut running = storage.load_session(session_id).await.unwrap().unwrap();
3416 let on_state = AgentRuntimeState {
3417 bypass_permissions: true,
3418 ..AgentRuntimeState::default()
3419 };
3420 running.agent_runtime_state = Some(on_state);
3421
3422 store.merge_save_runtime(&mut running).await.unwrap();
3423
3424 assert!(
3425 running
3426 .agent_runtime_state
3427 .as_ref()
3428 .is_some_and(|s| s.bypass_permissions),
3429 "a runtime-state-less disk copy must not force bypass OFF (#540)"
3430 );
3431 }
3432
3433 #[tokio::test]
3436 async fn merge_preserves_disk_title_when_versions_equal() {
3437 let (_temp, storage) = make_storage().await;
3438 let session_id = "merge-equal";
3439
3440 let mut on_disk = fresh(session_id);
3441 on_disk.title = "User Set This".to_string();
3442 on_disk.title_version = 0;
3443 on_disk.title_generated = true;
3444 on_disk.metadata_version = 0;
3445 storage.save_session(&on_disk).await.unwrap();
3446
3447 let mut runtime_copy = fresh(session_id);
3448 runtime_copy.title = "Stale Default".to_string();
3449 runtime_copy.title_version = 0;
3450 runtime_copy.title_generated = false;
3451 runtime_copy.metadata_version = 0;
3452 runtime_copy.messages = vec![];
3453
3454 merge_save_session(&storage, &mut runtime_copy)
3455 .await
3456 .unwrap();
3457
3458 let after = storage.load_session(session_id).await.unwrap().unwrap();
3459 assert_eq!(after.title, "User Set This");
3460 assert_eq!(after.title_version, 0);
3461 assert!(after.title_generated);
3462 assert_eq!(runtime_copy.title, "User Set This");
3463 assert!(runtime_copy.title_generated);
3464 }
3465
3466 #[tokio::test]
3467 async fn merge_preserves_disk_when_disk_version_higher() {
3468 let (_temp, storage) = make_storage().await;
3469 let session_id = "merge-higher";
3470
3471 let mut on_disk = fresh(session_id);
3472 on_disk.title = "User Title v3".to_string();
3473 on_disk.title_version = 3;
3474 on_disk.metadata_version = 5;
3475 storage.save_session(&on_disk).await.unwrap();
3476
3477 let mut runtime_copy = fresh(session_id);
3478 runtime_copy.title = "Stale".to_string();
3479 runtime_copy.title_version = 1;
3480 runtime_copy.metadata_version = 0;
3481
3482 merge_save_session(&storage, &mut runtime_copy)
3483 .await
3484 .unwrap();
3485
3486 let after = storage.load_session(session_id).await.unwrap().unwrap();
3487 assert_eq!(after.title, "User Title v3");
3488 assert_eq!(after.title_version, 3);
3489 assert_eq!(after.metadata_version, 5);
3490 }
3491
3492 #[tokio::test]
3493 async fn merge_now_preserves_disk_pinned_in_metadata_group() {
3494 let (_temp, storage) = make_storage().await;
3495 let session_id = "pinned-merge";
3496
3497 let mut on_disk = fresh(session_id);
3498 on_disk.pinned = true;
3499 on_disk.metadata_version = 2;
3500 storage.save_session(&on_disk).await.unwrap();
3501
3502 let mut runtime_copy = fresh(session_id);
3503 runtime_copy.pinned = false;
3504 runtime_copy.metadata_version = 0;
3505
3506 merge_save_session(&storage, &mut runtime_copy)
3507 .await
3508 .unwrap();
3509
3510 let after = storage.load_session(session_id).await.unwrap().unwrap();
3511 assert!(
3512 after.pinned,
3513 "disk pinned=true should win over runtime false"
3514 );
3515 assert_eq!(after.metadata_version, 2);
3516 }
3517
3518 #[tokio::test]
3519 async fn merge_keeps_in_memory_when_session_version_higher() {
3520 let (_temp, storage) = make_storage().await;
3521 let session_id = "merge-bumped";
3522
3523 let mut on_disk = fresh(session_id);
3524 on_disk.title = "Old".to_string();
3525 on_disk.title_version = 1;
3526 on_disk.metadata_version = 3;
3527 storage.save_session(&on_disk).await.unwrap();
3528
3529 let mut authoritative_copy = fresh(session_id);
3530 authoritative_copy.title = "New Authoritative".to_string();
3531 authoritative_copy.title_version = 2;
3532 authoritative_copy.metadata_version = 4;
3533 authoritative_copy.pinned = true;
3534
3535 merge_save_session(&storage, &mut authoritative_copy)
3536 .await
3537 .unwrap();
3538
3539 let after = storage.load_session(session_id).await.unwrap().unwrap();
3540 assert_eq!(after.title, "New Authoritative");
3541 assert_eq!(after.title_version, 2);
3542 assert_eq!(after.metadata_version, 4);
3543 assert!(after.pinned);
3544 }
3545
3546 #[tokio::test]
3547 async fn merge_keeps_runtime_messages_when_disk_only_changed_metadata() {
3548 let (_temp, storage) = make_storage().await;
3549 let session_id = "merge-messages";
3550
3551 let mut on_disk = fresh(session_id);
3552 on_disk.title = "Fresh Title".to_string();
3553 on_disk.title_version = 2;
3554 on_disk.metadata_version = 5;
3555 storage.save_session(&on_disk).await.unwrap();
3556
3557 let mut runtime_copy = fresh(session_id);
3558 runtime_copy.title = "Stale".to_string();
3559 runtime_copy.metadata_version = 0;
3560 runtime_copy.messages = vec![bamboo_domain::session::types::Message {
3561 role: bamboo_domain::session::types::Role::User,
3562 content: "keep me".to_string(),
3563 id: "msg-1".to_string(),
3564 created_at: chrono::Utc::now(),
3565 reasoning: None,
3566 reasoning_signature: None,
3567 content_parts: None,
3568 image_ocr: None,
3569 phase: None,
3570 tool_calls: None,
3571 tool_call_id: None,
3572 tool_success: None,
3573 compressed: false,
3574 compressed_by_event_id: None,
3575 never_compress: false,
3576 compression_level: 0,
3577 metadata: None,
3578 }];
3579
3580 merge_save_session(&storage, &mut runtime_copy)
3581 .await
3582 .unwrap();
3583
3584 let after = storage.load_session(session_id).await.unwrap().unwrap();
3585 assert_eq!(after.title, "Fresh Title");
3586 assert_eq!(after.metadata_version, 5);
3587 assert_eq!(after.messages.len(), 1);
3588 assert_eq!(after.messages[0].content, "keep me");
3589 }
3590
3591 #[tokio::test]
3592 async fn runtime_control_plane_port_uses_sidecar_without_rewriting_messages() {
3593 use bamboo_domain::session::types::Message;
3594
3595 let (_temp, storage) = make_storage().await;
3596 let store = LockedSessionStore::new(storage.clone());
3597 let session_id = "runtime-control-plane";
3598
3599 let mut durable = fresh(session_id);
3600 durable.add_message(Message::user("durable transcript"));
3601 storage.save_session(&durable).await.unwrap();
3602
3603 let mut runtime = durable.clone();
3604 runtime.model = "updated-control-plane-model".to_string();
3605 runtime.add_message(Message::assistant("uncheckpointed runtime message", None));
3606 RuntimeSessionPersistence::save_runtime_control_plane(&store, &mut runtime)
3607 .await
3608 .unwrap();
3609
3610 let control_plane =
3611 RuntimeSessionPersistence::load_runtime_control_plane(&store, session_id)
3612 .await
3613 .unwrap()
3614 .expect("control-plane exists");
3615 assert!(
3616 control_plane.messages.is_empty(),
3617 "LockedSessionStore must expose its message-free sidecar"
3618 );
3619 assert_eq!(control_plane.model, "updated-control-plane-model");
3620
3621 let reloaded = storage
3622 .load_session(session_id)
3623 .await
3624 .unwrap()
3625 .expect("session exists");
3626 assert_eq!(reloaded.model, "updated-control-plane-model");
3627 assert_eq!(
3628 reloaded.messages.len(),
3629 1,
3630 "control-plane save must not write the uncheckpointed message"
3631 );
3632 assert_eq!(reloaded.messages[0].content, "durable transcript");
3633 }
3634
3635 #[tokio::test]
3636 async fn atomic_task_patch_loads_inside_lock_and_preserves_interleaved_runtime_state() {
3637 let temp = tempfile::tempdir().unwrap();
3638 let inner = Arc::new(
3639 SessionStoreV2::new(temp.path().to_path_buf())
3640 .await
3641 .expect("storage init"),
3642 );
3643 let session_id = "atomic-task-patch";
3644 inner
3645 .save_session(&fresh(session_id))
3646 .await
3647 .expect("seed session");
3648
3649 let counted = Arc::new(CountingControlPlaneStorage {
3650 inner: inner.clone(),
3651 control_plane_loads: AtomicUsize::new(0),
3652 full_saves: AtomicUsize::new(0),
3653 runtime_state_saves: AtomicUsize::new(0),
3654 });
3655 let storage: Arc<dyn Storage> = counted.clone();
3656 let store = Arc::new(LockedSessionStore::new(storage));
3657 let guard = store.acquire_lock(session_id).await;
3658 let now = chrono::Utc::now();
3659 let task_list = bamboo_domain::TaskList {
3660 session_id: session_id.to_string(),
3661 title: "Atomic Task patch".to_string(),
3662 items: Vec::new(),
3663 created_at: now,
3664 updated_at: now,
3665 };
3666 let (started_tx, started_rx) = tokio::sync::oneshot::channel();
3667 let patch_store = store.clone();
3668 let patch = tokio::spawn(async move {
3669 let _ = started_tx.send(());
3670 RuntimeSessionPersistence::update_task_list_control_plane(
3671 patch_store.as_ref(),
3672 session_id,
3673 &task_list,
3674 "9",
3675 )
3676 .await
3677 });
3678 started_rx.await.expect("patch task started");
3679 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
3680 assert_eq!(
3681 counted.control_plane_loads.load(Ordering::SeqCst),
3682 0,
3683 "Task patch must acquire the session lock before loading its snapshot"
3684 );
3685
3686 let mut latest = inner
3690 .load_runtime_control_plane(session_id)
3691 .await
3692 .expect("load latest control-plane")
3693 .expect("control-plane exists");
3694 latest.agent_runtime_state = Some(bamboo_domain::AgentRuntimeState::new("latest-run"));
3695 latest
3696 .metadata
3697 .insert("concurrent.runtime".to_string(), "preserve".to_string());
3698 inner
3699 .save_runtime_state(&latest)
3700 .await
3701 .expect("publish concurrent runtime transition");
3702 drop(guard);
3703
3704 assert!(
3705 patch.await.expect("patch join").expect("patch succeeds"),
3706 "existing root must be patched"
3707 );
3708 assert_eq!(counted.control_plane_loads.load(Ordering::SeqCst), 1);
3709 let reloaded = inner
3710 .load_session(session_id)
3711 .await
3712 .expect("reload")
3713 .expect("session exists");
3714 assert_eq!(
3715 reloaded
3716 .agent_runtime_state
3717 .as_ref()
3718 .map(|state| state.run_id.as_str()),
3719 Some("latest-run")
3720 );
3721 assert_eq!(
3722 reloaded
3723 .metadata
3724 .get("concurrent.runtime")
3725 .map(String::as_str),
3726 Some("preserve")
3727 );
3728 assert_eq!(reloaded.task_list_version_meta().as_deref(), Some("9"));
3729 assert_eq!(
3730 reloaded.task_list.as_ref().map(|list| list.title.as_str()),
3731 Some("Atomic Task patch")
3732 );
3733 }
3734
3735 #[tokio::test]
3736 async fn paired_task_cas_conflict_cannot_overwrite_newer_root_or_child_state() {
3737 let (_temp, storage) = make_storage().await;
3738 let store = LockedSessionStore::new(storage.clone());
3739 let root_id = "task-cas-root";
3740 let child_id = "task-cas-child";
3741 let now = chrono::Utc::now();
3742 let task_list = |title: &str| bamboo_domain::TaskList {
3743 session_id: root_id.to_string(),
3744 title: title.to_string(),
3745 items: Vec::new(),
3746 created_at: now,
3747 updated_at: now,
3748 };
3749
3750 let mut root = fresh(root_id);
3751 root.set_task_list(task_list("newer root"));
3752 root.set_task_list_version_meta("2");
3753 storage.save_session(&root).await.expect("seed root");
3754 let mut child = Session::new_child(child_id, root_id, "model", "child");
3755 child.set_task_list(task_list("current child"));
3756 child.set_task_list_version_meta("1");
3757 storage.save_session(&child).await.expect("seed child");
3758
3759 let updated = RuntimeSessionPersistence::update_task_list_control_planes_if_version(
3760 &store,
3761 child_id,
3762 root_id,
3763 "1",
3764 &task_list("current child"),
3765 &task_list("stale evaluator"),
3766 "3",
3767 )
3768 .await
3769 .expect("CAS returns clean conflict");
3770 assert!(
3771 !updated,
3772 "mismatched root generation must reject both writes"
3773 );
3774
3775 let durable_root = storage
3776 .load_session(root_id)
3777 .await
3778 .expect("load root")
3779 .expect("root exists");
3780 let durable_child = storage
3781 .load_session(child_id)
3782 .await
3783 .expect("load child")
3784 .expect("child exists");
3785 assert_eq!(durable_root.task_list_version_meta().as_deref(), Some("2"));
3786 assert_eq!(
3787 durable_root
3788 .task_list
3789 .as_ref()
3790 .map(|list| list.title.as_str()),
3791 Some("newer root")
3792 );
3793 assert_eq!(durable_child.task_list_version_meta().as_deref(), Some("1"));
3794 assert_eq!(
3795 durable_child
3796 .task_list
3797 .as_ref()
3798 .map(|list| list.title.as_str()),
3799 Some("current child")
3800 );
3801 }
3802
3803 #[tokio::test]
3804 async fn paired_task_cas_success_uses_only_targeted_saves_and_preserves_both_transcripts() {
3805 use bamboo_domain::session::types::Message;
3806
3807 let temp = tempfile::tempdir().unwrap();
3808 let inner = Arc::new(
3809 SessionStoreV2::new(temp.path().to_path_buf())
3810 .await
3811 .expect("storage init"),
3812 );
3813 let root_id = "task-cas-success-root";
3814 let child_id = "task-cas-success-child";
3815 let now = chrono::Utc::now();
3816 let task_list = |title: &str| bamboo_domain::TaskList {
3817 session_id: root_id.to_string(),
3818 title: title.to_string(),
3819 items: Vec::new(),
3820 created_at: now,
3821 updated_at: now,
3822 };
3823
3824 let mut root = fresh(root_id);
3825 root.add_message(Message::user("root transcript"));
3826 root.metadata
3827 .insert("unrelated.root".to_string(), "preserve".to_string());
3828 root.agent_runtime_state = Some(bamboo_domain::AgentRuntimeState::new("root-run"));
3829 root.set_task_list(task_list("old shared"));
3830 root.set_task_list_version_meta("1");
3831 inner.save_session(&root).await.expect("seed root");
3832
3833 let mut child = Session::new_child(child_id, root_id, "model", "child");
3834 child.add_message(Message::user("child transcript"));
3835 child
3836 .metadata
3837 .insert("unrelated.child".to_string(), "preserve".to_string());
3838 child.agent_runtime_state = Some(bamboo_domain::AgentRuntimeState::new("child-run"));
3839 child.set_task_list(task_list("old shared"));
3840 child.set_task_list_version_meta("1");
3841 inner.save_session(&child).await.expect("seed child");
3842
3843 let counted = Arc::new(CountingControlPlaneStorage {
3844 inner: inner.clone(),
3845 control_plane_loads: AtomicUsize::new(0),
3846 full_saves: AtomicUsize::new(0),
3847 runtime_state_saves: AtomicUsize::new(0),
3848 });
3849 let storage: Arc<dyn Storage> = counted.clone();
3850 let store = LockedSessionStore::new(storage);
3851 assert!(
3852 RuntimeSessionPersistence::update_task_list_control_planes_if_version(
3853 &store,
3854 child_id,
3855 root_id,
3856 "1",
3857 &task_list("old shared"),
3858 &task_list("evaluated"),
3859 "2",
3860 )
3861 .await
3862 .expect("paired CAS succeeds")
3863 );
3864 assert_eq!(counted.control_plane_loads.load(Ordering::SeqCst), 2);
3865 assert_eq!(counted.runtime_state_saves.load(Ordering::SeqCst), 2);
3866 assert_eq!(
3867 counted.full_saves.load(Ordering::SeqCst),
3868 0,
3869 "evaluation CAS must not call full save_session for child or root"
3870 );
3871
3872 let durable_root = inner.load_session(root_id).await.unwrap().unwrap();
3873 let durable_child = inner.load_session(child_id).await.unwrap().unwrap();
3874 for (session, transcript, metadata_key, run_id) in [
3875 (
3876 &durable_root,
3877 "root transcript",
3878 "unrelated.root",
3879 "root-run",
3880 ),
3881 (
3882 &durable_child,
3883 "child transcript",
3884 "unrelated.child",
3885 "child-run",
3886 ),
3887 ] {
3888 assert_eq!(session.task_list_version_meta().as_deref(), Some("2"));
3889 assert_eq!(
3890 session.task_list.as_ref().map(|list| list.title.as_str()),
3891 Some("evaluated")
3892 );
3893 assert_eq!(session.messages.len(), 1);
3894 assert_eq!(session.messages[0].content, transcript);
3895 assert_eq!(
3896 session.metadata.get(metadata_key).map(String::as_str),
3897 Some("preserve")
3898 );
3899 assert_eq!(
3900 session
3901 .agent_runtime_state
3902 .as_ref()
3903 .map(|state| state.run_id.as_str()),
3904 Some(run_id)
3905 );
3906 }
3907 }
3908
3909 #[tokio::test]
3910 async fn locked_runtime_and_full_saves_adopt_task_conflicts_before_publish() {
3911 let temp = tempfile::tempdir().unwrap();
3912 let home = temp.path().to_path_buf();
3913 let first_storage = Arc::new(
3914 SessionStoreV2::new(home.clone())
3915 .await
3916 .expect("first storage init"),
3917 );
3918 let second_storage = Arc::new(
3919 SessionStoreV2::new(home)
3920 .await
3921 .expect("second storage init"),
3922 );
3923 let now = chrono::Utc::now();
3924 let task_list = |session_id: &str, title: &str| bamboo_domain::TaskList {
3925 session_id: session_id.to_string(),
3926 title: title.to_string(),
3927 items: Vec::new(),
3928 created_at: now,
3929 updated_at: now,
3930 };
3931
3932 let runtime_id = "ordinary-task-retry-runtime";
3933 let full_id = "ordinary-task-retry-full";
3934 let mut runtime_initial = fresh(runtime_id);
3935 runtime_initial.set_task_list(task_list(runtime_id, "runtime v1"));
3936 runtime_initial.set_task_list_version_meta("1");
3937 first_storage
3938 .save_session(&runtime_initial)
3939 .await
3940 .expect("seed runtime session");
3941 let mut full_initial = fresh(full_id);
3942 full_initial.set_task_list(task_list(full_id, "full v1"));
3943 full_initial.set_task_list_version_meta("1");
3944 first_storage
3945 .save_session(&full_initial)
3946 .await
3947 .expect("seed full session");
3948
3949 let mut runtime_advanced = runtime_initial.clone();
3950 runtime_advanced.set_task_list(task_list(runtime_id, "runtime v2"));
3951 runtime_advanced.set_task_list_version_meta("2");
3952 second_storage
3953 .save_runtime_state(&runtime_advanced)
3954 .await
3955 .expect("advance runtime Task generation");
3956 let mut full_advanced = full_initial.clone();
3957 full_advanced.set_task_list(task_list(full_id, "full v2"));
3958 full_advanced.set_task_list_version_meta("2");
3959 second_storage
3960 .save_runtime_state(&full_advanced)
3961 .await
3962 .expect("advance full Task generation");
3963
3964 let storage: Arc<dyn Storage> = first_storage.clone();
3965 let store = LockedSessionStore::new(storage);
3966 let runtime_published = Arc::new(std::sync::Mutex::new(None));
3967 let runtime_callback = runtime_published.clone();
3968 let mut runtime_stale = runtime_initial;
3969 runtime_stale
3970 .metadata
3971 .insert("runtime.non-task".to_string(), "preserved".to_string());
3972 store
3973 .save_runtime_only_and_publish(&mut runtime_stale, move |saved| {
3974 *runtime_callback.lock().expect("runtime publish lock") = Some(saved.clone());
3975 })
3976 .await
3977 .expect("locked runtime save rebases and retries");
3978
3979 let full_published = Arc::new(std::sync::Mutex::new(None));
3980 let full_callback = full_published.clone();
3981 let mut full_stale = full_initial;
3982 full_stale
3983 .metadata
3984 .insert("full.non-task".to_string(), "preserved".to_string());
3985 store
3986 .merge_save_runtime_and_publish(&mut full_stale, move |saved, committed| {
3987 assert!(committed);
3988 *full_callback.lock().expect("full publish lock") = Some(saved.clone());
3989 })
3990 .await
3991 .expect("locked full save rebases and retries");
3992
3993 let durable_runtime = first_storage
3994 .load_session(runtime_id)
3995 .await
3996 .unwrap()
3997 .expect("durable runtime session");
3998 let durable_full = first_storage
3999 .load_session(full_id)
4000 .await
4001 .unwrap()
4002 .expect("durable full session");
4003 let published_runtime = runtime_published
4004 .lock()
4005 .expect("runtime publish lock")
4006 .clone()
4007 .expect("runtime published snapshot");
4008 let published_full = full_published
4009 .lock()
4010 .expect("full publish lock")
4011 .clone()
4012 .expect("full published snapshot");
4013
4014 for (session, expected_title, metadata_key) in [
4015 (&runtime_stale, "runtime v2", "runtime.non-task"),
4016 (&durable_runtime, "runtime v2", "runtime.non-task"),
4017 (&published_runtime, "runtime v2", "runtime.non-task"),
4018 (&full_stale, "full v2", "full.non-task"),
4019 (&durable_full, "full v2", "full.non-task"),
4020 (&published_full, "full v2", "full.non-task"),
4021 ] {
4022 assert_eq!(session.task_list_version_meta().as_deref(), Some("2"));
4023 assert_eq!(
4024 session.task_list.as_ref().map(|list| list.title.as_str()),
4025 Some(expected_title)
4026 );
4027 assert_eq!(
4028 session.metadata.get(metadata_key).map(String::as_str),
4029 Some("preserved")
4030 );
4031 }
4032 }
4033
4034 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
4035 async fn single_task_cas_rejects_same_version_divergent_snapshot_without_callback() {
4036 let (_temp, storage) = make_storage().await;
4037 let store = LockedSessionStore::new(storage.clone());
4038 let session_id = "single-task-exact-snapshot";
4039 let now = chrono::Utc::now();
4040 let task_list = |title: &str| bamboo_domain::TaskList {
4041 session_id: session_id.to_string(),
4042 title: title.to_string(),
4043 items: Vec::new(),
4044 created_at: now,
4045 updated_at: now,
4046 };
4047 let durable_winner = task_list("durable winner");
4048 let stale_snapshot = task_list("stale same-version snapshot");
4049 let mut session = fresh(session_id);
4050 session.set_task_list(durable_winner.clone());
4051 session.set_task_list_version_meta("1");
4052 storage.save_session(&session).await.expect("seed session");
4053
4054 let published = Arc::new(AtomicBool::new(false));
4055 let callback = published.clone();
4056 assert!(!store
4057 .update_task_list_control_plane_if_version_and_publish(
4058 session_id,
4059 "1",
4060 &stale_snapshot,
4061 &task_list("stale evaluation"),
4062 "2",
4063 move |_| callback.store(true, Ordering::SeqCst),
4064 )
4065 .await
4066 .expect("same-version divergence is a clean stale result"));
4067 assert!(!published.load(Ordering::SeqCst));
4068 let durable = storage
4069 .load_session(session_id)
4070 .await
4071 .unwrap()
4072 .expect("session remains");
4073 assert_eq!(durable.task_list_version_meta().as_deref(), Some("1"));
4074 assert_eq!(
4075 serde_json::to_value(&durable.task_list).expect("serialize durable Task list"),
4076 serde_json::to_value(Some(&durable_winner)).expect("serialize expected Task list")
4077 );
4078 }
4079
4080 #[tokio::test]
4081 async fn paired_task_cas_rejects_same_version_divergent_snapshot_without_callback() {
4082 let (_temp, storage) = make_storage().await;
4083 let store = LockedSessionStore::new(storage.clone());
4084 let root_id = "paired-task-exact-snapshot-root";
4085 let child_id = "paired-task-exact-snapshot-child";
4086 let now = chrono::Utc::now();
4087 let task_list = |title: &str| bamboo_domain::TaskList {
4088 session_id: root_id.to_string(),
4089 title: title.to_string(),
4090 items: Vec::new(),
4091 created_at: now,
4092 updated_at: now,
4093 };
4094 let durable_winner = task_list("durable winner");
4095 let stale_snapshot = task_list("stale same-version snapshot");
4096 let mut root = fresh(root_id);
4097 root.set_task_list(durable_winner.clone());
4098 root.set_task_list_version_meta("1");
4099 storage.save_session(&root).await.expect("seed root");
4100 let mut child = Session::new_child(child_id, root_id, "model", "child");
4101 child.set_task_list(durable_winner.clone());
4102 child.set_task_list_version_meta("1");
4103 storage.save_session(&child).await.expect("seed child");
4104
4105 let published = Arc::new(AtomicBool::new(false));
4106 let callback = published.clone();
4107 assert!(!store
4108 .update_task_list_control_planes_if_version_and_publish(
4109 child_id,
4110 root_id,
4111 "1",
4112 &stale_snapshot,
4113 &task_list("stale evaluation"),
4114 "2",
4115 move |_, _| callback.store(true, Ordering::SeqCst),
4116 )
4117 .await
4118 .expect("same-version divergence is a clean stale result"));
4119 assert!(!published.load(Ordering::SeqCst));
4120 for id in [child_id, root_id] {
4121 let durable = storage
4122 .load_session(id)
4123 .await
4124 .unwrap()
4125 .expect("session remains");
4126 assert_eq!(durable.task_list_version_meta().as_deref(), Some("1"));
4127 assert_eq!(
4128 serde_json::to_value(&durable.task_list).expect("serialize durable Task list"),
4129 serde_json::to_value(Some(&durable_winner)).expect("serialize expected Task list")
4130 );
4131 }
4132 }
4133
4134 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
4135 async fn unconditional_root_task_patch_reports_final_cas_conflict_without_publishing() {
4136 let temp = tempfile::tempdir().unwrap();
4137 let home = temp.path().to_path_buf();
4138 let first_inner = Arc::new(
4139 SessionStoreV2::new(home.clone())
4140 .await
4141 .expect("first storage init"),
4142 );
4143 let root_id = "single-task-unconditional-loser";
4144 let now = chrono::Utc::now();
4145 let task_list = |title: &str| bamboo_domain::TaskList {
4146 session_id: root_id.to_string(),
4147 title: title.to_string(),
4148 items: Vec::new(),
4149 created_at: now,
4150 updated_at: now,
4151 };
4152 let mut root = fresh(root_id);
4153 root.task_list = Some(task_list("original"));
4154 root.set_task_list_version_meta("1");
4155 first_inner.save_session(&root).await.expect("seed root");
4156 let second_inner = Arc::new(
4157 SessionStoreV2::new(home)
4158 .await
4159 .expect("second storage init"),
4160 );
4161 let commit_reached = Arc::new(tokio::sync::Barrier::new(2));
4162 let release_commit = Arc::new(tokio::sync::Barrier::new(2));
4163 let storage: Arc<dyn Storage> = Arc::new(SingleCommitPauseStorage {
4164 inner: first_inner.clone(),
4165 commit_reached: commit_reached.clone(),
4166 release_commit: release_commit.clone(),
4167 });
4168 let store = LockedSessionStore::new(storage);
4169 let published = Arc::new(AtomicBool::new(false));
4170 let callback = published.clone();
4171 let loser_candidate = task_list("loser");
4172
4173 let loser = store.update_task_list_control_plane_and_publish(
4174 root_id,
4175 &loser_candidate,
4176 "2",
4177 move |_| callback.store(true, Ordering::SeqCst),
4178 );
4179 let winner = async {
4180 commit_reached.wait().await;
4181 let original = second_inner
4182 .load_runtime_control_plane(root_id)
4183 .await
4184 .expect("load winner original")
4185 .expect("winner original exists");
4186 let mut updated = original.clone();
4187 updated.task_list = Some(task_list("winner"));
4188 updated.set_task_list_version_meta("2");
4189 assert!(second_inner
4190 .save_task_control_plane_if_matches(&original, &updated)
4191 .await
4192 .expect("commit winner"));
4193 release_commit.wait().await;
4194 };
4195 let (loser_result, ()) = tokio::join!(loser, winner);
4196 let error = loser_result.expect_err("unconditional loser must be an explicit conflict");
4197 assert_eq!(error.kind(), std::io::ErrorKind::WouldBlock);
4198 assert!(!published.load(Ordering::SeqCst));
4199 let durable = first_inner
4200 .load_session(root_id)
4201 .await
4202 .unwrap()
4203 .expect("durable root");
4204 assert_eq!(durable.task_list_version_meta().as_deref(), Some("2"));
4205 assert_eq!(
4206 durable.task_list.as_ref().map(|list| list.title.as_str()),
4207 Some("winner")
4208 );
4209 }
4210
4211 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
4212 async fn independent_root_task_patches_have_one_final_cas_winner() {
4213 let temp = tempfile::tempdir().unwrap();
4214 let home = temp.path().to_path_buf();
4215 let first_inner = Arc::new(
4216 SessionStoreV2::new(home.clone())
4217 .await
4218 .expect("first storage init"),
4219 );
4220 let root_id = "single-task-cas-race-root";
4221 let now = chrono::Utc::now();
4222 let task_list = |title: &str| bamboo_domain::TaskList {
4223 session_id: root_id.to_string(),
4224 title: title.to_string(),
4225 items: Vec::new(),
4226 created_at: now,
4227 updated_at: now,
4228 };
4229 let mut root = fresh(root_id);
4230 root.set_task_list(task_list("original"));
4231 root.set_task_list_version_meta("1");
4232 first_inner.save_session(&root).await.expect("seed root");
4233 let second_inner = Arc::new(
4234 SessionStoreV2::new(home)
4235 .await
4236 .expect("second storage init"),
4237 );
4238 let before_commit = Arc::new(tokio::sync::Barrier::new(2));
4239 let first_storage: Arc<dyn Storage> = Arc::new(SingleCommitBarrierStorage {
4240 inner: first_inner.clone(),
4241 before_commit: before_commit.clone(),
4242 });
4243 let second_storage: Arc<dyn Storage> = Arc::new(SingleCommitBarrierStorage {
4244 inner: second_inner.clone(),
4245 before_commit,
4246 });
4247 let first_store = LockedSessionStore::new(first_storage);
4248 let second_store = LockedSessionStore::new(second_storage);
4249 let first_published = Arc::new(AtomicBool::new(false));
4250 let second_published = Arc::new(AtomicBool::new(false));
4251 let first_callback = first_published.clone();
4252 let second_callback = second_published.clone();
4253 let expected = task_list("original");
4254 let first_candidate = task_list("candidate one");
4255 let second_candidate = task_list("candidate two");
4256
4257 let first = first_store.update_task_list_control_plane_if_version_and_publish(
4260 root_id,
4261 "1",
4262 &expected,
4263 &first_candidate,
4264 "2",
4265 move |_| first_callback.store(true, Ordering::SeqCst),
4266 );
4267 let second = second_store.update_task_list_control_plane_if_version_and_publish(
4268 root_id,
4269 "1",
4270 &expected,
4271 &second_candidate,
4272 "2",
4273 move |_| second_callback.store(true, Ordering::SeqCst),
4274 );
4275 let (first_result, second_result) = tokio::join!(first, second);
4276 let first_won = first_result.expect("first root result");
4277 let second_won = second_result.expect("second root result");
4278 assert_ne!(first_won, second_won, "exactly one root candidate wins");
4279 assert_eq!(first_published.load(Ordering::SeqCst), first_won);
4280 assert_eq!(second_published.load(Ordering::SeqCst), second_won);
4281
4282 let durable = first_inner
4283 .load_session(root_id)
4284 .await
4285 .unwrap()
4286 .expect("root");
4287 assert_eq!(durable.task_list_version_meta().as_deref(), Some("2"));
4288 assert_eq!(
4289 durable.task_list.as_ref().map(|list| list.title.as_str()),
4290 Some(if first_won {
4291 "candidate one"
4292 } else {
4293 "candidate two"
4294 })
4295 );
4296 }
4297
4298 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
4299 async fn independent_locked_stores_revalidate_pair_cas_at_storage_commit_point() {
4300 let temp = tempfile::tempdir().unwrap();
4301 let home = temp.path().to_path_buf();
4302 let first_inner = Arc::new(
4303 SessionStoreV2::new(home.clone())
4304 .await
4305 .expect("first storage init"),
4306 );
4307 let root_id = "task-cas-race-root";
4308 let child_id = "task-cas-race-child";
4309 let now = chrono::Utc::now();
4310 let task_list = |title: &str| bamboo_domain::TaskList {
4311 session_id: root_id.to_string(),
4312 title: title.to_string(),
4313 items: Vec::new(),
4314 created_at: now,
4315 updated_at: now,
4316 };
4317
4318 let mut root = fresh(root_id);
4319 root.set_task_list(task_list("original shared"));
4320 root.set_task_list_version_meta("1");
4321 first_inner.save_session(&root).await.expect("seed root");
4322 let mut child = Session::new_child(child_id, root_id, "model", "child");
4323 child.set_task_list(task_list("original shared"));
4324 child.set_task_list_version_meta("1");
4325 first_inner.save_session(&child).await.expect("seed child");
4326
4327 let second_inner = Arc::new(
4328 SessionStoreV2::new(home)
4329 .await
4330 .expect("second storage init"),
4331 );
4332 let before_commit = Arc::new(tokio::sync::Barrier::new(2));
4333 let first_storage: Arc<dyn Storage> = Arc::new(PairCommitBarrierStorage {
4334 inner: first_inner.clone(),
4335 before_commit: before_commit.clone(),
4336 });
4337 let second_storage: Arc<dyn Storage> = Arc::new(PairCommitBarrierStorage {
4338 inner: second_inner.clone(),
4339 before_commit,
4340 });
4341 let first_store = LockedSessionStore::new(first_storage);
4342 let second_store = LockedSessionStore::new(second_storage);
4343 let first_published = Arc::new(AtomicBool::new(false));
4344 let second_published = Arc::new(AtomicBool::new(false));
4345 let first_published_callback = first_published.clone();
4346 let second_published_callback = second_published.clone();
4347 let expected = task_list("original shared");
4348 let first_candidate = task_list("candidate one");
4349 let second_candidate = task_list("candidate two");
4350
4351 let first = first_store.update_task_list_control_planes_if_version_and_publish(
4352 child_id,
4353 root_id,
4354 "1",
4355 &expected,
4356 &first_candidate,
4357 "2",
4358 move |_, _| first_published_callback.store(true, Ordering::SeqCst),
4359 );
4360 let second = second_store.update_task_list_control_planes_if_version_and_publish(
4361 child_id,
4362 root_id,
4363 "1",
4364 &expected,
4365 &second_candidate,
4366 "2",
4367 move |_, _| second_published_callback.store(true, Ordering::SeqCst),
4368 );
4369 let (first_result, second_result) = tokio::join!(first, second);
4370 let first_won = first_result.expect("first CAS result");
4371 let second_won = second_result.expect("second CAS result");
4372 assert_ne!(first_won, second_won, "exactly one staged v1 CAS may win");
4373 assert_eq!(first_published.load(Ordering::SeqCst), first_won);
4374 assert_eq!(second_published.load(Ordering::SeqCst), second_won);
4375
4376 let expected_title = if first_won {
4377 "candidate one"
4378 } else {
4379 "candidate two"
4380 };
4381 let durable_child = first_inner
4382 .load_session(child_id)
4383 .await
4384 .unwrap()
4385 .expect("child");
4386 let durable_root = second_inner
4387 .load_session(root_id)
4388 .await
4389 .unwrap()
4390 .expect("root");
4391 for session in [&durable_child, &durable_root] {
4392 assert_eq!(session.task_list_version_meta().as_deref(), Some("2"));
4393 assert_eq!(
4394 session.task_list.as_ref().map(|list| list.title.as_str()),
4395 Some(expected_title)
4396 );
4397 }
4398 }
4399
4400 #[tokio::test]
4401 async fn paired_task_second_write_failure_rolls_back_and_skips_publish_callback() {
4402 use bamboo_domain::session::types::Message;
4403
4404 let temp = tempfile::tempdir().unwrap();
4405 let inner = Arc::new(
4406 SessionStoreV2::new(temp.path().to_path_buf())
4407 .await
4408 .expect("storage init"),
4409 );
4410 let root_id = "task-cas-failure-root";
4411 let child_id = "task-cas-failure-child";
4412 let now = chrono::Utc::now();
4413 let task_list = |title: &str| bamboo_domain::TaskList {
4414 session_id: root_id.to_string(),
4415 title: title.to_string(),
4416 items: Vec::new(),
4417 created_at: now,
4418 updated_at: now,
4419 };
4420
4421 let mut root = fresh(root_id);
4422 root.add_message(Message::user("root transcript"));
4423 root.metadata
4424 .insert("unrelated.root".to_string(), "preserve".to_string());
4425 root.set_task_list(task_list("old shared"));
4426 root.set_task_list_version_meta("1");
4427 inner.save_session(&root).await.expect("seed root");
4428
4429 let mut child = Session::new_child(child_id, root_id, "model", "child");
4430 child.add_message(Message::user("child transcript"));
4431 child
4432 .metadata
4433 .insert("unrelated.child".to_string(), "preserve".to_string());
4434 child.set_task_list(task_list("old shared"));
4435 child.set_task_list_version_meta("1");
4436 inner.save_session(&child).await.expect("seed child");
4437
4438 inner
4439 .inject_runtime_task_transaction_fault(RuntimeTaskTransactionFault::SecondUpdatedWrite);
4440 let storage: Arc<dyn Storage> = inner.clone();
4441 let store = LockedSessionStore::new(storage);
4442 let published = Arc::new(AtomicBool::new(false));
4443 let published_for_callback = published.clone();
4444 let error = store
4445 .update_task_list_control_planes_if_version_and_publish(
4446 child_id,
4447 root_id,
4448 "1",
4449 &task_list("old shared"),
4450 &task_list("must roll back"),
4451 "2",
4452 move |_, _| published_for_callback.store(true, Ordering::SeqCst),
4453 )
4454 .await
4455 .expect_err("injected second write fails");
4456 assert!(error.to_string().contains("rolled back"), "{error}");
4457 assert!(
4458 !published.load(Ordering::SeqCst),
4459 "durable failure must not publish either cache snapshot"
4460 );
4461
4462 let durable_root = inner.load_session(root_id).await.unwrap().unwrap();
4463 let durable_child = inner.load_session(child_id).await.unwrap().unwrap();
4464 for (session, title, transcript, metadata_key) in [
4465 (
4466 &durable_root,
4467 "old shared",
4468 "root transcript",
4469 "unrelated.root",
4470 ),
4471 (
4472 &durable_child,
4473 "old shared",
4474 "child transcript",
4475 "unrelated.child",
4476 ),
4477 ] {
4478 assert_eq!(session.task_list_version_meta().as_deref(), Some("1"));
4479 assert_eq!(
4480 session.task_list.as_ref().map(|list| list.title.as_str()),
4481 Some(title)
4482 );
4483 assert_eq!(session.messages[0].content, transcript);
4484 assert_eq!(
4485 session.metadata.get(metadata_key).map(String::as_str),
4486 Some("preserve")
4487 );
4488 }
4489 }
4490
4491 #[tokio::test]
4494 async fn locked_merge_save_runtime_serialises_concurrent_writes() {
4495 let (_temp, storage) = make_storage().await;
4496 let store = Arc::new(LockedSessionStore::new(storage));
4497 let session_id = "lock-serial".to_string();
4498
4499 let base = fresh(&session_id);
4501 store.storage().save_session(&base).await.unwrap();
4502
4503 let store_a = store.clone();
4506 let store_b = store.clone();
4507 let sid_a = session_id.clone();
4508 let sid_b = session_id.clone();
4509
4510 let a = tokio::spawn(async move {
4511 let _guard = store_a.acquire_lock(&sid_a).await;
4512 let mut s = store_a
4513 .storage()
4514 .load_session(&sid_a)
4515 .await
4516 .unwrap()
4517 .unwrap();
4518 s.title = "Writer A".to_string();
4519 s.title_version = s.title_version.saturating_add(1);
4520 s.metadata_version = s.metadata_version.saturating_add(1);
4521 s.updated_at = chrono::Utc::now();
4522 store_a.storage().save_session(&s).await.unwrap();
4523 s.title_version
4524 });
4525
4526 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
4528
4529 let b = tokio::spawn(async move {
4530 let _guard = store_b.acquire_lock(&sid_b).await;
4531 let mut s = store_b
4532 .storage()
4533 .load_session(&sid_b)
4534 .await
4535 .unwrap()
4536 .unwrap();
4537 s.title = "Writer B".to_string();
4538 s.title_version = s.title_version.saturating_add(1);
4539 s.metadata_version = s.metadata_version.saturating_add(1);
4540 s.updated_at = chrono::Utc::now();
4541 store_b.storage().save_session(&s).await.unwrap();
4542 s.title_version
4543 });
4544
4545 let (ver_a, ver_b) = tokio::join!(a, b);
4546 let final_s = store
4547 .storage()
4548 .load_session(&session_id)
4549 .await
4550 .unwrap()
4551 .unwrap();
4552 assert!(
4553 ver_a.unwrap() != ver_b.unwrap(),
4554 "concurrent writers must produce distinct versions"
4555 );
4556 assert_eq!(final_s.metadata_version, 2);
4557 }
4558
4559 #[tokio::test]
4560 async fn commit_metadata_is_plain_save_inside_lock() {
4561 let (_temp, storage) = make_storage().await;
4562 let store = LockedSessionStore::new(storage);
4563 let session_id = "commit-plain";
4564
4565 let mut s = fresh(session_id);
4566 s.title = "Committed".to_string();
4567 s.metadata_version = 1;
4568 s.title_version = 2;
4569
4570 store.commit_metadata(&s).await.unwrap();
4571
4572 let after = store
4573 .storage()
4574 .load_session(session_id)
4575 .await
4576 .unwrap()
4577 .unwrap();
4578 assert_eq!(after.title, "Committed");
4579 assert_eq!(after.metadata_version, 1);
4580 assert_eq!(after.title_version, 2);
4581 }
4582
4583 #[tokio::test]
4586 async fn acquire_lock_self_evicts_when_no_other_holder() {
4587 let (_temp, storage) = make_storage().await;
4588 let store = LockedSessionStore::new(storage);
4589
4590 {
4591 let _guard = store.acquire_lock("solo").await;
4592 assert_eq!(store.locks.len(), 1, "entry present while the lock is held");
4593 }
4594 assert_eq!(
4597 store.locks.len(),
4598 0,
4599 "lock entry must be evicted once released with no other holder"
4600 );
4601 }
4602
4603 #[tokio::test]
4604 async fn acquire_lock_many_distinct_ids_do_not_accumulate() {
4605 let (_temp, storage) = make_storage().await;
4606 let store = LockedSessionStore::new(storage);
4607
4608 for i in 0..100 {
4610 let _guard = store.acquire_lock(&format!("sess-{i}")).await;
4611 }
4612 assert_eq!(
4613 store.locks.len(),
4614 0,
4615 "acquiring locks for many distinct ids must not grow the map"
4616 );
4617 }
4618
4619 #[tokio::test]
4620 async fn acquire_lock_concurrent_waiter_keeps_valid_lock_and_map_drains() {
4621 use std::sync::atomic::{AtomicUsize, Ordering};
4622
4623 let (_temp, storage) = make_storage().await;
4624 let store = Arc::new(LockedSessionStore::new(storage));
4625
4626 let active = Arc::new(AtomicUsize::new(0));
4628 let max_seen = Arc::new(AtomicUsize::new(0));
4629
4630 let mut handles = Vec::new();
4631 for _ in 0..8 {
4632 let store = store.clone();
4633 let active = active.clone();
4634 let max_seen = max_seen.clone();
4635 handles.push(tokio::spawn(async move {
4636 let _guard = store.acquire_lock("contended").await;
4637 let now = active.fetch_add(1, Ordering::SeqCst) + 1;
4638 max_seen.fetch_max(now, Ordering::SeqCst);
4639 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
4641 active.fetch_sub(1, Ordering::SeqCst);
4642 }));
4643 }
4644 for h in handles {
4645 h.await.unwrap();
4646 }
4647
4648 assert_eq!(
4653 max_seen.load(Ordering::SeqCst),
4654 1,
4655 "at most one holder of a given session lock at a time"
4656 );
4657 assert_eq!(
4658 store.locks.len(),
4659 0,
4660 "after all holders release, the contended entry must be fully evicted"
4661 );
4662 }
4663}