From 7a11cf0dffd2feb89c7767f47bc563b4b8793afc Mon Sep 17 00:00:00 2001 From: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com> Date: Tue, 25 Aug 2026 21:35:59 +0000 Subject: [PATCH] fix: preserve deferred query state provenance --- nodejs/__test__/table.test.ts | 24 +++++++++++++++ nodejs/lancedb/query.ts | 13 +++++++- rust/lancedb/src/table.rs | 3 +- rust/lancedb/src/table/dataset.rs | 50 +++++++++++++++++++++++++------ 4 files changed, 78 insertions(+), 12 deletions(-) diff --git a/nodejs/__test__/table.test.ts b/nodejs/__test__/table.test.ts index 345121ab1..a998f1e83 100644 --- a/nodejs/__test__/table.test.ts +++ b/nodejs/__test__/table.test.ts @@ -3424,6 +3424,30 @@ describe("column name options", () => { expect(results[1].query_index).toBe(1); }); + test("observes promised additional vectors while the query is pending", async () => { + const initialVector = new Promise(() => undefined); + const query = table.query().nearestTo(initialVector); + const unhandled: unknown[] = []; + const onUnhandled = (reason: unknown) => unhandled.push(reason); + process.on("unhandledRejection", onUnhandled); + + try { + query.addQueryVector(Promise.reject(new Error("extra vector failed"))); + await new Promise((resolve) => setImmediate(resolve)); + expect(unhandled).toEqual([]); + + const rejectedQuery = table + .query() + .nearestTo([0.1, 0.2]) + .addQueryVector(Promise.reject(new Error("consumed vector failed"))); + await expect(rejectedQuery.toArray()).rejects.toThrow( + "consumed vector failed", + ); + } finally { + process.off("unhandledRejection", onUnhandled); + } + }); + test("index and search multivectors", async () => { const db = await connect(tmpDir.name); const data = []; diff --git a/nodejs/lancedb/query.ts b/nodejs/lancedb/query.ts index 8532222f9..f1d31eae1 100644 --- a/nodejs/lancedb/query.ts +++ b/nodejs/lancedb/query.ts @@ -735,9 +735,20 @@ export class VectorQuery extends StandardQueryBase { */ addQueryVector(vector: IntoVector): VectorQuery { if (vector instanceof Promise) { + // Observe the promise as soon as it is accepted. The existing native + // query may still be pending, and delaying observation until it resolves + // can otherwise surface a fast rejection as unhandled. + const settledVector = vector.then( + (value) => ({ status: "fulfilled" as const, value }), + (reason) => ({ status: "rejected" as const, reason }), + ); const res = (async () => { const inner = await this.getInner(); - addQueryVectorToNative(inner, await vector); + const outcome = await settledVector; + if (outcome.status === "rejected") { + throw outcome.reason; + } + addQueryVectorToNative(inner, outcome.value); return inner; })(); return new VectorQuery(res); diff --git a/rust/lancedb/src/table.rs b/rust/lancedb/src/table.rs index 5d794319b..89897d5c7 100644 --- a/rust/lancedb/src/table.rs +++ b/rust/lancedb/src/table.rs @@ -3077,8 +3077,7 @@ impl BaseTable for NativeTable { } async fn query_snapshot(&self) -> Result> { - let dataset = self.dataset.get().await?; - let snapshot = self.dataset.new_query_snapshot(dataset); + let snapshot = self.dataset.new_query_snapshot().await?; let mut table = self.with_dataset(snapshot); // QueryTable requests do not carry a revision. A pinned snapshot must // execute locally until the namespace API can accept that revision. diff --git a/rust/lancedb/src/table/dataset.rs b/rust/lancedb/src/table/dataset.rs index 3fde6ef62..5e3733b85 100644 --- a/rust/lancedb/src/table/dataset.rs +++ b/rust/lancedb/src/table/dataset.rs @@ -98,17 +98,26 @@ impl DatasetConsistencyWrapper { wrapper } - /// Create an independent read-only wrapper pinned to `dataset` while - /// retaining this wrapper's live MemWAL read context. - pub fn new_query_snapshot(&self, dataset: Arc) -> Self { - let version = dataset.version().version; - let query_snapshot = { - let state = self.state.lock().unwrap_or_else(|e| e.into_inner()); + /// Create an independent read-only wrapper pinned to the current dataset + /// while retaining this wrapper's live MemWAL read context. + pub async fn new_query_snapshot(&self) -> Result { + // Apply the configured consistency policy before taking the snapshot. + // The returned dataset is intentionally discarded: a checkout may race + // after this await, so the dataset and its pin provenance must instead + // be cloned together from one authoritative state sample below. + self.get().await?; + + let (dataset, query_snapshot) = { + let state = self.state.lock()?; // Preserve user time travel so the MemWAL safety guard still sees // it. Latest and already-internal snapshots remain internal pins. - state.query_snapshot || state.pinned_version.is_none() + ( + state.dataset.clone(), + state.query_snapshot || state.pinned_version.is_none(), + ) }; - Self { + let version = dataset.version().version; + Ok(Self { state: Arc::new(Mutex::new(DatasetState { dataset, pinned_version: Some(version), @@ -116,7 +125,7 @@ impl DatasetConsistencyWrapper { })), consistency: ConsistencyMode::Lazy, shard_writer: self.shard_writer.clone(), - } + }) } /// The MemWAL `ShardWriter` cache co-located with this dataset. @@ -490,6 +499,29 @@ mod tests { assert_eq!(wrapper.time_travel_version(), Some(1)); } + #[tokio::test] + async fn test_query_snapshot_samples_dataset_and_pin_together() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().to_str().unwrap(); + let ds = create_test_dataset(uri).await; + + let wrapper = DatasetConsistencyWrapper::new_latest(ds, None); + wrapper.as_time_travel(1u64).await.unwrap(); + let stale_time_travel_dataset = wrapper.get().await.unwrap(); + + append_to_dataset(uri).await; + wrapper.as_latest().await.unwrap(); + + let snapshot = wrapper.new_query_snapshot().await.unwrap(); + let snapshot_dataset = snapshot.get().await.unwrap(); + assert_eq!(snapshot_dataset.version().version, 2); + assert_ne!( + snapshot_dataset.version().version, + stale_time_travel_dataset.version().version + ); + assert_eq!(snapshot.time_travel_version(), None); + } + #[tokio::test] async fn test_as_latest_from_time_travel() { let dir = tempfile::tempdir().unwrap();