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