Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
51 changes: 51 additions & 0 deletions crates/lance-context-core/src/rollout_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
31 changes: 24 additions & 7 deletions crates/lance-context-core/src/store_base.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<usize>,
/// 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<usize>,
/// Process-wide byte budget shared by every merge this process runs.
/// `merge_max_bytes` bounds one merge; this bounds all of them together,
Expand Down Expand Up @@ -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<Arc<MergeMemoryBudget>>,
Expand Down Expand Up @@ -1742,7 +1742,24 @@ impl StorageBase {
.try_collect()
.await?;

Ok(snapshots.into_iter().flatten().collect())
let snapshots: Vec<ShardSnapshot> = 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
Expand Down
Loading