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
63 changes: 63 additions & 0 deletions crates/lance-context-core/src/rollout_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -722,7 +722,21 @@ impl RolloutStore {
/// only on the shard this instance owns, so it is safe to call concurrently
/// with this instance's own appends but must not target another instance's
/// shard.
///
/// # Seals first
///
/// This flushes the active memtable before looking at the manifest. Without
/// that, a deployment with the periodic flush sweeper disabled
/// (`ROLLOUT_FLUSH_INTERVAL_SECS=0`) could never make progress: nothing
/// would seal the memtable, so `flushed_generations` would stay empty, so
/// the threshold check below would return `0` and never reach the merge —
/// leaving rows durable but permanently invisible until a process restart
/// replayed the WAL. Sealing here makes the cleanup path a genuine
/// standalone fallback, as `ROLLOUT_FLUSH_INTERVAL_SECS`'s documentation
/// already promised.
pub async fn cleanup_own_shard(&mut self) -> LanceResult<usize> {
// Materialize anything buffered so it is eligible for this pass.
self.flush().await?;
// Threshold `1`: merge whenever at least one generation is pending. The
// time trigger must not depend on the count threshold — that is what
// makes the two triggers a true OR.
Expand Down Expand Up @@ -3060,6 +3074,55 @@ mod tests {
});
}

#[test]
fn cleanup_own_shard_seals_before_merging() {
// Regression guard for the `ROLLOUT_FLUSH_INTERVAL_SECS=0` trap: with no
// periodic flush, nothing seals the active memtable, so
// `flushed_generations` stays empty and the threshold check in
// `merge_own_shard_if_ready` used to return 0 without ever merging —
// leaving rows durable but permanently invisible.
//
// `cleanup_own_shard` now flushes first, making it a genuine standalone
// fallback, which is what the config docs already claimed.
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 mut store = RolloutStore::open_with_options(
&uri,
RolloutStoreOptions {
storage_options: None,
session: None,
shard_id: Some("cleanup-seal-0".to_string()),
// Count trigger disabled: cleanup is the only path that can
// make this row visible, exactly as with flush interval 0.
merge_after_generations: Some(0),
},
)
.await
.unwrap();

store.add(&[assistant_record("c-0")]).await.unwrap();

// Nothing sealed yet: invisible, and no generation pending.
assert!(store.list(None, None).await.unwrap().is_empty());
assert_eq!(
store.observe().await.unwrap().pending_wal_generations,
0,
"precondition: the memtable is unsealed, so no generation exists"
);

// A single cleanup pass must seal, merge, and expose the row.
let reclaimed = store.cleanup_own_shard().await.unwrap();
assert_eq!(reclaimed, 1, "cleanup must seal then merge the generation");

let seen = store.list(None, None).await.unwrap();
assert_eq!(seen.len(), 1, "row must be visible after cleanup alone");
assert_eq!(seen[0].id, "c-0");
});
}

#[test]
fn distinct_shards_share_one_dataset() {
// Two instances writing distinct shards of the same dataset both
Expand Down
11 changes: 9 additions & 2 deletions crates/lance-context-server/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,8 +53,15 @@ pub struct ServerConfig {
/// storage) but are not visible to reads until the memtable is flushed, so
/// this interval bounds read-after-write latency. Decoupling flush from the
/// append path is what lets concurrent appends run without serializing behind
/// a per-append seal. Default `30`; `0` disables periodic flush (rows then
/// only become visible when the cleanup/merge path flushes them).
/// a per-append seal. Default `30`.
///
/// `0` disables periodic flush, leaving the cleanup sweeper
/// (`ROLLOUT_CLEANUP_INTERVAL_SECS`) as the only thing that seals memtables
/// — it flushes before merging, so it is a sufficient fallback, but
/// read-after-write latency is then bounded by the *cleanup* interval
/// instead. Setting **both** to `0` means nothing ever seals: appends stay
/// durable but invisible until the process restarts and replays the WAL.
/// The server warns at startup in that configuration.
#[arg(long, env = "ROLLOUT_FLUSH_INTERVAL_SECS", default_value = "30")]
pub rollout_flush_interval_secs: u64,

Expand Down
14 changes: 14 additions & 0 deletions crates/lance-context-server/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,20 @@ async fn main() {
// without serializing concurrent appends. Detached for the server lifetime.
let _flush_sweeper = state.spawn_flush_sweeper();

// With both sweepers off nothing ever seals the active memtable, so rollout
// appends stay durable-but-invisible until the process restarts and replays
// the WAL. `cleanup_own_shard` flushes before merging, so the cleanup
// sweeper alone is a sufficient fallback — but if neither runs, warn loudly
// rather than leaving the operator to discover it as missing rows.
if state.rollout_flush_interval_secs == 0 && state.rollout_cleanup_interval_secs == 0 {
tracing::warn!(
"ROLLOUT_FLUSH_INTERVAL_SECS=0 and ROLLOUT_CLEANUP_INTERVAL_SECS=0: \
nothing will seal MemWAL memtables, so rollout appends will be durable \
but invisible to reads until this process restarts. Set at least one \
of them to a non-zero interval."
);
}

// Install the Prometheus recorder once, before any metrics are emitted.
let metrics_handle = lance_context_metrics::install_recorder();

Expand Down
Loading