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 README.md
Original file line number Diff line number Diff line change
Expand Up @@ -362,6 +362,53 @@ auto ds = lance::Dataset::open_with_session(session, "data.lance");
auto stats = session.cache_stats();
```

### Prewarm an index synchronously

Prewarm a logical index before serving queries to move cold index reads out of
foreground query execution. The call blocks until the Rust SDK's
`Dataset::prewarm_index` finishes for all segments with that name in the dataset
handle's snapshot. It does not create an index or change the dataset version.

```c
LanceSession* session = lance_session_new(64 * 1024 * 1024, 16 * 1024 * 1024);
if (!session) return -1;
LanceDataset* ds = lance_dataset_open_with_session("data.lance", NULL, 0, session);
if (!ds) {
lance_session_close(session);
return -1;
}
/* Assumes an index named embedding_idx already exists. */
int32_t status = lance_dataset_prewarm_index(ds, "embedding_idx");
if (status != 0) {
const char* message = lance_last_error_message();
fprintf(stderr, "%s\n", message);
lance_free_string(message);
}
lance_dataset_close(ds);
/* Keep session alive and reuse it when opening subsequent query datasets. */
lance_session_close(session);
```

```cpp
lance::Session session(64 * 1024 * 1024, 16 * 1024 * 1024);
auto ds = lance::Dataset::open_with_session(session, "data.lance");
ds.prewarm_index("embedding_idx"); // Blocks; throws lance::Error on failure.
```

The destination is the dataset's process-local session index cache. Reusing the
same session across handles preserves warmed entries after a dataset closes.
Independent sessions and other processes do not share that cache. Entries remain
evictable: success guarantees completion, not that the entire index fits in the
cache or remains resident. Replacing an index creates new segments that must be
warmed separately; historical handles still use their own snapshots. A failure
can leave partially warmed entries in the cache; retrying is safe.

This binding uses the SDK's index-specific prewarm behavior without additional
options. Tests cover multi-segment IVF-Flat and B-tree indexes; other index types
and formats follow the underlying SDK's support and error behavior. It does not
prewarm ordinary data columns or eliminate result-row reads, and it adds no
persistent cache, automatic refresh, or background task management.

### Open at a specific version

`lance_dataset_open` takes a `version` argument — `0` means the latest, any
Expand Down
20 changes: 20 additions & 0 deletions include/lance/lance.h
Original file line number Diff line number Diff line change
Expand Up @@ -2051,6 +2051,26 @@ int32_t lance_dataset_commit_index_segments(
size_t segment_count
);

/**
* Synchronously prewarm a logical index in the dataset's session index cache.
*
* Delegates to the Rust SDK's prewarm_index for all segments with index_name
* in the handle's current snapshot. Does not change the dataset version.
* Index types/formats follow the SDK's prewarm support and error behavior.
*
* Reuse the same LanceSession when opening subsequent query datasets to reuse
* warmed entries across handles. Cache entries remain evictable; success does
* not guarantee full or permanent residency. Ordinary data columns are not
* prewarmed. An error may leave partially warmed entries; retrying is safe.
*
* dataset must be valid and index_name must be a non-empty, NUL-terminated
* UTF-8 string. Both must remain valid until this blocking call returns.
* Returns 0 on success, -1 on error (see lance_last_error_code/message).
* NULL/empty/invalid UTF-8 arguments yield LANCE_ERR_INVALID_ARGUMENT;
* an index absent from this snapshot yields LANCE_ERR_NOT_FOUND.
*/
int32_t lance_dataset_prewarm_index(const LanceDataset* dataset, const char* index_name);

/** Drop an index by name. Returns -1 (NOT_FOUND) if no such index. */
int32_t lance_dataset_drop_index(LanceDataset* dataset, const char* name);

Expand Down
8 changes: 8 additions & 0 deletions include/lance/lance.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -1017,6 +1017,14 @@ class Dataset {
check_error();
}

/// Synchronously prewarm all segments of an index in this snapshot.
/// Uses the dataset's session index cache; entries remain evictable.
/// Throws lance::Error on failure, including an absent or empty index name.
void prewarm_index(const std::string& index_name) const {
if (lance_dataset_prewarm_index(handle_.get(), index_name.c_str()) != 0)
check_error();
}

/// Number of user indexes (excludes system indexes).
uint64_t index_count() const {
uint64_t n = lance_dataset_index_count(handle_.get());
Expand Down
32 changes: 32 additions & 0 deletions src/index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -269,6 +269,38 @@ unsafe fn dataset_index_segments_inner(
Ok(0)
}

/// Synchronously prewarm all segments of a logical index in this snapshot.
///
/// Uses the dataset's session index cache. Success does not pin cache entries.
#[unsafe(no_mangle)]
pub unsafe extern "C" fn lance_dataset_prewarm_index(
dataset: *const LanceDataset,
index_name: *const c_char,
) -> i32 {
ffi_try!(unsafe { prewarm_index_inner(dataset, index_name) }, neg)
}

unsafe fn prewarm_index_inner(
dataset: *const LanceDataset,
index_name: *const c_char,
) -> Result<i32> {
if dataset.is_null() {
return Err(lance_core::Error::invalid_input_source(
"dataset must not be NULL".into(),
));
}
let index_name = unsafe { helpers::parse_c_string(index_name)? }
.filter(|name| !name.is_empty())
.ok_or_else(|| {
lance_core::Error::invalid_input_source("index_name must not be NULL or empty".into())
})?;
// Keep one snapshot and its shared session alive without holding the
// dataset lock across potentially long-running index IO.
let snapshot = unsafe { &*dataset }.snapshot();
block_on(snapshot.prewarm_index(index_name))?;
Ok(0)
}

/// Drop an index by name.
#[unsafe(no_mangle)]
pub unsafe extern "C" fn lance_dataset_drop_index(
Expand Down
176 changes: 176 additions & 0 deletions tests/c_api_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17804,3 +17804,179 @@ fn test_scanner_nearest_batch_null_arguments() {
lance_dataset_close(ds);
}
}

#[test]
fn test_prewarm_index_invalid_arguments_and_missing_index() {
let (_tmp, uri) = create_test_dataset();
unsafe {
let ds = lance_dataset_open(c_str(&uri).as_ptr(), ptr::null(), 0);
assert!(!ds.is_null());
let name = c_str("missing_idx");
assert_eq!(lance_dataset_prewarm_index(ptr::null(), name.as_ptr()), -1);
assert_eq!(lance_last_error_code(), LanceErrorCode::InvalidArgument);
for name in [ptr::null(), c"".as_ptr(), c"\xff".as_ptr()] {
assert_eq!(lance_dataset_prewarm_index(ds, name), -1);
assert_eq!(lance_last_error_code(), LanceErrorCode::InvalidArgument);
}
assert_eq!(lance_dataset_prewarm_index(ds, name.as_ptr()), -1);
assert_eq!(lance_last_error_code(), LanceErrorCode::NotFound);
let message = lance_last_error_message();
assert!(!message.is_null());
assert!(
std::ffi::CStr::from_ptr(message)
.to_string_lossy()
.contains("missing_idx")
);
lance_free_string(message);
assert_eq!(lance_dataset_count_rows(ds), 5);
lance_dataset_close(ds);
}
}

#[test]
fn test_prewarm_index_scalar_segments_and_snapshot() {
let (_tmp, uri, _) = create_scalar_segment_fixture(lance_index::IndexType::BTree, false);
unsafe {
let ds = lance_dataset_open(c_str(&uri).as_ptr(), ptr::null(), 0);
assert!(!ds.is_null());
assert_eq!(lance_dataset_index_count(ds), 2);
let version = lance_dataset_version(ds);
let historical = lance_dataset_open(c_str(&uri).as_ptr(), ptr::null(), version);
assert!(!historical.is_null());
assert_eq!(lance_dataset_prewarm_index(ds, c"key_idx".as_ptr()), 0);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Assert the scalar index was actually warmed

This test currently proves that the B-tree call succeeds for a multi-segment snapshot, but not that either segment was loaded into the session cache. Could we also query both segments through the same session and assert the relevant cache-hit/load statistics? That would prevent a no-op or existence-only implementation from satisfying the test.

assert_eq!(lance_dataset_version(ds), version);
assert_eq!(lance_dataset_drop_index(ds, c"key_idx".as_ptr()), 0);
assert_eq!(lance_dataset_prewarm_index(ds, c"key_idx".as_ptr()), -1);
assert_eq!(lance_last_error_code(), LanceErrorCode::NotFound);
assert_eq!(
lance_dataset_prewarm_index(historical, c"key_idx".as_ptr()),
0
);
assert_eq!(lance_last_error_code(), LanceErrorCode::Ok);
assert_eq!(lance_dataset_version(historical), version);
lance_dataset_close(historical);
lance_dataset_close(ds);
}
}

#[test]
fn test_prewarm_index_vector_segments_reuse_shared_session() {
use lance::index::{DatasetIndexExt, vector::VectorIndexParams};
let (_tmp, uri) = create_multi_fragment_vector_dataset(2, 64, 8, false);
lance_c::runtime::block_on(async {
let mut ds = Dataset::open(&uri).await.unwrap();
let params = VectorIndexParams::ivf_flat(2, lance_linalg::distance::MetricType::L2);
let mut segments = Vec::new();
for fragment in ds.get_fragments() {
segments.push(
ds.create_index_builder(&["embedding"], lance_index::IndexType::Vector, &params)
.name("vector_idx".into())
.fragments(vec![fragment.id() as u32])
.execute_uncommitted()
.await
.unwrap(),
);
}
ds.commit_existing_index_segments("vector_idx", "embedding", segments)
.await
.unwrap();
});

fn query(ds: *mut LanceDataset, row: i32) -> CapturedScanStatistics {
unsafe {
let scanner = lance_scanner_new(ds, ptr::null(), ptr::null());
assert!(!scanner.is_null());
let query: Vec<f32> = (0..8).map(|i| row as f32 + i as f32 / 8.0).collect();
assert_eq!(
lance_scanner_nearest(
scanner,
c"embedding".as_ptr(),
query.as_ptr().cast(),
8,
LanceDataType::Float32 as i32,
1
),
0
);
assert_eq!(lance_scanner_set_nprobes(scanner, 2), 0);
let mut captured = CapturedScanStatistics::default();
assert_eq!(
lance_scanner_set_statistics_callback(
scanner,
Some(capture_scan_statistics),
(&mut captured as *mut CapturedScanStatistics).cast()
),
0
);
let batches = scan_all_rows_from_scanner(scanner);
assert_eq!(batches.iter().map(|b| b.num_rows()).sum::<usize>(), 1);
assert_eq!(
batches[0]
.column_by_name("id")
.unwrap()
.as_any()
.downcast_ref::<Int32Array>()
.unwrap()
.value(0),
row
);
assert_eq!(captured.calls, 1);
lance_scanner_close(scanner);
captured
}
}

// Build the index with a different session, so only prewarm can populate
// the cache used by the first query on each subsequently opened handle.
unsafe {
let cold = lance_dataset_open(c_str(&uri).as_ptr(), ptr::null(), 0);
assert!(!cold.is_null());
assert!(query(cold, 5).index_partitions_loaded > 0);
lance_dataset_close(cold);

let session = lance_session_new(64 * 1024 * 1024, 16 * 1024 * 1024);
assert!(!session.is_null());
let ds = lance_dataset_open_with_session(c_str(&uri).as_ptr(), ptr::null(), 0, session);
assert!(!ds.is_null());
assert_eq!(
lance_dataset_index_segment_count(ds, c"vector_idx".as_ptr()),
2
);
let version = lance_dataset_version(ds);
assert_eq!(lance_dataset_prewarm_index(ds, c"vector_idx".as_ptr()), 0);
assert_eq!(lance_dataset_version(ds), version);
let mut warmed = LanceSessionCacheStats::default();
assert_eq!(lance_session_get_cache_stats(session, &mut warmed), 0);
assert!(warmed.index_cache_entries > 0);
assert!(warmed.index_cache_size_bytes > 0);
lance_dataset_close(ds);

let reopened =
lance_dataset_open_with_session(c_str(&uri).as_ptr(), ptr::null(), version, session);
assert!(!reopened.is_null());
let mut query_stats = Vec::new();
// Query both segments and all IVF partitions, rather than a single
// vector that could pass even if prewarm skipped one segment.
for row in [5, 100] {
query_stats.push(query(reopened, row));
}
let mut queried = LanceSessionCacheStats::default();
assert_eq!(lance_session_get_cache_stats(session, &mut queried), 0);
assert!(queried.index_cache_hits > warmed.index_cache_hits);
for stats in query_stats {
assert_eq!(
stats.index_partitions_loaded, 0,
"prewarmed partitions should not be loaded again"
);
}
// Repeating the operation is safe, and the dataset owns the session
// even after its C session handle is released.
lance_session_close(session);
assert_eq!(
lance_dataset_prewarm_index(reopened, c"vector_idx".as_ptr()),
0
);
assert_eq!(lance_dataset_version(reopened), version);
lance_dataset_close(reopened);
}
}
5 changes: 5 additions & 0 deletions tests/cpp/test_c_api.c
Original file line number Diff line number Diff line change
Expand Up @@ -1216,6 +1216,11 @@ static void test_commit_index_segments(const char *uri) {
"commit must bump the dataset version exactly once");
ASSERT(lance_dataset_index_segment_count(ds, "c_distributed_idx") == 2,
"committed index must have two segments");
ASSERT(lance_dataset_prewarm_index(ds, "c_distributed_idx") == 0,
"prewarm both index segments");
ASSERT(lance_dataset_prewarm_index(ds, "c_distributed_idx") == 0,
"repeat prewarm");

uint8_t committed_uuids[32] = {0};
uint64_t committed_count = 0;
ASSERT(lance_dataset_index_segments(ds, "c_distributed_idx",
Expand Down
12 changes: 12 additions & 0 deletions tests/cpp/test_cpp_api.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -528,6 +528,18 @@ static void test_index_lifecycle(const std::string& uri) {
assert(json.find("id_idx") != std::string::npos);
printf("listed: %s... ", json.c_str());

const auto& readonly_ds = ds;
readonly_ds.prewarm_index("id_idx");
readonly_ds.prewarm_index("id_idx");
bool missing_index = false;
try {
readonly_ds.prewarm_index("missing_idx");
} catch (const lance::Error& error) {
assert(error.code == LANCE_ERR_NOT_FOUND);
missing_index = true;
}
assert(missing_index);

ds.drop_index("id_idx");
assert(ds.index_count() == 0);

Expand Down
Loading