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
3 changes: 3 additions & 0 deletions crates/lance-context-core/src/datagen_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,8 @@ pub struct DatagenStoreOptions {
/// (256); `Some(0)` disables the warn. The `rollout_wal_pending_generations`
/// histogram is emitted regardless.
pub pending_generations_warn: Option<usize>,
/// See `StorageBaseOptions::pending_generations_max`.
pub pending_generations_max: Option<usize>,
/// Process-wide byte budget shared by every merge this process runs; a
/// merge that cannot fit waits for another to release. `None` disables
/// the bound. See [`crate::merge_budget`] for the design.
Expand Down Expand Up @@ -114,6 +116,7 @@ impl DatagenStore {
merge_max_generations: options.merge_max_generations,
merge_max_bytes: options.merge_max_bytes,
pending_generations_warn: options.pending_generations_warn,
pending_generations_max: options.pending_generations_max,
merge_budget: options.merge_budget.clone(),
session: None,
schema: Arc::new(datagen_log_schema()),
Expand Down
3 changes: 3 additions & 0 deletions crates/lance-context-core/src/generic_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,8 @@ pub struct GenericStoreOptions {
/// alarm. `None` uses the crate default (256); `Some(0)` disables the warn.
/// The `rollout_wal_pending_generations` histogram is emitted regardless.
pub pending_generations_warn: Option<usize>,
/// See `StorageBaseOptions::pending_generations_max`.
pub pending_generations_max: Option<usize>,
/// Process-wide byte budget shared by every merge this process runs; a
/// merge that cannot fit waits for another to release. `None` disables
/// the bound. See [`crate::merge_budget`] for the design.
Expand Down Expand Up @@ -186,6 +188,7 @@ impl GenericStore {
merge_max_generations: options.merge_max_generations,
merge_max_bytes: options.merge_max_bytes,
pending_generations_warn: options.pending_generations_warn,
pending_generations_max: options.pending_generations_max,
merge_budget: options.merge_budget.clone(),
session: options.session,
schema: create_schema,
Expand Down
3 changes: 3 additions & 0 deletions crates/lance-context-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,9 @@ pub mod serde;
mod storage;
mod store;
mod store_base;
pub use store_base::{
is_pending_generations_exceeded, DEFAULT_PENDING_GENERATIONS_MAX, PENDING_GENERATIONS_EXCEEDED,
};

// Request/DTO conversions, exported so the server does not keep its own copies.
// These were duplicated verbatim between here and `routes/`; see #214.
Expand Down
131 changes: 131 additions & 0 deletions crates/lance-context-core/src/rollout_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -389,6 +389,8 @@ pub struct RolloutStoreOptions {
/// (256); `Some(0)` disables the warn. The `rollout_wal_pending_generations`
/// histogram is emitted regardless.
pub pending_generations_warn: Option<usize>,
/// See `StorageBaseOptions::pending_generations_max`.
pub pending_generations_max: Option<usize>,
/// Process-wide byte budget shared by every merge this process runs; a
/// merge that cannot fit waits for another to release. `None` disables
/// the bound. See [`crate::merge_budget`] for the design.
Expand Down Expand Up @@ -480,6 +482,7 @@ impl RolloutStore {
merge_max_generations,
merge_max_bytes,
pending_generations_warn,
pending_generations_max,
merge_budget,
session,
} = options;
Expand All @@ -492,6 +495,7 @@ impl RolloutStore {
merge_max_generations,
merge_max_bytes,
pending_generations_warn,
pending_generations_max,
merge_budget,
session,
schema: Arc::new(rollout_schema()),
Expand Down Expand Up @@ -707,6 +711,11 @@ impl RolloutStore {
self.base.create_key_btree_index().await
}

/// See `StorageBase::extend_key_btree_index`.
pub async fn extend_id_btree_index(&mut self) -> LanceResult<usize> {
self.base.extend_key_btree_index().await
}

/// See `StorageBase::has_key_btree_index`.
pub async fn has_id_btree_index(&self) -> LanceResult<bool> {
self.base.has_key_btree_index().await
Expand Down Expand Up @@ -2525,6 +2534,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
session: None,
schema: legacy_schema.clone(),
Expand Down Expand Up @@ -2846,6 +2856,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -2892,6 +2903,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
};

Expand Down Expand Up @@ -2943,6 +2955,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -2996,6 +3009,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
};

Expand Down Expand Up @@ -3219,6 +3233,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
..Default::default()
},
Expand Down Expand Up @@ -3283,6 +3298,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
..Default::default()
},
Expand Down Expand Up @@ -3385,6 +3401,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -3445,6 +3462,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -3486,6 +3504,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -3623,6 +3642,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -3711,6 +3731,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -3799,6 +3820,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -3848,6 +3870,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -3919,6 +3942,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -3958,6 +3982,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -4018,6 +4043,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -4061,6 +4087,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -4272,6 +4299,107 @@ mod tests {
});
}

/// Every merge appends fragments the id BTree does not cover, and
/// `merge_insert` full-scans exactly those. `extend_id_btree_index`
/// appends an index delta over them; a fully covered table is a no-op
/// and a table without the BTree is left alone.
#[test]
fn extend_id_btree_index_covers_fragments_added_since_build() {
use lance::index::DatasetIndexInternalExt as _;
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(&uri).await.unwrap();
store.add(&[assistant_record("a-0")]).await.unwrap();
store.flush().await.unwrap();
store.cleanup_own_shard().await.unwrap();

// No BTree yet: nothing to extend.
assert_eq!(store.extend_id_btree_index().await.unwrap(), 0);

store.create_id_btree_index().await.unwrap();
assert_eq!(store.extend_id_btree_index().await.unwrap(), 0);

// Two merges land two fragments the index knows nothing about.
for id in ["a-1", "a-2"] {
store.add(&[assistant_record(id)]).await.unwrap();
store.flush().await.unwrap();
store.cleanup_own_shard().await.unwrap();
}
let unindexed = |s: &RolloutStore| {
let dataset = s.base.dataset.clone();
async move {
dataset
.unindexed_fragments(ROLLOUT_ID_INDEX_NAME)
.await
.unwrap()
.len()
}
};
assert_eq!(unindexed(&store).await, 2);

assert_eq!(store.extend_id_btree_index().await.unwrap(), 2);
assert_eq!(unindexed(&store).await, 0);
assert!(store.has_id_btree_index().await.unwrap());
assert_eq!(store.extend_id_btree_index().await.unwrap(), 0);

// Rows are still merged correctly through the extended index.
store.add(&[assistant_record("a-0")]).await.unwrap();
store.flush().await.unwrap();
store.cleanup_own_shard().await.unwrap();
assert_eq!(store.list(None, None).await.unwrap().len(), 3);
});
}

/// Reads open every flushed generation pending merge. Past the cap the
/// read is refused with a recognizable error (the server maps it to 503)
/// instead of holding a worker's memory hostage; at the cap it passes.
#[test]
fn reads_are_refused_past_the_pending_generation_cap() {
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("rollout-0".to_string()),
merge_after_generations: None,
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: Some(2),
merge_budget: None,
},
)
.await
.unwrap();

for id in ["a-0", "a-1"] {
store.add(&[assistant_record(id)]).await.unwrap();
store.flush().await.unwrap();
}
// At the cap: still readable.
assert_eq!(store.list(None, None).await.unwrap().len(), 2);

store.add(&[assistant_record("a-2")]).await.unwrap();
store.flush().await.unwrap();
let err = store.list(None, None).await.expect_err("over the cap");
assert!(
crate::store_base::is_pending_generations_exceeded(&err),
"{err}"
);
assert!(err.to_string().contains("3 flushed generations"), "{err}");

// The merge drains the backlog and reads resume.
assert_eq!(store.cleanup_own_shard().await.unwrap(), 3);
assert_eq!(store.list(None, None).await.unwrap().len(), 3);
});
}

#[test]
fn create_id_btree_index_builds_and_is_idempotent() {
// Building the BTree index on `id` must succeed even though the
Expand Down Expand Up @@ -4638,6 +4766,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
..Default::default()
},
Expand Down Expand Up @@ -4738,6 +4867,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down Expand Up @@ -4772,6 +4902,7 @@ mod tests {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
},
)
Expand Down
5 changes: 5 additions & 0 deletions crates/lance-context-core/src/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -329,6 +329,8 @@ pub struct ContextStoreOptions {
/// (256); `Some(0)` disables the warn. The `rollout_wal_pending_generations`
/// histogram is emitted regardless.
pub pending_generations_warn: Option<usize>,
/// Refuse reads once more than this many flushed generations are pending merge.
pub pending_generations_max: Option<usize>,
/// Process-wide byte budget shared by every merge this process runs; a
/// merge that cannot fit waits for another to release. `None` disables
/// the bound. See [`crate::merge_budget`] for the design.
Expand Down Expand Up @@ -365,6 +367,7 @@ impl Default for ContextStoreOptions {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
// Read-your-write by default; see the field docs.
seal_on_add: true,
Expand Down Expand Up @@ -648,6 +651,7 @@ impl ContextStore {
merge_max_generations: options.merge_max_generations,
merge_max_bytes: options.merge_max_bytes,
pending_generations_warn: options.pending_generations_warn,
pending_generations_max: options.pending_generations_max,
merge_budget: options.merge_budget.clone(),
session: None,
schema: Arc::new(arrow_schema.clone()),
Expand Down Expand Up @@ -2273,6 +2277,7 @@ impl ContextStore {
merge_max_generations: None,
merge_max_bytes: None,
pending_generations_warn: None,
pending_generations_max: None,
merge_budget: None,
// A compactor never appends, so the seal mode is irrelevant to it;
// deferring keeps it from ever emitting a generation.
Expand Down
Loading
Loading