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
47 changes: 47 additions & 0 deletions crates/lance-context-core/examples/codex_recovery_wal_status.rs
Original file line number Diff line number Diff line change
@@ -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}");
}
}
16 changes: 12 additions & 4 deletions crates/lance-context-server/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
),
);
}
Expand Down Expand Up @@ -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,
Expand Down
109 changes: 105 additions & 4 deletions crates/lance-context-server/src/sweeper.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
///
Expand Down Expand Up @@ -189,7 +189,11 @@ pub(crate) async fn resident<S: Clone>(cache: &Mutex<LruCache<String, S>>) -> 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<S: Sweepable>(stores: Vec<(String, S)>, pass_timeout: Duration) {
pub(crate) async fn flush_pass<S: Sweepable>(
stores: Vec<(String, S)>,
pass_timeout: Duration,
merge_slots: Option<Arc<Semaphore>>,
) {
let kind = S::kind();
for (name, store) in stores {
match tokio::time::timeout(pass_timeout, store.flush()).await {
Expand All @@ -199,11 +203,36 @@ pub(crate) async fn flush_pass<S: Sweepable>(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)) => {
Expand Down Expand Up @@ -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<RwLock<GenericStore>> {
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
);
}
}
Loading