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..3f04f44 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,36 @@ 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 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( &name, kind, - tokio::time::timeout(pass_timeout, store.merge_if_due()).await, + tokio::time::timeout(pass_timeout, merged).await, ); } Ok(Err(error)) => { @@ -389,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 + ); + } }