From af53e688cdf74c25fab2f4ef744e65d608dadfc4 Mon Sep 17 00:00:00 2001 From: Zixuan Chen Date: Fri, 4 Sep 2026 10:23:59 +0800 Subject: [PATCH 1/4] fix(loro-internal): bound the decoded container value cache Reading containers through JS handlers funnels through InnerStore::with_container_for_read, which pinned every decoded container value in the store for the lifetime of the doc (~4 KB per container). Walking a large document container by container retained memory superlinearly until doc.free(), trapping wasm32 at the 4 GiB limit around one million containers (loro-dev/loro#1092). Bound the cache with a second-chance FIFO over flushed lazy wrappers, which are pure caches over KV bytes: evicted wrappers are re-created from the KV store on the next read, so eviction only costs a re-decode. KV fallback lookups now run regardless of load_state so evicted entries stay reachable in AllLoaded mode as well. Refs #1092 --- Cargo.lock | 2 +- .../container_store/container_wrapper.rs | 35 ++++ .../src/state/container_store/inner_store.rs | 198 ++++++++++++------ crates/loro-internal/tests/mem_probe.rs | 121 +++++++++++ 4 files changed, 295 insertions(+), 61 deletions(-) create mode 100644 crates/loro-internal/tests/mem_probe.rs diff --git a/Cargo.lock b/Cargo.lock index f9f25b023..510a8f62a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1863,7 +1863,7 @@ checksum = "3f3d053a135388e6b1df14e8af1212af5064746e9b87a06a345a7a779ee9695a" [[package]] name = "loro-wasm" -version = "1.14.0" +version = "1.15.1" dependencies = [ "js-sys", "loro-common 1.13.1", diff --git a/crates/loro-internal/src/state/container_store/container_wrapper.rs b/crates/loro-internal/src/state/container_store/container_wrapper.rs index ba2a3a507..df35676ff 100644 --- a/crates/loro-internal/src/state/container_store/container_wrapper.rs +++ b/crates/loro-internal/src/state/container_store/container_wrapper.rs @@ -24,6 +24,13 @@ pub(crate) struct ContainerWrapper { parent: Option, data: ContainerData, flushed: bool, + /// Whether this wrapper is currently enqueued in the store's bounded + /// decoded-value cache (`InnerStore::value_cache_queue`). + in_value_cache_queue: bool, + /// Second-chance bit for the value-cache eviction clock: set on every read + /// that uses the cached value, cleared when the eviction pass considers + /// this wrapper. + value_cache_referenced: bool, } #[derive(Debug)] @@ -148,6 +155,8 @@ impl ContainerWrapper { kind: idx.get_type(), data: ContainerData::State(state), flushed: false, + in_value_cache_queue: false, + value_cache_referenced: false, } } @@ -544,6 +553,8 @@ impl ContainerWrapper { bytes_offset_for_state: None, })), flushed: true, + in_value_cache_queue: false, + value_cache_referenced: false, }) } @@ -577,6 +588,30 @@ impl ContainerWrapper { } } + /// Whether this wrapper holds an evictable cached decoded value: a flushed + /// lazy wrapper is a pure cache over the KV bytes, so dropping it loses + /// nothing (a later read re-creates it from the KV store). + pub(super) fn is_evictable_cached_value(&self) -> bool { + self.flushed && matches!(&self.data, ContainerData::Lazy(lazy) if lazy.value.is_some()) + } + + pub(super) fn is_in_value_cache_queue(&self) -> bool { + self.in_value_cache_queue + } + + pub(super) fn set_in_value_cache_queue(&mut self, value: bool) { + self.in_value_cache_queue = value; + } + + pub(super) fn mark_value_cache_referenced(&mut self) { + self.value_cache_referenced = true; + } + + /// Returns and clears the second-chance bit. + pub(super) fn take_value_cache_referenced(&mut self) -> bool { + std::mem::take(&mut self.value_cache_referenced) + } + #[cfg(test)] pub(super) fn has_cached_value_for_test(&self) -> bool { self.has_cached_value() diff --git a/crates/loro-internal/src/state/container_store/inner_store.rs b/crates/loro-internal/src/state/container_store/inner_store.rs index 5d1d71e14..afb4d3793 100644 --- a/crates/loro-internal/src/state/container_store/inner_store.rs +++ b/crates/loro-internal/src/state/container_store/inner_store.rs @@ -8,13 +8,32 @@ use crate::{ }; use bytes::Bytes; use loro_common::{ContainerID, LoroResult, LoroValue}; +use std::collections::VecDeque; use super::ContainerWrapper; +/// Upper bound on how many lazy containers keep their decoded value cached in +/// memory for the read paths (`map_get`, `list_get`, text reads, ...). +/// +/// Reading a container through a handler decodes its value into the wrapper +/// once so repeated reads of the same container are cheap. Without a bound, a +/// container-by-container walk over an imported document pins O(containers +/// ever read) of memory until the doc is dropped — about 4 KB per container, +/// which traps wasm32 at the 4 GiB limit around one million containers +/// (loro-dev/loro#1092). Evicted wrappers are re-created from the KV store on +/// the next read, so eviction only costs a re-decode. The queue is a +/// second-chance FIFO: a container whose cached value is hit again survives +/// one eviction pass, which keeps ancestors hot during deep walks. +const MAX_CACHED_CONTAINER_VALUES: usize = 2048; + /// The invariants about this struct: /// /// - `kv` is either the same or older than `store`. -/// - if `load_state` is `AllLoaded`, then `store` contains all the entries from `kv` +/// - `store` is a cache over `kv`: every container in `kv` is either present in +/// `store` or was evicted from it by the bounded decoded-value cache +/// ([`MAX_CACHED_CONTAINER_VALUES`]). Evicted entries are re-created from +/// `kv` on the next access, so lookups must always fall back to `kv` on a +/// `store` miss, regardless of `load_state`. /// /// Invariants: it should be agnostic to the users of this struct whether a container is stored in `kv` or `store` pub(crate) struct InnerStore { @@ -23,6 +42,11 @@ pub(crate) struct InnerStore { kv: KvWrapper, load_state: LoadState, config: Configure, + /// FIFO (with a second-chance bit on each wrapper) of containers whose + /// decoded value is currently cached in `store`. Bounded by + /// [`MAX_CACHED_CONTAINER_VALUES`]. Entries may be stale (the wrapper was + /// dropped or materialized into a full state); eviction skips those. + value_cache_queue: VecDeque, } #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -94,14 +118,11 @@ impl InnerStore { if self.get_entry_mut(idx).is_none() { let id = self.arena.get_container_id(idx).unwrap(); let key = id.to_bytes(); - let container = if self.load_state != LoadState::AllLoaded { - self.kv - .get(&key) - .map(ContainerWrapper::new_from_bytes) - .unwrap_or_else(f) - } else { - f() - }; + let container = self + .kv + .get(&key) + .map(ContainerWrapper::new_from_bytes) + .unwrap_or_else(f); Self::insert_entry(&mut self.store, idx, container); } @@ -117,14 +138,12 @@ impl InnerStore { return; } - if self.load_state != LoadState::AllLoaded { - let id = self.arena.get_container_id(idx).unwrap(); - let key = id.to_bytes(); - if let Some(v) = self.kv.get(&key) { - let c = ContainerWrapper::new_from_bytes(v); - Self::insert_entry(&mut self.store, idx, c); - return; - } + let id = self.arena.get_container_id(idx).unwrap(); + let key = id.to_bytes(); + if let Some(v) = self.kv.get(&key) { + let c = ContainerWrapper::new_from_bytes(v); + Self::insert_entry(&mut self.store, idx, c); + return; } let c = f(); @@ -132,7 +151,7 @@ impl InnerStore { } pub(crate) fn get_mut(&mut self, idx: ContainerIdx) -> Option<&mut ContainerWrapper> { - if self.get_entry_mut(idx).is_none() && self.load_state != LoadState::AllLoaded { + if self.get_entry_mut(idx).is_none() { let id = self.arena.get_container_id(idx).unwrap(); let key = id.to_bytes(); if let Some(v) = self.kv.get(&key) { @@ -150,25 +169,84 @@ impl InnerStore { f: impl FnOnce(&mut ContainerWrapper) -> R, ) -> Option { if let Some(entry) = self.get_entry_mut(idx) { - return Some(f(entry)); + let ans = f(entry); + self.track_value_cache(idx); + return Some(ans); } - if self.load_state != LoadState::AllLoaded { - let id = self.arena.get_container_id(idx).unwrap(); - let key = id.to_bytes(); - if let Some(v) = self.kv.get(&key) { - let mut container = ContainerWrapper::new_from_bytes(v); - let ans = f(&mut container); - if container.has_cached_value() { - Self::insert_entry(&mut self.store, idx, container); - } - return Some(ans); + let id = self.arena.get_container_id(idx).unwrap(); + let key = id.to_bytes(); + if let Some(v) = self.kv.get(&key) { + let mut container = ContainerWrapper::new_from_bytes(v); + let ans = f(&mut container); + if container.has_cached_value() { + Self::insert_entry(&mut self.store, idx, container); + self.track_value_cache(idx); } + return Some(ans); } None } + /// Track a wrapper whose read may have populated the decoded-value cache, + /// and evict the oldest cached values once the cache exceeds + /// [`MAX_CACHED_CONTAINER_VALUES`]. + /// + /// Eviction drops the whole wrapper from `store`; the wrapper is a pure + /// cache over the KV bytes (`ContainerWrapper::is_evictable_cached_value` + /// requires a flushed lazy wrapper), so the next read re-creates it from + /// the KV store. + fn track_value_cache(&mut self, idx: ContainerIdx) { + let Some(entry) = self.get_entry_mut(idx) else { + return; + }; + if !entry.is_evictable_cached_value() { + return; + } + if !entry.is_in_value_cache_queue() { + entry.set_in_value_cache_queue(true); + self.value_cache_queue.push_back(idx); + } + self.get_entry_mut(idx) + .unwrap() + .mark_value_cache_referenced(); + + while self.value_cache_queue.len() > MAX_CACHED_CONTAINER_VALUES { + let Some(victim) = self.value_cache_queue.pop_front() else { + break; + }; + enum Action { + Skip, + Requeue, + Evict, + } + let action = match self.get_entry_mut(victim) { + Some(entry) if entry.is_evictable_cached_value() => { + if entry.take_value_cache_referenced() { + Action::Requeue + } else { + entry.set_in_value_cache_queue(false); + Action::Evict + } + } + // Stale entry: the wrapper was dropped or materialized into a + // full state since it was enqueued. + _ => Action::Skip, + }; + match action { + Action::Skip => {} + Action::Requeue => self.value_cache_queue.push_back(victim), + Action::Evict => { + let slot = Self::slot(victim); + if let Some(entry) = self.store.get_mut(slot) { + *entry = None; + } + } + } + } + } + /// Read a container without retaining a wrapper loaded only for this read. pub(crate) fn try_with_container_for_ephemeral_read( &mut self, @@ -179,13 +257,11 @@ impl InnerStore { return Ok(Some(f(entry))); } - if self.load_state != LoadState::AllLoaded { - let id = self.arena.get_container_id(idx).unwrap(); - let key = id.to_bytes(); - if let Some(v) = self.kv.get(&key) { - let mut container = ContainerWrapper::try_new_from_bytes(v)?; - return Ok(Some(f(&mut container))); - } + let id = self.arena.get_container_id(idx).unwrap(); + let key = id.to_bytes(); + if let Some(v) = self.kv.get(&key) { + let mut container = ContainerWrapper::try_new_from_bytes(v)?; + return Ok(Some(f(&mut container))); } Ok(None) @@ -205,13 +281,11 @@ impl InnerStore { return entry.try_get_value_ephemeral(idx, ctx).map(Some); } - if self.load_state != LoadState::AllLoaded { - let id = self.arena.get_container_id(idx).unwrap(); - let key = id.to_bytes(); - if let Some(bytes) = self.kv.get(&key) { - let mut container = ContainerWrapper::try_new_from_bytes(bytes)?; - return container.try_get_value(idx, ctx).map(Some); - } + let id = self.arena.get_container_id(idx).unwrap(); + let key = id.to_bytes(); + if let Some(bytes) = self.kv.get(&key) { + let mut container = ContainerWrapper::try_new_from_bytes(bytes)?; + return container.try_get_value(idx, ctx).map(Some); } Ok(None) @@ -234,17 +308,15 @@ impl InnerStore { return Ok(Some((parent, value))); } - if self.load_state != LoadState::AllLoaded { - let id = self.arena.get_container_id(idx).unwrap(); - let key = id.to_bytes(); - if let Some(value) = self.kv.get(&key) { - let mut container = ContainerWrapper::try_new_from_bytes(value)?; - let parent = container.parent().cloned(); - // This wrapper is already temporary, so decoding into it does not retain state in - // the document and avoids constructing another temporary wrapper internally. - let value = container.try_get_value(idx, ctx)?; - return Ok(Some((parent, value))); - } + let id = self.arena.get_container_id(idx).unwrap(); + let key = id.to_bytes(); + if let Some(value) = self.kv.get(&key) { + let mut container = ContainerWrapper::try_new_from_bytes(value)?; + let parent = container.parent().cloned(); + // This wrapper is already temporary, so decoding into it does not retain state in + // the document and avoids constructing another temporary wrapper internally. + let value = container.try_get_value(idx, ctx)?; + return Ok(Some((parent, value))); } Ok(None) @@ -274,12 +346,8 @@ impl InnerStore { } } - if self.load_state != LoadState::AllLoaded { - let key = id.to_bytes(); - return self.kv.contains_key(&key); - } - - false + let key = id.to_bytes(); + self.kv.contains_key(&key) } pub(crate) fn iter_all_containers_mut( @@ -376,6 +444,7 @@ impl InnerStore { })); self.store.clear(); + self.value_cache_queue.clear(); self.load_state = LoadState::Lazy; Ok(fr) } @@ -409,6 +478,14 @@ impl InnerStore { } }); + // Entries were rebuilt from the KV store; drop any queue entries that + // refer to replaced wrappers so the queue only names wrappers that are + // actually enqueued. + self.value_cache_queue.clear(); + for entry in self.store.iter_mut().flatten() { + entry.set_in_value_cache_queue(false); + } + self.load_state = LoadState::AllLoaded; Ok(()) } @@ -489,6 +566,7 @@ impl InnerStore { kv: KvWrapper::new_mem(), load_state: LoadState::AllLoaded, config, + value_cache_queue: VecDeque::new(), } } diff --git a/crates/loro-internal/tests/mem_probe.rs b/crates/loro-internal/tests/mem_probe.rs new file mode 100644 index 000000000..d748e1a1e --- /dev/null +++ b/crates/loro-internal/tests/mem_probe.rs @@ -0,0 +1,121 @@ +//! Temporary memory probe for loro#1092: measures live Rust heap bytes retained +//! by the container-by-container handle walk vs the bulk deep-value path. +//! Run: cargo test -p loro-internal --test mem_probe -- --nocapture + +use dev_utils::get_mem_usage; +use loro_internal::handler::{Handler, ValueOrHandler}; +use loro_internal::cursor::PosType; +use loro_internal::LoroDoc; + +fn live_mib() -> f64 { + get_mem_usage().0 as f64 / 2f64.powi(20) +} + +const TEXT: &str = "The quick brown fox jumps over the lazy dog. \ +The quick brown fox jumps over the lazy dog. The quick brown fox jumps over the lazy dog."; + +fn build_snapshot(turns: usize, items_per_turn: usize) -> Vec { + let doc = LoroDoc::new_auto_commit(); + let history = doc.get_list("history"); + for t in 0..turns { + let turn = history + .push_container(loro_internal::handler::MapHandler::new_detached()) + .unwrap(); + turn.insert("id", format!("turn-{t}").as_str()).unwrap(); + turn.insert("role", if t % 2 == 0 { "user" } else { "assistant" }) + .unwrap(); + let items = turn + .insert_container("items", loro_internal::handler::ListHandler::new_detached()) + .unwrap(); + let count = if t % 2 == 0 { 1 } else { items_per_turn }; + for i in 0..count { + let item = items + .push_container(loro_internal::handler::MapHandler::new_detached()) + .unwrap(); + if i % 4 == 3 { + item.insert("type", "text").unwrap(); + let text = item + .insert_container("text", loro_internal::handler::TextHandler::new_detached()) + .unwrap(); + text.insert(0, &format!("{TEXT} turn {t} item {i}"), PosType::Unicode) + .unwrap(); + } else { + item.insert("type", "tool_call").unwrap(); + item.insert("toolCallId", format!("tc-{t}-{i}").as_str()).unwrap(); + let raw = item + .insert_container("rawInput", loro_internal::handler::MapHandler::new_detached()) + .unwrap(); + raw.insert("command", "pnpm test").unwrap(); + raw.insert("cwd", "/repo").unwrap(); + let text = item + .insert_container("title", loro_internal::handler::TextHandler::new_detached()) + .unwrap(); + text.insert(0, TEXT, PosType::Unicode).unwrap(); + } + } + } + doc.export(loro_internal::encoding::ExportMode::snapshot()).unwrap() +} + +fn walk(h: &Handler, handles: &mut usize) { + *handles += 1; + match h { + Handler::Map(m) => { + let keys: Vec<_> = m.keys().collect(); + for k in keys { + if let Some(v) = m.get_(&k) { + if let ValueOrHandler::Handler(child) = v { + walk(&child, handles); + } + } + } + } + Handler::List(l) => { + for i in 0..l.len() { + if let Some(v) = l.get_(i) { + if let ValueOrHandler::Handler(child) = v { + walk(&child, handles); + } + } + } + } + Handler::Text(t) => { + let _ = t.to_string(); + } + _ => {} + } +} + +#[test] +fn probe_handle_walk_memory() { + let snapshot = build_snapshot(95, 100); // ≈ 63k containers like the 190-turn repro + println!("snapshot: {:.1} MiB", snapshot.len() as f64 / 2f64.powi(20)); + + // Bulk path baseline + let doc = LoroDoc::new_auto_commit(); + let base = live_mib(); + doc.import(&snapshot).unwrap(); + let after_import = live_mib(); + let _ = doc.get_deep_value(); + let after_deep = live_mib(); + println!("bulk: import={after_import:.1} deep_value={after_deep:.1} MiB (base {base:.1})"); + drop(doc); + let after_drop = live_mib(); + println!("bulk: after doc drop: {after_drop:.1} MiB"); + + // Handle path + let doc = LoroDoc::new_auto_commit(); + doc.import(&snapshot).unwrap(); + let after_import2 = live_mib(); + let mut handles = 0; + let history = Handler::List(doc.get_list("history")); + walk(&history, &mut handles); + let after_walk1 = live_mib(); + println!("handle: import={after_import2:.1} walk#1={after_walk1:.1} MiB handles={handles}"); + walk(&history, &mut handles); + let after_walk2 = live_mib(); + println!("handle: walk#2={after_walk2:.1} MiB handles={handles}"); + drop(history); + drop(doc); + println!("handle: after doc drop: {:.1} MiB", live_mib()); +} From fc4b5d4385d8b836213d28f199ad74e3fb3f0af7 Mon Sep 17 00:00:00 2001 From: Zixuan Chen Date: Fri, 4 Sep 2026 10:51:03 +0800 Subject: [PATCH 2/4] test: regression coverage for the bounded container value cache - loro-internal: bounded-cache unit tests (bounded growth during a container walk, re-reads and edits after eviction, eviction under AllLoaded) with a small cache bound under cfg(test) - loro-wasm: walk_mem.test.ts asserts a ~100k-container handle walk keeps external memory within a small multiple of the toJSON() delta (fails at ~286 MB on the pre-fix build, ~13 MB after) - changeset, context/container-value-cache.md, AGENTS.md links - drop the temporary mem_probe diagnostic --- .changeset/bounded-container-value-cache.md | 24 ++++ AGENTS.md | 3 + context/container-value-cache.md | 64 +++++++++ crates/loro-internal/src/state/AGENTS.md | 5 +- .../src/state/container_store.rs | 121 ++++++++++++++++++ .../src/state/container_store/inner_store.rs | 17 +++ crates/loro-internal/tests/mem_probe.rs | 121 ------------------ crates/loro-wasm/tests/walk_mem.test.ts | 114 +++++++++++++++++ 8 files changed, 347 insertions(+), 122 deletions(-) create mode 100644 .changeset/bounded-container-value-cache.md create mode 100644 context/container-value-cache.md delete mode 100644 crates/loro-internal/tests/mem_probe.rs create mode 100644 crates/loro-wasm/tests/walk_mem.test.ts diff --git a/.changeset/bounded-container-value-cache.md b/.changeset/bounded-container-value-cache.md new file mode 100644 index 000000000..86c96b22b --- /dev/null +++ b/.changeset/bounded-container-value-cache.md @@ -0,0 +1,24 @@ +--- +"loro-crdt": patch +--- + +Fix wasm memory retention when reading a document container by container +(loro-dev/loro#1092). + +Every read through a container handle (`LoroMap.keys()`/`get()`, +`LoroList.get()`, `LoroText.toJSON()`, ...) decoded the container's value into +an in-memory cache that was pinned for the lifetime of the document — about +4 KB per container, released only by `doc.free()`. Walking a large document +this way (the pattern loro-mirror's initial state build uses) retained ~4 KB × +containers-ever-read and trapped wasm32 at the 4 GiB limit around one million +containers. + +The decoded-value cache is now bounded (2048 entries, second-chance FIFO). +Evicted entries are pure caches over the KV store and are re-decoded on the +next read, so this only changes memory behavior, not API semantics. +`doc.free()` semantics are unchanged. + +Measured on the issue's repro (570-turn document, 188k container handles, +release build): the handle walk retains no per-container memory (external +memory flat at ~84 MiB vs +641 MiB before) and runs ~10x faster +(1.28 s vs 12.5 s) thanks to the smaller working set. diff --git a/AGENTS.md b/AGENTS.md index 41c73e54c..58f657493 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -36,6 +36,9 @@ Loro is a Rust CRDT workspace with JS/WASM packaging and a MoonBit codec. [context/wasm-error-reporting.md](context/wasm-error-reporting.md). - WASM container id wrapper identity, lazy caching, and benchmark: [context/wasm-container-id-cache.md](context/wasm-container-id-cache.md). +- Bounded decoded-value cache in `InnerStore` (second-chance FIFO, eviction + safety contract, loro-dev/loro#1092): + [context/container-value-cache.md](context/container-value-cache.md). - Context backlog: [context/CONTEXT-GAPS.md](context/CONTEXT-GAPS.md). ## Commands diff --git a/context/container-value-cache.md b/context/container-value-cache.md new file mode 100644 index 000000000..7ef4af46d --- /dev/null +++ b/context/container-value-cache.md @@ -0,0 +1,64 @@ +# Container Value Cache (InnerStore) + +Verified against code 2026-09-04. + +## The funnel + +Every per-container read — handler reads like `map_get`, `list_get`, +`map_len`, text reads, `contains_id` — goes through +`InnerStore::with_container_for_read` +(`crates/loro-internal/src/state/container_store/inner_store.rs`). On a KV +miss the wrapper is created from the KV bytes and, if the read decoded the +value (`ContainerWrapper::has_cached_value`), inserted into +`InnerStore.store`. Bulk paths (`toJSON`, deep values, snapshot export) use +`try_get_value_ephemeral` / `try_with_container_for_ephemeral_read`, which +deliberately leave no residue. + +## The bounded cache (loro-dev/loro#1092) + +Before 2026-09, a decoded value stayed pinned in `store` until the doc was +dropped: ~4 KB per container ever read, superlinear growth on +container-by-container walks, wasm32 trap at 4 GiB near one million +containers. + +`track_value_cache` now bounds the number of cached decoded values to +`MAX_CACHED_CONTAINER_VALUES` (2048; 16 under `cfg(test)`) with a +second-chance FIFO (`value_cache_queue` + per-wrapper +`in_value_cache_queue` / `value_cache_referenced` bits on +`ContainerWrapper`). The second-chance bit exists so a hot ancestor (e.g. the +root list re-read once per item during a deep walk) survives eviction passes. + +## Eviction safety contract + +Only `flushed && Lazy && value.is_some()` wrappers are evictable +(`ContainerWrapper::is_evictable_cached_value`): a flushed lazy wrapper is a +pure cache over the KV bytes, so dropping it loses nothing. Once a container +is mutated, `get_state_mut` converts the wrapper to `State` and clears +`flushed`; `State` wrappers are never evicted, so unflushed edits can never be +lost. Evicted entries are re-created from KV on the next read — eviction only +costs a re-decode. + +Because of eviction, `store` is strictly a cache over `kv`: all lookup paths +(`get_or_insert_with`, `ensure_container`, `get_mut`, `with_container_for_read`, +the ephemeral reads, `contains_id`) fall back to `kv` on a `store` miss +regardless of `load_state`. Do not reintroduce `load_state != AllLoaded` gates +on those fallbacks — `load_all`/`decode_twice` (GC snapshot import) put the +store in `AllLoaded` mode, and evicted entries must stay reachable there. + +## Remaining pins + +- Tree containers materialize a full `State` on first value read + (`decode_value_from_bytes` returns `decoded_state` for trees), so tree walks + still pin state; they are not covered by the bound. +- `decode_twice`/`load_all` insert one lightweight lazy shell per container + into `store` (no decoded value); those shells are small and stay. + +## Tests + +- `crates/loro-internal/src/state/container_store.rs` (`mod test`): + `handle_reads_bound_cached_container_values`, + `evicted_containers_stay_readable_and_editable`, + `evicted_entries_stay_readable_when_all_loaded`. +- `crates/loro-wasm/tests/walk_mem.test.ts`: ~100k-container handle walk keeps + `process.memoryUsage().external` within a small multiple of the `toJSON()` + delta (fails at ~286 MB on the pre-fix build, ~13 MB after). diff --git a/crates/loro-internal/src/state/AGENTS.md b/crates/loro-internal/src/state/AGENTS.md index cc1d6399b..74c8f6680 100644 --- a/crates/loro-internal/src/state/AGENTS.md +++ b/crates/loro-internal/src/state/AGENTS.md @@ -11,7 +11,10 @@ before changing mergeable child behavior. - `../state.rs`: `DocState`, checkout/path/deep-value traversal, state replay, lifecycle, and alive-container discovery. - `container_store/`: persisted KV-backed container snapshots and - `ContainerWrapper` encoding. + `ContainerWrapper` encoding. The decoded-value cache in `InnerStore` is + bounded and evicted wrappers must stay re-creatable from KV; read + [../../../../context/container-value-cache.md](../../../../context/container-value-cache.md) + before changing read/caching paths there. - `map_state.rs`, `list_state.rs`, `richtext_state.rs`, `tree_state.rs`, `movable_list_state.rs`, `counter_state.rs`: per-container state and snapshot codecs. `richtext_state.rs` also hosts `redact_dead_style_values`, used by diff --git a/crates/loro-internal/src/state/container_store.rs b/crates/loro-internal/src/state/container_store.rs index 22b4b10d6..95e1b5e3b 100644 --- a/crates/loro-internal/src/state/container_store.rs +++ b/crates/loro-internal/src/state/container_store.rs @@ -746,4 +746,125 @@ mod test { let mut state = round_tripped.app_state().lock(); assert!(state.store.contains_id(&child_id)); } + + fn doc_with_many_child_maps(n: usize) -> LoroDoc { + let doc = LoroDoc::new_auto_commit(); + let list = doc.get_list("list"); + for i in 0..n { + let map = list + .insert_container(i, MapHandler::new_detached()) + .unwrap(); + map.insert("key", i as i64).unwrap(); + } + doc + } + + /// Regression test for loro-dev/loro#1092: walking an imported document + /// container by container must not pin every decoded value in memory. + #[test] + fn handle_reads_bound_cached_container_values() { + let n = inner_store::MAX_CACHED_CONTAINER_VALUES * 4; + let source = doc_with_many_child_maps(n); + let bytes = export_container_store(&source); + let mut store = decode_container_store(bytes); + + let list_idx = store + .arena + .register_container(&ContainerID::new_root("list", ContainerType::List)); + assert_eq!(store.list_len(list_idx), n); + for i in 0..n { + let value = store.list_get(list_idx, i).unwrap(); + let child_id = value.as_container().unwrap().clone(); + let child_idx = store.arena.register_container(&child_id); + assert_eq!(store.map_get(child_idx, "key"), Some((i as i64).into())); + } + + assert!( + store.store.cached_value_count_for_test() <= inner_store::MAX_CACHED_CONTAINER_VALUES, + "cached decoded values must stay bounded, got {}", + store.store.cached_value_count_for_test() + ); + } + + /// Evicted containers are pure KV caches: re-reading and editing them after + /// eviction must return and preserve the correct values. + #[test] + fn evicted_containers_stay_readable_and_editable() { + let n = inner_store::MAX_CACHED_CONTAINER_VALUES * 4; + let source = doc_with_many_child_maps(n); + let snapshot = source + .export(crate::encoding::ExportMode::Snapshot) + .unwrap(); + let imported = LoroDoc::new_auto_commit(); + imported.import(&snapshot).unwrap(); + + let read_all = |doc: &LoroDoc| { + let list = doc.get_list("list"); + for i in 0..n { + let value = list.get(i).unwrap(); + let child_id = value.as_container().unwrap().clone(); + let child = doc.get_map(child_id); + assert_eq!(child.get("key"), Some((i as i64).into())); + } + }; + + // The first walk evicts the early children; the second walk re-reads + // them through the KV fallback. + read_all(&imported); + read_all(&imported); + + // Mutating an evicted-then-reloaded container must persist. + let first_id = imported + .get_list("list") + .get(0) + .unwrap() + .as_container() + .unwrap() + .clone(); + imported.get_map(first_id).insert("edited", true).unwrap(); + + let snapshot = imported + .export(crate::encoding::ExportMode::Snapshot) + .unwrap(); + let round_tripped = LoroDoc::new(); + round_tripped.import(&snapshot).unwrap(); + assert_eq!(round_tripped.get_deep_value(), imported.get_deep_value()); + } + + /// `load_all` (and GC snapshot imports) put the store in `AllLoaded` mode; + /// evicted entries must still be re-created from the KV store there. + #[test] + fn evicted_entries_stay_readable_when_all_loaded() { + let n = inner_store::MAX_CACHED_CONTAINER_VALUES * 4; + let source = doc_with_many_child_maps(n); + let bytes = export_container_store(&source); + let mut store = decode_container_store(bytes); + + let list_idx = store + .arena + .register_container(&ContainerID::new_root("list", ContainerType::List)); + store.store.load_all(); + + for i in 0..n { + let value = store.list_get(list_idx, i).unwrap(); + let child_id = value.as_container().unwrap().clone(); + let child_idx = store.arena.register_container(&child_id); + assert_eq!(store.map_get(child_idx, "key"), Some((i as i64).into())); + } + + assert!( + store.store.cached_value_count_for_test() <= inner_store::MAX_CACHED_CONTAINER_VALUES, + "cached decoded values must stay bounded, got {}", + store.store.cached_value_count_for_test() + ); + + // The first children were evicted during the walk; re-reading them must + // fall back to the KV store even in AllLoaded mode. + for i in 0..n { + let value = store.list_get(list_idx, i).unwrap(); + let child_id = value.as_container().unwrap().clone(); + let child_idx = store.arena.register_container(&child_id); + assert_eq!(store.map_get(child_idx, "key"), Some((i as i64).into())); + } + } } diff --git a/crates/loro-internal/src/state/container_store/inner_store.rs b/crates/loro-internal/src/state/container_store/inner_store.rs index afb4d3793..9dffa5648 100644 --- a/crates/loro-internal/src/state/container_store/inner_store.rs +++ b/crates/loro-internal/src/state/container_store/inner_store.rs @@ -24,7 +24,12 @@ use super::ContainerWrapper; /// the next read, so eviction only costs a re-decode. The queue is a /// second-chance FIFO: a container whose cached value is hit again survives /// one eviction pass, which keeps ancestors hot during deep walks. +#[cfg(not(test))] const MAX_CACHED_CONTAINER_VALUES: usize = 2048; +/// Tests use a small bound so exercising eviction does not require building +/// thousands of containers. +#[cfg(test)] +pub(super) const MAX_CACHED_CONTAINER_VALUES: usize = 16; /// The invariants about this struct: /// @@ -551,6 +556,18 @@ impl InnerStore { .is_some_and(|entry| entry.has_cached_value_for_test()) } + /// Number of wrappers currently holding an evictable cached decoded value. + /// Must stay bounded by [`MAX_CACHED_CONTAINER_VALUES`] no matter how many + /// containers have been read (loro-dev/loro#1092). + #[cfg(test)] + pub(super) fn cached_value_count_for_test(&self) -> usize { + self.store + .iter() + .flatten() + .filter(|entry| entry.is_evictable_cached_value()) + .count() + } + #[cfg(test)] pub(super) fn has_materialized_map_value_for_test(&mut self, idx: ContainerIdx) -> bool { self.get_entry_mut(idx) diff --git a/crates/loro-internal/tests/mem_probe.rs b/crates/loro-internal/tests/mem_probe.rs deleted file mode 100644 index d748e1a1e..000000000 --- a/crates/loro-internal/tests/mem_probe.rs +++ /dev/null @@ -1,121 +0,0 @@ -//! Temporary memory probe for loro#1092: measures live Rust heap bytes retained -//! by the container-by-container handle walk vs the bulk deep-value path. -//! Run: cargo test -p loro-internal --test mem_probe -- --nocapture - -use dev_utils::get_mem_usage; -use loro_internal::handler::{Handler, ValueOrHandler}; -use loro_internal::cursor::PosType; -use loro_internal::LoroDoc; - -fn live_mib() -> f64 { - get_mem_usage().0 as f64 / 2f64.powi(20) -} - -const TEXT: &str = "The quick brown fox jumps over the lazy dog. \ -The quick brown fox jumps over the lazy dog. The quick brown fox jumps over the lazy dog."; - -fn build_snapshot(turns: usize, items_per_turn: usize) -> Vec { - let doc = LoroDoc::new_auto_commit(); - let history = doc.get_list("history"); - for t in 0..turns { - let turn = history - .push_container(loro_internal::handler::MapHandler::new_detached()) - .unwrap(); - turn.insert("id", format!("turn-{t}").as_str()).unwrap(); - turn.insert("role", if t % 2 == 0 { "user" } else { "assistant" }) - .unwrap(); - let items = turn - .insert_container("items", loro_internal::handler::ListHandler::new_detached()) - .unwrap(); - let count = if t % 2 == 0 { 1 } else { items_per_turn }; - for i in 0..count { - let item = items - .push_container(loro_internal::handler::MapHandler::new_detached()) - .unwrap(); - if i % 4 == 3 { - item.insert("type", "text").unwrap(); - let text = item - .insert_container("text", loro_internal::handler::TextHandler::new_detached()) - .unwrap(); - text.insert(0, &format!("{TEXT} turn {t} item {i}"), PosType::Unicode) - .unwrap(); - } else { - item.insert("type", "tool_call").unwrap(); - item.insert("toolCallId", format!("tc-{t}-{i}").as_str()).unwrap(); - let raw = item - .insert_container("rawInput", loro_internal::handler::MapHandler::new_detached()) - .unwrap(); - raw.insert("command", "pnpm test").unwrap(); - raw.insert("cwd", "/repo").unwrap(); - let text = item - .insert_container("title", loro_internal::handler::TextHandler::new_detached()) - .unwrap(); - text.insert(0, TEXT, PosType::Unicode).unwrap(); - } - } - } - doc.export(loro_internal::encoding::ExportMode::snapshot()).unwrap() -} - -fn walk(h: &Handler, handles: &mut usize) { - *handles += 1; - match h { - Handler::Map(m) => { - let keys: Vec<_> = m.keys().collect(); - for k in keys { - if let Some(v) = m.get_(&k) { - if let ValueOrHandler::Handler(child) = v { - walk(&child, handles); - } - } - } - } - Handler::List(l) => { - for i in 0..l.len() { - if let Some(v) = l.get_(i) { - if let ValueOrHandler::Handler(child) = v { - walk(&child, handles); - } - } - } - } - Handler::Text(t) => { - let _ = t.to_string(); - } - _ => {} - } -} - -#[test] -fn probe_handle_walk_memory() { - let snapshot = build_snapshot(95, 100); // ≈ 63k containers like the 190-turn repro - println!("snapshot: {:.1} MiB", snapshot.len() as f64 / 2f64.powi(20)); - - // Bulk path baseline - let doc = LoroDoc::new_auto_commit(); - let base = live_mib(); - doc.import(&snapshot).unwrap(); - let after_import = live_mib(); - let _ = doc.get_deep_value(); - let after_deep = live_mib(); - println!("bulk: import={after_import:.1} deep_value={after_deep:.1} MiB (base {base:.1})"); - drop(doc); - let after_drop = live_mib(); - println!("bulk: after doc drop: {after_drop:.1} MiB"); - - // Handle path - let doc = LoroDoc::new_auto_commit(); - doc.import(&snapshot).unwrap(); - let after_import2 = live_mib(); - let mut handles = 0; - let history = Handler::List(doc.get_list("history")); - walk(&history, &mut handles); - let after_walk1 = live_mib(); - println!("handle: import={after_import2:.1} walk#1={after_walk1:.1} MiB handles={handles}"); - walk(&history, &mut handles); - let after_walk2 = live_mib(); - println!("handle: walk#2={after_walk2:.1} MiB handles={handles}"); - drop(history); - drop(doc); - println!("handle: after doc drop: {:.1} MiB", live_mib()); -} diff --git a/crates/loro-wasm/tests/walk_mem.test.ts b/crates/loro-wasm/tests/walk_mem.test.ts new file mode 100644 index 000000000..37703f7fb --- /dev/null +++ b/crates/loro-wasm/tests/walk_mem.test.ts @@ -0,0 +1,114 @@ +import { describe, expect, it } from "vitest"; +import { LoroDoc, LoroList, LoroMap, LoroText } from "../bundler"; + +const MB = 1024 * 1024; + +/** + * Builds a document shaped like a chat history: a root list of turn maps, + * each with an items list of maps carrying nested maps/texts. With + * turns=1300 and itemsPerTurn=60 this yields ~100k containers. + */ +function buildSnapshot(turns: number, itemsPerTurn: number): Uint8Array { + const TEXT = "The quick brown fox jumps over the lazy dog. ".repeat(4); + const doc = new LoroDoc(); + const history = doc.getList("history"); + for (let t = 0; t < turns; t++) { + const turn = history.pushContainer(new LoroMap()); + turn.set("id", `turn-${t}`); + turn.set("role", t % 2 === 0 ? "user" : "assistant"); + const items = turn.setContainer("items", new LoroList()); + const count = t % 2 === 0 ? 1 : itemsPerTurn; + for (let i = 0; i < count; i++) { + const item = items.pushContainer(new LoroMap()); + if (i % 4 === 3) { + item.set("type", "text"); + item.setContainer("text", new LoroText()).insert(0, `${TEXT} ${t}/${i}`); + } else { + item.set("type", "tool_call"); + item.set("toolCallId", `tc-${t}-${i}`); + const raw = item.setContainer("rawInput", new LoroMap()); + raw.set("command", "pnpm test"); + raw.set("cwd", "/repo"); + item.setContainer("title", new LoroText()).insert(0, TEXT); + } + } + } + doc.commit(); + const snapshot = doc.export({ mode: "snapshot" }); + doc.free(); + return snapshot; +} + +type ContainerHandle = { kind(): string }; + +const isContainer = (v: unknown): v is ContainerHandle => + !!v && typeof v === "object" && typeof (v as ContainerHandle).kind === "function"; + +/** loro-mirror-style traversal: keys()/get() per map, get(i) per list, toJSON() per text. */ +function walk(container: ContainerHandle): unknown { + const kind = container.kind(); + if (kind === "Map") { + const map = container as unknown as { keys(): string[]; get(k: string): unknown }; + const out: Record = {}; + for (const key of map.keys()) { + const value = map.get(key); + out[key] = isContainer(value) ? walk(value) : value; + } + return out; + } + if (kind === "List" || kind === "MovableList") { + const list = container as unknown as { length: number; get(i: number): unknown }; + const out: unknown[] = []; + for (let i = 0; i < list.length; i++) { + const value = list.get(i); + out.push(isContainer(value) ? walk(value) : value); + } + return out; + } + return (container as unknown as { toJSON(): unknown }).toJSON(); +} + +describe("container handle walk memory", () => { + // Regression test for https://github.com/loro-dev/loro/issues/1092: + // reading an imported document container by container through JS handles + // used to pin ~4 KB of wasm linear memory per container until doc.free(). + it("keeps wasm memory bounded while walking every container", () => { + if (typeof global.gc !== "function") { + console.warn("Skipping memory test because --expose-gc was not provided."); + return; + } + + const snapshot = buildSnapshot(1300, 60); + + // Bulk-read baseline: toJSON() reads the same state without retaining + // per-container decoded values. + const docA = new LoroDoc(); + docA.import(snapshot); + global.gc(); + const beforeToJson = process.memoryUsage().external; + docA.toJSON(); + global.gc(); + const toJsonDelta = process.memoryUsage().external - beforeToJson; + docA.free(); + + // Handle walk over every container of an identical document. + const docB = new LoroDoc(); + docB.import(snapshot); + global.gc(); + const beforeWalk = process.memoryUsage().external; + walk(docB.getList("history") as unknown as ContainerHandle); + global.gc(); + const walkDelta = process.memoryUsage().external - beforeWalk; + docB.free(); + + console.log( + `external deltas: toJSON=${(toJsonDelta / MB).toFixed(1)}MB walk=${(walkDelta / MB).toFixed(1)}MB`, + ); + + // With the leak, the walk retained ~4 KB per container (~400 MB at + // this size). It must now stay within a small multiple of the bulk + // read, modulo allocator noise. + const allowance = Math.max(3 * Math.max(toJsonDelta, 0), 48 * MB); + expect(walkDelta).toBeLessThan(allowance); + }, 120_000); +}); From ddf7c6a3dfa39e704a829ffc11409c5248097065 Mon Sep 17 00:00:00 2001 From: Zixuan Chen Date: Fri, 4 Sep 2026 11:12:43 +0800 Subject: [PATCH 3/4] test(loro-internal): memory-based regression test for the bounded value cache - tests/handle_walk_memory.rs uses the dev_utils counting allocator to assert the loro-mirror-style handle walk retains bounded memory (walk number 2 adds ~0; total far below the unbounded ~1 KB/container) - extend first_lazy_read_caches_value with the new bounded-cache contract: values are still cached on first read, but reads beyond the bound evict, and evicted containers re-decode from KV transparently --- .../src/state/container_store.rs | 27 ++++ .../loro-internal/tests/handle_walk_memory.rs | 135 ++++++++++++++++++ 2 files changed, 162 insertions(+) create mode 100644 crates/loro-internal/tests/handle_walk_memory.rs diff --git a/crates/loro-internal/src/state/container_store.rs b/crates/loro-internal/src/state/container_store.rs index 95e1b5e3b..25cf48e3b 100644 --- a/crates/loro-internal/src/state/container_store.rs +++ b/crates/loro-internal/src/state/container_store.rs @@ -466,6 +466,33 @@ mod test { assert!(!store.store.has_cached_value_for_test(map_idx)); assert_eq!(store.map_len(map_idx), 2); assert!(store.store.has_cached_value_for_test(map_idx)); + + // The cache is bounded (loro-dev/loro#1092): reading more containers + // than the bound evicts older cached values, and re-reading an evicted + // container transparently re-decodes it from the KV store. + let n = inner_store::MAX_CACHED_CONTAINER_VALUES * 2; + let source = doc_with_many_child_maps(n); + let bytes = export_container_store(&source); + let mut store = decode_container_store(bytes); + let list_idx = store + .arena + .register_container(&ContainerID::new_root("list", ContainerType::List)); + let mut child_ids = Vec::new(); + for i in 0..n { + let value = store.list_get(list_idx, i).unwrap(); + child_ids.push(value.as_container().unwrap().clone()); + let child_idx = store.arena.register_container(child_ids.last().unwrap()); + assert_eq!(store.map_get(child_idx, "key"), Some((i as i64).into())); + } + + assert!( + store.store.cached_value_count_for_test() <= inner_store::MAX_CACHED_CONTAINER_VALUES, + "cached decoded values must stay bounded, got {}", + store.store.cached_value_count_for_test() + ); + // The first children were evicted; re-reading must fall back to KV. + let first_idx = store.arena.register_container(&child_ids[0]); + assert_eq!(store.map_get(first_idx, "key"), Some(0i64.into())); } #[test] diff --git a/crates/loro-internal/tests/handle_walk_memory.rs b/crates/loro-internal/tests/handle_walk_memory.rs new file mode 100644 index 000000000..74582a501 --- /dev/null +++ b/crates/loro-internal/tests/handle_walk_memory.rs @@ -0,0 +1,135 @@ +//! Regression test for loro#1092: a container-by-container handle walk must +//! retain bounded memory. The decoded-value cache in `InnerStore` is capped +//! (`MAX_CACHED_CONTAINER_VALUES`), so retained memory after a walk is +//! O(cache bound), not O(containers ever read). +//! +//! Run: cargo test -p loro-internal --test handle_walk_memory -- --nocapture + +use dev_utils::get_mem_usage; +use loro_internal::cursor::PosType; +use loro_internal::handler::{Handler, ValueOrHandler}; +use loro_internal::LoroDoc; + +fn live_bytes() -> usize { + get_mem_usage().0 +} + +const MIB: usize = 1 << 20; + +const TEXT: &str = "The quick brown fox jumps over the lazy dog. \ +The quick brown fox jumps over the lazy dog. The quick brown fox jumps over the lazy dog."; + +fn build_snapshot(turns: usize, items_per_turn: usize) -> Vec { + let doc = LoroDoc::new_auto_commit(); + let history = doc.get_list("history"); + for t in 0..turns { + let turn = history + .push_container(loro_internal::handler::MapHandler::new_detached()) + .unwrap(); + turn.insert("id", format!("turn-{t}").as_str()).unwrap(); + turn.insert("role", if t % 2 == 0 { "user" } else { "assistant" }) + .unwrap(); + let items = turn + .insert_container("items", loro_internal::handler::ListHandler::new_detached()) + .unwrap(); + let count = if t % 2 == 0 { 1 } else { items_per_turn }; + for i in 0..count { + let item = items + .push_container(loro_internal::handler::MapHandler::new_detached()) + .unwrap(); + if i % 4 == 3 { + item.insert("type", "text").unwrap(); + let text = item + .insert_container("text", loro_internal::handler::TextHandler::new_detached()) + .unwrap(); + text.insert(0, &format!("{TEXT} turn {t} item {i}"), PosType::Unicode) + .unwrap(); + } else { + item.insert("type", "tool_call").unwrap(); + item.insert("toolCallId", format!("tc-{t}-{i}").as_str()) + .unwrap(); + let raw = item + .insert_container( + "rawInput", + loro_internal::handler::MapHandler::new_detached(), + ) + .unwrap(); + raw.insert("command", "pnpm test").unwrap(); + raw.insert("cwd", "/repo").unwrap(); + let text = item + .insert_container( + "title", + loro_internal::handler::TextHandler::new_detached(), + ) + .unwrap(); + text.insert(0, TEXT, PosType::Unicode).unwrap(); + } + } + } + doc.export(loro_internal::encoding::ExportMode::snapshot()) + .unwrap() +} + +/// The loro-mirror-style traversal: keys()/get() per map, get(i) per list, +/// to_string() per text. +fn walk(h: &Handler, handles: &mut usize) { + *handles += 1; + match h { + Handler::Map(m) => { + let keys: Vec<_> = m.keys().collect(); + for k in keys { + if let Some(ValueOrHandler::Handler(child)) = m.get_(&k) { + walk(&child, handles); + } + } + } + Handler::List(l) => { + for i in 0..l.len() { + if let Some(ValueOrHandler::Handler(child)) = l.get_(i) { + walk(&child, handles); + } + } + } + Handler::Text(t) => { + let _ = t.to_string(); + } + _ => {} + } +} + +#[test] +fn handle_walk_retains_bounded_memory() { + let snapshot = build_snapshot(95, 100); // ≈ 63k containers + let doc = LoroDoc::new_auto_commit(); + doc.import(&snapshot).unwrap(); + let history = Handler::List(doc.get_list("history")); + + let after_import = live_bytes(); + let mut handles = 0; + walk(&history, &mut handles); + let after_walk1 = live_bytes(); + walk(&history, &mut handles); + let after_walk2 = live_bytes(); + + let walk1_retained = after_walk1.saturating_sub(after_import); + let walk2_growth = after_walk2.saturating_sub(after_walk1); + println!( + "handles={handles} retained after walk#1: {:.1} MiB, walk#2 growth: {:.1} MiB", + walk1_retained as f64 / MIB as f64, + walk2_growth as f64 / MIB as f64, + ); + + // Unbounded, this walk pinned ~1 KB per container (~60 MB at this size). + // The bounded cache (2048 entries) must retain far less than that. + assert!( + walk1_retained < 16 * MIB, + "walk retained {:.1} MiB; the decoded-value cache is not bounded", + walk1_retained as f64 / MIB as f64 + ); + // Re-walking hits the same bounded cache: no further growth. + assert!( + walk2_growth < 4 * MIB, + "second walk grew memory by {:.1} MiB; cached values must be reused or re-decoded, not accumulated", + walk2_growth as f64 / MIB as f64 + ); +} From e3266933d2edd1411adb66560a7e024c3d9516ea Mon Sep 17 00:00:00 2001 From: Zixuan Chen Date: Fri, 4 Sep 2026 18:07:40 +0800 Subject: [PATCH 4/4] fix(loro-internal): force kv re-scan in load_all after value-cache eviction MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review of the bounded value cache found that the AllLoaded short-circuit in load_all() became incorrect once eviction exists: AllLoaded no longer implies store contains every kv entry, so iter_all_container_ids() and iter_all_containers_mut() silently enumerated only the eviction survivors. The shallow-snapshot re-export path enumerates containers this way, which dropped evicted overlay containers from the export — silent data loss on the next import. InnerStore now tracks evicted_since_full_load (set on eviction, cleared by decode/decode_twice/a full load_all scan); load_all() re-scans kv while the flag is set instead of trusting AllLoaded. Cost: one bool check per call plus at most one re-scan after the first eviction following a full load. Regression tests: - iter_all_container_ids().count() stays complete after an evicting walk in AllLoaded mode (was 16 instead of 65 pre-fix). - shallow root + overlay-only containers -> import -> evicting walk -> re-export the same root -> re-import keeps every container (fails pre-fix). - mergeable child / tree-meta map: evict -> read via parent marker / tree path -> mutate -> snapshot round-trip. - stale queue entry after Lazy->State conversion under continued cache pressure. - both memory tests now assert walk completeness (exact handle count in handle_walk_memory.rs, deep-equal against toJSON() in walk_mem.test.ts). --- context/container-value-cache.md | 30 +++- .../src/encoding/shallow_snapshot.rs | 72 +++++++++ .../src/state/container_store.rs | 137 ++++++++++++++++++ .../src/state/container_store/inner_store.rs | 17 ++- .../loro-internal/tests/handle_walk_memory.rs | 14 +- crates/loro-wasm/tests/walk_mem.test.ts | 6 +- 6 files changed, 268 insertions(+), 8 deletions(-) diff --git a/context/container-value-cache.md b/context/container-value-cache.md index 7ef4af46d..bdca86533 100644 --- a/context/container-value-cache.md +++ b/context/container-value-cache.md @@ -45,6 +45,25 @@ regardless of `load_state`. Do not reintroduce `load_state != AllLoaded` gates on those fallbacks — `load_all`/`decode_twice` (GC snapshot import) put the store in `AllLoaded` mode, and evicted entries must stay reachable there. +## What `AllLoaded` means with eviction + +`LoadState::AllLoaded` now means "every `kv` entry was materialized into +`store` at some point", NOT "`store` is complete right now". Eviction can +remove entries afterwards, so `load_all()` tracks `evicted_since_full_load` +(set on every eviction, cleared by `decode`/`decode_twice`/a full `load_all` +scan) and re-scans `kv` when it is set instead of short-circuiting on +`AllLoaded`. Without this, `iter_all_container_ids()` / +`iter_all_containers_mut()` silently miss evicted containers; the shallow +snapshot re-export path (`encoding/shallow_snapshot.rs`) enumerates containers +that way and dropped them from the latest-state overlay — a silent data-loss +on re-import. The "content in `store` is newer than `kv`" skip inside +`load_all` stays sound: evicted entries are absent from `store`, so they are +rebuilt from `kv`. + +Cost: one bool check per `load_all()` call, plus at most one `kv` re-scan +after the first eviction following a full load (the scan's per-entry +"already in `store`" skip makes repeat scans cheap). + ## Remaining pins - Tree containers materialize a full `State` on first value read @@ -58,7 +77,14 @@ store in `AllLoaded` mode, and evicted entries must stay reachable there. - `crates/loro-internal/src/state/container_store.rs` (`mod test`): `handle_reads_bound_cached_container_values`, `evicted_containers_stay_readable_and_editable`, - `evicted_entries_stay_readable_when_all_loaded`. + `evicted_entries_stay_readable_when_all_loaded` (also asserts + `iter_all_container_ids` stays complete after evictions), + `evicted_mergeable_child_and_tree_meta_survive_round_trip`, + `stale_queue_entry_after_lazy_to_state_conversion`. +- `crates/loro-internal/src/encoding/shallow_snapshot.rs`: + `reexport_same_shallow_root_after_walk_eviction_keeps_overlay_containers` + (the silent data-loss regression). - `crates/loro-wasm/tests/walk_mem.test.ts`: ~100k-container handle walk keeps `process.memoryUsage().external` within a small multiple of the `toJSON()` - delta (fails at ~286 MB on the pre-fix build, ~13 MB after). + delta (fails at ~286 MB on the pre-fix build, ~13 MB after); also asserts + the walk result equals `toJSON()` so an incomplete walk cannot pass. diff --git a/crates/loro-internal/src/encoding/shallow_snapshot.rs b/crates/loro-internal/src/encoding/shallow_snapshot.rs index b3002eecf..f69c2dc16 100644 --- a/crates/loro-internal/src/encoding/shallow_snapshot.rs +++ b/crates/loro-internal/src/encoding/shallow_snapshot.rs @@ -549,6 +549,7 @@ mod tests { use crate::encoding::fast_snapshot::_decode_snapshot_bytes; use crate::encoding::EncodeMode; use crate::encoding::ExportMode; + use crate::handler::MapHandler; use crate::handler::TextHandler; use crate::state::{ContainerCreationContext, FastStateSnapshot, RichtextState}; use crate::HandlerTrait; @@ -765,4 +766,75 @@ mod tests { LoroValue::Null ); } + + /// Regression test for the P1 found in review of the bounded container + /// value cache (loro-dev/loro#1092): after a shallow-snapshot import the + /// store is `AllLoaded`, and a container walk evicts most decoded values. + /// Re-exporting the same shallow root must still enumerate every container + /// created after the root — otherwise the latest-state overlay silently + /// drops the evicted ones and the next import loses them. + #[test] + fn reexport_same_shallow_root_after_walk_eviction_keeps_overlay_containers() { + // Doc A: base content behind the shallow root. + let a = LoroDoc::new_auto_commit(); + a.set_peer_id(1).unwrap(); + a.get_text("text") + .insert(0, "base", PosType::Unicode) + .unwrap(); + a.commit_then_renew(); + let start = a.oplog_frontiers(); + + // Doc B: import the shallow snapshot, then create containers that live + // only in the latest-state overlay (they are not in the root KV). + let b = LoroDoc::new_auto_commit(); + b.set_peer_id(2).unwrap(); + b.import(&a.export(ExportMode::shallow_snapshot(&start)).unwrap()) + .unwrap(); + let n = 64; + let list = b.get_list("list"); + for i in 0..n { + let map = list + .insert_container(i, MapHandler::new_detached()) + .unwrap(); + map.insert("key", i as i64).unwrap(); + } + b.commit_then_renew(); + // 2 ops per container > MAX_OPS_NUM_TO_ENCODE_WITHOUT_LATEST_STATE + // (16 in test builds), so this export ships the latest-state overlay. + let blob = b.export(ExportMode::shallow_snapshot(&start)).unwrap(); + + // Doc C imports the blob (store is AllLoaded with lazy wrappers) and + // walks every overlay container, evicting most of them from the + // bounded value cache. + let c = LoroDoc::new_auto_commit(); + c.import(&blob).unwrap(); + let list = c.get_list("list"); + for i in 0..n { + let child_id = list.get(i).unwrap().as_container().unwrap().clone(); + assert_eq!(c.get_map(child_id).get("key"), Some((i as i64).into())); + } + + // Re-exporting the same shallow root reuses the stored root bytes and + // rebuilds the overlay by enumerating all containers. With the broken + // `load_all` short-circuit this dropped every evicted container. + let reexported = c.export(ExportMode::shallow_snapshot(&start)).unwrap(); + assert!( + shallow_sections(&reexported).state_bytes.is_some(), + "test setup must take the overlay export path" + ); + + let d = LoroDoc::new_auto_commit(); + d.import(&reexported).unwrap(); + let list = d.get_list("list"); + assert_eq!(list.len(), n); + for i in 0..n { + let child_id = list + .get(i) + .unwrap_or_else(|| panic!("container {i} lost after shallow re-export")) + .as_container() + .unwrap() + .clone(); + assert_eq!(d.get_map(child_id).get("key"), Some((i as i64).into())); + } + } } diff --git a/crates/loro-internal/src/state/container_store.rs b/crates/loro-internal/src/state/container_store.rs index 25cf48e3b..2d2dbc979 100644 --- a/crates/loro-internal/src/state/container_store.rs +++ b/crates/loro-internal/src/state/container_store.rs @@ -885,6 +885,11 @@ mod test { store.store.cached_value_count_for_test() ); + // `AllLoaded` must not blind full enumeration to evicted entries: once + // an eviction has happened, `load_all` re-scans the KV store instead of + // short-circuiting on `load_state`. + assert_eq!(store.store.iter_all_container_ids().count(), n + 1); + // The first children were evicted during the walk; re-reading them must // fall back to the KV store even in AllLoaded mode. for i in 0..n { @@ -894,4 +899,136 @@ mod test { assert_eq!(store.map_get(child_idx, "key"), Some((i as i64).into())); } } + + /// Mergeable children and tree meta maps are reachable through a parent + /// marker / a tree node rather than a plain value slot. Evicting their + /// cached values must not lose them: re-read through the parent/path, + /// mutate, and snapshot round-trip. + #[test] + fn evicted_mergeable_child_and_tree_meta_survive_round_trip() { + let filler = inner_store::MAX_CACHED_CONTAINER_VALUES * 2; + let source = LoroDoc::new_auto_commit(); + source + .get_map("state") + .ensure_mergeable_text("notes") + .unwrap() + .insert(0, "hello", PosType::Unicode) + .unwrap(); + let node = source.get_tree("tree").create(TreeParentId::Root).unwrap(); + source + .get_tree("tree") + .get_meta(node) + .unwrap() + .insert("k", "v") + .unwrap(); + let list = source.get_list("list"); + for i in 0..filler { + let map = list + .insert_container(i, MapHandler::new_detached()) + .unwrap(); + map.insert("key", i as i64).unwrap(); + } + let snapshot = source + .export(crate::encoding::ExportMode::Snapshot) + .unwrap(); + + let doc = LoroDoc::new_auto_commit(); + doc.import(&snapshot).unwrap(); + + // Cache the mergeable child and the tree meta values, then evict them + // by reading more containers than the cache bound. + let notes = doc.get_map("state").ensure_mergeable_text("notes").unwrap(); + assert_eq!(notes.to_string(), "hello"); + let meta = doc.get_tree("tree").get_meta(node).unwrap(); + assert_eq!(meta.get("k"), Some("v".into())); + let list = doc.get_list("list"); + for i in 0..filler { + let child_id = list.get(i).unwrap().as_container().unwrap().clone(); + assert_eq!(doc.get_map(child_id).get("key"), Some((i as i64).into())); + } + { + let mut state = doc.app_state().lock(); + assert!( + state.store.store.cached_value_count_for_test() + <= inner_store::MAX_CACHED_CONTAINER_VALUES, + "the walk must have evicted cached values, got {}", + state.store.store.cached_value_count_for_test() + ); + } + + // Re-read through the parent marker / tree path after eviction, then + // mutate both. + let notes = doc.get_map("state").ensure_mergeable_text("notes").unwrap(); + assert_eq!(notes.to_string(), "hello"); + notes.insert(5, "!", PosType::Unicode).unwrap(); + let meta = doc.get_tree("tree").get_meta(node).unwrap(); + assert_eq!(meta.get("k"), Some("v".into())); + meta.insert("k2", "v2").unwrap(); + + let snapshot = doc.export(crate::encoding::ExportMode::Snapshot).unwrap(); + let round_tripped = LoroDoc::new(); + round_tripped.import(&snapshot).unwrap(); + assert_eq!(round_tripped.get_deep_value(), doc.get_deep_value()); + assert_eq!( + round_tripped + .get_map("state") + .ensure_mergeable_text("notes") + .unwrap() + .to_string(), + "hello!" + ); + assert_eq!( + round_tripped + .get_tree("tree") + .get_meta(node) + .unwrap() + .get("k2"), + Some("v2".into()) + ); + } + + /// A wrapper that is mutated while enqueued in the value-cache queue + /// materializes into a full state, leaving a stale queue entry behind. + /// Eviction must skip the stale entry without losing the mutation or + /// wedging the queue. + #[test] + fn stale_queue_entry_after_lazy_to_state_conversion() { + let n = inner_store::MAX_CACHED_CONTAINER_VALUES * 2; + let source = doc_with_many_child_maps(n); + let snapshot = source + .export(crate::encoding::ExportMode::Snapshot) + .unwrap(); + let doc = LoroDoc::new_auto_commit(); + doc.import(&snapshot).unwrap(); + + // Read child 0 (cached + enqueued), then mutate it: the lazy wrapper + // materializes into a full state and its queue entry becomes stale. + let list = doc.get_list("list"); + let first_id = list.get(0).unwrap().as_container().unwrap().clone(); + let first_map = doc.get_map(first_id); + assert_eq!(first_map.get("key"), Some(0i64.into())); + first_map.insert("edited", true).unwrap(); + + // Keep reading past the bound so eviction pops the stale entry. + for i in 0..n { + let child_id = list.get(i).unwrap().as_container().unwrap().clone(); + assert_eq!(doc.get_map(child_id).get("key"), Some((i as i64).into())); + } + { + let mut state = doc.app_state().lock(); + assert!( + state.store.store.cached_value_count_for_test() + <= inner_store::MAX_CACHED_CONTAINER_VALUES, + "cached decoded values must stay bounded, got {}", + state.store.store.cached_value_count_for_test() + ); + } + + assert_eq!(first_map.get("edited"), Some(true.into())); + let round_tripped = LoroDoc::new(); + round_tripped + .import(&doc.export(crate::encoding::ExportMode::Snapshot).unwrap()) + .unwrap(); + assert_eq!(round_tripped.get_deep_value(), doc.get_deep_value()); + } } diff --git a/crates/loro-internal/src/state/container_store/inner_store.rs b/crates/loro-internal/src/state/container_store/inner_store.rs index 9dffa5648..e67293d67 100644 --- a/crates/loro-internal/src/state/container_store/inner_store.rs +++ b/crates/loro-internal/src/state/container_store/inner_store.rs @@ -39,6 +39,9 @@ pub(super) const MAX_CACHED_CONTAINER_VALUES: usize = 16; /// ([`MAX_CACHED_CONTAINER_VALUES`]). Evicted entries are re-created from /// `kv` on the next access, so lookups must always fall back to `kv` on a /// `store` miss, regardless of `load_state`. +/// - `load_state == AllLoaded` means every `kv` entry was materialized into +/// `store` at some point; it does NOT mean `store` is still complete. Once +/// `evicted_since_full_load` is set, `load_all()` must re-scan `kv`. /// /// Invariants: it should be agnostic to the users of this struct whether a container is stored in `kv` or `store` pub(crate) struct InnerStore { @@ -52,6 +55,11 @@ pub(crate) struct InnerStore { /// [`MAX_CACHED_CONTAINER_VALUES`]. Entries may be stale (the wrapper was /// dropped or materialized into a full state); eviction skips those. value_cache_queue: VecDeque, + /// True once the value cache has evicted at least one wrapper since the + /// last full KV scan. While this is set, `load_state == AllLoaded` no + /// longer implies `store` contains every entry in `kv`, so `load_all()` + /// must re-scan `kv` instead of short-circuiting. + evicted_since_full_load: bool, } #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -247,6 +255,9 @@ impl InnerStore { if let Some(entry) = self.store.get_mut(slot) { *entry = None; } + // `store` may no longer mirror `kv`: a full enumeration has + // to re-scan `kv` even in `AllLoaded` mode. + self.evicted_since_full_load = true; } } } @@ -450,6 +461,7 @@ impl InnerStore { self.store.clear(); self.value_cache_queue.clear(); + self.evicted_since_full_load = false; self.load_state = LoadState::Lazy; Ok(fr) } @@ -491,12 +503,13 @@ impl InnerStore { entry.set_in_value_cache_queue(false); } + self.evicted_since_full_load = false; self.load_state = LoadState::AllLoaded; Ok(()) } pub fn load_all(&mut self) { - if self.load_state == LoadState::AllLoaded { + if self.load_state == LoadState::AllLoaded && !self.evicted_since_full_load { return; } @@ -518,6 +531,7 @@ impl InnerStore { } }); + self.evicted_since_full_load = false; self.load_state = LoadState::AllLoaded; } @@ -584,6 +598,7 @@ impl InnerStore { load_state: LoadState::AllLoaded, config, value_cache_queue: VecDeque::new(), + evicted_since_full_load: false, } } diff --git a/crates/loro-internal/tests/handle_walk_memory.rs b/crates/loro-internal/tests/handle_walk_memory.rs index 74582a501..098e4836d 100644 --- a/crates/loro-internal/tests/handle_walk_memory.rs +++ b/crates/loro-internal/tests/handle_walk_memory.rs @@ -57,10 +57,7 @@ fn build_snapshot(turns: usize, items_per_turn: usize) -> Vec { raw.insert("command", "pnpm test").unwrap(); raw.insert("cwd", "/repo").unwrap(); let text = item - .insert_container( - "title", - loro_internal::handler::TextHandler::new_detached(), - ) + .insert_container("title", loro_internal::handler::TextHandler::new_detached()) .unwrap(); text.insert(0, TEXT, PosType::Unicode).unwrap(); } @@ -104,11 +101,20 @@ fn handle_walk_retains_bounded_memory() { doc.import(&snapshot).unwrap(); let history = Handler::List(doc.get_list("history")); + // Every walk must visit all 13,260 containers (48 even turns × 5 + // containers + 47 odd turns × 277 containers + the root list): a walk that + // silently visits fewer containers would pass the memory thresholds below + // more easily. + let expected_containers = 48 * 5 + 47 * 277 + 1; + let after_import = live_bytes(); let mut handles = 0; walk(&history, &mut handles); + assert_eq!(handles, expected_containers, "first walk is incomplete"); let after_walk1 = live_bytes(); + handles = 0; walk(&history, &mut handles); + assert_eq!(handles, expected_containers, "second walk is incomplete"); let after_walk2 = live_bytes(); let walk1_retained = after_walk1.saturating_sub(after_import); diff --git a/crates/loro-wasm/tests/walk_mem.test.ts b/crates/loro-wasm/tests/walk_mem.test.ts index 37703f7fb..94e421d30 100644 --- a/crates/loro-wasm/tests/walk_mem.test.ts +++ b/crates/loro-wasm/tests/walk_mem.test.ts @@ -96,9 +96,13 @@ describe("container handle walk memory", () => { docB.import(snapshot); global.gc(); const beforeWalk = process.memoryUsage().external; - walk(docB.getList("history") as unknown as ContainerHandle); + const walked = walk(docB.getList("history") as unknown as ContainerHandle); global.gc(); const walkDelta = process.memoryUsage().external - beforeWalk; + + // The walk must have visited every container: a walk that silently + // visits fewer containers would pass the memory threshold more easily. + expect(walked).toEqual(docB.toJSON().history); docB.free(); console.log(