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
8 changes: 8 additions & 0 deletions crates/lance-context-core/src/generic_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -382,6 +382,14 @@ impl GenericStore {
self.base.cleanup_own_shard().await
}

/// The shared-lock half of [`Self::maybe_merge_wal`]: `None` unless the
/// count trigger is configured and met.
pub async fn prepare_count_merge(
&self,
) -> LanceResult<Option<(ShardManifestStore, ShardManifest, PreparedMerge)>> {
self.base.prepare_count_merge().await
}

/// The shared-lock half of [`Self::cleanup_wal`]: seal, then read a
/// budgeted prefix of flushed generations into memory. Callers holding a
/// read lock run this while appends continue, then take the write lock
Expand Down
8 changes: 8 additions & 0 deletions crates/lance-context-core/src/rollout_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -629,6 +629,14 @@ impl RolloutStore {
self.base.prepare_merge_if_ready(threshold).await
}

/// The shared-lock half of [`Self::maybe_merge_own_shard`]: `None` unless
/// the count trigger is configured and met.
pub async fn prepare_count_merge(
&self,
) -> LanceResult<Option<(ShardManifestStore, ShardManifest, PreparedMerge)>> {
self.base.prepare_count_merge().await
}

/// [`Self::prepare_merge_if_ready`], but seals the active memtable first —
/// the time-triggered behavior of [`Self::cleanup_own_shard`].
pub async fn prepare_cleanup_merge(
Expand Down
16 changes: 16 additions & 0 deletions crates/lance-context-core/src/store_base.rs
Original file line number Diff line number Diff line change
Expand Up @@ -783,6 +783,22 @@ impl StorageBase {
self.prepare_merge_if_ready_inner(threshold, false).await
}

/// The shared-lock half of the *count*-triggered merge: prepare only if the
/// shard has at least `merge_after_generations` flushed generations pending
/// (`0` disables the trigger, so this returns `None`). Same threshold as
/// [`Self::maybe_merge_own_shard`], but split so the expensive generation
/// read runs under a read lock and only [`Self::commit_prepared_merge`]
/// needs the write lock.
pub async fn prepare_count_merge(
&self,
) -> LanceResult<Option<(ShardManifestStore, ShardManifest, PreparedMerge)>> {
if self.merge_after_generations == 0 {
return Ok(None);
}
self.prepare_merge_if_ready_inner(self.merge_after_generations, false)
.await
}

/// [`Self::prepare_merge_if_ready`], but seals the active memtable *before*
/// consulting the manifest — the time-triggered (`threshold = 1`) behavior
/// of [`Self::cleanup_own_shard`]. See that method for why the ordering is
Expand Down
21 changes: 19 additions & 2 deletions crates/lance-context-server/src/routes/generic.rs
Original file line number Diff line number Diff line change
Expand Up @@ -271,8 +271,25 @@ pub async fn merge_generic_wal(
Path(name): Path<String>,
) -> Result<Json<serde_json::Value>, AppError> {
let store = state.get_or_open_generic_store(&name).await?;
let mut guard = store.write().await;
let reclaimed = guard.cleanup_wal().await.map_err(AppError::from_lance)?;
// Same prepare/commit split as the sweeper: the object-storage read of the
// generations runs under the shared lock so the store keeps serving.
let prepared = {
let guard = store.read().await;
guard
.prepare_cleanup_merge()
.await
.map_err(AppError::from_lance)?
};
let reclaimed = match prepared {
Some((manifest_store, manifest, prepared)) => {
let mut guard = store.write().await;
guard
.commit_prepared_merge(&manifest_store, &manifest, prepared)
.await
.map_err(AppError::from_lance)?
}
None => 0,
};
Ok(Json(serde_json::json!({ "reclaimed": reclaimed })))
}

Expand Down
33 changes: 22 additions & 11 deletions crates/lance-context-server/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -945,8 +945,8 @@ impl AppState {
/// deferred seal and genuinely depend on this; datagen seals on each append,
/// so its pass is a no-op in steady state and is kept only for symmetry.
///
/// For rollout it also runs the count-triggered merge
/// ([`RolloutStore::maybe_merge_own_shard`]): the read-amplification bound
/// For rollout and generic it also runs the count-triggered merge
/// ([`sweeper::Sweepable::merge_if_due`]): the read-amplification bound
/// that formerly lived on the append path now rides this timer. The heavier
/// time-based cleanup/merge remains on [`Self::spawn_global_sweeper`].
///
Expand All @@ -969,15 +969,26 @@ impl AppState {
let Some(state) = weak.upgrade() else {
return;
};
sweeper::flush_pass(sweeper::resident(&state.rollout_stores).await, pass_timeout)
.await;
// Datagen seals on every append, so its pass is a no-op in
// steady state; generic stores default to a deferred seal and
// genuinely depend on this.
sweeper::flush_pass(sweeper::resident(&state.datagen_stores).await, pass_timeout)
.await;
sweeper::flush_pass(sweeper::resident(&state.generic_stores).await, pass_timeout)
.await;
// Kinds run concurrently for the same reason the cleanup sweeper
// does: this pass now carries the count-triggered merge, so a
// slow walk over hundreds of rollout stores must not delay
// generic's tick. Datagen seals on every append, so its pass is
// a no-op in steady state; generic stores default to a deferred
// seal and genuinely depend on this.
tokio::join!(
sweeper::flush_pass(
sweeper::resident(&state.rollout_stores).await,
pass_timeout
),
sweeper::flush_pass(
sweeper::resident(&state.datagen_stores).await,
pass_timeout
),
sweeper::flush_pass(
sweeper::resident(&state.generic_stores).await,
pass_timeout
),
);
}
}))
}
Expand Down
Loading
Loading