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