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    /// Persist an execute-boundary transcript checkpoint without allowing a
199    /// stale runner snapshot to shrink or rewrite the durable message log.
200    ///
201    /// The latest load, append-only message reconciliation, metadata merge and
202    /// save all happen while holding the same per-session lock.  Loading is
203    /// deliberately fail-closed: falling back to a blind full save when the
204    /// latest transcript cannot be read would reintroduce the SHRINK hazard
205    /// this checkpoint exists to prevent.
206    pub async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
207        let _guard = self.acquire_lock(&session.id).await;
208        let latest = self.storage.load_session(&session.id).await?;
209
210        if let Some(latest) = latest.as_ref() {
211            let incoming_count = session.messages.len();
212            let durable_count = latest.messages.len();
213            let appended = bamboo_domain::append_missing_runtime_messages(session, latest);
214            tracing::debug!(
215                "[{}] append-safe runtime checkpoint: durable={}, incoming={}, appended={}, saved={}",
216                session.id,
217                durable_count,
218                incoming_count,
219                appended,
220                session.messages.len(),
221            );
222            apply_authoritative_metadata(session, latest);
223            adopt_disk_bypass_permissions(session, latest);
224        }
225
226        self.storage.save_session(session).await
227    }
228
229    /// Like [`Self::merge_save_runtime`] but does NOT adopt the on-disk
230    /// `bypass_permissions` — the caller's in-memory value is authoritative and
231    /// persists as-is.
232    ///
233    /// For parent-side control writes to a child session (e.g. the #74
234    /// resident-reuse posture re-seed), which set the flag deliberately and must
235    /// not be reverted by the disk-wins protection meant for a running loop's
236    /// own stale saves. Still merges the authoritative metadata group.
237    pub async fn save_runtime_authoritative_flags(
238        &self,
239        session: &mut Session,
240    ) -> std::io::Result<()> {
241        self.merge_save_runtime_inner(session, false).await
242    }
243
244    async fn merge_save_runtime_inner(
245        &self,
246        session: &mut Session,
247        adopt_bypass: bool,
248    ) -> std::io::Result<()> {
249        let _guard = self.acquire_lock(&session.id).await;
250
251        // Single disk read serves BOTH the SHRINK diagnostic and the
252        // authoritative-metadata merge below. Previously this path loaded the
253        // session twice (once here, once inside the merge helper); on a parent
254        // session carrying the full conversation history that doubled the
255        // deserialization cost of every runtime save, which is the hot path
256        // during sub-agent spawn.
257        let latest = self.storage.load_session(&session.id).await.ok().flatten();
258
259        // DIAGNOSTIC: merge_save_runtime overwrites the whole `messages` array
260        // (it only merges authoritative metadata, not messages). If the incoming
261        // session is stale (fewer messages than what is already on disk), this save
262        // silently reverts a concurrent append (e.g. a just-persisted user message).
263        // Log a SHRINK warning so we can identify the stale writer.
264        let existing_message_count = latest.as_ref().map(|s| s.messages.len());
265        let incoming_message_count = session.messages.len();
266        if existing_message_count.is_some_and(|existing| existing > incoming_message_count) {
267            tracing::warn!(
268                "[{}] merge_save_runtime SHRINK: disk has {:?} messages, saving {} (last_role={:?}, updated_at={}); a stale writer is reverting a concurrent append",
269                session.id,
270                existing_message_count,
271                incoming_message_count,
272                session.messages.last().map(|m| format!("{:?}", m.role)),
273                session.updated_at,
274            );
275        } else {
276            tracing::debug!(
277                "[{}] merge_save_runtime: disk={:?} messages, saving {} (updated_at={})",
278                session.id,
279                existing_message_count,
280                incoming_message_count,
281                session.updated_at,
282            );
283        }
284
285        if let Some(latest) = latest.as_ref() {
286            apply_authoritative_metadata(session, latest);
287            // Never let a running loop's save revert a concurrent mid-run
288            // `PATCH /sessions {bypass_permissions}` flip. #540. Skipped for
289            // authoritative flag writers (`save_runtime_authoritative_flags`).
290            if adopt_bypass {
291                adopt_disk_bypass_permissions(session, latest);
292            }
293        }
294        self.storage.save_session(session).await
295    }
296
297    /// Apply a config-only mutation to a session without ever clobbering its
298    /// `messages` (or other concurrently-written state).
299    ///
300    /// Unlike [`Self::merge_save_runtime`], the caller does NOT pass a session
301    /// snapshot. Instead this loads the **latest** session from storage *inside*
302    /// the per-session lock, applies `mutate` (intended for small config fields
303    /// like `model_ref` / `reasoning_effort`), and saves. Because the load and
304    /// save both happen under the lock, a concurrent append (e.g. `POST /chat`
305    /// adding a user message) can never be reverted by this write.
306    ///
307    /// Returns the saved session, or `None` if it does not exist.
308    pub async fn update_runtime_config<F>(
309        &self,
310        session_id: &str,
311        mutate: F,
312    ) -> std::io::Result<Option<Session>>
313    where
314        F: FnOnce(&mut Session),
315    {
316        let _guard = self.acquire_lock(session_id).await;
317        let Some(mut session) = self.storage.load_session(session_id).await? else {
318            return Ok(None);
319        };
320        mutate(&mut session);
321        self.storage.save_session(&session).await?;
322        Ok(Some(session))
323    }
324}
325
326/// Infrastructure implementation of the domain runtime-persistence port.
327/// Server should assemble this as `Arc<dyn RuntimeSessionPersistence>` and must
328/// not define a separate adapter layer for the same behavior.
329#[async_trait::async_trait]
330impl RuntimeSessionPersistence for LockedSessionStore {
331    async fn save_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
332        self.merge_save_runtime(session).await
333    }
334
335    async fn checkpoint_runtime_session(&self, session: &mut Session) -> std::io::Result<()> {
336        LockedSessionStore::checkpoint_runtime_session(self, session).await
337    }
338
339    async fn load_runtime_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
340        self.storage.load_session(session_id).await
341    }
342}
343
344// ── Internal merge helper ─────────────────────────────────────────────
345
346/// Re-read the on-disk session and, when the disk copy carries a
347/// `metadata_version >= session.metadata_version`, overwrite the in-memory
348/// authoritative metadata fields with the disk values.
349///
350/// This is the core staleness-correction: non-authoritative writers call it
351/// before saving so they don't accidentally revert a concurrent UI edit.
352async fn merge_authoritative_metadata_into_stale(
353    storage: &Arc<dyn Storage>,
354    session: &mut Session,
355) {
356    if let Ok(Some(latest)) = storage.load_session(&session.id).await {
357        apply_authoritative_metadata(session, &latest);
358        adopt_disk_bypass_permissions(session, &latest);
359    }
360}
361
362/// Adopt the on-disk `agent_runtime_state.bypass_permissions` into the session
363/// about to be saved.
364///
365/// `PATCH /sessions {bypass_permissions}` is the SOLE authoritative writer of
366/// this flag (a running loop only carries it forward from run start). Without
367/// this, a runtime save from an in-flight run — which holds the run-start value
368/// — silently reverts a concurrent mid-run flip on disk. Unlike the metadata
369/// group this is NOT version-gated: the PATCH writes via `update_runtime_config`,
370/// which does not bump `metadata_version`. #540.
371fn adopt_disk_bypass_permissions(session: &mut Session, latest: &Session) {
372    // A disk copy with NO runtime state at all carries no authoritative bypass
373    // value — treat it as "unknown" and leave the in-memory flag untouched,
374    // rather than forcing it OFF (which would silently disable a legitimately
375    // bypassed run on any backend/path that doesn't round-trip the field). #540.
376    let Some(disk_bypass) = latest
377        .agent_runtime_state
378        .as_ref()
379        .map(|state| state.bypass_permissions)
380    else {
381        return;
382    };
383    match session.agent_runtime_state.as_mut() {
384        Some(state) => state.bypass_permissions = disk_bypass,
385        // No runtime state in memory and disk says "off" → nothing to adopt;
386        // avoid allocating a default state just to store `false`.
387        None if disk_bypass => {
388            session
389                .agent_runtime_state
390                .get_or_insert_with(bamboo_domain::AgentRuntimeState::default)
391                .bypass_permissions = true;
392        }
393        None => {}
394    }
395}
396
397/// Pure merge step: given a freshly-loaded on-disk copy, overwrite the
398/// in-memory authoritative metadata group when disk's `metadata_version` is at
399/// least the in-memory one. Split out so callers that have already loaded the
400/// disk copy (e.g. [`LockedSessionStore::merge_save_runtime`]) don't pay for a
401/// second read.
402fn apply_authoritative_metadata(session: &mut Session, latest: &Session) {
403    if latest.metadata_version >= session.metadata_version {
404        session.title = latest.title.clone();
405        session.title_version = latest.title_version;
406        session.pinned = latest.pinned;
407        for key in AUTHORITATIVE_METADATA_KEYS {
408            if let Some(value) = latest.metadata.get(*key) {
409                session.metadata.insert((*key).to_string(), value.clone());
410            } else {
411                session.metadata.remove(*key);
412            }
413        }
414        session.metadata_version = latest.metadata_version;
415    }
416}
417
418// ── Free merge-save function ──────────────────────────────────────────
419
420/// Save a session while preserving any concurrent UI edits to the
421/// authoritative metadata group.
422///
423/// Behaviour: if the on-disk session has `metadata_version >=
424/// session.metadata_version`, the on-disk `title`, `title_version`, `pinned`
425/// and `metadata_version` overwrite the in-memory values before writing.
426///
427/// This is the stateless variant (no per-session lock). Prefer
428/// [`LockedSessionStore::merge_save_runtime`] for server-side paths where an
429/// authoritative writer may race with this save.
430pub async fn merge_save_session(
431    storage: &Arc<dyn Storage>,
432    session: &mut Session,
433) -> std::io::Result<()> {
434    merge_authoritative_metadata_into_stale(storage, session).await;
435    storage.save_session(session).await
436}
437
438// ── Tests ─────────────────────────────────────────────────────────────
439
440#[cfg(test)]
441mod tests {
442    use super::*;
443    use crate::v2::SessionStoreV2;
444    use bamboo_domain::session::types::Session;
445
446    async fn make_storage() -> (tempfile::TempDir, Arc<dyn Storage>) {
447        let temp = tempfile::tempdir().unwrap();
448        let storage = SessionStoreV2::new(temp.path().to_path_buf())
449            .await
450            .expect("storage init");
451        (temp, Arc::new(storage) as Arc<dyn Storage>)
452    }
453
454    fn fresh(id: &str) -> Session {
455        Session::new(id.to_string(), "test-model".to_string())
456    }
457
458    // ── update_runtime_config: config patches must never clobber messages ──
459
460    #[tokio::test]
461    async fn update_runtime_config_preserves_concurrently_appended_messages() {
462        use bamboo_domain::session::types::Message;
463        use bamboo_domain::ReasoningEffort;
464
465        let (_temp, storage) = make_storage().await;
466        let store = LockedSessionStore::new(storage.clone());
467        let session_id = "cfg-preserve";
468
469        // Persisted baseline: one user + one assistant turn.
470        let mut initial = fresh(session_id);
471        initial.add_message(Message::user("hello"));
472        initial.add_message(Message::assistant("hi", None));
473        storage.save_session(&initial).await.unwrap();
474
475        // Simulate `POST /chat` appending a new user message to disk.
476        let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
477        after_chat.add_message(Message::user("second question"));
478        storage.save_session(&after_chat).await.unwrap();
479        assert_eq!(after_chat.messages.len(), 3);
480
481        // A config-only patch must load the freshest session and preserve the
482        // appended message (this is the regression that broke message sending on
483        // existing sessions).
484        let updated = store
485            .update_runtime_config(session_id, |s| {
486                s.reasoning_effort = Some(ReasoningEffort::Max);
487            })
488            .await
489            .unwrap()
490            .expect("session exists");
491
492        assert_eq!(updated.reasoning_effort, Some(ReasoningEffort::Max));
493        assert_eq!(
494            updated.messages.len(),
495            3,
496            "config patch must not revert a concurrently-appended message"
497        );
498
499        let on_disk = storage.load_session(session_id).await.unwrap().unwrap();
500        assert_eq!(on_disk.messages.len(), 3);
501        assert_eq!(on_disk.reasoning_effort, Some(ReasoningEffort::Max));
502    }
503
504    #[tokio::test]
505    async fn update_runtime_config_returns_none_for_missing_session() {
506        use bamboo_domain::ReasoningEffort;
507
508        let (_temp, storage) = make_storage().await;
509        let store = LockedSessionStore::new(storage);
510        let result = store
511            .update_runtime_config("does-not-exist", |s| {
512                s.reasoning_effort = Some(ReasoningEffort::Low);
513            })
514            .await
515            .unwrap();
516        assert!(result.is_none());
517    }
518
519    #[tokio::test]
520    async fn merge_save_runtime_overwrites_messages_from_stale_snapshot() {
521        // Characterization of the bug that motivated `update_runtime_config`:
522        // `merge_save_runtime` writes the caller's `messages` verbatim, so a
523        // stale snapshot reverts a concurrent append. Config-only writers must
524        // therefore use `update_runtime_config`, never `merge_save_runtime`.
525        use bamboo_domain::session::types::Message;
526
527        let (_temp, storage) = make_storage().await;
528        let store = LockedSessionStore::new(storage.clone());
529        let session_id = "stale-clobber";
530
531        // A handler loads the session (1 message) …
532        let mut baseline = fresh(session_id);
533        baseline.add_message(Message::user("hello"));
534        storage.save_session(&baseline).await.unwrap();
535        let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
536
537        // … then `POST /chat` appends a second message to disk …
538        let mut after_chat = storage.load_session(session_id).await.unwrap().unwrap();
539        after_chat.add_message(Message::user("second"));
540        storage.save_session(&after_chat).await.unwrap();
541        assert_eq!(
542            storage
543                .load_session(session_id)
544                .await
545                .unwrap()
546                .unwrap()
547                .messages
548                .len(),
549            2
550        );
551
552        // … and the stale handler saves via merge_save_runtime -> append reverted.
553        store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
554        let after = storage.load_session(session_id).await.unwrap().unwrap();
555        assert_eq!(
556            after.messages.len(),
557            1,
558            "merge_save_runtime clobbers concurrent appends — this is why config patches must use update_runtime_config"
559        );
560    }
561
562    #[tokio::test]
563    async fn checkpoint_runtime_session_preserves_disk_suffix_and_appends_live_messages() {
564        use bamboo_domain::session::types::Message;
565
566        let (_temp, storage) = make_storage().await;
567        let store = LockedSessionStore::new(storage.clone());
568        let session_id = "checkpoint-no-shrink";
569
570        let mut baseline = fresh(session_id);
571        baseline.add_message(Message::user("base"));
572        storage.save_session(&baseline).await.unwrap();
573        let mut runner_snapshot = baseline.clone();
574
575        let mut durable = baseline;
576        let mut disk_only = Message::user("concurrent injected message");
577        disk_only.id = "disk-only".to_string();
578        durable.add_message(disk_only);
579        storage.save_session(&durable).await.unwrap();
580
581        let mut live_only = Message::assistant("partial runner output", None);
582        live_only.id = "live-only".to_string();
583        runner_snapshot.add_message(live_only);
584
585        store
586            .checkpoint_runtime_session(&mut runner_snapshot)
587            .await
588            .unwrap();
589
590        let saved = storage.load_session(session_id).await.unwrap().unwrap();
591        let ids = saved
592            .messages
593            .iter()
594            .map(|message| message.id.as_str())
595            .collect::<Vec<_>>();
596        assert_eq!(
597            ids,
598            vec![durable.messages[0].id.as_str(), "disk-only", "live-only"]
599        );
600        assert_eq!(runner_snapshot.messages.len(), saved.messages.len());
601        assert_eq!(runner_snapshot.messages[1].id, saved.messages[1].id);
602        assert_eq!(runner_snapshot.messages[2].id, saved.messages[2].id);
603        assert_eq!(saved.messages[1].content, "concurrent injected message");
604        assert_eq!(saved.messages[2].content, "partial runner output");
605    }
606
607    #[tokio::test]
608    async fn merge_save_runtime_preserves_disk_authoritative_metadata_with_single_load() {
609        // Regression guard for the single-load refactor of `merge_save_runtime`:
610        // it must STILL pull the authoritative metadata group (title / pinned /
611        // metadata_version) from the freshest on-disk copy when disk's
612        // metadata_version >= the in-memory one, even though it now reads disk
613        // only once.
614        let (_temp, storage) = make_storage().await;
615        let store = LockedSessionStore::new(storage.clone());
616        let session_id = "runtime-merge-meta";
617
618        // Baseline persisted by a runtime writer (metadata_version 0).
619        let mut baseline = fresh(session_id);
620        baseline.title = "Auto Title".to_string();
621        baseline.metadata_version = 0;
622        storage.save_session(&baseline).await.unwrap();
623
624        // A stale runtime snapshot (still metadata_version 0, old title).
625        let mut stale_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
626
627        // An authoritative UI rename bumps metadata_version on disk.
628        let mut renamed = storage.load_session(session_id).await.unwrap().unwrap();
629        renamed.title = "User Renamed".to_string();
630        renamed.title_version = 1;
631        renamed.pinned = true;
632        renamed.metadata_version = 1;
633        store.commit_metadata(&renamed).await.unwrap();
634
635        // The stale runtime writer saves: it must adopt the disk title/pinned.
636        stale_snapshot.title = "Auto Title".to_string();
637        store.merge_save_runtime(&mut stale_snapshot).await.unwrap();
638
639        let after = storage.load_session(session_id).await.unwrap().unwrap();
640        assert_eq!(after.title, "User Renamed");
641        assert!(after.pinned);
642        assert_eq!(after.metadata_version, 1);
643        // And the in-memory copy was corrected by the merge too.
644        assert_eq!(stale_snapshot.title, "User Renamed");
645        assert_eq!(stale_snapshot.metadata_version, 1);
646    }
647
648    #[tokio::test]
649    async fn merge_save_runtime_preserves_durable_workflow_run_index_from_stale_runner() {
650        let (_temp, storage) = make_storage().await;
651        let store = LockedSessionStore::new(storage.clone());
652        let session_id = "runtime-workflow-run-index";
653
654        let baseline = fresh(session_id);
655        storage.save_session(&baseline).await.unwrap();
656        let mut stale_runner = storage.load_session(session_id).await.unwrap().unwrap();
657
658        store
659            .update_runtime_config(session_id, |session| {
660                session.metadata.insert(
661                    "workflow.run_ids.v1".to_string(),
662                    r#"["http-started-run"]"#.to_string(),
663                );
664            })
665            .await
666            .unwrap()
667            .expect("session exists");
668
669        store.merge_save_runtime(&mut stale_runner).await.unwrap();
670
671        assert_eq!(
672            stale_runner
673                .metadata
674                .get("workflow.run_ids.v1")
675                .map(String::as_str),
676            Some(r#"["http-started-run"]"#)
677        );
678        let durable = storage.load_session(session_id).await.unwrap().unwrap();
679        assert_eq!(
680            durable
681                .metadata
682                .get("workflow.run_ids.v1")
683                .map(String::as_str),
684            Some(r#"["http-started-run"]"#)
685        );
686    }
687
688    // #540: a running loop's `merge_save_runtime` (carrying the run-start bypass
689    // value) must NOT revert a concurrent mid-run `PATCH /sessions
690    // {bypass_permissions}` write on disk — disk is the authoritative writer.
691    #[tokio::test]
692    async fn merge_save_runtime_adopts_disk_bypass_permissions() {
693        use bamboo_domain::AgentRuntimeState;
694
695        let (_temp, storage) = make_storage().await;
696        let store = LockedSessionStore::new(storage.clone());
697        let session_id = "runtime-bypass";
698
699        // Baseline persisted with bypass OFF.
700        let baseline = fresh(session_id);
701        storage.save_session(&baseline).await.unwrap();
702
703        // The running loop holds a snapshot with bypass OFF (run-start value).
704        let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
705        loop_snapshot.agent_runtime_state = Some(AgentRuntimeState::default());
706
707        // A concurrent PATCH flips bypass ON on disk (via update_runtime_config).
708        store
709            .update_runtime_config(session_id, |s| {
710                s.agent_runtime_state
711                    .get_or_insert_with(AgentRuntimeState::default)
712                    .bypass_permissions = true;
713            })
714            .await
715            .unwrap()
716            .expect("session exists");
717
718        // The loop saves its stale snapshot: it must adopt disk's ON value, not
719        // revert to OFF.
720        store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
721
722        let after = storage.load_session(session_id).await.unwrap().unwrap();
723        assert!(
724            after
725                .agent_runtime_state
726                .as_ref()
727                .is_some_and(|s| s.bypass_permissions),
728            "disk bypass=ON must survive a stale runtime save (#540)"
729        );
730        // The in-memory copy is corrected too.
731        assert!(loop_snapshot
732            .agent_runtime_state
733            .as_ref()
734            .is_some_and(|s| s.bypass_permissions));
735    }
736
737    // The reverse direction: a PATCH turning bypass OFF must also stick against
738    // a stale loop snapshot that still has it ON.
739    #[tokio::test]
740    async fn merge_save_runtime_adopts_disk_bypass_off() {
741        use bamboo_domain::AgentRuntimeState;
742
743        let (_temp, storage) = make_storage().await;
744        let store = LockedSessionStore::new(storage.clone());
745        let session_id = "runtime-bypass-off";
746
747        // Baseline persisted with bypass ON.
748        let mut baseline = fresh(session_id);
749        let mut on_state = AgentRuntimeState::default();
750        on_state.bypass_permissions = true;
751        baseline.agent_runtime_state = Some(on_state);
752        storage.save_session(&baseline).await.unwrap();
753
754        // Loop snapshot still ON.
755        let mut loop_snapshot = storage.load_session(session_id).await.unwrap().unwrap();
756
757        // PATCH flips OFF on disk.
758        store
759            .update_runtime_config(session_id, |s| {
760                s.agent_runtime_state
761                    .get_or_insert_with(AgentRuntimeState::default)
762                    .bypass_permissions = false;
763            })
764            .await
765            .unwrap()
766            .expect("session exists");
767
768        store.merge_save_runtime(&mut loop_snapshot).await.unwrap();
769
770        let after = storage.load_session(session_id).await.unwrap().unwrap();
771        assert!(
772            !after
773                .agent_runtime_state
774                .as_ref()
775                .is_some_and(|s| s.bypass_permissions),
776            "disk bypass=OFF must survive a stale runtime save (#540)"
777        );
778    }
779
780    // #540 review: the authoritative flag writer (#74 child-reseed) must NOT be
781    // reverted by the disk-wins protection — its in-memory value persists as-is.
782    #[tokio::test]
783    async fn save_runtime_authoritative_flags_persists_in_memory_bypass() {
784        use bamboo_domain::AgentRuntimeState;
785
786        let (_temp, storage) = make_storage().await;
787        let store = LockedSessionStore::new(storage.clone());
788        let session_id = "child-reseed";
789
790        // Child on disk has bypass ON (created under a bypassed parent).
791        let mut baseline = fresh(session_id);
792        let mut on_state = AgentRuntimeState::default();
793        on_state.bypass_permissions = true;
794        baseline.agent_runtime_state = Some(on_state);
795        storage.save_session(&baseline).await.unwrap();
796
797        // Parent re-seeds the reused child to OFF (parent flipped bypass off),
798        // loading the child then setting the flag in memory.
799        let mut child = storage.load_session(session_id).await.unwrap().unwrap();
800        child
801            .agent_runtime_state
802            .get_or_insert_with(AgentRuntimeState::default)
803            .bypass_permissions = false;
804
805        // Authoritative write must persist OFF, not adopt the disk's stale ON.
806        store
807            .save_runtime_authoritative_flags(&mut child)
808            .await
809            .unwrap();
810
811        let after = storage.load_session(session_id).await.unwrap().unwrap();
812        assert!(
813            !after
814                .agent_runtime_state
815                .as_ref()
816                .is_some_and(|s| s.bypass_permissions),
817            "authoritative re-seed of bypass=OFF must persist, not be reverted (#540/#74)"
818        );
819    }
820
821    // A disk copy lacking runtime state must not force the in-memory bypass OFF.
822    #[tokio::test]
823    async fn merge_save_runtime_leaves_bypass_when_disk_has_no_runtime_state() {
824        use bamboo_domain::AgentRuntimeState;
825
826        let (_temp, storage) = make_storage().await;
827        let store = LockedSessionStore::new(storage.clone());
828        let session_id = "no-runtime-state";
829
830        // Disk copy with NO agent_runtime_state.
831        let baseline = fresh(session_id);
832        assert!(baseline.agent_runtime_state.is_none());
833        storage.save_session(&baseline).await.unwrap();
834
835        // A running loop legitimately carries bypass ON in memory.
836        let mut running = storage.load_session(session_id).await.unwrap().unwrap();
837        let mut on_state = AgentRuntimeState::default();
838        on_state.bypass_permissions = true;
839        running.agent_runtime_state = Some(on_state);
840
841        store.merge_save_runtime(&mut running).await.unwrap();
842
843        assert!(
844            running
845                .agent_runtime_state
846                .as_ref()
847                .is_some_and(|s| s.bypass_permissions),
848            "a runtime-state-less disk copy must not force bypass OFF (#540)"
849        );
850    }
851
852    // ── Free-function merge tests (updated for metadata-group) ──────
853
854    #[tokio::test]
855    async fn merge_preserves_disk_title_when_versions_equal() {
856        let (_temp, storage) = make_storage().await;
857        let session_id = "merge-equal";
858
859        let mut on_disk = fresh(session_id);
860        on_disk.title = "User Set This".to_string();
861        on_disk.title_version = 0;
862        on_disk.metadata_version = 0;
863        storage.save_session(&on_disk).await.unwrap();
864
865        let mut runtime_copy = fresh(session_id);
866        runtime_copy.title = "Stale Default".to_string();
867        runtime_copy.title_version = 0;
868        runtime_copy.metadata_version = 0;
869        runtime_copy.messages = vec![];
870
871        merge_save_session(&storage, &mut runtime_copy)
872            .await
873            .unwrap();
874
875        let after = storage.load_session(session_id).await.unwrap().unwrap();
876        assert_eq!(after.title, "User Set This");
877        assert_eq!(after.title_version, 0);
878        assert_eq!(runtime_copy.title, "User Set This");
879    }
880
881    #[tokio::test]
882    async fn merge_preserves_disk_when_disk_version_higher() {
883        let (_temp, storage) = make_storage().await;
884        let session_id = "merge-higher";
885
886        let mut on_disk = fresh(session_id);
887        on_disk.title = "User Title v3".to_string();
888        on_disk.title_version = 3;
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.title_version = 1;
895        runtime_copy.metadata_version = 0;
896
897        merge_save_session(&storage, &mut runtime_copy)
898            .await
899            .unwrap();
900
901        let after = storage.load_session(session_id).await.unwrap().unwrap();
902        assert_eq!(after.title, "User Title v3");
903        assert_eq!(after.title_version, 3);
904        assert_eq!(after.metadata_version, 5);
905    }
906
907    #[tokio::test]
908    async fn merge_now_preserves_disk_pinned_in_metadata_group() {
909        let (_temp, storage) = make_storage().await;
910        let session_id = "pinned-merge";
911
912        let mut on_disk = fresh(session_id);
913        on_disk.pinned = true;
914        on_disk.metadata_version = 2;
915        storage.save_session(&on_disk).await.unwrap();
916
917        let mut runtime_copy = fresh(session_id);
918        runtime_copy.pinned = false;
919        runtime_copy.metadata_version = 0;
920
921        merge_save_session(&storage, &mut runtime_copy)
922            .await
923            .unwrap();
924
925        let after = storage.load_session(session_id).await.unwrap().unwrap();
926        assert!(
927            after.pinned,
928            "disk pinned=true should win over runtime false"
929        );
930        assert_eq!(after.metadata_version, 2);
931    }
932
933    #[tokio::test]
934    async fn merge_keeps_in_memory_when_session_version_higher() {
935        let (_temp, storage) = make_storage().await;
936        let session_id = "merge-bumped";
937
938        let mut on_disk = fresh(session_id);
939        on_disk.title = "Old".to_string();
940        on_disk.title_version = 1;
941        on_disk.metadata_version = 3;
942        storage.save_session(&on_disk).await.unwrap();
943
944        let mut authoritative_copy = fresh(session_id);
945        authoritative_copy.title = "New Authoritative".to_string();
946        authoritative_copy.title_version = 2;
947        authoritative_copy.metadata_version = 4;
948        authoritative_copy.pinned = true;
949
950        merge_save_session(&storage, &mut authoritative_copy)
951            .await
952            .unwrap();
953
954        let after = storage.load_session(session_id).await.unwrap().unwrap();
955        assert_eq!(after.title, "New Authoritative");
956        assert_eq!(after.title_version, 2);
957        assert_eq!(after.metadata_version, 4);
958        assert!(after.pinned);
959    }
960
961    #[tokio::test]
962    async fn merge_keeps_runtime_messages_when_disk_only_changed_metadata() {
963        let (_temp, storage) = make_storage().await;
964        let session_id = "merge-messages";
965
966        let mut on_disk = fresh(session_id);
967        on_disk.title = "Fresh Title".to_string();
968        on_disk.title_version = 2;
969        on_disk.metadata_version = 5;
970        storage.save_session(&on_disk).await.unwrap();
971
972        let mut runtime_copy = fresh(session_id);
973        runtime_copy.title = "Stale".to_string();
974        runtime_copy.metadata_version = 0;
975        runtime_copy.messages = vec![bamboo_domain::session::types::Message {
976            role: bamboo_domain::session::types::Role::User,
977            content: "keep me".to_string(),
978            id: "msg-1".to_string(),
979            created_at: chrono::Utc::now(),
980            reasoning: None,
981            reasoning_signature: None,
982            content_parts: None,
983            image_ocr: None,
984            phase: None,
985            tool_calls: None,
986            tool_call_id: None,
987            tool_success: None,
988            compressed: false,
989            compressed_by_event_id: None,
990            never_compress: false,
991            compression_level: 0,
992            metadata: None,
993        }];
994
995        merge_save_session(&storage, &mut runtime_copy)
996            .await
997            .unwrap();
998
999        let after = storage.load_session(session_id).await.unwrap().unwrap();
1000        assert_eq!(after.title, "Fresh Title");
1001        assert_eq!(after.metadata_version, 5);
1002        assert_eq!(after.messages.len(), 1);
1003        assert_eq!(after.messages[0].content, "keep me");
1004    }
1005
1006    // ── LockedSessionStore tests ────────────────────────────────────
1007
1008    #[tokio::test]
1009    async fn locked_merge_save_runtime_serialises_concurrent_writes() {
1010        let (_temp, storage) = make_storage().await;
1011        let store = Arc::new(LockedSessionStore::new(storage));
1012        let session_id = "lock-serial".to_string();
1013
1014        // Seed with base version.
1015        let base = fresh(&session_id);
1016        store.storage().save_session(&base).await.unwrap();
1017
1018        // Two concurrent authorised writers each bump and commit.
1019        // We'll simulate via clone-and-bump-then-commit.
1020        let store_a = store.clone();
1021        let store_b = store.clone();
1022        let sid_a = session_id.clone();
1023        let sid_b = session_id.clone();
1024
1025        let a = tokio::spawn(async move {
1026            let _guard = store_a.acquire_lock(&sid_a).await;
1027            let mut s = store_a
1028                .storage()
1029                .load_session(&sid_a)
1030                .await
1031                .unwrap()
1032                .unwrap();
1033            s.title = "Writer A".to_string();
1034            s.title_version = s.title_version.saturating_add(1);
1035            s.metadata_version = s.metadata_version.saturating_add(1);
1036            s.updated_at = chrono::Utc::now();
1037            store_a.storage().save_session(&s).await.unwrap();
1038            s.title_version
1039        });
1040
1041        // Tiny yield so A goes first.
1042        tokio::time::sleep(std::time::Duration::from_millis(10)).await;
1043
1044        let b = tokio::spawn(async move {
1045            let _guard = store_b.acquire_lock(&sid_b).await;
1046            let mut s = store_b
1047                .storage()
1048                .load_session(&sid_b)
1049                .await
1050                .unwrap()
1051                .unwrap();
1052            s.title = "Writer B".to_string();
1053            s.title_version = s.title_version.saturating_add(1);
1054            s.metadata_version = s.metadata_version.saturating_add(1);
1055            s.updated_at = chrono::Utc::now();
1056            store_b.storage().save_session(&s).await.unwrap();
1057            s.title_version
1058        });
1059
1060        let (ver_a, ver_b) = tokio::join!(a, b);
1061        let final_s = store
1062            .storage()
1063            .load_session(&session_id)
1064            .await
1065            .unwrap()
1066            .unwrap();
1067        assert!(
1068            ver_a.unwrap() != ver_b.unwrap(),
1069            "concurrent writers must produce distinct versions"
1070        );
1071        assert_eq!(final_s.metadata_version, 2);
1072    }
1073
1074    #[tokio::test]
1075    async fn commit_metadata_is_plain_save_inside_lock() {
1076        let (_temp, storage) = make_storage().await;
1077        let store = LockedSessionStore::new(storage);
1078        let session_id = "commit-plain";
1079
1080        let mut s = fresh(session_id);
1081        s.title = "Committed".to_string();
1082        s.metadata_version = 1;
1083        s.title_version = 2;
1084
1085        store.commit_metadata(&s).await.unwrap();
1086
1087        let after = store
1088            .storage()
1089            .load_session(session_id)
1090            .await
1091            .unwrap()
1092            .unwrap();
1093        assert_eq!(after.title, "Committed");
1094        assert_eq!(after.metadata_version, 1);
1095        assert_eq!(after.title_version, 2);
1096    }
1097
1098    // ── Self-cleaning per-session lock (issue #346) ─────────────────
1099
1100    #[tokio::test]
1101    async fn acquire_lock_self_evicts_when_no_other_holder() {
1102        let (_temp, storage) = make_storage().await;
1103        let store = LockedSessionStore::new(storage);
1104
1105        {
1106            let _guard = store.acquire_lock("solo").await;
1107            assert_eq!(store.locks.len(), 1, "entry present while the lock is held");
1108        }
1109        // Dropping the guard runs the self-cleaning `remove_if`. Without the
1110        // eviction logic this stays at 1 forever (the pre-#346 leak).
1111        assert_eq!(
1112            store.locks.len(),
1113            0,
1114            "lock entry must be evicted once released with no other holder"
1115        );
1116    }
1117
1118    #[tokio::test]
1119    async fn acquire_lock_many_distinct_ids_do_not_accumulate() {
1120        let (_temp, storage) = make_storage().await;
1121        let store = LockedSessionStore::new(storage);
1122
1123        // Serially acquire+release for 100 distinct session ids.
1124        for i in 0..100 {
1125            let _guard = store.acquire_lock(&format!("sess-{i}")).await;
1126        }
1127        assert_eq!(
1128            store.locks.len(),
1129            0,
1130            "acquiring locks for many distinct ids must not grow the map"
1131        );
1132    }
1133
1134    #[tokio::test]
1135    async fn acquire_lock_concurrent_waiter_keeps_valid_lock_and_map_drains() {
1136        use std::sync::atomic::{AtomicUsize, Ordering};
1137
1138        let (_temp, storage) = make_storage().await;
1139        let store = Arc::new(LockedSessionStore::new(storage));
1140
1141        // Tracks concurrent holders of the SAME session lock; must never exceed 1.
1142        let active = Arc::new(AtomicUsize::new(0));
1143        let max_seen = Arc::new(AtomicUsize::new(0));
1144
1145        let mut handles = Vec::new();
1146        for _ in 0..8 {
1147            let store = store.clone();
1148            let active = active.clone();
1149            let max_seen = max_seen.clone();
1150            handles.push(tokio::spawn(async move {
1151                let _guard = store.acquire_lock("contended").await;
1152                let now = active.fetch_add(1, Ordering::SeqCst) + 1;
1153                max_seen.fetch_max(now, Ordering::SeqCst);
1154                // Hold briefly so the other tasks actually queue on the mutex.
1155                tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1156                active.fetch_sub(1, Ordering::SeqCst);
1157            }));
1158        }
1159        for h in handles {
1160            h.await.unwrap();
1161        }
1162
1163        // Mutual exclusion must hold: a self-cleaning removal that raced (removed
1164        // the entry a waiter had already cloned, letting a later task create and
1165        // lock a *second* mutex for the same id) would show 2 concurrent holders.
1166        // `remove_if`'s atomic strong-count check under the shard lock prevents it.
1167        assert_eq!(
1168            max_seen.load(Ordering::SeqCst),
1169            1,
1170            "at most one holder of a given session lock at a time"
1171        );
1172        assert_eq!(
1173            store.locks.len(),
1174            0,
1175            "after all holders release, the contended entry must be fully evicted"
1176        );
1177    }
1178}