fix: preserve deferred query state provenance

This commit is contained in:
Gatefixer
2026-08-25 21:35:59 +00:00
parent d83353680e
commit 7a11cf0dff
4 changed files with 78 additions and 12 deletions
+24
View File
@@ -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<number[]>(() => 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<void>((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 = [];
+12 -1
View File
@@ -735,9 +735,20 @@ export class VectorQuery extends StandardQueryBase<NativeVectorQuery> {
*/
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);
+1 -2
View File
@@ -3077,8 +3077,7 @@ impl BaseTable for NativeTable {
}
async fn query_snapshot(&self) -> Result<Arc<dyn BaseTable>> {
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.
+41 -9
View File
@@ -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<Dataset>) -> 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<Self> {
// 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();