Compare commits

..

1 Commits

Author SHA1 Message Date
Gatefixer a23ee5fdac fix(node): reject musl binaries with unresolved AVX-512 symbol 2026-08-05 18:51:18 +00:00
12 changed files with 134 additions and 150 deletions
+1
View File
@@ -15,6 +15,7 @@ renovate.json
src
lancedb
examples
scripts
nodejs-artifacts
Cargo.toml
biome.json
+47
View File
@@ -0,0 +1,47 @@
// 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",
);
});
+2 -1
View File
@@ -86,7 +86,8 @@
"postdocs": "node typedoc_post_process.js",
"lint": "biome check . && biome format .",
"lint-fix": "biome check --write . && biome format --write .",
"prepublishOnly": "napi prepublish -t npm",
"check:native-symbols": "node scripts/check-native-symbols.js",
"prepublishOnly": "pnpm check:native-symbols && napi prepublish -t npm",
"test": "jest --verbose",
"integration": "S3_TEST=1 pnpm test",
"universal": "napi universalize",
+77
View File
@@ -0,0 +1,77 @@
// 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 };
+3 -17
View File
@@ -707,9 +707,6 @@ 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
@@ -759,14 +756,11 @@ 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 self._read_consistency_interval
return LOOP.run(self._conn.get_read_consistency_interval())
@property
def session(self) -> Optional[Session]:
@@ -777,16 +771,8 @@ class LanceDBConnection(DBConnection):
return self._conn.uri
@classmethod
def from_inner(
cls,
inner: LanceDbConnection,
read_consistency_interval: Optional[timedelta],
):
return cls(
None,
read_consistency_interval=read_consistency_interval,
_inner=inner,
)
def from_inner(cls, inner: LanceDbConnection):
return cls(None, _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 await LanceTable.from_inner(inner_tbl)
return LanceTable.from_inner(inner_tbl)
return LOOP.run(do_execute())
+3 -7
View File
@@ -2182,15 +2182,11 @@ class LanceTable(Table):
return self.name
@classmethod
async def from_inner(cls, tbl: LanceDBTable):
from .db import AsyncConnection, LanceDBConnection
def from_inner(cls, tbl: LanceDBTable):
from .db import LanceDBConnection
async_tbl = AsyncTable(tbl)
inner_conn = tbl.database()
read_consistency_interval = await AsyncConnection(
inner_conn
).get_read_consistency_interval()
conn = LanceDBConnection.from_inner(inner_conn, read_consistency_interval)
conn = LanceDBConnection.from_inner(tbl.database())
return cls(
conn,
async_tbl.name,
-17
View File
@@ -77,23 +77,6 @@ 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,7 +6,6 @@ import math
import pytest
from lancedb import DBConnection, Table, connect
from lancedb.background_loop import LOOP
from lancedb.permutation import Permutation, Permutations, permutation_builder
@@ -32,25 +31,6 @@ 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,7 +6,6 @@ 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
@@ -2125,27 +2124,6 @@ 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,9 +745,6 @@ 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 {
-62
View File
@@ -304,68 +304,6 @@ mod tests {
assert_eq!(all_values, expected);
}
#[tokio::test]
async fn test_parallel_compaction_reserves_fragment_ids_once() {
let conn = connect("memory://").execute().await.unwrap();
let schema = Arc::new(Schema::new(vec![Field::new("i", DataType::Int32, false)]));
let batch =
RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from_iter_values(0..10))])
.unwrap();
let table = conn
.create_table("test_parallel_compaction", batch.clone())
.execute()
.await
.unwrap();
// Create 64 fragments. With a 20-row target, compaction plans 32 tasks,
// which is more than the commit retry limit that used to be exhausted
// when each parallel task reserved fragment IDs independently.
for _ in 1..64 {
table.add(batch.clone()).execute().await.unwrap();
}
// Legacy row IDs require fragment IDs before an index can be remapped.
assert!(
!table
.as_native()
.unwrap()
.manifest()
.await
.unwrap()
.uses_stable_row_ids()
);
table
.create_index(&["i"], Index::BTree(BTreeIndexBuilder::default()))
.execute()
.await
.unwrap();
let version_before = table.version().await.unwrap();
let stats = table
.optimize(OptimizeAction::Compact {
options: CompactionOptions {
target_rows_per_fragment: 20,
num_threads: Some(64),
..Default::default()
},
remap_options: None,
})
.await
.unwrap()
.compaction
.unwrap();
assert_eq!(stats.fragments_removed, 64);
assert_eq!(stats.fragments_added, 32);
assert_eq!(table.count_rows(None).await.unwrap(), 640);
assert_eq!(
table.version().await.unwrap(),
version_before + 2,
"parallel compaction should use one fragment reservation commit and one rewrite commit"
);
}
#[tokio::test]
async fn test_optimize_prune_versions() {
let conn = connect("memory://").execute().await.unwrap();