refactor(lsm): remove the index-catchup activation surface

Follows the lance change that removes the feature bit. With one set of
semantics there is no mode to switch into, so `require_mem_wal_index_catchup`
goes from the trait, `NativeTable`, and the LSM merge module.

`exclusion_watermarks` loses its `catchup_required` argument and keeps the
conservative branch: an index with no `index_catchup` entry is not known to
hold the compacted rows, so its generations stay readable from their
SSTables. That is what makes an existing table safe to read the moment the
new binary starts -- nothing is excluded until an index records that it
covers it.

Two tests changed because they encoded the branch that is gone, and they
had conflated two different things: an index that is *caught up* and an
index that is *untracked* both fell back to the compaction watermark.
Untracked now retains everything, so the tests assert that and cover the
genuinely-caught-up case separately.

484 lancedb lib tests pass.
This commit is contained in:
XYZhan
2026-08-20 18:29:56 -04:00
parent 41e0161067
commit 6c4269bbb0
3 changed files with 32 additions and 87 deletions
-25
View File
@@ -662,15 +662,6 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
message: "set_lsm_write_spec is not supported on this table type".into(),
})
}
/// Switch this table to required index catch-up, one way.
///
/// The default implementation returns `NotSupported`. Implementations
/// that support the MemWAL LSM write path must override this.
async fn require_mem_wal_index_catchup(&self) -> Result<()> {
Err(Error::NotSupported {
message: "require_mem_wal_index_catchup is not supported on this table type".into(),
})
}
/// Remove the [`LsmWriteSpec`] from this table.
///
/// This is a no-op if no spec is currently set.
@@ -1706,18 +1697,6 @@ impl Table {
self.inner.set_lsm_write_spec(spec).await
}
/// Switch this table to required index catch-up, one way.
///
/// Separate from [`Self::set_lsm_write_spec`] on purpose: a table carrying
/// the bit retains its SSTables until an index records that it holds the
/// compacted rows, so turn it on only once something can repair coverage.
///
/// Errors if no spec is set, or if the table already records SSTable
/// compaction progress from before this protocol.
pub async fn require_mem_wal_index_catchup(&self) -> Result<()> {
self.inner.require_mem_wal_index_catchup().await
}
/// Remove the [`LsmWriteSpec`] from this table, reverting to the standard
/// `merge_insert` write path.
///
@@ -3168,10 +3147,6 @@ impl BaseTable for NativeTable {
merge::lsm::set_lsm_write_spec(self, spec).await
}
async fn require_mem_wal_index_catchup(&self) -> Result<()> {
merge::lsm::require_mem_wal_index_catchup(self).await
}
async fn unset_lsm_write_spec(&self) -> Result<()> {
merge::lsm::unset_lsm_write_spec(self).await
}
-30
View File
@@ -118,36 +118,6 @@ pub(crate) async fn set_lsm_write_spec(table: &NativeTable, spec: LsmWriteSpec)
Ok(())
}
// =============================================================================
// require_mem_wal_index_catchup
// =============================================================================
/// Switch this table to required index catch-up, one way.
///
/// Deliberately **not** part of installing the write spec. Until something can
/// actually repair coverage, a table carrying the bit reports every index as
/// not known to hold the compacted rows, so its SSTables are retained
/// indefinitely -- and the WAL pod trims on the legacy rule meanwhile, leaving
/// readers pointed at files that are gone. Turn this on only once remote
/// maintenance owns the merge and the repair for the table.
///
/// Lance refuses the activation if the table already records SSTable
/// compaction progress: those numbers predate this protocol and cannot be
/// validated, so such a table must be drained rather than activated.
#[allow(clippy::redundant_pub_crate)]
pub(crate) async fn require_mem_wal_index_catchup(table: &NativeTable) -> Result<()> {
table.dataset.ensure_mutable()?;
let mut dataset = (*table.dataset.get().await?).clone();
if dataset.mem_wal_index_details().await?.is_none() {
return Err(Error::InvalidInput {
message: "require_mem_wal_index_catchup: no LSM write spec is set on this table".into(),
});
}
dataset.require_mem_wal_index_catchup().await?;
table.dataset.update(dataset);
Ok(())
}
// =============================================================================
// unset_lsm_write_spec
// =============================================================================
+32 -32
View File
@@ -248,7 +248,6 @@ fn pk_columns(dataset: &Dataset) -> Result<Vec<String>> {
fn exclusion_watermarks(
details: &MemWalIndexDetails,
index_names: &[String],
catchup_required: bool,
) -> HashMap<Uuid, u64> {
let mut exclude: HashMap<Uuid, u64> = HashMap::new();
for entry in &details.compacted_sstables {
@@ -261,13 +260,10 @@ fn exclusion_watermarks(
.and_then(|icp| icp.caught_up_generation_for_shard(&entry.shard_id))
{
Some(caught_up) => watermark = watermark.min(caught_up),
// No entry. On a table that requires catch-up this means the
// index is *not* known to hold these rows, and the base arm is
// index-only -- so every generation stays readable from its
// SSTable. Without the bit the field is not maintained at all,
// and absence carries no information.
None if catchup_required => watermark = 0,
None => {}
// No entry means the index is *not* known to hold these rows,
// and the base arm is index-only -- so every generation stays
// readable from its SSTable.
None => watermark = 0,
}
}
exclude.entry(entry.shard_id).or_insert(watermark);
@@ -288,8 +284,7 @@ async fn build_read_context(
details: &MemWalIndexDetails,
index_names: &[String],
) -> Result<(Vec<ShardSnapshot>, HashMap<Uuid, InMemoryMemTables>)> {
let catchup_required = dataset.requires_mem_wal_index_catchup();
let exclude = exclusion_watermarks(details, index_names, catchup_required);
let exclude = exclusion_watermarks(details, index_names);
let shard_ids = dataset.list_mem_wal_latest_shard_ids().await?;
// Use the dataset's own object store (not `ObjectStore::from_uri`, which
@@ -776,23 +771,34 @@ mod tests {
};
// Plain scan: drop every compacted generation (through 5).
assert_eq!(
exclusion_watermarks(&details, &[], false).get(&shard),
Some(&5)
);
assert_eq!(exclusion_watermarks(&details, &[]).get(&shard), Some(&5));
// FTS arm with a lagging index: exclusion is capped at the index catch-up
// (2), so SSTable generations 3..=5 are retained until the index covers
// them — otherwise those documents would silently vanish from FTS results.
assert_eq!(
exclusion_watermarks(&details, &["fts_idx".to_string()], false).get(&shard),
exclusion_watermarks(&details, &["fts_idx".to_string()]).get(&shard),
Some(&2)
);
// A caught-up index — or one untracked in index_catchup — falls back to the
// compaction watermark.
// An index with no entry has not recorded that it holds these rows, so
// nothing is excluded. This is the case a table written before catch-up
// was maintained lands in, and it errs toward reading the SSTables.
assert_eq!(
exclusion_watermarks(&details, &["caught_up_idx".to_string()], false).get(&shard),
exclusion_watermarks(&details, &["untracked_idx".to_string()]).get(&shard),
Some(&0)
);
// An index recorded as covering the compaction watermark excludes up to it.
let caught_up = MemWalIndexDetails {
index_catchup: vec![IndexCatchupProgress::new(
"caught_up_idx".to_string(),
vec![CompactedSsTable::new(shard, 5)],
)],
..details.clone()
};
assert_eq!(
exclusion_watermarks(&caught_up, &["caught_up_idx".to_string()]).get(&shard),
Some(&5)
);
}
@@ -821,31 +827,28 @@ mod tests {
// Each index alone stops at its own catch-up.
assert_eq!(
exclusion_watermarks(&details, &["vec_idx".to_string()], false).get(&shard),
exclusion_watermarks(&details, &["vec_idx".to_string()]).get(&shard),
Some(&7)
);
assert_eq!(
exclusion_watermarks(&details, &["fts_idx".to_string()], false).get(&shard),
exclusion_watermarks(&details, &["fts_idx".to_string()]).get(&shard),
Some(&4)
);
// Used together, the lower one governs regardless of order.
let both = ["vec_idx".to_string(), "fts_idx".to_string()];
assert_eq!(
exclusion_watermarks(&details, &both, false).get(&shard),
Some(&4)
);
assert_eq!(exclusion_watermarks(&details, &both).get(&shard), Some(&4));
let reversed = ["fts_idx".to_string(), "vec_idx".to_string()];
assert_eq!(
exclusion_watermarks(&details, &reversed, false).get(&shard),
exclusion_watermarks(&details, &reversed).get(&shard),
Some(&4)
);
}
/// An index with no catch-up entry contributes no cap today, so a lagging
/// sibling must still govern rather than being widened by the untracked one.
/// An index with no catch-up entry is not known to hold anything, so it
/// governs over a lagging sibling rather than the other way round.
#[test]
fn an_untracked_index_does_not_widen_a_lagging_sibling() {
fn an_untracked_index_retains_everything() {
let shard = Uuid::from_u128(1);
let details = MemWalIndexDetails {
compacted_sstables: vec![CompactedSsTable::new(shard, 9)],
@@ -858,10 +861,7 @@ mod tests {
};
let both = ["fts_idx".to_string(), "untracked_idx".to_string()];
assert_eq!(
exclusion_watermarks(&details, &both, false).get(&shard),
Some(&4)
);
assert_eq!(exclusion_watermarks(&details, &both).get(&shard), Some(&0));
}
#[test]