Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
320791d
feat: mask data overlay files in scalar and vector index queries
wjones127 Jun 23, 2026
59791fa
perf(index): avoid index-segment clone on the no-overlay fast path
wjones127 Jul 1, 2026
ba9999c
refactor(index): clarify overlay masking code comments
wjones127 Jul 1, 2026
e49e379
test(index): add compound-predicate overlay masking test + FTS gap note
wjones127 Jul 1, 2026
95967fb
test(bench): scale overlay benchmark to 1M rows on local disk
wjones127 Jul 1, 2026
f80a69e
fix(test): use dataset.base to resolve overlay file paths for file://…
wjones127 Jul 1, 2026
1d0aef5
perf(index): row-level overlay masking for vector ANN
wjones127 Jul 1, 2026
0271f39
fix(fts): mask stale FTS index segments when data overlays present
wjones127 Jul 1, 2026
ee04109
perf(index): row-level overlay masking for BTree scalar index
wjones127 Jul 1, 2026
91591fb
fix(index): row-level V2 scalar masking, phrase-query overlay masking…
wjones127 Jul 14, 2026
703060c
test(scanner): move overlay index-masking tests out of scanner.rs
wjones127 Jul 14, 2026
eba4c5c
fix(overlay): correct index masking under stable row ids + review fixes
wjones127 Jul 14, 2026
fc2d3ed
test(overlay): fast_search regression tests + self-review cleanups
wjones127 Jul 14, 2026
a57c860
build: raise recursion_limit for overlay-deepened scan futures
wjones127 Jul 15, 2026
ce6fc24
refactor(overlay): shorten diff — dedup test helpers, tighten comments
wjones127 Jul 15, 2026
89f202f
fix(overlay): box deep futures instead of raising recursion_limit; te…
wjones127 Jul 15, 2026
6499de0
fix: box create_plan futures at inline call sites to avoid recursion-…
wjones127 Jul 16, 2026
c8835ff
refactor(scanner): consolidate future boxing and address review nits
wjones127 Jul 21, 2026
9c17e73
Update rust/lance/src/dataset/rowids.rs
wjones127 Jul 22, 2026
5442369
fix(overlay): use Error::internal for row-id translation invariants
wjones127 Jul 22, 2026
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
7 changes: 6 additions & 1 deletion rust/lance/src/dataset/index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,12 @@ impl IndexRemapper for DatasetIndexRemapper {
.any(|frag_idx| affected_frag_ids.contains(&(frag_idx as u64))),
};
if needs_remapped {
let remap_result = self.remap_index(index, &mapping).await?;
// Box the remap future at the call site: inlining `remap_index` into this
// loop's async layout otherwise exceeds rustc's depth limit. It has to be
// boxed here, not inside `remap_index` — boxing internally turns the
// future's `Send` check into a `Box<Future>: Send` trait obligation that
// overflows the solver through the cache types (E0275 downstream).
let remap_result = Box::pin(self.remap_index(index, &mapping)).await?;
match remap_result {
RemapResult::Drop => continue,
RemapResult::Keep(id) => {
Expand Down
10 changes: 7 additions & 3 deletions rust/lance/src/dataset/mem_wal/scanner/point_lookup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -654,7 +654,11 @@ impl LsmPointLookupPlanner {
scanner.with_row_address();
}
scanner.filter_expr(filter.clone());
scanner.create_plan().await?
// Box at the call site: `create_plan`'s inlined async layout exceeds
// rustc's depth limit up this point-lookup chain, and boxing inside
// `create_plan` instead triggers a `Box<Future>: Send` solver overflow
// (E0275 downstream). Same for the other arms below.
Box::pin(scanner.create_plan()).await?
}
LsmDataSource::FlushedMemTable { path, .. } => {
let dataset = open_flushed_dataset(
Expand All @@ -672,7 +676,7 @@ impl LsmPointLookupPlanner {
let cols = cols_with_tombstone(&cols, dataset.schema().field(TOMBSTONE).is_some());
scanner.project(&cols.iter().map(|s| s.as_str()).collect::<Vec<_>>())?;
scanner.filter_expr(filter.clone());
scanner.create_plan().await?
Box::pin(scanner.create_plan()).await?
}
LsmDataSource::ActiveMemTable {
batch_store,
Expand All @@ -695,7 +699,7 @@ impl LsmPointLookupPlanner {
// over insert-ordered scan would return the *oldest* of
// multiple rows sharing the target primary key.
scanner.with_row_id();
let raw = scanner.create_plan().await?;
let raw = Box::pin(scanner.create_plan()).await?;
// The filter already restricts to the exact PK value, so the
// scan yields that key's insert history. Within the active
// memtable larger `_rowid` = newer insert, so sorting `_rowid`
Expand Down
292 changes: 291 additions & 1 deletion rust/lance/src/dataset/overlay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,11 +44,129 @@ use lance_core::datatypes::{Field, Schema};
use lance_core::{Error, Result};
use roaring::RoaringBitmap;

use lance_table::format::DataFile;
use lance_table::format::overlay::DataOverlayFile;
use lance_table::format::{DataFile, Fragment, IndexMetadata};
use lance_table::utils::stream::ReadBatchFut;

use crate::dataset::fragment::{FileFragment, FragReadConfig, GenericFileReader};

/// The physical offsets within a fragment whose value for an indexed field may be
/// stale relative to an index built at `index_version`, and so must be excluded
/// from that index's results and re-evaluated against current values on the flat
/// path.
///
/// The set is the union, over every overlay whose `committed_version` is newer
/// than `index_version`, of that overlay's coverage **restricted to the indexed
/// fields**. The restriction makes exclusion field-aware: an overlay that touches
/// only non-indexed fields contributes nothing. An overlay whose
/// `committed_version <= index_version` is already incorporated by the index and
/// is ignored.
pub fn overlay_exclusion_offsets(
overlays: &[DataOverlayFile],
indexed_field_ids: &[i32],
index_version: u64,
) -> Result<RoaringBitmap> {
let mut excluded = RoaringBitmap::new();
for overlay in overlays {
if overlay.committed_version <= index_version {
continue;
}
for (field_pos, field_id) in overlay.data_file.fields.iter().enumerate() {
if indexed_field_ids.contains(field_id) {
excluded |= &*overlay.coverage_for_field(field_pos)?;
}
}
}
Ok(excluded)
}

// Stale row offsets contributed by one fragment's overlays for a given index version.
// Applies a cheap version gate first: if every overlay predates the segment it is already
// incorporated by the index, so there is nothing stale and the field/bitmap work is skipped.
fn stale_offsets_for_fragment(
fragment: &Fragment,
fields: &[i32],
index_version: u64,
) -> Result<RoaringBitmap> {
if fragment
.overlays
.iter()
.all(|o| o.committed_version <= index_version)
{
return Ok(RoaringBitmap::new());
}
overlay_exclusion_offsets(&fragment.overlays, fields, index_version)
}

// A missing `fragment_bitmap` means the index predates fragment-bitmap tracking; treat it as
// covering every fragment (matching `DatasetPreFilter::new`) so overlay-stale rows can't slip
// through unmasked. Only skip fragments explicitly absent from a present bitmap.
fn covers_fragment(coverage: Option<&RoaringBitmap>, frag_id: u32) -> bool {
coverage.is_none_or(|c| c.contains(frag_id))
}

/// Index by fragment id the fragments that carry at least one overlay. Overlays are rare, so
/// this is empty on the common path, letting callers skip index loading entirely; when non-empty
/// it bounds the stale-collection loops to `O(overlaid fragments)`.
pub fn overlaid_fragments(fragments: &[Fragment]) -> HashMap<u32, &Fragment> {
fragments
.iter()
.filter(|f| !f.overlays.is_empty())
.map(|f| (f.id as u32, f))
.collect()
}

/// Insert into `stale` the ids of fragments covered by `segment` whose index entries may be
/// stale because an overlay committed after the segment was built touches a field the segment
/// indexes. Field-aware and version-gated via [`overlay_exclusion_offsets`].
///
/// `overlaid_frags` holds only the fragments that actually carry overlays (rare), so the loop is
/// `O(overlaid_frags)` rather than `O(fragments the segment covers)`.
pub fn collect_overlay_stale_frags(
segment: &IndexMetadata,
overlaid_frags: &HashMap<u32, &Fragment>,
stale: &mut RoaringBitmap,
) -> Result<()> {
let coverage = segment.fragment_bitmap.as_ref();
for (&frag_id, fragment) in overlaid_frags {
if stale.contains(frag_id) || !covers_fragment(coverage, frag_id) {
continue;
}
if !stale_offsets_for_fragment(fragment, &segment.fields, segment.dataset_version)?
.is_empty()
{
stale.insert(frag_id);
}
}
Ok(())
}

/// Like [`collect_overlay_stale_frags`] but with row-level granularity: instead of marking the
/// whole fragment stale, it computes exactly which row offsets within each covered fragment are
/// stale and accumulates them into `stale` (fragment_id → stale row offsets).
///
/// Used by the scalar and vector paths to block only the affected rows from index results and
/// re-evaluate only those rows on the flat path, keeping overhead proportional to the number of
/// overlaid rows rather than the whole fragment size.
pub fn collect_overlay_stale_rows_for_segment(
segment: &IndexMetadata,
overlaid_frags: &HashMap<u32, &Fragment>,
stale: &mut HashMap<u32, RoaringBitmap>,
) -> Result<()> {
let coverage = segment.fragment_bitmap.as_ref();
for (&frag_id, fragment) in overlaid_frags {
if !covers_fragment(coverage, frag_id) {
continue;
}
let excluded =
stale_offsets_for_fragment(fragment, &segment.fields, segment.dataset_version)?;
if !excluded.is_empty() {
*stale.entry(frag_id).or_default() |= &excluded;
}
}
Ok(())
}

/// The plan for merging one field's overlays into one batch: which source (base or
/// a particular overlay) supplies each output row, and which overlay values must be
/// fetched to do it.
Expand Down Expand Up @@ -759,6 +877,7 @@ async fn fetch_overlay_values(
mod tests {
use super::*;
use arrow_array::{Int32Array, StringArray, UInt32Array};
use lance_table::format::overlay::OverlayCoverage;
use std::sync::Arc;

fn i32_array(values: impl IntoIterator<Item = Option<i32>>) -> ArrayRef {
Expand Down Expand Up @@ -1061,4 +1180,175 @@ mod tests {
assert!(spliced.is_null(1));
assert!(!spliced.is_null(2));
}

/// A dense overlay covering `offsets` for `field_ids`, committed at `version`.
fn dense_overlay(
field_ids: Vec<i32>,
offsets: impl IntoIterator<Item = u32>,
version: u64,
) -> lance_table::format::overlay::DataOverlayFile {
DataOverlayFile {
data_file: DataFile::new_legacy_from_fields("o.lance", field_ids, None),
coverage: OverlayCoverage::dense(bitmap(offsets)),
committed_version: version,
}
}

#[test]
fn test_exclusion_offsets_version_gate() {
// index built at version 5; only overlays committed > 5 are excluded.
let overlays = vec![
dense_overlay(vec![3], [0, 1], 4),
dense_overlay(vec![3], [2, 7], 6),
];
let excluded = overlay_exclusion_offsets(&overlays, &[3], 5).unwrap();
assert_eq!(excluded, bitmap([2, 7]));
// An overlay exactly at the index version is already incorporated.
let overlays = vec![dense_overlay(vec![3], [9], 5)];
assert!(
overlay_exclusion_offsets(&overlays, &[3], 5)
.unwrap()
.is_empty()
);
}

#[test]
fn test_exclusion_offsets_is_field_aware() {
// An overlay touching only an unrelated field excludes nothing.
let overlays = vec![dense_overlay(vec![2], [0, 1, 2], 9)];
assert!(
overlay_exclusion_offsets(&overlays, &[3], 1)
.unwrap()
.is_empty()
);
// The union spans only the indexed fields the overlay actually carries.
let overlays = vec![dense_overlay(vec![2, 3], [4], 9)];
assert_eq!(
overlay_exclusion_offsets(&overlays, &[3], 1).unwrap(),
bitmap([4])
);
}

#[test]
fn test_exclusion_offsets_sparse_per_field() {
// Sparse overlay: field 2 covers {2,3}, field 4 covers {1}.
let overlay = DataOverlayFile {
data_file: DataFile::new_legacy_from_fields("o.lance", vec![2, 4], None),
coverage: OverlayCoverage::sparse(vec![bitmap([2, 3]), bitmap([1])]),
committed_version: 9,
};
let overlays = vec![overlay];
// Only the bitmap for the indexed field (4) contributes.
assert_eq!(
overlay_exclusion_offsets(&overlays, &[4], 1).unwrap(),
bitmap([1])
);
assert_eq!(
overlay_exclusion_offsets(&overlays, &[2], 1).unwrap(),
bitmap([2, 3])
);
}

#[test]
fn test_exclusion_offsets_unions_multiple_overlays() {
let overlays = vec![
dense_overlay(vec![3], [1], 6),
dense_overlay(vec![3], [4, 5], 7),
];
assert_eq!(
overlay_exclusion_offsets(&overlays, &[3], 1).unwrap(),
bitmap([1, 4, 5])
);
}
Comment thread
wjones127 marked this conversation as resolved.

/// An index segment covering `fields`, built at `dataset_version`, with the given
/// fragment coverage (`None` = legacy index predating fragment-bitmap tracking).
fn segment(
fields: Vec<i32>,
dataset_version: u64,
fragment_bitmap: Option<RoaringBitmap>,
) -> IndexMetadata {
IndexMetadata {
uuid: uuid::Uuid::new_v4(),
name: "idx".into(),
fields,
dataset_version,
fragment_bitmap,
index_details: None,
index_version: 0,
created_at: None,
base_id: None,
files: None,
}
}

fn fragment_with_overlay(id: u64, overlay: DataOverlayFile) -> Fragment {
let mut fragment = Fragment::new(id);
fragment.overlays.push(overlay);
fragment
}

#[test]
fn test_collect_frags_missing_bitmap_covers_all() {
// A segment with no fragment_bitmap (legacy index predating bitmap tracking) must treat
// every overlaid fragment as covered so stale rows can't leak past the index unmasked.
let fragment = fragment_with_overlay(3, dense_overlay(vec![3], [1, 2], 9));
let overlaid: HashMap<u32, &Fragment> = HashMap::from([(3u32, &fragment)]);

let mut stale = RoaringBitmap::new();
collect_overlay_stale_frags(&segment(vec![3], 1, None), &overlaid, &mut stale).unwrap();
assert_eq!(stale, bitmap([3]), "missing bitmap must cover fragment 3");

// A present bitmap that excludes fragment 3 leaves it untouched.
let mut stale = RoaringBitmap::new();
collect_overlay_stale_frags(
&segment(vec![3], 1, Some(bitmap([0]))),
&overlaid,
&mut stale,
)
.unwrap();
assert!(
stale.is_empty(),
"fragment absent from bitmap is not covered"
);

// A present bitmap that includes fragment 3 marks it stale.
let mut stale = RoaringBitmap::new();
collect_overlay_stale_frags(
&segment(vec![3], 1, Some(bitmap([3]))),
&overlaid,
&mut stale,
)
.unwrap();
assert_eq!(stale, bitmap([3]));
}

#[test]
fn test_collect_rows_missing_bitmap_covers_all() {
// Same covers-all guarantee at row-level granularity.
let fragment = fragment_with_overlay(3, dense_overlay(vec![3], [1, 2], 9));
let overlaid: HashMap<u32, &Fragment> = HashMap::from([(3u32, &fragment)]);

let mut stale = HashMap::new();
collect_overlay_stale_rows_for_segment(&segment(vec![3], 1, None), &overlaid, &mut stale)
.unwrap();
assert_eq!(
stale.get(&3),
Some(&bitmap([1, 2])),
"missing bitmap must cover fragment 3"
);

// A present bitmap that excludes fragment 3 yields no stale rows.
let mut stale = HashMap::new();
collect_overlay_stale_rows_for_segment(
&segment(vec![3], 1, Some(bitmap([0]))),
&overlaid,
&mut stale,
)
.unwrap();
assert!(
stale.is_empty(),
"fragment absent from bitmap contributes no rows"
);
}
}
Loading
Loading