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