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 ROOT_PROJECT_CONTEXT_KEYS: &[&str] = &[
49 "workspace_source",
50 "workspace_binding_status",
51 "project_context_rendered",
52 "project_resources_rendered",
53 "runtime_prompt_snapshot",
54];
55const RESPONSE_CONTROL_METADATA_KEYS: &[&str] = &[
56 CONSUMED_CLARIFICATION_IDS_KEY,
57 CONSUMED_RESPONSE_OCCURRENCES_KEY,
58 "runtime.suspend_reason",
59 "clarification_resume_pending",
60 "conclusion_with_options_resume_pending",
61 "execute.startup_handoff_at",
62 "permission.reexecute_tool_call_id",
63 "permission.reexecute_request_generation",
64 "retry_resume_pending",
65 "retry_resume_reason",
66 "provider_name",
67];
68const TASK_CONTROL_PLANE_CONFLICT_PREFIX: &str = "Task control-plane changed while saving session ";
69const MAX_TASK_CONTROL_PLANE_REBASE_RETRIES: usize = 3;
70
71fn may_publish_runtime_result(result: &std::io::Result<()>) -> bool {
72 !result.as_ref().err().is_some_and(|error| {
73 error
74 .get_ref()
75 .is_some_and(|cause| cause.is::<bamboo_domain::SessionAuthorityConflict>())
76 })
77}
78
79fn adopt_durable_consumed_clarification(session: &mut Session, durable: &Session) -> bool {
88 let Some(incoming_tool_call_id) = session
89 .pending_question
90 .as_ref()
91 .map(|pending| pending.tool_call_id.clone())
92 else {
93 return false;
94 };
95 let occurrence_ledger = durable
96 .metadata
97 .get(CONSUMED_RESPONSE_OCCURRENCES_KEY)
98 .map(|value| serde_json::from_str::<Vec<ResponseOccurrence>>(value));
99 let was_consumed = match occurrence_ledger {
100 Some(Ok(consumed)) => latest_response_occurrence(session, &incoming_tool_call_id)
101 .is_some_and(|incoming| consumed.iter().any(|entry| entry == &incoming)),
102 Some(Err(_)) => false,
105 None => {
106 let legacy_consumed = durable
107 .metadata
108 .get(CONSUMED_CLARIFICATION_IDS_KEY)
109 .and_then(|value| serde_json::from_str::<Vec<String>>(value).ok())
110 .unwrap_or_default()
111 .iter()
112 .any(|tool_call_id| tool_call_id == &incoming_tool_call_id);
113 legacy_consumed
117 && latest_response_occurrence(durable, &incoming_tool_call_id).is_some_and(
118 |durable_occurrence| {
119 latest_response_occurrence(session, &incoming_tool_call_id)
120 .is_some_and(|incoming| incoming == durable_occurrence)
121 },
122 )
123 }
124 };
125 if !was_consumed {
126 return false;
127 }
128
129 bamboo_domain::append_missing_runtime_messages(session, durable);
132 session
133 .pending_question
134 .clone_from(&durable.pending_question);
135 for key in RESPONSE_CONTROL_METADATA_KEYS {
136 if let Some(value) = durable.metadata.get(*key) {
137 session.metadata.insert((*key).to_string(), value.clone());
138 } else {
139 session.metadata.remove(*key);
140 }
141 }
142 session.model.clone_from(&durable.model);
143 session.model_ref.clone_from(&durable.model_ref);
144 session.reasoning_effort = durable.reasoning_effort;
145 session
146 .agent_runtime_state
147 .clone_from(&durable.agent_runtime_state);
148 if let Some(runtime_metadata) = session.runtime_metadata.as_mut() {
149 runtime_metadata.provider_name = durable
150 .runtime_metadata
151 .as_ref()
152 .and_then(|metadata| metadata.provider_name.clone());
153 } else if durable
154 .runtime_metadata
155 .as_ref()
156 .is_some_and(|metadata| metadata.provider_name.is_some())
157 {
158 session.runtime_metadata = durable.runtime_metadata.as_ref().map(|metadata| {
159 let mut response_metadata = bamboo_domain::SessionRuntimeMetadata::default();
160 response_metadata
161 .provider_name
162 .clone_from(&metadata.provider_name);
163 response_metadata
164 });
165 }
166 true
167}
168
169fn is_task_control_plane_save_conflict(error: &std::io::Error) -> bool {
170 error.kind() == std::io::ErrorKind::WouldBlock
171 && error
172 .to_string()
173 .starts_with(TASK_CONTROL_PLANE_CONFLICT_PREFIX)
174}
175
176fn adopt_durable_task_control_plane(session: &mut Session, durable: &Session) {
177 session.task_list = durable.task_list.clone();
178 session
179 .metadata
180 .remove(bamboo_domain::session::runtime_metadata::keys::TASK_LIST_VERSION);
181 if let Some(runtime_metadata) = session.runtime_metadata.as_mut() {
182 runtime_metadata.task_list_version = None;
183 }
184 if session
185 .runtime_metadata
186 .as_ref()
187 .is_some_and(bamboo_domain::session::SessionRuntimeMetadata::is_empty)
188 {
189 session.runtime_metadata = None;
190 }
191 if let Some(version) = durable.task_list_version_meta() {
192 session.set_task_list_version_meta(version);
193 }
194}
195
196fn task_list_snapshot_matches(
197 session: &Session,
198 expected_task_list: &bamboo_domain::TaskList,
199) -> std::io::Result<bool> {
200 Ok(serde_json::to_value(&session.task_list)
201 .map_err(|error| std::io::Error::other(error.to_string()))?
202 == serde_json::to_value(Some(expected_task_list))
203 .map_err(|error| std::io::Error::other(error.to_string()))?)
204}
205
206fn unconditional_task_patch_would_regress(
207 durable: &Session,
208 incoming_task_list: &bamboo_domain::TaskList,
209 incoming_version: &str,
210) -> std::io::Result<bool> {
211 let same_list = task_list_snapshot_matches(durable, incoming_task_list)?;
212 let Some(durable_version) = durable.task_list_version_meta() else {
213 return Ok(false);
214 };
215 match (
216 incoming_version.parse::<u64>(),
217 durable_version.parse::<u64>(),
218 ) {
219 (Ok(incoming), Ok(durable)) => {
220 Ok(incoming < durable || (incoming == durable && !same_list))
221 }
222 _ => Ok(incoming_version != durable_version || !same_list),
223 }
224}
225
226pub struct LockedSessionStore {
234 storage: Arc<dyn Storage>,
235 locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
236 task_pair_transaction_lock: Arc<Mutex<()>>,
240}
241
242pub struct SessionLockGuard {
265 guard: Option<OwnedMutexGuard<()>>,
267 locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
268 session_id: String,
269}
270
271impl Drop for SessionLockGuard {
272 fn drop(&mut self) {
273 self.guard.take();
276 self.locks
277 .remove_if(&self.session_id, |_, arc| Arc::strong_count(arc) == 1);
278 }
279}
280
281impl LockedSessionStore {
282 pub fn new(storage: Arc<dyn Storage>) -> Self {
284 Self {
285 storage,
286 locks: Arc::new(DashMap::new()),
287 task_pair_transaction_lock: Arc::new(Mutex::new(())),
288 }
289 }
290
291 pub fn storage(&self) -> &Arc<dyn Storage> {
293 &self.storage
294 }
295
296 pub async fn acquire_lock(&self, session_id: &str) -> SessionLockGuard {
308 let lock = self
313 .locks
314 .entry(session_id.to_string())
315 .or_insert_with(|| Arc::new(Mutex::new(())))
316 .clone();
317 let mut guard = SessionLockGuard {
321 guard: None,
322 locks: self.locks.clone(),
323 session_id: session_id.to_string(),
324 };
325 guard.guard = Some(lock.lock_owned().await);
326 guard
327 }
328
329 async fn save_session_rebasing_task_conflicts(
334 &self,
335 session: &mut Session,
336 ) -> std::io::Result<()> {
337 for attempt in 0..=MAX_TASK_CONTROL_PLANE_REBASE_RETRIES {
338 match self.storage.save_session(session).await {
339 Ok(()) => return Ok(()),
340 Err(error)
341 if is_task_control_plane_save_conflict(&error)
342 && attempt < MAX_TASK_CONTROL_PLANE_REBASE_RETRIES =>
343 {
344 let Some(durable) =
345 self.storage.load_runtime_control_plane(&session.id).await?
346 else {
347 return Err(error);
348 };
349 adopt_durable_task_control_plane(session, &durable);
350 }
351 Err(error) => {
352 if is_task_control_plane_save_conflict(&error) {
353 if let Some(durable) =
354 self.storage.load_runtime_control_plane(&session.id).await?
355 {
356 adopt_durable_task_control_plane(session, &durable);
357 }
358 }
359 return Err(error);
360 }
361 }
362 }
363 unreachable!("bounded Task conflict retry loop always returns")
364 }
365
366 async fn save_runtime_state_rebasing_task_conflicts(
369 &self,
370 session: &mut Session,
371 ) -> std::io::Result<()> {
372 for attempt in 0..=MAX_TASK_CONTROL_PLANE_REBASE_RETRIES {
373 match self.storage.save_runtime_state(session).await {
374 Ok(()) => return Ok(()),
375 Err(error)
376 if is_task_control_plane_save_conflict(&error)
377 && attempt < MAX_TASK_CONTROL_PLANE_REBASE_RETRIES =>
378 {
379 let Some(durable) =
380 self.storage.load_runtime_control_plane(&session.id).await?
381 else {
382 return Err(error);
383 };
384 adopt_durable_task_control_plane(session, &durable);
385 }
386 Err(error) => {
387 if is_task_control_plane_save_conflict(&error) {
388 if let Some(durable) =
389 self.storage.load_runtime_control_plane(&session.id).await?
390 {
391 adopt_durable_task_control_plane(session, &durable);
392 }
393 }
394 return Err(error);
395 }
396 }
397 }
398 unreachable!("bounded Task conflict retry loop always returns")
399 }
400
401 pub async fn save_runtime_only(&self, session: &mut Session) -> std::io::Result<()> {
419 self.save_runtime_only_and_publish(session, |_| {}).await
420 }
421
422 pub async fn save_runtime_only_and_publish<F>(
433 &self,
434 session: &mut Session,
435 publish: F,
436 ) -> std::io::Result<()>
437 where
438 F: FnOnce(&Session) + Send,
439 {
440 let _guard = self.acquire_lock(&session.id).await;
441 if let Some(latest) = self.storage.load_runtime_control_plane(&session.id).await? {
442 apply_authoritative_metadata(session, &latest);
443 adopt_fresher_disk_permission_posture(session, &latest);
446 adopt_durable_model_context_state(session, &latest);
451 }
452 let result = self
453 .save_runtime_state_rebasing_task_conflicts(session)
454 .await;
455 if may_publish_runtime_result(&result) {
456 publish(session);
457 }
458 result
459 }
460
461 pub async fn update_task_list_control_plane_and_publish<F>(
468 &self,
469 session_id: &str,
470 task_list: &bamboo_domain::TaskList,
471 version: &str,
472 publish: F,
473 ) -> std::io::Result<bool>
474 where
475 F: FnOnce(&Session) + Send,
476 {
477 let _guard = self.acquire_lock(session_id).await;
478 let Some(mut latest) = self.storage.load_runtime_control_plane(session_id).await? else {
479 return Ok(false);
480 };
481 if unconditional_task_patch_would_regress(&latest, task_list, version)? {
482 return Err(std::io::Error::new(
483 std::io::ErrorKind::WouldBlock,
484 format!("Task control-plane changed while patching session {session_id}"),
485 ));
486 }
487 let original = latest.clone();
488 latest.task_list = Some(task_list.clone());
489 latest.set_task_list_version_meta(version.to_string());
490 if !self
491 .storage
492 .save_task_control_plane_if_matches(&original, &latest)
493 .await?
494 {
495 return Err(std::io::Error::new(
496 std::io::ErrorKind::WouldBlock,
497 format!("Task control-plane changed while patching session {session_id}"),
498 ));
499 }
500 publish(&latest);
501 Ok(true)
502 }
503
504 pub async fn update_task_list_control_plane_if_version_and_publish<F>(
508 &self,
509 session_id: &str,
510 expected_version: &str,
511 expected_task_list: &bamboo_domain::TaskList,
512 task_list: &bamboo_domain::TaskList,
513 version: &str,
514 publish: F,
515 ) -> std::io::Result<bool>
516 where
517 F: FnOnce(&Session) + Send,
518 {
519 let _guard = self.acquire_lock(session_id).await;
520 let Some(mut latest) = self.storage.load_runtime_control_plane(session_id).await? else {
521 return Ok(false);
522 };
523 if latest.task_list_version_meta().as_deref() != Some(expected_version)
524 || !task_list_snapshot_matches(&latest, expected_task_list)?
525 {
526 return Ok(false);
527 }
528 let original = latest.clone();
529 latest.task_list = Some(task_list.clone());
530 latest.set_task_list_version_meta(version.to_string());
531 if !self
532 .storage
533 .save_task_control_plane_if_matches(&original, &latest)
534 .await?
535 {
536 return Ok(false);
537 }
538 publish(&latest);
539 Ok(true)
540 }
541
542 pub async fn update_task_list_control_planes_if_version_and_publish<F>(
548 &self,
549 session_id: &str,
550 shared_session_id: &str,
551 expected_version: &str,
552 expected_task_list: &bamboo_domain::TaskList,
553 task_list: &bamboo_domain::TaskList,
554 version: &str,
555 publish: F,
556 ) -> std::io::Result<bool>
557 where
558 F: FnOnce(&Session, &Session) + Send,
559 {
560 if session_id == shared_session_id {
561 return self
562 .update_task_list_control_plane_if_version_and_publish(
563 session_id,
564 expected_version,
565 expected_task_list,
566 task_list,
567 version,
568 |session| publish(session, session),
569 )
570 .await;
571 }
572
573 let (first_id, second_id) = if session_id < shared_session_id {
574 (session_id, shared_session_id)
575 } else {
576 (shared_session_id, session_id)
577 };
578 let _transaction_guard = self.task_pair_transaction_lock.lock().await;
579 let _first_guard = self.acquire_lock(first_id).await;
580 let _second_guard = self.acquire_lock(second_id).await;
581
582 self.storage
586 .recover_task_control_plane_transaction(first_id, second_id)
587 .await?;
588
589 let Some(mut local) = self.storage.load_runtime_control_plane(session_id).await? else {
590 return Ok(false);
591 };
592 let Some(mut shared) = self
593 .storage
594 .load_runtime_control_plane(shared_session_id)
595 .await?
596 else {
597 return Ok(false);
598 };
599 if local.task_list_version_meta().as_deref() != Some(expected_version)
600 || shared.task_list_version_meta().as_deref() != Some(expected_version)
601 || !task_list_snapshot_matches(&local, expected_task_list)?
602 || !task_list_snapshot_matches(&shared, expected_task_list)?
603 {
604 return Ok(false);
605 }
606
607 let local_original = local.clone();
608 let shared_original = shared.clone();
609 local.task_list = Some(task_list.clone());
613 local.set_task_list_version_meta(version.to_string());
614 shared.task_list = Some(task_list.clone());
615 shared.set_task_list_version_meta(version.to_string());
616 let (first_original, first_updated, second_original, second_updated) =
617 if session_id < shared_session_id {
618 (&local_original, &local, &shared_original, &shared)
619 } else {
620 (&shared_original, &shared, &local_original, &local)
621 };
622 let committed = self
623 .storage
624 .save_task_control_planes_atomically(
625 first_original,
626 first_updated,
627 second_original,
628 second_updated,
629 )
630 .await?;
631 if !committed {
632 return Ok(false);
633 }
634 publish(&local, &shared);
635 Ok(true)
636 }
637
638 pub async fn commit_metadata(&self, session: &Session) -> std::io::Result<()> {
648 let _guard = self.acquire_lock(&session.id).await;
649 let mut committed = session.clone();
650 if let Some(latest) = self.storage.load_runtime_control_plane(&session.id).await? {
651 adopt_durable_model_context_state(&mut committed, &latest);
652 }
653 self.save_session_rebasing_task_conflicts(&mut committed)
654 .await
655 }
656
657 pub async fn merge_save_runtime(&self, session: &mut Session) -> std::io::Result<()> {
674 self.merge_save_runtime_and_publish(session, |_, _| {})
675 .await
676 }
677
678 pub async fn merge_save_runtime_and_publish<F>(
686 &self,
687 session: &mut Session,
688 publish: F,
689 ) -> std::io::Result<()>
690 where
691 F: FnOnce(&Session, bool) + Send,
692 {
693 self.merge_save_runtime_inner_and_publish(session, true, publish)
694 .await
695 }
696
697 pub async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
706 self.checkpoint_runtime_session_and_publish(session, |_, _| {})
707 .await
708 }
709
710 pub async fn checkpoint_runtime_session_and_publish<F>(
717 &self,
718 session: &mut Session,
719 publish: F,
720 ) -> std::io::Result<()>
721 where
722 F: FnOnce(&Session, bool) + Send,
723 {
724 let _guard = self.acquire_lock(&session.id).await;
725 let latest = self.storage.load_session(&session.id).await?;
726
727 if let Some(latest) = latest.as_ref() {
728 ensure_model_context_checkpoint_is_current(session, latest)?;
729 let incoming_count = session.messages.len();
730 let durable_count = latest.messages.len();
731 let appended = bamboo_domain::append_missing_runtime_messages(session, latest);
732 bamboo_domain::merge_session_inbox_admission(session, latest);
733 let adopted_response = adopt_durable_consumed_clarification(session, latest);
734 tracing::debug!(
735 "[{}] append-safe runtime checkpoint: durable={}, incoming={}, appended={}, adopted_response={}, saved={}",
736 session.id,
737 durable_count,
738 incoming_count,
739 appended,
740 adopted_response,
741 session.messages.len(),
742 );
743 apply_authoritative_metadata(session, latest);
744 adopt_fresher_disk_permission_posture(session, latest);
745 }
746
747 let result = self.save_session_rebasing_task_conflicts(session).await;
748 if may_publish_runtime_result(&result) {
749 publish(session, result.is_ok());
750 }
751 result
752 }
753
754 pub async fn save_runtime_authoritative_flags(
763 &self,
764 session: &mut Session,
765 ) -> std::io::Result<()> {
766 self.merge_save_runtime_inner_and_publish(session, false, |_, _| {})
767 .await
768 }
769
770 async fn merge_save_runtime_inner_and_publish<F>(
771 &self,
772 session: &mut Session,
773 adopt_bypass: bool,
774 publish: F,
775 ) -> std::io::Result<()>
776 where
777 F: FnOnce(&Session, bool) + Send,
778 {
779 let _guard = self.acquire_lock(&session.id).await;
780
781 let latest = self.storage.load_session(&session.id).await?;
788
789 let existing_message_count = latest.as_ref().map(|s| s.messages.len());
795 let incoming_message_count = session.messages.len();
796 if existing_message_count.is_some_and(|existing| existing > incoming_message_count) {
797 tracing::warn!(
798 "[{}] merge_save_runtime SHRINK: disk has {:?} messages, saving {} (last_role={:?}, updated_at={}); a stale writer is reverting a concurrent append",
799 session.id,
800 existing_message_count,
801 incoming_message_count,
802 session.messages.last().map(|m| format!("{:?}", m.role)),
803 session.updated_at,
804 );
805 } else {
806 tracing::debug!(
807 "[{}] merge_save_runtime: disk={:?} messages, saving {} (updated_at={})",
808 session.id,
809 existing_message_count,
810 incoming_message_count,
811 session.updated_at,
812 );
813 }
814
815 if let Some(latest) = latest.as_ref() {
816 adopt_durable_consumed_clarification(session, latest);
817 apply_authoritative_metadata(session, latest);
818 let restored = bamboo_domain::restore_missing_admitted_inbox_messages(session, latest);
819 if restored > 0 {
820 tracing::warn!(
821 session_id = %session.id,
822 restored,
823 "restored durable SessionInbox transcript messages into stale runtime save"
824 );
825 }
826 bamboo_domain::merge_session_inbox_admission(session, latest);
827 if adopt_bypass {
832 adopt_fresher_disk_permission_posture(session, latest);
833 }
834 adopt_fresher_durable_model_context_state(session, latest);
835 }
836 let result = self.save_session_rebasing_task_conflicts(session).await;
837 if may_publish_runtime_result(&result) {
838 publish(session, result.is_ok());
839 }
840 result
841 }
842
843 pub async fn seed_runtime_activation_and_publish<F>(
853 &self,
854 session: &mut Session,
855 publish: F,
856 ) -> std::io::Result<()>
857 where
858 F: FnOnce(&Session, bool) + Send,
859 {
860 let _guard = self.acquire_lock(&session.id).await;
861 let mut incoming_audit = PermissionAuditSnapshot::from_metadata(&session.metadata)
862 .ok_or_else(|| {
863 std::io::Error::new(
864 std::io::ErrorKind::InvalidInput,
865 "activation seed requires a complete permission audit record",
866 )
867 })?;
868
869 if let Some(latest) = self.storage.load_session(&session.id).await? {
870 apply_authoritative_metadata(session, &latest);
871 bamboo_domain::restore_missing_admitted_inbox_messages(session, &latest);
872 bamboo_domain::merge_session_inbox_admission(session, &latest);
873 adopt_fresher_durable_model_context_state(session, &latest);
874
875 let durable_audit = PermissionAuditSnapshot::from_metadata(&latest.metadata);
876 let durable_floor = durable_audit
877 .as_ref()
878 .map(|snapshot| snapshot.audit_revision)
879 .unwrap_or_default();
880 if let Some(durable_audit) = durable_audit {
881 if durable_audit.resolution == incoming_audit.resolution {
882 incoming_audit.transitioned_at = durable_audit.transitioned_at;
883 }
884 }
885 incoming_audit.audit_revision = bamboo_domain::next_permission_audit_revision_after(
886 durable_floor.max(incoming_audit.audit_revision),
887 )
888 .map_err(|error| {
889 std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string())
890 })?;
891 }
892
893 session
894 .agent_runtime_state
895 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
896 .set_permission_mode(incoming_audit.resolution.requested);
897 incoming_audit.write_to(&mut session.metadata);
898
899 let result = self.save_session_rebasing_task_conflicts(session).await;
900 if may_publish_runtime_result(&result) {
901 publish(session, result.is_ok());
902 }
903 result
904 }
905
906 pub async fn update_authoritative_permission_posture_and_publish<M, P>(
912 &self,
913 session_id: &str,
914 seed: &PermissionAuditSeed,
915 mutate: M,
916 publish: P,
917 ) -> std::io::Result<Option<Session>>
918 where
919 M: FnOnce(&mut Session),
920 P: FnOnce(&Session),
921 {
922 let _guard = self.acquire_lock(session_id).await;
923 let Some(mut latest) = self.storage.load_session(session_id).await? else {
924 return Ok(None);
925 };
926 let previous_mode = latest
927 .agent_runtime_state
928 .as_ref()
929 .map(|state| state.effective_permission_mode())
930 .unwrap_or_default();
931 let previous_resolution = PermissionAuditSnapshot::from_metadata(&latest.metadata)
932 .map(|snapshot| snapshot.resolution);
933 mutate(&mut latest);
934 latest
935 .agent_runtime_state
936 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
937 .set_permission_mode(seed.resolution.requested);
938 let mode_changed = previous_mode != seed.resolution.requested;
939 let posture_changed = previous_resolution != Some(seed.resolution);
940 let transitioned_at = posture_changed.then(|| chrono::Utc::now().to_rfc3339());
941 bamboo_domain::record_permission_audit(
942 &mut latest.metadata,
943 seed,
944 transitioned_at.as_deref(),
945 )
946 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string()))?;
947 if mode_changed {
948 latest.metadata_version = latest.metadata_version.saturating_add(1);
949 }
950 self.save_session_rebasing_task_conflicts(&mut latest)
951 .await?;
952 publish(&latest);
953 Ok(Some(latest))
954 }
955
956 pub async fn record_permission_posture_activation_and_publish<P>(
961 &self,
962 session_id: &str,
963 expected_audit_revision: Option<u64>,
964 seed: &PermissionAuditSeed,
965 publish: P,
966 ) -> std::io::Result<Option<Session>>
967 where
968 P: FnOnce(&Session),
969 {
970 let _guard = self.acquire_lock(session_id).await;
971 let Some(mut latest) = self.storage.load_session(session_id).await? else {
972 return Ok(None);
973 };
974 let durable_audit = PermissionAuditSnapshot::from_metadata(&latest.metadata);
975 let durable_revision = durable_audit
976 .as_ref()
977 .map(|snapshot| snapshot.audit_revision);
978 if durable_revision != expected_audit_revision {
979 return Err(std::io::Error::new(
980 std::io::ErrorKind::InvalidData,
981 "stale permission posture activation: durable audit changed after dispatch",
982 ));
983 }
984 let durable_requested = latest
985 .agent_runtime_state
986 .as_ref()
987 .map(|state| state.effective_permission_mode())
988 .unwrap_or_default();
989 if durable_requested != seed.resolution.requested || !seed.resolution.is_consistent() {
990 return Err(std::io::Error::new(
991 std::io::ErrorKind::InvalidData,
992 "stale or inconsistent permission posture activation",
993 ));
994 }
995 bamboo_domain::record_permission_audit(&mut latest.metadata, seed, None).map_err(
996 |error| std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string()),
997 )?;
998 self.save_session_rebasing_task_conflicts(&mut latest)
999 .await?;
1000 publish(&latest);
1001 Ok(Some(latest))
1002 }
1003
1004 pub async fn update_runtime_config<F>(
1016 &self,
1017 session_id: &str,
1018 mutate: F,
1019 ) -> std::io::Result<Option<Session>>
1020 where
1021 F: FnOnce(&mut Session),
1022 {
1023 self.update_runtime_config_and_publish(session_id, mutate, |_| {})
1024 .await
1025 }
1026
1027 pub async fn update_runtime_config_and_publish<M, P>(
1030 &self,
1031 session_id: &str,
1032 mutate: M,
1033 publish: P,
1034 ) -> std::io::Result<Option<Session>>
1035 where
1036 M: FnOnce(&mut Session),
1037 P: FnOnce(&Session),
1038 {
1039 let _guard = self.acquire_lock(session_id).await;
1040 let Some(mut session) = self.storage.load_session(session_id).await? else {
1041 return Ok(None);
1042 };
1043 mutate(&mut session);
1044 self.save_session_rebasing_task_conflicts(&mut session)
1045 .await?;
1046 publish(&session);
1047 Ok(Some(session))
1048 }
1049
1050 async fn load_response_candidate<C>(
1055 &self,
1056 session_id: &str,
1057 load_cached: C,
1058 ) -> std::io::Result<Option<Session>>
1059 where
1060 C: FnOnce() -> Option<Session> + Send,
1061 {
1062 let cached_candidate = load_cached();
1063 let durable = self.storage.load_session(session_id).await?;
1064 let Some(mut session) = (match (cached_candidate, durable.as_ref()) {
1065 (Some(cached), Some(durable)) => {
1066 let prefer_durable = durable.updated_at > cached.updated_at
1067 || (durable.updated_at == cached.updated_at
1068 && cached.pending_question.is_none()
1069 && durable.pending_question.is_some());
1070 Some(if prefer_durable {
1071 durable.clone()
1072 } else {
1073 cached
1074 })
1075 }
1076 (Some(cached), None) => Some(cached),
1077 (None, durable) => durable.cloned(),
1078 }) else {
1079 return Ok(None);
1080 };
1081 if let Some(latest) = durable.as_ref() {
1082 adopt_durable_consumed_clarification(&mut session, latest);
1087 bamboo_domain::append_missing_runtime_messages(&mut session, latest);
1088 apply_authoritative_metadata(&mut session, latest);
1089 let restored =
1090 bamboo_domain::restore_missing_admitted_inbox_messages(&mut session, latest);
1091 if restored > 0 {
1092 tracing::warn!(
1093 session_id,
1094 restored,
1095 "restored durable SessionInbox transcript messages into response transaction"
1096 );
1097 }
1098 bamboo_domain::merge_session_inbox_admission(&mut session, latest);
1099 adopt_fresher_disk_permission_posture(&mut session, latest);
1100 adopt_fresher_durable_model_context_state(&mut session, latest);
1101 }
1102 Ok(Some(session))
1103 }
1104
1105 pub async fn inspect_runtime_session_for_response<C>(
1110 &self,
1111 session_id: &str,
1112 load_cached: C,
1113 ) -> std::io::Result<Option<Session>>
1114 where
1115 C: FnOnce() -> Option<Session> + Send,
1116 {
1117 let _guard = self.acquire_lock(session_id).await;
1118 self.load_response_candidate(session_id, load_cached).await
1119 }
1120
1121 pub async fn mutate_runtime_session_and_publish<C, M, P, E>(
1131 &self,
1132 session_id: &str,
1133 load_cached: C,
1134 mutate: M,
1135 publish: P,
1136 ) -> std::io::Result<Result<Option<Session>, E>>
1137 where
1138 C: FnOnce() -> Option<Session> + Send,
1139 M: FnOnce(&mut Session) -> Result<(), E> + Send,
1140 P: FnOnce(&Session) + Send,
1141 E: Send,
1142 {
1143 let _guard = self.acquire_lock(session_id).await;
1144 let Some(mut session) = self
1145 .load_response_candidate(session_id, load_cached)
1146 .await?
1147 else {
1148 return Ok(Ok(None));
1149 };
1150 if let Err(error) = mutate(&mut session) {
1151 return Ok(Err(error));
1152 }
1153 self.save_session_rebasing_task_conflicts(&mut session)
1154 .await?;
1155 publish(&session);
1156 Ok(Ok(Some(session)))
1157 }
1158
1159 pub async fn clear_legacy_pending_messages_and_publish<F>(
1162 &self,
1163 session_id: &str,
1164 expected: &[serde_json::Value],
1165 publish: F,
1166 ) -> std::io::Result<bool>
1167 where
1168 F: FnOnce(&Session) + Send,
1169 {
1170 let _guard = self.acquire_lock(session_id).await;
1171 let Some(mut latest) = self.storage.load_session(session_id).await? else {
1172 return Ok(false);
1173 };
1174 if latest.pending_injected_messages().as_deref() != Some(expected) {
1175 return Ok(false);
1176 }
1177 latest.clear_pending_injected_messages();
1178 self.save_runtime_state_rebasing_task_conflicts(&mut latest)
1179 .await?;
1180 publish(&latest);
1181 Ok(true)
1182 }
1183}
1184
1185#[async_trait::async_trait]
1189impl RuntimeSessionPersistence for LockedSessionStore {
1190 async fn save_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
1191 self.merge_save_runtime(session).await
1192 }
1193
1194 async fn seed_runtime_activation(&self, session: &mut Session) -> std::io::Result<()> {
1195 self.seed_runtime_activation_and_publish(session, |_, _| {})
1196 .await
1197 }
1198
1199 async fn record_permission_posture_activation(
1200 &self,
1201 session_id: &str,
1202 expected_audit_revision: Option<u64>,
1203 seed: &PermissionAuditSeed,
1204 ) -> std::io::Result<Option<Session>> {
1205 self.record_permission_posture_activation_and_publish(
1206 session_id,
1207 expected_audit_revision,
1208 seed,
1209 |_| {},
1210 )
1211 .await
1212 }
1213
1214 async fn save_runtime_control_plane(&self, session: &mut Session) -> std::io::Result<()> {
1215 self.save_runtime_only(session).await
1216 }
1217
1218 async fn load_runtime_control_plane(
1219 &self,
1220 session_id: &str,
1221 ) -> std::io::Result<Option<Session>> {
1222 self.storage.load_runtime_control_plane(session_id).await
1223 }
1224
1225 async fn update_task_list_control_plane(
1226 &self,
1227 session_id: &str,
1228 task_list: &bamboo_domain::TaskList,
1229 version: &str,
1230 ) -> std::io::Result<bool> {
1231 self.update_task_list_control_plane_and_publish(session_id, task_list, version, |_| {})
1232 .await
1233 }
1234
1235 async fn update_task_list_control_plane_if_version(
1236 &self,
1237 session_id: &str,
1238 expected_version: &str,
1239 expected_task_list: &bamboo_domain::TaskList,
1240 task_list: &bamboo_domain::TaskList,
1241 version: &str,
1242 ) -> std::io::Result<bool> {
1243 self.update_task_list_control_plane_if_version_and_publish(
1244 session_id,
1245 expected_version,
1246 expected_task_list,
1247 task_list,
1248 version,
1249 |_| {},
1250 )
1251 .await
1252 }
1253
1254 async fn update_task_list_control_planes_if_version(
1255 &self,
1256 session_id: &str,
1257 shared_session_id: &str,
1258 expected_version: &str,
1259 expected_task_list: &bamboo_domain::TaskList,
1260 task_list: &bamboo_domain::TaskList,
1261 version: &str,
1262 ) -> std::io::Result<bool> {
1263 self.update_task_list_control_planes_if_version_and_publish(
1264 session_id,
1265 shared_session_id,
1266 expected_version,
1267 expected_task_list,
1268 task_list,
1269 version,
1270 |_, _| {},
1271 )
1272 .await
1273 }
1274
1275 async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
1276 LockedSessionStore::checkpoint_runtime_session(self, session).await
1277 }
1278
1279 async fn load_runtime_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
1280 self.storage.load_session(session_id).await
1281 }
1282
1283 async fn clear_legacy_pending_messages(
1284 &self,
1285 session_id: &str,
1286 expected: &[serde_json::Value],
1287 ) -> std::io::Result<bool> {
1288 self.clear_legacy_pending_messages_and_publish(session_id, expected, |_| {})
1289 .await
1290 }
1291}
1292
1293async fn merge_authoritative_metadata_into_stale(
1302 storage: &Arc<dyn Storage>,
1303 session: &mut Session,
1304) -> std::io::Result<()> {
1305 if let Some(latest) = storage.load_session(&session.id).await? {
1306 adopt_durable_consumed_clarification(session, &latest);
1307 apply_authoritative_metadata(session, &latest);
1308 bamboo_domain::restore_missing_admitted_inbox_messages(session, &latest);
1309 bamboo_domain::merge_session_inbox_admission(session, &latest);
1310 adopt_fresher_disk_permission_posture(session, &latest);
1311 adopt_fresher_durable_model_context_state(session, &latest);
1312 }
1313 Ok(())
1314}
1315
1316fn adopt_durable_model_context_state(session: &mut Session, latest: &Session) {
1321 session
1322 .model_context_state
1323 .clone_from(&latest.model_context_state);
1324}
1325
1326fn adopt_fresher_durable_model_context_state(session: &mut Session, latest: &Session) {
1332 let adopt = match (
1333 session.model_context_state.as_ref(),
1334 latest.model_context_state.as_ref(),
1335 ) {
1336 (None, Some(_)) => true,
1337 (Some(incoming), Some(durable)) => {
1338 durable.state_revision > incoming.state_revision
1339 || (durable.state_revision == incoming.state_revision && durable != incoming)
1340 }
1341 _ => false,
1342 };
1343 if adopt {
1344 adopt_durable_model_context_state(session, latest);
1345 }
1346}
1347
1348fn ensure_model_context_checkpoint_is_current(
1353 session: &Session,
1354 latest: &Session,
1355) -> std::io::Result<()> {
1356 let stale_or_conflicting = match (
1357 session.model_context_state.as_ref(),
1358 latest.model_context_state.as_ref(),
1359 ) {
1360 (None, Some(_)) => true,
1361 (Some(incoming), Some(durable)) => {
1362 durable.state_revision > incoming.state_revision
1363 || (durable.state_revision == incoming.state_revision && durable != incoming)
1364 }
1365 _ => false,
1366 };
1367 if stale_or_conflicting {
1368 return Err(std::io::Error::new(
1369 std::io::ErrorKind::WouldBlock,
1370 "stale or conflicting model-context ledger checkpoint",
1371 ));
1372 }
1373 Ok(())
1374}
1375
1376fn adopt_fresher_disk_permission_posture(session: &mut Session, latest: &Session) {
1388 let Some(disk_mode) = latest
1393 .agent_runtime_state
1394 .as_ref()
1395 .map(|state| state.effective_permission_mode())
1396 else {
1397 return;
1398 };
1399 let current_mode = session
1400 .agent_runtime_state
1401 .as_ref()
1402 .map(|state| state.effective_permission_mode())
1403 .unwrap_or_default();
1404 let Some(disk_audit) = bamboo_domain::fresher_disk_permission_audit(
1405 current_mode,
1406 &session.metadata,
1407 disk_mode,
1408 &latest.metadata,
1409 ) else {
1410 return;
1411 };
1412
1413 match session.agent_runtime_state.as_mut() {
1414 Some(state) => state.set_permission_mode(disk_mode),
1415 None if disk_mode != bamboo_domain::SessionPermissionMode::Default => {
1418 let state = session
1419 .agent_runtime_state
1420 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default);
1421 state.set_permission_mode(disk_mode);
1422 }
1423 None => {}
1424 }
1425
1426 disk_audit.write_to(&mut session.metadata);
1428}
1429
1430fn apply_authoritative_metadata(session: &mut Session, latest: &Session) {
1436 if session.authority_identity.is_ordinary() && session.created_at == latest.created_at {
1443 session.authority_identity = latest.authority_identity.clone();
1444 }
1445 if session.kind == bamboo_domain::SessionKind::Root
1449 && latest.kind == bamboo_domain::SessionKind::Root
1450 && session.created_at == latest.created_at
1451 && session.authority_identity == latest.authority_identity
1452 {
1453 session.supervisor_management = latest.supervisor_management.clone();
1454 }
1455 if session.kind == bamboo_domain::SessionKind::Root
1461 && latest.kind == bamboo_domain::SessionKind::Root
1462 && session.created_at == latest.created_at
1463 && latest.metadata_version >= session.metadata_version
1464 && (latest.metadata_version > session.metadata_version
1465 || latest.project_id_meta() != session.project_id_meta())
1466 {
1467 match latest.project_id_meta() {
1468 Some(project) => session.set_project_id_meta(project),
1469 None => session.clear_project_id_meta(),
1470 }
1471 match latest.workspace_path_meta() {
1472 Some(workspace) => session.set_workspace_path_meta(workspace),
1473 None => {
1474 session.metadata.remove("workspace_path");
1475 if let Some(metadata) = session.runtime_metadata.as_mut() {
1476 metadata.workspace_path = None;
1477 }
1478 }
1479 }
1480 session.workspace.clone_from(&latest.workspace);
1481 for key in ROOT_PROJECT_CONTEXT_KEYS {
1482 match latest.metadata.get(*key) {
1483 Some(value) => {
1484 session.metadata.insert((*key).to_string(), value.clone());
1485 }
1486 None => {
1487 session.metadata.remove(*key);
1488 }
1489 }
1490 }
1491 session.prompt_snapshot.clone_from(&latest.prompt_snapshot);
1492 }
1493 if latest.metadata_version >= session.metadata_version {
1494 session.title = latest.title.clone();
1495 session.title_version = latest.title_version;
1496 session.title_generated = latest.title_generated;
1497 session.pinned = latest.pinned;
1498 for key in AUTHORITATIVE_METADATA_KEYS {
1499 if let Some(value) = latest.metadata.get(*key) {
1500 session.metadata.insert((*key).to_string(), value.clone());
1501 } else {
1502 session.metadata.remove(*key);
1503 }
1504 }
1505 session.metadata_version = latest.metadata_version;
1506 }
1507}
1508
1509pub async fn merge_save_session(
1522 storage: &Arc<dyn Storage>,
1523 session: &mut Session,
1524) -> std::io::Result<()> {
1525 merge_authoritative_metadata_into_stale(storage, session).await?;
1526 storage.save_session(session).await
1527}
1528
1529#[cfg(test)]
1532mod tests {
1533 use super::*;
1534 use crate::v2::{RuntimeTaskTransactionFault, SessionStoreV2};
1535 use bamboo_domain::{session::types::Session, PermissionMode};
1536 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
1537
1538 #[test]
1539 fn authority_merge_updates_ordinary_cache_but_does_not_hide_stale_incarnation() {
1540 let mut latest = Session::new(bamboo_domain::DEFAULT_SUPERVISOR_SESSION_ID, "model");
1541 latest.authority_identity = bamboo_domain::SessionAuthorityIdentity::Supervisor {
1542 incarnation_id: uuid::Uuid::new_v4(),
1543 };
1544 let mut stale = latest.clone();
1545 stale.authority_identity = bamboo_domain::SessionAuthorityIdentity::Ordinary;
1546 stale.metadata_version = 100;
1547 apply_authoritative_metadata(&mut stale, &latest);
1548 assert_eq!(stale.authority_identity, latest.authority_identity);
1549 let old_identity = bamboo_domain::SessionAuthorityIdentity::Supervisor {
1550 incarnation_id: uuid::Uuid::new_v4(),
1551 };
1552 stale.authority_identity = old_identity.clone();
1553 apply_authoritative_metadata(&mut stale, &latest);
1554 assert_eq!(stale.authority_identity, old_identity);
1555 }
1556
1557 struct AuthoritySavePauseStorage {
1558 inner: Arc<SessionStoreV2>,
1559 reached: tokio::sync::Barrier,
1560 release: tokio::sync::Barrier,
1561 }
1562
1563 #[tokio::test]
1564 async fn root_project_merge_publishes_project_revision_and_workspace_together() {
1565 for runtime_only in [false, true] {
1566 for equal_revision in [false, true] {
1567 let temp = tempfile::tempdir().unwrap();
1568 let storage = Arc::new(SessionStoreV2::new(temp.path().into()).await.unwrap());
1569 let mut stale = Session::new("root-project-merge", "model");
1570 stale.set_project_id_meta("project-a");
1571 stale.set_workspace_path_meta("/project-a");
1572 stale.metadata.insert(
1573 "runtime_prompt_snapshot".into(),
1574 "old Project A prompt".into(),
1575 );
1576 stale.add_message(bamboo_domain::Message::user("Keep transcript"));
1577 storage.save_session(&stale).await.unwrap();
1578 let mut current = stale.clone();
1579 current.metadata_version += 1;
1580 current.set_project_id_meta("project-b");
1581 current.set_workspace_path_meta("/project-b");
1582 current.metadata.remove("runtime_prompt_snapshot");
1583 current
1584 .metadata
1585 .insert("workspace_source".into(), "project_default".into());
1586 current
1587 .metadata
1588 .insert("project_context_rendered".into(), "Project B".into());
1589 storage.save_session(¤t).await.unwrap();
1590 if equal_revision {
1591 stale.metadata_version = current.metadata_version;
1594 }
1595 let locked = LockedSessionStore::new(storage.clone());
1596 let published = AtomicBool::new(false);
1597 let publish = |saved: &Session| {
1598 assert!(!saved.metadata.contains_key("runtime_prompt_snapshot"));
1599 assert_eq!(saved.project_id_meta().as_deref(), Some("project-b"));
1600 assert_eq!(saved.workspace_path_meta().as_deref(), Some("/project-b"));
1601 assert_eq!(saved.metadata_version, current.metadata_version);
1602 assert_eq!(
1603 saved.metadata.get("workspace_source").map(String::as_str),
1604 Some("project_default")
1605 );
1606 assert_eq!(
1607 saved
1608 .metadata
1609 .get("project_context_rendered")
1610 .map(String::as_str),
1611 Some("Project B")
1612 );
1613 published.store(true, Ordering::SeqCst);
1614 };
1615 if runtime_only {
1616 locked
1617 .save_runtime_only_and_publish(&mut stale, publish)
1618 .await
1619 .unwrap();
1620 } else {
1621 locked
1622 .merge_save_runtime_and_publish(&mut stale, |saved, committed| {
1623 assert!(committed);
1624 publish(saved);
1625 })
1626 .await
1627 .unwrap();
1628 }
1629 assert!(published.load(Ordering::SeqCst));
1630 assert_eq!(stale.project_id_meta(), current.project_id_meta());
1631 let loaded = storage.load_session(&stale.id).await.unwrap().unwrap();
1632 assert_eq!(loaded.project_id_meta(), current.project_id_meta());
1633 assert_eq!(loaded.messages.len(), 1);
1634 storage.flush_search_index().await;
1635 }
1636 }
1637 }
1638
1639 #[tokio::test]
1640 async fn root_project_change_after_merge_read_rejects_without_cache_publication() {
1641 for runtime_only in [false, true] {
1642 let temp = tempfile::tempdir().unwrap();
1643 let first = Arc::new(SessionStoreV2::new(temp.path().into()).await.unwrap());
1644 let mut stale = Session::new("root-project-race", "model");
1645 stale.set_project_id_meta("project-a");
1646 stale.add_message(bamboo_domain::Message::user("Keep transcript"));
1647 first.save_session(&stale).await.unwrap();
1648 let second = SessionStoreV2::new(temp.path().into()).await.unwrap();
1649 let mut current = stale.clone();
1650 current.metadata_version += 1;
1651 current.set_project_id_meta("project-b");
1652 let paused = Arc::new(AuthoritySavePauseStorage {
1653 inner: first.clone(),
1654 reached: tokio::sync::Barrier::new(2),
1655 release: tokio::sync::Barrier::new(2),
1656 });
1657 let locked = LockedSessionStore::new(paused.clone());
1658 let published = AtomicBool::new(false);
1659 let save = async {
1660 if runtime_only {
1661 locked
1662 .save_runtime_only_and_publish(&mut stale, |_| {
1663 published.store(true, Ordering::SeqCst);
1664 })
1665 .await
1666 } else {
1667 locked
1668 .merge_save_runtime_and_publish(&mut stale, |_, _| {
1669 published.store(true, Ordering::SeqCst);
1670 })
1671 .await
1672 }
1673 };
1674 let update = async {
1675 paused.reached.wait().await;
1676 second.save_session(¤t).await.unwrap();
1677 paused.release.wait().await;
1678 };
1679 let (result, ()) = tokio::time::timeout(std::time::Duration::from_secs(10), async {
1680 tokio::join!(save, update)
1681 })
1682 .await
1683 .expect("deterministic Project/save race completes");
1684 assert!(!may_publish_runtime_result(&Err(result.unwrap_err())));
1685 assert!(!published.load(Ordering::SeqCst));
1686 let loaded = first.load_session(¤t.id).await.unwrap().unwrap();
1687 assert_eq!(loaded.project_id_meta(), current.project_id_meta());
1688 assert_eq!(loaded.metadata_version, current.metadata_version);
1689 assert_eq!(loaded.messages.len(), 1);
1690 first.flush_search_index().await;
1691 second.flush_search_index().await;
1692 }
1693 }
1694
1695 #[async_trait::async_trait]
1696 impl Storage for AuthoritySavePauseStorage {
1697 async fn load_session(&self, id: &str) -> std::io::Result<Option<Session>> {
1698 self.inner.load_session(id).await
1699 }
1700 async fn load_runtime_control_plane(&self, id: &str) -> std::io::Result<Option<Session>> {
1701 self.inner.load_runtime_control_plane(id).await
1702 }
1703 async fn delete_session(&self, id: &str) -> std::io::Result<bool> {
1704 self.inner.delete_session(id).await
1705 }
1706 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
1707 self.reached.wait().await;
1708 self.release.wait().await;
1709 self.inner.save_session(session).await
1710 }
1711 async fn save_runtime_state(&self, session: &Session) -> std::io::Result<()> {
1712 self.reached.wait().await;
1713 self.release.wait().await;
1714 self.inner.save_runtime_state(session).await
1715 }
1716 }
1717
1718 #[tokio::test]
1719 async fn root_deleted_after_merge_read_rejects_without_cache_or_event_publication() {
1720 for runtime_only in [false, true] {
1721 let temp = tempfile::tempdir().unwrap();
1722 let first = Arc::new(SessionStoreV2::new(temp.path().into()).await.unwrap());
1723 let mut stale = Session::new("root-delete-race", "model");
1724 stale.set_project_id_meta("project-a");
1725 stale.add_message(bamboo_domain::Message::user("Old lifetime"));
1726 first.save_session(&stale).await.unwrap();
1727 let id = stale.id.clone();
1728 let second = SessionStoreV2::new(temp.path().into()).await.unwrap();
1729 let paused = Arc::new(AuthoritySavePauseStorage {
1730 inner: first.clone(),
1731 reached: tokio::sync::Barrier::new(2),
1732 release: tokio::sync::Barrier::new(2),
1733 });
1734 let locked = LockedSessionStore::new(paused.clone());
1735 let published = AtomicBool::new(false);
1736 let save = async {
1737 if runtime_only {
1738 locked
1739 .save_runtime_only_and_publish(&mut stale, |_| {
1740 published.store(true, Ordering::SeqCst);
1741 })
1742 .await
1743 } else {
1744 locked
1745 .merge_save_runtime_and_publish(&mut stale, |_, _| {
1746 published.store(true, Ordering::SeqCst);
1747 })
1748 .await
1749 }
1750 };
1751 let delete = async {
1752 paused.reached.wait().await;
1753 assert!(second.delete_session(&id).await.unwrap());
1754 paused.release.wait().await;
1755 };
1756 let (result, ()) = tokio::time::timeout(std::time::Duration::from_secs(10), async {
1757 tokio::join!(save, delete)
1758 })
1759 .await
1760 .expect("deterministic delete/save race completes");
1761 assert!(!may_publish_runtime_result(&Err(result.unwrap_err())));
1762 assert!(!published.load(Ordering::SeqCst));
1763 assert!(first.load_root_authority(&id).await.unwrap().is_none());
1764 assert!(!first.sessions_root_dir().join(&id).exists());
1765 first.flush_search_index().await;
1766 second.flush_search_index().await;
1767 }
1768 }
1769
1770 #[tokio::test]
1771 async fn supervisor_bootstrap_between_merge_read_and_save_rejects_without_publishing() {
1772 for runtime_only in [false, true] {
1773 let temp = tempfile::tempdir().unwrap();
1774 let first = Arc::new(SessionStoreV2::new(temp.path().into()).await.unwrap());
1775 let second = SessionStoreV2::new(temp.path().into()).await.unwrap();
1776 let paused = Arc::new(AuthoritySavePauseStorage {
1777 inner: first.clone(),
1778 reached: tokio::sync::Barrier::new(2),
1779 release: tokio::sync::Barrier::new(2),
1780 });
1781 let store = LockedSessionStore::new(paused.clone());
1782 let mut stale = Session::new(bamboo_domain::DEFAULT_SUPERVISOR_SESSION_ID, "stale");
1783 stale.add_message(bamboo_domain::Message::user(
1784 "must not enter new Supervisor",
1785 ));
1786 let published = AtomicBool::new(false);
1787 let save = async {
1788 if runtime_only {
1789 store
1790 .save_runtime_only_and_publish(&mut stale, |_| {
1791 published.store(true, Ordering::SeqCst);
1792 })
1793 .await
1794 } else {
1795 store
1796 .merge_save_runtime_and_publish(&mut stale, |_, _| {
1797 published.store(true, Ordering::SeqCst);
1798 })
1799 .await
1800 }
1801 };
1802 let bootstrap = async {
1803 paused.reached.wait().await;
1806 let receipt = second
1807 .get_or_create_default_supervisor("supervisor")
1808 .await
1809 .unwrap();
1810 paused.release.wait().await;
1811 receipt
1812 };
1813 let (result, receipt) =
1814 tokio::time::timeout(std::time::Duration::from_secs(10), async {
1815 tokio::join!(save, bootstrap)
1816 })
1817 .await
1818 .expect("deterministic bootstrap/save race completes");
1819 let error = result.unwrap_err();
1820 assert!(!may_publish_runtime_result(&Err(error)));
1821 assert!(!published.load(Ordering::SeqCst));
1822 assert!(stale.authority_identity.is_ordinary());
1823 let observed = first
1824 .load_root_authority(&receipt.session_id)
1825 .await
1826 .unwrap()
1827 .unwrap();
1828 let durable = second
1831 .load_session(&receipt.session_id)
1832 .await
1833 .unwrap()
1834 .unwrap();
1835 assert_eq!(durable.model, "supervisor");
1836 assert_eq!(observed.authority_identity, durable.authority_identity);
1837 assert!(durable.messages.is_empty());
1838 assert_eq!(
1839 durable.authority_identity,
1840 bamboo_domain::SessionAuthorityIdentity::Supervisor {
1841 incarnation_id: receipt.incarnation_id,
1842 }
1843 );
1844 }
1845 }
1846
1847 #[tokio::test]
1848 async fn supervisor_merge_publishes_adopted_identity_but_never_a_rejected_incarnation() {
1849 for runtime_only in [false, true] {
1850 let temp = tempfile::tempdir().unwrap();
1851 let storage = Arc::new(SessionStoreV2::new(temp.path().into()).await.unwrap());
1852 let receipt = storage
1853 .get_or_create_default_supervisor("model")
1854 .await
1855 .unwrap();
1856 let baseline = storage
1857 .load_session(&receipt.session_id)
1858 .await
1859 .unwrap()
1860 .unwrap();
1861 let expected = baseline.authority_identity.clone();
1862 let store = LockedSessionStore::new(storage.clone());
1863 for case in 0..3 {
1864 let rejected = case != 0;
1865 let mut snapshot = baseline.clone();
1866 snapshot.authority_identity = if case == 1 {
1867 bamboo_domain::SessionAuthorityIdentity::Supervisor {
1868 incarnation_id: uuid::Uuid::new_v4(),
1869 }
1870 } else {
1871 bamboo_domain::SessionAuthorityIdentity::Ordinary
1872 };
1873 if case == 2 {
1874 snapshot.created_at -= chrono::Duration::seconds(1);
1875 snapshot.model = "stale Ordinary instance".into();
1876 snapshot.add_message(bamboo_domain::Message::user("must not be rebound"));
1877 }
1878 let published = AtomicBool::new(false);
1879 let callback = |saved: &Session| {
1880 assert_eq!(saved.authority_identity, expected);
1881 published.store(true, Ordering::SeqCst);
1882 };
1883 let result = if runtime_only {
1884 store
1885 .save_runtime_only_and_publish(&mut snapshot, callback)
1886 .await
1887 } else {
1888 store
1889 .merge_save_runtime_and_publish(&mut snapshot, |saved, committed| {
1890 assert!(committed);
1891 callback(saved);
1892 })
1893 .await
1894 };
1895 assert_eq!(result.is_err(), rejected);
1896 assert_eq!(published.load(Ordering::SeqCst), !rejected);
1897 if !rejected {
1898 assert_eq!(snapshot.authority_identity, expected);
1899 }
1900 let durable = storage
1901 .load_session(&receipt.session_id)
1902 .await
1903 .unwrap()
1904 .unwrap();
1905 assert_eq!(durable.authority_identity, expected);
1906 assert_eq!(durable.created_at, baseline.created_at);
1907 assert_eq!(durable.model, baseline.model);
1908 assert!(durable.messages.is_empty());
1909 }
1910 }
1911 }
1912
1913 struct CountingControlPlaneStorage {
1914 inner: Arc<SessionStoreV2>,
1915 control_plane_loads: AtomicUsize,
1916 full_saves: AtomicUsize,
1917 runtime_state_saves: AtomicUsize,
1918 }
1919
1920 struct PairCommitBarrierStorage {
1921 inner: Arc<SessionStoreV2>,
1922 before_commit: Arc<tokio::sync::Barrier>,
1923 }
1924
1925 struct SingleCommitBarrierStorage {
1926 inner: Arc<SessionStoreV2>,
1927 before_commit: Arc<tokio::sync::Barrier>,
1928 }
1929
1930 struct SingleCommitPauseStorage {
1931 inner: Arc<SessionStoreV2>,
1932 commit_reached: Arc<tokio::sync::Barrier>,
1933 release_commit: Arc<tokio::sync::Barrier>,
1934 }
1935
1936 #[async_trait::async_trait]
1937 impl Storage for SingleCommitPauseStorage {
1938 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
1939 self.inner.save_session(session).await
1940 }
1941
1942 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
1943 self.inner.load_session(session_id).await
1944 }
1945
1946 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
1947 self.inner.delete_session(session_id).await
1948 }
1949
1950 async fn load_runtime_control_plane(
1951 &self,
1952 session_id: &str,
1953 ) -> std::io::Result<Option<Session>> {
1954 self.inner.load_runtime_control_plane(session_id).await
1955 }
1956
1957 async fn save_task_control_plane_if_matches(
1958 &self,
1959 original: &Session,
1960 updated: &Session,
1961 ) -> std::io::Result<bool> {
1962 self.commit_reached.wait().await;
1963 self.release_commit.wait().await;
1964 self.inner
1965 .save_task_control_plane_if_matches(original, updated)
1966 .await
1967 }
1968
1969 async fn save_task_control_planes_atomically(
1970 &self,
1971 first_original: &Session,
1972 first_updated: &Session,
1973 second_original: &Session,
1974 second_updated: &Session,
1975 ) -> std::io::Result<bool> {
1976 self.commit_reached.wait().await;
1977 self.release_commit.wait().await;
1978 self.inner
1979 .save_task_control_planes_atomically(
1980 first_original,
1981 first_updated,
1982 second_original,
1983 second_updated,
1984 )
1985 .await
1986 }
1987 }
1988
1989 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1990 async fn supervisor_management_race_rejects_staged_task_callbacks_and_fresh_invocation_succeeds(
1991 ) {
1992 use bamboo_domain::{
1993 SupervisorManagementMutation, SupervisorManagementRequest, SupervisorReference,
1994 };
1995
1996 for mode in 0..4 {
1997 let home = tempfile::tempdir().unwrap();
1998 let inner = Arc::new(SessionStoreV2::new(home.path().into()).await.unwrap());
1999 let receipt = inner
2000 .get_or_create_default_supervisor("model")
2001 .await
2002 .unwrap();
2003 let reference = SupervisorReference::from(&receipt);
2004 let mut root = inner
2005 .load_session(&reference.session_id)
2006 .await
2007 .unwrap()
2008 .unwrap();
2009 let list = bamboo_domain::TaskList {
2010 session_id: root.id.clone(),
2011 title: "original".into(),
2012 items: vec![],
2013 created_at: root.created_at,
2014 updated_at: root.created_at,
2015 };
2016 let updated = bamboo_domain::TaskList {
2017 title: "updated".into(),
2018 ..list.clone()
2019 };
2020 root.task_list = Some(list.clone());
2021 root.set_task_list_version_meta("1");
2022 inner.save_session(&root).await.unwrap();
2023 let child_id = if mode == 2 { "aaa-child" } else { "zzz-child" };
2024 let mut child = Session::new_child_of(child_id, &root, "model", "child");
2025 child.task_list = Some(list.clone());
2026 child.set_task_list_version_meta("1");
2027 inner.save_session(&child).await.unwrap();
2028 let independent = SessionStoreV2::new(home.path().into()).await.unwrap();
2029 let commit_reached = Arc::new(tokio::sync::Barrier::new(2));
2030 let release_commit = Arc::new(tokio::sync::Barrier::new(2));
2031 let paused = LockedSessionStore::new(Arc::new(SingleCommitPauseStorage {
2032 inner: inner.clone(),
2033 commit_reached: commit_reached.clone(),
2034 release_commit: release_commit.clone(),
2035 }));
2036 let published = AtomicBool::new(false);
2037 let loser = async {
2038 if mode == 0 {
2039 paused
2040 .update_task_list_control_plane_and_publish(&root.id, &updated, "2", |_| {
2041 published.store(true, Ordering::SeqCst)
2042 })
2043 .await
2044 } else if mode == 1 {
2045 paused
2046 .update_task_list_control_plane_if_version_and_publish(
2047 &root.id,
2048 "1",
2049 &list,
2050 &updated,
2051 "2",
2052 |_| published.store(true, Ordering::SeqCst),
2053 )
2054 .await
2055 } else {
2056 paused
2057 .update_task_list_control_planes_if_version_and_publish(
2058 child_id,
2059 &root.id,
2060 "1",
2061 &list,
2062 &updated,
2063 "2",
2064 |_, _| published.store(true, Ordering::SeqCst),
2065 )
2066 .await
2067 }
2068 };
2069 let winner = async {
2070 commit_reached.wait().await;
2071 independent
2072 .mutate_supervisor_management(&SupervisorManagementRequest {
2073 supervisor: reference.clone(),
2074 expected_state_revision: 0,
2075 mutation: SupervisorManagementMutation::ConfigureProjectScope {
2076 allowed_projects: ["project-a".parse().unwrap()].into(),
2077 },
2078 })
2079 .await
2080 .unwrap();
2081 release_commit.wait().await;
2082 };
2083 let (result, ()) = tokio::join!(loser, winner);
2084 if mode == 0 {
2085 assert_eq!(result.unwrap_err().kind(), std::io::ErrorKind::WouldBlock);
2086 } else {
2087 assert!(!result.unwrap());
2088 }
2089 assert!(!published.load(Ordering::SeqCst));
2090 for id in [&root.id, &child.id] {
2091 assert_eq!(
2092 inner
2093 .load_runtime_control_plane(id)
2094 .await
2095 .unwrap()
2096 .unwrap()
2097 .task_list_version_meta()
2098 .as_deref(),
2099 Some("1")
2100 );
2101 }
2102 let fresh = LockedSessionStore::new(inner.clone());
2103 let check = |saved: &Session| {
2104 assert_eq!(saved.supervisor_management.as_ref().unwrap().revision, 1);
2105 published.store(true, Ordering::SeqCst);
2106 };
2107 let retried = if mode == 0 {
2108 fresh
2109 .update_task_list_control_plane_and_publish(&root.id, &updated, "2", check)
2110 .await
2111 } else if mode == 1 {
2112 fresh
2113 .update_task_list_control_plane_if_version_and_publish(
2114 &root.id, "1", &list, &updated, "2", check,
2115 )
2116 .await
2117 } else {
2118 fresh
2119 .update_task_list_control_planes_if_version_and_publish(
2120 child_id,
2121 &root.id,
2122 "1",
2123 &list,
2124 &updated,
2125 "2",
2126 |_, shared| check(shared),
2127 )
2128 .await
2129 };
2130 assert!(retried.unwrap());
2131 assert!(published.load(Ordering::SeqCst));
2132 assert_eq!(
2133 inner
2134 .load_runtime_control_plane(&root.id)
2135 .await
2136 .unwrap()
2137 .unwrap()
2138 .task_list_version_meta()
2139 .as_deref(),
2140 Some("2")
2141 );
2142 }
2143 }
2144
2145 #[async_trait::async_trait]
2146 impl Storage for SingleCommitBarrierStorage {
2147 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
2148 self.inner.save_session(session).await
2149 }
2150
2151 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
2152 self.inner.load_session(session_id).await
2153 }
2154
2155 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
2156 self.inner.delete_session(session_id).await
2157 }
2158
2159 async fn load_runtime_control_plane(
2160 &self,
2161 session_id: &str,
2162 ) -> std::io::Result<Option<Session>> {
2163 self.inner.load_runtime_control_plane(session_id).await
2164 }
2165
2166 async fn save_task_control_plane_if_matches(
2167 &self,
2168 original: &Session,
2169 updated: &Session,
2170 ) -> std::io::Result<bool> {
2171 self.before_commit.wait().await;
2172 self.inner
2173 .save_task_control_plane_if_matches(original, updated)
2174 .await
2175 }
2176 }
2177
2178 #[async_trait::async_trait]
2179 impl Storage for PairCommitBarrierStorage {
2180 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
2181 self.inner.save_session(session).await
2182 }
2183
2184 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
2185 self.inner.load_session(session_id).await
2186 }
2187
2188 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
2189 self.inner.delete_session(session_id).await
2190 }
2191
2192 async fn save_runtime_state(&self, session: &Session) -> std::io::Result<()> {
2193 self.inner.save_runtime_state(session).await
2194 }
2195
2196 async fn load_runtime_control_plane(
2197 &self,
2198 session_id: &str,
2199 ) -> std::io::Result<Option<Session>> {
2200 self.inner.load_runtime_control_plane(session_id).await
2201 }
2202
2203 async fn recover_task_control_plane_transaction(
2204 &self,
2205 first_session_id: &str,
2206 second_session_id: &str,
2207 ) -> std::io::Result<()> {
2208 self.inner
2209 .recover_task_control_plane_transaction(first_session_id, second_session_id)
2210 .await
2211 }
2212
2213 async fn save_task_control_planes_atomically(
2214 &self,
2215 first_original: &Session,
2216 first_updated: &Session,
2217 second_original: &Session,
2218 second_updated: &Session,
2219 ) -> std::io::Result<bool> {
2220 self.before_commit.wait().await;
2225 self.inner
2226 .save_task_control_planes_atomically(
2227 first_original,
2228 first_updated,
2229 second_original,
2230 second_updated,
2231 )
2232 .await
2233 }
2234 }
2235
2236 #[async_trait::async_trait]
2237 impl Storage for CountingControlPlaneStorage {
2238 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
2239 self.full_saves.fetch_add(1, Ordering::SeqCst);
2240 self.inner.save_session(session).await
2241 }
2242
2243 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
2244 self.inner.load_session(session_id).await
2245 }
2246
2247 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
2248 self.inner.delete_session(session_id).await
2249 }
2250
2251 async fn save_runtime_state(&self, session: &Session) -> std::io::Result<()> {
2252 self.runtime_state_saves.fetch_add(1, Ordering::SeqCst);
2253 self.inner.save_runtime_state(session).await
2254 }
2255
2256 async fn load_runtime_control_plane(
2257 &self,
2258 session_id: &str,
2259 ) -> std::io::Result<Option<Session>> {
2260 self.control_plane_loads.fetch_add(1, Ordering::SeqCst);
2261 self.inner.load_runtime_control_plane(session_id).await
2262 }
2263
2264 async fn recover_task_control_plane_transaction(
2265 &self,
2266 first_session_id: &str,
2267 second_session_id: &str,
2268 ) -> std::io::Result<()> {
2269 self.inner
2270 .recover_task_control_plane_transaction(first_session_id, second_session_id)
2271 .await
2272 }
2273
2274 async fn save_task_control_plane_if_matches(
2275 &self,
2276 original: &Session,
2277 updated: &Session,
2278 ) -> std::io::Result<bool> {
2279 let committed = self
2280 .inner
2281 .save_task_control_plane_if_matches(original, updated)
2282 .await?;
2283 if committed {
2284 self.runtime_state_saves.fetch_add(1, Ordering::SeqCst);
2285 }
2286 Ok(committed)
2287 }
2288
2289 async fn save_task_control_planes_atomically(
2290 &self,
2291 first_original: &Session,
2292 first_updated: &Session,
2293 second_original: &Session,
2294 second_updated: &Session,
2295 ) -> std::io::Result<bool> {
2296 let committed = self
2297 .inner
2298 .save_task_control_planes_atomically(
2299 first_original,
2300 first_updated,
2301 second_original,
2302 second_updated,
2303 )
2304 .await?;
2305 if committed {
2306 self.runtime_state_saves.fetch_add(2, Ordering::SeqCst);
2307 }
2308 Ok(committed)
2309 }
2310 }
2311
2312 async fn make_storage() -> (tempfile::TempDir, Arc<dyn Storage>) {
2313 let temp = tempfile::tempdir().unwrap();
2314 let storage = SessionStoreV2::new(temp.path().to_path_buf())
2315 .await
2316 .expect("storage init");
2317 (temp, Arc::new(storage) as Arc<dyn Storage>)
2318 }
2319
2320 fn fresh(id: &str) -> Session {
2321 Session::new(id.to_string(), "test-model".to_string())
2322 }
2323
2324 fn typed_permission_result(
2325 tool_call_id: &str,
2326 message_id: &str,
2327 generation: &str,
2328 content: &str,
2329 ) -> bamboo_domain::session::types::Message {
2330 let mut message =
2331 bamboo_domain::session::types::Message::tool_result(tool_call_id, content);
2332 message.id = message_id.to_string();
2333 message.metadata = Some(serde_json::json!({
2334 "permission_request": {
2335 "request_generation": generation,
2336 }
2337 }));
2338 message
2339 }
2340
2341 fn ledger_state(state_revision: u64, marker: &str) -> bamboo_domain::ModelContextState {
2342 bamboo_domain::ModelContextState {
2343 state_revision,
2344 prefix_epoch: state_revision,
2345 cache_scope_sha256: Some("scope".to_string()),
2346 transcript_item_sha256: vec![marker.to_string()],
2347 ..bamboo_domain::ModelContextState::default()
2348 }
2349 }
2350
2351 fn set_permission_audit(
2352 session: &mut Session,
2353 requested: bamboo_domain::SessionPermissionMode,
2354 policy_revision: u64,
2355 mapping: &str,
2356 transitioned_at: &str,
2357 ) -> u64 {
2358 let resolution = bamboo_domain::resolve_permission_mode(requested, PermissionMode::Default);
2359 session
2360 .agent_runtime_state
2361 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
2362 .set_permission_mode(requested);
2363 bamboo_domain::record_permission_audit(
2364 &mut session.metadata,
2365 &PermissionAuditSeed::new(policy_revision, resolution, mapping),
2366 Some(transitioned_at),
2367 )
2368 .unwrap()
2369 }
2370
2371 #[tokio::test]
2372 async fn checked_runtime_mutation_persists_a_cache_only_pending_session() {
2373 let (_temp, storage) = make_storage().await;
2374 let store = LockedSessionStore::new(storage.clone());
2375 let mut cached = fresh("cache-only-response");
2376 cached.set_pending_question(
2377 "tool-1".to_string(),
2378 "ConclusionWithOptions".to_string(),
2379 "Choose".to_string(),
2380 vec!["A".to_string()],
2381 false,
2382 );
2383
2384 let saved = store
2385 .mutate_runtime_session_and_publish(
2386 &cached.id.clone(),
2387 move || Some(cached),
2388 |session| {
2389 assert!(session.pending_question.is_some());
2390 session.clear_pending_question();
2391 Ok::<_, ()>(())
2392 },
2393 |_| {},
2394 )
2395 .await
2396 .unwrap()
2397 .unwrap()
2398 .expect("cache-only session should be created durably");
2399
2400 assert!(saved.pending_question.is_none());
2401 assert!(storage
2402 .load_session("cache-only-response")
2403 .await
2404 .unwrap()
2405 .unwrap()
2406 .pending_question
2407 .is_none());
2408 }
2409
2410 #[tokio::test]
2411 async fn checked_runtime_mutation_preserves_durable_authorities_for_newer_cache() {
2412 let (_temp, storage) = make_storage().await;
2413 let store = LockedSessionStore::new(storage.clone());
2414 let session_id = "cached-response-authorities";
2415 let mut durable = fresh(session_id);
2416 durable.title = "Durable title".to_string();
2417 durable.title_version = 4;
2418 durable.title_generated = true;
2419 durable.metadata_version = 9;
2420 set_permission_audit(
2421 &mut durable,
2422 bamboo_domain::SessionPermissionMode::Auto,
2423 7,
2424 "bamboo_runtime:durable-auto",
2425 "2026-08-10T09:00:00Z",
2426 );
2427 storage.save_session(&durable).await.unwrap();
2428
2429 let mut cached = fresh(session_id);
2430 cached.created_at = durable.created_at;
2431 cached.title = "Stale cached title".to_string();
2432 cached.updated_at = durable.updated_at + chrono::Duration::seconds(1);
2433 cached.set_pending_question(
2434 "tool-1".to_string(),
2435 "ConclusionWithOptions".to_string(),
2436 "Choose".to_string(),
2437 vec!["A".to_string()],
2438 false,
2439 );
2440
2441 store
2442 .mutate_runtime_session_and_publish(
2443 session_id,
2444 move || Some(cached),
2445 |session| {
2446 session.clear_pending_question();
2447 Ok::<_, ()>(())
2448 },
2449 |_| {},
2450 )
2451 .await
2452 .unwrap()
2453 .unwrap()
2454 .expect("session should exist");
2455
2456 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2457 assert_eq!(saved.title, "Durable title");
2458 assert_eq!(saved.title_version, 4);
2459 assert_eq!(saved.metadata_version, 9);
2460 assert_eq!(
2461 saved
2462 .agent_runtime_state
2463 .as_ref()
2464 .unwrap()
2465 .effective_permission_mode(),
2466 bamboo_domain::SessionPermissionMode::Auto
2467 );
2468 let audit = bamboo_domain::PermissionAuditSnapshot::from_metadata(&saved.metadata).unwrap();
2469 assert_eq!(audit.policy_revision, 7);
2470 assert_eq!(audit.executor_mapping, "bamboo_runtime:durable-auto");
2471 }
2472
2473 #[tokio::test]
2474 async fn checked_runtime_mutation_cannot_resurrect_consumed_ask_from_newer_cache() {
2475 use bamboo_domain::session::types::Message;
2476
2477 let (_temp, storage) = make_storage().await;
2478 let store = LockedSessionStore::new(storage.clone());
2479 let session_id = "cached-consumed-response";
2480 let mut stale_cached = fresh(session_id);
2481 stale_cached.add_message(Message::tool_result("call-1", "waiting"));
2482 stale_cached.set_pending_question(
2483 "call-1".to_string(),
2484 "ConclusionWithOptions".to_string(),
2485 "Choose".to_string(),
2486 vec!["A".to_string()],
2487 false,
2488 );
2489
2490 let mut durable = stale_cached.clone();
2491 durable.clear_pending_question();
2492 durable.metadata.insert(
2493 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
2494 r#"["call-1"]"#.to_string(),
2495 );
2496 durable.messages[0].content = "Selected response: A".to_string();
2497 durable.add_message(Message::user("durable concurrent message"));
2498 storage.save_session(&durable).await.unwrap();
2499
2500 stale_cached.updated_at = durable.updated_at + chrono::Duration::seconds(1);
2503 let saved = store
2504 .mutate_runtime_session_and_publish(
2505 session_id,
2506 move || Some(stale_cached),
2507 |session| {
2508 assert!(session.pending_question.is_none());
2509 Ok::<_, ()>(())
2510 },
2511 |_| {},
2512 )
2513 .await
2514 .unwrap()
2515 .unwrap()
2516 .unwrap();
2517
2518 assert!(saved.pending_question.is_none());
2519 assert_eq!(saved.messages.len(), 2);
2520 assert_eq!(saved.messages[0].content, "Selected response: A");
2521 assert_eq!(saved.messages[1].content, "durable concurrent message");
2522 }
2523
2524 #[tokio::test]
2525 async fn response_inspection_adopts_durable_consumption_without_writing() {
2526 let (_temp, storage) = make_storage().await;
2527 let store = LockedSessionStore::new(storage.clone());
2528 let session_id = "inspect-consumed-response";
2529 let mut stale_cached = fresh(session_id);
2530 stale_cached.add_message(bamboo_domain::session::types::Message::tool_result(
2535 "call-1", "waiting",
2536 ));
2537 stale_cached.set_pending_question(
2538 "call-1".to_string(),
2539 "ConclusionWithOptions".to_string(),
2540 "Choose".to_string(),
2541 vec!["A".to_string()],
2542 false,
2543 );
2544
2545 let mut durable = stale_cached.clone();
2546 durable.clear_pending_question();
2547 durable.metadata.insert(
2548 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
2549 r#"["call-1"]"#.to_string(),
2550 );
2551 storage.save_session(&durable).await.unwrap();
2552 stale_cached.updated_at = durable.updated_at + chrono::Duration::seconds(1);
2553
2554 let inspected = store
2555 .inspect_runtime_session_for_response(session_id, move || Some(stale_cached))
2556 .await
2557 .unwrap()
2558 .expect("session should be inspectable");
2559 assert!(inspected.pending_question.is_none());
2560
2561 let unchanged = storage.load_session(session_id).await.unwrap().unwrap();
2562 assert_eq!(unchanged.updated_at, durable.updated_at);
2563 assert!(unchanged.pending_question.is_none());
2564 }
2565
2566 #[tokio::test]
2567 async fn same_mode_newer_run_start_audit_survives_every_runtime_save_path() {
2568 for path in ["merge", "checkpoint", "control-plane"] {
2569 let (_temp, storage) = make_storage().await;
2570 let store = LockedSessionStore::new(storage.clone());
2571 let session_id = format!("same-mode-newer-{path}");
2572 let mut durable = fresh(&session_id);
2573 set_permission_audit(
2574 &mut durable,
2575 bamboo_domain::SessionPermissionMode::Default,
2576 1,
2577 "bamboo_runtime:old-policy",
2578 "2026-07-31T12:00:00Z",
2579 );
2580 storage.save_session(&durable).await.unwrap();
2581
2582 let mut run_start = durable.clone();
2583 let old_revision = PermissionAuditSnapshot::from_metadata(&durable.metadata)
2584 .unwrap()
2585 .audit_revision;
2586 let new_revision = set_permission_audit(
2587 &mut run_start,
2588 bamboo_domain::SessionPermissionMode::Default,
2589 2,
2590 "bamboo_runtime:new-policy",
2591 "2026-07-31T12:00:00Z",
2592 );
2593 assert!(new_revision > old_revision);
2594
2595 match path {
2596 "merge" => store.merge_save_runtime(&mut run_start).await.unwrap(),
2597 "checkpoint" => store
2598 .checkpoint_runtime_session(&mut run_start)
2599 .await
2600 .unwrap(),
2601 "control-plane" => store.save_runtime_only(&mut run_start).await.unwrap(),
2602 _ => unreachable!(),
2603 }
2604
2605 let saved = storage.load_session(&session_id).await.unwrap().unwrap();
2606 let audit = PermissionAuditSnapshot::from_metadata(&saved.metadata).unwrap();
2607 assert_eq!(audit.audit_revision, new_revision, "path={path}");
2608 assert_eq!(audit.policy_revision, 2, "path={path}");
2609 assert_eq!(audit.executor_mapping, "bamboo_runtime:new-policy");
2610 }
2611 }
2612
2613 #[tokio::test]
2614 async fn newer_disk_transition_wins_after_mode_cycles_back() {
2615 let (_temp, storage) = make_storage().await;
2616 let store = LockedSessionStore::new(storage.clone());
2617 let session_id = "permission-cycle-back";
2618 let mut baseline = fresh(session_id);
2619 let stale_revision = set_permission_audit(
2620 &mut baseline,
2621 bamboo_domain::SessionPermissionMode::Default,
2622 1,
2623 "bamboo_runtime:initial",
2624 "2026-07-31T12:00:00Z",
2625 );
2626 storage.save_session(&baseline).await.unwrap();
2627 let mut stale_runtime = baseline.clone();
2628
2629 let mut durable = baseline;
2630 set_permission_audit(
2631 &mut durable,
2632 bamboo_domain::SessionPermissionMode::Auto,
2633 2,
2634 "bamboo_runtime:auto",
2635 "2026-07-31T12:01:00Z",
2636 );
2637 let durable_revision = set_permission_audit(
2638 &mut durable,
2639 bamboo_domain::SessionPermissionMode::Default,
2640 3,
2641 "bamboo_runtime:cycled-default",
2642 "2026-07-31T12:02:00Z",
2643 );
2644 assert!(durable_revision > stale_revision);
2645 storage.save_session(&durable).await.unwrap();
2646
2647 store.merge_save_runtime(&mut stale_runtime).await.unwrap();
2648 let saved = storage.load_session(session_id).await.unwrap().unwrap();
2649 let audit = PermissionAuditSnapshot::from_metadata(&saved.metadata).unwrap();
2650 assert_eq!(audit.audit_revision, durable_revision);
2651 assert_eq!(audit.policy_revision, 3);
2652 assert_eq!(audit.executor_mapping, "bamboo_runtime:cycled-default");
2653 }
2654
2655 #[tokio::test]
2656 async fn authoritative_activation_seed_replaces_every_warm_worker_posture() {
2657 let (_temp, storage) = make_storage().await;
2658 let store = LockedSessionStore::new(storage.clone());
2659 let session_id = "warm-permission-matrix";
2660 let cases = [
2661 (
2662 bamboo_domain::SessionPermissionMode::Auto,
2663 PermissionMode::Default,
2664 PermissionMode::Auto,
2665 ),
2666 (
2667 bamboo_domain::SessionPermissionMode::Default,
2668 PermissionMode::Default,
2669 PermissionMode::Default,
2670 ),
2671 (
2672 bamboo_domain::SessionPermissionMode::Auto,
2673 PermissionMode::Default,
2674 PermissionMode::Auto,
2675 ),
2676 (
2677 bamboo_domain::SessionPermissionMode::Bypass,
2678 PermissionMode::Auto,
2679 PermissionMode::BypassPermissions,
2680 ),
2681 ];
2682 let mut previous_revision = 0;
2683 let created_at = fresh(session_id).created_at;
2684
2685 for (index, (requested, configured, expected_effective)) in cases.into_iter().enumerate() {
2686 let mut activation = fresh(session_id);
2687 activation.created_at = created_at;
2688 activation
2689 .agent_runtime_state
2690 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
2691 .set_permission_mode(requested);
2692 let resolution = bamboo_domain::resolve_permission_mode(requested, configured);
2693 bamboo_domain::record_permission_audit(
2694 &mut activation.metadata,
2695 &PermissionAuditSeed::new(
2696 index as u64 + 1,
2697 resolution,
2698 format!("bamboo_worker:{}", resolution.effective.as_str()),
2699 ),
2700 Some("2026-07-31T12:00:00Z"),
2701 )
2702 .unwrap();
2703
2704 RuntimeSessionPersistence::seed_runtime_activation(&store, &mut activation)
2705 .await
2706 .unwrap();
2707 let durable = storage.load_session(session_id).await.unwrap().unwrap();
2708 assert_eq!(
2709 durable
2710 .agent_runtime_state
2711 .as_ref()
2712 .unwrap()
2713 .effective_permission_mode(),
2714 requested,
2715 "activation {index}"
2716 );
2717 let audit = PermissionAuditSnapshot::from_metadata(&durable.metadata).unwrap();
2718 assert_eq!(audit.resolution.requested, requested);
2719 assert_eq!(audit.resolution.effective, expected_effective);
2720 assert!(audit.audit_revision > previous_revision);
2721 previous_revision = audit.audit_revision;
2722 }
2723 }
2724
2725 #[tokio::test]
2726 async fn resident_reseed_bumps_etag_only_for_typed_transition() {
2727 let (_temp, storage) = make_storage().await;
2728 let store = LockedSessionStore::new(storage.clone());
2729 let session_id = "resident-atomic-permission";
2730 let mut baseline = fresh(session_id);
2731 baseline.metadata_version = 7;
2732 set_permission_audit(
2733 &mut baseline,
2734 bamboo_domain::SessionPermissionMode::Auto,
2735 1,
2736 "bamboo_runtime:auto",
2737 "2026-07-31T12:00:00Z",
2738 );
2739 storage.save_session(&baseline).await.unwrap();
2740 let initial_audit = PermissionAuditSnapshot::from_metadata(&baseline.metadata).unwrap();
2741
2742 let same_mode_seed = PermissionAuditSeed::bamboo_runtime(
2743 2,
2744 bamboo_domain::resolve_permission_mode(
2745 bamboo_domain::SessionPermissionMode::Auto,
2746 PermissionMode::Default,
2747 ),
2748 );
2749 let refreshed = store
2750 .update_authoritative_permission_posture_and_publish(
2751 session_id,
2752 &same_mode_seed,
2753 |session| {
2754 session
2755 .metadata
2756 .insert("resident.marker".to_string(), "same-mode".to_string());
2757 },
2758 |_| {},
2759 )
2760 .await
2761 .unwrap()
2762 .unwrap();
2763 let refreshed_audit = PermissionAuditSnapshot::from_metadata(&refreshed.metadata).unwrap();
2764 assert_eq!(refreshed.metadata_version, 7);
2765 assert!(refreshed_audit.audit_revision > initial_audit.audit_revision);
2766 assert_eq!(refreshed_audit.policy_revision, 2);
2767
2768 let transition_seed = PermissionAuditSeed::bamboo_runtime(
2769 3,
2770 bamboo_domain::resolve_permission_mode(
2771 bamboo_domain::SessionPermissionMode::Default,
2772 PermissionMode::Default,
2773 ),
2774 );
2775 let transitioned = store
2776 .update_authoritative_permission_posture_and_publish(
2777 session_id,
2778 &transition_seed,
2779 |session| {
2780 session
2781 .metadata
2782 .insert("resident.marker".to_string(), "transition".to_string());
2783 },
2784 |_| {},
2785 )
2786 .await
2787 .unwrap()
2788 .unwrap();
2789 let transitioned_audit =
2790 PermissionAuditSnapshot::from_metadata(&transitioned.metadata).unwrap();
2791 assert_eq!(transitioned.metadata_version, 8, "old ETag must be invalid");
2792 assert_eq!(
2793 transitioned
2794 .agent_runtime_state
2795 .as_ref()
2796 .unwrap()
2797 .effective_permission_mode(),
2798 bamboo_domain::SessionPermissionMode::Default
2799 );
2800 assert_eq!(
2801 transitioned_audit.resolution.requested,
2802 bamboo_domain::SessionPermissionMode::Default
2803 );
2804 assert!(transitioned_audit.audit_revision > refreshed_audit.audit_revision);
2805 assert_eq!(
2806 transitioned
2807 .metadata
2808 .get("resident.marker")
2809 .map(String::as_str),
2810 Some("transition")
2811 );
2812 }
2813
2814 #[tokio::test]
2815 async fn worker_activation_cas_cannot_overwrite_concurrent_permission_patch() {
2816 let (_temp, storage) = make_storage().await;
2817 let store = LockedSessionStore::new(storage.clone());
2818 let session_id = "permission-activation-cas";
2819 let mut baseline = fresh(session_id);
2820 set_permission_audit(
2821 &mut baseline,
2822 bamboo_domain::SessionPermissionMode::Default,
2823 1,
2824 "bamboo_runtime:default",
2825 "2026-07-31T12:00:00Z",
2826 );
2827 storage.save_session(&baseline).await.unwrap();
2828 let dispatched_revision = PermissionAuditSnapshot::from_metadata(&baseline.metadata)
2829 .unwrap()
2830 .audit_revision;
2831
2832 let patched_resolution = bamboo_domain::resolve_permission_mode(
2833 bamboo_domain::SessionPermissionMode::Auto,
2834 PermissionMode::Default,
2835 );
2836 let patched = store
2837 .update_authoritative_permission_posture_and_publish(
2838 session_id,
2839 &PermissionAuditSeed::new(2, patched_resolution, "patch:auto"),
2840 |_| {},
2841 |_| {},
2842 )
2843 .await
2844 .unwrap()
2845 .unwrap();
2846 let patched_audit = PermissionAuditSnapshot::from_metadata(&patched.metadata).unwrap();
2847 assert!(patched_audit.audit_revision > dispatched_revision);
2848
2849 let stale_worker_seed = PermissionAuditSeed::new(
2850 1,
2851 bamboo_domain::resolve_permission_mode(
2852 bamboo_domain::SessionPermissionMode::Default,
2853 PermissionMode::Default,
2854 ),
2855 "worker:stale-default",
2856 );
2857 let error = store
2858 .record_permission_posture_activation_and_publish(
2859 session_id,
2860 Some(dispatched_revision),
2861 &stale_worker_seed,
2862 |_| {},
2863 )
2864 .await
2865 .unwrap_err();
2866 assert!(error.to_string().contains("durable audit changed"));
2867
2868 let durable = storage.load_session(session_id).await.unwrap().unwrap();
2869 let durable_audit = PermissionAuditSnapshot::from_metadata(&durable.metadata).unwrap();
2870 assert_eq!(durable_audit, patched_audit);
2871 assert_eq!(durable_audit.executor_mapping, "patch:auto");
2872 }
2873
2874 #[tokio::test]
2877 async fn update_runtime_config_preserves_concurrently_appended_messages() {
2878 use bamboo_domain::session::types::Message;
2879 use bamboo_domain::ReasoningEffort;
2880
2881 let (_temp, storage) = make_storage().await;
2882 let store = LockedSessionStore::new(storage.clone());
2883 let session_id = "cfg-preserve";
2884
2885 let mut initial = fresh(session_id);
2887 initial.add_message(Message::user("hello"));
2888 initial.add_message(Message::assistant("hi", None));
2889 storage.save_session(&initial).await.unwrap();
2890
2891 let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
2893 after_chat.add_message(Message::user("second question"));
2894 storage.save_session(&after_chat).await.unwrap();
2895 assert_eq!(after_chat.messages.len(), 3);
2896
2897 let updated = store
2901 .update_runtime_config(session_id, |s| {
2902 s.reasoning_effort = Some(ReasoningEffort::Max);
2903 })
2904 .await
2905 .unwrap()
2906 .expect("session exists");
2907
2908 assert_eq!(updated.reasoning_effort, Some(ReasoningEffort::Max));
2909 assert_eq!(
2910 updated.messages.len(),
2911 3,
2912 "config patch must not revert a concurrently-appended message"
2913 );
2914
2915 let on_disk = storage.load_session(session_id).await.unwrap().unwrap();
2916 assert_eq!(on_disk.messages.len(), 3);
2917 assert_eq!(on_disk.reasoning_effort, Some(ReasoningEffort::Max));
2918 }
2919
2920 #[tokio::test]
2921 async fn update_runtime_config_returns_none_for_missing_session() {
2922 use bamboo_domain::ReasoningEffort;
2923
2924 let (_temp, storage) = make_storage().await;
2925 let store = LockedSessionStore::new(storage);
2926 let result = store
2927 .update_runtime_config("does-not-exist", |s| {
2928 s.reasoning_effort = Some(ReasoningEffort::Low);
2929 })
2930 .await
2931 .unwrap();
2932 assert!(result.is_none());
2933 }
2934
2935 #[tokio::test]
2936 async fn merge_save_runtime_overwrites_messages_from_stale_snapshot() {
2937 use bamboo_domain::session::types::Message;
2942
2943 let (_temp, storage) = make_storage().await;
2944 let store = LockedSessionStore::new(storage.clone());
2945 let session_id = "stale-clobber";
2946
2947 let mut baseline = fresh(session_id);
2949 baseline.add_message(Message::user("hello"));
2950 storage.save_session(&baseline).await.unwrap();
2951 let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
2952
2953 let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
2955 after_chat.add_message(Message::user("second"));
2956 storage.save_session(&after_chat).await.unwrap();
2957 assert_eq!(
2958 storage
2959 .load_session(session_id)
2960 .await
2961 .unwrap()
2962 .unwrap()
2963 .messages
2964 .len(),
2965 2
2966 );
2967
2968 store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
2970 let after = storage.load_session(session_id).await.unwrap().unwrap();
2971 assert_eq!(
2972 after.messages.len(),
2973 1,
2974 "merge_save_runtime clobbers concurrent appends — this is why config patches must use update_runtime_config"
2975 );
2976 }
2977
2978 #[tokio::test]
2979 async fn stale_runtime_save_cannot_remove_admitted_inbox_transcript() {
2980 use bamboo_domain::session::types::Message;
2981 use bamboo_domain::SessionMessageId;
2982
2983 let (_temp, storage) = make_storage().await;
2984 let store = LockedSessionStore::new(storage.clone());
2985 let session_id = "stale-inbox-preserve";
2986
2987 let mut baseline = fresh(session_id);
2988 let mut base = Message::user("base");
2989 base.id = "base".to_string();
2990 baseline.add_message(base);
2991 storage.save_session(&baseline).await.unwrap();
2992 let mut stale = baseline.clone();
2993 let mut later_assistant = Message::assistant("runner output", None);
2994 later_assistant.id = "later-assistant".to_string();
2995 stale.add_message(later_assistant);
2996
2997 let mut durable = baseline;
2998 let inbox_id = SessionMessageId::parse("durable-inbox-id").unwrap();
2999 let mut admitted = Message::user("durable inbox message");
3000 admitted.id = inbox_id.as_str().to_string();
3001 durable.add_message(admitted);
3002 durable
3003 .session_inbox_admission_mut()
3004 .record(inbox_id.clone(), 7);
3005 storage.save_session(&durable).await.unwrap();
3006
3007 store.merge_save_runtime(&mut stale).await.unwrap();
3008 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3009 let ids = saved
3010 .messages
3011 .iter()
3012 .map(|message| message.id.as_str())
3013 .collect::<Vec<_>>();
3014 assert_eq!(ids, vec!["base", "durable-inbox-id", "later-assistant"]);
3015 assert_eq!(ids.iter().filter(|id| **id == inbox_id.as_str()).count(), 1);
3016 assert!(saved
3017 .session_inbox_admission()
3018 .is_some_and(|state| state.contains(&inbox_id)));
3019 }
3020
3021 #[tokio::test]
3022 async fn stale_runtime_save_preserves_typed_inbox_message_after_cursor_eviction() {
3023 use bamboo_domain::{
3024 SessionMessageEnvelope, SessionMessageId, SESSION_INBOX_ADMITTED_CAPACITY,
3025 };
3026
3027 let (_temp, storage) = make_storage().await;
3028 let store = LockedSessionStore::new(storage.clone());
3029 let session_id = "evicted-inbox-preserve";
3030 let mut durable = fresh(session_id);
3031 let mut envelope = SessionMessageEnvelope::user_input(session_id, "old durable inbox");
3032 envelope.id = SessionMessageId::parse("old-inbox-id").unwrap();
3033 durable.add_message(envelope.to_provider_message().unwrap());
3034 durable
3035 .session_inbox_admission_mut()
3036 .record(envelope.id.clone(), 1);
3037 for sequence in 2..=(SESSION_INBOX_ADMITTED_CAPACITY as u64 + 1) {
3038 durable.session_inbox_admission_mut().record(
3039 SessionMessageId::parse(format!("newer-{sequence}")).unwrap(),
3040 sequence,
3041 );
3042 }
3043 assert!(!durable
3044 .session_inbox_admission()
3045 .unwrap()
3046 .contains(&envelope.id));
3047 storage.save_session(&durable).await.unwrap();
3048
3049 let mut stale = fresh(session_id);
3050 stale.created_at = durable.created_at;
3051 store.merge_save_runtime(&mut stale).await.unwrap();
3052 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3053 assert_eq!(
3054 saved
3055 .messages
3056 .iter()
3057 .filter(|message| message.id == envelope.id.as_str())
3058 .count(),
3059 1
3060 );
3061 }
3062
3063 #[tokio::test]
3064 async fn runtime_final_save_cannot_resurrect_a_consumed_clarification() {
3065 use bamboo_domain::session::types::Message;
3066
3067 let (_temp, storage) = make_storage().await;
3068 let store = LockedSessionStore::new(storage.clone());
3069 let session_id = "consumed-clarification-final-save";
3070
3071 let mut suspended = fresh(session_id);
3072 suspended.add_message(Message::tool_result(
3073 "call-1",
3074 r#"{"status":"awaiting_clarification"}"#,
3075 ));
3076 suspended.set_pending_question(
3077 "call-1".to_string(),
3078 "ConclusionWithOptions".to_string(),
3079 "Choose".to_string(),
3080 vec!["A".to_string()],
3081 false,
3082 );
3083 suspended.metadata.insert(
3084 "runtime.suspend_reason".to_string(),
3085 "awaiting_clarification".to_string(),
3086 );
3087 storage.save_session(&suspended).await.unwrap();
3088 let mut stale_runner = suspended.clone();
3089
3090 let mut answered = suspended;
3091 answered.clear_pending_question();
3092 answered.metadata.remove("runtime.suspend_reason");
3093 answered.metadata.insert(
3094 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
3095 r#"["call-1"]"#.to_string(),
3096 );
3097 answered.metadata.insert(
3098 "clarification_resume_pending".to_string(),
3099 "true".to_string(),
3100 );
3101 answered.metadata.insert(
3102 "conclusion_with_options_resume_pending".to_string(),
3103 "true".to_string(),
3104 );
3105 answered.metadata.insert(
3106 "execute.startup_handoff_at".to_string(),
3107 "2026-08-10T09:00:00.000Z".to_string(),
3108 );
3109 let answer = answered
3110 .messages
3111 .iter_mut()
3112 .find(|message| message.tool_call_id.as_deref() == Some("call-1"))
3113 .unwrap();
3114 answer.content = "Selected response: A".to_string();
3115 storage.save_session(&answered).await.unwrap();
3116
3117 store.merge_save_runtime(&mut stale_runner).await.unwrap();
3118
3119 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3120 assert!(saved.pending_question.is_none());
3121 assert!(!saved.metadata.contains_key("runtime.suspend_reason"));
3122 assert_eq!(
3123 saved
3124 .metadata
3125 .get("clarification_resume_pending")
3126 .map(String::as_str),
3127 Some("true")
3128 );
3129 assert_eq!(
3130 saved
3131 .metadata
3132 .get("execute.startup_handoff_at")
3133 .map(String::as_str),
3134 Some("2026-08-10T09:00:00.000Z")
3135 );
3136 let answers = saved
3137 .messages
3138 .iter()
3139 .filter(|message| message.tool_call_id.as_deref() == Some("call-1"))
3140 .collect::<Vec<_>>();
3141 assert_eq!(answers.len(), 1);
3142 assert_eq!(answers[0].content, "Selected response: A");
3143 }
3144
3145 #[tokio::test]
3146 async fn runtime_checkpoint_cannot_resurrect_a_consumed_clarification() {
3147 use bamboo_domain::session::types::Message;
3148
3149 let (_temp, storage) = make_storage().await;
3150 let store = LockedSessionStore::new(storage.clone());
3151 let session_id = "consumed-clarification-checkpoint";
3152 let mut stale_runner = fresh(session_id);
3153 stale_runner.add_message(Message::tool_result("call-1", "waiting"));
3154 stale_runner.set_pending_question(
3155 "call-1".to_string(),
3156 "ConclusionWithOptions".to_string(),
3157 "Choose".to_string(),
3158 vec!["A".to_string()],
3159 false,
3160 );
3161 stale_runner.metadata.insert(
3162 "runtime.suspend_reason".to_string(),
3163 "awaiting_clarification".to_string(),
3164 );
3165
3166 let mut answered = stale_runner.clone();
3167 answered.clear_pending_question();
3168 answered.metadata.remove("runtime.suspend_reason");
3169 answered.metadata.insert(
3170 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
3171 r#"["call-1"]"#.to_string(),
3172 );
3173 answered.metadata.insert(
3174 "clarification_resume_pending".to_string(),
3175 "true".to_string(),
3176 );
3177 answered.messages[0].content = "Selected response: A".to_string();
3178 storage.save_session(&answered).await.unwrap();
3179
3180 store
3181 .checkpoint_runtime_session(&mut stale_runner)
3182 .await
3183 .unwrap();
3184
3185 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3186 assert!(saved.pending_question.is_none());
3187 assert!(!saved.metadata.contains_key("runtime.suspend_reason"));
3188 assert_eq!(saved.messages.len(), 1);
3189 assert_eq!(saved.messages[0].content, "Selected response: A");
3190 }
3191
3192 #[tokio::test]
3193 async fn runtime_final_save_does_not_consume_a_new_reused_permission_occurrence() {
3194 let (_temp, storage) = make_storage().await;
3195 let store = LockedSessionStore::new(storage.clone());
3196 let session_id = "reused-permission-final-save";
3197
3198 let mut durable = fresh(session_id);
3199 durable.add_message(typed_permission_result(
3200 "reused-call",
3201 "old-result",
3202 "generation-old",
3203 "Selected response: Approve",
3204 ));
3205 durable.metadata.insert(
3206 CONSUMED_RESPONSE_OCCURRENCES_KEY.to_string(),
3207 serde_json::to_string(&vec![ResponseOccurrence {
3208 tool_call_id: "reused-call".to_string(),
3209 tool_result_message_id: "old-result".to_string(),
3210 permission_generation: Some("generation-old".to_string()),
3211 }])
3212 .unwrap(),
3213 );
3214 storage.save_session(&durable).await.unwrap();
3215
3216 let mut new_runner = durable;
3217 new_runner.add_message(typed_permission_result(
3218 "reused-call",
3219 "new-result",
3220 "generation-new",
3221 "waiting for the new decision",
3222 ));
3223 new_runner.set_pending_question(
3224 "reused-call".to_string(),
3225 "Permission".to_string(),
3226 "Approve new operation?".to_string(),
3227 vec!["Approve".to_string(), "Deny".to_string()],
3228 false,
3229 );
3230 new_runner.metadata.insert(
3231 "runtime.suspend_reason".to_string(),
3232 "awaiting_permission_approval".to_string(),
3233 );
3234
3235 store.merge_save_runtime(&mut new_runner).await.unwrap();
3236
3237 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3238 assert_eq!(
3239 saved
3240 .pending_question
3241 .as_ref()
3242 .map(|pending| pending.tool_call_id.as_str()),
3243 Some("reused-call")
3244 );
3245 assert_eq!(
3246 saved.messages.last().map(|message| message.id.as_str()),
3247 Some("new-result")
3248 );
3249 assert_eq!(
3250 saved.messages.last().unwrap().content,
3251 "waiting for the new decision"
3252 );
3253 assert_eq!(
3254 latest_response_occurrence(&saved, "reused-call")
3255 .and_then(|occurrence| occurrence.permission_generation),
3256 Some("generation-new".to_string())
3257 );
3258 }
3259
3260 #[tokio::test]
3261 async fn legacy_consumed_id_does_not_consume_a_new_reused_occurrence_after_upgrade() {
3262 let (_temp, storage) = make_storage().await;
3263 let store = LockedSessionStore::new(storage.clone());
3264 let session_id = "legacy-reused-permission-final-save";
3265
3266 let mut durable = fresh(session_id);
3267 durable.add_message(typed_permission_result(
3268 "reused-call",
3269 "old-result",
3270 "generation-old",
3271 "Selected response: Approve",
3272 ));
3273 durable.metadata.insert(
3274 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
3275 r#"["reused-call"]"#.to_string(),
3276 );
3277 storage.save_session(&durable).await.unwrap();
3278
3279 let mut new_runner = durable;
3280 new_runner.add_message(typed_permission_result(
3281 "reused-call",
3282 "new-result",
3283 "generation-new",
3284 "waiting for the new decision",
3285 ));
3286 new_runner.set_pending_question(
3287 "reused-call".to_string(),
3288 "Permission".to_string(),
3289 "Approve new operation?".to_string(),
3290 vec!["Approve".to_string(), "Deny".to_string()],
3291 false,
3292 );
3293
3294 store.merge_save_runtime(&mut new_runner).await.unwrap();
3295
3296 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3297 assert_eq!(
3298 saved
3299 .pending_question
3300 .as_ref()
3301 .map(|pending| pending.tool_call_id.as_str()),
3302 Some("reused-call")
3303 );
3304 assert_eq!(
3305 saved.messages.last().map(|message| message.id.as_str()),
3306 Some("new-result")
3307 );
3308 assert_eq!(
3309 latest_response_occurrence(&saved, "reused-call")
3310 .and_then(|occurrence| occurrence.permission_generation),
3311 Some("generation-new".to_string())
3312 );
3313 }
3314
3315 #[tokio::test]
3316 async fn runtime_checkpoint_does_not_consume_a_new_reused_permission_occurrence() {
3317 let (_temp, storage) = make_storage().await;
3318 let store = LockedSessionStore::new(storage.clone());
3319 let session_id = "reused-permission-checkpoint";
3320
3321 let mut durable = fresh(session_id);
3322 durable.add_message(typed_permission_result(
3323 "reused-call",
3324 "old-result",
3325 "generation-old",
3326 "Selected response: Deny",
3327 ));
3328 durable.metadata.insert(
3329 CONSUMED_RESPONSE_OCCURRENCES_KEY.to_string(),
3330 serde_json::to_string(&vec![ResponseOccurrence {
3331 tool_call_id: "reused-call".to_string(),
3332 tool_result_message_id: "old-result".to_string(),
3333 permission_generation: Some("generation-old".to_string()),
3334 }])
3335 .unwrap(),
3336 );
3337 storage.save_session(&durable).await.unwrap();
3338
3339 let mut new_runner = durable;
3340 new_runner.add_message(typed_permission_result(
3341 "reused-call",
3342 "new-result",
3343 "generation-new",
3344 "waiting for the new decision",
3345 ));
3346 new_runner.set_pending_question(
3347 "reused-call".to_string(),
3348 "Permission".to_string(),
3349 "Approve new operation?".to_string(),
3350 vec!["Approve".to_string(), "Deny".to_string()],
3351 false,
3352 );
3353
3354 store
3355 .checkpoint_runtime_session(&mut new_runner)
3356 .await
3357 .unwrap();
3358
3359 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3360 assert_eq!(
3361 saved
3362 .pending_question
3363 .as_ref()
3364 .map(|pending| pending.tool_call_id.as_str()),
3365 Some("reused-call")
3366 );
3367 assert_eq!(
3368 saved.messages.last().map(|message| message.id.as_str()),
3369 Some("new-result")
3370 );
3371 assert_eq!(
3372 latest_response_occurrence(&saved, "reused-call")
3373 .and_then(|occurrence| occurrence.permission_generation),
3374 Some("generation-new".to_string())
3375 );
3376 }
3377
3378 #[tokio::test]
3379 async fn consumed_permission_adoption_keeps_reexecute_id_and_generation_paired() {
3380 let (_temp, storage) = make_storage().await;
3381 let store = LockedSessionStore::new(storage.clone());
3382 let session_id = "consumed-permission-control-pair";
3383
3384 let mut stale_runner = fresh(session_id);
3385 stale_runner.add_message(typed_permission_result(
3386 "call-1",
3387 "result-1",
3388 "generation-1",
3389 "waiting",
3390 ));
3391 stale_runner.set_pending_question(
3392 "call-1".to_string(),
3393 "Permission".to_string(),
3394 "Approve?".to_string(),
3395 vec!["Approve".to_string(), "Deny".to_string()],
3396 false,
3397 );
3398
3399 let mut answered = stale_runner.clone();
3400 answered.clear_pending_question();
3401 answered.messages[0].content = "Selected response: Approve".to_string();
3402 answered.metadata.insert(
3403 CONSUMED_RESPONSE_OCCURRENCES_KEY.to_string(),
3404 serde_json::to_string(&vec![ResponseOccurrence {
3405 tool_call_id: "call-1".to_string(),
3406 tool_result_message_id: "result-1".to_string(),
3407 permission_generation: Some("generation-1".to_string()),
3408 }])
3409 .unwrap(),
3410 );
3411 answered.metadata.insert(
3412 "permission.reexecute_tool_call_id".to_string(),
3413 "call-1".to_string(),
3414 );
3415 answered.metadata.insert(
3416 "permission.reexecute_request_generation".to_string(),
3417 "generation-1".to_string(),
3418 );
3419 storage.save_session(&answered).await.unwrap();
3420
3421 store.merge_save_runtime(&mut stale_runner).await.unwrap();
3422
3423 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3424 assert!(saved.pending_question.is_none());
3425 assert_eq!(
3426 saved
3427 .metadata
3428 .get("permission.reexecute_tool_call_id")
3429 .map(String::as_str),
3430 Some("call-1")
3431 );
3432 assert_eq!(
3433 saved
3434 .metadata
3435 .get("permission.reexecute_request_generation")
3436 .map(String::as_str),
3437 Some("generation-1")
3438 );
3439 }
3440
3441 #[tokio::test]
3442 async fn checkpoint_runtime_session_preserves_disk_suffix_and_appends_live_messages() {
3443 use bamboo_domain::session::types::Message;
3444
3445 let (_temp, storage) = make_storage().await;
3446 let store = LockedSessionStore::new(storage.clone());
3447 let session_id = "checkpoint-no-shrink";
3448
3449 let mut baseline = fresh(session_id);
3450 baseline.add_message(Message::user("base"));
3451 storage.save_session(&baseline).await.unwrap();
3452 let mut runner_snapshot = baseline.clone();
3453
3454 let mut durable = baseline;
3455 let mut disk_only = Message::user("concurrent injected message");
3456 disk_only.id = "disk-only".to_string();
3457 durable.add_message(disk_only);
3458 storage.save_session(&durable).await.unwrap();
3459
3460 let mut live_only = Message::assistant("partial runner output", None);
3461 live_only.id = "live-only".to_string();
3462 runner_snapshot.add_message(live_only);
3463
3464 store
3465 .checkpoint_runtime_session(&mut runner_snapshot)
3466 .await
3467 .unwrap();
3468
3469 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3470 let ids = saved
3471 .messages
3472 .iter()
3473 .map(|message| message.id.as_str())
3474 .collect::<Vec<_>>();
3475 assert_eq!(
3476 ids,
3477 vec![durable.messages[0].id.as_str(), "disk-only", "live-only"]
3478 );
3479 assert_eq!(runner_snapshot.messages.len(), saved.messages.len());
3480 assert_eq!(runner_snapshot.messages[1].id, saved.messages[1].id);
3481 assert_eq!(runner_snapshot.messages[2].id, saved.messages[2].id);
3482 assert_eq!(saved.messages[1].content, "concurrent injected message");
3483 assert_eq!(saved.messages[2].content, "partial runner output");
3484 }
3485
3486 #[tokio::test]
3487 async fn runtime_only_save_preserves_checkpointed_ledger_and_publishes_merged_state() {
3488 let (_temp, storage) = make_storage().await;
3489 let store = LockedSessionStore::new(storage.clone());
3490 let session_id = "runtime-only-ledger-race";
3491 let baseline = fresh(session_id);
3492 storage.save_session(&baseline).await.unwrap();
3493 let mut stale_control = storage.load_session(session_id).await.unwrap().unwrap();
3494
3495 let mut runner = baseline;
3496 runner.model_context_state = Some(ledger_state(1, "runner-l1"));
3497 store.checkpoint_runtime_session(&mut runner).await.unwrap();
3498
3499 stale_control.metadata.insert(
3500 "runtime.suspend_reason".to_string(),
3501 "waiting_for_children".to_string(),
3502 );
3503 let published = std::sync::Arc::new(std::sync::Mutex::new(None));
3504 let published_clone = published.clone();
3505 store
3506 .save_runtime_only_and_publish(&mut stale_control, move |saved| {
3507 *published_clone.lock().unwrap() = Some(saved.clone());
3508 })
3509 .await
3510 .unwrap();
3511
3512 let expected = runner.model_context_state.clone();
3513 assert_eq!(stale_control.model_context_state, expected);
3514 assert_eq!(
3515 published
3516 .lock()
3517 .unwrap()
3518 .as_ref()
3519 .unwrap()
3520 .model_context_state,
3521 expected
3522 );
3523 let sidecar = storage
3524 .load_runtime_control_plane(session_id)
3525 .await
3526 .unwrap()
3527 .unwrap();
3528 assert_eq!(sidecar.model_context_state, expected);
3529 assert_eq!(
3530 sidecar
3531 .metadata
3532 .get("runtime.suspend_reason")
3533 .map(String::as_str),
3534 Some("waiting_for_children")
3535 );
3536 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
3537 assert_eq!(reloaded.model_context_state, expected);
3538 }
3539
3540 #[tokio::test]
3541 async fn full_runtime_save_preserves_newer_ledger_but_commits_control_mutation() {
3542 let (_temp, storage) = make_storage().await;
3543 let store = LockedSessionStore::new(storage.clone());
3544 let session_id = "full-save-ledger-race";
3545 let baseline = fresh(session_id);
3546 storage.save_session(&baseline).await.unwrap();
3547 let mut stale = storage.load_session(session_id).await.unwrap().unwrap();
3548
3549 let mut runner = baseline;
3550 runner.model_context_state = Some(ledger_state(1, "runner-l1"));
3551 store.checkpoint_runtime_session(&mut runner).await.unwrap();
3552
3553 stale
3554 .metadata
3555 .insert("activated_tools".to_string(), "[\"search\"]".to_string());
3556 let published = std::sync::Arc::new(std::sync::Mutex::new(None));
3557 let published_clone = published.clone();
3558 store
3559 .merge_save_runtime_and_publish(&mut stale, move |saved, committed| {
3560 assert!(committed);
3561 *published_clone.lock().unwrap() = Some(saved.clone());
3562 })
3563 .await
3564 .unwrap();
3565
3566 let expected = runner.model_context_state.clone();
3567 assert_eq!(stale.model_context_state, expected);
3568 assert_eq!(
3569 published
3570 .lock()
3571 .unwrap()
3572 .as_ref()
3573 .unwrap()
3574 .model_context_state,
3575 expected
3576 );
3577 let sidecar = storage
3578 .load_runtime_control_plane(session_id)
3579 .await
3580 .unwrap()
3581 .unwrap();
3582 assert_eq!(sidecar.model_context_state, expected);
3583 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
3584 assert_eq!(reloaded.model_context_state, expected);
3585 assert_eq!(
3586 reloaded.metadata.get("activated_tools").map(String::as_str),
3587 Some("[\"search\"]")
3588 );
3589 }
3590
3591 #[tokio::test]
3592 async fn newer_explicit_epoch_reset_wins_an_ordinary_full_runtime_save() {
3593 let (_temp, storage) = make_storage().await;
3594 let store = LockedSessionStore::new(storage.clone());
3595 let session_id = "full-save-ledger-reset";
3596 let mut durable = fresh(session_id);
3597 durable.model_context_state = Some(ledger_state(1, "runner-l1"));
3598 storage.save_session(&durable).await.unwrap();
3599
3600 let mut compression = storage.load_session(session_id).await.unwrap().unwrap();
3601 compression.reset_model_context_epoch(bamboo_domain::ModelContextResetReason::Compression);
3602 let reset = compression.model_context_state.clone();
3603 assert_eq!(reset.as_ref().unwrap().state_revision, 2);
3604 store.merge_save_runtime(&mut compression).await.unwrap();
3605
3606 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
3607 assert_eq!(reloaded.model_context_state, reset);
3608 assert_eq!(
3609 reloaded
3610 .model_context_state
3611 .as_ref()
3612 .and_then(|state| state.last_reset_reason),
3613 Some(bamboo_domain::ModelContextResetReason::Compression)
3614 );
3615 }
3616
3617 #[tokio::test]
3618 async fn checkpoint_rejects_equal_revision_divergence_without_overwriting_disk() {
3619 let (_temp, storage) = make_storage().await;
3620 let store = LockedSessionStore::new(storage.clone());
3621 let session_id = "ledger-checkpoint-cas";
3622 let mut baseline = fresh(session_id);
3623 baseline.model_context_state = Some(ledger_state(1, "runner-l1"));
3624 storage.save_session(&baseline).await.unwrap();
3625
3626 let mut first = baseline.clone();
3627 first.reset_model_context_epoch(bamboo_domain::ModelContextResetReason::Compression);
3628 let mut conflicting = baseline;
3629 conflicting.reset_model_context_epoch(bamboo_domain::ModelContextResetReason::Rollback);
3630 assert_eq!(
3631 first.model_context_state.as_ref().unwrap().state_revision,
3632 conflicting
3633 .model_context_state
3634 .as_ref()
3635 .unwrap()
3636 .state_revision
3637 );
3638
3639 store.checkpoint_runtime_session(&mut first).await.unwrap();
3640 let error = store
3641 .checkpoint_runtime_session(&mut conflicting)
3642 .await
3643 .unwrap_err();
3644 assert_eq!(error.kind(), std::io::ErrorKind::WouldBlock);
3645 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
3646 assert_eq!(reloaded.model_context_state, first.model_context_state);
3647 }
3648
3649 #[tokio::test]
3650 async fn activation_checkpoint_clears_presentation_without_shrinking_concurrent_turn() {
3651 use bamboo_domain::session::runtime_state::{
3652 AgentRuntimeState, AgentStatusState, WaitingForChildrenState,
3653 };
3654 use bamboo_domain::session::types::Message;
3655
3656 let (_temp, storage) = make_storage().await;
3657 let store = LockedSessionStore::new(storage.clone());
3658 let session_id = "activation-no-shrink";
3659 let mut baseline = fresh(session_id);
3660 baseline.add_message(Message::user("base"));
3661 let mut state = AgentRuntimeState::new("activation-run");
3662 state.status = AgentStatusState::Suspended;
3663 state.waiting_for_children = Some(WaitingForChildrenState::for_children(
3664 vec!["child-1".to_string()],
3665 bamboo_domain::session::runtime_state::ChildWaitPolicy::All,
3666 chrono::Utc::now(),
3667 ));
3668 baseline.agent_runtime_state = Some(state);
3669 baseline.metadata.insert(
3670 "runtime.suspend_reason".to_string(),
3671 "waiting_for_children".to_string(),
3672 );
3673 storage.save_session(&baseline).await.unwrap();
3674 let mut activation_snapshot = baseline.clone();
3675
3676 let mut concurrent = baseline;
3677 let mut normal = Message::assistant("normal concurrent answer", None);
3678 normal.id = "normal-concurrent".to_string();
3679 concurrent.add_message(normal);
3680 storage.save_session(&concurrent).await.unwrap();
3681
3682 let state = activation_snapshot.agent_runtime_state.as_mut().unwrap();
3683 state.status = AgentStatusState::Idle;
3684 state.suspension = None;
3685 activation_snapshot
3686 .metadata
3687 .remove("runtime.suspend_reason");
3688 store
3689 .checkpoint_runtime_session(&mut activation_snapshot)
3690 .await
3691 .unwrap();
3692
3693 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3694 assert!(saved
3695 .messages
3696 .iter()
3697 .any(|message| message.id == "normal-concurrent"));
3698 let state = saved.agent_runtime_state.unwrap();
3699 assert_eq!(state.status, AgentStatusState::Idle);
3700 assert!(state.waiting_for_children.is_some());
3701 assert!(!saved.metadata.contains_key("runtime.suspend_reason"));
3702 }
3703
3704 #[tokio::test]
3705 async fn merge_save_runtime_preserves_disk_authoritative_metadata_with_single_load() {
3706 let (_temp, storage) = make_storage().await;
3712 let store = LockedSessionStore::new(storage.clone());
3713 let session_id = "runtime-merge-meta";
3714
3715 let mut baseline = fresh(session_id);
3717 baseline.title = "Auto Title".to_string();
3718 baseline.metadata_version = 0;
3719 storage.save_session(&baseline).await.unwrap();
3720
3721 let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
3723
3724 let mut renamed = storage.load_session(session_id).await.unwrap().unwrap();
3726 renamed.title = "User Renamed".to_string();
3727 renamed.title_version = 1;
3728 renamed.pinned = true;
3729 renamed.metadata_version = 1;
3730 store.commit_metadata(&renamed).await.unwrap();
3731
3732 stale_snapshot.title = "Auto Title".to_string();
3734 store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
3735
3736 let after = storage.load_session(session_id).await.unwrap().unwrap();
3737 assert_eq!(after.title, "User Renamed");
3738 assert!(after.pinned);
3739 assert_eq!(after.metadata_version, 1);
3740 assert_eq!(stale_snapshot.title, "User Renamed");
3742 assert_eq!(stale_snapshot.metadata_version, 1);
3743 }
3744
3745 #[tokio::test]
3746 async fn merge_save_runtime_preserves_durable_workflow_run_index_from_stale_runner() {
3747 let (_temp, storage) = make_storage().await;
3748 let store = LockedSessionStore::new(storage.clone());
3749 let session_id = "runtime-workflow-run-index";
3750
3751 let baseline = fresh(session_id);
3752 storage.save_session(&baseline).await.unwrap();
3753 let mut stale_runner = storage.load_session(session_id).await.unwrap().unwrap();
3754
3755 store
3756 .update_runtime_config(session_id, |session| {
3757 session.metadata.insert(
3758 "workflow.run_ids.v1".to_string(),
3759 r#"["http-started-run"]"#.to_string(),
3760 );
3761 })
3762 .await
3763 .unwrap()
3764 .expect("session exists");
3765
3766 store.merge_save_runtime(&mut stale_runner).await.unwrap();
3767
3768 assert_eq!(
3769 stale_runner
3770 .metadata
3771 .get("workflow.run_ids.v1")
3772 .map(String::as_str),
3773 Some(r#"["http-started-run"]"#)
3774 );
3775 let durable = storage.load_session(session_id).await.unwrap().unwrap();
3776 assert_eq!(
3777 durable
3778 .metadata
3779 .get("workflow.run_ids.v1")
3780 .map(String::as_str),
3781 Some(r#"["http-started-run"]"#)
3782 );
3783 }
3784
3785 #[tokio::test]
3789 async fn merge_save_runtime_adopts_disk_bypass_permissions() {
3790 use bamboo_domain::AgentRuntimeState;
3791
3792 let (_temp, storage) = make_storage().await;
3793 let store = LockedSessionStore::new(storage.clone());
3794 let session_id = "runtime-bypass";
3795
3796 let baseline = fresh(session_id);
3798 storage.save_session(&baseline).await.unwrap();
3799
3800 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
3802 loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
3803
3804 store
3806 .update_runtime_config(session_id, |s| {
3807 s.agent_runtime_state
3808 .get_or_insert_with(AgentRuntimeState::default)
3809 .bypass_permissions = true;
3810 })
3811 .await
3812 .unwrap()
3813 .expect("session exists");
3814
3815 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
3818
3819 let after = storage.load_session(session_id).await.unwrap().unwrap();
3820 assert!(
3821 after
3822 .agent_runtime_state
3823 .as_ref()
3824 .is_some_and(|s| s.bypass_permissions),
3825 "disk bypass=ON must survive a stale runtime save (#540)"
3826 );
3827 assert!(loop_snapshot
3829 .agent_runtime_state
3830 .as_ref()
3831 .is_some_and(|s| s.bypass_permissions));
3832 }
3833
3834 #[tokio::test]
3837 async fn merge_save_runtime_adopts_disk_auto_permission_mode() {
3838 use bamboo_domain::{AgentRuntimeState, SessionPermissionMode};
3839
3840 let (_temp, storage) = make_storage().await;
3841 let store = LockedSessionStore::new(storage.clone());
3842 let session_id = "runtime-auto";
3843
3844 storage.save_session(&fresh(session_id)).await.unwrap();
3845 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
3846 loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
3847 loop_snapshot.metadata.insert(
3848 "permission.requested_mode".to_string(),
3849 "default".to_string(),
3850 );
3851 loop_snapshot.metadata.insert(
3852 "permission.effective_mode".to_string(),
3853 "default".to_string(),
3854 );
3855 loop_snapshot.metadata.insert(
3856 "permission.executor_mapping".to_string(),
3857 "bamboo_runtime:default".to_string(),
3858 );
3859
3860 store
3861 .update_runtime_config(session_id, |session| {
3862 session
3863 .agent_runtime_state
3864 .get_or_insert_with(AgentRuntimeState::default)
3865 .set_permission_mode(SessionPermissionMode::Auto);
3866 session
3867 .metadata
3868 .insert("permission.policy_revision".to_string(), "12".to_string());
3869 session
3870 .metadata
3871 .insert("permission.requested_mode".to_string(), "auto".to_string());
3872 session
3873 .metadata
3874 .insert("permission.effective_mode".to_string(), "auto".to_string());
3875 session.metadata.insert(
3876 "permission.executor_mapping".to_string(),
3877 "bamboo_runtime:auto".to_string(),
3878 );
3879 session.metadata.insert(
3880 "permission.transitioned_at".to_string(),
3881 "2026-07-31T12:00:00Z".to_string(),
3882 );
3883 session.metadata_version = session.metadata_version.saturating_add(1);
3884 })
3885 .await
3886 .unwrap()
3887 .expect("session exists");
3888
3889 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
3890
3891 let durable = storage.load_session(session_id).await.unwrap().unwrap();
3892 for state in [
3893 durable.agent_runtime_state.as_ref(),
3894 loop_snapshot.agent_runtime_state.as_ref(),
3895 ] {
3896 assert_eq!(
3897 state.map(AgentRuntimeState::effective_permission_mode),
3898 Some(SessionPermissionMode::Auto)
3899 );
3900 }
3901 for session in [&durable, &loop_snapshot] {
3902 assert_eq!(
3903 session.metadata.get("permission.policy_revision"),
3904 Some(&"12".to_string())
3905 );
3906 assert_eq!(
3907 session.metadata.get("permission.requested_mode"),
3908 Some(&"auto".to_string())
3909 );
3910 assert_eq!(
3911 session.metadata.get("permission.effective_mode"),
3912 Some(&"auto".to_string())
3913 );
3914 assert_eq!(
3915 session.metadata.get("permission.executor_mapping"),
3916 Some(&"bamboo_runtime:auto".to_string())
3917 );
3918 assert_eq!(
3919 session.metadata.get("permission.transitioned_at"),
3920 Some(&"2026-07-31T12:00:00Z".to_string())
3921 );
3922 }
3923 }
3924
3925 #[tokio::test]
3928 async fn merge_save_runtime_adopts_disk_bypass_off() {
3929 use bamboo_domain::AgentRuntimeState;
3930
3931 let (_temp, storage) = make_storage().await;
3932 let store = LockedSessionStore::new(storage.clone());
3933 let session_id = "runtime-bypass-off";
3934
3935 let mut baseline = fresh(session_id);
3937 let on_state = AgentRuntimeState {
3938 bypass_permissions: true,
3939 ..AgentRuntimeState::default()
3940 };
3941 baseline.agent_runtime_state = Some(on_state);
3942 storage.save_session(&baseline).await.unwrap();
3943
3944 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
3946
3947 store
3949 .update_runtime_config(session_id, |s| {
3950 s.agent_runtime_state
3951 .get_or_insert_with(AgentRuntimeState::default)
3952 .bypass_permissions = false;
3953 })
3954 .await
3955 .unwrap()
3956 .expect("session exists");
3957
3958 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
3959
3960 let after = storage.load_session(session_id).await.unwrap().unwrap();
3961 assert!(
3962 !after
3963 .agent_runtime_state
3964 .as_ref()
3965 .is_some_and(|s| s.bypass_permissions),
3966 "disk bypass=OFF must survive a stale runtime save (#540)"
3967 );
3968 }
3969
3970 #[tokio::test]
3973 async fn save_runtime_authoritative_flags_persists_in_memory_posture_and_audit() {
3974 use bamboo_domain::AgentRuntimeState;
3975
3976 let (_temp, storage) = make_storage().await;
3977 let store = LockedSessionStore::new(storage.clone());
3978 let session_id = "child-reseed";
3979
3980 let mut baseline = fresh(session_id);
3982 let on_state = AgentRuntimeState {
3983 bypass_permissions: true,
3984 ..AgentRuntimeState::default()
3985 };
3986 baseline.agent_runtime_state = Some(on_state);
3987 for (key, value) in [
3988 ("permission.policy_revision", "12"),
3989 ("permission.requested_mode", "bypass"),
3990 ("permission.effective_mode", "bypass"),
3991 ("permission.executor_mapping", "bamboo_runtime:bypass"),
3992 ("permission.transitioned_at", "2026-07-31T12:00:00Z"),
3993 ] {
3994 baseline.metadata.insert(key.to_string(), value.to_string());
3995 }
3996 storage.save_session(&baseline).await.unwrap();
3997
3998 let mut child = storage.load_session(session_id).await.unwrap().unwrap();
4001 child
4002 .agent_runtime_state
4003 .get_or_insert_with(AgentRuntimeState::default)
4004 .bypass_permissions = false;
4005 for (key, value) in [
4006 ("permission.policy_revision", "13"),
4007 ("permission.requested_mode", "default"),
4008 ("permission.effective_mode", "default"),
4009 ("permission.executor_mapping", "bamboo_runtime:default"),
4010 ("permission.transitioned_at", "2026-07-31T12:01:00Z"),
4011 ] {
4012 child.metadata.insert(key.to_string(), value.to_string());
4013 }
4014
4015 store
4017 .save_runtime_authoritative_flags(&mut child)
4018 .await
4019 .unwrap();
4020
4021 let after = storage.load_session(session_id).await.unwrap().unwrap();
4022 assert!(
4023 !after
4024 .agent_runtime_state
4025 .as_ref()
4026 .is_some_and(|s| s.bypass_permissions),
4027 "authoritative re-seed of bypass=OFF must persist, not be reverted (#540/#74)"
4028 );
4029 for (key, value) in [
4030 ("permission.policy_revision", "13"),
4031 ("permission.requested_mode", "default"),
4032 ("permission.effective_mode", "default"),
4033 ("permission.executor_mapping", "bamboo_runtime:default"),
4034 ("permission.transitioned_at", "2026-07-31T12:01:00Z"),
4035 ] {
4036 assert_eq!(after.metadata.get(key).map(String::as_str), Some(value));
4037 }
4038 }
4039
4040 #[tokio::test]
4042 async fn merge_save_runtime_leaves_bypass_when_disk_has_no_runtime_state() {
4043 use bamboo_domain::AgentRuntimeState;
4044
4045 let (_temp, storage) = make_storage().await;
4046 let store = LockedSessionStore::new(storage.clone());
4047 let session_id = "no-runtime-state";
4048
4049 let baseline = fresh(session_id);
4051 assert!(baseline.agent_runtime_state.is_none());
4052 storage.save_session(&baseline).await.unwrap();
4053
4054 let mut running = storage.load_session(session_id).await.unwrap().unwrap();
4056 let on_state = AgentRuntimeState {
4057 bypass_permissions: true,
4058 ..AgentRuntimeState::default()
4059 };
4060 running.agent_runtime_state = Some(on_state);
4061
4062 store.merge_save_runtime(&mut running).await.unwrap();
4063
4064 assert!(
4065 running
4066 .agent_runtime_state
4067 .as_ref()
4068 .is_some_and(|s| s.bypass_permissions),
4069 "a runtime-state-less disk copy must not force bypass OFF (#540)"
4070 );
4071 }
4072
4073 #[tokio::test]
4076 async fn merge_preserves_disk_title_when_versions_equal() {
4077 let (_temp, storage) = make_storage().await;
4078 let session_id = "merge-equal";
4079
4080 let mut on_disk = fresh(session_id);
4081 on_disk.title = "User Set This".to_string();
4082 on_disk.title_version = 0;
4083 on_disk.title_generated = true;
4084 on_disk.metadata_version = 0;
4085 storage.save_session(&on_disk).await.unwrap();
4086
4087 let mut runtime_copy = fresh(session_id);
4088 runtime_copy.created_at = on_disk.created_at;
4089 runtime_copy.title = "Stale Default".to_string();
4090 runtime_copy.title_version = 0;
4091 runtime_copy.title_generated = false;
4092 runtime_copy.metadata_version = 0;
4093 runtime_copy.messages = vec![];
4094
4095 merge_save_session(&storage, &mut runtime_copy)
4096 .await
4097 .unwrap();
4098
4099 let after = storage.load_session(session_id).await.unwrap().unwrap();
4100 assert_eq!(after.title, "User Set This");
4101 assert_eq!(after.title_version, 0);
4102 assert!(after.title_generated);
4103 assert_eq!(runtime_copy.title, "User Set This");
4104 assert!(runtime_copy.title_generated);
4105 }
4106
4107 #[tokio::test]
4108 async fn merge_preserves_disk_when_disk_version_higher() {
4109 let (_temp, storage) = make_storage().await;
4110 let session_id = "merge-higher";
4111
4112 let mut on_disk = fresh(session_id);
4113 on_disk.title = "User Title v3".to_string();
4114 on_disk.title_version = 3;
4115 on_disk.metadata_version = 5;
4116 storage.save_session(&on_disk).await.unwrap();
4117
4118 let mut runtime_copy = fresh(session_id);
4119 runtime_copy.created_at = on_disk.created_at;
4120 runtime_copy.title = "Stale".to_string();
4121 runtime_copy.title_version = 1;
4122 runtime_copy.metadata_version = 0;
4123
4124 merge_save_session(&storage, &mut runtime_copy)
4125 .await
4126 .unwrap();
4127
4128 let after = storage.load_session(session_id).await.unwrap().unwrap();
4129 assert_eq!(after.title, "User Title v3");
4130 assert_eq!(after.title_version, 3);
4131 assert_eq!(after.metadata_version, 5);
4132 }
4133
4134 #[tokio::test]
4135 async fn merge_now_preserves_disk_pinned_in_metadata_group() {
4136 let (_temp, storage) = make_storage().await;
4137 let session_id = "pinned-merge";
4138
4139 let mut on_disk = fresh(session_id);
4140 on_disk.pinned = true;
4141 on_disk.metadata_version = 2;
4142 storage.save_session(&on_disk).await.unwrap();
4143
4144 let mut runtime_copy = fresh(session_id);
4145 runtime_copy.created_at = on_disk.created_at;
4146 runtime_copy.pinned = false;
4147 runtime_copy.metadata_version = 0;
4148
4149 merge_save_session(&storage, &mut runtime_copy)
4150 .await
4151 .unwrap();
4152
4153 let after = storage.load_session(session_id).await.unwrap().unwrap();
4154 assert!(
4155 after.pinned,
4156 "disk pinned=true should win over runtime false"
4157 );
4158 assert_eq!(after.metadata_version, 2);
4159 }
4160
4161 #[tokio::test]
4162 async fn merge_keeps_in_memory_when_session_version_higher() {
4163 let (_temp, storage) = make_storage().await;
4164 let session_id = "merge-bumped";
4165
4166 let mut on_disk = fresh(session_id);
4167 on_disk.title = "Old".to_string();
4168 on_disk.title_version = 1;
4169 on_disk.metadata_version = 3;
4170 storage.save_session(&on_disk).await.unwrap();
4171
4172 let mut authoritative_copy = fresh(session_id);
4173 authoritative_copy.created_at = on_disk.created_at;
4174 authoritative_copy.title = "New Authoritative".to_string();
4175 authoritative_copy.title_version = 2;
4176 authoritative_copy.metadata_version = 4;
4177 authoritative_copy.pinned = true;
4178
4179 merge_save_session(&storage, &mut authoritative_copy)
4180 .await
4181 .unwrap();
4182
4183 let after = storage.load_session(session_id).await.unwrap().unwrap();
4184 assert_eq!(after.title, "New Authoritative");
4185 assert_eq!(after.title_version, 2);
4186 assert_eq!(after.metadata_version, 4);
4187 assert!(after.pinned);
4188 }
4189
4190 #[tokio::test]
4191 async fn merge_keeps_runtime_messages_when_disk_only_changed_metadata() {
4192 let (_temp, storage) = make_storage().await;
4193 let session_id = "merge-messages";
4194
4195 let mut on_disk = fresh(session_id);
4196 on_disk.title = "Fresh Title".to_string();
4197 on_disk.title_version = 2;
4198 on_disk.metadata_version = 5;
4199 storage.save_session(&on_disk).await.unwrap();
4200
4201 let mut runtime_copy = fresh(session_id);
4202 runtime_copy.created_at = on_disk.created_at;
4203 runtime_copy.title = "Stale".to_string();
4204 runtime_copy.metadata_version = 0;
4205 runtime_copy.messages = vec![bamboo_domain::session::types::Message {
4206 role: bamboo_domain::session::types::Role::User,
4207 content: "keep me".to_string(),
4208 id: "msg-1".to_string(),
4209 created_at: chrono::Utc::now(),
4210 reasoning: None,
4211 reasoning_signature: None,
4212 content_parts: None,
4213 image_ocr: None,
4214 phase: None,
4215 tool_calls: None,
4216 tool_call_id: None,
4217 tool_success: None,
4218 compressed: false,
4219 compressed_by_event_id: None,
4220 never_compress: false,
4221 compression_level: 0,
4222 metadata: None,
4223 }];
4224
4225 merge_save_session(&storage, &mut runtime_copy)
4226 .await
4227 .unwrap();
4228
4229 let after = storage.load_session(session_id).await.unwrap().unwrap();
4230 assert_eq!(after.title, "Fresh Title");
4231 assert_eq!(after.metadata_version, 5);
4232 assert_eq!(after.messages.len(), 1);
4233 assert_eq!(after.messages[0].content, "keep me");
4234 }
4235
4236 #[tokio::test]
4237 async fn runtime_control_plane_port_uses_sidecar_without_rewriting_messages() {
4238 use bamboo_domain::session::types::Message;
4239
4240 let (_temp, storage) = make_storage().await;
4241 let store = LockedSessionStore::new(storage.clone());
4242 let session_id = "runtime-control-plane";
4243
4244 let mut durable = fresh(session_id);
4245 durable.add_message(Message::user("durable transcript"));
4246 storage.save_session(&durable).await.unwrap();
4247
4248 let mut runtime = durable.clone();
4249 runtime.model = "updated-control-plane-model".to_string();
4250 runtime.add_message(Message::assistant("uncheckpointed runtime message", None));
4251 RuntimeSessionPersistence::save_runtime_control_plane(&store, &mut runtime)
4252 .await
4253 .unwrap();
4254
4255 let control_plane =
4256 RuntimeSessionPersistence::load_runtime_control_plane(&store, session_id)
4257 .await
4258 .unwrap()
4259 .expect("control-plane exists");
4260 assert!(
4261 control_plane.messages.is_empty(),
4262 "LockedSessionStore must expose its message-free sidecar"
4263 );
4264 assert_eq!(control_plane.model, "updated-control-plane-model");
4265
4266 let reloaded = storage
4267 .load_session(session_id)
4268 .await
4269 .unwrap()
4270 .expect("session exists");
4271 assert_eq!(reloaded.model, "updated-control-plane-model");
4272 assert_eq!(
4273 reloaded.messages.len(),
4274 1,
4275 "control-plane save must not write the uncheckpointed message"
4276 );
4277 assert_eq!(reloaded.messages[0].content, "durable transcript");
4278 }
4279
4280 #[tokio::test]
4281 async fn atomic_task_patch_loads_inside_lock_and_preserves_interleaved_runtime_state() {
4282 let temp = tempfile::tempdir().unwrap();
4283 let inner = Arc::new(
4284 SessionStoreV2::new(temp.path().to_path_buf())
4285 .await
4286 .expect("storage init"),
4287 );
4288 let session_id = "atomic-task-patch";
4289 inner
4290 .save_session(&fresh(session_id))
4291 .await
4292 .expect("seed session");
4293
4294 let counted = Arc::new(CountingControlPlaneStorage {
4295 inner: inner.clone(),
4296 control_plane_loads: AtomicUsize::new(0),
4297 full_saves: AtomicUsize::new(0),
4298 runtime_state_saves: AtomicUsize::new(0),
4299 });
4300 let storage: Arc<dyn Storage> = counted.clone();
4301 let store = Arc::new(LockedSessionStore::new(storage));
4302 let guard = store.acquire_lock(session_id).await;
4303 let now = chrono::Utc::now();
4304 let task_list = bamboo_domain::TaskList {
4305 session_id: session_id.to_string(),
4306 title: "Atomic Task patch".to_string(),
4307 items: Vec::new(),
4308 created_at: now,
4309 updated_at: now,
4310 };
4311 let (started_tx, started_rx) = tokio::sync::oneshot::channel();
4312 let patch_store = store.clone();
4313 let patch = tokio::spawn(async move {
4314 let _ = started_tx.send(());
4315 RuntimeSessionPersistence::update_task_list_control_plane(
4316 patch_store.as_ref(),
4317 session_id,
4318 &task_list,
4319 "9",
4320 )
4321 .await
4322 });
4323 started_rx.await.expect("patch task started");
4324 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
4325 assert_eq!(
4326 counted.control_plane_loads.load(Ordering::SeqCst),
4327 0,
4328 "Task patch must acquire the session lock before loading its snapshot"
4329 );
4330
4331 let mut latest = inner
4335 .load_runtime_control_plane(session_id)
4336 .await
4337 .expect("load latest control-plane")
4338 .expect("control-plane exists");
4339 latest.agent_runtime_state = Some(bamboo_domain::AgentRuntimeState::new("latest-run"));
4340 latest
4341 .metadata
4342 .insert("concurrent.runtime".to_string(), "preserve".to_string());
4343 inner
4344 .save_runtime_state(&latest)
4345 .await
4346 .expect("publish concurrent runtime transition");
4347 drop(guard);
4348
4349 assert!(
4350 patch.await.expect("patch join").expect("patch succeeds"),
4351 "existing root must be patched"
4352 );
4353 assert_eq!(counted.control_plane_loads.load(Ordering::SeqCst), 1);
4354 let reloaded = inner
4355 .load_session(session_id)
4356 .await
4357 .expect("reload")
4358 .expect("session exists");
4359 assert_eq!(
4360 reloaded
4361 .agent_runtime_state
4362 .as_ref()
4363 .map(|state| state.run_id.as_str()),
4364 Some("latest-run")
4365 );
4366 assert_eq!(
4367 reloaded
4368 .metadata
4369 .get("concurrent.runtime")
4370 .map(String::as_str),
4371 Some("preserve")
4372 );
4373 assert_eq!(reloaded.task_list_version_meta().as_deref(), Some("9"));
4374 assert_eq!(
4375 reloaded.task_list.as_ref().map(|list| list.title.as_str()),
4376 Some("Atomic Task patch")
4377 );
4378 }
4379
4380 #[tokio::test]
4381 async fn paired_task_cas_conflict_cannot_overwrite_newer_root_or_child_state() {
4382 let (_temp, storage) = make_storage().await;
4383 let store = LockedSessionStore::new(storage.clone());
4384 let root_id = "task-cas-root";
4385 let child_id = "task-cas-child";
4386 let now = chrono::Utc::now();
4387 let task_list = |title: &str| bamboo_domain::TaskList {
4388 session_id: root_id.to_string(),
4389 title: title.to_string(),
4390 items: Vec::new(),
4391 created_at: now,
4392 updated_at: now,
4393 };
4394
4395 let mut root = fresh(root_id);
4396 root.set_task_list(task_list("newer root"));
4397 root.set_task_list_version_meta("2");
4398 storage.save_session(&root).await.expect("seed root");
4399 let mut child = Session::new_child(child_id, root_id, "model", "child");
4400 child.set_task_list(task_list("current child"));
4401 child.set_task_list_version_meta("1");
4402 storage.save_session(&child).await.expect("seed child");
4403
4404 let updated = RuntimeSessionPersistence::update_task_list_control_planes_if_version(
4405 &store,
4406 child_id,
4407 root_id,
4408 "1",
4409 &task_list("current child"),
4410 &task_list("stale evaluator"),
4411 "3",
4412 )
4413 .await
4414 .expect("CAS returns clean conflict");
4415 assert!(
4416 !updated,
4417 "mismatched root generation must reject both writes"
4418 );
4419
4420 let durable_root = storage
4421 .load_session(root_id)
4422 .await
4423 .expect("load root")
4424 .expect("root exists");
4425 let durable_child = storage
4426 .load_session(child_id)
4427 .await
4428 .expect("load child")
4429 .expect("child exists");
4430 assert_eq!(durable_root.task_list_version_meta().as_deref(), Some("2"));
4431 assert_eq!(
4432 durable_root
4433 .task_list
4434 .as_ref()
4435 .map(|list| list.title.as_str()),
4436 Some("newer root")
4437 );
4438 assert_eq!(durable_child.task_list_version_meta().as_deref(), Some("1"));
4439 assert_eq!(
4440 durable_child
4441 .task_list
4442 .as_ref()
4443 .map(|list| list.title.as_str()),
4444 Some("current child")
4445 );
4446 }
4447
4448 #[tokio::test]
4449 async fn paired_task_cas_success_uses_only_targeted_saves_and_preserves_both_transcripts() {
4450 use bamboo_domain::session::types::Message;
4451
4452 let temp = tempfile::tempdir().unwrap();
4453 let inner = Arc::new(
4454 SessionStoreV2::new(temp.path().to_path_buf())
4455 .await
4456 .expect("storage init"),
4457 );
4458 let root_id = "task-cas-success-root";
4459 let child_id = "task-cas-success-child";
4460 let now = chrono::Utc::now();
4461 let task_list = |title: &str| bamboo_domain::TaskList {
4462 session_id: root_id.to_string(),
4463 title: title.to_string(),
4464 items: Vec::new(),
4465 created_at: now,
4466 updated_at: now,
4467 };
4468
4469 let mut root = fresh(root_id);
4470 root.add_message(Message::user("root transcript"));
4471 root.metadata
4472 .insert("unrelated.root".to_string(), "preserve".to_string());
4473 root.agent_runtime_state = Some(bamboo_domain::AgentRuntimeState::new("root-run"));
4474 root.set_task_list(task_list("old shared"));
4475 root.set_task_list_version_meta("1");
4476 inner.save_session(&root).await.expect("seed root");
4477
4478 let mut child = Session::new_child(child_id, root_id, "model", "child");
4479 child.add_message(Message::user("child transcript"));
4480 child
4481 .metadata
4482 .insert("unrelated.child".to_string(), "preserve".to_string());
4483 child.agent_runtime_state = Some(bamboo_domain::AgentRuntimeState::new("child-run"));
4484 child.set_task_list(task_list("old shared"));
4485 child.set_task_list_version_meta("1");
4486 inner.save_session(&child).await.expect("seed child");
4487
4488 let counted = Arc::new(CountingControlPlaneStorage {
4489 inner: inner.clone(),
4490 control_plane_loads: AtomicUsize::new(0),
4491 full_saves: AtomicUsize::new(0),
4492 runtime_state_saves: AtomicUsize::new(0),
4493 });
4494 let storage: Arc<dyn Storage> = counted.clone();
4495 let store = LockedSessionStore::new(storage);
4496 assert!(
4497 RuntimeSessionPersistence::update_task_list_control_planes_if_version(
4498 &store,
4499 child_id,
4500 root_id,
4501 "1",
4502 &task_list("old shared"),
4503 &task_list("evaluated"),
4504 "2",
4505 )
4506 .await
4507 .expect("paired CAS succeeds")
4508 );
4509 assert_eq!(counted.control_plane_loads.load(Ordering::SeqCst), 2);
4510 assert_eq!(counted.runtime_state_saves.load(Ordering::SeqCst), 2);
4511 assert_eq!(
4512 counted.full_saves.load(Ordering::SeqCst),
4513 0,
4514 "evaluation CAS must not call full save_session for child or root"
4515 );
4516
4517 let durable_root = inner.load_session(root_id).await.unwrap().unwrap();
4518 let durable_child = inner.load_session(child_id).await.unwrap().unwrap();
4519 for (session, transcript, metadata_key, run_id) in [
4520 (
4521 &durable_root,
4522 "root transcript",
4523 "unrelated.root",
4524 "root-run",
4525 ),
4526 (
4527 &durable_child,
4528 "child transcript",
4529 "unrelated.child",
4530 "child-run",
4531 ),
4532 ] {
4533 assert_eq!(session.task_list_version_meta().as_deref(), Some("2"));
4534 assert_eq!(
4535 session.task_list.as_ref().map(|list| list.title.as_str()),
4536 Some("evaluated")
4537 );
4538 assert_eq!(session.messages.len(), 1);
4539 assert_eq!(session.messages[0].content, transcript);
4540 assert_eq!(
4541 session.metadata.get(metadata_key).map(String::as_str),
4542 Some("preserve")
4543 );
4544 assert_eq!(
4545 session
4546 .agent_runtime_state
4547 .as_ref()
4548 .map(|state| state.run_id.as_str()),
4549 Some(run_id)
4550 );
4551 }
4552 }
4553
4554 #[tokio::test]
4555 async fn locked_runtime_and_full_saves_adopt_task_conflicts_before_publish() {
4556 let temp = tempfile::tempdir().unwrap();
4557 let home = temp.path().to_path_buf();
4558 let first_storage = Arc::new(
4559 SessionStoreV2::new(home.clone())
4560 .await
4561 .expect("first storage init"),
4562 );
4563 let second_storage = Arc::new(
4564 SessionStoreV2::new(home)
4565 .await
4566 .expect("second storage init"),
4567 );
4568 let now = chrono::Utc::now();
4569 let task_list = |session_id: &str, title: &str| bamboo_domain::TaskList {
4570 session_id: session_id.to_string(),
4571 title: title.to_string(),
4572 items: Vec::new(),
4573 created_at: now,
4574 updated_at: now,
4575 };
4576
4577 let runtime_id = "ordinary-task-retry-runtime";
4578 let full_id = "ordinary-task-retry-full";
4579 let mut runtime_initial = fresh(runtime_id);
4580 runtime_initial.set_task_list(task_list(runtime_id, "runtime v1"));
4581 runtime_initial.set_task_list_version_meta("1");
4582 first_storage
4583 .save_session(&runtime_initial)
4584 .await
4585 .expect("seed runtime session");
4586 let mut full_initial = fresh(full_id);
4587 full_initial.set_task_list(task_list(full_id, "full v1"));
4588 full_initial.set_task_list_version_meta("1");
4589 first_storage
4590 .save_session(&full_initial)
4591 .await
4592 .expect("seed full session");
4593
4594 let mut runtime_advanced = runtime_initial.clone();
4595 runtime_advanced.set_task_list(task_list(runtime_id, "runtime v2"));
4596 runtime_advanced.set_task_list_version_meta("2");
4597 second_storage
4598 .save_runtime_state(&runtime_advanced)
4599 .await
4600 .expect("advance runtime Task generation");
4601 let mut full_advanced = full_initial.clone();
4602 full_advanced.set_task_list(task_list(full_id, "full v2"));
4603 full_advanced.set_task_list_version_meta("2");
4604 second_storage
4605 .save_runtime_state(&full_advanced)
4606 .await
4607 .expect("advance full Task generation");
4608
4609 let storage: Arc<dyn Storage> = first_storage.clone();
4610 let store = LockedSessionStore::new(storage);
4611 let runtime_published = Arc::new(std::sync::Mutex::new(None));
4612 let runtime_callback = runtime_published.clone();
4613 let mut runtime_stale = runtime_initial;
4614 runtime_stale
4615 .metadata
4616 .insert("runtime.non-task".to_string(), "preserved".to_string());
4617 store
4618 .save_runtime_only_and_publish(&mut runtime_stale, move |saved| {
4619 *runtime_callback.lock().expect("runtime publish lock") = Some(saved.clone());
4620 })
4621 .await
4622 .expect("locked runtime save rebases and retries");
4623
4624 let full_published = Arc::new(std::sync::Mutex::new(None));
4625 let full_callback = full_published.clone();
4626 let mut full_stale = full_initial;
4627 full_stale
4628 .metadata
4629 .insert("full.non-task".to_string(), "preserved".to_string());
4630 store
4631 .merge_save_runtime_and_publish(&mut full_stale, move |saved, committed| {
4632 assert!(committed);
4633 *full_callback.lock().expect("full publish lock") = Some(saved.clone());
4634 })
4635 .await
4636 .expect("locked full save rebases and retries");
4637
4638 let durable_runtime = first_storage
4639 .load_session(runtime_id)
4640 .await
4641 .unwrap()
4642 .expect("durable runtime session");
4643 let durable_full = first_storage
4644 .load_session(full_id)
4645 .await
4646 .unwrap()
4647 .expect("durable full session");
4648 let published_runtime = runtime_published
4649 .lock()
4650 .expect("runtime publish lock")
4651 .clone()
4652 .expect("runtime published snapshot");
4653 let published_full = full_published
4654 .lock()
4655 .expect("full publish lock")
4656 .clone()
4657 .expect("full published snapshot");
4658
4659 for (session, expected_title, metadata_key) in [
4660 (&runtime_stale, "runtime v2", "runtime.non-task"),
4661 (&durable_runtime, "runtime v2", "runtime.non-task"),
4662 (&published_runtime, "runtime v2", "runtime.non-task"),
4663 (&full_stale, "full v2", "full.non-task"),
4664 (&durable_full, "full v2", "full.non-task"),
4665 (&published_full, "full v2", "full.non-task"),
4666 ] {
4667 assert_eq!(session.task_list_version_meta().as_deref(), Some("2"));
4668 assert_eq!(
4669 session.task_list.as_ref().map(|list| list.title.as_str()),
4670 Some(expected_title)
4671 );
4672 assert_eq!(
4673 session.metadata.get(metadata_key).map(String::as_str),
4674 Some("preserved")
4675 );
4676 }
4677 }
4678
4679 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
4680 async fn single_task_cas_rejects_same_version_divergent_snapshot_without_callback() {
4681 let (_temp, storage) = make_storage().await;
4682 let store = LockedSessionStore::new(storage.clone());
4683 let session_id = "single-task-exact-snapshot";
4684 let now = chrono::Utc::now();
4685 let task_list = |title: &str| bamboo_domain::TaskList {
4686 session_id: session_id.to_string(),
4687 title: title.to_string(),
4688 items: Vec::new(),
4689 created_at: now,
4690 updated_at: now,
4691 };
4692 let durable_winner = task_list("durable winner");
4693 let stale_snapshot = task_list("stale same-version snapshot");
4694 let mut session = fresh(session_id);
4695 session.set_task_list(durable_winner.clone());
4696 session.set_task_list_version_meta("1");
4697 storage.save_session(&session).await.expect("seed session");
4698
4699 let published = Arc::new(AtomicBool::new(false));
4700 let callback = published.clone();
4701 assert!(!store
4702 .update_task_list_control_plane_if_version_and_publish(
4703 session_id,
4704 "1",
4705 &stale_snapshot,
4706 &task_list("stale evaluation"),
4707 "2",
4708 move |_| callback.store(true, Ordering::SeqCst),
4709 )
4710 .await
4711 .expect("same-version divergence is a clean stale result"));
4712 assert!(!published.load(Ordering::SeqCst));
4713 let durable = storage
4714 .load_session(session_id)
4715 .await
4716 .unwrap()
4717 .expect("session remains");
4718 assert_eq!(durable.task_list_version_meta().as_deref(), Some("1"));
4719 assert_eq!(
4720 serde_json::to_value(&durable.task_list).expect("serialize durable Task list"),
4721 serde_json::to_value(Some(&durable_winner)).expect("serialize expected Task list")
4722 );
4723 }
4724
4725 #[tokio::test]
4726 async fn paired_task_cas_rejects_same_version_divergent_snapshot_without_callback() {
4727 let (_temp, storage) = make_storage().await;
4728 let store = LockedSessionStore::new(storage.clone());
4729 let root_id = "paired-task-exact-snapshot-root";
4730 let child_id = "paired-task-exact-snapshot-child";
4731 let now = chrono::Utc::now();
4732 let task_list = |title: &str| bamboo_domain::TaskList {
4733 session_id: root_id.to_string(),
4734 title: title.to_string(),
4735 items: Vec::new(),
4736 created_at: now,
4737 updated_at: now,
4738 };
4739 let durable_winner = task_list("durable winner");
4740 let stale_snapshot = task_list("stale same-version snapshot");
4741 let mut root = fresh(root_id);
4742 root.set_task_list(durable_winner.clone());
4743 root.set_task_list_version_meta("1");
4744 storage.save_session(&root).await.expect("seed root");
4745 let mut child = Session::new_child(child_id, root_id, "model", "child");
4746 child.set_task_list(durable_winner.clone());
4747 child.set_task_list_version_meta("1");
4748 storage.save_session(&child).await.expect("seed child");
4749
4750 let published = Arc::new(AtomicBool::new(false));
4751 let callback = published.clone();
4752 assert!(!store
4753 .update_task_list_control_planes_if_version_and_publish(
4754 child_id,
4755 root_id,
4756 "1",
4757 &stale_snapshot,
4758 &task_list("stale evaluation"),
4759 "2",
4760 move |_, _| callback.store(true, Ordering::SeqCst),
4761 )
4762 .await
4763 .expect("same-version divergence is a clean stale result"));
4764 assert!(!published.load(Ordering::SeqCst));
4765 for id in [child_id, root_id] {
4766 let durable = storage
4767 .load_session(id)
4768 .await
4769 .unwrap()
4770 .expect("session remains");
4771 assert_eq!(durable.task_list_version_meta().as_deref(), Some("1"));
4772 assert_eq!(
4773 serde_json::to_value(&durable.task_list).expect("serialize durable Task list"),
4774 serde_json::to_value(Some(&durable_winner)).expect("serialize expected Task list")
4775 );
4776 }
4777 }
4778
4779 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
4780 async fn unconditional_root_task_patch_reports_final_cas_conflict_without_publishing() {
4781 let temp = tempfile::tempdir().unwrap();
4782 let home = temp.path().to_path_buf();
4783 let first_inner = Arc::new(
4784 SessionStoreV2::new(home.clone())
4785 .await
4786 .expect("first storage init"),
4787 );
4788 let root_id = "single-task-unconditional-loser";
4789 let now = chrono::Utc::now();
4790 let task_list = |title: &str| bamboo_domain::TaskList {
4791 session_id: root_id.to_string(),
4792 title: title.to_string(),
4793 items: Vec::new(),
4794 created_at: now,
4795 updated_at: now,
4796 };
4797 let mut root = fresh(root_id);
4798 root.task_list = Some(task_list("original"));
4799 root.set_task_list_version_meta("1");
4800 first_inner.save_session(&root).await.expect("seed root");
4801 let second_inner = Arc::new(
4802 SessionStoreV2::new(home)
4803 .await
4804 .expect("second storage init"),
4805 );
4806 let commit_reached = Arc::new(tokio::sync::Barrier::new(2));
4807 let release_commit = Arc::new(tokio::sync::Barrier::new(2));
4808 let storage: Arc<dyn Storage> = Arc::new(SingleCommitPauseStorage {
4809 inner: first_inner.clone(),
4810 commit_reached: commit_reached.clone(),
4811 release_commit: release_commit.clone(),
4812 });
4813 let store = LockedSessionStore::new(storage);
4814 let published = Arc::new(AtomicBool::new(false));
4815 let callback = published.clone();
4816 let loser_candidate = task_list("loser");
4817
4818 let loser = store.update_task_list_control_plane_and_publish(
4819 root_id,
4820 &loser_candidate,
4821 "2",
4822 move |_| callback.store(true, Ordering::SeqCst),
4823 );
4824 let winner = async {
4825 commit_reached.wait().await;
4826 let original = second_inner
4827 .load_runtime_control_plane(root_id)
4828 .await
4829 .expect("load winner original")
4830 .expect("winner original exists");
4831 let mut updated = original.clone();
4832 updated.task_list = Some(task_list("winner"));
4833 updated.set_task_list_version_meta("2");
4834 assert!(second_inner
4835 .save_task_control_plane_if_matches(&original, &updated)
4836 .await
4837 .expect("commit winner"));
4838 release_commit.wait().await;
4839 };
4840 let (loser_result, ()) = tokio::join!(loser, winner);
4841 let error = loser_result.expect_err("unconditional loser must be an explicit conflict");
4842 assert_eq!(error.kind(), std::io::ErrorKind::WouldBlock);
4843 assert!(!published.load(Ordering::SeqCst));
4844 let durable = first_inner
4845 .load_session(root_id)
4846 .await
4847 .unwrap()
4848 .expect("durable root");
4849 assert_eq!(durable.task_list_version_meta().as_deref(), Some("2"));
4850 assert_eq!(
4851 durable.task_list.as_ref().map(|list| list.title.as_str()),
4852 Some("winner")
4853 );
4854 }
4855
4856 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
4857 async fn independent_root_task_patches_have_one_final_cas_winner() {
4858 let temp = tempfile::tempdir().unwrap();
4859 let home = temp.path().to_path_buf();
4860 let first_inner = Arc::new(
4861 SessionStoreV2::new(home.clone())
4862 .await
4863 .expect("first storage init"),
4864 );
4865 let root_id = "single-task-cas-race-root";
4866 let now = chrono::Utc::now();
4867 let task_list = |title: &str| bamboo_domain::TaskList {
4868 session_id: root_id.to_string(),
4869 title: title.to_string(),
4870 items: Vec::new(),
4871 created_at: now,
4872 updated_at: now,
4873 };
4874 let mut root = fresh(root_id);
4875 root.set_task_list(task_list("original"));
4876 root.set_task_list_version_meta("1");
4877 first_inner.save_session(&root).await.expect("seed root");
4878 let second_inner = Arc::new(
4879 SessionStoreV2::new(home)
4880 .await
4881 .expect("second storage init"),
4882 );
4883 let before_commit = Arc::new(tokio::sync::Barrier::new(2));
4884 let first_storage: Arc<dyn Storage> = Arc::new(SingleCommitBarrierStorage {
4885 inner: first_inner.clone(),
4886 before_commit: before_commit.clone(),
4887 });
4888 let second_storage: Arc<dyn Storage> = Arc::new(SingleCommitBarrierStorage {
4889 inner: second_inner.clone(),
4890 before_commit,
4891 });
4892 let first_store = LockedSessionStore::new(first_storage);
4893 let second_store = LockedSessionStore::new(second_storage);
4894 let first_published = Arc::new(AtomicBool::new(false));
4895 let second_published = Arc::new(AtomicBool::new(false));
4896 let first_callback = first_published.clone();
4897 let second_callback = second_published.clone();
4898 let expected = task_list("original");
4899 let first_candidate = task_list("candidate one");
4900 let second_candidate = task_list("candidate two");
4901
4902 let first = first_store.update_task_list_control_plane_if_version_and_publish(
4905 root_id,
4906 "1",
4907 &expected,
4908 &first_candidate,
4909 "2",
4910 move |_| first_callback.store(true, Ordering::SeqCst),
4911 );
4912 let second = second_store.update_task_list_control_plane_if_version_and_publish(
4913 root_id,
4914 "1",
4915 &expected,
4916 &second_candidate,
4917 "2",
4918 move |_| second_callback.store(true, Ordering::SeqCst),
4919 );
4920 let (first_result, second_result) = tokio::join!(first, second);
4921 let first_won = first_result.expect("first root result");
4922 let second_won = second_result.expect("second root result");
4923 assert_ne!(first_won, second_won, "exactly one root candidate wins");
4924 assert_eq!(first_published.load(Ordering::SeqCst), first_won);
4925 assert_eq!(second_published.load(Ordering::SeqCst), second_won);
4926
4927 let durable = first_inner
4928 .load_session(root_id)
4929 .await
4930 .unwrap()
4931 .expect("root");
4932 assert_eq!(durable.task_list_version_meta().as_deref(), Some("2"));
4933 assert_eq!(
4934 durable.task_list.as_ref().map(|list| list.title.as_str()),
4935 Some(if first_won {
4936 "candidate one"
4937 } else {
4938 "candidate two"
4939 })
4940 );
4941 }
4942
4943 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
4944 async fn independent_locked_stores_revalidate_pair_cas_at_storage_commit_point() {
4945 let temp = tempfile::tempdir().unwrap();
4946 let home = temp.path().to_path_buf();
4947 let first_inner = Arc::new(
4948 SessionStoreV2::new(home.clone())
4949 .await
4950 .expect("first storage init"),
4951 );
4952 let root_id = "task-cas-race-root";
4953 let child_id = "task-cas-race-child";
4954 let now = chrono::Utc::now();
4955 let task_list = |title: &str| bamboo_domain::TaskList {
4956 session_id: root_id.to_string(),
4957 title: title.to_string(),
4958 items: Vec::new(),
4959 created_at: now,
4960 updated_at: now,
4961 };
4962
4963 let mut root = fresh(root_id);
4964 root.set_task_list(task_list("original shared"));
4965 root.set_task_list_version_meta("1");
4966 first_inner.save_session(&root).await.expect("seed root");
4967 let mut child = Session::new_child(child_id, root_id, "model", "child");
4968 child.set_task_list(task_list("original shared"));
4969 child.set_task_list_version_meta("1");
4970 first_inner.save_session(&child).await.expect("seed child");
4971
4972 let second_inner = Arc::new(
4973 SessionStoreV2::new(home)
4974 .await
4975 .expect("second storage init"),
4976 );
4977 let before_commit = Arc::new(tokio::sync::Barrier::new(2));
4978 let first_storage: Arc<dyn Storage> = Arc::new(PairCommitBarrierStorage {
4979 inner: first_inner.clone(),
4980 before_commit: before_commit.clone(),
4981 });
4982 let second_storage: Arc<dyn Storage> = Arc::new(PairCommitBarrierStorage {
4983 inner: second_inner.clone(),
4984 before_commit,
4985 });
4986 let first_store = LockedSessionStore::new(first_storage);
4987 let second_store = LockedSessionStore::new(second_storage);
4988 let first_published = Arc::new(AtomicBool::new(false));
4989 let second_published = Arc::new(AtomicBool::new(false));
4990 let first_published_callback = first_published.clone();
4991 let second_published_callback = second_published.clone();
4992 let expected = task_list("original shared");
4993 let first_candidate = task_list("candidate one");
4994 let second_candidate = task_list("candidate two");
4995
4996 let first = first_store.update_task_list_control_planes_if_version_and_publish(
4997 child_id,
4998 root_id,
4999 "1",
5000 &expected,
5001 &first_candidate,
5002 "2",
5003 move |_, _| first_published_callback.store(true, Ordering::SeqCst),
5004 );
5005 let second = second_store.update_task_list_control_planes_if_version_and_publish(
5006 child_id,
5007 root_id,
5008 "1",
5009 &expected,
5010 &second_candidate,
5011 "2",
5012 move |_, _| second_published_callback.store(true, Ordering::SeqCst),
5013 );
5014 let (first_result, second_result) = tokio::join!(first, second);
5015 let first_won = first_result.expect("first CAS result");
5016 let second_won = second_result.expect("second CAS result");
5017 assert_ne!(first_won, second_won, "exactly one staged v1 CAS may win");
5018 assert_eq!(first_published.load(Ordering::SeqCst), first_won);
5019 assert_eq!(second_published.load(Ordering::SeqCst), second_won);
5020
5021 let expected_title = if first_won {
5022 "candidate one"
5023 } else {
5024 "candidate two"
5025 };
5026 let durable_child = first_inner
5027 .load_session(child_id)
5028 .await
5029 .unwrap()
5030 .expect("child");
5031 let durable_root = second_inner
5032 .load_session(root_id)
5033 .await
5034 .unwrap()
5035 .expect("root");
5036 for session in [&durable_child, &durable_root] {
5037 assert_eq!(session.task_list_version_meta().as_deref(), Some("2"));
5038 assert_eq!(
5039 session.task_list.as_ref().map(|list| list.title.as_str()),
5040 Some(expected_title)
5041 );
5042 }
5043 }
5044
5045 #[tokio::test]
5046 async fn paired_task_second_write_failure_rolls_back_and_skips_publish_callback() {
5047 use bamboo_domain::session::types::Message;
5048
5049 let temp = tempfile::tempdir().unwrap();
5050 let inner = Arc::new(
5051 SessionStoreV2::new(temp.path().to_path_buf())
5052 .await
5053 .expect("storage init"),
5054 );
5055 let root_id = "task-cas-failure-root";
5056 let child_id = "task-cas-failure-child";
5057 let now = chrono::Utc::now();
5058 let task_list = |title: &str| bamboo_domain::TaskList {
5059 session_id: root_id.to_string(),
5060 title: title.to_string(),
5061 items: Vec::new(),
5062 created_at: now,
5063 updated_at: now,
5064 };
5065
5066 let mut root = fresh(root_id);
5067 root.add_message(Message::user("root transcript"));
5068 root.metadata
5069 .insert("unrelated.root".to_string(), "preserve".to_string());
5070 root.set_task_list(task_list("old shared"));
5071 root.set_task_list_version_meta("1");
5072 inner.save_session(&root).await.expect("seed root");
5073
5074 let mut child = Session::new_child(child_id, root_id, "model", "child");
5075 child.add_message(Message::user("child transcript"));
5076 child
5077 .metadata
5078 .insert("unrelated.child".to_string(), "preserve".to_string());
5079 child.set_task_list(task_list("old shared"));
5080 child.set_task_list_version_meta("1");
5081 inner.save_session(&child).await.expect("seed child");
5082
5083 inner
5084 .inject_runtime_task_transaction_fault(RuntimeTaskTransactionFault::SecondUpdatedWrite);
5085 let storage: Arc<dyn Storage> = inner.clone();
5086 let store = LockedSessionStore::new(storage);
5087 let published = Arc::new(AtomicBool::new(false));
5088 let published_for_callback = published.clone();
5089 let error = store
5090 .update_task_list_control_planes_if_version_and_publish(
5091 child_id,
5092 root_id,
5093 "1",
5094 &task_list("old shared"),
5095 &task_list("must roll back"),
5096 "2",
5097 move |_, _| published_for_callback.store(true, Ordering::SeqCst),
5098 )
5099 .await
5100 .expect_err("injected second write fails");
5101 assert!(error.to_string().contains("rolled back"), "{error}");
5102 assert!(
5103 !published.load(Ordering::SeqCst),
5104 "durable failure must not publish either cache snapshot"
5105 );
5106
5107 let durable_root = inner.load_session(root_id).await.unwrap().unwrap();
5108 let durable_child = inner.load_session(child_id).await.unwrap().unwrap();
5109 for (session, title, transcript, metadata_key) in [
5110 (
5111 &durable_root,
5112 "old shared",
5113 "root transcript",
5114 "unrelated.root",
5115 ),
5116 (
5117 &durable_child,
5118 "old shared",
5119 "child transcript",
5120 "unrelated.child",
5121 ),
5122 ] {
5123 assert_eq!(session.task_list_version_meta().as_deref(), Some("1"));
5124 assert_eq!(
5125 session.task_list.as_ref().map(|list| list.title.as_str()),
5126 Some(title)
5127 );
5128 assert_eq!(session.messages[0].content, transcript);
5129 assert_eq!(
5130 session.metadata.get(metadata_key).map(String::as_str),
5131 Some("preserve")
5132 );
5133 }
5134 }
5135
5136 #[tokio::test]
5139 async fn locked_merge_save_runtime_serialises_concurrent_writes() {
5140 let (_temp, storage) = make_storage().await;
5141 let store = Arc::new(LockedSessionStore::new(storage));
5142 let session_id = "lock-serial".to_string();
5143
5144 let base = fresh(&session_id);
5146 store.storage().save_session(&base).await.unwrap();
5147
5148 let store_a = store.clone();
5151 let store_b = store.clone();
5152 let sid_a = session_id.clone();
5153 let sid_b = session_id.clone();
5154
5155 let a = tokio::spawn(async move {
5156 let _guard = store_a.acquire_lock(&sid_a).await;
5157 let mut s = store_a
5158 .storage()
5159 .load_session(&sid_a)
5160 .await
5161 .unwrap()
5162 .unwrap();
5163 s.title = "Writer A".to_string();
5164 s.title_version = s.title_version.saturating_add(1);
5165 s.metadata_version = s.metadata_version.saturating_add(1);
5166 s.updated_at = chrono::Utc::now();
5167 store_a.storage().save_session(&s).await.unwrap();
5168 s.title_version
5169 });
5170
5171 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
5173
5174 let b = tokio::spawn(async move {
5175 let _guard = store_b.acquire_lock(&sid_b).await;
5176 let mut s = store_b
5177 .storage()
5178 .load_session(&sid_b)
5179 .await
5180 .unwrap()
5181 .unwrap();
5182 s.title = "Writer B".to_string();
5183 s.title_version = s.title_version.saturating_add(1);
5184 s.metadata_version = s.metadata_version.saturating_add(1);
5185 s.updated_at = chrono::Utc::now();
5186 store_b.storage().save_session(&s).await.unwrap();
5187 s.title_version
5188 });
5189
5190 let (ver_a, ver_b) = tokio::join!(a, b);
5191 let final_s = store
5192 .storage()
5193 .load_session(&session_id)
5194 .await
5195 .unwrap()
5196 .unwrap();
5197 assert!(
5198 ver_a.unwrap() != ver_b.unwrap(),
5199 "concurrent writers must produce distinct versions"
5200 );
5201 assert_eq!(final_s.metadata_version, 2);
5202 }
5203
5204 #[tokio::test]
5205 async fn commit_metadata_is_plain_save_inside_lock() {
5206 let (_temp, storage) = make_storage().await;
5207 let store = LockedSessionStore::new(storage);
5208 let session_id = "commit-plain";
5209
5210 let mut s = fresh(session_id);
5211 s.title = "Committed".to_string();
5212 s.metadata_version = 1;
5213 s.title_version = 2;
5214
5215 store.commit_metadata(&s).await.unwrap();
5216
5217 let after = store
5218 .storage()
5219 .load_session(session_id)
5220 .await
5221 .unwrap()
5222 .unwrap();
5223 assert_eq!(after.title, "Committed");
5224 assert_eq!(after.metadata_version, 1);
5225 assert_eq!(after.title_version, 2);
5226 }
5227
5228 #[tokio::test]
5231 async fn acquire_lock_self_evicts_when_no_other_holder() {
5232 let (_temp, storage) = make_storage().await;
5233 let store = LockedSessionStore::new(storage);
5234
5235 {
5236 let _guard = store.acquire_lock("solo").await;
5237 assert_eq!(store.locks.len(), 1, "entry present while the lock is held");
5238 }
5239 assert_eq!(
5242 store.locks.len(),
5243 0,
5244 "lock entry must be evicted once released with no other holder"
5245 );
5246 }
5247
5248 #[tokio::test]
5249 async fn acquire_lock_many_distinct_ids_do_not_accumulate() {
5250 let (_temp, storage) = make_storage().await;
5251 let store = LockedSessionStore::new(storage);
5252
5253 for i in 0..100 {
5255 let _guard = store.acquire_lock(&format!("sess-{i}")).await;
5256 }
5257 assert_eq!(
5258 store.locks.len(),
5259 0,
5260 "acquiring locks for many distinct ids must not grow the map"
5261 );
5262 }
5263
5264 #[tokio::test]
5265 async fn cancelled_last_waiter_reclaims_hundreds_of_session_locks() {
5266 let (_temp, storage) = make_storage().await;
5267 let store = LockedSessionStore::new(storage);
5268 for index in 0..512 {
5269 let id = format!("cancelled-child-{index}");
5270 let held = store.acquire_lock(&id).await;
5271 let mut waiter = Box::pin(store.acquire_lock(&id));
5272 assert!(
5273 std::future::poll_fn(|cx| std::task::Poll::Ready(
5274 std::future::Future::poll(waiter.as_mut(), cx).is_pending()
5275 ))
5276 .await
5277 );
5278 drop(held);
5279 drop(waiter);
5282 }
5283 assert!(store.locks.is_empty());
5284 }
5285
5286 #[tokio::test]
5287 async fn acquire_lock_concurrent_waiter_keeps_valid_lock_and_map_drains() {
5288 use std::sync::atomic::{AtomicUsize, Ordering};
5289
5290 let (_temp, storage) = make_storage().await;
5291 let store = Arc::new(LockedSessionStore::new(storage));
5292
5293 let active = Arc::new(AtomicUsize::new(0));
5295 let max_seen = Arc::new(AtomicUsize::new(0));
5296
5297 let mut handles = Vec::new();
5298 for _ in 0..8 {
5299 let store = store.clone();
5300 let active = active.clone();
5301 let max_seen = max_seen.clone();
5302 handles.push(tokio::spawn(async move {
5303 let _guard = store.acquire_lock("contended").await;
5304 let now = active.fetch_add(1, Ordering::SeqCst) + 1;
5305 max_seen.fetch_max(now, Ordering::SeqCst);
5306 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
5308 active.fetch_sub(1, Ordering::SeqCst);
5309 }));
5310 }
5311 for h in handles {
5312 h.await.unwrap();
5313 }
5314
5315 assert_eq!(
5320 max_seen.load(Ordering::SeqCst),
5321 1,
5322 "at most one holder of a given session lock at a time"
5323 );
5324 assert_eq!(
5325 store.locks.len(),
5326 0,
5327 "after all holders release, the contended entry must be fully evicted"
5328 );
5329 }
5330}