Compare commits

..

3 Commits

Author SHA1 Message Date
Gatefixer 72ac16ba76 fix: support remote tables in storage root 2026-08-06 06:17:33 +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
14 changed files with 143 additions and 259 deletions
+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())
+7 -3
View File
@@ -2182,11 +2182,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,
+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(
+22
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",
+3
View File
@@ -745,6 +745,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 {
+31 -5
View File
@@ -420,11 +420,6 @@ impl Connection {
///
/// * `name` - The name of the table
/// * `initial_data` - The initial data to write to the table
///
/// Floating-point `List` columns named `vec`, or with `vector` or `embedding`
/// in their name, are inferred as vector columns when the first batch has a
/// uniform, non-zero list length. The inferred dimension is validated for all
/// subsequent batches and stored as a `FixedSizeList`.
pub fn create_table<T: Scannable + 'static>(
&self,
name: impl Into<String>,
@@ -661,6 +656,7 @@ pub struct ConnectRequest {
/// - `/path/to/database` - local database on file system.
/// - `s3://bucket/path/to/database` or `gs://bucket/path/to/database` - database on cloud object store
/// - `db://dbname` - LanceDB Cloud
/// - `db://` with a host override - remote tables in the storage root
pub uri: String,
#[cfg(feature = "remote")]
@@ -773,6 +769,8 @@ impl ConnectBuilder {
///
/// This option is only used when connecting to LanceDB Cloud (db:// URIs)
/// and will be ignored for other URIs.
/// Use the URI `db://` together with a host override to connect to remote
/// tables stored directly in the storage root.
///
/// # Arguments
///
@@ -1359,6 +1357,34 @@ mod tests {
}
}
#[cfg(feature = "remote")]
#[tokio::test]
async fn test_connect_remote_storage_root() {
let conn = ConnectBuilder::new("db://")
.region("us-east-1")
.api_key("my-api-key")
.host_override("https://example.com")
.execute()
.await
.unwrap();
let (impl_name, properties) = conn.namespace_client_config().await.unwrap();
assert_eq!(impl_name, "rest");
assert_eq!(properties["uri"], "https://example.com");
assert_eq!(properties["header.x-lancedb-database"], "");
let result = ConnectBuilder::new("db://")
.region("us-east-1")
.api_key("my-api-key")
.execute()
.await;
assert!(matches!(
result,
Err(Error::InvalidInput { message })
if message.contains("A host override is required")
));
}
#[cfg(feature = "remote")]
#[tokio::test]
async fn test_connect_rejects_header_provider_with_oauth_config() {
+1 -3
View File
@@ -8,7 +8,7 @@ use lance_io::object_store::StorageOptionsProvider;
use crate::{
Error, Result, Table,
connection::{merge_storage_options, set_storage_options_provider},
data::scannable::{Scannable, WithEmbeddingsScannable, maybe_infer_vector_schema},
data::scannable::{Scannable, WithEmbeddingsScannable},
database::{CreateTableMode, CreateTableRequest, Database},
embeddings::{EmbeddingDefinition, EmbeddingFunction, EmbeddingRegistry},
table::WriteOptions,
@@ -147,8 +147,6 @@ impl CreateTableBuilder {
let embedding_registry = self.embedding_registry.clone();
let parent = self.parent.clone();
self.request.data = maybe_infer_vector_schema(self.request.data).await?;
// If embeddings were configured via add_embedding(), wrap the data
if !self.embeddings.is_empty() {
let wrapped_data: Box<dyn Scannable> = Box::new(WithEmbeddingsScannable::try_new(
+2 -145
View File
@@ -18,16 +18,13 @@ use crate::embeddings::{
};
use crate::table::{ColumnDefinition, ColumnKind, TableDefinition};
use crate::{Error, Result};
use arrow_array::{ArrayRef, RecordBatch, RecordBatchIterator, RecordBatchReader, cast::AsArray};
use arrow_cast::{CastOptions, cast_with_options};
use arrow_schema::{ArrowError, DataType, Schema, SchemaRef};
use arrow_array::{ArrayRef, RecordBatch, RecordBatchIterator, RecordBatchReader};
use arrow_schema::{ArrowError, SchemaRef};
use async_trait::async_trait;
use futures::StreamExt;
use futures::stream::once;
use lance_datafusion::utils::StreamingWriteSource;
use super::inspect::infer_dimension;
pub trait Scannable: Send {
/// Returns the schema of the data.
fn schema(&self) -> SchemaRef;
@@ -500,146 +497,6 @@ impl Scannable for PeekedScannable {
}
}
fn name_suggests_vector_column(name: &str) -> bool {
let name = name.to_ascii_lowercase();
name == "vec" || name.contains("vector") || name.contains("embedding")
}
/// Infer fixed dimensions for vector-like floating-point list columns.
///
/// Lance vector search requires `FixedSizeList` columns, but Arrow data assembled
/// from runtime embedding models is often represented as `List`. For vector-like
/// column names, inspect the first batch and convert uniform, non-empty lists to a
/// fixed-size schema. Every subsequent batch is cast with strict length checking.
pub(crate) async fn maybe_infer_vector_schema(
data: Box<dyn Scannable>,
) -> Result<Box<dyn Scannable>> {
let input_schema = data.schema();
let candidates = input_schema
.fields()
.iter()
.enumerate()
.filter(|(_, field)| {
name_suggests_vector_column(field.name())
&& matches!(
field.data_type(),
DataType::List(item) | DataType::LargeList(item)
if item.data_type().is_floating()
)
})
.map(|(index, _)| index)
.collect::<Vec<_>>();
if candidates.is_empty() {
return Ok(data);
}
let mut peeked = PeekedScannable::new(data);
let Some(first_batch) = peeked.peek().await else {
return Ok(Box::new(peeked));
};
let mut fields = input_schema.fields().iter().cloned().collect::<Vec<_>>();
let mut changed = false;
for index in candidates {
let array = first_batch.column(index);
let dimension = match array.data_type() {
DataType::List(_) => {
infer_dimension::<arrow_array::types::Int32Type>(array.as_list::<i32>())?
.map(i64::from)
}
DataType::LargeList(_) => {
infer_dimension::<arrow_array::types::Int64Type>(array.as_list::<i64>())?
}
_ => unreachable!(),
};
let Some(dimension) = dimension.filter(|dimension| *dimension > 0) else {
continue;
};
let dimension = i32::try_from(dimension).map_err(|_| Error::InvalidInput {
message: format!(
"Vector column '{}' has a dimension larger than i32::MAX",
fields[index].name()
),
})?;
let item = match fields[index].data_type() {
DataType::List(item) | DataType::LargeList(item) => item.clone(),
_ => unreachable!(),
};
fields[index] = Arc::new(
fields[index]
.as_ref()
.clone()
.with_data_type(DataType::FixedSizeList(item, dimension)),
);
changed = true;
}
if !changed {
return Ok(Box::new(peeked));
}
let output_schema = Arc::new(Schema::new_with_metadata(
fields,
input_schema.metadata().clone(),
));
Ok(Box::new(InferredVectorScannable {
inner: peeked,
output_schema,
}))
}
struct InferredVectorScannable {
inner: PeekedScannable,
output_schema: SchemaRef,
}
impl Scannable for InferredVectorScannable {
fn schema(&self) -> SchemaRef {
self.output_schema.clone()
}
fn scan_as_stream(&mut self) -> SendableRecordBatchStream {
let output_schema = self.output_schema.clone();
let stream_schema = output_schema.clone();
let stream = self.inner.scan_as_stream().map(move |batch| {
let batch = batch?;
let columns = batch
.columns()
.iter()
.zip(output_schema.fields())
.map(|(array, field)| {
if array.data_type() == field.data_type() {
Ok(array.clone())
} else {
cast_with_options(
array,
field.data_type(),
&CastOptions {
safe: false,
..Default::default()
},
)
.map_err(Error::from)
}
})
.collect::<Result<Vec<_>>>()?;
Ok(RecordBatch::try_new(output_schema.clone(), columns)?)
});
Box::pin(SimpleRecordBatchStream {
schema: stream_schema,
stream,
})
}
fn num_rows(&self) -> Option<usize> {
self.inner.num_rows()
}
fn rescannable(&self) -> bool {
self.inner.rescannable()
}
}
/// Compute the number of write partitions based on data size estimates.
///
/// `sample_bytes` and `sample_rows` come from a representative batch and are
+2 -5
View File
@@ -54,6 +54,7 @@
//! - `/path/to/database` - local database on file system.
//! - `s3://bucket/path/to/database` or `gs://bucket/path/to/database` - database on cloud object store
//! - `db://dbname` - Lance Cloud
//! - `db://` with a host override - remote tables in the storage root
//!
//! You can also use [`ConnectBuilder`] to configure the connection to the database.
//!
@@ -72,9 +73,7 @@
//!
//! LanceDB uses [arrow-rs](https://github.com/apache/arrow-rs) to define schema, data types and array itself.
//! It treats [`FixedSizeList<Float16/Float32>`](https://docs.rs/arrow/latest/arrow/array/struct.FixedSizeListArray.html)
//! columns as vector columns. When creating a table with a floating-point `List`
//! column named `vec`, or with `vector` or `embedding` in its name, LanceDB infers
//! a uniform dimension from the first batch and stores it as a `FixedSizeList`.
//! columns as vector columns.
//!
//! For more details, please refer to the [LanceDB documentation](https://docs.lancedb.com).
//!
@@ -84,8 +83,6 @@
//! schema of the `RecordBatch` determines the schema of the table.
//!
//! Vector columns should be represented as `FixedSizeList<Float16/Float32>` data type.
//! A vector-like `List<Float16/Float32>` input is also accepted when every vector
//! has the same runtime dimension.
//!
//! ```rust
//! # use std::sync::Arc;
+2 -87
View File
@@ -1648,8 +1648,8 @@ mod tests {
use super::*;
use arrow::{array::downcast_array, compute::concat_batches, datatypes::Int32Type};
use arrow_array::{
FixedSizeListArray, Float32Array, Int32Array, ListArray, RecordBatch, StringArray,
cast::AsArray, types::Float32Type,
FixedSizeListArray, Float32Array, Int32Array, RecordBatch, StringArray, cast::AsArray,
types::Float32Type,
};
use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema};
use futures::{StreamExt, TryStreamExt};
@@ -2282,91 +2282,6 @@ mod tests {
);
}
#[tokio::test]
async fn vector_search_infers_dimension_from_list_array() {
let tmp_dir = tempdir().unwrap();
let schema = Arc::new(ArrowSchema::new(vec![
ArrowField::new("id", DataType::Int32, false),
ArrowField::new(
"vec",
DataType::List(Arc::new(ArrowField::new("item", DataType::Float32, true))),
true,
),
]));
let vectors = ListArray::from_iter_primitive::<Float32Type, _, _>([
Some([Some(0.0), Some(0.0)]),
Some([Some(1.0), Some(1.0)]),
]);
let batch = RecordBatch::try_new(
schema,
vec![Arc::new(Int32Array::from(vec![0, 1])), Arc::new(vectors)],
)
.unwrap();
let table = connect(tmp_dir.path().to_str().unwrap())
.execute()
.await
.unwrap()
.create_table("vectors", batch)
.execute()
.await
.unwrap();
assert!(matches!(
table.schema().await.unwrap().field(1).data_type(),
DataType::FixedSizeList(_, 2)
));
let results = table
.vector_search(&[0.0, 0.0])
.unwrap()
.limit(1)
.execute()
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap();
assert_eq!(results[0]["id"].as_primitive::<Int32Type>().value(0), 0);
}
#[tokio::test]
async fn inferred_vector_dimension_is_validated_across_batches() {
let tmp_dir = tempdir().unwrap();
let schema = Arc::new(ArrowSchema::new(vec![ArrowField::new(
"vec",
DataType::List(Arc::new(ArrowField::new("item", DataType::Float32, true))),
true,
)]));
let first = RecordBatch::try_new(
schema.clone(),
vec![Arc::new(
ListArray::from_iter_primitive::<Float32Type, _, _>([Some([Some(0.0), Some(0.0)])]),
)],
)
.unwrap();
let wrong_dimension = RecordBatch::try_new(
schema,
vec![Arc::new(
ListArray::from_iter_primitive::<Float32Type, _, _>([Some([
Some(1.0),
Some(1.0),
Some(1.0),
])]),
)],
)
.unwrap();
let result = connect(tmp_dir.path().to_str().unwrap())
.execute()
.await
.unwrap()
.create_table("vectors", vec![first, wrong_dimension])
.execute()
.await;
assert!(result.is_err());
}
#[tokio::test]
async fn test_fast_search_plan() {
let tmp_dir = tempdir().unwrap();
+11 -7
View File
@@ -349,18 +349,22 @@ pub struct ParsedDbUrl {
/// Parse a database URL and extract the database name and optional prefix.
///
/// Expected format: `db://db_name` or `db://db_name/prefix`
/// Expected format: `db://db_name`, `db://db_name/prefix`, or `db://` when
/// connecting to the storage root through a host override.
pub fn parse_db_url(db_url: &str) -> Result<ParsedDbUrl> {
let parsed_url = url::Url::parse(db_url).map_err(|err| Error::InvalidInput {
message: format!("db_url is not a valid URL. '{db_url}'. Error: {err}"),
})?;
debug_assert_eq!(parsed_url.scheme(), "db");
if !parsed_url.has_host() {
return Err(Error::InvalidInput {
message: format!("Invalid database URL (missing host) '{}'", db_url),
});
}
let db_name = parsed_url.host_str().unwrap().to_string();
let db_name = match parsed_url.host_str() {
Some(db_name) => db_name.to_string(),
None if matches!(parsed_url.path(), "" | "/") => String::new(),
None => {
return Err(Error::InvalidInput {
message: format!("Invalid database URL (missing host) '{}'", db_url),
});
}
};
let db_prefix = {
let prefix = parsed_url.path().trim_start_matches('/');
if prefix.is_empty() {
+7
View File
@@ -272,6 +272,13 @@ impl RemoteDatabase {
read_consistency_interval: Option<std::time::Duration>,
) -> Result<Self> {
let parsed = super::client::parse_db_url(uri)?;
if parsed.db_name.is_empty() && host_override.is_none() {
return Err(Error::InvalidInput {
message:
"A host override is required when connecting to the storage root with 'db://'"
.to_string(),
});
}
let header_map = RestfulLanceDbClient::<Sender>::default_headers(
api_key,
region,