From f93aa0dc28c46ba27f001ec9c629ddaa11bc422a Mon Sep 17 00:00:00 2001 From: Beinan Wang Date: Wed, 30 Sep 2026 11:23:16 +0000 Subject: [PATCH] fix: the pending-generation read cap counts every shard, not the worst one A read opens every pending generation of every writer shard, so the cap introduced in #278 must hold for the store as a whole. A store with 20 shards at ~1,400 generations each (25,588 total) passed the per-shard check and one point lookup opened all of them, OOMKilling the worker. The per-shard check stays (it fails fast, before the other shards are even listed); the total is checked once the snapshots are collected. Co-Authored-By: Claude Fable 5 --- .../lance-context-core/src/rollout_store.rs | 51 +++++++++++++++++++ crates/lance-context-core/src/store_base.rs | 31 ++++++++--- 2 files changed, 75 insertions(+), 7 deletions(-) diff --git a/crates/lance-context-core/src/rollout_store.rs b/crates/lance-context-core/src/rollout_store.rs index 97c782f..d12c378 100644 --- a/crates/lance-context-core/src/rollout_store.rs +++ b/crates/lance-context-core/src/rollout_store.rs @@ -4355,6 +4355,57 @@ mod tests { /// Reads open every flushed generation pending merge. Past the cap the /// read is refused with a recognizable error (the server maps it to 503) /// instead of holding a worker's memory hostage; at the cap it passes. + /// The cap holds for the store as a whole: two writer shards each under + /// it still add up to a read that opens every generation of both. + #[test] + fn pending_generation_cap_counts_every_shard() { + let dir = TempDir::new().unwrap(); + let uri = dir.path().to_string_lossy().to_string(); + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async { + let open = |shard: &str| { + RolloutStore::open_with_options( + &uri, + RolloutStoreOptions { + storage_options: None, + session: None, + shard_id: Some(shard.to_string()), + merge_after_generations: None, + merge_max_generations: None, + merge_max_bytes: None, + pending_generations_warn: None, + pending_generations_max: Some(3), + merge_budget: None, + }, + ) + }; + let mut a = open("rollout-a").await.unwrap(); + let b = open("rollout-b").await.unwrap(); + for id in ["a-0", "a-1"] { + a.add(&[assistant_record(id)]).await.unwrap(); + a.flush().await.unwrap(); + } + b.add(&[assistant_record("b-0")]).await.unwrap(); + b.flush().await.unwrap(); + // 2 + 1 = 3: at the cap, every shard well under it. + assert_eq!(a.list(None, None).await.unwrap().len(), 3); + + b.add(&[assistant_record("b-1")]).await.unwrap(); + b.flush().await.unwrap(); + // 2 + 2 = 4: no shard exceeds the cap, the store does. + let err = a.list(None, None).await.expect_err("over the cap in total"); + assert!( + crate::store_base::is_pending_generations_exceeded(&err), + "{err}" + ); + assert!(err.to_string().contains("across 2 shards"), "{err}"); + + // Merging one shard brings the total back under. + assert_eq!(a.cleanup_own_shard().await.unwrap(), 2); + assert_eq!(b.list(None, None).await.unwrap().len(), 4); + }); + } + #[test] fn reads_are_refused_past_the_pending_generation_cap() { let dir = TempDir::new().unwrap(); diff --git a/crates/lance-context-core/src/store_base.rs b/crates/lance-context-core/src/store_base.rs index f5cf1be..8a0e6bb 100644 --- a/crates/lance-context-core/src/store_base.rs +++ b/crates/lance-context-core/src/store_base.rs @@ -257,10 +257,10 @@ pub(crate) struct StorageBaseOptions { /// pending merge (sampled on every LSM read). `None` uses the crate default /// (256); `Some(0)` disables the warning. The metric is always emitted. pub pending_generations_warn: Option, - /// Refuse a read whose shard has more flushed generations pending than - /// this, instead of opening every one of them. `None` uses the crate - /// default (4096); `Some(0)` disables the cap. See - /// [`StorageBase::wal_shard_snapshots`]. + /// Refuse a read that would open more flushed generations than this, in + /// any one shard or across all shards, instead of opening every one of + /// them. `None` uses the crate default (4096); `Some(0)` disables the + /// cap. See [`StorageBase::wal_shard_snapshots`]. pub pending_generations_max: Option, /// Process-wide byte budget shared by every merge this process runs. /// `merge_max_bytes` bounds one merge; this bounds all of them together, @@ -330,8 +330,8 @@ pub(crate) struct StorageBase { merge_max_bytes: usize, /// Per-shard pending-generation count at which reads warn; `0` disables. pending_generations_warn: usize, - /// Refuse reads that would open more flushed generations than this per - /// shard (0 = unbounded). See [`Self::wal_shard_snapshots`]. + /// Refuse reads that would open more flushed generations than this, per + /// shard or in total (0 = unbounded). See [`Self::wal_shard_snapshots`]. pending_generations_max: usize, /// Process-wide merge byte budget; `None` means unbounded. merge_budget: Option>, @@ -1742,7 +1742,24 @@ impl StorageBase { .try_collect() .await?; - Ok(snapshots.into_iter().flatten().collect()) + let snapshots: Vec = snapshots.into_iter().flatten().collect(); + // A read opens every pending generation of every shard, so the cap + // has to hold for the store as a whole, not just the worst shard: 20 + // writer shards at ~1,400 generations each passed the per-shard check + // and still opened 25k datasets on one lookup (OOMKilled the worker). + let total: usize = snapshots + .iter() + .map(|snapshot| snapshot.flushed_generations.len()) + .sum(); + if max_at != 0 && total > max_at { + metrics::counter!("rollout_reads_refused_pending_total").increment(1); + return Err(LanceError::io(format!( + "{PENDING_GENERATIONS_EXCEEDED}: {total} flushed generations pending merge \ + across {} shards (cap {max_at}); retry after the merge catches up", + snapshots.len() + ))); + } + Ok(snapshots) } /// Number of flushed MemWAL generations pending merge into the base table