mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-03 12:08:52 +00:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 72ac16ba76 | |||
| 7357d63e87 | |||
| 624a75edf7 |
@@ -15,7 +15,6 @@ renovate.json
|
||||
src
|
||||
lancedb
|
||||
examples
|
||||
scripts
|
||||
nodejs-artifacts
|
||||
Cargo.toml
|
||||
biome.json
|
||||
|
||||
@@ -1,47 +0,0 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
const {
|
||||
checkNativeBinary,
|
||||
findForbiddenUndefinedSymbols,
|
||||
} = require("../scripts/check-native-symbols.js");
|
||||
|
||||
test("detects the unresolved AVX-512 symbol from broken musl binaries", () => {
|
||||
const output = [
|
||||
" U napi_create_function",
|
||||
" U sum_4bit_dist_table_32bytes_batch_avx512",
|
||||
" U strlen",
|
||||
].join("\n");
|
||||
|
||||
expect(findForbiddenUndefinedSymbols(output)).toEqual([
|
||||
"sum_4bit_dist_table_32bytes_batch_avx512",
|
||||
]);
|
||||
});
|
||||
|
||||
test("accepts native binaries without unresolved internal Lance symbols", () => {
|
||||
const runNm = jest.fn(() => ({
|
||||
status: 0,
|
||||
stdout:
|
||||
" U napi_create_function\n U strlen\n",
|
||||
stderr: "",
|
||||
}));
|
||||
|
||||
expect(() => checkNativeBinary("lancedb.node", runNm)).not.toThrow();
|
||||
expect(runNm).toHaveBeenCalledWith(
|
||||
"nm",
|
||||
["-D", "--undefined-only", "lancedb.node"],
|
||||
{ encoding: "utf8" },
|
||||
);
|
||||
});
|
||||
|
||||
test("rejects native binaries with the unresolved AVX-512 symbol", () => {
|
||||
const runNm = jest.fn(() => ({
|
||||
status: 0,
|
||||
stdout: " U sum_4bit_dist_table_32bytes_batch_avx512\n",
|
||||
stderr: "",
|
||||
}));
|
||||
|
||||
expect(() => checkNativeBinary("lancedb.node", runNm)).toThrow(
|
||||
"lancedb.node contains unresolved internal Lance symbols: sum_4bit_dist_table_32bytes_batch_avx512",
|
||||
);
|
||||
});
|
||||
+1
-2
@@ -86,8 +86,7 @@
|
||||
"postdocs": "node typedoc_post_process.js",
|
||||
"lint": "biome check . && biome format .",
|
||||
"lint-fix": "biome check --write . && biome format --write .",
|
||||
"check:native-symbols": "node scripts/check-native-symbols.js",
|
||||
"prepublishOnly": "pnpm check:native-symbols && napi prepublish -t npm",
|
||||
"prepublishOnly": "napi prepublish -t npm",
|
||||
"test": "jest --verbose",
|
||||
"integration": "S3_TEST=1 pnpm test",
|
||||
"universal": "napi universalize",
|
||||
|
||||
@@ -1,77 +0,0 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
const { spawnSync } = require("node:child_process");
|
||||
const { existsSync } = require("node:fs");
|
||||
const path = require("node:path");
|
||||
|
||||
const FORBIDDEN_UNDEFINED_SYMBOLS = new Set([
|
||||
"sum_4bit_dist_table_32bytes_batch_avx512",
|
||||
]);
|
||||
|
||||
function findForbiddenUndefinedSymbols(output) {
|
||||
const found = new Set();
|
||||
|
||||
for (const line of output.split(/\r?\n/)) {
|
||||
const columns = line.trim().split(/\s+/);
|
||||
const symbol = columns.at(-1)?.split("@")[0];
|
||||
if (symbol && FORBIDDEN_UNDEFINED_SYMBOLS.has(symbol)) {
|
||||
found.add(symbol);
|
||||
}
|
||||
}
|
||||
|
||||
return [...found].sort();
|
||||
}
|
||||
|
||||
function checkNativeBinary(binaryPath, runNm = spawnSync) {
|
||||
const result = runNm("nm", ["-D", "--undefined-only", binaryPath], {
|
||||
encoding: "utf8",
|
||||
});
|
||||
|
||||
if (result.error) {
|
||||
throw new Error(`Unable to inspect ${binaryPath}: ${result.error.message}`);
|
||||
}
|
||||
if (result.status !== 0) {
|
||||
throw new Error(
|
||||
`Unable to inspect ${binaryPath}: nm exited with status ${result.status}\n${result.stderr}`,
|
||||
);
|
||||
}
|
||||
|
||||
const forbidden = findForbiddenUndefinedSymbols(result.stdout);
|
||||
if (forbidden.length > 0) {
|
||||
throw new Error(
|
||||
`${binaryPath} contains unresolved internal Lance symbols: ${forbidden.join(
|
||||
", ",
|
||||
)}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
function main() {
|
||||
const binaryPath = path.resolve(
|
||||
__dirname,
|
||||
"..",
|
||||
"npm",
|
||||
"linux-x64-musl",
|
||||
"lancedb.linux-x64-musl.node",
|
||||
);
|
||||
|
||||
if (!existsSync(binaryPath)) {
|
||||
throw new Error(
|
||||
`Missing ${binaryPath}; assemble the native artifacts before publishing`,
|
||||
);
|
||||
}
|
||||
|
||||
checkNativeBinary(binaryPath);
|
||||
}
|
||||
|
||||
if (require.main === module) {
|
||||
try {
|
||||
main();
|
||||
} catch (error) {
|
||||
console.error(error instanceof Error ? error.message : error);
|
||||
process.exitCode = 1;
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = { checkNativeBinary, findForbiddenUndefinedSymbols };
|
||||
@@ -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})"
|
||||
|
||||
@@ -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())
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -656,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")]
|
||||
@@ -768,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
|
||||
///
|
||||
@@ -1354,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() {
|
||||
|
||||
@@ -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.
|
||||
//!
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user