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`, `pinned`,
5//! `metadata_version`) before writing the runtime-modified session to storage.
6//! Re-reads the latest persisted copy and only takes in-memory values when the
7//! caller's `metadata_version` strictly exceeds disk's.
8//!
9//! ## Field-by-field merge policy
10//!
11//! All authoritative metadata fields are grouped under `metadata_version`:
12//! when `disk.metadata_version >= session.metadata_version`, the on-disk
13//! `title`, `title_version`, `pinned`, and `metadata_version` overwrite the
14//! in-memory values before writing. Authoritative writers bump
15//! `metadata_version` (and `title_version` for title edits) before calling so
16//! their values survive the merge; non-authoritative writers don't bump and so
17//! are overwritten by any later disk changes.
18//!
19//! ## Two save primitives
20//!
21//! - **`merge_save_session`** — stateless merge+save. Still works for
22//!   non-authoritative writers that hold `Arc<dyn Storage>` directly.
23//! - **`LockedSessionStore::merge_save_runtime`** — per-session-locked variant
24//!   that additionally serializes writes for the same session. Prefer this for
25//!   server-side paths where an authoritative writer may race with a runtime
26//!   save.
27//! - **`LockedSessionStore::commit_metadata`** — plain save inside a per-session
28//!   lock. For authoritative writers that have already performed
29//!   load→mutate→bump inside the lock; no merge needed (they hold the latest).
30//!
31//! Bare [`Storage::save_session`] is reserved for first-write paths (e.g. new
32//! session creation) where there is no prior on-disk copy to merge against.
33
34use std::sync::Arc;
35
36use bamboo_domain::session::types::Session;
37use bamboo_domain::storage::Storage;
38use bamboo_domain::RuntimeSessionPersistence;
39use dashmap::DashMap;
40use tokio::sync::{Mutex, OwnedMutexGuard};
41
42const AUTHORITATIVE_METADATA_KEYS: &[&str] = &["gold_config", "workflow.run_ids.v1"];
43
44// ── LockedSessionStore ────────────────────────────────────────────────
45
46/// Wraps a [`Storage`] implementation with per-session write serialization.
47///
48/// Under the hood it maintains a `DashMap<String, Arc<Mutex<()>>>` so that
49/// only writes targeting the *same* session are serialised; different
50/// sessions proceed concurrently.
51pub struct LockedSessionStore {
52    storage: Arc<dyn Storage>,
53    locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
54}
55
56/// Self-cleaning guard returned by [`LockedSessionStore::acquire_lock`].
57///
58/// Holds the `OwnedMutexGuard` for the session's serialization mutex. On drop it
59/// releases the mutex **first** (so this guard's `Arc` clone is gone before the
60/// count is read) and then removes the map entry iff `Arc::strong_count == 1` —
61/// i.e. only the map's own reference remains, no other task holds or is waiting
62/// on this session's lock.
63///
64/// ## Race freedom
65///
66/// The strong-count check and the removal execute atomically under DashMap's
67/// per-shard lock via [`DashMap::remove_if`]. A waiter that clones the `Arc`
68/// (through `acquire_lock`'s `entry()`) does so under the same shard lock, so it
69/// either:
70/// - clones **before** our `remove_if` → `strong_count >= 2` → we skip removal,
71///   the waiter keeps a live, map-resident lock; or
72/// - clones **after** our `remove_if` → the entry is gone → it inserts a fresh
73///   `Arc<Mutex<()>>`; since our guard had already been released, the two tasks
74///   never overlapped and needed no mutual exclusion.
75///
76/// There is therefore no interleaving in which a waiter observes a lock that we
77/// then delete out from under it.
78pub struct SessionLockGuard {
79    /// `Option` so `Drop` can release the mutex before evaluating strong-count.
80    guard: Option<OwnedMutexGuard<()>>,
81    locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
82    session_id: String,
83}
84
85impl Drop for SessionLockGuard {
86    fn drop(&mut self) {
87        // Release the mutex (drops this guard's `Arc` clone) BEFORE reading the
88        // strong count, otherwise the count can never reach 1.
89        self.guard.take();
90        self.locks
91            .remove_if(&self.session_id, |_, arc| Arc::strong_count(arc) == 1);
92    }
93}
94
95impl LockedSessionStore {
96    /// Wrap an existing storage backend.
97    pub fn new(storage: Arc<dyn Storage>) -> Self {
98        Self {
99            storage,
100            locks: Arc::new(DashMap::new()),
101        }
102    }
103
104    /// Borrow the inner storage for read-only access.
105    pub fn storage(&self) -> &Arc<dyn Storage> {
106        &self.storage
107    }
108
109    /// Acquire a per-session serialization guard.
110    ///
111    /// Only writes for the **same** session are serialised; writes for
112    /// different sessions can proceed concurrently.
113    ///
114    /// The returned [`SessionLockGuard`] is **self-cleaning**: when it drops it
115    /// releases the mutex and then removes the map entry iff no other holder
116    /// remains. Without this the `locks` map grew by one entry for every session
117    /// id ever written and never shrank (issue #346), so a long-lived server
118    /// leaked one `Arc<Mutex<()>>` per session-ever-persisted. See
119    /// [`SessionLockGuard`] for the race-freedom argument.
120    pub async fn acquire_lock(&self, session_id: &str) -> SessionLockGuard {
121        // `entry().or_insert_with().clone()` releases the DashMap shard lock at
122        // the end of THIS statement, before the `.await` below — never hold a
123        // shard lock across the async lock acquisition (it would deadlock the
124        // self-cleaning `remove_if` on drop, which also takes the shard lock).
125        let lock = self
126            .locks
127            .entry(session_id.to_string())
128            .or_insert_with(|| Arc::new(Mutex::new(())))
129            .clone();
130        let guard = lock.lock_owned().await;
131        SessionLockGuard {
132            guard: Some(guard),
133            locks: self.locks.clone(),
134            session_id: session_id.to_string(),
135        }
136    }
137
138    /// Runtime-only save: persist the control-plane (`agent_runtime_state`,
139    /// metadata, …) without rewriting the message history.
140    ///
141    /// This is the fast path for runtime-state mutations that do NOT change
142    /// `messages` — e.g. registering a parent's wait for spawned children. It
143    /// delegates to [`Storage::save_runtime_state`], which writes a small
144    /// sidecar (or falls back to a full save on backends without one).
145    ///
146    /// Like [`Self::merge_save_runtime`], it merges newer authoritative metadata
147    /// from disk so a concurrent UI title/pin edit is never clobbered — but it
148    /// reads only the lightweight control-plane snapshot (no message history) to
149    /// do so.
150    ///
151    /// Callers MUST NOT use this when they have appended messages: the in-memory
152    /// `messages` are ignored by the sidecar and would not be persisted.
153    pub async fn save_runtime_only(&self, session: &mut Session) -> std::io::Result<()> {
154        let _guard = self.acquire_lock(&session.id).await;
155        if let Ok(Some(latest)) = self.storage.load_runtime_control_plane(&session.id).await {
156            apply_authoritative_metadata(session, &latest);
157            // The control-plane sidecar carries `agent_runtime_state`, so a
158            // concurrent mid-run bypass flip is here too — don't revert it. #540.
159            adopt_disk_bypass_permissions(session, &latest);
160        }
161        self.storage.save_runtime_state(session).await
162    }
163
164    /// Authoritative metadata commit.
165    ///
166    /// The caller must have already loaded the latest session, mutated the
167    /// metadata fields, and bumped `metadata_version` (and `title_version` if
168    /// applicable).  This method simply acquires the per-session lock and
169    /// performs a plain `storage.save_session`.
170    ///
171    /// The lock guarantees that no other write for this session interleaves
172    /// between the caller's load and this save, so merge is unnecessary.
173    pub async fn commit_metadata(&self, session: &Session) -> std::io::Result<()> {
174        let _guard = self.acquire_lock(&session.id).await;
175        self.storage.save_session(session).await
176    }
177
178    /// Runtime / non-authoritative save with per-session lock.
179    ///
180    /// Inside the lock: reload disk, merge the authoritative metadata group
181    /// (`title`, `title_version`, `pinned`, `metadata_version`) from disk into
182    /// the in-memory copy if disk's `metadata_version >= session.metadata_version`,
183    /// then save.
184    ///
185    /// This is the locked equivalent of [`merge_save_session`]; prefer it for
186    /// server-side paths where an authoritative write may race with this save.
187    ///
188    /// Adopts the on-disk `bypass_permissions` so a running loop's save can't
189    /// revert a concurrent `PATCH /sessions` flip (#540). Callers that are
190    /// themselves the authoritative writer of that flag — the parent seeding a
191    /// child's posture (#74) — must use
192    /// [`Self::save_runtime_authoritative_flags`] instead, which persists the
193    /// in-memory flag as-is.
194    pub async fn merge_save_runtime(&self, session: &mut Session) -> std::io::Result<()> {
195        self.merge_save_runtime_inner(session, true).await
196    }
197
198    /// Like [`Self::merge_save_runtime`] but does NOT adopt the on-disk
199    /// `bypass_permissions` — the caller's in-memory value is authoritative and
200    /// persists as-is.
201    ///
202    /// For parent-side control writes to a child session (e.g. the #74
203    /// resident-reuse posture re-seed), which set the flag deliberately and must
204    /// not be reverted by the disk-wins protection meant for a running loop's
205    /// own stale saves. Still merges the authoritative metadata group.
206    pub async fn save_runtime_authoritative_flags(
207        &self,
208        session: &mut Session,
209    ) -> std::io::Result<()> {
210        self.merge_save_runtime_inner(session, false).await
211    }
212
213    async fn merge_save_runtime_inner(
214        &self,
215        session: &mut Session,
216        adopt_bypass: bool,
217    ) -> std::io::Result<()> {
218        let _guard = self.acquire_lock(&session.id).await;
219
220        // Single disk read serves BOTH the SHRINK diagnostic and the
221        // authoritative-metadata merge below. Previously this path loaded the
222        // session twice (once here, once inside the merge helper); on a parent
223        // session carrying the full conversation history that doubled the
224        // deserialization cost of every runtime save, which is the hot path
225        // during sub-agent spawn.
226        let latest = self.storage.load_session(&session.id).await.ok().flatten();
227
228        // DIAGNOSTIC: merge_save_runtime overwrites the whole `messages` array
229        // (it only merges authoritative metadata, not messages). If the incoming
230        // session is stale (fewer messages than what is already on disk), this save
231        // silently reverts a concurrent append (e.g. a just-persisted user message).
232        // Log a SHRINK warning so we can identify the stale writer.
233        let existing_message_count = latest.as_ref().map(|s| s.messages.len());
234        let incoming_message_count = session.messages.len();
235        if existing_message_count.is_some_and(|existing| existing > incoming_message_count) {
236            tracing::warn!(
237                "[{}] merge_save_runtime SHRINK: disk has {:?} messages, saving {} (last_role={:?}, updated_at={}); a stale writer is reverting a concurrent append",
238                session.id,
239                existing_message_count,
240                incoming_message_count,
241                session.messages.last().map(|m| format!("{:?}", m.role)),
242                session.updated_at,
243            );
244        } else {
245            tracing::debug!(
246                "[{}] merge_save_runtime: disk={:?} messages, saving {} (updated_at={})",
247                session.id,
248                existing_message_count,
249                incoming_message_count,
250                session.updated_at,
251            );
252        }
253
254        if let Some(latest) = latest.as_ref() {
255            apply_authoritative_metadata(session, latest);
256            // Never let a running loop's save revert a concurrent mid-run
257            // `PATCH /sessions {bypass_permissions}` flip. #540. Skipped for
258            // authoritative flag writers (`save_runtime_authoritative_flags`).
259            if adopt_bypass {
260                adopt_disk_bypass_permissions(session, latest);
261            }
262        }
263        self.storage.save_session(session).await
264    }
265
266    /// Apply a config-only mutation to a session without ever clobbering its
267    /// `messages` (or other concurrently-written state).
268    ///
269    /// Unlike [`Self::merge_save_runtime`], the caller does NOT pass a session
270    /// snapshot. Instead this loads the **latest** session from storage *inside*
271    /// the per-session lock, applies `mutate` (intended for small config fields
272    /// like `model_ref` / `reasoning_effort`), and saves. Because the load and
273    /// save both happen under the lock, a concurrent append (e.g. `POST /chat`
274    /// adding a user message) can never be reverted by this write.
275    ///
276    /// Returns the saved session, or `None` if it does not exist.
277    pub async fn update_runtime_config<F>(
278        &self,
279        session_id: &str,
280        mutate: F,
281    ) -> std::io::Result<Option<Session>>
282    where
283        F: FnOnce(&mut Session),
284    {
285        let _guard = self.acquire_lock(session_id).await;
286        let Some(mut session) = self.storage.load_session(session_id).await? else {
287            return Ok(None);
288        };
289        mutate(&mut session);
290        self.storage.save_session(&session).await?;
291        Ok(Some(session))
292    }
293}
294
295/// Infrastructure implementation of the domain runtime-persistence port.
296/// Server should assemble this as `Arc<dyn RuntimeSessionPersistence>` and must
297/// not define a separate adapter layer for the same behavior.
298#[async_trait::async_trait]
299impl RuntimeSessionPersistence for LockedSessionStore {
300    async fn save_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
301        self.merge_save_runtime(session).await
302    }
303
304    async fn load_runtime_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
305        self.storage.load_session(session_id).await
306    }
307}
308
309// ── Internal merge helper ─────────────────────────────────────────────
310
311/// Re-read the on-disk session and, when the disk copy carries a
312/// `metadata_version >= session.metadata_version`, overwrite the in-memory
313/// authoritative metadata fields with the disk values.
314///
315/// This is the core staleness-correction: non-authoritative writers call it
316/// before saving so they don't accidentally revert a concurrent UI edit.
317async fn merge_authoritative_metadata_into_stale(
318    storage: &Arc<dyn Storage>,
319    session: &mut Session,
320) {
321    if let Ok(Some(latest)) = storage.load_session(&session.id).await {
322        apply_authoritative_metadata(session, &latest);
323        adopt_disk_bypass_permissions(session, &latest);
324    }
325}
326
327/// Adopt the on-disk `agent_runtime_state.bypass_permissions` into the session
328/// about to be saved.
329///
330/// `PATCH /sessions {bypass_permissions}` is the SOLE authoritative writer of
331/// this flag (a running loop only carries it forward from run start). Without
332/// this, a runtime save from an in-flight run — which holds the run-start value
333/// — silently reverts a concurrent mid-run flip on disk. Unlike the metadata
334/// group this is NOT version-gated: the PATCH writes via `update_runtime_config`,
335/// which does not bump `metadata_version`. #540.
336fn adopt_disk_bypass_permissions(session: &mut Session, latest: &Session) {
337    // A disk copy with NO runtime state at all carries no authoritative bypass
338    // value — treat it as "unknown" and leave the in-memory flag untouched,
339    // rather than forcing it OFF (which would silently disable a legitimately
340    // bypassed run on any backend/path that doesn't round-trip the field). #540.
341    let Some(disk_bypass) = latest
342        .agent_runtime_state
343        .as_ref()
344        .map(|state| state.bypass_permissions)
345    else {
346        return;
347    };
348    match session.agent_runtime_state.as_mut() {
349        Some(state) => state.bypass_permissions = disk_bypass,
350        // No runtime state in memory and disk says "off" → nothing to adopt;
351        // avoid allocating a default state just to store `false`.
352        None if disk_bypass => {
353            session
354                .agent_runtime_state
355                .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
356                .bypass_permissions = true;
357        }
358        None => {}
359    }
360}
361
362/// Pure merge step: given a freshly-loaded on-disk copy, overwrite the
363/// in-memory authoritative metadata group when disk's `metadata_version` is at
364/// least the in-memory one. Split out so callers that have already loaded the
365/// disk copy (e.g. [`LockedSessionStore::merge_save_runtime`]) don't pay for a
366/// second read.
367fn apply_authoritative_metadata(session: &mut Session, latest: &Session) {
368    if latest.metadata_version >= session.metadata_version {
369        session.title = latest.title.clone();
370        session.title_version = latest.title_version;
371        session.pinned = latest.pinned;
372        for key in AUTHORITATIVE_METADATA_KEYS {
373            if let Some(value) = latest.metadata.get(*key) {
374                session.metadata.insert((*key).to_string(), value.clone());
375            } else {
376                session.metadata.remove(*key);
377            }
378        }
379        session.metadata_version = latest.metadata_version;
380    }
381}
382
383// ── Free merge-save function ──────────────────────────────────────────
384
385/// Save a session while preserving any concurrent UI edits to the
386/// authoritative metadata group.
387///
388/// Behaviour: if the on-disk session has `metadata_version >=
389/// session.metadata_version`, the on-disk `title`, `title_version`, `pinned`
390/// and `metadata_version` overwrite the in-memory values before writing.
391///
392/// This is the stateless variant (no per-session lock). Prefer
393/// [`LockedSessionStore::merge_save_runtime`] for server-side paths where an
394/// authoritative writer may race with this save.
395pub async fn merge_save_session(
396    storage: &Arc<dyn Storage>,
397    session: &mut Session,
398) -> std::io::Result<()> {
399    merge_authoritative_metadata_into_stale(storage, session).await;
400    storage.save_session(session).await
401}
402
403// ── Tests ─────────────────────────────────────────────────────────────
404
405#[cfg(test)]
406mod tests {
407    use super::*;
408    use crate::v2::SessionStoreV2;
409    use bamboo_domain::session::types::Session;
410
411    async fn make_storage() -> (tempfile::TempDir, Arc<dyn Storage>) {
412        let temp = tempfile::tempdir().unwrap();
413        let storage = SessionStoreV2::new(temp.path().to_path_buf())
414            .await
415            .expect("storage init");
416        (temp, Arc::new(storage) as Arc<dyn Storage>)
417    }
418
419    fn fresh(id: &str) -> Session {
420        Session::new(id.to_string(), "test-model".to_string())
421    }
422
423    // ── update_runtime_config: config patches must never clobber messages ──
424
425    #[tokio::test]
426    async fn update_runtime_config_preserves_concurrently_appended_messages() {
427        use bamboo_domain::session::types::Message;
428        use bamboo_domain::ReasoningEffort;
429
430        let (_temp, storage) = make_storage().await;
431        let store = LockedSessionStore::new(storage.clone());
432        let session_id = "cfg-preserve";
433
434        // Persisted baseline: one user + one assistant turn.
435        let mut initial = fresh(session_id);
436        initial.add_message(Message::user("hello"));
437        initial.add_message(Message::assistant("hi", None));
438        storage.save_session(&initial).await.unwrap();
439
440        // Simulate `POST /chat` appending a new user message to disk.
441        let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
442        after_chat.add_message(Message::user("second question"));
443        storage.save_session(&after_chat).await.unwrap();
444        assert_eq!(after_chat.messages.len(), 3);
445
446        // A config-only patch must load the freshest session and preserve the
447        // appended message (this is the regression that broke message sending on
448        // existing sessions).
449        let updated = store
450            .update_runtime_config(session_id, |s| {
451                s.reasoning_effort = Some(ReasoningEffort::Max);
452            })
453            .await
454            .unwrap()
455            .expect("session exists");
456
457        assert_eq!(updated.reasoning_effort, Some(ReasoningEffort::Max));
458        assert_eq!(
459            updated.messages.len(),
460            3,
461            "config patch must not revert a concurrently-appended message"
462        );
463
464        let on_disk = storage.load_session(session_id).await.unwrap().unwrap();
465        assert_eq!(on_disk.messages.len(), 3);
466        assert_eq!(on_disk.reasoning_effort, Some(ReasoningEffort::Max));
467    }
468
469    #[tokio::test]
470    async fn update_runtime_config_returns_none_for_missing_session() {
471        use bamboo_domain::ReasoningEffort;
472
473        let (_temp, storage) = make_storage().await;
474        let store = LockedSessionStore::new(storage);
475        let result = store
476            .update_runtime_config("does-not-exist", |s| {
477                s.reasoning_effort = Some(ReasoningEffort::Low);
478            })
479            .await
480            .unwrap();
481        assert!(result.is_none());
482    }
483
484    #[tokio::test]
485    async fn merge_save_runtime_overwrites_messages_from_stale_snapshot() {
486        // Characterization of the bug that motivated `update_runtime_config`:
487        // `merge_save_runtime` writes the caller's `messages` verbatim, so a
488        // stale snapshot reverts a concurrent append. Config-only writers must
489        // therefore use `update_runtime_config`, never `merge_save_runtime`.
490        use bamboo_domain::session::types::Message;
491
492        let (_temp, storage) = make_storage().await;
493        let store = LockedSessionStore::new(storage.clone());
494        let session_id = "stale-clobber";
495
496        // A handler loads the session (1 message) …
497        let mut baseline = fresh(session_id);
498        baseline.add_message(Message::user("hello"));
499        storage.save_session(&baseline).await.unwrap();
500        let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
501
502        // … then `POST /chat` appends a second message to disk …
503        let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
504        after_chat.add_message(Message::user("second"));
505        storage.save_session(&after_chat).await.unwrap();
506        assert_eq!(
507            storage
508                .load_session(session_id)
509                .await
510                .unwrap()
511                .unwrap()
512                .messages
513                .len(),
514            2
515        );
516
517        // … and the stale handler saves via merge_save_runtime -> append reverted.
518        store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
519        let after = storage.load_session(session_id).await.unwrap().unwrap();
520        assert_eq!(
521            after.messages.len(),
522            1,
523            "merge_save_runtime clobbers concurrent appends — this is why config patches must use update_runtime_config"
524        );
525    }
526
527    #[tokio::test]
528    async fn merge_save_runtime_preserves_disk_authoritative_metadata_with_single_load() {
529        // Regression guard for the single-load refactor of `merge_save_runtime`:
530        // it must STILL pull the authoritative metadata group (title / pinned /
531        // metadata_version) from the freshest on-disk copy when disk's
532        // metadata_version >= the in-memory one, even though it now reads disk
533        // only once.
534        let (_temp, storage) = make_storage().await;
535        let store = LockedSessionStore::new(storage.clone());
536        let session_id = "runtime-merge-meta";
537
538        // Baseline persisted by a runtime writer (metadata_version 0).
539        let mut baseline = fresh(session_id);
540        baseline.title = "Auto Title".to_string();
541        baseline.metadata_version = 0;
542        storage.save_session(&baseline).await.unwrap();
543
544        // A stale runtime snapshot (still metadata_version 0, old title).
545        let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
546
547        // An authoritative UI rename bumps metadata_version on disk.
548        let mut renamed = storage.load_session(session_id).await.unwrap().unwrap();
549        renamed.title = "User Renamed".to_string();
550        renamed.title_version = 1;
551        renamed.pinned = true;
552        renamed.metadata_version = 1;
553        store.commit_metadata(&renamed).await.unwrap();
554
555        // The stale runtime writer saves: it must adopt the disk title/pinned.
556        stale_snapshot.title = "Auto Title".to_string();
557        store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
558
559        let after = storage.load_session(session_id).await.unwrap().unwrap();
560        assert_eq!(after.title, "User Renamed");
561        assert!(after.pinned);
562        assert_eq!(after.metadata_version, 1);
563        // And the in-memory copy was corrected by the merge too.
564        assert_eq!(stale_snapshot.title, "User Renamed");
565        assert_eq!(stale_snapshot.metadata_version, 1);
566    }
567
568    #[tokio::test]
569    async fn merge_save_runtime_preserves_durable_workflow_run_index_from_stale_runner() {
570        let (_temp, storage) = make_storage().await;
571        let store = LockedSessionStore::new(storage.clone());
572        let session_id = "runtime-workflow-run-index";
573
574        let baseline = fresh(session_id);
575        storage.save_session(&baseline).await.unwrap();
576        let mut stale_runner = storage.load_session(session_id).await.unwrap().unwrap();
577
578        store
579            .update_runtime_config(session_id, |session| {
580                session.metadata.insert(
581                    "workflow.run_ids.v1".to_string(),
582                    r#"["http-started-run"]"#.to_string(),
583                );
584            })
585            .await
586            .unwrap()
587            .expect("session exists");
588
589        store.merge_save_runtime(&mut stale_runner).await.unwrap();
590
591        assert_eq!(
592            stale_runner
593                .metadata
594                .get("workflow.run_ids.v1")
595                .map(String::as_str),
596            Some(r#"["http-started-run"]"#)
597        );
598        let durable = storage.load_session(session_id).await.unwrap().unwrap();
599        assert_eq!(
600            durable
601                .metadata
602                .get("workflow.run_ids.v1")
603                .map(String::as_str),
604            Some(r#"["http-started-run"]"#)
605        );
606    }
607
608    // #540: a running loop's `merge_save_runtime` (carrying the run-start bypass
609    // value) must NOT revert a concurrent mid-run `PATCH /sessions
610    // {bypass_permissions}` write on disk — disk is the authoritative writer.
611    #[tokio::test]
612    async fn merge_save_runtime_adopts_disk_bypass_permissions() {
613        use bamboo_domain::AgentRuntimeState;
614
615        let (_temp, storage) = make_storage().await;
616        let store = LockedSessionStore::new(storage.clone());
617        let session_id = "runtime-bypass";
618
619        // Baseline persisted with bypass OFF.
620        let baseline = fresh(session_id);
621        storage.save_session(&baseline).await.unwrap();
622
623        // The running loop holds a snapshot with bypass OFF (run-start value).
624        let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
625        loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
626
627        // A concurrent PATCH flips bypass ON on disk (via update_runtime_config).
628        store
629            .update_runtime_config(session_id, |s| {
630                s.agent_runtime_state
631                    .get_or_insert_with(AgentRuntimeState::default)
632                    .bypass_permissions = true;
633            })
634            .await
635            .unwrap()
636            .expect("session exists");
637
638        // The loop saves its stale snapshot: it must adopt disk's ON value, not
639        // revert to OFF.
640        store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
641
642        let after = storage.load_session(session_id).await.unwrap().unwrap();
643        assert!(
644            after
645                .agent_runtime_state
646                .as_ref()
647                .is_some_and(|s| s.bypass_permissions),
648            "disk bypass=ON must survive a stale runtime save (#540)"
649        );
650        // The in-memory copy is corrected too.
651        assert!(loop_snapshot
652            .agent_runtime_state
653            .as_ref()
654            .is_some_and(|s| s.bypass_permissions));
655    }
656
657    // The reverse direction: a PATCH turning bypass OFF must also stick against
658    // a stale loop snapshot that still has it ON.
659    #[tokio::test]
660    async fn merge_save_runtime_adopts_disk_bypass_off() {
661        use bamboo_domain::AgentRuntimeState;
662
663        let (_temp, storage) = make_storage().await;
664        let store = LockedSessionStore::new(storage.clone());
665        let session_id = "runtime-bypass-off";
666
667        // Baseline persisted with bypass ON.
668        let mut baseline = fresh(session_id);
669        let mut on_state = AgentRuntimeState::default();
670        on_state.bypass_permissions = true;
671        baseline.agent_runtime_state = Some(on_state);
672        storage.save_session(&baseline).await.unwrap();
673
674        // Loop snapshot still ON.
675        let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
676
677        // PATCH flips OFF on disk.
678        store
679            .update_runtime_config(session_id, |s| {
680                s.agent_runtime_state
681                    .get_or_insert_with(AgentRuntimeState::default)
682                    .bypass_permissions = false;
683            })
684            .await
685            .unwrap()
686            .expect("session exists");
687
688        store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
689
690        let after = storage.load_session(session_id).await.unwrap().unwrap();
691        assert!(
692            !after
693                .agent_runtime_state
694                .as_ref()
695                .is_some_and(|s| s.bypass_permissions),
696            "disk bypass=OFF must survive a stale runtime save (#540)"
697        );
698    }
699
700    // #540 review: the authoritative flag writer (#74 child-reseed) must NOT be
701    // reverted by the disk-wins protection — its in-memory value persists as-is.
702    #[tokio::test]
703    async fn save_runtime_authoritative_flags_persists_in_memory_bypass() {
704        use bamboo_domain::AgentRuntimeState;
705
706        let (_temp, storage) = make_storage().await;
707        let store = LockedSessionStore::new(storage.clone());
708        let session_id = "child-reseed";
709
710        // Child on disk has bypass ON (created under a bypassed parent).
711        let mut baseline = fresh(session_id);
712        let mut on_state = AgentRuntimeState::default();
713        on_state.bypass_permissions = true;
714        baseline.agent_runtime_state = Some(on_state);
715        storage.save_session(&baseline).await.unwrap();
716
717        // Parent re-seeds the reused child to OFF (parent flipped bypass off),
718        // loading the child then setting the flag in memory.
719        let mut child = storage.load_session(session_id).await.unwrap().unwrap();
720        child
721            .agent_runtime_state
722            .get_or_insert_with(AgentRuntimeState::default)
723            .bypass_permissions = false;
724
725        // Authoritative write must persist OFF, not adopt the disk's stale ON.
726        store
727            .save_runtime_authoritative_flags(&mut child)
728            .await
729            .unwrap();
730
731        let after = storage.load_session(session_id).await.unwrap().unwrap();
732        assert!(
733            !after
734                .agent_runtime_state
735                .as_ref()
736                .is_some_and(|s| s.bypass_permissions),
737            "authoritative re-seed of bypass=OFF must persist, not be reverted (#540/#74)"
738        );
739    }
740
741    // A disk copy lacking runtime state must not force the in-memory bypass OFF.
742    #[tokio::test]
743    async fn merge_save_runtime_leaves_bypass_when_disk_has_no_runtime_state() {
744        use bamboo_domain::AgentRuntimeState;
745
746        let (_temp, storage) = make_storage().await;
747        let store = LockedSessionStore::new(storage.clone());
748        let session_id = "no-runtime-state";
749
750        // Disk copy with NO agent_runtime_state.
751        let baseline = fresh(session_id);
752        assert!(baseline.agent_runtime_state.is_none());
753        storage.save_session(&baseline).await.unwrap();
754
755        // A running loop legitimately carries bypass ON in memory.
756        let mut running = storage.load_session(session_id).await.unwrap().unwrap();
757        let mut on_state = AgentRuntimeState::default();
758        on_state.bypass_permissions = true;
759        running.agent_runtime_state = Some(on_state);
760
761        store.merge_save_runtime(&mut running).await.unwrap();
762
763        assert!(
764            running
765                .agent_runtime_state
766                .as_ref()
767                .is_some_and(|s| s.bypass_permissions),
768            "a runtime-state-less disk copy must not force bypass OFF (#540)"
769        );
770    }
771
772    // ── Free-function merge tests (updated for metadata-group) ──────
773
774    #[tokio::test]
775    async fn merge_preserves_disk_title_when_versions_equal() {
776        let (_temp, storage) = make_storage().await;
777        let session_id = "merge-equal";
778
779        let mut on_disk = fresh(session_id);
780        on_disk.title = "User Set This".to_string();
781        on_disk.title_version = 0;
782        on_disk.metadata_version = 0;
783        storage.save_session(&on_disk).await.unwrap();
784
785        let mut runtime_copy = fresh(session_id);
786        runtime_copy.title = "Stale Default".to_string();
787        runtime_copy.title_version = 0;
788        runtime_copy.metadata_version = 0;
789        runtime_copy.messages = vec![];
790
791        merge_save_session(&storage, &mut runtime_copy)
792            .await
793            .unwrap();
794
795        let after = storage.load_session(session_id).await.unwrap().unwrap();
796        assert_eq!(after.title, "User Set This");
797        assert_eq!(after.title_version, 0);
798        assert_eq!(runtime_copy.title, "User Set This");
799    }
800
801    #[tokio::test]
802    async fn merge_preserves_disk_when_disk_version_higher() {
803        let (_temp, storage) = make_storage().await;
804        let session_id = "merge-higher";
805
806        let mut on_disk = fresh(session_id);
807        on_disk.title = "User Title v3".to_string();
808        on_disk.title_version = 3;
809        on_disk.metadata_version = 5;
810        storage.save_session(&on_disk).await.unwrap();
811
812        let mut runtime_copy = fresh(session_id);
813        runtime_copy.title = "Stale".to_string();
814        runtime_copy.title_version = 1;
815        runtime_copy.metadata_version = 0;
816
817        merge_save_session(&storage, &mut runtime_copy)
818            .await
819            .unwrap();
820
821        let after = storage.load_session(session_id).await.unwrap().unwrap();
822        assert_eq!(after.title, "User Title v3");
823        assert_eq!(after.title_version, 3);
824        assert_eq!(after.metadata_version, 5);
825    }
826
827    #[tokio::test]
828    async fn merge_now_preserves_disk_pinned_in_metadata_group() {
829        let (_temp, storage) = make_storage().await;
830        let session_id = "pinned-merge";
831
832        let mut on_disk = fresh(session_id);
833        on_disk.pinned = true;
834        on_disk.metadata_version = 2;
835        storage.save_session(&on_disk).await.unwrap();
836
837        let mut runtime_copy = fresh(session_id);
838        runtime_copy.pinned = false;
839        runtime_copy.metadata_version = 0;
840
841        merge_save_session(&storage, &mut runtime_copy)
842            .await
843            .unwrap();
844
845        let after = storage.load_session(session_id).await.unwrap().unwrap();
846        assert!(
847            after.pinned,
848            "disk pinned=true should win over runtime false"
849        );
850        assert_eq!(after.metadata_version, 2);
851    }
852
853    #[tokio::test]
854    async fn merge_keeps_in_memory_when_session_version_higher() {
855        let (_temp, storage) = make_storage().await;
856        let session_id = "merge-bumped";
857
858        let mut on_disk = fresh(session_id);
859        on_disk.title = "Old".to_string();
860        on_disk.title_version = 1;
861        on_disk.metadata_version = 3;
862        storage.save_session(&on_disk).await.unwrap();
863
864        let mut authoritative_copy = fresh(session_id);
865        authoritative_copy.title = "New Authoritative".to_string();
866        authoritative_copy.title_version = 2;
867        authoritative_copy.metadata_version = 4;
868        authoritative_copy.pinned = true;
869
870        merge_save_session(&storage, &mut authoritative_copy)
871            .await
872            .unwrap();
873
874        let after = storage.load_session(session_id).await.unwrap().unwrap();
875        assert_eq!(after.title, "New Authoritative");
876        assert_eq!(after.title_version, 2);
877        assert_eq!(after.metadata_version, 4);
878        assert!(after.pinned);
879    }
880
881    #[tokio::test]
882    async fn merge_keeps_runtime_messages_when_disk_only_changed_metadata() {
883        let (_temp, storage) = make_storage().await;
884        let session_id = "merge-messages";
885
886        let mut on_disk = fresh(session_id);
887        on_disk.title = "Fresh Title".to_string();
888        on_disk.title_version = 2;
889        on_disk.metadata_version = 5;
890        storage.save_session(&on_disk).await.unwrap();
891
892        let mut runtime_copy = fresh(session_id);
893        runtime_copy.title = "Stale".to_string();
894        runtime_copy.metadata_version = 0;
895        runtime_copy.messages = vec![bamboo_domain::session::types::Message {
896            role: bamboo_domain::session::types::Role::User,
897            content: "keep me".to_string(),
898            id: "msg-1".to_string(),
899            created_at: chrono::Utc::now(),
900            reasoning: None,
901            reasoning_signature: None,
902            content_parts: None,
903            image_ocr: None,
904            phase: None,
905            tool_calls: None,
906            tool_call_id: None,
907            tool_success: None,
908            compressed: false,
909            compressed_by_event_id: None,
910            never_compress: false,
911            compression_level: 0,
912            metadata: None,
913        }];
914
915        merge_save_session(&storage, &mut runtime_copy)
916            .await
917            .unwrap();
918
919        let after = storage.load_session(session_id).await.unwrap().unwrap();
920        assert_eq!(after.title, "Fresh Title");
921        assert_eq!(after.metadata_version, 5);
922        assert_eq!(after.messages.len(), 1);
923        assert_eq!(after.messages[0].content, "keep me");
924    }
925
926    // ── LockedSessionStore tests ────────────────────────────────────
927
928    #[tokio::test]
929    async fn locked_merge_save_runtime_serialises_concurrent_writes() {
930        let (_temp, storage) = make_storage().await;
931        let store = Arc::new(LockedSessionStore::new(storage));
932        let session_id = "lock-serial".to_string();
933
934        // Seed with base version.
935        let base = fresh(&session_id);
936        store.storage().save_session(&base).await.unwrap();
937
938        // Two concurrent authorised writers each bump and commit.
939        // We'll simulate via clone-and-bump-then-commit.
940        let store_a = store.clone();
941        let store_b = store.clone();
942        let sid_a = session_id.clone();
943        let sid_b = session_id.clone();
944
945        let a = tokio::spawn(async move {
946            let _guard = store_a.acquire_lock(&sid_a).await;
947            let mut s = store_a
948                .storage()
949                .load_session(&sid_a)
950                .await
951                .unwrap()
952                .unwrap();
953            s.title = "Writer A".to_string();
954            s.title_version = s.title_version.saturating_add(1);
955            s.metadata_version = s.metadata_version.saturating_add(1);
956            s.updated_at = chrono::Utc::now();
957            store_a.storage().save_session(&s).await.unwrap();
958            s.title_version
959        });
960
961        // Tiny yield so A goes first.
962        tokio::time::sleep(std::time::Duration::from_millis(10)).await;
963
964        let b = tokio::spawn(async move {
965            let _guard = store_b.acquire_lock(&sid_b).await;
966            let mut s = store_b
967                .storage()
968                .load_session(&sid_b)
969                .await
970                .unwrap()
971                .unwrap();
972            s.title = "Writer B".to_string();
973            s.title_version = s.title_version.saturating_add(1);
974            s.metadata_version = s.metadata_version.saturating_add(1);
975            s.updated_at = chrono::Utc::now();
976            store_b.storage().save_session(&s).await.unwrap();
977            s.title_version
978        });
979
980        let (ver_a, ver_b) = tokio::join!(a, b);
981        let final_s = store
982            .storage()
983            .load_session(&session_id)
984            .await
985            .unwrap()
986            .unwrap();
987        assert!(
988            ver_a.unwrap() != ver_b.unwrap(),
989            "concurrent writers must produce distinct versions"
990        );
991        assert_eq!(final_s.metadata_version, 2);
992    }
993
994    #[tokio::test]
995    async fn commit_metadata_is_plain_save_inside_lock() {
996        let (_temp, storage) = make_storage().await;
997        let store = LockedSessionStore::new(storage);
998        let session_id = "commit-plain";
999
1000        let mut s = fresh(session_id);
1001        s.title = "Committed".to_string();
1002        s.metadata_version = 1;
1003        s.title_version = 2;
1004
1005        store.commit_metadata(&s).await.unwrap();
1006
1007        let after = store
1008            .storage()
1009            .load_session(session_id)
1010            .await
1011            .unwrap()
1012            .unwrap();
1013        assert_eq!(after.title, "Committed");
1014        assert_eq!(after.metadata_version, 1);
1015        assert_eq!(after.title_version, 2);
1016    }
1017
1018    // ── Self-cleaning per-session lock (issue #346) ─────────────────
1019
1020    #[tokio::test]
1021    async fn acquire_lock_self_evicts_when_no_other_holder() {
1022        let (_temp, storage) = make_storage().await;
1023        let store = LockedSessionStore::new(storage);
1024
1025        {
1026            let _guard = store.acquire_lock("solo").await;
1027            assert_eq!(store.locks.len(), 1, "entry present while the lock is held");
1028        }
1029        // Dropping the guard runs the self-cleaning `remove_if`. Without the
1030        // eviction logic this stays at 1 forever (the pre-#346 leak).
1031        assert_eq!(
1032            store.locks.len(),
1033            0,
1034            "lock entry must be evicted once released with no other holder"
1035        );
1036    }
1037
1038    #[tokio::test]
1039    async fn acquire_lock_many_distinct_ids_do_not_accumulate() {
1040        let (_temp, storage) = make_storage().await;
1041        let store = LockedSessionStore::new(storage);
1042
1043        // Serially acquire+release for 100 distinct session ids.
1044        for i in 0..100 {
1045            let _guard = store.acquire_lock(&format!("sess-{i}")).await;
1046        }
1047        assert_eq!(
1048            store.locks.len(),
1049            0,
1050            "acquiring locks for many distinct ids must not grow the map"
1051        );
1052    }
1053
1054    #[tokio::test]
1055    async fn acquire_lock_concurrent_waiter_keeps_valid_lock_and_map_drains() {
1056        use std::sync::atomic::{AtomicUsize, Ordering};
1057
1058        let (_temp, storage) = make_storage().await;
1059        let store = Arc::new(LockedSessionStore::new(storage));
1060
1061        // Tracks concurrent holders of the SAME session lock; must never exceed 1.
1062        let active = Arc::new(AtomicUsize::new(0));
1063        let max_seen = Arc::new(AtomicUsize::new(0));
1064
1065        let mut handles = Vec::new();
1066        for _ in 0..8 {
1067            let store = store.clone();
1068            let active = active.clone();
1069            let max_seen = max_seen.clone();
1070            handles.push(tokio::spawn(async move {
1071                let _guard = store.acquire_lock("contended").await;
1072                let now = active.fetch_add(1, Ordering::SeqCst) + 1;
1073                max_seen.fetch_max(now, Ordering::SeqCst);
1074                // Hold briefly so the other tasks actually queue on the mutex.
1075                tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1076                active.fetch_sub(1, Ordering::SeqCst);
1077            }));
1078        }
1079        for h in handles {
1080            h.await.unwrap();
1081        }
1082
1083        // Mutual exclusion must hold: a self-cleaning removal that raced (removed
1084        // the entry a waiter had already cloned, letting a later task create and
1085        // lock a *second* mutex for the same id) would show 2 concurrent holders.
1086        // `remove_if`'s atomic strong-count check under the shard lock prevents it.
1087        assert_eq!(
1088            max_seen.load(Ordering::SeqCst),
1089            1,
1090            "at most one holder of a given session lock at a time"
1091        );
1092        assert_eq!(
1093            store.locks.len(),
1094            0,
1095            "after all holders release, the contended entry must be fully evicted"
1096        );
1097    }
1098}