Compare commits

...

4 Commits

Author SHA1 Message Date
Gatefixer 7cfc8d3848 Merge remote-tracking branch 'origin/main' into gatekeeper/fix-3086-1 2026-08-06 04:00:59 +00:00
Gatefixer d8c5695b7b fix: reject unsafe version cleanup windows 2026-08-06 04:00:56 +00:00
lancedb-gatefixer[bot] 7357d63e87 fix(python): guard concurrent table deletes (#3787)
<!-- lance-gatekeeper-fix:v1 agent=5c80c44c083b3b8ad0da595419d468fc
generation=1 -->

## Root cause

The legacy synchronous Python table called `delete` on a shared, mutable
`lance.Dataset`. Concurrent table operations could hold a PyO3 borrow
while delete requested an exclusive borrow, producing `RuntimeError:
Already borrowed`. The current async-backed binding fixes this by
cloning its thread-safe Rust table handle before awaiting, but that
concurrency contract had no regression coverage.

## Fix

- Document why delete must clone the Rust table handle before entering
its async future.
- Add a barrier-synchronized regression test that deletes distinct rows
through one shared table from eight Python threads.
- Verify every delete commits exactly one row, every commit gets a
distinct version, and no rows remain.

## Validation

- `cargo check --quiet --features remote --tests --examples`
- `cargo fmt --all -- --check`
- `uv run --extra tests --extra dev ruff format --check
python/tests/test_table.py`
- `uv run --extra tests --extra dev ruff check
python/tests/test_table.py`
- `uv run --extra tests --extra dev pytest
python/tests/test_table.py::test_concurrent_deletes_are_thread_safe
python/tests/test_table.py::test_delete
python/tests/test_table.py::test_delete_expr
python/tests/test_table.py::test_delete_expr_async -q` (4 passed)
- Manual stress reproduction: 100 concurrent deletes on one table
completed at versions 2–101 with zero rows remaining.

Fixes #530

Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com>
2026-08-05 15:17:04 -07:00
lancedb-gatefixer[bot] 624a75edf7 fix(python): avoid debugger deadlock during connection inspection (#3788)
## Summary

- cache the immutable read consistency interval on synchronous
connection wrappers
- keep debugger property expansion from dispatching to the background
event loop
- cover direct connections and wrappers reconstructed from native
connections

## Root cause

The debugger expands connection variables by evaluating properties after
suspending all Python threads.
`LanceDBConnection.read_consistency_interval` dispatched a coroutine to
`LanceDBBackgroundEventLoop` and synchronously waited for it, but that
loop thread was also suspended, causing a deadlock.

## Validation

- `uv run --no-sync pytest python/tests/test_db.py -q` (48 passed)
- `ruff format --check python/python/lancedb/db.py
python/python/tests/test_db.py`
- `ruff check .`
- `git diff --check`

Fixes #3773

<!-- lance-gatekeeper-fix:v1 agent=e2e612236d722d926f64245d3f682bbc
generation=1 -->

---------

Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com>
2026-08-05 15:15:49 -07:00
13 changed files with 250 additions and 29 deletions
+3
View File
@@ -587,6 +587,9 @@ Modeled after ``VACUUM`` in PostgreSQL.
you have added or modified 100,000 or more records or run more than 20 data
modification operations.
Cleanup retention must exceed the longest expected write. Retention shorter
than 10 minutes requires `deleteUnverified: true` and exclusive write access.
#### Parameters
* **options?**: `Partial`&lt;[`OptimizeOptions`](../interfaces/OptimizeOptions.md)&gt;
+6 -3
View File
@@ -16,7 +16,10 @@ cleanupOlderThan: Date;
If set then all versions older than the given date
be removed. The current version will never be removed.
The default is 7 days
The default is 7 days. The resulting retention period must be longer than
the longest expected write. Values shorter than 10 minutes require
`deleteUnverified: true` and are only safe when no other process can write
to the dataset.
#### Example
@@ -26,8 +29,8 @@ const olderThan = new Date();
olderThan.setDate(olderThan.getDate() - 1));
tbl.optimize({cleanupOlderThan: olderThan});
// Delete all versions except the current version
tbl.optimize({cleanupOlderThan: new Date()});
// With exclusive access, delete all versions except the current version
tbl.optimize({cleanupOlderThan: new Date(), deleteUnverified: true});
```
***
+9 -1
View File
@@ -2196,7 +2196,15 @@ describe("when optimizing a dataset", () => {
});
it("cleanups old versions", async () => {
const stats = await table.optimize({ cleanupOlderThan: new Date() });
await expect(
table.optimize({ cleanupOlderThan: new Date() }),
).rejects.toThrow("at least 10 minutes");
expect(await table.version()).toBe(2);
const stats = await table.optimize({
cleanupOlderThan: new Date(),
deleteUnverified: true,
});
expect(stats.prune.bytesRemoved).toBeGreaterThan(0);
expect(stats.prune.oldVersionsRemoved).toBe(3);
});
+9 -3
View File
@@ -129,15 +129,18 @@ export interface OptimizeOptions {
/**
* If set then all versions older than the given date
* be removed. The current version will never be removed.
* The default is 7 days
* The default is 7 days. The resulting retention period must be longer than
* the longest expected write. Values shorter than 10 minutes require
* `deleteUnverified: true` and are only safe when no other process can write
* to the dataset.
* @example
* // Delete all versions older than 1 day
* const olderThan = new Date();
* olderThan.setDate(olderThan.getDate() - 1));
* tbl.optimize({cleanupOlderThan: olderThan});
*
* // Delete all versions except the current version
* tbl.optimize({cleanupOlderThan: new Date()});
* // With exclusive access, delete all versions except the current version
* tbl.optimize({cleanupOlderThan: new Date(), deleteUnverified: true});
*/
cleanupOlderThan: Date;
/**
@@ -744,6 +747,9 @@ export abstract class Table {
* optimize should be run frequently. A good rule of thumb is to run optimize if
* you have added or modified 100,000 or more records or run more than 20 data
* modification operations.
*
* Cleanup retention must exceed the longest expected write. Retention shorter
* than 10 minutes requires `deleteUnverified: true` and exclusive write access.
*/
abstract optimize(options?: Partial<OptimizeOptions>): Promise<OptimizeStats>;
/** List all indices that have been created with {@link Table.createIndex} */
+2
View File
@@ -6,6 +6,7 @@ use std::collections::HashMap;
use chrono::{DateTime, Utc};
use lancedb::ipc::{ipc_file_to_batches, ipc_file_to_schema};
use lancedb::table::optimize::validate_cleanup_options;
use lancedb::table::{
AddDataMode, ColumnAlteration as LanceColumnAlteration, Duration,
FieldMetadataUpdate as LanceFieldMetadataUpdate, FtsToken as LanceDbFtsToken,
@@ -559,6 +560,7 @@ impl Table {
} else {
None
};
validate_cleanup_options(older_than, delete_unverified).default_error()?;
let compaction_stats = inner
.optimize(OptimizeAction::Compact {
+17 -3
View File
@@ -707,6 +707,9 @@ class LanceDBConnection(DBConnection):
self._namespace_client_properties = namespace_client_properties
if _inner is not None:
self._conn = _inner
# Native-derived wrappers resolve this in their async reconstruction
# path so construction never synchronously re-enters LOOP.
self._read_consistency_interval = read_consistency_interval
self._cached_namespace_client = None
return
@@ -756,11 +759,14 @@ class LanceDBConnection(DBConnection):
# storage_options. Also, this class really shouldn't be holding any state
# beyond _conn.
self._conn = AsyncConnection(LOOP.run(do_connect()))
# Keep property access synchronous so debugger introspection cannot wait on
# the background loop while that thread is suspended at a breakpoint.
self._read_consistency_interval = read_consistency_interval
self._cached_namespace_client: Optional[LanceNamespace] = None
@property
def read_consistency_interval(self) -> Optional[timedelta]:
return LOOP.run(self._conn.get_read_consistency_interval())
return self._read_consistency_interval
@property
def session(self) -> Optional[Session]:
@@ -771,8 +777,16 @@ class LanceDBConnection(DBConnection):
return self._conn.uri
@classmethod
def from_inner(cls, inner: LanceDbConnection):
return cls(None, _inner=inner)
def from_inner(
cls,
inner: LanceDbConnection,
read_consistency_interval: Optional[timedelta],
):
return cls(
None,
read_consistency_interval=read_consistency_interval,
_inner=inner,
)
def __repr__(self) -> str:
return f"{self.__class__.__name__}(uri={self._conn.uri!r})"
+1 -1
View File
@@ -226,7 +226,7 @@ class PermutationBuilder:
async def do_execute():
inner_tbl = await self._async.execute()
return LanceTable.from_inner(inner_tbl)
return await LanceTable.from_inner(inner_tbl)
return LOOP.run(do_execute())
+46 -14
View File
@@ -117,6 +117,23 @@ _MODEL_BACKED_TOKENIZER_ERRORS = (
"Failed to initialize default tokenizer",
)
_MIN_SAFE_CLEANUP_AGE = timedelta(minutes=10)
def _validate_cleanup_options(
older_than: Optional[timedelta], delete_unverified: bool
) -> None:
if (
older_than is not None
and older_than < _MIN_SAFE_CLEANUP_AGE
and not delete_unverified
):
raise ValueError(
"cleanup age must be at least 10 minutes unless delete_unverified is "
"true; short cleanup windows can remove a manifest still needed by an "
"in-progress write"
)
def _add_unique_note(exception: BaseException, note: str) -> None:
existing_notes = getattr(exception, "__notes__", ()) or ()
@@ -1767,7 +1784,9 @@ class Table(ABC):
----------
older_than: timedelta, default None
The minimum age of the version to delete. If None, then this defaults
to two weeks.
to two weeks. This must be longer than the longest expected write.
Values shorter than 10 minutes require `delete_unverified=True` and
are only safe when no other process can write to the dataset.
delete_unverified: bool, default False
Because they may be part of an in-progress transaction, files newer
than 7 days old are not deleted by default. If you are sure that
@@ -1835,9 +1854,11 @@ class Table(ABC):
Parameters
----------
cleanup_older_than: timedelta, optional default 7 days
All files belonging to versions older than this will be removed. Set
to 0 days to remove all versions except the latest. The latest version
is never removed.
All files belonging to versions older than this will be removed. The
latest version is never removed. This must be longer than the longest
expected write. Values shorter than 10 minutes require
`delete_unverified=True` and are only safe when no other process can
write to the dataset.
delete_unverified: bool, default False
Files leftover from a failed transaction may appear to be part of an
in-progress operation (e.g. appending new data) and these files will not
@@ -2182,11 +2203,15 @@ class LanceTable(Table):
return self.name
@classmethod
def from_inner(cls, tbl: LanceDBTable):
from .db import LanceDBConnection
async def from_inner(cls, tbl: LanceDBTable):
from .db import AsyncConnection, LanceDBConnection
async_tbl = AsyncTable(tbl)
conn = LanceDBConnection.from_inner(tbl.database())
inner_conn = tbl.database()
read_consistency_interval = await AsyncConnection(
inner_conn
).get_read_consistency_interval()
conn = LanceDBConnection.from_inner(inner_conn, read_consistency_interval)
return cls(
conn,
async_tbl.name,
@@ -3782,7 +3807,9 @@ class LanceTable(Table):
----------
older_than: timedelta, default None
The minimum age of the version to delete. If None, then this defaults
to two weeks.
to two weeks. This must be longer than the longest expected write.
Values shorter than 10 minutes require `delete_unverified=True` and
are only safe when no other process can write to the dataset.
delete_unverified: bool, default False
Because they may be part of an in-progress transaction, files newer
than 7 days old are not deleted by default. If you are sure that
@@ -3795,6 +3822,7 @@ class LanceTable(Table):
The stats of the cleanup operation, including how many bytes were
freed.
"""
_validate_cleanup_options(older_than, delete_unverified)
return self.to_lance().cleanup_old_versions(
older_than, delete_unverified=delete_unverified
)
@@ -3840,9 +3868,11 @@ class LanceTable(Table):
Parameters
----------
cleanup_older_than: timedelta, optional default 7 days
All files belonging to versions older than this will be removed. Set
to 0 days to remove all versions except the latest. The latest version
is never removed.
All files belonging to versions older than this will be removed. The
latest version is never removed. This must be longer than the longest
expected write. Values shorter than 10 minutes require
`delete_unverified=True` and are only safe when no other process can
write to the dataset.
delete_unverified: bool, default False
Files leftover from a failed transaction may appear to be part of an
in-progress operation (e.g. appending new data) and these files will not
@@ -6034,9 +6064,11 @@ class AsyncTable:
Parameters
----------
cleanup_older_than: timedelta, optional default 7 days
All files belonging to versions older than this will be removed. Set
to 0 days to remove all versions except the latest. The latest version
is never removed.
All files belonging to versions older than this will be removed. The
latest version is never removed. This must be longer than the longest
expected write. Values shorter than 10 minutes require
`delete_unverified=True` and are only safe when no other process can
write to the dataset.
delete_unverified: bool, default False
Files leftover from a failed transaction may appear to be part of an
in-progress operation (e.g. appending new data) and these files will not
+17
View File
@@ -77,6 +77,23 @@ def test_sync_repr_does_not_use_background_loop(tmp_path, monkeypatch):
assert repr(table) == f"LanceTable(name='test', _conn={db!r})"
def test_read_consistency_interval_does_not_use_background_loop(tmp_path, monkeypatch):
from lancedb.background_loop import LOOP
from lancedb.db import LanceDBConnection
consistency_interval = timedelta(seconds=5)
db = lancedb.connect(tmp_path, read_consistency_interval=consistency_interval)
db_from_inner = LanceDBConnection.from_inner(db._inner, consistency_interval)
def fail_run(*args, **kwargs):
raise AssertionError("properties should not use the Python background loop")
monkeypatch.setattr(LOOP, "run", fail_run)
assert db.read_consistency_interval == consistency_interval
assert db_from_inner.read_consistency_interval == consistency_interval
def test_ingest_pd(tmp_path):
db = lancedb.connect(tmp_path)
+20
View File
@@ -6,6 +6,7 @@ import math
import pytest
from lancedb import DBConnection, Table, connect
from lancedb.background_loop import LOOP
from lancedb.permutation import Permutation, Permutations, permutation_builder
@@ -31,6 +32,25 @@ def test_split_random_ratios(mem_db):
assert 65 <= split_1_count <= 75 # ~70% ± tolerance
def test_execute_does_not_reenter_background_loop(tmp_path, monkeypatch):
import threading
db = connect(tmp_path)
tbl = db.create_table("test_table", pa.table({"x": range(10)}))
original_run = LOOP.run
def fail_on_reentry(future):
assert threading.current_thread() is not LOOP.thread
return original_run(future)
monkeypatch.setattr(LOOP, "run", fail_on_reentry)
permutation_tbl = permutation_builder(tbl).execute()
assert permutation_tbl.count_rows() == 10
assert permutation_tbl._conn.read_consistency_interval is None
def test_split_random_counts(mem_db):
"""Test random splitting with absolute counts."""
tbl = mem_db.create_table(
+33 -1
View File
@@ -6,6 +6,7 @@ import os
import sys
import threading
import warnings
from concurrent.futures import ThreadPoolExecutor
from datetime import date, datetime, timedelta
from time import sleep
from typing import List
@@ -2124,6 +2125,27 @@ def test_delete(mem_db: DBConnection):
assert table.to_arrow()["id"].to_pylist() == [1]
def test_concurrent_deletes_are_thread_safe(mem_db: DBConnection):
num_workers = 8
table = mem_db.create_table(
"my_table", data=[{"id": row_id} for row_id in range(num_workers)]
)
barrier = threading.Barrier(num_workers)
def delete(row_id: int):
barrier.wait()
return table.delete(f"id = {row_id}")
with ThreadPoolExecutor(max_workers=num_workers) as pool:
results = list(pool.map(delete, range(num_workers)))
assert all(result.num_deleted_rows == 1 for result in results)
assert sorted(result.version for result in results) == list(
range(2, num_workers + 2)
)
assert table.count_rows() == 0
def test_delete_expr(mem_db: DBConnection):
table = mem_db.create_table(
"my_table",
@@ -2958,6 +2980,9 @@ def test_compact_cleanup(tmp_db: DBConnection):
stats = table.cleanup_old_versions()
assert stats.bytes_removed == 0
with pytest.raises(ValueError, match="at least 10 minutes"):
table.cleanup_old_versions(older_than=timedelta(0))
stats = table.cleanup_old_versions(older_than=timedelta(0), delete_unverified=True)
assert stats.bytes_removed > 0
assert table.version == 4
@@ -3374,7 +3399,14 @@ async def test_optimize(mem_db_async: AsyncConnection):
assert stats.prune.bytes_removed == 0
assert stats.prune.old_versions_removed == 0
stats = await table.optimize(cleanup_older_than=timedelta(seconds=0))
version_before_rejected_cleanup = await table.version()
with pytest.raises(ValueError, match="at least 10 minutes"):
await table.optimize(cleanup_older_than=timedelta(seconds=0))
assert await table.version() == version_before_rejected_cleanup
stats = await table.optimize(
cleanup_older_than=timedelta(seconds=0), delete_unverified=True
)
assert stats.prune.bytes_removed > 0
assert stats.prune.old_versions_removed == 3
+5
View File
@@ -19,6 +19,7 @@ use arrow::{
};
use lancedb::blob::{BlobFile, BlobRangeRequest};
use lancedb::index::scalar::FtsIndexBuilder;
use lancedb::table::optimize::validate_cleanup_options;
use lancedb::table::{
AddDataMode, ColumnAlteration, Duration, FieldMetadataUpdate, FtsToken as LanceDbFtsToken,
NewColumnTransform, OptimizeAction, OptimizeOptions, Ref, Table as LanceDbTable,
@@ -745,6 +746,9 @@ impl Table {
#[allow(private_interfaces)]
pub fn delete(self_: PyRef<'_, Self>, condition: PredicateArg) -> PyResult<Bound<'_, PyAny>> {
// Do not hold the Python borrow across the await. The cloned Rust table
// handle is thread-safe and allows deletes on the same Python table to
// run concurrently without PyO3 reporting "Already borrowed".
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
let result = match &condition {
@@ -1214,6 +1218,7 @@ impl Table {
} else {
None
};
validate_cleanup_options(older_than, delete_unverified).infer_error()?;
future_into_py(self_.py(), async move {
let compaction_stats = inner
.optimize(OptimizeAction::Compact {
+82 -3
View File
@@ -18,7 +18,33 @@ pub use chrono::Duration;
pub use lance::dataset::optimize::CompactionOptions;
use super::NativeTable;
use crate::error::Result;
use crate::error::{Error, Result};
const MIN_SAFE_CLEANUP_AGE_MINUTES: i64 = 10;
/// Validate version-cleanup options before starting an optimization operation.
///
/// This is public for the language bindings, which run compaction and cleanup as
/// separate operations and must reject unsafe cleanup options before compaction
/// changes the table.
#[doc(hidden)]
pub fn validate_cleanup_options(
older_than: Option<Duration>,
delete_unverified: Option<bool>,
) -> Result<()> {
let minimum_age =
Duration::try_minutes(MIN_SAFE_CLEANUP_AGE_MINUTES).expect("minimum cleanup age is valid");
if older_than.is_some_and(|age| age < minimum_age) && delete_unverified != Some(true) {
return Err(Error::InvalidInput {
message: format!(
"cleanup age must be at least {MIN_SAFE_CLEANUP_AGE_MINUTES} minutes unless \
delete_unverified is true; short cleanup windows can remove a manifest still \
needed by an in-progress write"
),
});
}
Ok(())
}
/// Optimize the dataset.
///
@@ -60,7 +86,9 @@ pub enum OptimizeAction {
///
/// Once a version is pruned it can no longer be checked out.
Prune {
/// The duration of time to keep versions of the dataset.
/// The duration of time to keep versions of the dataset. This should be longer than the
/// longest expected write. Values shorter than 10 minutes require `delete_unverified` to
/// be true and are only safe when no other process can write to the dataset.
older_than: Option<Duration>,
/// Because they may be part of an in-progress transaction, files newer than 7 days old are not deleted by default.
/// If you are sure that there are no in-progress transactions, then you can set this to True to delete all files older than `older_than`.
@@ -164,6 +192,15 @@ pub(crate) async fn execute_optimize(
table: &NativeTable,
action: OptimizeAction,
) -> Result<OptimizeStats> {
if let OptimizeAction::Prune {
older_than,
delete_unverified,
..
} = &action
{
validate_cleanup_options(*older_than, *delete_unverified)?;
}
let mut stats = OptimizeStats {
compaction: None,
prune: None,
@@ -222,7 +259,7 @@ mod tests {
use crate::connect;
use crate::index::{Index, scalar::BTreeIndexBuilder};
use crate::query::ExecutableQuery;
use crate::table::{CompactionOptions, OptimizeAction, OptimizeStats};
use crate::table::{CompactionOptions, Duration, OptimizeAction, OptimizeStats};
use futures::TryStreamExt;
#[tokio::test]
@@ -380,6 +417,48 @@ mod tests {
assert_eq!(all_values, expected);
}
#[tokio::test]
async fn test_optimize_rejects_unsafe_cleanup_age() {
let conn = connect("memory://").execute().await.unwrap();
let batch = RecordBatch::try_new(
Arc::new(Schema::new(vec![Field::new("i", DataType::Int32, false)])),
vec![Arc::new(Int32Array::from_iter_values(0..10))],
)
.unwrap();
let table = conn
.create_table("test_unsafe_prune", batch.clone())
.execute()
.await
.unwrap();
table.add(batch).execute().await.unwrap();
let versions_before = table
.list_versions()
.await
.unwrap()
.into_iter()
.map(|version| version.version)
.collect::<Vec<_>>();
let err = table
.optimize(OptimizeAction::Prune {
older_than: Some(Duration::zero()),
delete_unverified: Some(false),
error_if_tagged_old_versions: None,
})
.await
.unwrap_err();
assert!(err.to_string().contains("at least 10 minutes"));
let versions_after = table
.list_versions()
.await
.unwrap()
.into_iter()
.map(|version| version.version)
.collect::<Vec<_>>();
assert_eq!(versions_after, versions_before);
}
#[tokio::test]
async fn test_optimize_index() {
let conn = connect("memory://").execute().await.unwrap();