mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-30 00:45:37 +00:00
fix(python): prevent BlobFile finalizer panics on Tokio workers (#4291)
## Cause `io.RawIOBase.__del__` checks `closed` and calls `close()` during GC. Both PyO3 methods entered `block_on`, which panics when GC runs on a LanceDB Tokio worker. ## Fix Track Python blob handle closure atomically so `closed` is synchronous. `close()` marks the handle closed immediately and, when called on a Tokio thread, schedules the underlying async cleanup without blocking that thread. Reads and seeks after close fail immediately; ordinary callers still wait for cleanup. ## Validation - Rebuilt the Python extension with `uv run --extra tests --extra dev maturin develop --extras tests,dev`. - Ran the blob test file and remote blob handle test: 69 passed, 1 skipped. - Ran `cargo fmt --all`, `ruff format .`, and `ruff check .`. Fixes #4283 <!-- lance-gatekeeper-fix:v1 agent=15a7bd545e63dea4b69faa0429f0a684 generation=1 --> Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com>
This commit is contained in:
@@ -921,6 +921,70 @@ def test_fetch_blob_files_lazy_read():
|
||||
handles = table.fetch_blob_files("image", [by_id[1]])
|
||||
assert len(handles) == 1
|
||||
assert handles[0].read() == payload
|
||||
handles[0].close()
|
||||
assert handles[0].closed
|
||||
with pytest.raises(RuntimeError, match="already closed"):
|
||||
handles[0].read()
|
||||
|
||||
|
||||
def test_blob_file_finalizer_and_close_on_runtime_worker():
|
||||
# A tokio panic in IOBase.__del__ is printed but does not fail the add call.
|
||||
# Run in a child process so the assertion catches that panic on stderr.
|
||||
script = textwrap.dedent(
|
||||
"""\
|
||||
import gc
|
||||
import threading
|
||||
|
||||
import lancedb
|
||||
import pyarrow as pa
|
||||
|
||||
db = lancedb.connect("memory:///")
|
||||
schema = pa.schema([pa.field("id", pa.int64()), lancedb.blob("image")])
|
||||
table = db.create_table("finalize", schema=schema)
|
||||
table.add([{"id": 1, "image": b"first"}])
|
||||
hits = table.search().with_row_id(True).limit(1).to_arrow()
|
||||
row_id = hits["_rowid"][0].as_py()
|
||||
gc_handle, explicit_handle = table.fetch_blob_files("image", [row_id, row_id])
|
||||
|
||||
gc.disable()
|
||||
cycle = [gc_handle]
|
||||
cycle.append(cycle)
|
||||
del gc_handle, cycle
|
||||
|
||||
observed = []
|
||||
def on_progress(_):
|
||||
if not observed:
|
||||
observed.append((
|
||||
threading.current_thread() is not threading.main_thread(),
|
||||
explicit_handle.closed,
|
||||
gc.collect(),
|
||||
))
|
||||
explicit_handle.close()
|
||||
|
||||
table.add([{"id": 2, "image": b"second"}], progress=on_progress)
|
||||
gc.enable()
|
||||
|
||||
assert len(observed) == 1, observed
|
||||
assert observed[0][0] is True, observed
|
||||
assert observed[0][1] is False, observed
|
||||
assert observed[0][2] > 0, observed
|
||||
assert explicit_handle.closed
|
||||
explicit_handle.close() # Close remains idempotent after the worker call.
|
||||
try:
|
||||
explicit_handle.read()
|
||||
except RuntimeError as error:
|
||||
assert "already closed" in str(error)
|
||||
else:
|
||||
raise AssertionError("a closed blob file was readable")
|
||||
print("worker close completed")
|
||||
"""
|
||||
)
|
||||
result = subprocess.run(
|
||||
[sys.executable, "-c", script], capture_output=True, text=True, check=False
|
||||
)
|
||||
assert result.returncode == 0, result.stderr
|
||||
assert "Cannot start a runtime from within a runtime" not in result.stderr
|
||||
assert "worker close completed" in result.stdout
|
||||
|
||||
|
||||
def test_fetch_blob_files_null_alignment():
|
||||
|
||||
@@ -2470,6 +2470,10 @@ def test_remote_blob_files_are_lazy_seekable_handles():
|
||||
assert alpha.read_range(1, 3) == b"lph"
|
||||
gamma.seek(2)
|
||||
assert gamma.read() == b"mma"
|
||||
alpha.close()
|
||||
assert alpha.closed
|
||||
with pytest.raises(RuntimeError, match="already closed"):
|
||||
alpha.read_range(0, 1)
|
||||
|
||||
|
||||
def test_remote_blob_fetch_accepts_query_table():
|
||||
|
||||
@@ -197,6 +197,27 @@ pub fn block_on<F: std::future::Future>(fut: F) -> F::Output {
|
||||
get_runtime().block_on(fut)
|
||||
}
|
||||
|
||||
/// Run a detached task, including when the caller is already on a runtime worker.
|
||||
/// Keep it visible to [`shutdown`] until it finishes.
|
||||
pub fn spawn_background<F>(fut: F)
|
||||
where
|
||||
F: Future<Output = ()> + Send + 'static,
|
||||
{
|
||||
let guard = OutstandingGuard::new();
|
||||
let task = async move {
|
||||
let _guard = guard;
|
||||
fut.await;
|
||||
};
|
||||
match runtime::Handle::try_current() {
|
||||
Ok(handle) => {
|
||||
let _ = handle.spawn(task);
|
||||
}
|
||||
Err(_) => {
|
||||
let _ = get_runtime().spawn(task);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Gracefully quiesce the shared runtime, meant to run at normal process exit.
|
||||
///
|
||||
/// Waits (bounded by `timeout`) for [`OUTSTANDING`] to reach zero -- i.e.
|
||||
|
||||
+40
-6
@@ -1,8 +1,14 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
use std::{collections::HashMap, sync::Arc};
|
||||
use std::{
|
||||
collections::HashMap,
|
||||
sync::{
|
||||
Arc,
|
||||
atomic::{AtomicBool, Ordering},
|
||||
},
|
||||
};
|
||||
|
||||
use crate::runtime::{block_on, future_into_py};
|
||||
use crate::runtime::{block_on, future_into_py, spawn_background};
|
||||
use crate::{
|
||||
connection::Connection,
|
||||
error::PythonErrorExt,
|
||||
@@ -680,11 +686,23 @@ impl From<lancedb::table::DropColumnsResult> for DropColumnsResult {
|
||||
#[pyclass(name = "BlobFile")]
|
||||
pub struct PyBlobFile {
|
||||
inner: Arc<BlobFile>,
|
||||
closed: AtomicBool,
|
||||
}
|
||||
|
||||
impl PyBlobFile {
|
||||
fn ensure_open(&self) -> PyResult<()> {
|
||||
if self.closed.load(Ordering::Acquire) {
|
||||
Err(PyRuntimeError::new_err("blob file is already closed"))
|
||||
} else {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[pymethods]
|
||||
impl PyBlobFile {
|
||||
fn read_bytes(self_: PyRef<'_, Self>) -> PyResult<Py<PyBytes>> {
|
||||
self_.ensure_open()?;
|
||||
let inner = self_.inner.clone();
|
||||
let py = self_.py();
|
||||
let bytes = py
|
||||
@@ -694,6 +712,7 @@ impl PyBlobFile {
|
||||
}
|
||||
|
||||
pub fn read(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
|
||||
self_.ensure_open()?;
|
||||
let inner = self_.inner.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
let bytes = inner
|
||||
@@ -705,7 +724,20 @@ impl PyBlobFile {
|
||||
}
|
||||
|
||||
fn close(self_: PyRef<'_, Self>) -> PyResult<()> {
|
||||
if self_.closed.swap(true, Ordering::AcqRel) {
|
||||
return Ok(());
|
||||
}
|
||||
let inner = self_.inner.clone();
|
||||
if tokio::runtime::Handle::try_current().is_ok() {
|
||||
// IOBase.__del__ can call close while cyclic GC runs on a worker.
|
||||
// The status changes immediately; release Lance's async state there.
|
||||
spawn_background(async move {
|
||||
if let Err(error) = inner.close().await {
|
||||
log::warn!("blob close failed: {error}");
|
||||
}
|
||||
});
|
||||
return Ok(());
|
||||
}
|
||||
self_
|
||||
.py()
|
||||
.detach(move || block_on(async move { inner.close().await }))
|
||||
@@ -713,13 +745,11 @@ impl PyBlobFile {
|
||||
}
|
||||
|
||||
fn is_closed(self_: PyRef<'_, Self>) -> bool {
|
||||
let inner = self_.inner.clone();
|
||||
self_
|
||||
.py()
|
||||
.detach(move || block_on(async move { inner.is_closed().await }))
|
||||
self_.closed.load(Ordering::Acquire)
|
||||
}
|
||||
|
||||
fn seek(self_: PyRef<'_, Self>, position: u64) -> PyResult<()> {
|
||||
self_.ensure_open()?;
|
||||
let inner = self_.inner.clone();
|
||||
self_
|
||||
.py()
|
||||
@@ -728,6 +758,7 @@ impl PyBlobFile {
|
||||
}
|
||||
|
||||
fn tell(self_: PyRef<'_, Self>) -> PyResult<u64> {
|
||||
self_.ensure_open()?;
|
||||
let inner = self_.inner.clone();
|
||||
self_
|
||||
.py()
|
||||
@@ -741,6 +772,7 @@ impl PyBlobFile {
|
||||
|
||||
/// Read a blob-local byte range without moving the cursor.
|
||||
fn read_range(self_: PyRef<'_, Self>, offset: u64, length: usize) -> PyResult<Py<PyBytes>> {
|
||||
self_.ensure_open()?;
|
||||
let end = offset
|
||||
.checked_add(length as u64)
|
||||
.ok_or_else(|| PyValueError::new_err("offset + length overflowed"))?;
|
||||
@@ -753,6 +785,7 @@ impl PyBlobFile {
|
||||
}
|
||||
|
||||
fn read_up_to(self_: PyRef<'_, Self>, length: usize) -> PyResult<Py<PyBytes>> {
|
||||
self_.ensure_open()?;
|
||||
let inner = self_.inner.clone();
|
||||
let py = self_.py();
|
||||
let bytes = py
|
||||
@@ -1455,6 +1488,7 @@ impl Table {
|
||||
.map(|handle| {
|
||||
handle.map(|file| PyBlobFile {
|
||||
inner: Arc::new(file),
|
||||
closed: AtomicBool::new(false),
|
||||
})
|
||||
})
|
||||
.collect::<Vec<_>>())
|
||||
|
||||
Reference in New Issue
Block a user