Skip to main content

bamboo_storage/
session_merge.rs

1//! Merge-aware session save helper.
2//!
3//! Provides [`merge_save_session`], which preserves any concurrent UI edits to
4//! the authoritative metadata group (`title`, `title_version`,
5//! `title_generated`, `pinned`, `metadata_version`) before writing the
6//! runtime-modified session to storage.
7//! Re-reads the latest persisted copy and only takes in-memory values when the
8//! caller's `metadata_version` strictly exceeds disk's.
9//!
10//! ## Field-by-field merge policy
11//!
12//! All authoritative metadata fields are grouped under `metadata_version`:
13//! when `disk.metadata_version >= session.metadata_version`, the on-disk
14//! `title`, `title_version`, `title_generated`, `pinned`, and
15//! `metadata_version` overwrite the in-memory values before writing.
16//! Authoritative writers bump
17//! `metadata_version` (and `title_version` for title edits) before calling so
18//! their values survive the merge; non-authoritative writers don't bump and so
19//! are overwritten by any later disk changes.
20//!
21//! ## Two save primitives
22//!
23//! - **`merge_save_session`** — stateless merge+save. Still works for
24//!   non-authoritative writers that hold `Arc<dyn Storage>` directly.
25//! - **`LockedSessionStore::merge_save_runtime`** — per-session-locked variant
26//!   that additionally serializes writes for the same session. Prefer this for
27//!   server-side paths where an authoritative writer may race with a runtime
28//!   save.
29//! - **`LockedSessionStore::commit_metadata`** — plain save inside a per-session
30//!   lock. For authoritative writers that have already performed
31//!   load→mutate→bump inside the lock; no merge needed (they hold the latest).
32//!
33//! Bare [`Storage::save_session`] is reserved for first-write paths (e.g. new
34//! session creation) where there is no prior on-disk copy to merge against.
35
36use 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
85/// A pending response is an authoritative compare-and-consume transaction.
86/// When a stale runner still carries the consumed ask, its terminal save must
87/// adopt the durable response control plane instead of resurrecting the ask or
88/// erasing the resume handoff. The bounded consumed-id ledger distinguishes a
89/// genuinely new pending question from an old snapshot after `pending_question`
90/// itself has been cleared. New sessions bind that ledger to the concrete
91/// tool-result message and permission generation; the id-only key is a
92/// compatibility fallback only when no occurrence ledger exists.
93fn 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        // A present but malformed v1 ledger is never widened to the legacy
109        // id-only matcher: doing so could consume a later reused provider id.
110        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            // Upgrade compatibility must still bind the old id-only ledger to
120            // its concrete durable tool-result. Otherwise a provider-reused id
121            // in the next round would be mistaken for the consumed occurrence.
122            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    // Durable ordering/content wins (including the selected-response rewrite),
136    // while any truly new runner-only suffix remains append-only.
137    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
232// ── LockedSessionStore ────────────────────────────────────────────────
233
234/// Wraps a [`Storage`] implementation with per-session write serialization.
235///
236/// Under the hood it maintains a `DashMap<String, Arc<Mutex<()>>>` so that
237/// only writes targeting the *same* session are serialised; different
238/// sessions proceed concurrently.
239pub struct LockedSessionStore {
240    storage: Arc<dyn Storage>,
241    locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
242    /// Serializes recoverable child/root Task transactions. Per-session locks
243    /// still provide the data isolation; this gate ensures a retained recovery
244    /// journal is resolved before another pair attempts to commit.
245    task_pair_transaction_lock: Arc<Mutex<()>>,
246}
247
248/// Self-cleaning guard returned by [`LockedSessionStore::acquire_lock`].
249///
250/// Holds the `OwnedMutexGuard` for the session's serialization mutex. On drop it
251/// releases the mutex **first** (so this guard's `Arc` clone is gone before the
252/// count is read) and then removes the map entry iff `Arc::strong_count == 1` —
253/// i.e. only the map's own reference remains, no other task holds or is waiting
254/// on this session's lock.
255///
256/// ## Race freedom
257///
258/// The strong-count check and the removal execute atomically under DashMap's
259/// per-shard lock via [`DashMap::remove_if`]. A waiter that clones the `Arc`
260/// (through `acquire_lock`'s `entry()`) does so under the same shard lock, so it
261/// either:
262/// - clones **before** our `remove_if` → `strong_count >= 2` → we skip removal,
263///   the waiter keeps a live, map-resident lock; or
264/// - clones **after** our `remove_if` → the entry is gone → it inserts a fresh
265///   `Arc<Mutex<()>>`; since our guard had already been released, the two tasks
266///   never overlapped and needed no mutual exclusion.
267///
268/// There is therefore no interleaving in which a waiter observes a lock that we
269/// then delete out from under it.
270pub struct SessionLockGuard {
271    /// `Option` so `Drop` can release the mutex before evaluating strong-count.
272    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        // Release the mutex (drops this guard's `Arc` clone) BEFORE reading the
280        // strong count, otherwise the count can never reach 1.
281        self.guard.take();
282        self.locks
283            .remove_if(&self.session_id, |_, arc| Arc::strong_count(arc) == 1);
284    }
285}
286
287impl LockedSessionStore {
288    /// Wrap an existing storage backend.
289    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    /// Borrow the inner storage for read-only access.
298    pub fn storage(&self) -> &Arc<dyn Storage> {
299        &self.storage
300    }
301
302    /// Acquire a per-session serialization guard.
303    ///
304    /// Only writes for the **same** session are serialised; writes for
305    /// different sessions can proceed concurrently.
306    ///
307    /// The returned [`SessionLockGuard`] is **self-cleaning**: when it drops it
308    /// releases the mutex and then removes the map entry iff no other holder
309    /// remains. Without this the `locks` map grew by one entry for every session
310    /// id ever written and never shrank (issue #346), so a long-lived server
311    /// leaked one `Arc<Mutex<()>>` per session-ever-persisted. See
312    /// [`SessionLockGuard`] for the race-freedom argument.
313    pub async fn acquire_lock(&self, session_id: &str) -> SessionLockGuard {
314        // `entry().or_insert_with().clone()` releases the DashMap shard lock at
315        // the end of THIS statement, before the `.await` below — never hold a
316        // shard lock across the async lock acquisition (it would deadlock the
317        // self-cleaning `remove_if` on drop, which also takes the shard lock).
318        let lock = self
319            .locks
320            .entry(session_id.to_string())
321            .or_insert_with(|| Arc::new(Mutex::new(())))
322            .clone();
323        // Arm cleanup before the cancellable wait. The previous holder may
324        // drop while this waiter owns the last extra Arc; cancelling that waiter
325        // must still reclaim the map entry without needing another acquisition.
326        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    /// Save a full snapshot while preserving Task generations advanced by an
336    /// independent store instance. V2 rejects such a stale write before any
337    /// mutation; reloading and adopting only Task-owned fields makes the retry,
338    /// caller snapshot, and subsequent cache publication agree exactly.
339    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    /// Sidecar-only counterpart of
373    /// [`Self::save_session_rebasing_task_conflicts`].
374    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    /// Runtime-only save: persist the control-plane (`agent_runtime_state`,
408    /// metadata, …) without rewriting the message history.
409    ///
410    /// This is the fast path for runtime-state mutations that do NOT change
411    /// `messages` — e.g. registering a parent's wait for spawned children. It
412    /// delegates to [`Storage::save_runtime_state`], which writes a small
413    /// sidecar (or falls back to a full save on backends without one).
414    ///
415    /// Like [`Self::merge_save_runtime`], it merges newer authoritative metadata
416    /// from disk so a concurrent UI title/pin edit is never clobbered — but it
417    /// reads only the lightweight control-plane snapshot (no message history) to
418    /// do so.
419    ///
420    /// Callers MUST NOT use this when they have appended messages: the in-memory
421    /// `messages` are ignored by the sidecar and would not be persisted. They
422    /// also MUST NOT use it to author `model_context_state`; the durable ledger
423    /// loaded under the lock always wins this narrow save.
424    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    /// Save the runtime control-plane and synchronously publish the committed
429    /// snapshot before releasing this session's serialization lock.
430    ///
431    /// `publish` must remain a short, non-blocking operation. Its synchronous
432    /// shape intentionally prevents callers from holding an in-memory cache
433    /// guard across an await. The callback also runs when the durable save
434    /// fails, preserving [`RuntimeSessionPersistence::save_runtime_control_plane`]
435    /// implementations that publish current runtime authorization state while
436    /// still returning the storage error. Authority conflicts are excluded:
437    /// publishing a rejected identity would disagree with durable authority.
438    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            // The control-plane sidecar carries `agent_runtime_state`, so a
450            // concurrent mid-run bypass flip is here too — don't revert it. #540.
451            adopt_fresher_disk_permission_posture(session, &latest);
452            // Runtime-only callers own narrow control-plane fields (Task list,
453            // parent wait state, …), never the model-context ledger. The sidecar
454            // load and save share this lock, so disk is unconditionally
455            // authoritative for the ledger and the published snapshot.
456            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    /// Atomically patch Task-owned control-plane fields and publish the saved
468    /// value before releasing this session's serialization lock.
469    ///
470    /// This couples durable commit order to cache publication order for
471    /// repository callers. The synchronous callback may take a cache guard but
472    /// cannot hold one across an await.
473    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    /// Compare-and-patch variant used by asynchronous evaluators. The durable
511    /// generation check, narrow Task mutation, save, and cache publication all
512    /// occur under the same per-session lock.
513    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    /// Recoverably validate and narrowly patch both the executing session and
549    /// its shared root. Locks are acquired in lexical id order to prevent two
550    /// child evaluators from deadlocking while sharing a root. The underlying
551    /// storage transaction must either commit both Task generations or retain
552    /// an undo journal and fail closed until both originals are restored.
553    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        // A prior rollback may have failed after one sidecar was published.
589        // Recover under the same lexical pair locks before observing either
590        // generation; an unresolved journal fails this access closed.
591        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        // The paired persistence port owns only Task list/generation. Keep the
616        // surrounding session snapshot byte-stable so its minimal undo journal
617        // never needs transcript or unrelated metadata.
618        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    /// Authoritative metadata commit.
645    ///
646    /// The caller must have already loaded the latest session, mutated the
647    /// metadata fields, and bumped `metadata_version` (and `title_version` if
648    /// applicable).  This method simply acquires the per-session lock and
649    /// performs a plain `storage.save_session`.
650    ///
651    /// The lock guarantees that no other write for this session interleaves
652    /// between the caller's load and this save, so merge is unnecessary.
653    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    /// Runtime / non-authoritative save with per-session lock.
664    ///
665    /// Inside the lock: reload disk, merge the authoritative metadata group
666    /// (`title`, `title_version`, `title_generated`, `pinned`, `metadata_version`) from disk into
667    /// the in-memory copy if disk's `metadata_version >= session.metadata_version`,
668    /// then save.
669    ///
670    /// This is the locked equivalent of [`merge_save_session`]; prefer it for
671    /// server-side paths where an authoritative write may race with this save.
672    ///
673    /// Adopts the on-disk typed permission mode so a running loop's save can't
674    /// revert a concurrent `PATCH /sessions` transition (#540/#770). Callers
675    /// that are themselves the authoritative writer of that posture — the
676    /// parent seeding a child's mode (#74) — must use
677    /// [`Self::save_runtime_authoritative_flags`] instead, which persists the
678    /// in-memory mode as-is.
679    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    /// Merge-save a runtime session and synchronously publish the resulting
685    /// snapshot before releasing its per-session serialization lock.
686    ///
687    /// The callback receives whether the durable save committed. It runs
688    /// after the save attempt so repository callers can preserve their existing
689    /// cache-on-failure policy without reopening a durable-to-cache race,
690    /// except when authority validation rejects the snapshot.
691    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    /// Persist an execute-boundary transcript checkpoint without allowing a
704    /// stale runner snapshot to shrink or rewrite the durable message log.
705    ///
706    /// The latest load, append-only message reconciliation, metadata merge and
707    /// save all happen while holding the same per-session lock.  Loading is
708    /// deliberately fail-closed: falling back to a blind full save when the
709    /// latest transcript cannot be read would reintroduce the SHRINK hazard
710    /// this checkpoint exists to prevent.
711    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    /// Checkpoint a runtime session and publish its reconciled snapshot before
717    /// releasing the same per-session serialization lock.
718    ///
719    /// The callback runs after the durable save attempt and receives its commit
720    /// status. A load failure returns before publication, matching the
721    /// checkpoint's fail-closed behavior.
722    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    /// Commit a retrieval-window rewrite without passing its message mutations
761    /// through the ordinary append-only reconciliation path.
762    ///
763    /// A concurrent durable transcript change returns a clean rebased Session
764    /// without writing. The engine must replan from that value before retrying.
765    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    /// Commit a bounded System-prompt rewrite without routing it through the
798    /// append-safe checkpoint or pretending that an archive event occurred.
799    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    /// Commit one bounded `archive_context` rejection rewrite.
832    ///
833    /// Unlike the append-safe runtime checkpoint, this preserves the staged
834    /// Tool result content while fencing the write with the exact durable base.
835    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    /// Commit the consumed marker for a successful no-op `archive_context`.
868    ///
869    /// The marker is a metadata rewrite, so the append-safe runtime checkpoint
870    /// is insufficient: it could full-save stale open metadata after a
871    /// concurrent runtime writer. Fence the complete retrieval base and let the
872    /// engine restage from the latest durable snapshot on any conflict.
873    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    /// Like [`Self::merge_save_runtime`] but does NOT adopt the on-disk
906    /// permission mode — the caller's in-memory value is authoritative and
907    /// persists as-is.
908    ///
909    /// For parent-side control writes to a child session (e.g. the #74
910    /// resident-reuse posture re-seed), which set the flag deliberately and must
911    /// not be reverted by the disk-wins protection meant for a running loop's
912    /// own stale saves. Still merges the authoritative metadata group.
913    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        // Single disk read serves BOTH the SHRINK diagnostic and the
933        // authoritative-metadata merge below. Previously this path loaded the
934        // session twice (once here, once inside the merge helper); on a parent
935        // session carrying the full conversation history that doubled the
936        // deserialization cost of every runtime save, which is the hot path
937        // during sub-agent spawn.
938        let latest = self.storage.load_session(&session.id).await?;
939
940        // DIAGNOSTIC: merge_save_runtime overwrites the whole `messages` array
941        // (it only merges authoritative metadata, not messages). If the incoming
942        // session is stale (fewer messages than what is already on disk), this save
943        // silently reverts a concurrent append (e.g. a just-persisted user message).
944        // Log a SHRINK warning so we can identify the stale writer.
945        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            // Never let a running loop's save revert a concurrent mid-run
979            // `PATCH /sessions {permission_mode|bypass_permissions}` transition.
980            // #540/#770. Skipped for
981            // authoritative flag writers (`save_runtime_authoritative_flags`).
982            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    /// Persist one validated RunSpec activation as the exact authority for the
995    /// worker's requested posture and complete audit record.
996    ///
997    /// Warm workers reuse a durable session id. An ordinary runtime save is
998    /// intentionally disk-adopting, so using it here would let the previous
999    /// activation's posture stick. This dedicated transaction preserves only
1000    /// durable UI metadata and SessionInbox admission/transcript proof, then
1001    /// writes the incoming posture with an audit revision above the durable
1002    /// floor while holding the same per-session lock.
1003    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    /// Atomically re-seed a resident child from its parent posture.
1058    ///
1059    /// The latest session load, typed-mode comparison, complete audit refresh,
1060    /// metadata CAS bump (only for a true typed transition), narrow companion
1061    /// mutation, save, and cache publication share one session lock.
1062    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    /// Persist a worker's bounded executor mapping only when the exact host
1108    /// posture observed before dispatch is still current. The remote event does
1109    /// not contribute an audit revision or transition timestamp: both are
1110    /// allocated from the latest durable record while this session lock is held.
1111    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    /// Apply a config-only mutation to a session without ever clobbering its
1156    /// `messages` (or other concurrently-written state).
1157    ///
1158    /// Unlike [`Self::merge_save_runtime`], the caller does NOT pass a session
1159    /// snapshot. Instead this loads the **latest** session from storage *inside*
1160    /// the per-session lock, applies `mutate` (intended for small config fields
1161    /// like `model_ref` / `reasoning_effort`), and saves. Because the load and
1162    /// save both happen under the lock, a concurrent append (e.g. `POST /chat`
1163    /// adding a user message) can never be reverted by this write.
1164    ///
1165    /// Returns the saved session, or `None` if it does not exist.
1166    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    /// Apply a config-only mutation and synchronously publish the saved
1179    /// snapshot before releasing the session lock.
1180    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    /// Load the authoritative response transaction candidate while the caller
1202    /// holds this store's per-session lock. The merge rules intentionally
1203    /// match [`Self::mutate_runtime_session_and_publish`], including adoption
1204    /// of a clarification that durable storage already consumed.
1205    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            // A response transaction starts from an append-safe transcript.
1234            // In particular, a cache snapshot that is timestamp-newer but
1235            // still carries an already-consumed ask must not resurrect that
1236            // ask or discard ordinary messages committed with the answer.
1237            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    /// Inspect the same authoritative snapshot a response mutation would use,
1257    /// without persisting or publishing it. Callers use this immediately before
1258    /// reserving a successor so an already-consumed response cannot allocate a
1259    /// replacement runner.
1260    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    /// Atomically load, validate/mutate, persist, and publish one runtime
1273    /// session under its canonical per-session lock.
1274    ///
1275    /// The nested result keeps validation errors distinct from storage errors:
1276    /// `Ok(Err(error))` means `mutate` rejected the latest snapshot and no save
1277    /// occurred. This is suitable for compare-and-consume operations such as a
1278    /// typed pending-question response, where separate load/save calls would
1279    /// allow another writer to replace the question between validation and
1280    /// persistence.
1281    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    /// Clear the legacy compatibility queue using durable CAS and publish the
1311    /// saved full snapshot before releasing the same session lock.
1312    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/// Infrastructure implementation of the domain runtime-persistence port.
1337/// Server should assemble this as `Arc<dyn RuntimeSessionPersistence>` and must
1338/// not define a separate adapter layer for the same behavior.
1339#[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
1480// ── Internal merge helper ─────────────────────────────────────────────
1481
1482/// Re-read the on-disk session and, when the disk copy carries a
1483/// `metadata_version >= session.metadata_version`, overwrite the in-memory
1484/// authoritative metadata fields with the disk values.
1485///
1486/// This is the core staleness-correction: non-authoritative writers call it
1487/// before saving so they don't accidentally revert a concurrent UI edit.
1488async 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
1503/// Runtime-only writes never own the ledger, so the snapshot loaded under the
1504/// session lock is authoritative even if the caller happens to carry a larger
1505/// revision. This keeps both the sidecar and synchronous cache publication from
1506/// regressing after a concurrent engine checkpoint.
1507fn 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
1513/// Merge policy for ordinary full runtime saves. Legitimate ledger writers
1514/// (compression/rollback/reconciliation) advance `state_revision`; stale or
1515/// conflicting non-ledger writers adopt the already-committed state. Equal
1516/// revisions with different bytes are concurrent children of the same base, so
1517/// durable state wins instead of allowing last-writer-wins corruption.
1518fn 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
1535/// A provider-bound ledger checkpoint must never silently substitute a newer or
1536/// conflicting disk ledger after the request body was prepared. Reject it so
1537/// the engine can roll back the in-memory candidate and retry from a fresh
1538/// session; only a strictly newer incoming revision may replace disk.
1539fn 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    // A staged OCR cache may legitimately be newer than disk and is committed
1573    // with the archive. Every other message field must still describe the exact
1574    // durable prefix used for planning.
1575    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    // Narrow runtime writers commit arbitrary metadata keys under the same
2033    // per-session lock (for example pending background-completion injections,
2034    // workflow indexes, or skill activation state). Authoritative PATCHes also
2035    // change execution-profile fields under `metadata_version`. A
2036    // prompt/archive rewrite planned before either commit must rebase instead
2037    // of full-saving its stale metadata or model configuration snapshot.
2038    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
2078/// Adopt the on-disk typed permission posture into the session about to be
2079/// saved when the durable posture is semantically fresher.
2080///
2081/// `PATCH /sessions {permission_mode|bypass_permissions}` is the authoritative
2082/// writer of this posture (a running loop only carries it forward from run
2083/// start). Without this, a runtime save from an in-flight run — which holds the
2084/// run-start value — silently reverts a concurrent mid-run transition on disk.
2085/// A true typed-mode difference always represents an authoritative durable
2086/// transition. When the modes are equal, the complete audit revision is the
2087/// ordering fence: an older/missing disk audit must never delete a newer
2088/// run-start policy/mapping refresh. #540/#770.
2089fn adopt_fresher_disk_permission_posture(session: &mut Session, latest: &Session) {
2090    // A disk copy with NO runtime state at all carries no authoritative mode
2091    // value — treat it as "unknown" and leave the in-memory flag untouched,
2092    // rather than forcing it OFF (which would silently disable a legitimately
2093    // bypassed run on any backend/path that doesn't round-trip the field). #540.
2094    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        // No runtime state in memory and disk says "off" → nothing to adopt;
2118        // avoid allocating a default state just to store `false`.
2119        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    // The typed posture and its complete bounded audit record move together.
2129    disk_audit.write_to(&mut session.metadata);
2130}
2131
2132/// Pure merge step: given a freshly-loaded on-disk copy, overwrite the
2133/// in-memory authoritative metadata group when disk's `metadata_version` is at
2134/// least the in-memory one. Split out so callers that have already loaded the
2135/// disk copy (e.g. [`LockedSessionStore::merge_save_runtime`]) don't pay for a
2136/// second read.
2137fn apply_authoritative_metadata(session: &mut Session, latest: &Session) {
2138    // Identity is independent of the UI metadata revision. Preserve it in the
2139    // caller snapshot too, so a successful merge-save cannot downgrade the cache.
2140    // Never replace an explicit Supervisor incarnation: the final storage guard
2141    // must reject stale identities after deletion/recreation rather than hiding them.
2142    // A newly constructed or previously deleted Ordinary session with this ID
2143    // is not a snapshot of the current Root and must not be rebound to it.
2144    if session.authority_identity.is_ordinary() && session.created_at == latest.created_at {
2145        session.authority_identity = latest.authority_identity.clone();
2146    }
2147    // Relationships are a separate monotonic authority, independent of UI
2148    // metadata. Adopt the canonical state into the actual caller snapshot only
2149    // for the same Root lifetime/incarnation; never rebind stale identities.
2150    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    // Project and its revision are one fence. Never stamp a newer disk revision
2158    // onto the caller's old Project; that would manufacture a fresh-looking
2159    // stale assignment. Equal-revision runtime workspace refreshes within the
2160    // same Project remain valid, while an actual reassignment adopts its whole
2161    // workspace context before the caller can be published to a cache.
2162    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
2211// ── Free merge-save function ──────────────────────────────────────────
2212
2213/// Save a session while preserving any concurrent UI edits to the
2214/// authoritative metadata group.
2215///
2216/// Behaviour: if the on-disk session has `metadata_version >=
2217/// session.metadata_version`, the on-disk `title`, `title_version`, `title_generated`, `pinned`
2218/// and `metadata_version` overwrite the in-memory values before writing.
2219///
2220/// This is the stateless variant (no per-session lock). Prefer
2221/// [`LockedSessionStore::merge_save_runtime`] for server-side paths where an
2222/// authoritative writer may race with this save.
2223pub 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// ── Tests ─────────────────────────────────────────────────────────────
2232
2233#[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(&current).await.unwrap();
2292                if equal_revision {
2293                    // The pre-fix merge could already have copied only the
2294                    // version. Reconcile this equal-version divergent Project.
2295                    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(&current).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(&current.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                // The merge read saw None. Publish from an independent V2 store
2506                // before allowing the loser to take its final filesystem lock.
2507                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            // The independent store's ordinary lookup index can remain stale;
2531            // strict authority must observe canonical publication without it.
2532            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            // Both independent LockedSessionStores have already recovered,
2923            // loaded, and validated v1 when they meet here. Their per-instance
2924            // lexical locks cannot serialize each other; the V2 commit-point
2925            // revalidation must select exactly one winner.
2926            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        // A cache write after the durable response can have a newer wall-clock
3225        // timestamp while still containing the old runner snapshot.
3226        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        // Legacy id-only adoption is still bound to the concrete response
3255        // occurrence. Real pending questions always have a paired tool-result;
3256        // keep that identity in this compatibility fixture so it cannot model
3257        // an unsafe id-only consume.
3258        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    // ── update_runtime_config: config patches must never clobber messages ──
3599
3600    #[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        // Persisted baseline: one user + one assistant turn.
3610        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        // Simulate `POST /chat` appending a new user message to disk.
3616        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        // A config-only patch must load the freshest session and preserve the
3622        // appended message (this is the regression that broke message sending on
3623        // existing sessions).
3624        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        // Characterization of the bug that motivated `update_runtime_config`:
3662        // `merge_save_runtime` writes the caller's `messages` verbatim, so a
3663        // stale snapshot reverts a concurrent append. Config-only writers must
3664        // therefore use `update_runtime_config`, never `merge_save_runtime`.
3665        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        // A handler loads the session (1 message) …
3672        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        // … then `POST /chat` appends a second message to disk …
3678        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        // … and the stale handler saves via merge_save_runtime -> append reverted.
3693        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        // Regression guard for the single-load refactor of `merge_save_runtime`:
4956        // it must STILL pull the authoritative metadata group (title / pinned /
4957        // metadata_version) from the freshest on-disk copy when disk's
4958        // metadata_version >= the in-memory one, even though it now reads disk
4959        // only once.
4960        let (_temp, storage) = make_storage().await;
4961        let store = LockedSessionStore::new(storage.clone());
4962        let session_id = "runtime-merge-meta";
4963
4964        // Baseline persisted by a runtime writer (metadata_version 0).
4965        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        // A stale runtime snapshot (still metadata_version 0, old title).
4971        let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
4972
4973        // An authoritative UI rename bumps metadata_version on disk.
4974        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        // The stale runtime writer saves: it must adopt the disk title/pinned.
4982        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        // And the in-memory copy was corrected by the merge too.
4990        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    // #540: a running loop's `merge_save_runtime` (carrying the run-start bypass
5035    // value) must NOT revert a concurrent mid-run `PATCH /sessions
5036    // {bypass_permissions}` write on disk — disk is the authoritative writer.
5037    #[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        // Baseline persisted with bypass OFF.
5046        let baseline = fresh(session_id);
5047        storage.save_session(&baseline).await.unwrap();
5048
5049        // The running loop holds a snapshot with bypass OFF (run-start value).
5050        let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
5051        loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
5052
5053        // A concurrent PATCH flips bypass ON on disk (via update_runtime_config).
5054        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        // The loop saves its stale snapshot: it must adopt disk's ON value, not
5065        // revert to OFF.
5066        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        // The in-memory copy is corrected too.
5077        assert!(loop_snapshot
5078            .agent_runtime_state
5079            .as_ref()
5080            .is_some_and(|s| s.bypass_permissions));
5081    }
5082
5083    // #770: the generalized disk-wins path must preserve Auto as a distinct
5084    // typed mode rather than collapsing it into the legacy bypass boolean.
5085    #[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    // The reverse direction: a PATCH turning bypass OFF must also stick against
5175    // a stale loop snapshot that still has it ON.
5176    #[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        // Baseline persisted with bypass ON.
5185        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        // Loop snapshot still ON.
5194        let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
5195
5196        // PATCH flips OFF on disk.
5197        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    // #540 review: the authoritative flag writer (#74 child-reseed) must NOT be
5220    // reverted by the disk-wins protection — its in-memory value persists as-is.
5221    #[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        // Child on disk has bypass ON (created under a bypassed parent).
5230        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        // Parent re-seeds the reused child to OFF (parent flipped bypass off),
5248        // loading the child then setting the flag in memory.
5249        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        // Authoritative write must persist OFF, not adopt the disk's stale ON.
5265        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    // A disk copy lacking runtime state must not force the in-memory bypass OFF.
5290    #[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        // Disk copy with NO agent_runtime_state.
5299        let baseline = fresh(session_id);
5300        assert!(baseline.agent_runtime_state.is_none());
5301        storage.save_session(&baseline).await.unwrap();
5302
5303        // A running loop legitimately carries bypass ON in memory.
5304        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    // ── Free-function merge tests (updated for metadata-group) ──────
5323
5324    #[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        // Publish a newer unrelated runtime transition while the Task patch is
5581        // queued on the same session lock. Once the guard releases, the patch
5582        // must load this latest snapshot and change only Task-owned fields.
5583        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        // Two expected-version patches stage from v1, but the storage final-CAS
6152        // admits only one regardless of process-local lock ownership.
6153        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    // ── LockedSessionStore tests ────────────────────────────────────
6386
6387    #[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        // Seed with base version.
6394        let base = fresh(&session_id);
6395        store.storage().save_session(&base).await.unwrap();
6396
6397        // Two concurrent authorised writers each bump and commit.
6398        // We'll simulate via clone-and-bump-then-commit.
6399        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        // Tiny yield so A goes first.
6421        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    // ── Self-cleaning per-session lock (issue #346) ─────────────────
6478
6479    #[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        // Dropping the guard runs the self-cleaning `remove_if`. Without the
6489        // eviction logic this stays at 1 forever (the pre-#346 leak).
6490        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        // Serially acquire+release for 100 distinct session ids.
6503        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            // Cancel after the previous holder handed ownership to the waiter,
6529            // but before the waiter is polled again to construct its guard.
6530            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        // Tracks concurrent holders of the SAME session lock; must never exceed 1.
6543        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                // Hold briefly so the other tasks actually queue on the mutex.
6556                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        // Mutual exclusion must hold: a self-cleaning removal that raced (removed
6565        // the entry a waiter had already cloned, letting a later task create and
6566        // lock a *second* mutex for the same id) would show 2 concurrent holders.
6567        // `remove_if`'s atomic strong-count check under the shard lock prevents it.
6568        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}