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 RetrievalWindowCheckpointOutcome, RuntimeSessionPersistence, CONSUMED_CLARIFICATION_IDS_KEY,
43 CONSUMED_RESPONSE_OCCURRENCES_KEY,
44};
45use dashmap::DashMap;
46use tokio::sync::{Mutex, OwnedMutexGuard};
47
48const AUTHORITATIVE_METADATA_KEYS: &[&str] = &["gold_config", "workflow.run_ids.v1"];
49const ROOT_PROJECT_CONTEXT_KEYS: &[&str] = &[
50 "workspace_source",
51 "workspace_binding_status",
52 "project_context_rendered",
53 "project_resources_rendered",
54 "runtime_prompt_snapshot",
55];
56const RESPONSE_CONTROL_METADATA_KEYS: &[&str] = &[
57 CONSUMED_CLARIFICATION_IDS_KEY,
58 CONSUMED_RESPONSE_OCCURRENCES_KEY,
59 "runtime.suspend_reason",
60 "clarification_resume_pending",
61 "conclusion_with_options_resume_pending",
62 "execute.startup_handoff_at",
63 "permission.reexecute_tool_call_id",
64 "permission.reexecute_request_generation",
65 "retry_resume_pending",
66 "retry_resume_reason",
67 "provider_name",
68];
69const TASK_CONTROL_PLANE_CONFLICT_PREFIX: &str = "Task control-plane changed while saving session ";
70const MAX_TASK_CONTROL_PLANE_REBASE_RETRIES: usize = 3;
71const LAST_MANUAL_ARCHIVE_OCCURRENCE_KEY: &str =
72 "context_management.last_manual_archive_occurrence.v1";
73const MANUAL_ARCHIVE_REJECTIONS_KEY: &str = "context_management.manual_archive_rejections.v1";
74const MAX_MANUAL_ARCHIVE_REJECTIONS: usize = 64;
75const RESPONSES_PREVIOUS_RESPONSE_ID_KEY: &str = "responses.previous_response_id";
76
77fn may_publish_runtime_result(result: &std::io::Result<()>) -> bool {
78 !result.as_ref().err().is_some_and(|error| {
79 error
80 .get_ref()
81 .is_some_and(|cause| cause.is::<bamboo_domain::SessionAuthorityConflict>())
82 })
83}
84
85fn adopt_durable_consumed_clarification(session: &mut Session, durable: &Session) -> bool {
94 let Some(incoming_tool_call_id) = session
95 .pending_question
96 .as_ref()
97 .map(|pending| pending.tool_call_id.clone())
98 else {
99 return false;
100 };
101 let occurrence_ledger = durable
102 .metadata
103 .get(CONSUMED_RESPONSE_OCCURRENCES_KEY)
104 .map(|value| serde_json::from_str::<Vec<ResponseOccurrence>>(value));
105 let was_consumed = match occurrence_ledger {
106 Some(Ok(consumed)) => latest_response_occurrence(session, &incoming_tool_call_id)
107 .is_some_and(|incoming| consumed.iter().any(|entry| entry == &incoming)),
108 Some(Err(_)) => false,
111 None => {
112 let legacy_consumed = durable
113 .metadata
114 .get(CONSUMED_CLARIFICATION_IDS_KEY)
115 .and_then(|value| serde_json::from_str::<Vec<String>>(value).ok())
116 .unwrap_or_default()
117 .iter()
118 .any(|tool_call_id| tool_call_id == &incoming_tool_call_id);
119 legacy_consumed
123 && latest_response_occurrence(durable, &incoming_tool_call_id).is_some_and(
124 |durable_occurrence| {
125 latest_response_occurrence(session, &incoming_tool_call_id)
126 .is_some_and(|incoming| incoming == durable_occurrence)
127 },
128 )
129 }
130 };
131 if !was_consumed {
132 return false;
133 }
134
135 bamboo_domain::append_missing_runtime_messages(session, durable);
138 session
139 .pending_question
140 .clone_from(&durable.pending_question);
141 for key in RESPONSE_CONTROL_METADATA_KEYS {
142 if let Some(value) = durable.metadata.get(*key) {
143 session.metadata.insert((*key).to_string(), value.clone());
144 } else {
145 session.metadata.remove(*key);
146 }
147 }
148 session.model.clone_from(&durable.model);
149 session.model_ref.clone_from(&durable.model_ref);
150 session.reasoning_effort = durable.reasoning_effort;
151 session
152 .agent_runtime_state
153 .clone_from(&durable.agent_runtime_state);
154 if let Some(runtime_metadata) = session.runtime_metadata.as_mut() {
155 runtime_metadata.provider_name = durable
156 .runtime_metadata
157 .as_ref()
158 .and_then(|metadata| metadata.provider_name.clone());
159 } else if durable
160 .runtime_metadata
161 .as_ref()
162 .is_some_and(|metadata| metadata.provider_name.is_some())
163 {
164 session.runtime_metadata = durable.runtime_metadata.as_ref().map(|metadata| {
165 let mut response_metadata = bamboo_domain::SessionRuntimeMetadata::default();
166 response_metadata
167 .provider_name
168 .clone_from(&metadata.provider_name);
169 response_metadata
170 });
171 }
172 true
173}
174
175fn is_task_control_plane_save_conflict(error: &std::io::Error) -> bool {
176 error.kind() == std::io::ErrorKind::WouldBlock
177 && error
178 .to_string()
179 .starts_with(TASK_CONTROL_PLANE_CONFLICT_PREFIX)
180}
181
182fn adopt_durable_task_control_plane(session: &mut Session, durable: &Session) {
183 session.task_list = durable.task_list.clone();
184 session
185 .metadata
186 .remove(bamboo_domain::session::runtime_metadata::keys::TASK_LIST_VERSION);
187 if let Some(runtime_metadata) = session.runtime_metadata.as_mut() {
188 runtime_metadata.task_list_version = None;
189 }
190 if session
191 .runtime_metadata
192 .as_ref()
193 .is_some_and(bamboo_domain::session::SessionRuntimeMetadata::is_empty)
194 {
195 session.runtime_metadata = None;
196 }
197 if let Some(version) = durable.task_list_version_meta() {
198 session.set_task_list_version_meta(version);
199 }
200}
201
202fn task_list_snapshot_matches(
203 session: &Session,
204 expected_task_list: &bamboo_domain::TaskList,
205) -> std::io::Result<bool> {
206 Ok(serde_json::to_value(&session.task_list)
207 .map_err(|error| std::io::Error::other(error.to_string()))?
208 == serde_json::to_value(Some(expected_task_list))
209 .map_err(|error| std::io::Error::other(error.to_string()))?)
210}
211
212fn unconditional_task_patch_would_regress(
213 durable: &Session,
214 incoming_task_list: &bamboo_domain::TaskList,
215 incoming_version: &str,
216) -> std::io::Result<bool> {
217 let same_list = task_list_snapshot_matches(durable, incoming_task_list)?;
218 let Some(durable_version) = durable.task_list_version_meta() else {
219 return Ok(false);
220 };
221 match (
222 incoming_version.parse::<u64>(),
223 durable_version.parse::<u64>(),
224 ) {
225 (Ok(incoming), Ok(durable)) => {
226 Ok(incoming < durable || (incoming == durable && !same_list))
227 }
228 _ => Ok(incoming_version != durable_version || !same_list),
229 }
230}
231
232pub struct LockedSessionStore {
240 storage: Arc<dyn Storage>,
241 locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
242 task_pair_transaction_lock: Arc<Mutex<()>>,
246}
247
248pub struct SessionLockGuard {
271 guard: Option<OwnedMutexGuard<()>>,
273 locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
274 session_id: String,
275}
276
277impl Drop for SessionLockGuard {
278 fn drop(&mut self) {
279 self.guard.take();
282 self.locks
283 .remove_if(&self.session_id, |_, arc| Arc::strong_count(arc) == 1);
284 }
285}
286
287impl LockedSessionStore {
288 pub fn new(storage: Arc<dyn Storage>) -> Self {
290 Self {
291 storage,
292 locks: Arc::new(DashMap::new()),
293 task_pair_transaction_lock: Arc::new(Mutex::new(())),
294 }
295 }
296
297 pub fn storage(&self) -> &Arc<dyn Storage> {
299 &self.storage
300 }
301
302 pub async fn acquire_lock(&self, session_id: &str) -> SessionLockGuard {
314 let lock = self
319 .locks
320 .entry(session_id.to_string())
321 .or_insert_with(|| Arc::new(Mutex::new(())))
322 .clone();
323 let mut guard = SessionLockGuard {
327 guard: None,
328 locks: self.locks.clone(),
329 session_id: session_id.to_string(),
330 };
331 guard.guard = Some(lock.lock_owned().await);
332 guard
333 }
334
335 async fn save_session_rebasing_task_conflicts(
340 &self,
341 session: &mut Session,
342 ) -> std::io::Result<()> {
343 for attempt in 0..=MAX_TASK_CONTROL_PLANE_REBASE_RETRIES {
344 match self.storage.save_session(session).await {
345 Ok(()) => return Ok(()),
346 Err(error)
347 if is_task_control_plane_save_conflict(&error)
348 && attempt < MAX_TASK_CONTROL_PLANE_REBASE_RETRIES =>
349 {
350 let Some(durable) =
351 self.storage.load_runtime_control_plane(&session.id).await?
352 else {
353 return Err(error);
354 };
355 adopt_durable_task_control_plane(session, &durable);
356 }
357 Err(error) => {
358 if is_task_control_plane_save_conflict(&error) {
359 if let Some(durable) =
360 self.storage.load_runtime_control_plane(&session.id).await?
361 {
362 adopt_durable_task_control_plane(session, &durable);
363 }
364 }
365 return Err(error);
366 }
367 }
368 }
369 unreachable!("bounded Task conflict retry loop always returns")
370 }
371
372 async fn save_runtime_state_rebasing_task_conflicts(
375 &self,
376 session: &mut Session,
377 ) -> std::io::Result<()> {
378 for attempt in 0..=MAX_TASK_CONTROL_PLANE_REBASE_RETRIES {
379 match self.storage.save_runtime_state(session).await {
380 Ok(()) => return Ok(()),
381 Err(error)
382 if is_task_control_plane_save_conflict(&error)
383 && attempt < MAX_TASK_CONTROL_PLANE_REBASE_RETRIES =>
384 {
385 let Some(durable) =
386 self.storage.load_runtime_control_plane(&session.id).await?
387 else {
388 return Err(error);
389 };
390 adopt_durable_task_control_plane(session, &durable);
391 }
392 Err(error) => {
393 if is_task_control_plane_save_conflict(&error) {
394 if let Some(durable) =
395 self.storage.load_runtime_control_plane(&session.id).await?
396 {
397 adopt_durable_task_control_plane(session, &durable);
398 }
399 }
400 return Err(error);
401 }
402 }
403 }
404 unreachable!("bounded Task conflict retry loop always returns")
405 }
406
407 pub async fn save_runtime_only(&self, session: &mut Session) -> std::io::Result<()> {
425 self.save_runtime_only_and_publish(session, |_| {}).await
426 }
427
428 pub async fn save_runtime_only_and_publish<F>(
439 &self,
440 session: &mut Session,
441 publish: F,
442 ) -> std::io::Result<()>
443 where
444 F: FnOnce(&Session) + Send,
445 {
446 let _guard = self.acquire_lock(&session.id).await;
447 if let Some(latest) = self.storage.load_runtime_control_plane(&session.id).await? {
448 apply_authoritative_metadata(session, &latest);
449 adopt_fresher_disk_permission_posture(session, &latest);
452 adopt_durable_model_context_state(session, &latest);
457 }
458 let result = self
459 .save_runtime_state_rebasing_task_conflicts(session)
460 .await;
461 if may_publish_runtime_result(&result) {
462 publish(session);
463 }
464 result
465 }
466
467 pub async fn update_task_list_control_plane_and_publish<F>(
474 &self,
475 session_id: &str,
476 task_list: &bamboo_domain::TaskList,
477 version: &str,
478 publish: F,
479 ) -> std::io::Result<bool>
480 where
481 F: FnOnce(&Session) + Send,
482 {
483 let _guard = self.acquire_lock(session_id).await;
484 let Some(mut latest) = self.storage.load_runtime_control_plane(session_id).await? else {
485 return Ok(false);
486 };
487 if unconditional_task_patch_would_regress(&latest, task_list, version)? {
488 return Err(std::io::Error::new(
489 std::io::ErrorKind::WouldBlock,
490 format!("Task control-plane changed while patching session {session_id}"),
491 ));
492 }
493 let original = latest.clone();
494 latest.task_list = Some(task_list.clone());
495 latest.set_task_list_version_meta(version.to_string());
496 if !self
497 .storage
498 .save_task_control_plane_if_matches(&original, &latest)
499 .await?
500 {
501 return Err(std::io::Error::new(
502 std::io::ErrorKind::WouldBlock,
503 format!("Task control-plane changed while patching session {session_id}"),
504 ));
505 }
506 publish(&latest);
507 Ok(true)
508 }
509
510 pub async fn update_task_list_control_plane_if_version_and_publish<F>(
514 &self,
515 session_id: &str,
516 expected_version: &str,
517 expected_task_list: &bamboo_domain::TaskList,
518 task_list: &bamboo_domain::TaskList,
519 version: &str,
520 publish: F,
521 ) -> std::io::Result<bool>
522 where
523 F: FnOnce(&Session) + Send,
524 {
525 let _guard = self.acquire_lock(session_id).await;
526 let Some(mut latest) = self.storage.load_runtime_control_plane(session_id).await? else {
527 return Ok(false);
528 };
529 if latest.task_list_version_meta().as_deref() != Some(expected_version)
530 || !task_list_snapshot_matches(&latest, expected_task_list)?
531 {
532 return Ok(false);
533 }
534 let original = latest.clone();
535 latest.task_list = Some(task_list.clone());
536 latest.set_task_list_version_meta(version.to_string());
537 if !self
538 .storage
539 .save_task_control_plane_if_matches(&original, &latest)
540 .await?
541 {
542 return Ok(false);
543 }
544 publish(&latest);
545 Ok(true)
546 }
547
548 pub async fn update_task_list_control_planes_if_version_and_publish<F>(
554 &self,
555 session_id: &str,
556 shared_session_id: &str,
557 expected_version: &str,
558 expected_task_list: &bamboo_domain::TaskList,
559 task_list: &bamboo_domain::TaskList,
560 version: &str,
561 publish: F,
562 ) -> std::io::Result<bool>
563 where
564 F: FnOnce(&Session, &Session) + Send,
565 {
566 if session_id == shared_session_id {
567 return self
568 .update_task_list_control_plane_if_version_and_publish(
569 session_id,
570 expected_version,
571 expected_task_list,
572 task_list,
573 version,
574 |session| publish(session, session),
575 )
576 .await;
577 }
578
579 let (first_id, second_id) = if session_id < shared_session_id {
580 (session_id, shared_session_id)
581 } else {
582 (shared_session_id, session_id)
583 };
584 let _transaction_guard = self.task_pair_transaction_lock.lock().await;
585 let _first_guard = self.acquire_lock(first_id).await;
586 let _second_guard = self.acquire_lock(second_id).await;
587
588 self.storage
592 .recover_task_control_plane_transaction(first_id, second_id)
593 .await?;
594
595 let Some(mut local) = self.storage.load_runtime_control_plane(session_id).await? else {
596 return Ok(false);
597 };
598 let Some(mut shared) = self
599 .storage
600 .load_runtime_control_plane(shared_session_id)
601 .await?
602 else {
603 return Ok(false);
604 };
605 if local.task_list_version_meta().as_deref() != Some(expected_version)
606 || shared.task_list_version_meta().as_deref() != Some(expected_version)
607 || !task_list_snapshot_matches(&local, expected_task_list)?
608 || !task_list_snapshot_matches(&shared, expected_task_list)?
609 {
610 return Ok(false);
611 }
612
613 let local_original = local.clone();
614 let shared_original = shared.clone();
615 local.task_list = Some(task_list.clone());
619 local.set_task_list_version_meta(version.to_string());
620 shared.task_list = Some(task_list.clone());
621 shared.set_task_list_version_meta(version.to_string());
622 let (first_original, first_updated, second_original, second_updated) =
623 if session_id < shared_session_id {
624 (&local_original, &local, &shared_original, &shared)
625 } else {
626 (&shared_original, &shared, &local_original, &local)
627 };
628 let committed = self
629 .storage
630 .save_task_control_planes_atomically(
631 first_original,
632 first_updated,
633 second_original,
634 second_updated,
635 )
636 .await?;
637 if !committed {
638 return Ok(false);
639 }
640 publish(&local, &shared);
641 Ok(true)
642 }
643
644 pub async fn commit_metadata(&self, session: &Session) -> std::io::Result<()> {
654 let _guard = self.acquire_lock(&session.id).await;
655 let mut committed = session.clone();
656 if let Some(latest) = self.storage.load_runtime_control_plane(&session.id).await? {
657 adopt_durable_model_context_state(&mut committed, &latest);
658 }
659 self.save_session_rebasing_task_conflicts(&mut committed)
660 .await
661 }
662
663 pub async fn merge_save_runtime(&self, session: &mut Session) -> std::io::Result<()> {
680 self.merge_save_runtime_and_publish(session, |_, _| {})
681 .await
682 }
683
684 pub async fn merge_save_runtime_and_publish<F>(
692 &self,
693 session: &mut Session,
694 publish: F,
695 ) -> std::io::Result<()>
696 where
697 F: FnOnce(&Session, bool) + Send,
698 {
699 self.merge_save_runtime_inner_and_publish(session, true, publish)
700 .await
701 }
702
703 pub async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
712 self.checkpoint_runtime_session_and_publish(session, |_, _| {})
713 .await
714 }
715
716 pub async fn checkpoint_runtime_session_and_publish<F>(
723 &self,
724 session: &mut Session,
725 publish: F,
726 ) -> std::io::Result<()>
727 where
728 F: FnOnce(&Session, bool) + Send,
729 {
730 let _guard = self.acquire_lock(&session.id).await;
731 let latest = self.storage.load_session(&session.id).await?;
732
733 if let Some(latest) = latest.as_ref() {
734 ensure_model_context_checkpoint_is_current(session, latest)?;
735 let incoming_count = session.messages.len();
736 let durable_count = latest.messages.len();
737 let appended = bamboo_domain::append_missing_runtime_messages(session, latest);
738 bamboo_domain::merge_session_inbox_admission(session, latest);
739 let adopted_response = adopt_durable_consumed_clarification(session, latest);
740 tracing::debug!(
741 "[{}] append-safe runtime checkpoint: durable={}, incoming={}, appended={}, adopted_response={}, saved={}",
742 session.id,
743 durable_count,
744 incoming_count,
745 appended,
746 adopted_response,
747 session.messages.len(),
748 );
749 apply_authoritative_metadata(session, latest);
750 adopt_fresher_disk_permission_posture(session, latest);
751 }
752
753 let result = self.save_session_rebasing_task_conflicts(session).await;
754 if may_publish_runtime_result(&result) {
755 publish(session, result.is_ok());
756 }
757 result
758 }
759
760 pub async fn checkpoint_retrieval_window_and_publish<F>(
766 &self,
767 expected_base: &Session,
768 staged: &mut Session,
769 publish: F,
770 ) -> std::io::Result<RetrievalWindowCheckpointOutcome>
771 where
772 F: FnOnce(&Session) + Send,
773 {
774 validate_staged_retrieval_window_transition(expected_base, staged)?;
775 let _guard = self.acquire_lock(&staged.id).await;
776 let latest = self.storage.load_session(&staged.id).await?;
777
778 if let Some(latest) = latest.as_ref() {
779 if !retrieval_window_base_matches(expected_base, latest)? {
780 *staged = rebase_retrieval_window_base(expected_base, latest);
781 apply_authoritative_metadata(staged, latest);
782 adopt_fresher_disk_permission_posture(staged, latest);
783 return Ok(RetrievalWindowCheckpointOutcome::Rebased);
784 }
785
786 ensure_model_context_checkpoint_is_current(staged, latest)?;
787 bamboo_domain::merge_session_inbox_admission(staged, latest);
788 apply_authoritative_metadata(staged, latest);
789 adopt_fresher_disk_permission_posture(staged, latest);
790 }
791
792 self.save_session_rebasing_task_conflicts(staged).await?;
793 publish(staged);
794 Ok(RetrievalWindowCheckpointOutcome::Committed)
795 }
796
797 pub async fn checkpoint_prompt_rewrite_and_publish<F>(
800 &self,
801 expected_base: &Session,
802 staged: &mut Session,
803 publish: F,
804 ) -> std::io::Result<RetrievalWindowCheckpointOutcome>
805 where
806 F: FnOnce(&Session) + Send,
807 {
808 validate_staged_prompt_rewrite_transition(expected_base, staged)?;
809 let _guard = self.acquire_lock(&staged.id).await;
810 let latest = self.storage.load_session(&staged.id).await?;
811
812 if let Some(latest) = latest.as_ref() {
813 if !retrieval_window_base_matches(expected_base, latest)? {
814 *staged = rebase_retrieval_window_base(expected_base, latest);
815 apply_authoritative_metadata(staged, latest);
816 adopt_fresher_disk_permission_posture(staged, latest);
817 return Ok(RetrievalWindowCheckpointOutcome::Rebased);
818 }
819
820 ensure_model_context_checkpoint_is_current(staged, latest)?;
821 bamboo_domain::merge_session_inbox_admission(staged, latest);
822 apply_authoritative_metadata(staged, latest);
823 adopt_fresher_disk_permission_posture(staged, latest);
824 }
825
826 self.save_session_rebasing_task_conflicts(staged).await?;
827 publish(staged);
828 Ok(RetrievalWindowCheckpointOutcome::Committed)
829 }
830
831 pub async fn checkpoint_manual_archive_rejection_and_publish<F>(
836 &self,
837 expected_base: &Session,
838 staged: &mut Session,
839 publish: F,
840 ) -> std::io::Result<RetrievalWindowCheckpointOutcome>
841 where
842 F: FnOnce(&Session) + Send,
843 {
844 validate_staged_manual_archive_rejection_transition(expected_base, staged)?;
845 let _guard = self.acquire_lock(&staged.id).await;
846 let latest = self.storage.load_session(&staged.id).await?;
847
848 if let Some(latest) = latest.as_ref() {
849 if !retrieval_window_base_matches(expected_base, latest)? {
850 *staged = latest.clone();
851 apply_authoritative_metadata(staged, latest);
852 adopt_fresher_disk_permission_posture(staged, latest);
853 return Ok(RetrievalWindowCheckpointOutcome::Rebased);
854 }
855
856 ensure_model_context_checkpoint_is_current(staged, latest)?;
857 bamboo_domain::merge_session_inbox_admission(staged, latest);
858 apply_authoritative_metadata(staged, latest);
859 adopt_fresher_disk_permission_posture(staged, latest);
860 }
861
862 self.save_session_rebasing_task_conflicts(staged).await?;
863 publish(staged);
864 Ok(RetrievalWindowCheckpointOutcome::Committed)
865 }
866
867 pub async fn checkpoint_manual_archive_consumption_and_publish<F>(
874 &self,
875 expected_base: &Session,
876 staged: &mut Session,
877 publish: F,
878 ) -> std::io::Result<RetrievalWindowCheckpointOutcome>
879 where
880 F: FnOnce(&Session) + Send,
881 {
882 validate_staged_manual_archive_consumption_transition(expected_base, staged)?;
883 let _guard = self.acquire_lock(&staged.id).await;
884 let latest = self.storage.load_session(&staged.id).await?;
885
886 if let Some(latest) = latest.as_ref() {
887 if !retrieval_window_base_matches(expected_base, latest)? {
888 *staged = latest.clone();
889 apply_authoritative_metadata(staged, latest);
890 adopt_fresher_disk_permission_posture(staged, latest);
891 return Ok(RetrievalWindowCheckpointOutcome::Rebased);
892 }
893
894 ensure_model_context_checkpoint_is_current(staged, latest)?;
895 bamboo_domain::merge_session_inbox_admission(staged, latest);
896 apply_authoritative_metadata(staged, latest);
897 adopt_fresher_disk_permission_posture(staged, latest);
898 }
899
900 self.save_session_rebasing_task_conflicts(staged).await?;
901 publish(staged);
902 Ok(RetrievalWindowCheckpointOutcome::Committed)
903 }
904
905 pub async fn save_runtime_authoritative_flags(
914 &self,
915 session: &mut Session,
916 ) -> std::io::Result<()> {
917 self.merge_save_runtime_inner_and_publish(session, false, |_, _| {})
918 .await
919 }
920
921 async fn merge_save_runtime_inner_and_publish<F>(
922 &self,
923 session: &mut Session,
924 adopt_bypass: bool,
925 publish: F,
926 ) -> std::io::Result<()>
927 where
928 F: FnOnce(&Session, bool) + Send,
929 {
930 let _guard = self.acquire_lock(&session.id).await;
931
932 let latest = self.storage.load_session(&session.id).await?;
939
940 let existing_message_count = latest.as_ref().map(|s| s.messages.len());
946 let incoming_message_count = session.messages.len();
947 if existing_message_count.is_some_and(|existing| existing > incoming_message_count) {
948 tracing::warn!(
949 "[{}] merge_save_runtime SHRINK: disk has {:?} messages, saving {} (last_role={:?}, updated_at={}); a stale writer is reverting a concurrent append",
950 session.id,
951 existing_message_count,
952 incoming_message_count,
953 session.messages.last().map(|m| format!("{:?}", m.role)),
954 session.updated_at,
955 );
956 } else {
957 tracing::debug!(
958 "[{}] merge_save_runtime: disk={:?} messages, saving {} (updated_at={})",
959 session.id,
960 existing_message_count,
961 incoming_message_count,
962 session.updated_at,
963 );
964 }
965
966 if let Some(latest) = latest.as_ref() {
967 adopt_durable_consumed_clarification(session, latest);
968 apply_authoritative_metadata(session, latest);
969 let restored = bamboo_domain::restore_missing_admitted_inbox_messages(session, latest);
970 if restored > 0 {
971 tracing::warn!(
972 session_id = %session.id,
973 restored,
974 "restored durable SessionInbox transcript messages into stale runtime save"
975 );
976 }
977 bamboo_domain::merge_session_inbox_admission(session, latest);
978 if adopt_bypass {
983 adopt_fresher_disk_permission_posture(session, latest);
984 }
985 adopt_fresher_durable_model_context_state(session, latest);
986 }
987 let result = self.save_session_rebasing_task_conflicts(session).await;
988 if may_publish_runtime_result(&result) {
989 publish(session, result.is_ok());
990 }
991 result
992 }
993
994 pub async fn seed_runtime_activation_and_publish<F>(
1004 &self,
1005 session: &mut Session,
1006 publish: F,
1007 ) -> std::io::Result<()>
1008 where
1009 F: FnOnce(&Session, bool) + Send,
1010 {
1011 let _guard = self.acquire_lock(&session.id).await;
1012 let mut incoming_audit = PermissionAuditSnapshot::from_metadata(&session.metadata)
1013 .ok_or_else(|| {
1014 std::io::Error::new(
1015 std::io::ErrorKind::InvalidInput,
1016 "activation seed requires a complete permission audit record",
1017 )
1018 })?;
1019
1020 if let Some(latest) = self.storage.load_session(&session.id).await? {
1021 apply_authoritative_metadata(session, &latest);
1022 bamboo_domain::restore_missing_admitted_inbox_messages(session, &latest);
1023 bamboo_domain::merge_session_inbox_admission(session, &latest);
1024 adopt_fresher_durable_model_context_state(session, &latest);
1025
1026 let durable_audit = PermissionAuditSnapshot::from_metadata(&latest.metadata);
1027 let durable_floor = durable_audit
1028 .as_ref()
1029 .map(|snapshot| snapshot.audit_revision)
1030 .unwrap_or_default();
1031 if let Some(durable_audit) = durable_audit {
1032 if durable_audit.resolution == incoming_audit.resolution {
1033 incoming_audit.transitioned_at = durable_audit.transitioned_at;
1034 }
1035 }
1036 incoming_audit.audit_revision = bamboo_domain::next_permission_audit_revision_after(
1037 durable_floor.max(incoming_audit.audit_revision),
1038 )
1039 .map_err(|error| {
1040 std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string())
1041 })?;
1042 }
1043
1044 session
1045 .agent_runtime_state
1046 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
1047 .set_permission_mode(incoming_audit.resolution.requested);
1048 incoming_audit.write_to(&mut session.metadata);
1049
1050 let result = self.save_session_rebasing_task_conflicts(session).await;
1051 if may_publish_runtime_result(&result) {
1052 publish(session, result.is_ok());
1053 }
1054 result
1055 }
1056
1057 pub async fn update_authoritative_permission_posture_and_publish<M, P>(
1063 &self,
1064 session_id: &str,
1065 seed: &PermissionAuditSeed,
1066 mutate: M,
1067 publish: P,
1068 ) -> std::io::Result<Option<Session>>
1069 where
1070 M: FnOnce(&mut Session),
1071 P: FnOnce(&Session),
1072 {
1073 let _guard = self.acquire_lock(session_id).await;
1074 let Some(mut latest) = self.storage.load_session(session_id).await? else {
1075 return Ok(None);
1076 };
1077 let previous_mode = latest
1078 .agent_runtime_state
1079 .as_ref()
1080 .map(|state| state.effective_permission_mode())
1081 .unwrap_or_default();
1082 let previous_resolution = PermissionAuditSnapshot::from_metadata(&latest.metadata)
1083 .map(|snapshot| snapshot.resolution);
1084 mutate(&mut latest);
1085 latest
1086 .agent_runtime_state
1087 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
1088 .set_permission_mode(seed.resolution.requested);
1089 let mode_changed = previous_mode != seed.resolution.requested;
1090 let posture_changed = previous_resolution != Some(seed.resolution);
1091 let transitioned_at = posture_changed.then(|| chrono::Utc::now().to_rfc3339());
1092 bamboo_domain::record_permission_audit(
1093 &mut latest.metadata,
1094 seed,
1095 transitioned_at.as_deref(),
1096 )
1097 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string()))?;
1098 if mode_changed {
1099 latest.metadata_version = latest.metadata_version.saturating_add(1);
1100 }
1101 self.save_session_rebasing_task_conflicts(&mut latest)
1102 .await?;
1103 publish(&latest);
1104 Ok(Some(latest))
1105 }
1106
1107 pub async fn record_permission_posture_activation_and_publish<P>(
1112 &self,
1113 session_id: &str,
1114 expected_audit_revision: Option<u64>,
1115 seed: &PermissionAuditSeed,
1116 publish: P,
1117 ) -> std::io::Result<Option<Session>>
1118 where
1119 P: FnOnce(&Session),
1120 {
1121 let _guard = self.acquire_lock(session_id).await;
1122 let Some(mut latest) = self.storage.load_session(session_id).await? else {
1123 return Ok(None);
1124 };
1125 let durable_audit = PermissionAuditSnapshot::from_metadata(&latest.metadata);
1126 let durable_revision = durable_audit
1127 .as_ref()
1128 .map(|snapshot| snapshot.audit_revision);
1129 if durable_revision != expected_audit_revision {
1130 return Err(std::io::Error::new(
1131 std::io::ErrorKind::InvalidData,
1132 "stale permission posture activation: durable audit changed after dispatch",
1133 ));
1134 }
1135 let durable_requested = latest
1136 .agent_runtime_state
1137 .as_ref()
1138 .map(|state| state.effective_permission_mode())
1139 .unwrap_or_default();
1140 if durable_requested != seed.resolution.requested || !seed.resolution.is_consistent() {
1141 return Err(std::io::Error::new(
1142 std::io::ErrorKind::InvalidData,
1143 "stale or inconsistent permission posture activation",
1144 ));
1145 }
1146 bamboo_domain::record_permission_audit(&mut latest.metadata, seed, None).map_err(
1147 |error| std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string()),
1148 )?;
1149 self.save_session_rebasing_task_conflicts(&mut latest)
1150 .await?;
1151 publish(&latest);
1152 Ok(Some(latest))
1153 }
1154
1155 pub async fn update_runtime_config<F>(
1167 &self,
1168 session_id: &str,
1169 mutate: F,
1170 ) -> std::io::Result<Option<Session>>
1171 where
1172 F: FnOnce(&mut Session),
1173 {
1174 self.update_runtime_config_and_publish(session_id, mutate, |_| {})
1175 .await
1176 }
1177
1178 pub async fn update_runtime_config_and_publish<M, P>(
1181 &self,
1182 session_id: &str,
1183 mutate: M,
1184 publish: P,
1185 ) -> std::io::Result<Option<Session>>
1186 where
1187 M: FnOnce(&mut Session),
1188 P: FnOnce(&Session),
1189 {
1190 let _guard = self.acquire_lock(session_id).await;
1191 let Some(mut session) = self.storage.load_session(session_id).await? else {
1192 return Ok(None);
1193 };
1194 mutate(&mut session);
1195 self.save_session_rebasing_task_conflicts(&mut session)
1196 .await?;
1197 publish(&session);
1198 Ok(Some(session))
1199 }
1200
1201 async fn load_response_candidate<C>(
1206 &self,
1207 session_id: &str,
1208 load_cached: C,
1209 ) -> std::io::Result<Option<Session>>
1210 where
1211 C: FnOnce() -> Option<Session> + Send,
1212 {
1213 let cached_candidate = load_cached();
1214 let durable = self.storage.load_session(session_id).await?;
1215 let Some(mut session) = (match (cached_candidate, durable.as_ref()) {
1216 (Some(cached), Some(durable)) => {
1217 let prefer_durable = durable.updated_at > cached.updated_at
1218 || (durable.updated_at == cached.updated_at
1219 && cached.pending_question.is_none()
1220 && durable.pending_question.is_some());
1221 Some(if prefer_durable {
1222 durable.clone()
1223 } else {
1224 cached
1225 })
1226 }
1227 (Some(cached), None) => Some(cached),
1228 (None, durable) => durable.cloned(),
1229 }) else {
1230 return Ok(None);
1231 };
1232 if let Some(latest) = durable.as_ref() {
1233 adopt_durable_consumed_clarification(&mut session, latest);
1238 bamboo_domain::append_missing_runtime_messages(&mut session, latest);
1239 apply_authoritative_metadata(&mut session, latest);
1240 let restored =
1241 bamboo_domain::restore_missing_admitted_inbox_messages(&mut session, latest);
1242 if restored > 0 {
1243 tracing::warn!(
1244 session_id,
1245 restored,
1246 "restored durable SessionInbox transcript messages into response transaction"
1247 );
1248 }
1249 bamboo_domain::merge_session_inbox_admission(&mut session, latest);
1250 adopt_fresher_disk_permission_posture(&mut session, latest);
1251 adopt_fresher_durable_model_context_state(&mut session, latest);
1252 }
1253 Ok(Some(session))
1254 }
1255
1256 pub async fn inspect_runtime_session_for_response<C>(
1261 &self,
1262 session_id: &str,
1263 load_cached: C,
1264 ) -> std::io::Result<Option<Session>>
1265 where
1266 C: FnOnce() -> Option<Session> + Send,
1267 {
1268 let _guard = self.acquire_lock(session_id).await;
1269 self.load_response_candidate(session_id, load_cached).await
1270 }
1271
1272 pub async fn mutate_runtime_session_and_publish<C, M, P, E>(
1282 &self,
1283 session_id: &str,
1284 load_cached: C,
1285 mutate: M,
1286 publish: P,
1287 ) -> std::io::Result<Result<Option<Session>, E>>
1288 where
1289 C: FnOnce() -> Option<Session> + Send,
1290 M: FnOnce(&mut Session) -> Result<(), E> + Send,
1291 P: FnOnce(&Session) + Send,
1292 E: Send,
1293 {
1294 let _guard = self.acquire_lock(session_id).await;
1295 let Some(mut session) = self
1296 .load_response_candidate(session_id, load_cached)
1297 .await?
1298 else {
1299 return Ok(Ok(None));
1300 };
1301 if let Err(error) = mutate(&mut session) {
1302 return Ok(Err(error));
1303 }
1304 self.save_session_rebasing_task_conflicts(&mut session)
1305 .await?;
1306 publish(&session);
1307 Ok(Ok(Some(session)))
1308 }
1309
1310 pub async fn clear_legacy_pending_messages_and_publish<F>(
1313 &self,
1314 session_id: &str,
1315 expected: &[serde_json::Value],
1316 publish: F,
1317 ) -> std::io::Result<bool>
1318 where
1319 F: FnOnce(&Session) + Send,
1320 {
1321 let _guard = self.acquire_lock(session_id).await;
1322 let Some(mut latest) = self.storage.load_session(session_id).await? else {
1323 return Ok(false);
1324 };
1325 if latest.pending_injected_messages().as_deref() != Some(expected) {
1326 return Ok(false);
1327 }
1328 latest.clear_pending_injected_messages();
1329 self.save_runtime_state_rebasing_task_conflicts(&mut latest)
1330 .await?;
1331 publish(&latest);
1332 Ok(true)
1333 }
1334}
1335
1336#[async_trait::async_trait]
1340impl RuntimeSessionPersistence for LockedSessionStore {
1341 async fn save_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
1342 self.merge_save_runtime(session).await
1343 }
1344
1345 async fn seed_runtime_activation(&self, session: &mut Session) -> std::io::Result<()> {
1346 self.seed_runtime_activation_and_publish(session, |_, _| {})
1347 .await
1348 }
1349
1350 async fn record_permission_posture_activation(
1351 &self,
1352 session_id: &str,
1353 expected_audit_revision: Option<u64>,
1354 seed: &PermissionAuditSeed,
1355 ) -> std::io::Result<Option<Session>> {
1356 self.record_permission_posture_activation_and_publish(
1357 session_id,
1358 expected_audit_revision,
1359 seed,
1360 |_| {},
1361 )
1362 .await
1363 }
1364
1365 async fn save_runtime_control_plane(&self, session: &mut Session) -> std::io::Result<()> {
1366 self.save_runtime_only(session).await
1367 }
1368
1369 async fn load_runtime_control_plane(
1370 &self,
1371 session_id: &str,
1372 ) -> std::io::Result<Option<Session>> {
1373 self.storage.load_runtime_control_plane(session_id).await
1374 }
1375
1376 async fn update_task_list_control_plane(
1377 &self,
1378 session_id: &str,
1379 task_list: &bamboo_domain::TaskList,
1380 version: &str,
1381 ) -> std::io::Result<bool> {
1382 self.update_task_list_control_plane_and_publish(session_id, task_list, version, |_| {})
1383 .await
1384 }
1385
1386 async fn update_task_list_control_plane_if_version(
1387 &self,
1388 session_id: &str,
1389 expected_version: &str,
1390 expected_task_list: &bamboo_domain::TaskList,
1391 task_list: &bamboo_domain::TaskList,
1392 version: &str,
1393 ) -> std::io::Result<bool> {
1394 self.update_task_list_control_plane_if_version_and_publish(
1395 session_id,
1396 expected_version,
1397 expected_task_list,
1398 task_list,
1399 version,
1400 |_| {},
1401 )
1402 .await
1403 }
1404
1405 async fn update_task_list_control_planes_if_version(
1406 &self,
1407 session_id: &str,
1408 shared_session_id: &str,
1409 expected_version: &str,
1410 expected_task_list: &bamboo_domain::TaskList,
1411 task_list: &bamboo_domain::TaskList,
1412 version: &str,
1413 ) -> std::io::Result<bool> {
1414 self.update_task_list_control_planes_if_version_and_publish(
1415 session_id,
1416 shared_session_id,
1417 expected_version,
1418 expected_task_list,
1419 task_list,
1420 version,
1421 |_, _| {},
1422 )
1423 .await
1424 }
1425
1426 async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
1427 LockedSessionStore::checkpoint_runtime_session(self, session).await
1428 }
1429
1430 async fn checkpoint_retrieval_window(
1431 &self,
1432 expected_base: &Session,
1433 staged: &mut Session,
1434 ) -> std::io::Result<RetrievalWindowCheckpointOutcome> {
1435 self.checkpoint_retrieval_window_and_publish(expected_base, staged, |_| {})
1436 .await
1437 }
1438
1439 async fn checkpoint_prompt_rewrite(
1440 &self,
1441 expected_base: &Session,
1442 staged: &mut Session,
1443 ) -> std::io::Result<RetrievalWindowCheckpointOutcome> {
1444 self.checkpoint_prompt_rewrite_and_publish(expected_base, staged, |_| {})
1445 .await
1446 }
1447
1448 async fn checkpoint_manual_archive_rejection(
1449 &self,
1450 expected_base: &Session,
1451 staged: &mut Session,
1452 ) -> std::io::Result<RetrievalWindowCheckpointOutcome> {
1453 self.checkpoint_manual_archive_rejection_and_publish(expected_base, staged, |_| {})
1454 .await
1455 }
1456
1457 async fn checkpoint_manual_archive_consumption(
1458 &self,
1459 expected_base: &Session,
1460 staged: &mut Session,
1461 ) -> std::io::Result<RetrievalWindowCheckpointOutcome> {
1462 self.checkpoint_manual_archive_consumption_and_publish(expected_base, staged, |_| {})
1463 .await
1464 }
1465
1466 async fn load_runtime_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
1467 self.storage.load_session(session_id).await
1468 }
1469
1470 async fn clear_legacy_pending_messages(
1471 &self,
1472 session_id: &str,
1473 expected: &[serde_json::Value],
1474 ) -> std::io::Result<bool> {
1475 self.clear_legacy_pending_messages_and_publish(session_id, expected, |_| {})
1476 .await
1477 }
1478}
1479
1480async fn merge_authoritative_metadata_into_stale(
1489 storage: &Arc<dyn Storage>,
1490 session: &mut Session,
1491) -> std::io::Result<()> {
1492 if let Some(latest) = storage.load_session(&session.id).await? {
1493 adopt_durable_consumed_clarification(session, &latest);
1494 apply_authoritative_metadata(session, &latest);
1495 bamboo_domain::restore_missing_admitted_inbox_messages(session, &latest);
1496 bamboo_domain::merge_session_inbox_admission(session, &latest);
1497 adopt_fresher_disk_permission_posture(session, &latest);
1498 adopt_fresher_durable_model_context_state(session, &latest);
1499 }
1500 Ok(())
1501}
1502
1503fn adopt_durable_model_context_state(session: &mut Session, latest: &Session) {
1508 session
1509 .model_context_state
1510 .clone_from(&latest.model_context_state);
1511}
1512
1513fn adopt_fresher_durable_model_context_state(session: &mut Session, latest: &Session) {
1519 let adopt = match (
1520 session.model_context_state.as_ref(),
1521 latest.model_context_state.as_ref(),
1522 ) {
1523 (None, Some(_)) => true,
1524 (Some(incoming), Some(durable)) => {
1525 durable.state_revision > incoming.state_revision
1526 || (durable.state_revision == incoming.state_revision && durable != incoming)
1527 }
1528 _ => false,
1529 };
1530 if adopt {
1531 adopt_durable_model_context_state(session, latest);
1532 }
1533}
1534
1535fn ensure_model_context_checkpoint_is_current(
1540 session: &Session,
1541 latest: &Session,
1542) -> std::io::Result<()> {
1543 let stale_or_conflicting = match (
1544 session.model_context_state.as_ref(),
1545 latest.model_context_state.as_ref(),
1546 ) {
1547 (None, Some(_)) => true,
1548 (Some(incoming), Some(durable)) => {
1549 durable.state_revision > incoming.state_revision
1550 || (durable.state_revision == incoming.state_revision && durable != incoming)
1551 }
1552 _ => false,
1553 };
1554 if stale_or_conflicting {
1555 return Err(std::io::Error::new(
1556 std::io::ErrorKind::WouldBlock,
1557 "stale or conflicting model-context ledger checkpoint",
1558 ));
1559 }
1560 Ok(())
1561}
1562
1563fn message_matches_retrieval_window_base(
1564 expected: &bamboo_domain::Message,
1565 durable: &bamboo_domain::Message,
1566) -> std::io::Result<bool> {
1567 if durable.image_ocr.is_some() && durable.image_ocr != expected.image_ocr {
1568 return Ok(false);
1569 }
1570 let mut expected = expected.clone();
1571 let mut durable = durable.clone();
1572 expected.image_ocr = None;
1576 durable.image_ocr = None;
1577 let expected = serde_json::to_vec(&expected)
1578 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1579 let durable = serde_json::to_vec(&durable)
1580 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1581 Ok(expected == durable)
1582}
1583
1584fn invalid_retrieval_window_checkpoint(message: impl Into<String>) -> std::io::Error {
1585 std::io::Error::new(std::io::ErrorKind::InvalidInput, message.into())
1586}
1587
1588fn validate_staged_manual_archive_consumption_transition(
1589 expected: &Session,
1590 staged: &Session,
1591) -> std::io::Result<()> {
1592 if expected.id != staged.id || expected.messages.len() != staged.messages.len() {
1593 return Err(invalid_retrieval_window_checkpoint(
1594 "manual archive-consumption checkpoint must preserve the Session ID and message array length",
1595 ));
1596 }
1597
1598 let occurrence = staged
1599 .metadata
1600 .get(LAST_MANUAL_ARCHIVE_OCCURRENCE_KEY)
1601 .ok_or_else(|| {
1602 invalid_retrieval_window_checkpoint(
1603 "manual archive-consumption checkpoint is missing its consumed occurrence",
1604 )
1605 })
1606 .and_then(|value| {
1607 serde_json::from_str::<ResponseOccurrence>(value).map_err(|error| {
1608 invalid_retrieval_window_checkpoint(format!(
1609 "manual archive-consumption checkpoint has an invalid consumed occurrence: {error}"
1610 ))
1611 })
1612 })?;
1613 let result_index = expected
1614 .messages
1615 .iter()
1616 .position(|message| {
1617 message.id == occurrence.tool_result_message_id
1618 && message.tool_call_id.as_deref() == Some(occurrence.tool_call_id.as_str())
1619 && matches!(message.role, bamboo_domain::Role::Tool)
1620 })
1621 .ok_or_else(|| {
1622 invalid_retrieval_window_checkpoint(
1623 "manual archive-consumption checkpoint occurrence does not identify a Tool result",
1624 )
1625 })?;
1626 let correlated_archive_call = expected.messages[..result_index]
1627 .iter()
1628 .rev()
1629 .find(|message| !matches!(message.role, bamboo_domain::Role::Tool))
1630 .filter(|message| matches!(message.role, bamboo_domain::Role::Assistant))
1631 .and_then(|message| message.tool_calls.as_ref())
1632 .into_iter()
1633 .flatten()
1634 .any(|call| {
1635 call.id == occurrence.tool_call_id
1636 && bamboo_domain::canonical_tool_name(&call.function.name) == "archive_context"
1637 });
1638 if !correlated_archive_call {
1639 return Err(invalid_retrieval_window_checkpoint(
1640 "manual archive-consumption checkpoint is not correlated to the current archive_context batch",
1641 ));
1642 }
1643
1644 let mut canonical = expected.clone();
1645 canonical.metadata.insert(
1646 LAST_MANUAL_ARCHIVE_OCCURRENCE_KEY.to_string(),
1647 staged
1648 .metadata
1649 .get(LAST_MANUAL_ARCHIVE_OCCURRENCE_KEY)
1650 .expect("validated occurrence metadata")
1651 .clone(),
1652 );
1653 let canonical = serde_json::to_vec(&canonical)
1654 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1655 let staged = serde_json::to_vec(staged)
1656 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1657 if canonical != staged {
1658 return Err(invalid_retrieval_window_checkpoint(
1659 "manual archive-consumption checkpoint contains mutations outside the correlated consumed marker",
1660 ));
1661 }
1662 Ok(())
1663}
1664
1665fn validate_staged_manual_archive_rejection_transition(
1666 expected: &Session,
1667 staged: &Session,
1668) -> std::io::Result<()> {
1669 if expected.id != staged.id || expected.messages.len() != staged.messages.len() {
1670 return Err(invalid_retrieval_window_checkpoint(
1671 "manual archive-rejection checkpoint must preserve the Session ID and message array length",
1672 ));
1673 }
1674
1675 let occurrence = staged
1676 .metadata
1677 .get(LAST_MANUAL_ARCHIVE_OCCURRENCE_KEY)
1678 .ok_or_else(|| {
1679 invalid_retrieval_window_checkpoint(
1680 "manual archive-rejection checkpoint is missing its consumed occurrence",
1681 )
1682 })
1683 .and_then(|value| {
1684 serde_json::from_str::<ResponseOccurrence>(value).map_err(|error| {
1685 invalid_retrieval_window_checkpoint(format!(
1686 "manual archive-rejection checkpoint has an invalid consumed occurrence: {error}"
1687 ))
1688 })
1689 })?;
1690 let result_index = expected
1691 .messages
1692 .iter()
1693 .position(|message| {
1694 message.id == occurrence.tool_result_message_id
1695 && message.tool_call_id.as_deref() == Some(occurrence.tool_call_id.as_str())
1696 && matches!(message.role, bamboo_domain::Role::Tool)
1697 })
1698 .ok_or_else(|| {
1699 invalid_retrieval_window_checkpoint(
1700 "manual archive-rejection checkpoint occurrence does not identify a Tool result",
1701 )
1702 })?;
1703 let before = &expected.messages[result_index];
1704 let after = &staged.messages[result_index];
1705 let rejection_reason = after
1706 .content
1707 .strip_prefix("archive_context rejected: ")
1708 .filter(|reason| !reason.trim().is_empty())
1709 .ok_or_else(|| {
1710 invalid_retrieval_window_checkpoint(
1711 "manual archive-rejection checkpoint has an invalid Tool result message",
1712 )
1713 })?;
1714 if after.tool_success != Some(false)
1715 || (before.content == after.content && before.tool_success == after.tool_success)
1716 {
1717 return Err(invalid_retrieval_window_checkpoint(
1718 "manual archive-rejection checkpoint did not reject the Tool result",
1719 ));
1720 }
1721
1722 let correlated_archive_call = expected.messages[..result_index]
1723 .iter()
1724 .rev()
1725 .find(|message| !matches!(message.role, bamboo_domain::Role::Tool))
1726 .filter(|message| matches!(message.role, bamboo_domain::Role::Assistant))
1727 .and_then(|message| message.tool_calls.as_ref())
1728 .into_iter()
1729 .flatten()
1730 .any(|call| {
1731 call.id == occurrence.tool_call_id
1732 && bamboo_domain::canonical_tool_name(&call.function.name) == "archive_context"
1733 });
1734 if !correlated_archive_call {
1735 return Err(invalid_retrieval_window_checkpoint(
1736 "manual archive-rejection checkpoint is not correlated to the current archive_context batch",
1737 ));
1738 }
1739
1740 let staged_rejections = staged
1741 .metadata
1742 .get(MANUAL_ARCHIVE_REJECTIONS_KEY)
1743 .ok_or_else(|| {
1744 invalid_retrieval_window_checkpoint(
1745 "manual archive-rejection checkpoint is missing its rejection ledger",
1746 )
1747 })
1748 .and_then(|value| {
1749 serde_json::from_str::<Vec<serde_json::Value>>(value).map_err(|error| {
1750 invalid_retrieval_window_checkpoint(format!(
1751 "manual archive-rejection checkpoint has an invalid rejection ledger: {error}"
1752 ))
1753 })
1754 })?;
1755 let occurrence_value = serde_json::to_value(&occurrence)
1756 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1757 let Some(last_rejection) = staged_rejections.last() else {
1758 return Err(invalid_retrieval_window_checkpoint(
1759 "manual archive-rejection checkpoint has an empty rejection ledger",
1760 ));
1761 };
1762 if last_rejection.get("occurrence") != Some(&occurrence_value)
1763 || last_rejection
1764 .get("reason")
1765 .and_then(serde_json::Value::as_str)
1766 != Some(rejection_reason)
1767 {
1768 return Err(invalid_retrieval_window_checkpoint(
1769 "manual archive-rejection checkpoint ledger does not match its Tool result",
1770 ));
1771 }
1772 let mut expected_rejections = expected
1773 .metadata
1774 .get(MANUAL_ARCHIVE_REJECTIONS_KEY)
1775 .and_then(|value| serde_json::from_str::<Vec<serde_json::Value>>(value).ok())
1776 .unwrap_or_default();
1777 expected_rejections.retain(|rejection| rejection.get("occurrence") != Some(&occurrence_value));
1778 expected_rejections.push(last_rejection.clone());
1779 if expected_rejections.len() > MAX_MANUAL_ARCHIVE_REJECTIONS {
1780 expected_rejections.drain(..expected_rejections.len() - MAX_MANUAL_ARCHIVE_REJECTIONS);
1781 }
1782 if expected_rejections != staged_rejections {
1783 return Err(invalid_retrieval_window_checkpoint(
1784 "manual archive-rejection checkpoint rewrote unrelated rejection-ledger entries",
1785 ));
1786 }
1787
1788 let mut canonical = expected.clone();
1789 canonical.messages[result_index]
1790 .content
1791 .clone_from(&after.content);
1792 canonical.messages[result_index].tool_success = Some(false);
1793 canonical.metadata.insert(
1794 LAST_MANUAL_ARCHIVE_OCCURRENCE_KEY.to_string(),
1795 staged
1796 .metadata
1797 .get(LAST_MANUAL_ARCHIVE_OCCURRENCE_KEY)
1798 .expect("validated occurrence metadata")
1799 .clone(),
1800 );
1801 canonical.metadata.insert(
1802 MANUAL_ARCHIVE_REJECTIONS_KEY.to_string(),
1803 staged
1804 .metadata
1805 .get(MANUAL_ARCHIVE_REJECTIONS_KEY)
1806 .expect("validated rejection metadata")
1807 .clone(),
1808 );
1809 canonical
1810 .metadata
1811 .remove(RESPONSES_PREVIOUS_RESPONSE_ID_KEY);
1812 canonical
1813 .reset_model_context_epoch(bamboo_domain::ModelContextResetReason::ExplicitHistoryRewrite);
1814 let canonical = serde_json::to_vec(&canonical)
1815 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1816 let staged = serde_json::to_vec(staged)
1817 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1818 if canonical != staged {
1819 return Err(invalid_retrieval_window_checkpoint(
1820 "manual archive-rejection checkpoint contains mutations outside the correlated Tool result, bounded metadata, and provider reset",
1821 ));
1822 }
1823 Ok(())
1824}
1825
1826fn validate_staged_prompt_rewrite_transition(
1827 expected: &Session,
1828 staged: &Session,
1829) -> std::io::Result<()> {
1830 if expected.id != staged.id || expected.messages.len() != staged.messages.len() {
1831 return Err(invalid_retrieval_window_checkpoint(
1832 "prompt-rewrite checkpoint must preserve the Session ID and message array length",
1833 ));
1834 }
1835
1836 let mut canonical = expected.clone();
1837 let mut rewrote_system_prompt = false;
1838 for ((before, after), canonical_message) in expected
1839 .messages
1840 .iter()
1841 .zip(&staged.messages)
1842 .zip(&mut canonical.messages)
1843 {
1844 if before.id != after.id {
1845 return Err(invalid_retrieval_window_checkpoint(
1846 "prompt-rewrite checkpoint reordered or replaced a message",
1847 ));
1848 }
1849 if before.content != after.content {
1850 if !matches!(before.role, bamboo_domain::Role::System)
1851 || !matches!(after.role, bamboo_domain::Role::System)
1852 {
1853 return Err(invalid_retrieval_window_checkpoint(
1854 "prompt-rewrite checkpoint changed non-System message content",
1855 ));
1856 }
1857 canonical_message.content.clone_from(&after.content);
1858 rewrote_system_prompt = true;
1859 }
1860 }
1861 if !rewrote_system_prompt {
1862 return Err(invalid_retrieval_window_checkpoint(
1863 "prompt-rewrite checkpoint did not change System message content",
1864 ));
1865 }
1866
1867 canonical.metadata.remove("responses.previous_response_id");
1868 canonical
1869 .reset_model_context_epoch(bamboo_domain::ModelContextResetReason::ExplicitHistoryRewrite);
1870 let canonical = serde_json::to_vec(&canonical)
1871 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1872 let staged = serde_json::to_vec(staged)
1873 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1874 if canonical != staged {
1875 return Err(invalid_retrieval_window_checkpoint(
1876 "prompt-rewrite checkpoint contains mutations outside the System prompt and provider reset",
1877 ));
1878 }
1879 Ok(())
1880}
1881
1882fn validate_staged_retrieval_window_transition(
1883 expected: &Session,
1884 staged: &Session,
1885) -> std::io::Result<()> {
1886 if expected.id != staged.id {
1887 return Err(invalid_retrieval_window_checkpoint(
1888 "retrieval-window checkpoint Session IDs differ",
1889 ));
1890 }
1891 if expected.conversation_summary.is_some() || staged.conversation_summary.is_some() {
1892 return Err(invalid_retrieval_window_checkpoint(
1893 "retrieval-window checkpoint cannot contain a conversation summary",
1894 ));
1895 }
1896 if expected.messages.len() != staged.messages.len() {
1897 return Err(invalid_retrieval_window_checkpoint(
1898 "retrieval-window checkpoint must preserve the message array",
1899 ));
1900 }
1901 if staged.compression_events.len() != expected.compression_events.len().saturating_add(1) {
1902 return Err(invalid_retrieval_window_checkpoint(
1903 "retrieval-window checkpoint must append exactly one compression event",
1904 ));
1905 }
1906
1907 let expected_events = serde_json::to_vec(&expected.compression_events)
1908 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1909 let staged_prefix =
1910 serde_json::to_vec(&staged.compression_events[..expected.compression_events.len()])
1911 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1912 if expected_events != staged_prefix {
1913 return Err(invalid_retrieval_window_checkpoint(
1914 "retrieval-window checkpoint rewrote an existing compression event",
1915 ));
1916 }
1917
1918 let event = staged
1919 .compression_events
1920 .last()
1921 .expect("length check guarantees a staged compression event");
1922 if event.kind != bamboo_domain::CompressionEventKind::RetrievalWindow
1923 || event.id.is_empty()
1924 || expected
1925 .compression_events
1926 .iter()
1927 .any(|existing| existing.id == event.id)
1928 {
1929 return Err(invalid_retrieval_window_checkpoint(
1930 "retrieval-window checkpoint has an invalid archive event",
1931 ));
1932 }
1933
1934 let mut newly_archived = 0usize;
1935 for (before, after) in expected.messages.iter().zip(&staged.messages) {
1936 if before.id != after.id {
1937 return Err(invalid_retrieval_window_checkpoint(
1938 "retrieval-window checkpoint reordered or replaced a message",
1939 ));
1940 }
1941
1942 let is_new_archive = !before.compressed && after.compressed;
1943 if is_new_archive {
1944 if before.compressed_by_event_id.is_some()
1945 || after.compressed_by_event_id.as_deref() != Some(event.id.as_str())
1946 {
1947 return Err(invalid_retrieval_window_checkpoint(
1948 "retrieval-window checkpoint has an invalid message correlation",
1949 ));
1950 }
1951 newly_archived = newly_archived.saturating_add(1);
1952 } else if before.compressed != after.compressed
1953 || before.compressed_by_event_id != after.compressed_by_event_id
1954 {
1955 return Err(invalid_retrieval_window_checkpoint(
1956 "retrieval-window checkpoint contains an unsupported archive mutation",
1957 ));
1958 }
1959
1960 let mut normalized_after = after.clone();
1961 normalized_after.compressed = before.compressed;
1962 normalized_after
1963 .compressed_by_event_id
1964 .clone_from(&before.compressed_by_event_id);
1965 let before = serde_json::to_vec(before)
1966 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1967 let after = serde_json::to_vec(&normalized_after)
1968 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
1969 if before != after {
1970 return Err(invalid_retrieval_window_checkpoint(
1971 "retrieval-window checkpoint mutated message content",
1972 ));
1973 }
1974 }
1975
1976 if newly_archived == 0 || newly_archived != event.messages_compressed {
1977 return Err(invalid_retrieval_window_checkpoint(
1978 "retrieval-window checkpoint event count does not match message correlations",
1979 ));
1980 }
1981 if staged
1982 .model_context_state
1983 .as_ref()
1984 .and_then(|state| state.last_reset_reason)
1985 != Some(bamboo_domain::ModelContextResetReason::Compression)
1986 {
1987 return Err(invalid_retrieval_window_checkpoint(
1988 "retrieval-window checkpoint is missing its model-context reset",
1989 ));
1990 }
1991
1992 Ok(())
1993}
1994
1995fn retrieval_window_base_matches(expected: &Session, durable: &Session) -> std::io::Result<bool> {
1996 if expected.id != durable.id || durable.messages.len() > expected.messages.len() {
1997 return Ok(false);
1998 }
1999 for (durable_message, expected_message) in durable.messages.iter().zip(expected.messages.iter())
2000 {
2001 if !message_matches_retrieval_window_base(expected_message, durable_message)? {
2002 return Ok(false);
2003 }
2004 }
2005
2006 if durable.compression_events.len() > expected.compression_events.len() {
2007 return Ok(false);
2008 }
2009 for (durable_event, expected_event) in durable
2010 .compression_events
2011 .iter()
2012 .zip(expected.compression_events.iter())
2013 {
2014 let durable_event = serde_json::to_vec(durable_event)
2015 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
2016 let expected_event = serde_json::to_vec(expected_event)
2017 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
2018 if durable_event != expected_event {
2019 return Ok(false);
2020 }
2021 }
2022
2023 let summary_matches = serde_json::to_vec(&expected.conversation_summary)
2024 .and_then(|expected| {
2025 serde_json::to_vec(&durable.conversation_summary).map(|durable| expected == durable)
2026 })
2027 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
2028 if !summary_matches {
2029 return Ok(false);
2030 }
2031
2032 if expected.metadata != durable.metadata
2039 || expected.runtime_metadata != durable.runtime_metadata
2040 || expected.metadata_version != durable.metadata_version
2041 || expected.model != durable.model
2042 || expected.model_ref != durable.model_ref
2043 || expected.reasoning_effort != durable.reasoning_effort
2044 {
2045 return Ok(false);
2046 }
2047
2048 Ok(ensure_model_context_checkpoint_is_current(expected, durable).is_ok())
2049}
2050
2051fn rebase_retrieval_window_base(expected: &Session, durable: &Session) -> Session {
2052 let mut rebased = expected.clone();
2053 bamboo_domain::append_missing_runtime_messages(&mut rebased, durable);
2054 rebased.metadata.clone_from(&durable.metadata);
2055 rebased
2056 .runtime_metadata
2057 .clone_from(&durable.runtime_metadata);
2058 rebased.model.clone_from(&durable.model);
2059 rebased.model_ref.clone_from(&durable.model_ref);
2060 rebased.reasoning_effort = durable.reasoning_effort;
2061 bamboo_domain::merge_session_inbox_admission(&mut rebased, durable);
2062 rebased
2063 .conversation_summary
2064 .clone_from(&durable.conversation_summary);
2065 rebased
2066 .compression_events
2067 .clone_from(&durable.compression_events);
2068 rebased.token_usage.clone_from(&durable.token_usage);
2069 rebased
2070 .model_context_state
2071 .clone_from(&durable.model_context_state);
2072 if durable.updated_at > rebased.updated_at {
2073 rebased.updated_at = durable.updated_at;
2074 }
2075 rebased
2076}
2077
2078fn adopt_fresher_disk_permission_posture(session: &mut Session, latest: &Session) {
2090 let Some(disk_mode) = latest
2095 .agent_runtime_state
2096 .as_ref()
2097 .map(|state| state.effective_permission_mode())
2098 else {
2099 return;
2100 };
2101 let current_mode = session
2102 .agent_runtime_state
2103 .as_ref()
2104 .map(|state| state.effective_permission_mode())
2105 .unwrap_or_default();
2106 let Some(disk_audit) = bamboo_domain::fresher_disk_permission_audit(
2107 current_mode,
2108 &session.metadata,
2109 disk_mode,
2110 &latest.metadata,
2111 ) else {
2112 return;
2113 };
2114
2115 match session.agent_runtime_state.as_mut() {
2116 Some(state) => state.set_permission_mode(disk_mode),
2117 None if disk_mode != bamboo_domain::SessionPermissionMode::Default => {
2120 let state = session
2121 .agent_runtime_state
2122 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default);
2123 state.set_permission_mode(disk_mode);
2124 }
2125 None => {}
2126 }
2127
2128 disk_audit.write_to(&mut session.metadata);
2130}
2131
2132fn apply_authoritative_metadata(session: &mut Session, latest: &Session) {
2138 if session.authority_identity.is_ordinary() && session.created_at == latest.created_at {
2145 session.authority_identity = latest.authority_identity.clone();
2146 }
2147 if session.kind == bamboo_domain::SessionKind::Root
2151 && latest.kind == bamboo_domain::SessionKind::Root
2152 && session.created_at == latest.created_at
2153 && session.authority_identity == latest.authority_identity
2154 {
2155 session.supervisor_management = latest.supervisor_management.clone();
2156 }
2157 if session.kind == bamboo_domain::SessionKind::Root
2163 && latest.kind == bamboo_domain::SessionKind::Root
2164 && session.created_at == latest.created_at
2165 && latest.metadata_version >= session.metadata_version
2166 && (latest.metadata_version > session.metadata_version
2167 || latest.project_id_meta() != session.project_id_meta())
2168 {
2169 match latest.project_id_meta() {
2170 Some(project) => session.set_project_id_meta(project),
2171 None => session.clear_project_id_meta(),
2172 }
2173 match latest.workspace_path_meta() {
2174 Some(workspace) => session.set_workspace_path_meta(workspace),
2175 None => {
2176 session.metadata.remove("workspace_path");
2177 if let Some(metadata) = session.runtime_metadata.as_mut() {
2178 metadata.workspace_path = None;
2179 }
2180 }
2181 }
2182 session.workspace.clone_from(&latest.workspace);
2183 for key in ROOT_PROJECT_CONTEXT_KEYS {
2184 match latest.metadata.get(*key) {
2185 Some(value) => {
2186 session.metadata.insert((*key).to_string(), value.clone());
2187 }
2188 None => {
2189 session.metadata.remove(*key);
2190 }
2191 }
2192 }
2193 session.prompt_snapshot.clone_from(&latest.prompt_snapshot);
2194 }
2195 if latest.metadata_version >= session.metadata_version {
2196 session.title = latest.title.clone();
2197 session.title_version = latest.title_version;
2198 session.title_generated = latest.title_generated;
2199 session.pinned = latest.pinned;
2200 for key in AUTHORITATIVE_METADATA_KEYS {
2201 if let Some(value) = latest.metadata.get(*key) {
2202 session.metadata.insert((*key).to_string(), value.clone());
2203 } else {
2204 session.metadata.remove(*key);
2205 }
2206 }
2207 session.metadata_version = latest.metadata_version;
2208 }
2209}
2210
2211pub async fn merge_save_session(
2224 storage: &Arc<dyn Storage>,
2225 session: &mut Session,
2226) -> std::io::Result<()> {
2227 merge_authoritative_metadata_into_stale(storage, session).await?;
2228 storage.save_session(session).await
2229}
2230
2231#[cfg(test)]
2234mod tests {
2235 use super::*;
2236 use crate::v2::{RuntimeTaskTransactionFault, SessionStoreV2};
2237 use bamboo_domain::{session::types::Session, PermissionMode};
2238 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
2239
2240 #[test]
2241 fn authority_merge_updates_ordinary_cache_but_does_not_hide_stale_incarnation() {
2242 let mut latest = Session::new(bamboo_domain::DEFAULT_SUPERVISOR_SESSION_ID, "model");
2243 latest.authority_identity = bamboo_domain::SessionAuthorityIdentity::Supervisor {
2244 incarnation_id: uuid::Uuid::new_v4(),
2245 };
2246 let mut stale = latest.clone();
2247 stale.authority_identity = bamboo_domain::SessionAuthorityIdentity::Ordinary;
2248 stale.metadata_version = 100;
2249 apply_authoritative_metadata(&mut stale, &latest);
2250 assert_eq!(stale.authority_identity, latest.authority_identity);
2251 let old_identity = bamboo_domain::SessionAuthorityIdentity::Supervisor {
2252 incarnation_id: uuid::Uuid::new_v4(),
2253 };
2254 stale.authority_identity = old_identity.clone();
2255 apply_authoritative_metadata(&mut stale, &latest);
2256 assert_eq!(stale.authority_identity, old_identity);
2257 }
2258
2259 struct AuthoritySavePauseStorage {
2260 inner: Arc<SessionStoreV2>,
2261 reached: tokio::sync::Barrier,
2262 release: tokio::sync::Barrier,
2263 }
2264
2265 #[tokio::test]
2266 async fn root_project_merge_publishes_project_revision_and_workspace_together() {
2267 for runtime_only in [false, true] {
2268 for equal_revision in [false, true] {
2269 let temp = tempfile::tempdir().unwrap();
2270 let storage = Arc::new(SessionStoreV2::new(temp.path().into()).await.unwrap());
2271 let mut stale = Session::new("root-project-merge", "model");
2272 stale.set_project_id_meta("project-a");
2273 stale.set_workspace_path_meta("/project-a");
2274 stale.metadata.insert(
2275 "runtime_prompt_snapshot".into(),
2276 "old Project A prompt".into(),
2277 );
2278 stale.add_message(bamboo_domain::Message::user("Keep transcript"));
2279 storage.save_session(&stale).await.unwrap();
2280 let mut current = stale.clone();
2281 current.metadata_version += 1;
2282 current.set_project_id_meta("project-b");
2283 current.set_workspace_path_meta("/project-b");
2284 current.metadata.remove("runtime_prompt_snapshot");
2285 current
2286 .metadata
2287 .insert("workspace_source".into(), "project_default".into());
2288 current
2289 .metadata
2290 .insert("project_context_rendered".into(), "Project B".into());
2291 storage.save_session(¤t).await.unwrap();
2292 if equal_revision {
2293 stale.metadata_version = current.metadata_version;
2296 }
2297 let locked = LockedSessionStore::new(storage.clone());
2298 let published = AtomicBool::new(false);
2299 let publish = |saved: &Session| {
2300 assert!(!saved.metadata.contains_key("runtime_prompt_snapshot"));
2301 assert_eq!(saved.project_id_meta().as_deref(), Some("project-b"));
2302 assert_eq!(saved.workspace_path_meta().as_deref(), Some("/project-b"));
2303 assert_eq!(saved.metadata_version, current.metadata_version);
2304 assert_eq!(
2305 saved.metadata.get("workspace_source").map(String::as_str),
2306 Some("project_default")
2307 );
2308 assert_eq!(
2309 saved
2310 .metadata
2311 .get("project_context_rendered")
2312 .map(String::as_str),
2313 Some("Project B")
2314 );
2315 published.store(true, Ordering::SeqCst);
2316 };
2317 if runtime_only {
2318 locked
2319 .save_runtime_only_and_publish(&mut stale, publish)
2320 .await
2321 .unwrap();
2322 } else {
2323 locked
2324 .merge_save_runtime_and_publish(&mut stale, |saved, committed| {
2325 assert!(committed);
2326 publish(saved);
2327 })
2328 .await
2329 .unwrap();
2330 }
2331 assert!(published.load(Ordering::SeqCst));
2332 assert_eq!(stale.project_id_meta(), current.project_id_meta());
2333 let loaded = storage.load_session(&stale.id).await.unwrap().unwrap();
2334 assert_eq!(loaded.project_id_meta(), current.project_id_meta());
2335 assert_eq!(loaded.messages.len(), 1);
2336 storage.flush_search_index().await;
2337 }
2338 }
2339 }
2340
2341 #[tokio::test]
2342 async fn root_project_change_after_merge_read_rejects_without_cache_publication() {
2343 for runtime_only in [false, true] {
2344 let temp = tempfile::tempdir().unwrap();
2345 let first = Arc::new(SessionStoreV2::new(temp.path().into()).await.unwrap());
2346 let mut stale = Session::new("root-project-race", "model");
2347 stale.set_project_id_meta("project-a");
2348 stale.add_message(bamboo_domain::Message::user("Keep transcript"));
2349 first.save_session(&stale).await.unwrap();
2350 let second = SessionStoreV2::new(temp.path().into()).await.unwrap();
2351 let mut current = stale.clone();
2352 current.metadata_version += 1;
2353 current.set_project_id_meta("project-b");
2354 let paused = Arc::new(AuthoritySavePauseStorage {
2355 inner: first.clone(),
2356 reached: tokio::sync::Barrier::new(2),
2357 release: tokio::sync::Barrier::new(2),
2358 });
2359 let locked = LockedSessionStore::new(paused.clone());
2360 let published = AtomicBool::new(false);
2361 let save = async {
2362 if runtime_only {
2363 locked
2364 .save_runtime_only_and_publish(&mut stale, |_| {
2365 published.store(true, Ordering::SeqCst);
2366 })
2367 .await
2368 } else {
2369 locked
2370 .merge_save_runtime_and_publish(&mut stale, |_, _| {
2371 published.store(true, Ordering::SeqCst);
2372 })
2373 .await
2374 }
2375 };
2376 let update = async {
2377 paused.reached.wait().await;
2378 second.save_session(¤t).await.unwrap();
2379 paused.release.wait().await;
2380 };
2381 let (result, ()) = tokio::time::timeout(std::time::Duration::from_secs(10), async {
2382 tokio::join!(save, update)
2383 })
2384 .await
2385 .expect("deterministic Project/save race completes");
2386 assert!(!may_publish_runtime_result(&Err(result.unwrap_err())));
2387 assert!(!published.load(Ordering::SeqCst));
2388 let loaded = first.load_session(¤t.id).await.unwrap().unwrap();
2389 assert_eq!(loaded.project_id_meta(), current.project_id_meta());
2390 assert_eq!(loaded.metadata_version, current.metadata_version);
2391 assert_eq!(loaded.messages.len(), 1);
2392 first.flush_search_index().await;
2393 second.flush_search_index().await;
2394 }
2395 }
2396
2397 #[async_trait::async_trait]
2398 impl Storage for AuthoritySavePauseStorage {
2399 async fn load_session(&self, id: &str) -> std::io::Result<Option<Session>> {
2400 self.inner.load_session(id).await
2401 }
2402 async fn load_runtime_control_plane(&self, id: &str) -> std::io::Result<Option<Session>> {
2403 self.inner.load_runtime_control_plane(id).await
2404 }
2405 async fn delete_session(&self, id: &str) -> std::io::Result<bool> {
2406 self.inner.delete_session(id).await
2407 }
2408 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
2409 self.reached.wait().await;
2410 self.release.wait().await;
2411 self.inner.save_session(session).await
2412 }
2413 async fn save_runtime_state(&self, session: &Session) -> std::io::Result<()> {
2414 self.reached.wait().await;
2415 self.release.wait().await;
2416 self.inner.save_runtime_state(session).await
2417 }
2418 }
2419
2420 #[tokio::test]
2421 async fn root_deleted_after_merge_read_rejects_without_cache_or_event_publication() {
2422 for runtime_only in [false, true] {
2423 let temp = tempfile::tempdir().unwrap();
2424 let first = Arc::new(SessionStoreV2::new(temp.path().into()).await.unwrap());
2425 let mut stale = Session::new("root-delete-race", "model");
2426 stale.set_project_id_meta("project-a");
2427 stale.add_message(bamboo_domain::Message::user("Old lifetime"));
2428 first.save_session(&stale).await.unwrap();
2429 let id = stale.id.clone();
2430 let second = SessionStoreV2::new(temp.path().into()).await.unwrap();
2431 let paused = Arc::new(AuthoritySavePauseStorage {
2432 inner: first.clone(),
2433 reached: tokio::sync::Barrier::new(2),
2434 release: tokio::sync::Barrier::new(2),
2435 });
2436 let locked = LockedSessionStore::new(paused.clone());
2437 let published = AtomicBool::new(false);
2438 let save = async {
2439 if runtime_only {
2440 locked
2441 .save_runtime_only_and_publish(&mut stale, |_| {
2442 published.store(true, Ordering::SeqCst);
2443 })
2444 .await
2445 } else {
2446 locked
2447 .merge_save_runtime_and_publish(&mut stale, |_, _| {
2448 published.store(true, Ordering::SeqCst);
2449 })
2450 .await
2451 }
2452 };
2453 let delete = async {
2454 paused.reached.wait().await;
2455 assert!(second.delete_session(&id).await.unwrap());
2456 paused.release.wait().await;
2457 };
2458 let (result, ()) = tokio::time::timeout(std::time::Duration::from_secs(10), async {
2459 tokio::join!(save, delete)
2460 })
2461 .await
2462 .expect("deterministic delete/save race completes");
2463 assert!(!may_publish_runtime_result(&Err(result.unwrap_err())));
2464 assert!(!published.load(Ordering::SeqCst));
2465 assert!(first.load_root_authority(&id).await.unwrap().is_none());
2466 assert!(!first.sessions_root_dir().join(&id).exists());
2467 first.flush_search_index().await;
2468 second.flush_search_index().await;
2469 }
2470 }
2471
2472 #[tokio::test]
2473 async fn supervisor_bootstrap_between_merge_read_and_save_rejects_without_publishing() {
2474 for runtime_only in [false, true] {
2475 let temp = tempfile::tempdir().unwrap();
2476 let first = Arc::new(SessionStoreV2::new(temp.path().into()).await.unwrap());
2477 let second = SessionStoreV2::new(temp.path().into()).await.unwrap();
2478 let paused = Arc::new(AuthoritySavePauseStorage {
2479 inner: first.clone(),
2480 reached: tokio::sync::Barrier::new(2),
2481 release: tokio::sync::Barrier::new(2),
2482 });
2483 let store = LockedSessionStore::new(paused.clone());
2484 let mut stale = Session::new(bamboo_domain::DEFAULT_SUPERVISOR_SESSION_ID, "stale");
2485 stale.add_message(bamboo_domain::Message::user(
2486 "must not enter new Supervisor",
2487 ));
2488 let published = AtomicBool::new(false);
2489 let save = async {
2490 if runtime_only {
2491 store
2492 .save_runtime_only_and_publish(&mut stale, |_| {
2493 published.store(true, Ordering::SeqCst);
2494 })
2495 .await
2496 } else {
2497 store
2498 .merge_save_runtime_and_publish(&mut stale, |_, _| {
2499 published.store(true, Ordering::SeqCst);
2500 })
2501 .await
2502 }
2503 };
2504 let bootstrap = async {
2505 paused.reached.wait().await;
2508 let receipt = second
2509 .get_or_create_default_supervisor("supervisor")
2510 .await
2511 .unwrap();
2512 paused.release.wait().await;
2513 receipt
2514 };
2515 let (result, receipt) =
2516 tokio::time::timeout(std::time::Duration::from_secs(10), async {
2517 tokio::join!(save, bootstrap)
2518 })
2519 .await
2520 .expect("deterministic bootstrap/save race completes");
2521 let error = result.unwrap_err();
2522 assert!(!may_publish_runtime_result(&Err(error)));
2523 assert!(!published.load(Ordering::SeqCst));
2524 assert!(stale.authority_identity.is_ordinary());
2525 let observed = first
2526 .load_root_authority(&receipt.session_id)
2527 .await
2528 .unwrap()
2529 .unwrap();
2530 let durable = second
2533 .load_session(&receipt.session_id)
2534 .await
2535 .unwrap()
2536 .unwrap();
2537 assert_eq!(durable.model, "supervisor");
2538 assert_eq!(observed.authority_identity, durable.authority_identity);
2539 assert!(durable.messages.is_empty());
2540 assert_eq!(
2541 durable.authority_identity,
2542 bamboo_domain::SessionAuthorityIdentity::Supervisor {
2543 incarnation_id: receipt.incarnation_id,
2544 }
2545 );
2546 }
2547 }
2548
2549 #[tokio::test]
2550 async fn supervisor_merge_publishes_adopted_identity_but_never_a_rejected_incarnation() {
2551 for runtime_only in [false, true] {
2552 let temp = tempfile::tempdir().unwrap();
2553 let storage = Arc::new(SessionStoreV2::new(temp.path().into()).await.unwrap());
2554 let receipt = storage
2555 .get_or_create_default_supervisor("model")
2556 .await
2557 .unwrap();
2558 let baseline = storage
2559 .load_session(&receipt.session_id)
2560 .await
2561 .unwrap()
2562 .unwrap();
2563 let expected = baseline.authority_identity.clone();
2564 let store = LockedSessionStore::new(storage.clone());
2565 for case in 0..3 {
2566 let rejected = case != 0;
2567 let mut snapshot = baseline.clone();
2568 snapshot.authority_identity = if case == 1 {
2569 bamboo_domain::SessionAuthorityIdentity::Supervisor {
2570 incarnation_id: uuid::Uuid::new_v4(),
2571 }
2572 } else {
2573 bamboo_domain::SessionAuthorityIdentity::Ordinary
2574 };
2575 if case == 2 {
2576 snapshot.created_at -= chrono::Duration::seconds(1);
2577 snapshot.model = "stale Ordinary instance".into();
2578 snapshot.add_message(bamboo_domain::Message::user("must not be rebound"));
2579 }
2580 let published = AtomicBool::new(false);
2581 let callback = |saved: &Session| {
2582 assert_eq!(saved.authority_identity, expected);
2583 published.store(true, Ordering::SeqCst);
2584 };
2585 let result = if runtime_only {
2586 store
2587 .save_runtime_only_and_publish(&mut snapshot, callback)
2588 .await
2589 } else {
2590 store
2591 .merge_save_runtime_and_publish(&mut snapshot, |saved, committed| {
2592 assert!(committed);
2593 callback(saved);
2594 })
2595 .await
2596 };
2597 assert_eq!(result.is_err(), rejected);
2598 assert_eq!(published.load(Ordering::SeqCst), !rejected);
2599 if !rejected {
2600 assert_eq!(snapshot.authority_identity, expected);
2601 }
2602 let durable = storage
2603 .load_session(&receipt.session_id)
2604 .await
2605 .unwrap()
2606 .unwrap();
2607 assert_eq!(durable.authority_identity, expected);
2608 assert_eq!(durable.created_at, baseline.created_at);
2609 assert_eq!(durable.model, baseline.model);
2610 assert!(durable.messages.is_empty());
2611 }
2612 }
2613 }
2614
2615 struct CountingControlPlaneStorage {
2616 inner: Arc<SessionStoreV2>,
2617 control_plane_loads: AtomicUsize,
2618 full_saves: AtomicUsize,
2619 runtime_state_saves: AtomicUsize,
2620 }
2621
2622 struct PairCommitBarrierStorage {
2623 inner: Arc<SessionStoreV2>,
2624 before_commit: Arc<tokio::sync::Barrier>,
2625 }
2626
2627 struct SingleCommitBarrierStorage {
2628 inner: Arc<SessionStoreV2>,
2629 before_commit: Arc<tokio::sync::Barrier>,
2630 }
2631
2632 struct SingleCommitPauseStorage {
2633 inner: Arc<SessionStoreV2>,
2634 commit_reached: Arc<tokio::sync::Barrier>,
2635 release_commit: Arc<tokio::sync::Barrier>,
2636 }
2637
2638 #[async_trait::async_trait]
2639 impl Storage for SingleCommitPauseStorage {
2640 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
2641 self.inner.save_session(session).await
2642 }
2643
2644 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
2645 self.inner.load_session(session_id).await
2646 }
2647
2648 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
2649 self.inner.delete_session(session_id).await
2650 }
2651
2652 async fn load_runtime_control_plane(
2653 &self,
2654 session_id: &str,
2655 ) -> std::io::Result<Option<Session>> {
2656 self.inner.load_runtime_control_plane(session_id).await
2657 }
2658
2659 async fn save_task_control_plane_if_matches(
2660 &self,
2661 original: &Session,
2662 updated: &Session,
2663 ) -> std::io::Result<bool> {
2664 self.commit_reached.wait().await;
2665 self.release_commit.wait().await;
2666 self.inner
2667 .save_task_control_plane_if_matches(original, updated)
2668 .await
2669 }
2670
2671 async fn save_task_control_planes_atomically(
2672 &self,
2673 first_original: &Session,
2674 first_updated: &Session,
2675 second_original: &Session,
2676 second_updated: &Session,
2677 ) -> std::io::Result<bool> {
2678 self.commit_reached.wait().await;
2679 self.release_commit.wait().await;
2680 self.inner
2681 .save_task_control_planes_atomically(
2682 first_original,
2683 first_updated,
2684 second_original,
2685 second_updated,
2686 )
2687 .await
2688 }
2689 }
2690
2691 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2692 async fn supervisor_management_race_rejects_staged_task_callbacks_and_fresh_invocation_succeeds(
2693 ) {
2694 use bamboo_domain::{
2695 SupervisorManagementMutation, SupervisorManagementRequest, SupervisorReference,
2696 };
2697
2698 for mode in 0..4 {
2699 let home = tempfile::tempdir().unwrap();
2700 let inner = Arc::new(SessionStoreV2::new(home.path().into()).await.unwrap());
2701 let receipt = inner
2702 .get_or_create_default_supervisor("model")
2703 .await
2704 .unwrap();
2705 let reference = SupervisorReference::from(&receipt);
2706 let mut root = inner
2707 .load_session(&reference.session_id)
2708 .await
2709 .unwrap()
2710 .unwrap();
2711 let list = bamboo_domain::TaskList {
2712 session_id: root.id.clone(),
2713 title: "original".into(),
2714 items: vec![],
2715 created_at: root.created_at,
2716 updated_at: root.created_at,
2717 };
2718 let updated = bamboo_domain::TaskList {
2719 title: "updated".into(),
2720 ..list.clone()
2721 };
2722 root.task_list = Some(list.clone());
2723 root.set_task_list_version_meta("1");
2724 inner.save_session(&root).await.unwrap();
2725 let child_id = if mode == 2 { "aaa-child" } else { "zzz-child" };
2726 let mut child = Session::new_child_of(child_id, &root, "model", "child");
2727 child.task_list = Some(list.clone());
2728 child.set_task_list_version_meta("1");
2729 inner.save_session(&child).await.unwrap();
2730 let independent = SessionStoreV2::new(home.path().into()).await.unwrap();
2731 let commit_reached = Arc::new(tokio::sync::Barrier::new(2));
2732 let release_commit = Arc::new(tokio::sync::Barrier::new(2));
2733 let paused = LockedSessionStore::new(Arc::new(SingleCommitPauseStorage {
2734 inner: inner.clone(),
2735 commit_reached: commit_reached.clone(),
2736 release_commit: release_commit.clone(),
2737 }));
2738 let published = AtomicBool::new(false);
2739 let loser = async {
2740 if mode == 0 {
2741 paused
2742 .update_task_list_control_plane_and_publish(&root.id, &updated, "2", |_| {
2743 published.store(true, Ordering::SeqCst)
2744 })
2745 .await
2746 } else if mode == 1 {
2747 paused
2748 .update_task_list_control_plane_if_version_and_publish(
2749 &root.id,
2750 "1",
2751 &list,
2752 &updated,
2753 "2",
2754 |_| published.store(true, Ordering::SeqCst),
2755 )
2756 .await
2757 } else {
2758 paused
2759 .update_task_list_control_planes_if_version_and_publish(
2760 child_id,
2761 &root.id,
2762 "1",
2763 &list,
2764 &updated,
2765 "2",
2766 |_, _| published.store(true, Ordering::SeqCst),
2767 )
2768 .await
2769 }
2770 };
2771 let winner = async {
2772 commit_reached.wait().await;
2773 independent
2774 .mutate_supervisor_management(&SupervisorManagementRequest {
2775 supervisor: reference.clone(),
2776 expected_state_revision: 0,
2777 mutation: SupervisorManagementMutation::ConfigureProjectScope {
2778 allowed_projects: ["project-a".parse().unwrap()].into(),
2779 },
2780 })
2781 .await
2782 .unwrap();
2783 release_commit.wait().await;
2784 };
2785 let (result, ()) = tokio::join!(loser, winner);
2786 if mode == 0 {
2787 assert_eq!(result.unwrap_err().kind(), std::io::ErrorKind::WouldBlock);
2788 } else {
2789 assert!(!result.unwrap());
2790 }
2791 assert!(!published.load(Ordering::SeqCst));
2792 for id in [&root.id, &child.id] {
2793 assert_eq!(
2794 inner
2795 .load_runtime_control_plane(id)
2796 .await
2797 .unwrap()
2798 .unwrap()
2799 .task_list_version_meta()
2800 .as_deref(),
2801 Some("1")
2802 );
2803 }
2804 let fresh = LockedSessionStore::new(inner.clone());
2805 let check = |saved: &Session| {
2806 assert_eq!(saved.supervisor_management.as_ref().unwrap().revision, 1);
2807 published.store(true, Ordering::SeqCst);
2808 };
2809 let retried = if mode == 0 {
2810 fresh
2811 .update_task_list_control_plane_and_publish(&root.id, &updated, "2", check)
2812 .await
2813 } else if mode == 1 {
2814 fresh
2815 .update_task_list_control_plane_if_version_and_publish(
2816 &root.id, "1", &list, &updated, "2", check,
2817 )
2818 .await
2819 } else {
2820 fresh
2821 .update_task_list_control_planes_if_version_and_publish(
2822 child_id,
2823 &root.id,
2824 "1",
2825 &list,
2826 &updated,
2827 "2",
2828 |_, shared| check(shared),
2829 )
2830 .await
2831 };
2832 assert!(retried.unwrap());
2833 assert!(published.load(Ordering::SeqCst));
2834 assert_eq!(
2835 inner
2836 .load_runtime_control_plane(&root.id)
2837 .await
2838 .unwrap()
2839 .unwrap()
2840 .task_list_version_meta()
2841 .as_deref(),
2842 Some("2")
2843 );
2844 }
2845 }
2846
2847 #[async_trait::async_trait]
2848 impl Storage for SingleCommitBarrierStorage {
2849 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
2850 self.inner.save_session(session).await
2851 }
2852
2853 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
2854 self.inner.load_session(session_id).await
2855 }
2856
2857 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
2858 self.inner.delete_session(session_id).await
2859 }
2860
2861 async fn load_runtime_control_plane(
2862 &self,
2863 session_id: &str,
2864 ) -> std::io::Result<Option<Session>> {
2865 self.inner.load_runtime_control_plane(session_id).await
2866 }
2867
2868 async fn save_task_control_plane_if_matches(
2869 &self,
2870 original: &Session,
2871 updated: &Session,
2872 ) -> std::io::Result<bool> {
2873 self.before_commit.wait().await;
2874 self.inner
2875 .save_task_control_plane_if_matches(original, updated)
2876 .await
2877 }
2878 }
2879
2880 #[async_trait::async_trait]
2881 impl Storage for PairCommitBarrierStorage {
2882 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
2883 self.inner.save_session(session).await
2884 }
2885
2886 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
2887 self.inner.load_session(session_id).await
2888 }
2889
2890 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
2891 self.inner.delete_session(session_id).await
2892 }
2893
2894 async fn save_runtime_state(&self, session: &Session) -> std::io::Result<()> {
2895 self.inner.save_runtime_state(session).await
2896 }
2897
2898 async fn load_runtime_control_plane(
2899 &self,
2900 session_id: &str,
2901 ) -> std::io::Result<Option<Session>> {
2902 self.inner.load_runtime_control_plane(session_id).await
2903 }
2904
2905 async fn recover_task_control_plane_transaction(
2906 &self,
2907 first_session_id: &str,
2908 second_session_id: &str,
2909 ) -> std::io::Result<()> {
2910 self.inner
2911 .recover_task_control_plane_transaction(first_session_id, second_session_id)
2912 .await
2913 }
2914
2915 async fn save_task_control_planes_atomically(
2916 &self,
2917 first_original: &Session,
2918 first_updated: &Session,
2919 second_original: &Session,
2920 second_updated: &Session,
2921 ) -> std::io::Result<bool> {
2922 self.before_commit.wait().await;
2927 self.inner
2928 .save_task_control_planes_atomically(
2929 first_original,
2930 first_updated,
2931 second_original,
2932 second_updated,
2933 )
2934 .await
2935 }
2936 }
2937
2938 #[async_trait::async_trait]
2939 impl Storage for CountingControlPlaneStorage {
2940 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
2941 self.full_saves.fetch_add(1, Ordering::SeqCst);
2942 self.inner.save_session(session).await
2943 }
2944
2945 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
2946 self.inner.load_session(session_id).await
2947 }
2948
2949 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
2950 self.inner.delete_session(session_id).await
2951 }
2952
2953 async fn save_runtime_state(&self, session: &Session) -> std::io::Result<()> {
2954 self.runtime_state_saves.fetch_add(1, Ordering::SeqCst);
2955 self.inner.save_runtime_state(session).await
2956 }
2957
2958 async fn load_runtime_control_plane(
2959 &self,
2960 session_id: &str,
2961 ) -> std::io::Result<Option<Session>> {
2962 self.control_plane_loads.fetch_add(1, Ordering::SeqCst);
2963 self.inner.load_runtime_control_plane(session_id).await
2964 }
2965
2966 async fn recover_task_control_plane_transaction(
2967 &self,
2968 first_session_id: &str,
2969 second_session_id: &str,
2970 ) -> std::io::Result<()> {
2971 self.inner
2972 .recover_task_control_plane_transaction(first_session_id, second_session_id)
2973 .await
2974 }
2975
2976 async fn save_task_control_plane_if_matches(
2977 &self,
2978 original: &Session,
2979 updated: &Session,
2980 ) -> std::io::Result<bool> {
2981 let committed = self
2982 .inner
2983 .save_task_control_plane_if_matches(original, updated)
2984 .await?;
2985 if committed {
2986 self.runtime_state_saves.fetch_add(1, Ordering::SeqCst);
2987 }
2988 Ok(committed)
2989 }
2990
2991 async fn save_task_control_planes_atomically(
2992 &self,
2993 first_original: &Session,
2994 first_updated: &Session,
2995 second_original: &Session,
2996 second_updated: &Session,
2997 ) -> std::io::Result<bool> {
2998 let committed = self
2999 .inner
3000 .save_task_control_planes_atomically(
3001 first_original,
3002 first_updated,
3003 second_original,
3004 second_updated,
3005 )
3006 .await?;
3007 if committed {
3008 self.runtime_state_saves.fetch_add(2, Ordering::SeqCst);
3009 }
3010 Ok(committed)
3011 }
3012 }
3013
3014 async fn make_storage() -> (tempfile::TempDir, Arc<dyn Storage>) {
3015 let temp = tempfile::tempdir().unwrap();
3016 let storage = SessionStoreV2::new(temp.path().to_path_buf())
3017 .await
3018 .expect("storage init");
3019 (temp, Arc::new(storage) as Arc<dyn Storage>)
3020 }
3021
3022 fn fresh(id: &str) -> Session {
3023 Session::new(id.to_string(), "test-model".to_string())
3024 }
3025
3026 fn stage_retrieval_window_archive(expected: &Session, message_index: usize) -> Session {
3027 let mut staged = expected.clone();
3028 let mut event = bamboo_domain::CompressionEvent::new(
3029 1,
3030 1,
3031 80.0,
3032 60.0,
3033 0,
3034 bamboo_domain::CompressionTriggerType::Auto,
3035 0.0,
3036 None,
3037 0,
3038 );
3039 event.kind = bamboo_domain::CompressionEventKind::RetrievalWindow;
3040 let event_id = event.id.clone();
3041 staged.messages[message_index].compressed = true;
3042 staged.messages[message_index].compressed_by_event_id = Some(event_id);
3043 staged.compression_events.push(event);
3044 staged.reset_model_context_epoch(bamboo_domain::ModelContextResetReason::Compression);
3045 staged
3046 }
3047
3048 fn typed_permission_result(
3049 tool_call_id: &str,
3050 message_id: &str,
3051 generation: &str,
3052 content: &str,
3053 ) -> bamboo_domain::session::types::Message {
3054 let mut message =
3055 bamboo_domain::session::types::Message::tool_result(tool_call_id, content);
3056 message.id = message_id.to_string();
3057 message.metadata = Some(serde_json::json!({
3058 "permission_request": {
3059 "request_generation": generation,
3060 }
3061 }));
3062 message
3063 }
3064
3065 fn ledger_state(state_revision: u64, marker: &str) -> bamboo_domain::ModelContextState {
3066 bamboo_domain::ModelContextState {
3067 state_revision,
3068 prefix_epoch: state_revision,
3069 cache_scope_sha256: Some("scope".to_string()),
3070 transcript_item_sha256: vec![marker.to_string()],
3071 ..bamboo_domain::ModelContextState::default()
3072 }
3073 }
3074
3075 fn set_permission_audit(
3076 session: &mut Session,
3077 requested: bamboo_domain::SessionPermissionMode,
3078 policy_revision: u64,
3079 mapping: &str,
3080 transitioned_at: &str,
3081 ) -> u64 {
3082 let resolution = bamboo_domain::resolve_permission_mode(requested, PermissionMode::Default);
3083 session
3084 .agent_runtime_state
3085 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
3086 .set_permission_mode(requested);
3087 bamboo_domain::record_permission_audit(
3088 &mut session.metadata,
3089 &PermissionAuditSeed::new(policy_revision, resolution, mapping),
3090 Some(transitioned_at),
3091 )
3092 .unwrap()
3093 }
3094
3095 #[tokio::test]
3096 async fn checked_runtime_mutation_persists_a_cache_only_pending_session() {
3097 let (_temp, storage) = make_storage().await;
3098 let store = LockedSessionStore::new(storage.clone());
3099 let mut cached = fresh("cache-only-response");
3100 cached.set_pending_question(
3101 "tool-1".to_string(),
3102 "ConclusionWithOptions".to_string(),
3103 "Choose".to_string(),
3104 vec!["A".to_string()],
3105 false,
3106 );
3107
3108 let saved = store
3109 .mutate_runtime_session_and_publish(
3110 &cached.id.clone(),
3111 move || Some(cached),
3112 |session| {
3113 assert!(session.pending_question.is_some());
3114 session.clear_pending_question();
3115 Ok::<_, ()>(())
3116 },
3117 |_| {},
3118 )
3119 .await
3120 .unwrap()
3121 .unwrap()
3122 .expect("cache-only session should be created durably");
3123
3124 assert!(saved.pending_question.is_none());
3125 assert!(storage
3126 .load_session("cache-only-response")
3127 .await
3128 .unwrap()
3129 .unwrap()
3130 .pending_question
3131 .is_none());
3132 }
3133
3134 #[tokio::test]
3135 async fn checked_runtime_mutation_preserves_durable_authorities_for_newer_cache() {
3136 let (_temp, storage) = make_storage().await;
3137 let store = LockedSessionStore::new(storage.clone());
3138 let session_id = "cached-response-authorities";
3139 let mut durable = fresh(session_id);
3140 durable.title = "Durable title".to_string();
3141 durable.title_version = 4;
3142 durable.title_generated = true;
3143 durable.metadata_version = 9;
3144 set_permission_audit(
3145 &mut durable,
3146 bamboo_domain::SessionPermissionMode::Auto,
3147 7,
3148 "bamboo_runtime:durable-auto",
3149 "2026-08-10T09:00:00Z",
3150 );
3151 storage.save_session(&durable).await.unwrap();
3152
3153 let mut cached = fresh(session_id);
3154 cached.created_at = durable.created_at;
3155 cached.title = "Stale cached title".to_string();
3156 cached.updated_at = durable.updated_at + chrono::Duration::seconds(1);
3157 cached.set_pending_question(
3158 "tool-1".to_string(),
3159 "ConclusionWithOptions".to_string(),
3160 "Choose".to_string(),
3161 vec!["A".to_string()],
3162 false,
3163 );
3164
3165 store
3166 .mutate_runtime_session_and_publish(
3167 session_id,
3168 move || Some(cached),
3169 |session| {
3170 session.clear_pending_question();
3171 Ok::<_, ()>(())
3172 },
3173 |_| {},
3174 )
3175 .await
3176 .unwrap()
3177 .unwrap()
3178 .expect("session should exist");
3179
3180 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3181 assert_eq!(saved.title, "Durable title");
3182 assert_eq!(saved.title_version, 4);
3183 assert_eq!(saved.metadata_version, 9);
3184 assert_eq!(
3185 saved
3186 .agent_runtime_state
3187 .as_ref()
3188 .unwrap()
3189 .effective_permission_mode(),
3190 bamboo_domain::SessionPermissionMode::Auto
3191 );
3192 let audit = bamboo_domain::PermissionAuditSnapshot::from_metadata(&saved.metadata).unwrap();
3193 assert_eq!(audit.policy_revision, 7);
3194 assert_eq!(audit.executor_mapping, "bamboo_runtime:durable-auto");
3195 }
3196
3197 #[tokio::test]
3198 async fn checked_runtime_mutation_cannot_resurrect_consumed_ask_from_newer_cache() {
3199 use bamboo_domain::session::types::Message;
3200
3201 let (_temp, storage) = make_storage().await;
3202 let store = LockedSessionStore::new(storage.clone());
3203 let session_id = "cached-consumed-response";
3204 let mut stale_cached = fresh(session_id);
3205 stale_cached.add_message(Message::tool_result("call-1", "waiting"));
3206 stale_cached.set_pending_question(
3207 "call-1".to_string(),
3208 "ConclusionWithOptions".to_string(),
3209 "Choose".to_string(),
3210 vec!["A".to_string()],
3211 false,
3212 );
3213
3214 let mut durable = stale_cached.clone();
3215 durable.clear_pending_question();
3216 durable.metadata.insert(
3217 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
3218 r#"["call-1"]"#.to_string(),
3219 );
3220 durable.messages[0].content = "Selected response: A".to_string();
3221 durable.add_message(Message::user("durable concurrent message"));
3222 storage.save_session(&durable).await.unwrap();
3223
3224 stale_cached.updated_at = durable.updated_at + chrono::Duration::seconds(1);
3227 let saved = store
3228 .mutate_runtime_session_and_publish(
3229 session_id,
3230 move || Some(stale_cached),
3231 |session| {
3232 assert!(session.pending_question.is_none());
3233 Ok::<_, ()>(())
3234 },
3235 |_| {},
3236 )
3237 .await
3238 .unwrap()
3239 .unwrap()
3240 .unwrap();
3241
3242 assert!(saved.pending_question.is_none());
3243 assert_eq!(saved.messages.len(), 2);
3244 assert_eq!(saved.messages[0].content, "Selected response: A");
3245 assert_eq!(saved.messages[1].content, "durable concurrent message");
3246 }
3247
3248 #[tokio::test]
3249 async fn response_inspection_adopts_durable_consumption_without_writing() {
3250 let (_temp, storage) = make_storage().await;
3251 let store = LockedSessionStore::new(storage.clone());
3252 let session_id = "inspect-consumed-response";
3253 let mut stale_cached = fresh(session_id);
3254 stale_cached.add_message(bamboo_domain::session::types::Message::tool_result(
3259 "call-1", "waiting",
3260 ));
3261 stale_cached.set_pending_question(
3262 "call-1".to_string(),
3263 "ConclusionWithOptions".to_string(),
3264 "Choose".to_string(),
3265 vec!["A".to_string()],
3266 false,
3267 );
3268
3269 let mut durable = stale_cached.clone();
3270 durable.clear_pending_question();
3271 durable.metadata.insert(
3272 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
3273 r#"["call-1"]"#.to_string(),
3274 );
3275 storage.save_session(&durable).await.unwrap();
3276 stale_cached.updated_at = durable.updated_at + chrono::Duration::seconds(1);
3277
3278 let inspected = store
3279 .inspect_runtime_session_for_response(session_id, move || Some(stale_cached))
3280 .await
3281 .unwrap()
3282 .expect("session should be inspectable");
3283 assert!(inspected.pending_question.is_none());
3284
3285 let unchanged = storage.load_session(session_id).await.unwrap().unwrap();
3286 assert_eq!(unchanged.updated_at, durable.updated_at);
3287 assert!(unchanged.pending_question.is_none());
3288 }
3289
3290 #[tokio::test]
3291 async fn same_mode_newer_run_start_audit_survives_every_runtime_save_path() {
3292 for path in ["merge", "checkpoint", "control-plane"] {
3293 let (_temp, storage) = make_storage().await;
3294 let store = LockedSessionStore::new(storage.clone());
3295 let session_id = format!("same-mode-newer-{path}");
3296 let mut durable = fresh(&session_id);
3297 set_permission_audit(
3298 &mut durable,
3299 bamboo_domain::SessionPermissionMode::Default,
3300 1,
3301 "bamboo_runtime:old-policy",
3302 "2026-07-31T12:00:00Z",
3303 );
3304 storage.save_session(&durable).await.unwrap();
3305
3306 let mut run_start = durable.clone();
3307 let old_revision = PermissionAuditSnapshot::from_metadata(&durable.metadata)
3308 .unwrap()
3309 .audit_revision;
3310 let new_revision = set_permission_audit(
3311 &mut run_start,
3312 bamboo_domain::SessionPermissionMode::Default,
3313 2,
3314 "bamboo_runtime:new-policy",
3315 "2026-07-31T12:00:00Z",
3316 );
3317 assert!(new_revision > old_revision);
3318
3319 match path {
3320 "merge" => store.merge_save_runtime(&mut run_start).await.unwrap(),
3321 "checkpoint" => store
3322 .checkpoint_runtime_session(&mut run_start)
3323 .await
3324 .unwrap(),
3325 "control-plane" => store.save_runtime_only(&mut run_start).await.unwrap(),
3326 _ => unreachable!(),
3327 }
3328
3329 let saved = storage.load_session(&session_id).await.unwrap().unwrap();
3330 let audit = PermissionAuditSnapshot::from_metadata(&saved.metadata).unwrap();
3331 assert_eq!(audit.audit_revision, new_revision, "path={path}");
3332 assert_eq!(audit.policy_revision, 2, "path={path}");
3333 assert_eq!(audit.executor_mapping, "bamboo_runtime:new-policy");
3334 }
3335 }
3336
3337 #[tokio::test]
3338 async fn newer_disk_transition_wins_after_mode_cycles_back() {
3339 let (_temp, storage) = make_storage().await;
3340 let store = LockedSessionStore::new(storage.clone());
3341 let session_id = "permission-cycle-back";
3342 let mut baseline = fresh(session_id);
3343 let stale_revision = set_permission_audit(
3344 &mut baseline,
3345 bamboo_domain::SessionPermissionMode::Default,
3346 1,
3347 "bamboo_runtime:initial",
3348 "2026-07-31T12:00:00Z",
3349 );
3350 storage.save_session(&baseline).await.unwrap();
3351 let mut stale_runtime = baseline.clone();
3352
3353 let mut durable = baseline;
3354 set_permission_audit(
3355 &mut durable,
3356 bamboo_domain::SessionPermissionMode::Auto,
3357 2,
3358 "bamboo_runtime:auto",
3359 "2026-07-31T12:01:00Z",
3360 );
3361 let durable_revision = set_permission_audit(
3362 &mut durable,
3363 bamboo_domain::SessionPermissionMode::Default,
3364 3,
3365 "bamboo_runtime:cycled-default",
3366 "2026-07-31T12:02:00Z",
3367 );
3368 assert!(durable_revision > stale_revision);
3369 storage.save_session(&durable).await.unwrap();
3370
3371 store.merge_save_runtime(&mut stale_runtime).await.unwrap();
3372 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3373 let audit = PermissionAuditSnapshot::from_metadata(&saved.metadata).unwrap();
3374 assert_eq!(audit.audit_revision, durable_revision);
3375 assert_eq!(audit.policy_revision, 3);
3376 assert_eq!(audit.executor_mapping, "bamboo_runtime:cycled-default");
3377 }
3378
3379 #[tokio::test]
3380 async fn authoritative_activation_seed_replaces_every_warm_worker_posture() {
3381 let (_temp, storage) = make_storage().await;
3382 let store = LockedSessionStore::new(storage.clone());
3383 let session_id = "warm-permission-matrix";
3384 let cases = [
3385 (
3386 bamboo_domain::SessionPermissionMode::Auto,
3387 PermissionMode::Default,
3388 PermissionMode::Auto,
3389 ),
3390 (
3391 bamboo_domain::SessionPermissionMode::Default,
3392 PermissionMode::Default,
3393 PermissionMode::Default,
3394 ),
3395 (
3396 bamboo_domain::SessionPermissionMode::Auto,
3397 PermissionMode::Default,
3398 PermissionMode::Auto,
3399 ),
3400 (
3401 bamboo_domain::SessionPermissionMode::Bypass,
3402 PermissionMode::Auto,
3403 PermissionMode::BypassPermissions,
3404 ),
3405 ];
3406 let mut previous_revision = 0;
3407 let created_at = fresh(session_id).created_at;
3408
3409 for (index, (requested, configured, expected_effective)) in cases.into_iter().enumerate() {
3410 let mut activation = fresh(session_id);
3411 activation.created_at = created_at;
3412 activation
3413 .agent_runtime_state
3414 .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
3415 .set_permission_mode(requested);
3416 let resolution = bamboo_domain::resolve_permission_mode(requested, configured);
3417 bamboo_domain::record_permission_audit(
3418 &mut activation.metadata,
3419 &PermissionAuditSeed::new(
3420 index as u64 + 1,
3421 resolution,
3422 format!("bamboo_worker:{}", resolution.effective.as_str()),
3423 ),
3424 Some("2026-07-31T12:00:00Z"),
3425 )
3426 .unwrap();
3427
3428 RuntimeSessionPersistence::seed_runtime_activation(&store, &mut activation)
3429 .await
3430 .unwrap();
3431 let durable = storage.load_session(session_id).await.unwrap().unwrap();
3432 assert_eq!(
3433 durable
3434 .agent_runtime_state
3435 .as_ref()
3436 .unwrap()
3437 .effective_permission_mode(),
3438 requested,
3439 "activation {index}"
3440 );
3441 let audit = PermissionAuditSnapshot::from_metadata(&durable.metadata).unwrap();
3442 assert_eq!(audit.resolution.requested, requested);
3443 assert_eq!(audit.resolution.effective, expected_effective);
3444 assert!(audit.audit_revision > previous_revision);
3445 previous_revision = audit.audit_revision;
3446 }
3447 }
3448
3449 #[tokio::test]
3450 async fn resident_reseed_bumps_etag_only_for_typed_transition() {
3451 let (_temp, storage) = make_storage().await;
3452 let store = LockedSessionStore::new(storage.clone());
3453 let session_id = "resident-atomic-permission";
3454 let mut baseline = fresh(session_id);
3455 baseline.metadata_version = 7;
3456 set_permission_audit(
3457 &mut baseline,
3458 bamboo_domain::SessionPermissionMode::Auto,
3459 1,
3460 "bamboo_runtime:auto",
3461 "2026-07-31T12:00:00Z",
3462 );
3463 storage.save_session(&baseline).await.unwrap();
3464 let initial_audit = PermissionAuditSnapshot::from_metadata(&baseline.metadata).unwrap();
3465
3466 let same_mode_seed = PermissionAuditSeed::bamboo_runtime(
3467 2,
3468 bamboo_domain::resolve_permission_mode(
3469 bamboo_domain::SessionPermissionMode::Auto,
3470 PermissionMode::Default,
3471 ),
3472 );
3473 let refreshed = store
3474 .update_authoritative_permission_posture_and_publish(
3475 session_id,
3476 &same_mode_seed,
3477 |session| {
3478 session
3479 .metadata
3480 .insert("resident.marker".to_string(), "same-mode".to_string());
3481 },
3482 |_| {},
3483 )
3484 .await
3485 .unwrap()
3486 .unwrap();
3487 let refreshed_audit = PermissionAuditSnapshot::from_metadata(&refreshed.metadata).unwrap();
3488 assert_eq!(refreshed.metadata_version, 7);
3489 assert!(refreshed_audit.audit_revision > initial_audit.audit_revision);
3490 assert_eq!(refreshed_audit.policy_revision, 2);
3491
3492 let transition_seed = PermissionAuditSeed::bamboo_runtime(
3493 3,
3494 bamboo_domain::resolve_permission_mode(
3495 bamboo_domain::SessionPermissionMode::Default,
3496 PermissionMode::Default,
3497 ),
3498 );
3499 let transitioned = store
3500 .update_authoritative_permission_posture_and_publish(
3501 session_id,
3502 &transition_seed,
3503 |session| {
3504 session
3505 .metadata
3506 .insert("resident.marker".to_string(), "transition".to_string());
3507 },
3508 |_| {},
3509 )
3510 .await
3511 .unwrap()
3512 .unwrap();
3513 let transitioned_audit =
3514 PermissionAuditSnapshot::from_metadata(&transitioned.metadata).unwrap();
3515 assert_eq!(transitioned.metadata_version, 8, "old ETag must be invalid");
3516 assert_eq!(
3517 transitioned
3518 .agent_runtime_state
3519 .as_ref()
3520 .unwrap()
3521 .effective_permission_mode(),
3522 bamboo_domain::SessionPermissionMode::Default
3523 );
3524 assert_eq!(
3525 transitioned_audit.resolution.requested,
3526 bamboo_domain::SessionPermissionMode::Default
3527 );
3528 assert!(transitioned_audit.audit_revision > refreshed_audit.audit_revision);
3529 assert_eq!(
3530 transitioned
3531 .metadata
3532 .get("resident.marker")
3533 .map(String::as_str),
3534 Some("transition")
3535 );
3536 }
3537
3538 #[tokio::test]
3539 async fn worker_activation_cas_cannot_overwrite_concurrent_permission_patch() {
3540 let (_temp, storage) = make_storage().await;
3541 let store = LockedSessionStore::new(storage.clone());
3542 let session_id = "permission-activation-cas";
3543 let mut baseline = fresh(session_id);
3544 set_permission_audit(
3545 &mut baseline,
3546 bamboo_domain::SessionPermissionMode::Default,
3547 1,
3548 "bamboo_runtime:default",
3549 "2026-07-31T12:00:00Z",
3550 );
3551 storage.save_session(&baseline).await.unwrap();
3552 let dispatched_revision = PermissionAuditSnapshot::from_metadata(&baseline.metadata)
3553 .unwrap()
3554 .audit_revision;
3555
3556 let patched_resolution = bamboo_domain::resolve_permission_mode(
3557 bamboo_domain::SessionPermissionMode::Auto,
3558 PermissionMode::Default,
3559 );
3560 let patched = store
3561 .update_authoritative_permission_posture_and_publish(
3562 session_id,
3563 &PermissionAuditSeed::new(2, patched_resolution, "patch:auto"),
3564 |_| {},
3565 |_| {},
3566 )
3567 .await
3568 .unwrap()
3569 .unwrap();
3570 let patched_audit = PermissionAuditSnapshot::from_metadata(&patched.metadata).unwrap();
3571 assert!(patched_audit.audit_revision > dispatched_revision);
3572
3573 let stale_worker_seed = PermissionAuditSeed::new(
3574 1,
3575 bamboo_domain::resolve_permission_mode(
3576 bamboo_domain::SessionPermissionMode::Default,
3577 PermissionMode::Default,
3578 ),
3579 "worker:stale-default",
3580 );
3581 let error = store
3582 .record_permission_posture_activation_and_publish(
3583 session_id,
3584 Some(dispatched_revision),
3585 &stale_worker_seed,
3586 |_| {},
3587 )
3588 .await
3589 .unwrap_err();
3590 assert!(error.to_string().contains("durable audit changed"));
3591
3592 let durable = storage.load_session(session_id).await.unwrap().unwrap();
3593 let durable_audit = PermissionAuditSnapshot::from_metadata(&durable.metadata).unwrap();
3594 assert_eq!(durable_audit, patched_audit);
3595 assert_eq!(durable_audit.executor_mapping, "patch:auto");
3596 }
3597
3598 #[tokio::test]
3601 async fn update_runtime_config_preserves_concurrently_appended_messages() {
3602 use bamboo_domain::session::types::Message;
3603 use bamboo_domain::ReasoningEffort;
3604
3605 let (_temp, storage) = make_storage().await;
3606 let store = LockedSessionStore::new(storage.clone());
3607 let session_id = "cfg-preserve";
3608
3609 let mut initial = fresh(session_id);
3611 initial.add_message(Message::user("hello"));
3612 initial.add_message(Message::assistant("hi", None));
3613 storage.save_session(&initial).await.unwrap();
3614
3615 let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
3617 after_chat.add_message(Message::user("second question"));
3618 storage.save_session(&after_chat).await.unwrap();
3619 assert_eq!(after_chat.messages.len(), 3);
3620
3621 let updated = store
3625 .update_runtime_config(session_id, |s| {
3626 s.reasoning_effort = Some(ReasoningEffort::Max);
3627 })
3628 .await
3629 .unwrap()
3630 .expect("session exists");
3631
3632 assert_eq!(updated.reasoning_effort, Some(ReasoningEffort::Max));
3633 assert_eq!(
3634 updated.messages.len(),
3635 3,
3636 "config patch must not revert a concurrently-appended message"
3637 );
3638
3639 let on_disk = storage.load_session(session_id).await.unwrap().unwrap();
3640 assert_eq!(on_disk.messages.len(), 3);
3641 assert_eq!(on_disk.reasoning_effort, Some(ReasoningEffort::Max));
3642 }
3643
3644 #[tokio::test]
3645 async fn update_runtime_config_returns_none_for_missing_session() {
3646 use bamboo_domain::ReasoningEffort;
3647
3648 let (_temp, storage) = make_storage().await;
3649 let store = LockedSessionStore::new(storage);
3650 let result = store
3651 .update_runtime_config("does-not-exist", |s| {
3652 s.reasoning_effort = Some(ReasoningEffort::Low);
3653 })
3654 .await
3655 .unwrap();
3656 assert!(result.is_none());
3657 }
3658
3659 #[tokio::test]
3660 async fn merge_save_runtime_overwrites_messages_from_stale_snapshot() {
3661 use bamboo_domain::session::types::Message;
3666
3667 let (_temp, storage) = make_storage().await;
3668 let store = LockedSessionStore::new(storage.clone());
3669 let session_id = "stale-clobber";
3670
3671 let mut baseline = fresh(session_id);
3673 baseline.add_message(Message::user("hello"));
3674 storage.save_session(&baseline).await.unwrap();
3675 let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
3676
3677 let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
3679 after_chat.add_message(Message::user("second"));
3680 storage.save_session(&after_chat).await.unwrap();
3681 assert_eq!(
3682 storage
3683 .load_session(session_id)
3684 .await
3685 .unwrap()
3686 .unwrap()
3687 .messages
3688 .len(),
3689 2
3690 );
3691
3692 store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
3694 let after = storage.load_session(session_id).await.unwrap().unwrap();
3695 assert_eq!(
3696 after.messages.len(),
3697 1,
3698 "merge_save_runtime clobbers concurrent appends — this is why config patches must use update_runtime_config"
3699 );
3700 }
3701
3702 #[tokio::test]
3703 async fn stale_runtime_save_cannot_remove_admitted_inbox_transcript() {
3704 use bamboo_domain::session::types::Message;
3705 use bamboo_domain::SessionMessageId;
3706
3707 let (_temp, storage) = make_storage().await;
3708 let store = LockedSessionStore::new(storage.clone());
3709 let session_id = "stale-inbox-preserve";
3710
3711 let mut baseline = fresh(session_id);
3712 let mut base = Message::user("base");
3713 base.id = "base".to_string();
3714 baseline.add_message(base);
3715 storage.save_session(&baseline).await.unwrap();
3716 let mut stale = baseline.clone();
3717 let mut later_assistant = Message::assistant("runner output", None);
3718 later_assistant.id = "later-assistant".to_string();
3719 stale.add_message(later_assistant);
3720
3721 let mut durable = baseline;
3722 let inbox_id = SessionMessageId::parse("durable-inbox-id").unwrap();
3723 let mut admitted = Message::user("durable inbox message");
3724 admitted.id = inbox_id.as_str().to_string();
3725 durable.add_message(admitted);
3726 durable
3727 .session_inbox_admission_mut()
3728 .record(inbox_id.clone(), 7);
3729 storage.save_session(&durable).await.unwrap();
3730
3731 store.merge_save_runtime(&mut stale).await.unwrap();
3732 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3733 let ids = saved
3734 .messages
3735 .iter()
3736 .map(|message| message.id.as_str())
3737 .collect::<Vec<_>>();
3738 assert_eq!(ids, vec!["base", "durable-inbox-id", "later-assistant"]);
3739 assert_eq!(ids.iter().filter(|id| **id == inbox_id.as_str()).count(), 1);
3740 assert!(saved
3741 .session_inbox_admission()
3742 .is_some_and(|state| state.contains(&inbox_id)));
3743 }
3744
3745 #[tokio::test]
3746 async fn stale_runtime_save_preserves_typed_inbox_message_after_cursor_eviction() {
3747 use bamboo_domain::{
3748 SessionMessageEnvelope, SessionMessageId, SESSION_INBOX_ADMITTED_CAPACITY,
3749 };
3750
3751 let (_temp, storage) = make_storage().await;
3752 let store = LockedSessionStore::new(storage.clone());
3753 let session_id = "evicted-inbox-preserve";
3754 let mut durable = fresh(session_id);
3755 let mut envelope = SessionMessageEnvelope::user_input(session_id, "old durable inbox");
3756 envelope.id = SessionMessageId::parse("old-inbox-id").unwrap();
3757 durable.add_message(envelope.to_provider_message().unwrap());
3758 durable
3759 .session_inbox_admission_mut()
3760 .record(envelope.id.clone(), 1);
3761 for sequence in 2..=(SESSION_INBOX_ADMITTED_CAPACITY as u64 + 1) {
3762 durable.session_inbox_admission_mut().record(
3763 SessionMessageId::parse(format!("newer-{sequence}")).unwrap(),
3764 sequence,
3765 );
3766 }
3767 assert!(!durable
3768 .session_inbox_admission()
3769 .unwrap()
3770 .contains(&envelope.id));
3771 storage.save_session(&durable).await.unwrap();
3772
3773 let mut stale = fresh(session_id);
3774 stale.created_at = durable.created_at;
3775 store.merge_save_runtime(&mut stale).await.unwrap();
3776 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3777 assert_eq!(
3778 saved
3779 .messages
3780 .iter()
3781 .filter(|message| message.id == envelope.id.as_str())
3782 .count(),
3783 1
3784 );
3785 }
3786
3787 #[tokio::test]
3788 async fn runtime_final_save_cannot_resurrect_a_consumed_clarification() {
3789 use bamboo_domain::session::types::Message;
3790
3791 let (_temp, storage) = make_storage().await;
3792 let store = LockedSessionStore::new(storage.clone());
3793 let session_id = "consumed-clarification-final-save";
3794
3795 let mut suspended = fresh(session_id);
3796 suspended.add_message(Message::tool_result(
3797 "call-1",
3798 r#"{"status":"awaiting_clarification"}"#,
3799 ));
3800 suspended.set_pending_question(
3801 "call-1".to_string(),
3802 "ConclusionWithOptions".to_string(),
3803 "Choose".to_string(),
3804 vec!["A".to_string()],
3805 false,
3806 );
3807 suspended.metadata.insert(
3808 "runtime.suspend_reason".to_string(),
3809 "awaiting_clarification".to_string(),
3810 );
3811 storage.save_session(&suspended).await.unwrap();
3812 let mut stale_runner = suspended.clone();
3813
3814 let mut answered = suspended;
3815 answered.clear_pending_question();
3816 answered.metadata.remove("runtime.suspend_reason");
3817 answered.metadata.insert(
3818 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
3819 r#"["call-1"]"#.to_string(),
3820 );
3821 answered.metadata.insert(
3822 "clarification_resume_pending".to_string(),
3823 "true".to_string(),
3824 );
3825 answered.metadata.insert(
3826 "conclusion_with_options_resume_pending".to_string(),
3827 "true".to_string(),
3828 );
3829 answered.metadata.insert(
3830 "execute.startup_handoff_at".to_string(),
3831 "2026-08-10T09:00:00.000Z".to_string(),
3832 );
3833 let answer = answered
3834 .messages
3835 .iter_mut()
3836 .find(|message| message.tool_call_id.as_deref() == Some("call-1"))
3837 .unwrap();
3838 answer.content = "Selected response: A".to_string();
3839 storage.save_session(&answered).await.unwrap();
3840
3841 store.merge_save_runtime(&mut stale_runner).await.unwrap();
3842
3843 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3844 assert!(saved.pending_question.is_none());
3845 assert!(!saved.metadata.contains_key("runtime.suspend_reason"));
3846 assert_eq!(
3847 saved
3848 .metadata
3849 .get("clarification_resume_pending")
3850 .map(String::as_str),
3851 Some("true")
3852 );
3853 assert_eq!(
3854 saved
3855 .metadata
3856 .get("execute.startup_handoff_at")
3857 .map(String::as_str),
3858 Some("2026-08-10T09:00:00.000Z")
3859 );
3860 let answers = saved
3861 .messages
3862 .iter()
3863 .filter(|message| message.tool_call_id.as_deref() == Some("call-1"))
3864 .collect::<Vec<_>>();
3865 assert_eq!(answers.len(), 1);
3866 assert_eq!(answers[0].content, "Selected response: A");
3867 }
3868
3869 #[tokio::test]
3870 async fn runtime_checkpoint_cannot_resurrect_a_consumed_clarification() {
3871 use bamboo_domain::session::types::Message;
3872
3873 let (_temp, storage) = make_storage().await;
3874 let store = LockedSessionStore::new(storage.clone());
3875 let session_id = "consumed-clarification-checkpoint";
3876 let mut stale_runner = fresh(session_id);
3877 stale_runner.add_message(Message::tool_result("call-1", "waiting"));
3878 stale_runner.set_pending_question(
3879 "call-1".to_string(),
3880 "ConclusionWithOptions".to_string(),
3881 "Choose".to_string(),
3882 vec!["A".to_string()],
3883 false,
3884 );
3885 stale_runner.metadata.insert(
3886 "runtime.suspend_reason".to_string(),
3887 "awaiting_clarification".to_string(),
3888 );
3889
3890 let mut answered = stale_runner.clone();
3891 answered.clear_pending_question();
3892 answered.metadata.remove("runtime.suspend_reason");
3893 answered.metadata.insert(
3894 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
3895 r#"["call-1"]"#.to_string(),
3896 );
3897 answered.metadata.insert(
3898 "clarification_resume_pending".to_string(),
3899 "true".to_string(),
3900 );
3901 answered.messages[0].content = "Selected response: A".to_string();
3902 storage.save_session(&answered).await.unwrap();
3903
3904 store
3905 .checkpoint_runtime_session(&mut stale_runner)
3906 .await
3907 .unwrap();
3908
3909 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3910 assert!(saved.pending_question.is_none());
3911 assert!(!saved.metadata.contains_key("runtime.suspend_reason"));
3912 assert_eq!(saved.messages.len(), 1);
3913 assert_eq!(saved.messages[0].content, "Selected response: A");
3914 }
3915
3916 #[tokio::test]
3917 async fn runtime_final_save_does_not_consume_a_new_reused_permission_occurrence() {
3918 let (_temp, storage) = make_storage().await;
3919 let store = LockedSessionStore::new(storage.clone());
3920 let session_id = "reused-permission-final-save";
3921
3922 let mut durable = fresh(session_id);
3923 durable.add_message(typed_permission_result(
3924 "reused-call",
3925 "old-result",
3926 "generation-old",
3927 "Selected response: Approve",
3928 ));
3929 durable.metadata.insert(
3930 CONSUMED_RESPONSE_OCCURRENCES_KEY.to_string(),
3931 serde_json::to_string(&vec![ResponseOccurrence {
3932 tool_call_id: "reused-call".to_string(),
3933 tool_result_message_id: "old-result".to_string(),
3934 permission_generation: Some("generation-old".to_string()),
3935 }])
3936 .unwrap(),
3937 );
3938 storage.save_session(&durable).await.unwrap();
3939
3940 let mut new_runner = durable;
3941 new_runner.add_message(typed_permission_result(
3942 "reused-call",
3943 "new-result",
3944 "generation-new",
3945 "waiting for the new decision",
3946 ));
3947 new_runner.set_pending_question(
3948 "reused-call".to_string(),
3949 "Permission".to_string(),
3950 "Approve new operation?".to_string(),
3951 vec!["Approve".to_string(), "Deny".to_string()],
3952 false,
3953 );
3954 new_runner.metadata.insert(
3955 "runtime.suspend_reason".to_string(),
3956 "awaiting_permission_approval".to_string(),
3957 );
3958
3959 store.merge_save_runtime(&mut new_runner).await.unwrap();
3960
3961 let saved = storage.load_session(session_id).await.unwrap().unwrap();
3962 assert_eq!(
3963 saved
3964 .pending_question
3965 .as_ref()
3966 .map(|pending| pending.tool_call_id.as_str()),
3967 Some("reused-call")
3968 );
3969 assert_eq!(
3970 saved.messages.last().map(|message| message.id.as_str()),
3971 Some("new-result")
3972 );
3973 assert_eq!(
3974 saved.messages.last().unwrap().content,
3975 "waiting for the new decision"
3976 );
3977 assert_eq!(
3978 latest_response_occurrence(&saved, "reused-call")
3979 .and_then(|occurrence| occurrence.permission_generation),
3980 Some("generation-new".to_string())
3981 );
3982 }
3983
3984 #[tokio::test]
3985 async fn legacy_consumed_id_does_not_consume_a_new_reused_occurrence_after_upgrade() {
3986 let (_temp, storage) = make_storage().await;
3987 let store = LockedSessionStore::new(storage.clone());
3988 let session_id = "legacy-reused-permission-final-save";
3989
3990 let mut durable = fresh(session_id);
3991 durable.add_message(typed_permission_result(
3992 "reused-call",
3993 "old-result",
3994 "generation-old",
3995 "Selected response: Approve",
3996 ));
3997 durable.metadata.insert(
3998 CONSUMED_CLARIFICATION_IDS_KEY.to_string(),
3999 r#"["reused-call"]"#.to_string(),
4000 );
4001 storage.save_session(&durable).await.unwrap();
4002
4003 let mut new_runner = durable;
4004 new_runner.add_message(typed_permission_result(
4005 "reused-call",
4006 "new-result",
4007 "generation-new",
4008 "waiting for the new decision",
4009 ));
4010 new_runner.set_pending_question(
4011 "reused-call".to_string(),
4012 "Permission".to_string(),
4013 "Approve new operation?".to_string(),
4014 vec!["Approve".to_string(), "Deny".to_string()],
4015 false,
4016 );
4017
4018 store.merge_save_runtime(&mut new_runner).await.unwrap();
4019
4020 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4021 assert_eq!(
4022 saved
4023 .pending_question
4024 .as_ref()
4025 .map(|pending| pending.tool_call_id.as_str()),
4026 Some("reused-call")
4027 );
4028 assert_eq!(
4029 saved.messages.last().map(|message| message.id.as_str()),
4030 Some("new-result")
4031 );
4032 assert_eq!(
4033 latest_response_occurrence(&saved, "reused-call")
4034 .and_then(|occurrence| occurrence.permission_generation),
4035 Some("generation-new".to_string())
4036 );
4037 }
4038
4039 #[tokio::test]
4040 async fn runtime_checkpoint_does_not_consume_a_new_reused_permission_occurrence() {
4041 let (_temp, storage) = make_storage().await;
4042 let store = LockedSessionStore::new(storage.clone());
4043 let session_id = "reused-permission-checkpoint";
4044
4045 let mut durable = fresh(session_id);
4046 durable.add_message(typed_permission_result(
4047 "reused-call",
4048 "old-result",
4049 "generation-old",
4050 "Selected response: Deny",
4051 ));
4052 durable.metadata.insert(
4053 CONSUMED_RESPONSE_OCCURRENCES_KEY.to_string(),
4054 serde_json::to_string(&vec![ResponseOccurrence {
4055 tool_call_id: "reused-call".to_string(),
4056 tool_result_message_id: "old-result".to_string(),
4057 permission_generation: Some("generation-old".to_string()),
4058 }])
4059 .unwrap(),
4060 );
4061 storage.save_session(&durable).await.unwrap();
4062
4063 let mut new_runner = durable;
4064 new_runner.add_message(typed_permission_result(
4065 "reused-call",
4066 "new-result",
4067 "generation-new",
4068 "waiting for the new decision",
4069 ));
4070 new_runner.set_pending_question(
4071 "reused-call".to_string(),
4072 "Permission".to_string(),
4073 "Approve new operation?".to_string(),
4074 vec!["Approve".to_string(), "Deny".to_string()],
4075 false,
4076 );
4077
4078 store
4079 .checkpoint_runtime_session(&mut new_runner)
4080 .await
4081 .unwrap();
4082
4083 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4084 assert_eq!(
4085 saved
4086 .pending_question
4087 .as_ref()
4088 .map(|pending| pending.tool_call_id.as_str()),
4089 Some("reused-call")
4090 );
4091 assert_eq!(
4092 saved.messages.last().map(|message| message.id.as_str()),
4093 Some("new-result")
4094 );
4095 assert_eq!(
4096 latest_response_occurrence(&saved, "reused-call")
4097 .and_then(|occurrence| occurrence.permission_generation),
4098 Some("generation-new".to_string())
4099 );
4100 }
4101
4102 #[tokio::test]
4103 async fn consumed_permission_adoption_keeps_reexecute_id_and_generation_paired() {
4104 let (_temp, storage) = make_storage().await;
4105 let store = LockedSessionStore::new(storage.clone());
4106 let session_id = "consumed-permission-control-pair";
4107
4108 let mut stale_runner = fresh(session_id);
4109 stale_runner.add_message(typed_permission_result(
4110 "call-1",
4111 "result-1",
4112 "generation-1",
4113 "waiting",
4114 ));
4115 stale_runner.set_pending_question(
4116 "call-1".to_string(),
4117 "Permission".to_string(),
4118 "Approve?".to_string(),
4119 vec!["Approve".to_string(), "Deny".to_string()],
4120 false,
4121 );
4122
4123 let mut answered = stale_runner.clone();
4124 answered.clear_pending_question();
4125 answered.messages[0].content = "Selected response: Approve".to_string();
4126 answered.metadata.insert(
4127 CONSUMED_RESPONSE_OCCURRENCES_KEY.to_string(),
4128 serde_json::to_string(&vec![ResponseOccurrence {
4129 tool_call_id: "call-1".to_string(),
4130 tool_result_message_id: "result-1".to_string(),
4131 permission_generation: Some("generation-1".to_string()),
4132 }])
4133 .unwrap(),
4134 );
4135 answered.metadata.insert(
4136 "permission.reexecute_tool_call_id".to_string(),
4137 "call-1".to_string(),
4138 );
4139 answered.metadata.insert(
4140 "permission.reexecute_request_generation".to_string(),
4141 "generation-1".to_string(),
4142 );
4143 storage.save_session(&answered).await.unwrap();
4144
4145 store.merge_save_runtime(&mut stale_runner).await.unwrap();
4146
4147 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4148 assert!(saved.pending_question.is_none());
4149 assert_eq!(
4150 saved
4151 .metadata
4152 .get("permission.reexecute_tool_call_id")
4153 .map(String::as_str),
4154 Some("call-1")
4155 );
4156 assert_eq!(
4157 saved
4158 .metadata
4159 .get("permission.reexecute_request_generation")
4160 .map(String::as_str),
4161 Some("generation-1")
4162 );
4163 }
4164
4165 #[tokio::test]
4166 async fn checkpoint_runtime_session_preserves_disk_suffix_and_appends_live_messages() {
4167 use bamboo_domain::session::types::Message;
4168
4169 let (_temp, storage) = make_storage().await;
4170 let store = LockedSessionStore::new(storage.clone());
4171 let session_id = "checkpoint-no-shrink";
4172
4173 let mut baseline = fresh(session_id);
4174 baseline.add_message(Message::user("base"));
4175 storage.save_session(&baseline).await.unwrap();
4176 let mut runner_snapshot = baseline.clone();
4177
4178 let mut durable = baseline;
4179 let mut disk_only = Message::user("concurrent injected message");
4180 disk_only.id = "disk-only".to_string();
4181 durable.add_message(disk_only);
4182 storage.save_session(&durable).await.unwrap();
4183
4184 let mut live_only = Message::assistant("partial runner output", None);
4185 live_only.id = "live-only".to_string();
4186 runner_snapshot.add_message(live_only);
4187
4188 store
4189 .checkpoint_runtime_session(&mut runner_snapshot)
4190 .await
4191 .unwrap();
4192
4193 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4194 let ids = saved
4195 .messages
4196 .iter()
4197 .map(|message| message.id.as_str())
4198 .collect::<Vec<_>>();
4199 assert_eq!(
4200 ids,
4201 vec![durable.messages[0].id.as_str(), "disk-only", "live-only"]
4202 );
4203 assert_eq!(runner_snapshot.messages.len(), saved.messages.len());
4204 assert_eq!(runner_snapshot.messages[1].id, saved.messages[1].id);
4205 assert_eq!(runner_snapshot.messages[2].id, saved.messages[2].id);
4206 assert_eq!(saved.messages[1].content, "concurrent injected message");
4207 assert_eq!(saved.messages[2].content, "partial runner output");
4208 }
4209
4210 #[tokio::test]
4211 async fn prompt_rewrite_checkpoint_commits_without_an_archive_event() {
4212 use bamboo_domain::session::types::Message;
4213
4214 let (_temp, storage) = make_storage().await;
4215 let store = LockedSessionStore::new(storage.clone());
4216 let session_id = "prompt-rewrite-checkpoint";
4217 let mut expected = fresh(session_id);
4218 expected.add_message(Message::system("Base\n\nDEGRADABLE TOOL GUIDE"));
4219 expected.add_message(Message::user("continue"));
4220 expected.metadata.insert(
4221 "responses.previous_response_id".to_string(),
4222 "response-before-rewrite".to_string(),
4223 );
4224 storage.save_session(&expected).await.unwrap();
4225
4226 let mut staged = expected.clone();
4227 staged.messages[0].content = "Base".to_string();
4228 staged.metadata.remove("responses.previous_response_id");
4229 staged.reset_model_context_epoch(
4230 bamboo_domain::ModelContextResetReason::ExplicitHistoryRewrite,
4231 );
4232
4233 let outcome = store
4234 .checkpoint_prompt_rewrite_and_publish(&expected, &mut staged, |_| {})
4235 .await
4236 .unwrap();
4237
4238 assert_eq!(outcome, RetrievalWindowCheckpointOutcome::Committed);
4239 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4240 assert_eq!(saved.messages[0].content, "Base");
4241 assert!(saved.compression_events.is_empty());
4242 assert!(!saved
4243 .metadata
4244 .contains_key("responses.previous_response_id"));
4245 assert_eq!(
4246 saved
4247 .model_context_state
4248 .as_ref()
4249 .and_then(|state| state.last_reset_reason),
4250 Some(bamboo_domain::ModelContextResetReason::ExplicitHistoryRewrite)
4251 );
4252 assert_eq!(
4253 serde_json::to_value(&staged).unwrap(),
4254 serde_json::to_value(&saved).unwrap()
4255 );
4256 }
4257
4258 #[tokio::test]
4259 async fn prompt_rewrite_checkpoint_rebases_a_concurrent_suffix_without_writing() {
4260 use bamboo_domain::session::types::Message;
4261
4262 let (_temp, storage) = make_storage().await;
4263 let store = LockedSessionStore::new(storage.clone());
4264 let session_id = "prompt-rewrite-rebase";
4265 let mut expected = fresh(session_id);
4266 expected.add_message(Message::system("Base\n\nDEGRADABLE TOOL GUIDE"));
4267 storage.save_session(&expected).await.unwrap();
4268
4269 let mut staged = expected.clone();
4270 staged.messages[0].content = "Base".to_string();
4271 staged.reset_model_context_epoch(
4272 bamboo_domain::ModelContextResetReason::ExplicitHistoryRewrite,
4273 );
4274
4275 let mut durable = expected.clone();
4276 let mut concurrent = Message::user("concurrent durable suffix");
4277 concurrent.id = "prompt-rewrite-concurrent-suffix".to_string();
4278 durable.add_message(concurrent);
4279 storage.save_session(&durable).await.unwrap();
4280
4281 let outcome = store
4282 .checkpoint_prompt_rewrite_and_publish(&expected, &mut staged, |_| {})
4283 .await
4284 .unwrap();
4285
4286 assert_eq!(outcome, RetrievalWindowCheckpointOutcome::Rebased);
4287 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4288 assert_eq!(
4289 serde_json::to_value(&saved).unwrap(),
4290 serde_json::to_value(&durable).unwrap()
4291 );
4292 assert_eq!(staged.messages[0].content, expected.messages[0].content);
4293 assert_eq!(staged.messages[1].id, "prompt-rewrite-concurrent-suffix");
4294 assert!(staged.model_context_state.is_none());
4295 }
4296
4297 #[tokio::test]
4298 async fn prompt_rewrite_checkpoint_rebases_a_concurrent_pending_injection() {
4299 use bamboo_domain::session::types::Message;
4300
4301 let (_temp, storage) = make_storage().await;
4302 let store = LockedSessionStore::new(storage.clone());
4303 let session_id = "prompt-rewrite-pending-injection";
4304 let mut expected = fresh(session_id);
4305 expected.add_message(Message::system("Base\n\nDEGRADABLE TOOL GUIDE"));
4306 storage.save_session(&expected).await.unwrap();
4307
4308 let mut staged = expected.clone();
4309 staged.messages[0].content = "Base".to_string();
4310 staged.reset_model_context_epoch(
4311 bamboo_domain::ModelContextResetReason::ExplicitHistoryRewrite,
4312 );
4313
4314 let mut durable = expected.clone();
4315 durable.set_pending_injected_messages(vec![serde_json::json!({
4316 "content": "background shell completed",
4317 })]);
4318 storage.save_session(&durable).await.unwrap();
4319
4320 let outcome = store
4321 .checkpoint_prompt_rewrite_and_publish(&expected, &mut staged, |_| {})
4322 .await
4323 .unwrap();
4324
4325 assert_eq!(outcome, RetrievalWindowCheckpointOutcome::Rebased);
4326 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4327 assert_eq!(
4328 serde_json::to_value(&saved).unwrap(),
4329 serde_json::to_value(&durable).unwrap()
4330 );
4331 assert_eq!(staged.messages[0].content, expected.messages[0].content);
4332 assert_eq!(
4333 staged.pending_injected_messages(),
4334 durable.pending_injected_messages()
4335 );
4336 assert!(staged.model_context_state.is_none());
4337 }
4338
4339 #[tokio::test]
4340 async fn prompt_rewrite_checkpoint_rebases_concurrent_open_runtime_metadata() {
4341 use bamboo_domain::session::types::Message;
4342
4343 let (_temp, storage) = make_storage().await;
4344 let store = LockedSessionStore::new(storage.clone());
4345 let session_id = "prompt-rewrite-open-runtime-metadata";
4346 let mut expected = fresh(session_id);
4347 expected.add_message(Message::system("Base\n\nDEGRADABLE TOOL GUIDE"));
4348 storage.save_session(&expected).await.unwrap();
4349
4350 let mut staged = expected.clone();
4351 staged.messages[0].content = "Base".to_string();
4352 staged.reset_model_context_epoch(
4353 bamboo_domain::ModelContextResetReason::ExplicitHistoryRewrite,
4354 );
4355
4356 let mut durable = expected.clone();
4357 durable.metadata.insert(
4358 "concurrent.runtime.marker".to_string(),
4359 "latest".to_string(),
4360 );
4361 storage.save_session(&durable).await.unwrap();
4362
4363 let outcome = store
4364 .checkpoint_prompt_rewrite_and_publish(&expected, &mut staged, |_| {})
4365 .await
4366 .unwrap();
4367
4368 assert_eq!(outcome, RetrievalWindowCheckpointOutcome::Rebased);
4369 assert_eq!(
4370 staged
4371 .metadata
4372 .get("concurrent.runtime.marker")
4373 .map(String::as_str),
4374 Some("latest")
4375 );
4376 assert_eq!(staged.messages[0].content, expected.messages[0].content);
4377 assert!(staged.model_context_state.is_none());
4378 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4379 assert_eq!(
4380 serde_json::to_value(&saved).unwrap(),
4381 serde_json::to_value(&durable).unwrap()
4382 );
4383 }
4384
4385 #[tokio::test]
4386 async fn prompt_rewrite_checkpoint_rebases_concurrent_reasoning_update() {
4387 use bamboo_domain::{reasoning::ReasoningEffort, session::types::Message};
4388
4389 let (_temp, storage) = make_storage().await;
4390 let store = LockedSessionStore::new(storage.clone());
4391 let session_id = "prompt-rewrite-reasoning-update";
4392 let mut expected = fresh(session_id);
4393 expected.add_message(Message::system("Base\n\nDEGRADABLE TOOL GUIDE"));
4394 storage.save_session(&expected).await.unwrap();
4395
4396 let mut staged = expected.clone();
4397 staged.messages[0].content = "Base".to_string();
4398 staged.reset_model_context_epoch(
4399 bamboo_domain::ModelContextResetReason::ExplicitHistoryRewrite,
4400 );
4401
4402 let mut durable = expected.clone();
4403 durable.reasoning_effort = Some(ReasoningEffort::High);
4404 durable.metadata_version = 1;
4405 storage.save_session(&durable).await.unwrap();
4406
4407 let outcome = store
4408 .checkpoint_prompt_rewrite_and_publish(&expected, &mut staged, |_| {})
4409 .await
4410 .unwrap();
4411
4412 assert_eq!(outcome, RetrievalWindowCheckpointOutcome::Rebased);
4413 assert_eq!(staged.reasoning_effort, Some(ReasoningEffort::High));
4414 assert_eq!(staged.metadata_version, 1);
4415 assert_eq!(staged.messages[0].content, expected.messages[0].content);
4416 assert!(staged.model_context_state.is_none());
4417 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4418 assert_eq!(
4419 serde_json::to_value(&saved).unwrap(),
4420 serde_json::to_value(&durable).unwrap()
4421 );
4422 }
4423
4424 #[tokio::test]
4425 async fn manual_archive_rejection_checkpoint_persists_the_rewritten_tool_result() {
4426 use bamboo_domain::{FunctionCall, Message, ToolCall};
4427
4428 let (_temp, storage) = make_storage().await;
4429 let store = LockedSessionStore::new(storage.clone());
4430 let session_id = "manual-archive-rejection-checkpoint";
4431 let mut expected = fresh(session_id);
4432 expected.add_message(Message::user("archive older context"));
4433 let mut assistant = Message::assistant("", None);
4434 assistant.tool_calls = Some(vec![ToolCall {
4435 id: "archive-call".to_string(),
4436 tool_type: "function".to_string(),
4437 function: FunctionCall {
4438 name: "archive_context".to_string(),
4439 arguments: "{}".to_string(),
4440 },
4441 }]);
4442 expected.add_message(assistant);
4443 let mut result = Message::tool_result("archive-call", "Retrieval-window archive requested");
4444 result.id = "archive-result".to_string();
4445 result.tool_success = Some(true);
4446 expected.add_message(result);
4447 expected.metadata.insert(
4448 RESPONSES_PREVIOUS_RESPONSE_ID_KEY.to_string(),
4449 "response-before-rejection".to_string(),
4450 );
4451 storage.save_session(&expected).await.unwrap();
4452
4453 let occurrence = ResponseOccurrence {
4454 tool_call_id: "archive-call".to_string(),
4455 tool_result_message_id: "archive-result".to_string(),
4456 permission_generation: None,
4457 };
4458 let reason = "no eligible active logical group can be archived";
4459 let stage_rejection = |base: &Session| {
4460 let mut staged = base.clone();
4461 let staged_result = staged
4462 .messages
4463 .iter_mut()
4464 .find(|message| message.id == "archive-result")
4465 .unwrap();
4466 staged_result.tool_success = Some(false);
4467 staged_result.content = format!("archive_context rejected: {reason}");
4468 staged.metadata.insert(
4469 LAST_MANUAL_ARCHIVE_OCCURRENCE_KEY.to_string(),
4470 serde_json::to_string(&occurrence).unwrap(),
4471 );
4472 staged.metadata.insert(
4473 MANUAL_ARCHIVE_REJECTIONS_KEY.to_string(),
4474 serde_json::to_string(&vec![serde_json::json!({
4475 "occurrence": occurrence.clone(),
4476 "reason": reason,
4477 })])
4478 .unwrap(),
4479 );
4480 staged.metadata.remove(RESPONSES_PREVIOUS_RESPONSE_ID_KEY);
4481 staged.reset_model_context_epoch(
4482 bamboo_domain::ModelContextResetReason::ExplicitHistoryRewrite,
4483 );
4484 staged
4485 };
4486
4487 let mut staged = stage_rejection(&expected);
4488 let mut durable = expected.clone();
4489 durable.set_pending_injected_messages(vec![serde_json::json!({
4490 "content": "background shell completed",
4491 })]);
4492 storage.save_session(&durable).await.unwrap();
4493
4494 let outcome = store
4495 .checkpoint_manual_archive_rejection_and_publish(&expected, &mut staged, |_| {})
4496 .await
4497 .unwrap();
4498
4499 assert_eq!(outcome, RetrievalWindowCheckpointOutcome::Rebased);
4500 assert_eq!(
4501 staged.pending_injected_messages(),
4502 durable.pending_injected_messages()
4503 );
4504 let unchanged = storage.load_session(session_id).await.unwrap().unwrap();
4505 assert_eq!(
4506 serde_json::to_value(&unchanged).unwrap(),
4507 serde_json::to_value(&durable).unwrap()
4508 );
4509
4510 expected = staged;
4511 let mut staged = stage_rejection(&expected);
4512
4513 let outcome = store
4514 .checkpoint_manual_archive_rejection_and_publish(&expected, &mut staged, |_| {})
4515 .await
4516 .unwrap();
4517
4518 assert_eq!(outcome, RetrievalWindowCheckpointOutcome::Committed);
4519 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4520 let saved_result = saved
4521 .messages
4522 .iter()
4523 .find(|message| message.id == "archive-result")
4524 .unwrap();
4525 assert_eq!(saved_result.tool_success, Some(false));
4526 assert_eq!(
4527 saved_result.content,
4528 "archive_context rejected: no eligible active logical group can be archived"
4529 );
4530 assert!(saved
4531 .metadata
4532 .contains_key(LAST_MANUAL_ARCHIVE_OCCURRENCE_KEY));
4533 assert!(saved.metadata.contains_key(MANUAL_ARCHIVE_REJECTIONS_KEY));
4534 assert!(!saved
4535 .metadata
4536 .contains_key(RESPONSES_PREVIOUS_RESPONSE_ID_KEY));
4537 assert_eq!(
4538 saved.pending_injected_messages(),
4539 durable.pending_injected_messages()
4540 );
4541 assert_eq!(
4542 saved
4543 .model_context_state
4544 .as_ref()
4545 .and_then(|state| state.last_reset_reason),
4546 Some(bamboo_domain::ModelContextResetReason::ExplicitHistoryRewrite)
4547 );
4548 assert_eq!(
4549 serde_json::to_value(&saved).unwrap(),
4550 serde_json::to_value(&staged).unwrap()
4551 );
4552 }
4553
4554 #[tokio::test]
4555 async fn manual_archive_consumption_checkpoint_rebases_concurrent_open_metadata() {
4556 use bamboo_domain::{FunctionCall, Message, ToolCall};
4557
4558 let (_temp, storage) = make_storage().await;
4559 let store = LockedSessionStore::new(storage.clone());
4560 let session_id = "manual-archive-consumption-checkpoint";
4561 let mut expected = fresh(session_id);
4562 expected.add_message(Message::user("archive older context"));
4563 let mut assistant = Message::assistant("", None);
4564 assistant.tool_calls = Some(vec![ToolCall {
4565 id: "archive-call".to_string(),
4566 tool_type: "function".to_string(),
4567 function: FunctionCall {
4568 name: "archive_context".to_string(),
4569 arguments: "{}".to_string(),
4570 },
4571 }]);
4572 expected.add_message(assistant);
4573 let mut result = Message::tool_result("archive-call", "Retrieval-window archive requested");
4574 result.id = "archive-result".to_string();
4575 result.tool_success = Some(true);
4576 expected.add_message(result);
4577 storage.save_session(&expected).await.unwrap();
4578
4579 let occurrence = ResponseOccurrence {
4580 tool_call_id: "archive-call".to_string(),
4581 tool_result_message_id: "archive-result".to_string(),
4582 permission_generation: None,
4583 };
4584 let stage_consumption = |base: &Session| {
4585 let mut staged = base.clone();
4586 staged.metadata.insert(
4587 LAST_MANUAL_ARCHIVE_OCCURRENCE_KEY.to_string(),
4588 serde_json::to_string(&occurrence).unwrap(),
4589 );
4590 staged
4591 };
4592
4593 let mut staged = stage_consumption(&expected);
4594 let mut durable = expected.clone();
4595 durable.metadata.insert(
4596 "concurrent.workflow_index".to_string(),
4597 "workflow-42".to_string(),
4598 );
4599 durable.set_pending_injected_messages(vec![serde_json::json!({
4600 "id": "bash-1",
4601 "status": "completed",
4602 })]);
4603 storage.save_session(&durable).await.unwrap();
4604
4605 let outcome = store
4606 .checkpoint_manual_archive_consumption_and_publish(&expected, &mut staged, |_| {})
4607 .await
4608 .unwrap();
4609
4610 assert_eq!(outcome, RetrievalWindowCheckpointOutcome::Rebased);
4611 assert_eq!(
4612 serde_json::to_value(&staged).unwrap(),
4613 serde_json::to_value(&durable).unwrap()
4614 );
4615 let unchanged = storage.load_session(session_id).await.unwrap().unwrap();
4616 assert_eq!(
4617 serde_json::to_value(&unchanged).unwrap(),
4618 serde_json::to_value(&durable).unwrap()
4619 );
4620
4621 expected = staged;
4622 let mut staged = stage_consumption(&expected);
4623 let outcome = store
4624 .checkpoint_manual_archive_consumption_and_publish(&expected, &mut staged, |_| {})
4625 .await
4626 .unwrap();
4627
4628 assert_eq!(outcome, RetrievalWindowCheckpointOutcome::Committed);
4629 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4630 assert_eq!(
4631 saved.metadata.get("concurrent.workflow_index"),
4632 Some(&"workflow-42".to_string())
4633 );
4634 assert_eq!(
4635 saved.pending_injected_messages(),
4636 Some(vec![serde_json::json!({
4637 "id": "bash-1",
4638 "status": "completed",
4639 })])
4640 );
4641 assert_eq!(
4642 saved.metadata.get(LAST_MANUAL_ARCHIVE_OCCURRENCE_KEY),
4643 Some(&serde_json::to_string(&occurrence).unwrap())
4644 );
4645 }
4646
4647 #[tokio::test]
4648 async fn retrieval_window_checkpoint_preserves_staged_archive_flags() {
4649 use bamboo_domain::session::types::Message;
4650
4651 let (_temp, storage) = make_storage().await;
4652 let store = LockedSessionStore::new(storage.clone());
4653 let session_id = "retrieval-checkpoint-archive-flags";
4654 let mut expected = fresh(session_id);
4655 expected.add_message(Message::user("archive me"));
4656 expected.add_message(Message::assistant("retain me", None));
4657 storage.save_session(&expected).await.unwrap();
4658 let mut staged = stage_retrieval_window_archive(&expected, 0);
4659 let event_id = staged.compression_events[0].id.clone();
4660 let published = Arc::new(std::sync::Mutex::new(None));
4661 let published_clone = Arc::clone(&published);
4662
4663 let outcome = store
4664 .checkpoint_retrieval_window_and_publish(&expected, &mut staged, move |saved| {
4665 *published_clone.lock().unwrap() = Some(saved.clone());
4666 })
4667 .await
4668 .unwrap();
4669
4670 assert_eq!(outcome, RetrievalWindowCheckpointOutcome::Committed);
4671 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4672 assert!(saved.messages[0].compressed);
4673 assert_eq!(
4674 saved.messages[0].compressed_by_event_id.as_deref(),
4675 Some(event_id.as_str())
4676 );
4677 assert!(!saved.messages[1].compressed);
4678 assert_eq!(saved.compression_events.len(), 1);
4679 assert_eq!(
4680 saved.compression_events[0].kind,
4681 bamboo_domain::CompressionEventKind::RetrievalWindow
4682 );
4683 assert_eq!(
4684 serde_json::to_value(published.lock().unwrap().as_ref().unwrap()).unwrap(),
4685 serde_json::to_value(&saved).unwrap()
4686 );
4687 assert_eq!(
4688 serde_json::to_value(&staged).unwrap(),
4689 serde_json::to_value(&saved).unwrap()
4690 );
4691 }
4692
4693 #[tokio::test]
4694 async fn retrieval_window_checkpoint_rebases_concurrent_suffix_without_writing() {
4695 use bamboo_domain::session::types::Message;
4696
4697 let (_temp, storage) = make_storage().await;
4698 let store = LockedSessionStore::new(storage.clone());
4699 let session_id = "retrieval-checkpoint-concurrent-suffix";
4700 let mut expected = fresh(session_id);
4701 expected.add_message(Message::user("archive candidate"));
4702 storage.save_session(&expected).await.unwrap();
4703 let mut staged = stage_retrieval_window_archive(&expected, 0);
4704
4705 let mut durable = expected.clone();
4706 let mut concurrent = Message::user("concurrent durable suffix");
4707 concurrent.id = "concurrent-durable-suffix".to_string();
4708 durable.add_message(concurrent);
4709 storage.save_session(&durable).await.unwrap();
4710 let published = Arc::new(AtomicBool::new(false));
4711 let published_clone = Arc::clone(&published);
4712
4713 let outcome = store
4714 .checkpoint_retrieval_window_and_publish(&expected, &mut staged, move |_| {
4715 published_clone.store(true, Ordering::SeqCst);
4716 })
4717 .await
4718 .unwrap();
4719
4720 assert_eq!(outcome, RetrievalWindowCheckpointOutcome::Rebased);
4721 assert!(!published.load(Ordering::SeqCst));
4722 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4723 assert_eq!(
4724 serde_json::to_value(&saved).unwrap(),
4725 serde_json::to_value(&durable).unwrap(),
4726 "a conflict must not write the staged archive"
4727 );
4728 assert_eq!(staged.messages.len(), 2);
4729 assert_eq!(staged.messages[1].id, "concurrent-durable-suffix");
4730 assert!(staged.messages.iter().all(|message| !message.compressed));
4731 assert!(staged.compression_events.is_empty());
4732 assert!(staged.model_context_state.is_none());
4733 }
4734
4735 #[tokio::test]
4736 async fn runtime_only_save_preserves_checkpointed_ledger_and_publishes_merged_state() {
4737 let (_temp, storage) = make_storage().await;
4738 let store = LockedSessionStore::new(storage.clone());
4739 let session_id = "runtime-only-ledger-race";
4740 let baseline = fresh(session_id);
4741 storage.save_session(&baseline).await.unwrap();
4742 let mut stale_control = storage.load_session(session_id).await.unwrap().unwrap();
4743
4744 let mut runner = baseline;
4745 runner.model_context_state = Some(ledger_state(1, "runner-l1"));
4746 store.checkpoint_runtime_session(&mut runner).await.unwrap();
4747
4748 stale_control.metadata.insert(
4749 "runtime.suspend_reason".to_string(),
4750 "waiting_for_children".to_string(),
4751 );
4752 let published = std::sync::Arc::new(std::sync::Mutex::new(None));
4753 let published_clone = published.clone();
4754 store
4755 .save_runtime_only_and_publish(&mut stale_control, move |saved| {
4756 *published_clone.lock().unwrap() = Some(saved.clone());
4757 })
4758 .await
4759 .unwrap();
4760
4761 let expected = runner.model_context_state.clone();
4762 assert_eq!(stale_control.model_context_state, expected);
4763 assert_eq!(
4764 published
4765 .lock()
4766 .unwrap()
4767 .as_ref()
4768 .unwrap()
4769 .model_context_state,
4770 expected
4771 );
4772 let sidecar = storage
4773 .load_runtime_control_plane(session_id)
4774 .await
4775 .unwrap()
4776 .unwrap();
4777 assert_eq!(sidecar.model_context_state, expected);
4778 assert_eq!(
4779 sidecar
4780 .metadata
4781 .get("runtime.suspend_reason")
4782 .map(String::as_str),
4783 Some("waiting_for_children")
4784 );
4785 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
4786 assert_eq!(reloaded.model_context_state, expected);
4787 }
4788
4789 #[tokio::test]
4790 async fn full_runtime_save_preserves_newer_ledger_but_commits_control_mutation() {
4791 let (_temp, storage) = make_storage().await;
4792 let store = LockedSessionStore::new(storage.clone());
4793 let session_id = "full-save-ledger-race";
4794 let baseline = fresh(session_id);
4795 storage.save_session(&baseline).await.unwrap();
4796 let mut stale = storage.load_session(session_id).await.unwrap().unwrap();
4797
4798 let mut runner = baseline;
4799 runner.model_context_state = Some(ledger_state(1, "runner-l1"));
4800 store.checkpoint_runtime_session(&mut runner).await.unwrap();
4801
4802 stale
4803 .metadata
4804 .insert("activated_tools".to_string(), "[\"search\"]".to_string());
4805 let published = std::sync::Arc::new(std::sync::Mutex::new(None));
4806 let published_clone = published.clone();
4807 store
4808 .merge_save_runtime_and_publish(&mut stale, move |saved, committed| {
4809 assert!(committed);
4810 *published_clone.lock().unwrap() = Some(saved.clone());
4811 })
4812 .await
4813 .unwrap();
4814
4815 let expected = runner.model_context_state.clone();
4816 assert_eq!(stale.model_context_state, expected);
4817 assert_eq!(
4818 published
4819 .lock()
4820 .unwrap()
4821 .as_ref()
4822 .unwrap()
4823 .model_context_state,
4824 expected
4825 );
4826 let sidecar = storage
4827 .load_runtime_control_plane(session_id)
4828 .await
4829 .unwrap()
4830 .unwrap();
4831 assert_eq!(sidecar.model_context_state, expected);
4832 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
4833 assert_eq!(reloaded.model_context_state, expected);
4834 assert_eq!(
4835 reloaded.metadata.get("activated_tools").map(String::as_str),
4836 Some("[\"search\"]")
4837 );
4838 }
4839
4840 #[tokio::test]
4841 async fn newer_explicit_epoch_reset_wins_an_ordinary_full_runtime_save() {
4842 let (_temp, storage) = make_storage().await;
4843 let store = LockedSessionStore::new(storage.clone());
4844 let session_id = "full-save-ledger-reset";
4845 let mut durable = fresh(session_id);
4846 durable.model_context_state = Some(ledger_state(1, "runner-l1"));
4847 storage.save_session(&durable).await.unwrap();
4848
4849 let mut compression = storage.load_session(session_id).await.unwrap().unwrap();
4850 compression.reset_model_context_epoch(bamboo_domain::ModelContextResetReason::Compression);
4851 let reset = compression.model_context_state.clone();
4852 assert_eq!(reset.as_ref().unwrap().state_revision, 2);
4853 store.merge_save_runtime(&mut compression).await.unwrap();
4854
4855 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
4856 assert_eq!(reloaded.model_context_state, reset);
4857 assert_eq!(
4858 reloaded
4859 .model_context_state
4860 .as_ref()
4861 .and_then(|state| state.last_reset_reason),
4862 Some(bamboo_domain::ModelContextResetReason::Compression)
4863 );
4864 }
4865
4866 #[tokio::test]
4867 async fn checkpoint_rejects_equal_revision_divergence_without_overwriting_disk() {
4868 let (_temp, storage) = make_storage().await;
4869 let store = LockedSessionStore::new(storage.clone());
4870 let session_id = "ledger-checkpoint-cas";
4871 let mut baseline = fresh(session_id);
4872 baseline.model_context_state = Some(ledger_state(1, "runner-l1"));
4873 storage.save_session(&baseline).await.unwrap();
4874
4875 let mut first = baseline.clone();
4876 first.reset_model_context_epoch(bamboo_domain::ModelContextResetReason::Compression);
4877 let mut conflicting = baseline;
4878 conflicting.reset_model_context_epoch(bamboo_domain::ModelContextResetReason::Rollback);
4879 assert_eq!(
4880 first.model_context_state.as_ref().unwrap().state_revision,
4881 conflicting
4882 .model_context_state
4883 .as_ref()
4884 .unwrap()
4885 .state_revision
4886 );
4887
4888 store.checkpoint_runtime_session(&mut first).await.unwrap();
4889 let error = store
4890 .checkpoint_runtime_session(&mut conflicting)
4891 .await
4892 .unwrap_err();
4893 assert_eq!(error.kind(), std::io::ErrorKind::WouldBlock);
4894 let reloaded = storage.load_session(session_id).await.unwrap().unwrap();
4895 assert_eq!(reloaded.model_context_state, first.model_context_state);
4896 }
4897
4898 #[tokio::test]
4899 async fn activation_checkpoint_clears_presentation_without_shrinking_concurrent_turn() {
4900 use bamboo_domain::session::runtime_state::{
4901 AgentRuntimeState, AgentStatusState, WaitingForChildrenState,
4902 };
4903 use bamboo_domain::session::types::Message;
4904
4905 let (_temp, storage) = make_storage().await;
4906 let store = LockedSessionStore::new(storage.clone());
4907 let session_id = "activation-no-shrink";
4908 let mut baseline = fresh(session_id);
4909 baseline.add_message(Message::user("base"));
4910 let mut state = AgentRuntimeState::new("activation-run");
4911 state.status = AgentStatusState::Suspended;
4912 state.waiting_for_children = Some(WaitingForChildrenState::for_children(
4913 vec!["child-1".to_string()],
4914 bamboo_domain::session::runtime_state::ChildWaitPolicy::All,
4915 chrono::Utc::now(),
4916 ));
4917 baseline.agent_runtime_state = Some(state);
4918 baseline.metadata.insert(
4919 "runtime.suspend_reason".to_string(),
4920 "waiting_for_children".to_string(),
4921 );
4922 storage.save_session(&baseline).await.unwrap();
4923 let mut activation_snapshot = baseline.clone();
4924
4925 let mut concurrent = baseline;
4926 let mut normal = Message::assistant("normal concurrent answer", None);
4927 normal.id = "normal-concurrent".to_string();
4928 concurrent.add_message(normal);
4929 storage.save_session(&concurrent).await.unwrap();
4930
4931 let state = activation_snapshot.agent_runtime_state.as_mut().unwrap();
4932 state.status = AgentStatusState::Idle;
4933 state.suspension = None;
4934 activation_snapshot
4935 .metadata
4936 .remove("runtime.suspend_reason");
4937 store
4938 .checkpoint_runtime_session(&mut activation_snapshot)
4939 .await
4940 .unwrap();
4941
4942 let saved = storage.load_session(session_id).await.unwrap().unwrap();
4943 assert!(saved
4944 .messages
4945 .iter()
4946 .any(|message| message.id == "normal-concurrent"));
4947 let state = saved.agent_runtime_state.unwrap();
4948 assert_eq!(state.status, AgentStatusState::Idle);
4949 assert!(state.waiting_for_children.is_some());
4950 assert!(!saved.metadata.contains_key("runtime.suspend_reason"));
4951 }
4952
4953 #[tokio::test]
4954 async fn merge_save_runtime_preserves_disk_authoritative_metadata_with_single_load() {
4955 let (_temp, storage) = make_storage().await;
4961 let store = LockedSessionStore::new(storage.clone());
4962 let session_id = "runtime-merge-meta";
4963
4964 let mut baseline = fresh(session_id);
4966 baseline.title = "Auto Title".to_string();
4967 baseline.metadata_version = 0;
4968 storage.save_session(&baseline).await.unwrap();
4969
4970 let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
4972
4973 let mut renamed = storage.load_session(session_id).await.unwrap().unwrap();
4975 renamed.title = "User Renamed".to_string();
4976 renamed.title_version = 1;
4977 renamed.pinned = true;
4978 renamed.metadata_version = 1;
4979 store.commit_metadata(&renamed).await.unwrap();
4980
4981 stale_snapshot.title = "Auto Title".to_string();
4983 store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
4984
4985 let after = storage.load_session(session_id).await.unwrap().unwrap();
4986 assert_eq!(after.title, "User Renamed");
4987 assert!(after.pinned);
4988 assert_eq!(after.metadata_version, 1);
4989 assert_eq!(stale_snapshot.title, "User Renamed");
4991 assert_eq!(stale_snapshot.metadata_version, 1);
4992 }
4993
4994 #[tokio::test]
4995 async fn merge_save_runtime_preserves_durable_workflow_run_index_from_stale_runner() {
4996 let (_temp, storage) = make_storage().await;
4997 let store = LockedSessionStore::new(storage.clone());
4998 let session_id = "runtime-workflow-run-index";
4999
5000 let baseline = fresh(session_id);
5001 storage.save_session(&baseline).await.unwrap();
5002 let mut stale_runner = storage.load_session(session_id).await.unwrap().unwrap();
5003
5004 store
5005 .update_runtime_config(session_id, |session| {
5006 session.metadata.insert(
5007 "workflow.run_ids.v1".to_string(),
5008 r#"["http-started-run"]"#.to_string(),
5009 );
5010 })
5011 .await
5012 .unwrap()
5013 .expect("session exists");
5014
5015 store.merge_save_runtime(&mut stale_runner).await.unwrap();
5016
5017 assert_eq!(
5018 stale_runner
5019 .metadata
5020 .get("workflow.run_ids.v1")
5021 .map(String::as_str),
5022 Some(r#"["http-started-run"]"#)
5023 );
5024 let durable = storage.load_session(session_id).await.unwrap().unwrap();
5025 assert_eq!(
5026 durable
5027 .metadata
5028 .get("workflow.run_ids.v1")
5029 .map(String::as_str),
5030 Some(r#"["http-started-run"]"#)
5031 );
5032 }
5033
5034 #[tokio::test]
5038 async fn merge_save_runtime_adopts_disk_bypass_permissions() {
5039 use bamboo_domain::AgentRuntimeState;
5040
5041 let (_temp, storage) = make_storage().await;
5042 let store = LockedSessionStore::new(storage.clone());
5043 let session_id = "runtime-bypass";
5044
5045 let baseline = fresh(session_id);
5047 storage.save_session(&baseline).await.unwrap();
5048
5049 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
5051 loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
5052
5053 store
5055 .update_runtime_config(session_id, |s| {
5056 s.agent_runtime_state
5057 .get_or_insert_with(AgentRuntimeState::default)
5058 .bypass_permissions = true;
5059 })
5060 .await
5061 .unwrap()
5062 .expect("session exists");
5063
5064 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
5067
5068 let after = storage.load_session(session_id).await.unwrap().unwrap();
5069 assert!(
5070 after
5071 .agent_runtime_state
5072 .as_ref()
5073 .is_some_and(|s| s.bypass_permissions),
5074 "disk bypass=ON must survive a stale runtime save (#540)"
5075 );
5076 assert!(loop_snapshot
5078 .agent_runtime_state
5079 .as_ref()
5080 .is_some_and(|s| s.bypass_permissions));
5081 }
5082
5083 #[tokio::test]
5086 async fn merge_save_runtime_adopts_disk_auto_permission_mode() {
5087 use bamboo_domain::{AgentRuntimeState, SessionPermissionMode};
5088
5089 let (_temp, storage) = make_storage().await;
5090 let store = LockedSessionStore::new(storage.clone());
5091 let session_id = "runtime-auto";
5092
5093 storage.save_session(&fresh(session_id)).await.unwrap();
5094 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
5095 loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
5096 loop_snapshot.metadata.insert(
5097 "permission.requested_mode".to_string(),
5098 "default".to_string(),
5099 );
5100 loop_snapshot.metadata.insert(
5101 "permission.effective_mode".to_string(),
5102 "default".to_string(),
5103 );
5104 loop_snapshot.metadata.insert(
5105 "permission.executor_mapping".to_string(),
5106 "bamboo_runtime:default".to_string(),
5107 );
5108
5109 store
5110 .update_runtime_config(session_id, |session| {
5111 session
5112 .agent_runtime_state
5113 .get_or_insert_with(AgentRuntimeState::default)
5114 .set_permission_mode(SessionPermissionMode::Auto);
5115 session
5116 .metadata
5117 .insert("permission.policy_revision".to_string(), "12".to_string());
5118 session
5119 .metadata
5120 .insert("permission.requested_mode".to_string(), "auto".to_string());
5121 session
5122 .metadata
5123 .insert("permission.effective_mode".to_string(), "auto".to_string());
5124 session.metadata.insert(
5125 "permission.executor_mapping".to_string(),
5126 "bamboo_runtime:auto".to_string(),
5127 );
5128 session.metadata.insert(
5129 "permission.transitioned_at".to_string(),
5130 "2026-07-31T12:00:00Z".to_string(),
5131 );
5132 session.metadata_version = session.metadata_version.saturating_add(1);
5133 })
5134 .await
5135 .unwrap()
5136 .expect("session exists");
5137
5138 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
5139
5140 let durable = storage.load_session(session_id).await.unwrap().unwrap();
5141 for state in [
5142 durable.agent_runtime_state.as_ref(),
5143 loop_snapshot.agent_runtime_state.as_ref(),
5144 ] {
5145 assert_eq!(
5146 state.map(AgentRuntimeState::effective_permission_mode),
5147 Some(SessionPermissionMode::Auto)
5148 );
5149 }
5150 for session in [&durable, &loop_snapshot] {
5151 assert_eq!(
5152 session.metadata.get("permission.policy_revision"),
5153 Some(&"12".to_string())
5154 );
5155 assert_eq!(
5156 session.metadata.get("permission.requested_mode"),
5157 Some(&"auto".to_string())
5158 );
5159 assert_eq!(
5160 session.metadata.get("permission.effective_mode"),
5161 Some(&"auto".to_string())
5162 );
5163 assert_eq!(
5164 session.metadata.get("permission.executor_mapping"),
5165 Some(&"bamboo_runtime:auto".to_string())
5166 );
5167 assert_eq!(
5168 session.metadata.get("permission.transitioned_at"),
5169 Some(&"2026-07-31T12:00:00Z".to_string())
5170 );
5171 }
5172 }
5173
5174 #[tokio::test]
5177 async fn merge_save_runtime_adopts_disk_bypass_off() {
5178 use bamboo_domain::AgentRuntimeState;
5179
5180 let (_temp, storage) = make_storage().await;
5181 let store = LockedSessionStore::new(storage.clone());
5182 let session_id = "runtime-bypass-off";
5183
5184 let mut baseline = fresh(session_id);
5186 let on_state = AgentRuntimeState {
5187 bypass_permissions: true,
5188 ..AgentRuntimeState::default()
5189 };
5190 baseline.agent_runtime_state = Some(on_state);
5191 storage.save_session(&baseline).await.unwrap();
5192
5193 let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
5195
5196 store
5198 .update_runtime_config(session_id, |s| {
5199 s.agent_runtime_state
5200 .get_or_insert_with(AgentRuntimeState::default)
5201 .bypass_permissions = false;
5202 })
5203 .await
5204 .unwrap()
5205 .expect("session exists");
5206
5207 store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
5208
5209 let after = storage.load_session(session_id).await.unwrap().unwrap();
5210 assert!(
5211 !after
5212 .agent_runtime_state
5213 .as_ref()
5214 .is_some_and(|s| s.bypass_permissions),
5215 "disk bypass=OFF must survive a stale runtime save (#540)"
5216 );
5217 }
5218
5219 #[tokio::test]
5222 async fn save_runtime_authoritative_flags_persists_in_memory_posture_and_audit() {
5223 use bamboo_domain::AgentRuntimeState;
5224
5225 let (_temp, storage) = make_storage().await;
5226 let store = LockedSessionStore::new(storage.clone());
5227 let session_id = "child-reseed";
5228
5229 let mut baseline = fresh(session_id);
5231 let on_state = AgentRuntimeState {
5232 bypass_permissions: true,
5233 ..AgentRuntimeState::default()
5234 };
5235 baseline.agent_runtime_state = Some(on_state);
5236 for (key, value) in [
5237 ("permission.policy_revision", "12"),
5238 ("permission.requested_mode", "bypass"),
5239 ("permission.effective_mode", "bypass"),
5240 ("permission.executor_mapping", "bamboo_runtime:bypass"),
5241 ("permission.transitioned_at", "2026-07-31T12:00:00Z"),
5242 ] {
5243 baseline.metadata.insert(key.to_string(), value.to_string());
5244 }
5245 storage.save_session(&baseline).await.unwrap();
5246
5247 let mut child = storage.load_session(session_id).await.unwrap().unwrap();
5250 child
5251 .agent_runtime_state
5252 .get_or_insert_with(AgentRuntimeState::default)
5253 .bypass_permissions = false;
5254 for (key, value) in [
5255 ("permission.policy_revision", "13"),
5256 ("permission.requested_mode", "default"),
5257 ("permission.effective_mode", "default"),
5258 ("permission.executor_mapping", "bamboo_runtime:default"),
5259 ("permission.transitioned_at", "2026-07-31T12:01:00Z"),
5260 ] {
5261 child.metadata.insert(key.to_string(), value.to_string());
5262 }
5263
5264 store
5266 .save_runtime_authoritative_flags(&mut child)
5267 .await
5268 .unwrap();
5269
5270 let after = storage.load_session(session_id).await.unwrap().unwrap();
5271 assert!(
5272 !after
5273 .agent_runtime_state
5274 .as_ref()
5275 .is_some_and(|s| s.bypass_permissions),
5276 "authoritative re-seed of bypass=OFF must persist, not be reverted (#540/#74)"
5277 );
5278 for (key, value) in [
5279 ("permission.policy_revision", "13"),
5280 ("permission.requested_mode", "default"),
5281 ("permission.effective_mode", "default"),
5282 ("permission.executor_mapping", "bamboo_runtime:default"),
5283 ("permission.transitioned_at", "2026-07-31T12:01:00Z"),
5284 ] {
5285 assert_eq!(after.metadata.get(key).map(String::as_str), Some(value));
5286 }
5287 }
5288
5289 #[tokio::test]
5291 async fn merge_save_runtime_leaves_bypass_when_disk_has_no_runtime_state() {
5292 use bamboo_domain::AgentRuntimeState;
5293
5294 let (_temp, storage) = make_storage().await;
5295 let store = LockedSessionStore::new(storage.clone());
5296 let session_id = "no-runtime-state";
5297
5298 let baseline = fresh(session_id);
5300 assert!(baseline.agent_runtime_state.is_none());
5301 storage.save_session(&baseline).await.unwrap();
5302
5303 let mut running = storage.load_session(session_id).await.unwrap().unwrap();
5305 let on_state = AgentRuntimeState {
5306 bypass_permissions: true,
5307 ..AgentRuntimeState::default()
5308 };
5309 running.agent_runtime_state = Some(on_state);
5310
5311 store.merge_save_runtime(&mut running).await.unwrap();
5312
5313 assert!(
5314 running
5315 .agent_runtime_state
5316 .as_ref()
5317 .is_some_and(|s| s.bypass_permissions),
5318 "a runtime-state-less disk copy must not force bypass OFF (#540)"
5319 );
5320 }
5321
5322 #[tokio::test]
5325 async fn merge_preserves_disk_title_when_versions_equal() {
5326 let (_temp, storage) = make_storage().await;
5327 let session_id = "merge-equal";
5328
5329 let mut on_disk = fresh(session_id);
5330 on_disk.title = "User Set This".to_string();
5331 on_disk.title_version = 0;
5332 on_disk.title_generated = true;
5333 on_disk.metadata_version = 0;
5334 storage.save_session(&on_disk).await.unwrap();
5335
5336 let mut runtime_copy = fresh(session_id);
5337 runtime_copy.created_at = on_disk.created_at;
5338 runtime_copy.title = "Stale Default".to_string();
5339 runtime_copy.title_version = 0;
5340 runtime_copy.title_generated = false;
5341 runtime_copy.metadata_version = 0;
5342 runtime_copy.messages = vec![];
5343
5344 merge_save_session(&storage, &mut runtime_copy)
5345 .await
5346 .unwrap();
5347
5348 let after = storage.load_session(session_id).await.unwrap().unwrap();
5349 assert_eq!(after.title, "User Set This");
5350 assert_eq!(after.title_version, 0);
5351 assert!(after.title_generated);
5352 assert_eq!(runtime_copy.title, "User Set This");
5353 assert!(runtime_copy.title_generated);
5354 }
5355
5356 #[tokio::test]
5357 async fn merge_preserves_disk_when_disk_version_higher() {
5358 let (_temp, storage) = make_storage().await;
5359 let session_id = "merge-higher";
5360
5361 let mut on_disk = fresh(session_id);
5362 on_disk.title = "User Title v3".to_string();
5363 on_disk.title_version = 3;
5364 on_disk.metadata_version = 5;
5365 storage.save_session(&on_disk).await.unwrap();
5366
5367 let mut runtime_copy = fresh(session_id);
5368 runtime_copy.created_at = on_disk.created_at;
5369 runtime_copy.title = "Stale".to_string();
5370 runtime_copy.title_version = 1;
5371 runtime_copy.metadata_version = 0;
5372
5373 merge_save_session(&storage, &mut runtime_copy)
5374 .await
5375 .unwrap();
5376
5377 let after = storage.load_session(session_id).await.unwrap().unwrap();
5378 assert_eq!(after.title, "User Title v3");
5379 assert_eq!(after.title_version, 3);
5380 assert_eq!(after.metadata_version, 5);
5381 }
5382
5383 #[tokio::test]
5384 async fn merge_now_preserves_disk_pinned_in_metadata_group() {
5385 let (_temp, storage) = make_storage().await;
5386 let session_id = "pinned-merge";
5387
5388 let mut on_disk = fresh(session_id);
5389 on_disk.pinned = true;
5390 on_disk.metadata_version = 2;
5391 storage.save_session(&on_disk).await.unwrap();
5392
5393 let mut runtime_copy = fresh(session_id);
5394 runtime_copy.created_at = on_disk.created_at;
5395 runtime_copy.pinned = false;
5396 runtime_copy.metadata_version = 0;
5397
5398 merge_save_session(&storage, &mut runtime_copy)
5399 .await
5400 .unwrap();
5401
5402 let after = storage.load_session(session_id).await.unwrap().unwrap();
5403 assert!(
5404 after.pinned,
5405 "disk pinned=true should win over runtime false"
5406 );
5407 assert_eq!(after.metadata_version, 2);
5408 }
5409
5410 #[tokio::test]
5411 async fn merge_keeps_in_memory_when_session_version_higher() {
5412 let (_temp, storage) = make_storage().await;
5413 let session_id = "merge-bumped";
5414
5415 let mut on_disk = fresh(session_id);
5416 on_disk.title = "Old".to_string();
5417 on_disk.title_version = 1;
5418 on_disk.metadata_version = 3;
5419 storage.save_session(&on_disk).await.unwrap();
5420
5421 let mut authoritative_copy = fresh(session_id);
5422 authoritative_copy.created_at = on_disk.created_at;
5423 authoritative_copy.title = "New Authoritative".to_string();
5424 authoritative_copy.title_version = 2;
5425 authoritative_copy.metadata_version = 4;
5426 authoritative_copy.pinned = true;
5427
5428 merge_save_session(&storage, &mut authoritative_copy)
5429 .await
5430 .unwrap();
5431
5432 let after = storage.load_session(session_id).await.unwrap().unwrap();
5433 assert_eq!(after.title, "New Authoritative");
5434 assert_eq!(after.title_version, 2);
5435 assert_eq!(after.metadata_version, 4);
5436 assert!(after.pinned);
5437 }
5438
5439 #[tokio::test]
5440 async fn merge_keeps_runtime_messages_when_disk_only_changed_metadata() {
5441 let (_temp, storage) = make_storage().await;
5442 let session_id = "merge-messages";
5443
5444 let mut on_disk = fresh(session_id);
5445 on_disk.title = "Fresh Title".to_string();
5446 on_disk.title_version = 2;
5447 on_disk.metadata_version = 5;
5448 storage.save_session(&on_disk).await.unwrap();
5449
5450 let mut runtime_copy = fresh(session_id);
5451 runtime_copy.created_at = on_disk.created_at;
5452 runtime_copy.title = "Stale".to_string();
5453 runtime_copy.metadata_version = 0;
5454 runtime_copy.messages = vec![bamboo_domain::session::types::Message {
5455 role: bamboo_domain::session::types::Role::User,
5456 content: "keep me".to_string(),
5457 id: "msg-1".to_string(),
5458 created_at: chrono::Utc::now(),
5459 reasoning: None,
5460 reasoning_signature: None,
5461 content_parts: None,
5462 image_ocr: None,
5463 phase: None,
5464 tool_calls: None,
5465 tool_call_id: None,
5466 tool_success: None,
5467 compressed: false,
5468 compressed_by_event_id: None,
5469 never_compress: false,
5470 compression_level: 0,
5471 metadata: None,
5472 }];
5473
5474 merge_save_session(&storage, &mut runtime_copy)
5475 .await
5476 .unwrap();
5477
5478 let after = storage.load_session(session_id).await.unwrap().unwrap();
5479 assert_eq!(after.title, "Fresh Title");
5480 assert_eq!(after.metadata_version, 5);
5481 assert_eq!(after.messages.len(), 1);
5482 assert_eq!(after.messages[0].content, "keep me");
5483 }
5484
5485 #[tokio::test]
5486 async fn runtime_control_plane_port_uses_sidecar_without_rewriting_messages() {
5487 use bamboo_domain::session::types::Message;
5488
5489 let (_temp, storage) = make_storage().await;
5490 let store = LockedSessionStore::new(storage.clone());
5491 let session_id = "runtime-control-plane";
5492
5493 let mut durable = fresh(session_id);
5494 durable.add_message(Message::user("durable transcript"));
5495 storage.save_session(&durable).await.unwrap();
5496
5497 let mut runtime = durable.clone();
5498 runtime.model = "updated-control-plane-model".to_string();
5499 runtime.add_message(Message::assistant("uncheckpointed runtime message", None));
5500 RuntimeSessionPersistence::save_runtime_control_plane(&store, &mut runtime)
5501 .await
5502 .unwrap();
5503
5504 let control_plane =
5505 RuntimeSessionPersistence::load_runtime_control_plane(&store, session_id)
5506 .await
5507 .unwrap()
5508 .expect("control-plane exists");
5509 assert!(
5510 control_plane.messages.is_empty(),
5511 "LockedSessionStore must expose its message-free sidecar"
5512 );
5513 assert_eq!(control_plane.model, "updated-control-plane-model");
5514
5515 let reloaded = storage
5516 .load_session(session_id)
5517 .await
5518 .unwrap()
5519 .expect("session exists");
5520 assert_eq!(reloaded.model, "updated-control-plane-model");
5521 assert_eq!(
5522 reloaded.messages.len(),
5523 1,
5524 "control-plane save must not write the uncheckpointed message"
5525 );
5526 assert_eq!(reloaded.messages[0].content, "durable transcript");
5527 }
5528
5529 #[tokio::test]
5530 async fn atomic_task_patch_loads_inside_lock_and_preserves_interleaved_runtime_state() {
5531 let temp = tempfile::tempdir().unwrap();
5532 let inner = Arc::new(
5533 SessionStoreV2::new(temp.path().to_path_buf())
5534 .await
5535 .expect("storage init"),
5536 );
5537 let session_id = "atomic-task-patch";
5538 inner
5539 .save_session(&fresh(session_id))
5540 .await
5541 .expect("seed session");
5542
5543 let counted = Arc::new(CountingControlPlaneStorage {
5544 inner: inner.clone(),
5545 control_plane_loads: AtomicUsize::new(0),
5546 full_saves: AtomicUsize::new(0),
5547 runtime_state_saves: AtomicUsize::new(0),
5548 });
5549 let storage: Arc<dyn Storage> = counted.clone();
5550 let store = Arc::new(LockedSessionStore::new(storage));
5551 let guard = store.acquire_lock(session_id).await;
5552 let now = chrono::Utc::now();
5553 let task_list = bamboo_domain::TaskList {
5554 session_id: session_id.to_string(),
5555 title: "Atomic Task patch".to_string(),
5556 items: Vec::new(),
5557 created_at: now,
5558 updated_at: now,
5559 };
5560 let (started_tx, started_rx) = tokio::sync::oneshot::channel();
5561 let patch_store = store.clone();
5562 let patch = tokio::spawn(async move {
5563 let _ = started_tx.send(());
5564 RuntimeSessionPersistence::update_task_list_control_plane(
5565 patch_store.as_ref(),
5566 session_id,
5567 &task_list,
5568 "9",
5569 )
5570 .await
5571 });
5572 started_rx.await.expect("patch task started");
5573 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
5574 assert_eq!(
5575 counted.control_plane_loads.load(Ordering::SeqCst),
5576 0,
5577 "Task patch must acquire the session lock before loading its snapshot"
5578 );
5579
5580 let mut latest = inner
5584 .load_runtime_control_plane(session_id)
5585 .await
5586 .expect("load latest control-plane")
5587 .expect("control-plane exists");
5588 latest.agent_runtime_state = Some(bamboo_domain::AgentRuntimeState::new("latest-run"));
5589 latest
5590 .metadata
5591 .insert("concurrent.runtime".to_string(), "preserve".to_string());
5592 inner
5593 .save_runtime_state(&latest)
5594 .await
5595 .expect("publish concurrent runtime transition");
5596 drop(guard);
5597
5598 assert!(
5599 patch.await.expect("patch join").expect("patch succeeds"),
5600 "existing root must be patched"
5601 );
5602 assert_eq!(counted.control_plane_loads.load(Ordering::SeqCst), 1);
5603 let reloaded = inner
5604 .load_session(session_id)
5605 .await
5606 .expect("reload")
5607 .expect("session exists");
5608 assert_eq!(
5609 reloaded
5610 .agent_runtime_state
5611 .as_ref()
5612 .map(|state| state.run_id.as_str()),
5613 Some("latest-run")
5614 );
5615 assert_eq!(
5616 reloaded
5617 .metadata
5618 .get("concurrent.runtime")
5619 .map(String::as_str),
5620 Some("preserve")
5621 );
5622 assert_eq!(reloaded.task_list_version_meta().as_deref(), Some("9"));
5623 assert_eq!(
5624 reloaded.task_list.as_ref().map(|list| list.title.as_str()),
5625 Some("Atomic Task patch")
5626 );
5627 }
5628
5629 #[tokio::test]
5630 async fn paired_task_cas_conflict_cannot_overwrite_newer_root_or_child_state() {
5631 let (_temp, storage) = make_storage().await;
5632 let store = LockedSessionStore::new(storage.clone());
5633 let root_id = "task-cas-root";
5634 let child_id = "task-cas-child";
5635 let now = chrono::Utc::now();
5636 let task_list = |title: &str| bamboo_domain::TaskList {
5637 session_id: root_id.to_string(),
5638 title: title.to_string(),
5639 items: Vec::new(),
5640 created_at: now,
5641 updated_at: now,
5642 };
5643
5644 let mut root = fresh(root_id);
5645 root.set_task_list(task_list("newer root"));
5646 root.set_task_list_version_meta("2");
5647 storage.save_session(&root).await.expect("seed root");
5648 let mut child = Session::new_child(child_id, root_id, "model", "child");
5649 child.set_task_list(task_list("current child"));
5650 child.set_task_list_version_meta("1");
5651 storage.save_session(&child).await.expect("seed child");
5652
5653 let updated = RuntimeSessionPersistence::update_task_list_control_planes_if_version(
5654 &store,
5655 child_id,
5656 root_id,
5657 "1",
5658 &task_list("current child"),
5659 &task_list("stale evaluator"),
5660 "3",
5661 )
5662 .await
5663 .expect("CAS returns clean conflict");
5664 assert!(
5665 !updated,
5666 "mismatched root generation must reject both writes"
5667 );
5668
5669 let durable_root = storage
5670 .load_session(root_id)
5671 .await
5672 .expect("load root")
5673 .expect("root exists");
5674 let durable_child = storage
5675 .load_session(child_id)
5676 .await
5677 .expect("load child")
5678 .expect("child exists");
5679 assert_eq!(durable_root.task_list_version_meta().as_deref(), Some("2"));
5680 assert_eq!(
5681 durable_root
5682 .task_list
5683 .as_ref()
5684 .map(|list| list.title.as_str()),
5685 Some("newer root")
5686 );
5687 assert_eq!(durable_child.task_list_version_meta().as_deref(), Some("1"));
5688 assert_eq!(
5689 durable_child
5690 .task_list
5691 .as_ref()
5692 .map(|list| list.title.as_str()),
5693 Some("current child")
5694 );
5695 }
5696
5697 #[tokio::test]
5698 async fn paired_task_cas_success_uses_only_targeted_saves_and_preserves_both_transcripts() {
5699 use bamboo_domain::session::types::Message;
5700
5701 let temp = tempfile::tempdir().unwrap();
5702 let inner = Arc::new(
5703 SessionStoreV2::new(temp.path().to_path_buf())
5704 .await
5705 .expect("storage init"),
5706 );
5707 let root_id = "task-cas-success-root";
5708 let child_id = "task-cas-success-child";
5709 let now = chrono::Utc::now();
5710 let task_list = |title: &str| bamboo_domain::TaskList {
5711 session_id: root_id.to_string(),
5712 title: title.to_string(),
5713 items: Vec::new(),
5714 created_at: now,
5715 updated_at: now,
5716 };
5717
5718 let mut root = fresh(root_id);
5719 root.add_message(Message::user("root transcript"));
5720 root.metadata
5721 .insert("unrelated.root".to_string(), "preserve".to_string());
5722 root.agent_runtime_state = Some(bamboo_domain::AgentRuntimeState::new("root-run"));
5723 root.set_task_list(task_list("old shared"));
5724 root.set_task_list_version_meta("1");
5725 inner.save_session(&root).await.expect("seed root");
5726
5727 let mut child = Session::new_child(child_id, root_id, "model", "child");
5728 child.add_message(Message::user("child transcript"));
5729 child
5730 .metadata
5731 .insert("unrelated.child".to_string(), "preserve".to_string());
5732 child.agent_runtime_state = Some(bamboo_domain::AgentRuntimeState::new("child-run"));
5733 child.set_task_list(task_list("old shared"));
5734 child.set_task_list_version_meta("1");
5735 inner.save_session(&child).await.expect("seed child");
5736
5737 let counted = Arc::new(CountingControlPlaneStorage {
5738 inner: inner.clone(),
5739 control_plane_loads: AtomicUsize::new(0),
5740 full_saves: AtomicUsize::new(0),
5741 runtime_state_saves: AtomicUsize::new(0),
5742 });
5743 let storage: Arc<dyn Storage> = counted.clone();
5744 let store = LockedSessionStore::new(storage);
5745 assert!(
5746 RuntimeSessionPersistence::update_task_list_control_planes_if_version(
5747 &store,
5748 child_id,
5749 root_id,
5750 "1",
5751 &task_list("old shared"),
5752 &task_list("evaluated"),
5753 "2",
5754 )
5755 .await
5756 .expect("paired CAS succeeds")
5757 );
5758 assert_eq!(counted.control_plane_loads.load(Ordering::SeqCst), 2);
5759 assert_eq!(counted.runtime_state_saves.load(Ordering::SeqCst), 2);
5760 assert_eq!(
5761 counted.full_saves.load(Ordering::SeqCst),
5762 0,
5763 "evaluation CAS must not call full save_session for child or root"
5764 );
5765
5766 let durable_root = inner.load_session(root_id).await.unwrap().unwrap();
5767 let durable_child = inner.load_session(child_id).await.unwrap().unwrap();
5768 for (session, transcript, metadata_key, run_id) in [
5769 (
5770 &durable_root,
5771 "root transcript",
5772 "unrelated.root",
5773 "root-run",
5774 ),
5775 (
5776 &durable_child,
5777 "child transcript",
5778 "unrelated.child",
5779 "child-run",
5780 ),
5781 ] {
5782 assert_eq!(session.task_list_version_meta().as_deref(), Some("2"));
5783 assert_eq!(
5784 session.task_list.as_ref().map(|list| list.title.as_str()),
5785 Some("evaluated")
5786 );
5787 assert_eq!(session.messages.len(), 1);
5788 assert_eq!(session.messages[0].content, transcript);
5789 assert_eq!(
5790 session.metadata.get(metadata_key).map(String::as_str),
5791 Some("preserve")
5792 );
5793 assert_eq!(
5794 session
5795 .agent_runtime_state
5796 .as_ref()
5797 .map(|state| state.run_id.as_str()),
5798 Some(run_id)
5799 );
5800 }
5801 }
5802
5803 #[tokio::test]
5804 async fn locked_runtime_and_full_saves_adopt_task_conflicts_before_publish() {
5805 let temp = tempfile::tempdir().unwrap();
5806 let home = temp.path().to_path_buf();
5807 let first_storage = Arc::new(
5808 SessionStoreV2::new(home.clone())
5809 .await
5810 .expect("first storage init"),
5811 );
5812 let second_storage = Arc::new(
5813 SessionStoreV2::new(home)
5814 .await
5815 .expect("second storage init"),
5816 );
5817 let now = chrono::Utc::now();
5818 let task_list = |session_id: &str, title: &str| bamboo_domain::TaskList {
5819 session_id: session_id.to_string(),
5820 title: title.to_string(),
5821 items: Vec::new(),
5822 created_at: now,
5823 updated_at: now,
5824 };
5825
5826 let runtime_id = "ordinary-task-retry-runtime";
5827 let full_id = "ordinary-task-retry-full";
5828 let mut runtime_initial = fresh(runtime_id);
5829 runtime_initial.set_task_list(task_list(runtime_id, "runtime v1"));
5830 runtime_initial.set_task_list_version_meta("1");
5831 first_storage
5832 .save_session(&runtime_initial)
5833 .await
5834 .expect("seed runtime session");
5835 let mut full_initial = fresh(full_id);
5836 full_initial.set_task_list(task_list(full_id, "full v1"));
5837 full_initial.set_task_list_version_meta("1");
5838 first_storage
5839 .save_session(&full_initial)
5840 .await
5841 .expect("seed full session");
5842
5843 let mut runtime_advanced = runtime_initial.clone();
5844 runtime_advanced.set_task_list(task_list(runtime_id, "runtime v2"));
5845 runtime_advanced.set_task_list_version_meta("2");
5846 second_storage
5847 .save_runtime_state(&runtime_advanced)
5848 .await
5849 .expect("advance runtime Task generation");
5850 let mut full_advanced = full_initial.clone();
5851 full_advanced.set_task_list(task_list(full_id, "full v2"));
5852 full_advanced.set_task_list_version_meta("2");
5853 second_storage
5854 .save_runtime_state(&full_advanced)
5855 .await
5856 .expect("advance full Task generation");
5857
5858 let storage: Arc<dyn Storage> = first_storage.clone();
5859 let store = LockedSessionStore::new(storage);
5860 let runtime_published = Arc::new(std::sync::Mutex::new(None));
5861 let runtime_callback = runtime_published.clone();
5862 let mut runtime_stale = runtime_initial;
5863 runtime_stale
5864 .metadata
5865 .insert("runtime.non-task".to_string(), "preserved".to_string());
5866 store
5867 .save_runtime_only_and_publish(&mut runtime_stale, move |saved| {
5868 *runtime_callback.lock().expect("runtime publish lock") = Some(saved.clone());
5869 })
5870 .await
5871 .expect("locked runtime save rebases and retries");
5872
5873 let full_published = Arc::new(std::sync::Mutex::new(None));
5874 let full_callback = full_published.clone();
5875 let mut full_stale = full_initial;
5876 full_stale
5877 .metadata
5878 .insert("full.non-task".to_string(), "preserved".to_string());
5879 store
5880 .merge_save_runtime_and_publish(&mut full_stale, move |saved, committed| {
5881 assert!(committed);
5882 *full_callback.lock().expect("full publish lock") = Some(saved.clone());
5883 })
5884 .await
5885 .expect("locked full save rebases and retries");
5886
5887 let durable_runtime = first_storage
5888 .load_session(runtime_id)
5889 .await
5890 .unwrap()
5891 .expect("durable runtime session");
5892 let durable_full = first_storage
5893 .load_session(full_id)
5894 .await
5895 .unwrap()
5896 .expect("durable full session");
5897 let published_runtime = runtime_published
5898 .lock()
5899 .expect("runtime publish lock")
5900 .clone()
5901 .expect("runtime published snapshot");
5902 let published_full = full_published
5903 .lock()
5904 .expect("full publish lock")
5905 .clone()
5906 .expect("full published snapshot");
5907
5908 for (session, expected_title, metadata_key) in [
5909 (&runtime_stale, "runtime v2", "runtime.non-task"),
5910 (&durable_runtime, "runtime v2", "runtime.non-task"),
5911 (&published_runtime, "runtime v2", "runtime.non-task"),
5912 (&full_stale, "full v2", "full.non-task"),
5913 (&durable_full, "full v2", "full.non-task"),
5914 (&published_full, "full v2", "full.non-task"),
5915 ] {
5916 assert_eq!(session.task_list_version_meta().as_deref(), Some("2"));
5917 assert_eq!(
5918 session.task_list.as_ref().map(|list| list.title.as_str()),
5919 Some(expected_title)
5920 );
5921 assert_eq!(
5922 session.metadata.get(metadata_key).map(String::as_str),
5923 Some("preserved")
5924 );
5925 }
5926 }
5927
5928 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
5929 async fn single_task_cas_rejects_same_version_divergent_snapshot_without_callback() {
5930 let (_temp, storage) = make_storage().await;
5931 let store = LockedSessionStore::new(storage.clone());
5932 let session_id = "single-task-exact-snapshot";
5933 let now = chrono::Utc::now();
5934 let task_list = |title: &str| bamboo_domain::TaskList {
5935 session_id: session_id.to_string(),
5936 title: title.to_string(),
5937 items: Vec::new(),
5938 created_at: now,
5939 updated_at: now,
5940 };
5941 let durable_winner = task_list("durable winner");
5942 let stale_snapshot = task_list("stale same-version snapshot");
5943 let mut session = fresh(session_id);
5944 session.set_task_list(durable_winner.clone());
5945 session.set_task_list_version_meta("1");
5946 storage.save_session(&session).await.expect("seed session");
5947
5948 let published = Arc::new(AtomicBool::new(false));
5949 let callback = published.clone();
5950 assert!(!store
5951 .update_task_list_control_plane_if_version_and_publish(
5952 session_id,
5953 "1",
5954 &stale_snapshot,
5955 &task_list("stale evaluation"),
5956 "2",
5957 move |_| callback.store(true, Ordering::SeqCst),
5958 )
5959 .await
5960 .expect("same-version divergence is a clean stale result"));
5961 assert!(!published.load(Ordering::SeqCst));
5962 let durable = storage
5963 .load_session(session_id)
5964 .await
5965 .unwrap()
5966 .expect("session remains");
5967 assert_eq!(durable.task_list_version_meta().as_deref(), Some("1"));
5968 assert_eq!(
5969 serde_json::to_value(&durable.task_list).expect("serialize durable Task list"),
5970 serde_json::to_value(Some(&durable_winner)).expect("serialize expected Task list")
5971 );
5972 }
5973
5974 #[tokio::test]
5975 async fn paired_task_cas_rejects_same_version_divergent_snapshot_without_callback() {
5976 let (_temp, storage) = make_storage().await;
5977 let store = LockedSessionStore::new(storage.clone());
5978 let root_id = "paired-task-exact-snapshot-root";
5979 let child_id = "paired-task-exact-snapshot-child";
5980 let now = chrono::Utc::now();
5981 let task_list = |title: &str| bamboo_domain::TaskList {
5982 session_id: root_id.to_string(),
5983 title: title.to_string(),
5984 items: Vec::new(),
5985 created_at: now,
5986 updated_at: now,
5987 };
5988 let durable_winner = task_list("durable winner");
5989 let stale_snapshot = task_list("stale same-version snapshot");
5990 let mut root = fresh(root_id);
5991 root.set_task_list(durable_winner.clone());
5992 root.set_task_list_version_meta("1");
5993 storage.save_session(&root).await.expect("seed root");
5994 let mut child = Session::new_child(child_id, root_id, "model", "child");
5995 child.set_task_list(durable_winner.clone());
5996 child.set_task_list_version_meta("1");
5997 storage.save_session(&child).await.expect("seed child");
5998
5999 let published = Arc::new(AtomicBool::new(false));
6000 let callback = published.clone();
6001 assert!(!store
6002 .update_task_list_control_planes_if_version_and_publish(
6003 child_id,
6004 root_id,
6005 "1",
6006 &stale_snapshot,
6007 &task_list("stale evaluation"),
6008 "2",
6009 move |_, _| callback.store(true, Ordering::SeqCst),
6010 )
6011 .await
6012 .expect("same-version divergence is a clean stale result"));
6013 assert!(!published.load(Ordering::SeqCst));
6014 for id in [child_id, root_id] {
6015 let durable = storage
6016 .load_session(id)
6017 .await
6018 .unwrap()
6019 .expect("session remains");
6020 assert_eq!(durable.task_list_version_meta().as_deref(), Some("1"));
6021 assert_eq!(
6022 serde_json::to_value(&durable.task_list).expect("serialize durable Task list"),
6023 serde_json::to_value(Some(&durable_winner)).expect("serialize expected Task list")
6024 );
6025 }
6026 }
6027
6028 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
6029 async fn unconditional_root_task_patch_reports_final_cas_conflict_without_publishing() {
6030 let temp = tempfile::tempdir().unwrap();
6031 let home = temp.path().to_path_buf();
6032 let first_inner = Arc::new(
6033 SessionStoreV2::new(home.clone())
6034 .await
6035 .expect("first storage init"),
6036 );
6037 let root_id = "single-task-unconditional-loser";
6038 let now = chrono::Utc::now();
6039 let task_list = |title: &str| bamboo_domain::TaskList {
6040 session_id: root_id.to_string(),
6041 title: title.to_string(),
6042 items: Vec::new(),
6043 created_at: now,
6044 updated_at: now,
6045 };
6046 let mut root = fresh(root_id);
6047 root.task_list = Some(task_list("original"));
6048 root.set_task_list_version_meta("1");
6049 first_inner.save_session(&root).await.expect("seed root");
6050 let second_inner = Arc::new(
6051 SessionStoreV2::new(home)
6052 .await
6053 .expect("second storage init"),
6054 );
6055 let commit_reached = Arc::new(tokio::sync::Barrier::new(2));
6056 let release_commit = Arc::new(tokio::sync::Barrier::new(2));
6057 let storage: Arc<dyn Storage> = Arc::new(SingleCommitPauseStorage {
6058 inner: first_inner.clone(),
6059 commit_reached: commit_reached.clone(),
6060 release_commit: release_commit.clone(),
6061 });
6062 let store = LockedSessionStore::new(storage);
6063 let published = Arc::new(AtomicBool::new(false));
6064 let callback = published.clone();
6065 let loser_candidate = task_list("loser");
6066
6067 let loser = store.update_task_list_control_plane_and_publish(
6068 root_id,
6069 &loser_candidate,
6070 "2",
6071 move |_| callback.store(true, Ordering::SeqCst),
6072 );
6073 let winner = async {
6074 commit_reached.wait().await;
6075 let original = second_inner
6076 .load_runtime_control_plane(root_id)
6077 .await
6078 .expect("load winner original")
6079 .expect("winner original exists");
6080 let mut updated = original.clone();
6081 updated.task_list = Some(task_list("winner"));
6082 updated.set_task_list_version_meta("2");
6083 assert!(second_inner
6084 .save_task_control_plane_if_matches(&original, &updated)
6085 .await
6086 .expect("commit winner"));
6087 release_commit.wait().await;
6088 };
6089 let (loser_result, ()) = tokio::join!(loser, winner);
6090 let error = loser_result.expect_err("unconditional loser must be an explicit conflict");
6091 assert_eq!(error.kind(), std::io::ErrorKind::WouldBlock);
6092 assert!(!published.load(Ordering::SeqCst));
6093 let durable = first_inner
6094 .load_session(root_id)
6095 .await
6096 .unwrap()
6097 .expect("durable root");
6098 assert_eq!(durable.task_list_version_meta().as_deref(), Some("2"));
6099 assert_eq!(
6100 durable.task_list.as_ref().map(|list| list.title.as_str()),
6101 Some("winner")
6102 );
6103 }
6104
6105 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
6106 async fn independent_root_task_patches_have_one_final_cas_winner() {
6107 let temp = tempfile::tempdir().unwrap();
6108 let home = temp.path().to_path_buf();
6109 let first_inner = Arc::new(
6110 SessionStoreV2::new(home.clone())
6111 .await
6112 .expect("first storage init"),
6113 );
6114 let root_id = "single-task-cas-race-root";
6115 let now = chrono::Utc::now();
6116 let task_list = |title: &str| bamboo_domain::TaskList {
6117 session_id: root_id.to_string(),
6118 title: title.to_string(),
6119 items: Vec::new(),
6120 created_at: now,
6121 updated_at: now,
6122 };
6123 let mut root = fresh(root_id);
6124 root.set_task_list(task_list("original"));
6125 root.set_task_list_version_meta("1");
6126 first_inner.save_session(&root).await.expect("seed root");
6127 let second_inner = Arc::new(
6128 SessionStoreV2::new(home)
6129 .await
6130 .expect("second storage init"),
6131 );
6132 let before_commit = Arc::new(tokio::sync::Barrier::new(2));
6133 let first_storage: Arc<dyn Storage> = Arc::new(SingleCommitBarrierStorage {
6134 inner: first_inner.clone(),
6135 before_commit: before_commit.clone(),
6136 });
6137 let second_storage: Arc<dyn Storage> = Arc::new(SingleCommitBarrierStorage {
6138 inner: second_inner.clone(),
6139 before_commit,
6140 });
6141 let first_store = LockedSessionStore::new(first_storage);
6142 let second_store = LockedSessionStore::new(second_storage);
6143 let first_published = Arc::new(AtomicBool::new(false));
6144 let second_published = Arc::new(AtomicBool::new(false));
6145 let first_callback = first_published.clone();
6146 let second_callback = second_published.clone();
6147 let expected = task_list("original");
6148 let first_candidate = task_list("candidate one");
6149 let second_candidate = task_list("candidate two");
6150
6151 let first = first_store.update_task_list_control_plane_if_version_and_publish(
6154 root_id,
6155 "1",
6156 &expected,
6157 &first_candidate,
6158 "2",
6159 move |_| first_callback.store(true, Ordering::SeqCst),
6160 );
6161 let second = second_store.update_task_list_control_plane_if_version_and_publish(
6162 root_id,
6163 "1",
6164 &expected,
6165 &second_candidate,
6166 "2",
6167 move |_| second_callback.store(true, Ordering::SeqCst),
6168 );
6169 let (first_result, second_result) = tokio::join!(first, second);
6170 let first_won = first_result.expect("first root result");
6171 let second_won = second_result.expect("second root result");
6172 assert_ne!(first_won, second_won, "exactly one root candidate wins");
6173 assert_eq!(first_published.load(Ordering::SeqCst), first_won);
6174 assert_eq!(second_published.load(Ordering::SeqCst), second_won);
6175
6176 let durable = first_inner
6177 .load_session(root_id)
6178 .await
6179 .unwrap()
6180 .expect("root");
6181 assert_eq!(durable.task_list_version_meta().as_deref(), Some("2"));
6182 assert_eq!(
6183 durable.task_list.as_ref().map(|list| list.title.as_str()),
6184 Some(if first_won {
6185 "candidate one"
6186 } else {
6187 "candidate two"
6188 })
6189 );
6190 }
6191
6192 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
6193 async fn independent_locked_stores_revalidate_pair_cas_at_storage_commit_point() {
6194 let temp = tempfile::tempdir().unwrap();
6195 let home = temp.path().to_path_buf();
6196 let first_inner = Arc::new(
6197 SessionStoreV2::new(home.clone())
6198 .await
6199 .expect("first storage init"),
6200 );
6201 let root_id = "task-cas-race-root";
6202 let child_id = "task-cas-race-child";
6203 let now = chrono::Utc::now();
6204 let task_list = |title: &str| bamboo_domain::TaskList {
6205 session_id: root_id.to_string(),
6206 title: title.to_string(),
6207 items: Vec::new(),
6208 created_at: now,
6209 updated_at: now,
6210 };
6211
6212 let mut root = fresh(root_id);
6213 root.set_task_list(task_list("original shared"));
6214 root.set_task_list_version_meta("1");
6215 first_inner.save_session(&root).await.expect("seed root");
6216 let mut child = Session::new_child(child_id, root_id, "model", "child");
6217 child.set_task_list(task_list("original shared"));
6218 child.set_task_list_version_meta("1");
6219 first_inner.save_session(&child).await.expect("seed child");
6220
6221 let second_inner = Arc::new(
6222 SessionStoreV2::new(home)
6223 .await
6224 .expect("second storage init"),
6225 );
6226 let before_commit = Arc::new(tokio::sync::Barrier::new(2));
6227 let first_storage: Arc<dyn Storage> = Arc::new(PairCommitBarrierStorage {
6228 inner: first_inner.clone(),
6229 before_commit: before_commit.clone(),
6230 });
6231 let second_storage: Arc<dyn Storage> = Arc::new(PairCommitBarrierStorage {
6232 inner: second_inner.clone(),
6233 before_commit,
6234 });
6235 let first_store = LockedSessionStore::new(first_storage);
6236 let second_store = LockedSessionStore::new(second_storage);
6237 let first_published = Arc::new(AtomicBool::new(false));
6238 let second_published = Arc::new(AtomicBool::new(false));
6239 let first_published_callback = first_published.clone();
6240 let second_published_callback = second_published.clone();
6241 let expected = task_list("original shared");
6242 let first_candidate = task_list("candidate one");
6243 let second_candidate = task_list("candidate two");
6244
6245 let first = first_store.update_task_list_control_planes_if_version_and_publish(
6246 child_id,
6247 root_id,
6248 "1",
6249 &expected,
6250 &first_candidate,
6251 "2",
6252 move |_, _| first_published_callback.store(true, Ordering::SeqCst),
6253 );
6254 let second = second_store.update_task_list_control_planes_if_version_and_publish(
6255 child_id,
6256 root_id,
6257 "1",
6258 &expected,
6259 &second_candidate,
6260 "2",
6261 move |_, _| second_published_callback.store(true, Ordering::SeqCst),
6262 );
6263 let (first_result, second_result) = tokio::join!(first, second);
6264 let first_won = first_result.expect("first CAS result");
6265 let second_won = second_result.expect("second CAS result");
6266 assert_ne!(first_won, second_won, "exactly one staged v1 CAS may win");
6267 assert_eq!(first_published.load(Ordering::SeqCst), first_won);
6268 assert_eq!(second_published.load(Ordering::SeqCst), second_won);
6269
6270 let expected_title = if first_won {
6271 "candidate one"
6272 } else {
6273 "candidate two"
6274 };
6275 let durable_child = first_inner
6276 .load_session(child_id)
6277 .await
6278 .unwrap()
6279 .expect("child");
6280 let durable_root = second_inner
6281 .load_session(root_id)
6282 .await
6283 .unwrap()
6284 .expect("root");
6285 for session in [&durable_child, &durable_root] {
6286 assert_eq!(session.task_list_version_meta().as_deref(), Some("2"));
6287 assert_eq!(
6288 session.task_list.as_ref().map(|list| list.title.as_str()),
6289 Some(expected_title)
6290 );
6291 }
6292 }
6293
6294 #[tokio::test]
6295 async fn paired_task_second_write_failure_rolls_back_and_skips_publish_callback() {
6296 use bamboo_domain::session::types::Message;
6297
6298 let temp = tempfile::tempdir().unwrap();
6299 let inner = Arc::new(
6300 SessionStoreV2::new(temp.path().to_path_buf())
6301 .await
6302 .expect("storage init"),
6303 );
6304 let root_id = "task-cas-failure-root";
6305 let child_id = "task-cas-failure-child";
6306 let now = chrono::Utc::now();
6307 let task_list = |title: &str| bamboo_domain::TaskList {
6308 session_id: root_id.to_string(),
6309 title: title.to_string(),
6310 items: Vec::new(),
6311 created_at: now,
6312 updated_at: now,
6313 };
6314
6315 let mut root = fresh(root_id);
6316 root.add_message(Message::user("root transcript"));
6317 root.metadata
6318 .insert("unrelated.root".to_string(), "preserve".to_string());
6319 root.set_task_list(task_list("old shared"));
6320 root.set_task_list_version_meta("1");
6321 inner.save_session(&root).await.expect("seed root");
6322
6323 let mut child = Session::new_child(child_id, root_id, "model", "child");
6324 child.add_message(Message::user("child transcript"));
6325 child
6326 .metadata
6327 .insert("unrelated.child".to_string(), "preserve".to_string());
6328 child.set_task_list(task_list("old shared"));
6329 child.set_task_list_version_meta("1");
6330 inner.save_session(&child).await.expect("seed child");
6331
6332 inner
6333 .inject_runtime_task_transaction_fault(RuntimeTaskTransactionFault::SecondUpdatedWrite);
6334 let storage: Arc<dyn Storage> = inner.clone();
6335 let store = LockedSessionStore::new(storage);
6336 let published = Arc::new(AtomicBool::new(false));
6337 let published_for_callback = published.clone();
6338 let error = store
6339 .update_task_list_control_planes_if_version_and_publish(
6340 child_id,
6341 root_id,
6342 "1",
6343 &task_list("old shared"),
6344 &task_list("must roll back"),
6345 "2",
6346 move |_, _| published_for_callback.store(true, Ordering::SeqCst),
6347 )
6348 .await
6349 .expect_err("injected second write fails");
6350 assert!(error.to_string().contains("rolled back"), "{error}");
6351 assert!(
6352 !published.load(Ordering::SeqCst),
6353 "durable failure must not publish either cache snapshot"
6354 );
6355
6356 let durable_root = inner.load_session(root_id).await.unwrap().unwrap();
6357 let durable_child = inner.load_session(child_id).await.unwrap().unwrap();
6358 for (session, title, transcript, metadata_key) in [
6359 (
6360 &durable_root,
6361 "old shared",
6362 "root transcript",
6363 "unrelated.root",
6364 ),
6365 (
6366 &durable_child,
6367 "old shared",
6368 "child transcript",
6369 "unrelated.child",
6370 ),
6371 ] {
6372 assert_eq!(session.task_list_version_meta().as_deref(), Some("1"));
6373 assert_eq!(
6374 session.task_list.as_ref().map(|list| list.title.as_str()),
6375 Some(title)
6376 );
6377 assert_eq!(session.messages[0].content, transcript);
6378 assert_eq!(
6379 session.metadata.get(metadata_key).map(String::as_str),
6380 Some("preserve")
6381 );
6382 }
6383 }
6384
6385 #[tokio::test]
6388 async fn locked_merge_save_runtime_serialises_concurrent_writes() {
6389 let (_temp, storage) = make_storage().await;
6390 let store = Arc::new(LockedSessionStore::new(storage));
6391 let session_id = "lock-serial".to_string();
6392
6393 let base = fresh(&session_id);
6395 store.storage().save_session(&base).await.unwrap();
6396
6397 let store_a = store.clone();
6400 let store_b = store.clone();
6401 let sid_a = session_id.clone();
6402 let sid_b = session_id.clone();
6403
6404 let a = tokio::spawn(async move {
6405 let _guard = store_a.acquire_lock(&sid_a).await;
6406 let mut s = store_a
6407 .storage()
6408 .load_session(&sid_a)
6409 .await
6410 .unwrap()
6411 .unwrap();
6412 s.title = "Writer A".to_string();
6413 s.title_version = s.title_version.saturating_add(1);
6414 s.metadata_version = s.metadata_version.saturating_add(1);
6415 s.updated_at = chrono::Utc::now();
6416 store_a.storage().save_session(&s).await.unwrap();
6417 s.title_version
6418 });
6419
6420 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
6422
6423 let b = tokio::spawn(async move {
6424 let _guard = store_b.acquire_lock(&sid_b).await;
6425 let mut s = store_b
6426 .storage()
6427 .load_session(&sid_b)
6428 .await
6429 .unwrap()
6430 .unwrap();
6431 s.title = "Writer B".to_string();
6432 s.title_version = s.title_version.saturating_add(1);
6433 s.metadata_version = s.metadata_version.saturating_add(1);
6434 s.updated_at = chrono::Utc::now();
6435 store_b.storage().save_session(&s).await.unwrap();
6436 s.title_version
6437 });
6438
6439 let (ver_a, ver_b) = tokio::join!(a, b);
6440 let final_s = store
6441 .storage()
6442 .load_session(&session_id)
6443 .await
6444 .unwrap()
6445 .unwrap();
6446 assert!(
6447 ver_a.unwrap() != ver_b.unwrap(),
6448 "concurrent writers must produce distinct versions"
6449 );
6450 assert_eq!(final_s.metadata_version, 2);
6451 }
6452
6453 #[tokio::test]
6454 async fn commit_metadata_is_plain_save_inside_lock() {
6455 let (_temp, storage) = make_storage().await;
6456 let store = LockedSessionStore::new(storage);
6457 let session_id = "commit-plain";
6458
6459 let mut s = fresh(session_id);
6460 s.title = "Committed".to_string();
6461 s.metadata_version = 1;
6462 s.title_version = 2;
6463
6464 store.commit_metadata(&s).await.unwrap();
6465
6466 let after = store
6467 .storage()
6468 .load_session(session_id)
6469 .await
6470 .unwrap()
6471 .unwrap();
6472 assert_eq!(after.title, "Committed");
6473 assert_eq!(after.metadata_version, 1);
6474 assert_eq!(after.title_version, 2);
6475 }
6476
6477 #[tokio::test]
6480 async fn acquire_lock_self_evicts_when_no_other_holder() {
6481 let (_temp, storage) = make_storage().await;
6482 let store = LockedSessionStore::new(storage);
6483
6484 {
6485 let _guard = store.acquire_lock("solo").await;
6486 assert_eq!(store.locks.len(), 1, "entry present while the lock is held");
6487 }
6488 assert_eq!(
6491 store.locks.len(),
6492 0,
6493 "lock entry must be evicted once released with no other holder"
6494 );
6495 }
6496
6497 #[tokio::test]
6498 async fn acquire_lock_many_distinct_ids_do_not_accumulate() {
6499 let (_temp, storage) = make_storage().await;
6500 let store = LockedSessionStore::new(storage);
6501
6502 for i in 0..100 {
6504 let _guard = store.acquire_lock(&format!("sess-{i}")).await;
6505 }
6506 assert_eq!(
6507 store.locks.len(),
6508 0,
6509 "acquiring locks for many distinct ids must not grow the map"
6510 );
6511 }
6512
6513 #[tokio::test]
6514 async fn cancelled_last_waiter_reclaims_hundreds_of_session_locks() {
6515 let (_temp, storage) = make_storage().await;
6516 let store = LockedSessionStore::new(storage);
6517 for index in 0..512 {
6518 let id = format!("cancelled-child-{index}");
6519 let held = store.acquire_lock(&id).await;
6520 let mut waiter = Box::pin(store.acquire_lock(&id));
6521 assert!(
6522 std::future::poll_fn(|cx| std::task::Poll::Ready(
6523 std::future::Future::poll(waiter.as_mut(), cx).is_pending()
6524 ))
6525 .await
6526 );
6527 drop(held);
6528 drop(waiter);
6531 }
6532 assert!(store.locks.is_empty());
6533 }
6534
6535 #[tokio::test]
6536 async fn acquire_lock_concurrent_waiter_keeps_valid_lock_and_map_drains() {
6537 use std::sync::atomic::{AtomicUsize, Ordering};
6538
6539 let (_temp, storage) = make_storage().await;
6540 let store = Arc::new(LockedSessionStore::new(storage));
6541
6542 let active = Arc::new(AtomicUsize::new(0));
6544 let max_seen = Arc::new(AtomicUsize::new(0));
6545
6546 let mut handles = Vec::new();
6547 for _ in 0..8 {
6548 let store = store.clone();
6549 let active = active.clone();
6550 let max_seen = max_seen.clone();
6551 handles.push(tokio::spawn(async move {
6552 let _guard = store.acquire_lock("contended").await;
6553 let now = active.fetch_add(1, Ordering::SeqCst) + 1;
6554 max_seen.fetch_max(now, Ordering::SeqCst);
6555 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
6557 active.fetch_sub(1, Ordering::SeqCst);
6558 }));
6559 }
6560 for h in handles {
6561 h.await.unwrap();
6562 }
6563
6564 assert_eq!(
6569 max_seen.load(Ordering::SeqCst),
6570 1,
6571 "at most one holder of a given session lock at a time"
6572 );
6573 assert_eq!(
6574 store.locks.len(),
6575 0,
6576 "after all holders release, the contended entry must be fully evicted"
6577 );
6578 }
6579}