From a9a41c123e8578ca2406fff2d138b7e7a4364849 Mon Sep 17 00:00:00 2001 From: Beinan Wang Date: Thu, 1 Oct 2026 08:26:19 +0000 Subject: [PATCH 1/2] fix(server): the count-triggered self-merge takes the worker's merge slot ROLLOUT_MERGE_AFTER_GENERATIONS merges a shard as soon as it has N pending generations, on the 30 s flush sweeper. That merge bypassed the per-worker merge slot (ROLLOUT_MERGE_CONCURRENCY), so turning the count trigger on would stack sweeper merges on top of the master's requests with no bound. It now acquires the same slot, so the bound holds no matter who starts the merge. Needed to enable pending-driven merging in production: hot stores get merged when they accumulate generations, not when the master's 600 s sweep comes around, which is what keeps reads over them fast. Co-Authored-By: Claude Fable 5 --- .../examples/codex_recovery_wal_status.rs | 47 +++++++++++++++++++ crates/lance-context-server/src/state.rs | 16 +++++-- crates/lance-context-server/src/sweeper.rs | 28 +++++++++-- 3 files changed, 83 insertions(+), 8 deletions(-) create mode 100644 crates/lance-context-core/examples/codex_recovery_wal_status.rs diff --git a/crates/lance-context-core/examples/codex_recovery_wal_status.rs b/crates/lance-context-core/examples/codex_recovery_wal_status.rs new file mode 100644 index 0000000..f89c6ac --- /dev/null +++ b/crates/lance-context-core/examples/codex_recovery_wal_status.rs @@ -0,0 +1,47 @@ +//! Read only base + shard manifests. Never opens a generation dataset or payload. +use lance_context_core::{RolloutStore, RolloutStoreOptions}; +use std::{ + io::{self, BufRead}, + time::Duration, +}; +fn main() { + let rt = tokio::runtime::Builder::new_multi_thread() + .worker_threads(2) + .enable_all() + .build() + .unwrap(); + let session = RolloutStore::build_session(96 * 1024 * 1024, 32 * 1024 * 1024); + for line in io::stdin().lock().lines() { + let uri = line.unwrap(); + if uri.is_empty() { + continue; + } + let start = std::time::Instant::now(); + let options = RolloutStoreOptions { + session: Some(session.clone()), + pending_generations_max: Some(0), + pending_generations_warn: Some(0), + ..Default::default() + }; + let result = rt.block_on(async { + tokio::time::timeout(Duration::from_secs(40), async { + let store = RolloutStore::open_existing_with_options(&uri, options).await?; + let pending = store.pending_wal_generations().await?; + Ok::<_, lance::Error>(( + pending, + store.version(), + store.compaction_stats().total_fragments, + )) + }) + .await + }); + let row = match result { + Ok(Ok((pending, version, fragments))) => { + serde_json::json!({"uri":uri,"pending":pending,"version":version,"fragments":fragments,"seconds":start.elapsed().as_secs_f64()}) + } + Ok(Err(e)) => serde_json::json!({"uri":uri,"error":e.to_string()}), + Err(e) => serde_json::json!({"uri":uri,"error":e.to_string()}), + }; + println!("{row}"); + } +} diff --git a/crates/lance-context-server/src/state.rs b/crates/lance-context-server/src/state.rs index a34a88d..bc3ce77 100644 --- a/crates/lance-context-server/src/state.rs +++ b/crates/lance-context-server/src/state.rs @@ -1016,15 +1016,18 @@ impl AppState { tokio::join!( sweeper::flush_pass( sweeper::resident(&state.rollout_stores).await, - pass_timeout + pass_timeout, + state.merge_slots.clone(), ), sweeper::flush_pass( sweeper::resident(&state.datagen_stores).await, - pass_timeout + pass_timeout, + state.merge_slots.clone(), ), sweeper::flush_pass( sweeper::resident(&state.generic_stores).await, - pass_timeout + pass_timeout, + state.merge_slots.clone(), ), ); } @@ -1264,7 +1267,12 @@ mod tests { .await .unwrap(); - sweeper::flush_pass(sweeper::resident(&state.generic_stores).await, timeout).await; + sweeper::flush_pass( + sweeper::resident(&state.generic_stores).await, + timeout, + state.merge_slots.clone(), + ) + .await; assert_eq!( generic.read().await.list(None, None).await.unwrap().len(), 1, diff --git a/crates/lance-context-server/src/sweeper.rs b/crates/lance-context-server/src/sweeper.rs index c13cc07..e10327c 100644 --- a/crates/lance-context-server/src/sweeper.rs +++ b/crates/lance-context-server/src/sweeper.rs @@ -20,7 +20,7 @@ use std::time::Duration; use lance_context_core::{DatagenStore, GenericStore, RolloutStore}; use lru::LruCache; -use tokio::sync::{Mutex, RwLock}; +use tokio::sync::{Mutex, RwLock, Semaphore}; /// A store the sweepers can maintain. /// @@ -189,7 +189,11 @@ pub(crate) async fn resident(cache: &Mutex>) -> Ve /// and alerts keep working; the new `kind` label is what distinguishes the /// store types. Renaming them would be a silent breakage for anyone graphing /// these today. -pub(crate) async fn flush_pass(stores: Vec<(String, S)>, pass_timeout: Duration) { +pub(crate) async fn flush_pass( + stores: Vec<(String, S)>, + pass_timeout: Duration, + merge_slots: Option>, +) { let kind = S::kind(); for (name, store) in stores { match tokio::time::timeout(pass_timeout, store.flush()).await { @@ -199,11 +203,27 @@ pub(crate) async fn flush_pass(stores: Vec<(String, S)>, pass_time // The count-triggered merge rides this timer, but it is a merge, // not a flush: its outcome is reported under the cleanup // counters so a failing merge cannot masquerade as a failing - // flush on the dashboards. + // flush on the dashboards. It is also a merge for memory + // purposes: it takes the same per-worker slot the master's + // requests take (ROLLOUT_MERGE_CONCURRENCY), so enabling the + // count trigger cannot stack merges past that bound. + let merged = async { + let _slot = match &merge_slots { + Some(slots) => Some( + slots + .clone() + .acquire_owned() + .await + .expect("merge slot semaphore is never closed"), + ), + None => None, + }; + store.merge_if_due().await + }; report_merge( &name, kind, - tokio::time::timeout(pass_timeout, store.merge_if_due()).await, + tokio::time::timeout(pass_timeout, merged).await, ); } Ok(Err(error)) => { From 2b0ada99ca2256ec7df007beb1dfe578db88f609 Mon Sep 17 00:00:00 2001 From: Beinan Wang Date: Thu, 1 Oct 2026 08:43:37 +0000 Subject: [PATCH 2/2] fix(server): try the merge slot, never wait on it, in the flush pass Review on #289: awaiting the slot before merge_if_due stalled the flush of every store behind the first one whenever all slots were busy, even with the count trigger off. The pass now try_acquires; when the slots are full the self-merge is skipped (rollout_wal_self_merge_skipped_total) and flushing continues. Regression tests: two stores flush with all slots held and the trigger off; a store merges when a slot is free. Co-Authored-By: Claude Fable 5 --- crates/lance-context-server/src/sweeper.rs | 109 ++++++++++++++++++--- 1 file changed, 95 insertions(+), 14 deletions(-) diff --git a/crates/lance-context-server/src/sweeper.rs b/crates/lance-context-server/src/sweeper.rs index e10327c..3f04f44 100644 --- a/crates/lance-context-server/src/sweeper.rs +++ b/crates/lance-context-server/src/sweeper.rs @@ -204,20 +204,29 @@ pub(crate) async fn flush_pass( // not a flush: its outcome is reported under the cleanup // counters so a failing merge cannot masquerade as a failing // flush on the dashboards. It is also a merge for memory - // purposes: it takes the same per-worker slot the master's - // requests take (ROLLOUT_MERGE_CONCURRENCY), so enabling the - // count trigger cannot stack merges past that bound. - let merged = async { - let _slot = match &merge_slots { - Some(slots) => Some( - slots - .clone() - .acquire_owned() - .await - .expect("merge slot semaphore is never closed"), - ), - None => None, - }; + // purposes: it shares the per-worker slot with the master's + // merge requests (ROLLOUT_MERGE_CONCURRENCY). The slot is only + // *tried*, never awaited: this pass flushes every resident + // store in sequence, and blocking on a busy slot here would + // hold up the flush of every store behind this one. When the + // slots are full the merge is skipped; the next pass (30 s) + // tries again, and the master's sweep covers the store anyway. + let slot = match &merge_slots { + Some(slots) => match slots.clone().try_acquire_owned() { + Ok(permit) => Some(permit), + Err(_) => { + metrics::counter!( + "rollout_wal_self_merge_skipped_total", + "kind" => kind + ) + .increment(1); + continue; + } + }, + None => None, + }; + let merged = async move { + let _slot = slot; store.merge_if_due().await }; report_merge( @@ -409,4 +418,76 @@ mod tests { .unwrap(); assert_eq!(reclaimed, 3); } + + /// The flush pass must flush every store even when every merge slot is + /// taken: the self-merge only *tries* the slot and skips, it never waits. + /// Two stores, count trigger off, slots exhausted — both must flush, and + /// nothing must merge. + #[tokio::test] + async fn flush_pass_never_waits_on_a_busy_merge_slot() { + let dir_a = TempDir::new().unwrap(); + let dir_b = TempDir::new().unwrap(); + // seal_on_add is on in `generic_with_pending`, so build stores with + // unsealed rows by hand: count trigger off, rows left in the memtable. + async fn unsealed(dir: &TempDir) -> Arc> { + let uri = dir.path().to_string_lossy().to_string(); + let store = GenericStore::open( + &uri, + spec(), + GenericStoreOptions { + merge_after_generations: Some(0), + seal_on_add: false, + ..Default::default() + }, + ) + .await + .unwrap(); + store + .add(&[json!({"id": "r0"}).as_object().unwrap().clone()]) + .await + .unwrap(); + assert_eq!(store.pending_wal_generations().await.unwrap(), 0); + Arc::new(RwLock::new(store)) + } + let a = unsealed(&dir_a).await; + let b = unsealed(&dir_b).await; + let slots = Arc::new(Semaphore::new(1)); + let _held = slots.clone().acquire_owned().await.unwrap(); + + tokio::time::timeout( + Duration::from_secs(10), + flush_pass( + vec![("a".to_string(), a.clone()), ("b".to_string(), b.clone())], + Duration::from_secs(5), + Some(slots), + ), + ) + .await + .expect("flush pass must not block on the merge slot"); + + // Both flushed: each now has exactly one sealed generation, unmerged. + for store in [&a, &b] { + assert_eq!( + store.read().await.pending_wal_generations().await.unwrap(), + 1 + ); + } + } + + /// With a free slot and the count trigger set, the pass still merges. + #[tokio::test] + async fn flush_pass_merges_when_a_slot_is_free() { + let dir = TempDir::new().unwrap(); + let store = generic_with_pending(&dir, 2, 2).await; + flush_pass( + vec![("s".to_string(), store.clone())], + Duration::from_secs(5), + Some(Arc::new(Semaphore::new(1))), + ) + .await; + assert_eq!( + store.read().await.pending_wal_generations().await.unwrap(), + 0 + ); + } }