surrealdb-core 3.3.1

A scalable, distributed, collaborative, document-graph database, for the realtime web
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
//! DiskANN behaviour observed through the query driver.

use std::sync::Arc;

use anyhow::Result;

use crate::catalog::providers::CatalogProvider;
use crate::dbs::{NewPlannerStrategy, Session};
use crate::key::schema::RecordKey;
use crate::kvs::{Datastore, QueryRequest, TransactionType};
use crate::val::RecordIdKey;

/// DiskANN counterpart of the HNSW filtered-KNN batching test. A filtered KNN
/// evaluates the residual `WHERE` against each visited candidate's record; the
/// prefetch batches those fetches so the search issues far fewer KV *get
/// operations* than the records it reads (`ops_get` well below `keys_read`).
///
/// It also prints the committed-path over-fetch under a NON-selective filter:
/// the candidate list is prefetched up front, but the result builder fills
/// quickly so most candidates are never evaluated — the gap between
/// `keys_read` and the rows actually needed quantifies the over-fetch the
/// windowed prefetch bounds. DiskANN graph construction is deterministic, so
/// no build seed is needed; the assertions are structural so they hold for any
/// graph.
#[tokio::test(flavor = "multi_thread")]
async fn test_diskann_filtered_knn_batches_record_fetches() -> Result<()> {
	// Compaction happens only where this test calls it: the assertions compare
	// the pending set against the compacted graph, which a background
	// compactor would fold away before the first phase reads it.
	let ds = Datastore::builder().without_maintenance_tasks().build_with_path("memory").await?;
	{
		let tx = ds.transaction(TransactionType::Write).await?;
		tx.ensure_ns_db(None, "test", "test").await?;
		tx.commit().await?;
	}
	let session = Session::owner()
		.with_ns("test")
		.with_db("test")
		.new_planner_strategy(NewPlannerStrategy::AllReadOnlyStatements);

	// 500 deterministic 8-d points with a low-cardinality `category`, plus a
	// DiskANN index. A selective filter makes the search visit many candidates
	// before finding K matches — the case batching helps.
	let n = 500u32;
	let cats = 20u32;
	let mut setup = String::from(
		"DEFINE INDEX emb ON pts FIELDS vec DISKANN DIMENSION 8 DIST EUCLIDEAN TYPE F32;\n",
	);
	for i in 0..n {
		let mut v = String::new();
		for j in 0..8u32 {
			if j > 0 {
				v.push_str(", ");
			}
			let f = ((i.wrapping_mul(7).wrapping_add(j.wrapping_mul(131))) % 1000) as f32 / 1000.0;
			v.push_str(&format!("{f}f"));
		}
		setup.push_str(&format!("CREATE pts:{i} SET vec = [{v}], category = {};\n", i % cats));
	}
	for response in ds.execute(&setup, &session, None).await? {
		response.result?;
	}

	// Run a query on an owned read transaction so we can read its KV metrics.
	async fn run(
		ds: &Arc<Datastore>,
		session: &Session,
		query: &str,
	) -> Result<(usize, crate::observe::TransactionMetricsSnapshot)> {
		let tx = Arc::new(ds.transaction(TransactionType::Read).await?);
		let mut response =
			ds.run(QueryRequest::new(query, session).with_transaction(Arc::clone(&tx))).await?;
		let len = match response.remove(0).result? {
			surrealdb_types::Value::Array(a) => a.len(),
			_ => 0,
		};
		Ok((len, tx.metrics_snapshot_for_test()))
	}

	// Selective filter (category = 7, 1-in-20).
	let selective = "SELECT id FROM pts \
		WHERE vec <|10,400|> [0.5f,0.5f,0.5f,0.5f,0.5f,0.5f,0.5f,0.5f] AND category = 7;";

	// Before compaction the data is in the pending set (`search_pendings`).
	let (pending_len, pending_m) = run(&ds, &session, selective).await?;
	eprintln!(
		"PENDING      ops_get={} keys_read={} value_bytes_read={} results={pending_len}",
		pending_m.ops_get, pending_m.keys_read, pending_m.value_bytes_read
	);

	// Compact into the committed graph, then query again (`search_graph`).
	Datastore::index_compaction(
		Arc::clone(&ds),
		std::time::Duration::from_secs(1),
		tokio_util::sync::CancellationToken::new(),
	)
	.await?;
	let (committed_len, committed_m) = run(&ds, &session, selective).await?;
	eprintln!(
		"COMMITTED    ops_get={} keys_read={} value_bytes_read={} results={committed_len}",
		committed_m.ops_get, committed_m.keys_read, committed_m.value_bytes_read
	);

	// Non-selective filter (category < 10, ~50%) with small k: the builder
	// fills early so most prefetched candidates are never evaluated. Prints the
	// over-fetch (keys_read >> rows needed) that the windowed prefetch bounds.
	let nonselective = "SELECT id FROM pts \
		WHERE vec <|5,400|> [0.5f,0.5f,0.5f,0.5f,0.5f,0.5f,0.5f,0.5f] AND category < 10;";
	let (ns_len, ns_m) = run(&ds, &session, nonselective).await?;
	eprintln!(
		"NONSELECTIVE ops_get={} keys_read={} value_bytes_read={} results={ns_len}",
		ns_m.ops_get, ns_m.keys_read, ns_m.value_bytes_read
	);

	// Both selective paths return K matching records...
	assert_eq!(pending_len, 10, "pending filtered KNN should return K matches");
	assert_eq!(committed_len, 10, "committed filtered KNN should return K matches");
	// ...and both batch their record fetches: one-per-get gives `ops_get` ~
	// `keys_read`; batching pulls `ops_get` well below it.
	assert!(
		u64::from(pending_m.ops_get) * 4 < pending_m.keys_read * 3,
		"pending path should batch: ops_get={} keys_read={}",
		pending_m.ops_get,
		pending_m.keys_read
	);
	assert!(
		u64::from(committed_m.ops_get) * 4 < committed_m.keys_read * 3,
		"committed path should batch: ops_get={} keys_read={}",
		committed_m.ops_get,
		committed_m.keys_read
	);
	assert_eq!(ns_len, 5, "non-selective filtered KNN should return K matches");
	// The windowed prefetch bounds the committed-path over-fetch: a
	// non-selective filter fills the result builder inside the first window,
	// so the search stops fetching once the distance gate closes and reads far
	// fewer keys than a full candidate-list walk (the selective query above,
	// whose builder rarely fills). Before windowing the non-selective path
	// prefetched ~all candidates (measured keys_read ~4x today's); this guards
	// against reintroducing that without pinning a golden number.
	assert!(
		ns_m.keys_read * 4 < committed_m.keys_read,
		"windowed prefetch should bound non-selective over-fetch: \
		 non-selective keys_read={} vs selective keys_read={}",
		ns_m.keys_read,
		committed_m.keys_read
	);
	Ok(())
}

/// C4: a committed candidate whose underlying record row is missing (deleted
/// out from under a not-yet-recompacted graph) must be skipped — `is_record_truthy`
/// already returns not-truthy for a nullish record, and `prefetch_records` marks
/// such ids not-found during the batch warm so the eval loop never issues a
/// redundant per-candidate `get_record`. Here we force the scenario by deleting
/// a record's KV row directly (bypassing the index, which would otherwise drop
/// the graph entry too) and assert the filtered KNN excludes it and backfills.
#[tokio::test(flavor = "multi_thread")]
async fn test_diskann_filtered_knn_skips_missing_record() -> Result<()> {
	// Compaction happens only where this test calls it: the assertions compare
	// the pending set against the compacted graph, which a background
	// compactor would fold away before the first phase reads it.
	let ds = Datastore::builder().without_maintenance_tasks().build_with_path("memory").await?;
	let db_def = {
		let tx = ds.transaction(TransactionType::Write).await?;
		let db = tx.ensure_ns_db(None, "test", "test").await?;
		tx.commit().await?;
		db
	};
	let session = Session::owner()
		.with_ns("test")
		.with_db("test")
		.new_planner_strategy(NewPlannerStrategy::AllReadOnlyStatements);

	// 1-D points 10,20,…,120 with alternating category; "a" selects the odd-id
	// points (10,30,50,…). Nearest "a" matches to query [0] are pts:1 (10),
	// pts:3 (30), pts:5 (50).
	let mut setup = String::from(
		"DEFINE INDEX pt ON pts FIELDS point DISKANN DIMENSION 1 DIST EUCLIDEAN TYPE F32;\n",
	);
	for i in 1..=12u32 {
		let cat = if i % 2 == 1 {
			"a"
		} else {
			"b"
		};
		setup.push_str(&format!("CREATE pts:{i} SET point = [{}f], category = '{cat}';\n", i * 10));
	}
	for response in ds.execute(&setup, &session, None).await? {
		response.result?;
	}

	// Compact so the search runs over the committed graph (`search_graph`).
	Datastore::index_compaction(
		Arc::clone(&ds),
		std::time::Duration::from_secs(1),
		tokio_util::sync::CancellationToken::new(),
	)
	.await?;

	// Delete pts:1's record ROW at the KV layer, leaving the graph entry
	// intact — the committed candidate now resolves to a missing record.
	{
		let tx = ds.transaction(TransactionType::Write).await?;
		let tb = surrealdb_strand::TableName::from("pts");
		let key = RecordKey {
			ns: db_def.namespace_id,
			db: db_def.database_id,
			tb: std::borrow::Cow::Borrowed(&tb),
			id: std::borrow::Cow::Owned(RecordIdKey::Number(1)),
		};
		tx.del_key(&key).await?;
		tx.commit().await?;
	}

	// Top-2 "a" matches to [0]: pts:1 (distance 10) is now missing, so the
	// result must skip it and backfill with pts:3 (30) then pts:5 (50) — and
	// must not error. If the missing record leaked in, the nearest distance
	// would be 10.
	let query = "SELECT VALUE vector::distance::knn() FROM pts \
		WHERE point <|2,40|> [0f] AND category = 'a';";
	let mut dists: Vec<f64> =
		ds.execute(query, &session, None).await?.remove(0).result?.into_t::<Vec<f64>>()?;
	dists.sort_by(f64::total_cmp);
	assert_eq!(
		dists,
		vec![30.0, 50.0],
		"missing pts:1 (dist 10) must be excluded and backfilled, got {dists:?}"
	);
	Ok(())
}

/// A filtered KNN at a `VERSION` reads each visited candidate's record once.
/// Versioned reads bypass the transaction record cache, so the records the
/// windowed prefetch reads must be the ones the filter evaluates; a
/// second, per-candidate read would double the record keys the query reads.
/// The same query unversioned reads each record once through the cache. The
/// versioned run additionally re-reads its K results when it materialises
/// them and reads its catalog at the version, so it may exceed the
/// unversioned count by at most 2K keys — far below the hundreds a second
/// read of every visited candidate adds.
#[cfg(feature = "kv-surrealkv")]
#[tokio::test(flavor = "multi_thread")]
async fn diskann_versioned_filtered_knn_reads_each_candidate_once() -> Result<()> {
	let dir = temp_dir::TempDir::new()?;
	let path = format!("surrealkv://{}?versioned=true&retention=1h", dir.path().to_string_lossy());
	// Compaction happens only where this test calls it.
	let ds = Datastore::builder().without_maintenance_tasks().build_with_path(&path).await?;
	{
		let tx = ds.transaction(TransactionType::Write).await?;
		tx.ensure_ns_db(None, "test", "test").await?;
		tx.commit().await?;
	}
	let session = Session::owner()
		.with_ns("test")
		.with_db("test")
		.new_planner_strategy(NewPlannerStrategy::AllReadOnlyStatements);

	let mut setup = String::from(
		"DEFINE INDEX emb ON pts FIELDS vec DISKANN DIMENSION 8 DIST EUCLIDEAN TYPE F32;\n",
	);
	for i in 0..500u32 {
		let mut v = String::new();
		for j in 0..8u32 {
			if j > 0 {
				v.push_str(", ");
			}
			let f = ((i.wrapping_mul(7).wrapping_add(j.wrapping_mul(131))) % 1000) as f32 / 1000.0;
			v.push_str(&format!("{f}f"));
		}
		setup.push_str(&format!("CREATE pts:{i} SET vec = [{v}], category = {};\n", i % 20));
	}
	for response in ds.execute(&setup, &session, None).await? {
		response.result?;
	}
	Datastore::index_compaction(
		Arc::clone(&ds),
		std::time::Duration::from_secs(1),
		tokio_util::sync::CancellationToken::new(),
	)
	.await?;
	let mut response = ds.execute("RETURN <string> time::now();", &session, None).await?;
	let surrealdb_types::Value::String(stamp) = response.remove(0).result? else {
		panic!("Expected a datetime string");
	};

	async fn run(ds: &Arc<Datastore>, session: &Session, query: &str) -> Result<(usize, u64)> {
		let tx = Arc::new(ds.transaction(TransactionType::Read).await?);
		let mut response =
			ds.run(QueryRequest::new(query, session).with_transaction(Arc::clone(&tx))).await?;
		let len = match response.remove(0).result? {
			surrealdb_types::Value::Array(a) => a.len(),
			_ => 0,
		};
		Ok((len, tx.metrics_snapshot_for_test().keys_read))
	}

	const K: u64 = 10;
	let query = format!(
		"SELECT id FROM pts \
		 WHERE vec <|{K},400|> [0.5f,0.5f,0.5f,0.5f,0.5f,0.5f,0.5f,0.5f] AND category = 7"
	);
	let current = format!("{query};");
	let versioned = format!("{query} VERSION d'{stamp}';");
	// The first search loads the graph into the process-wide DiskANN cache; warm
	// it so both measured runs read only candidate records and doc-id maps.
	run(&ds, &session, &current).await?;
	run(&ds, &session, &versioned).await?;
	let (current_len, current_keys) = run(&ds, &session, &current).await?;
	let (versioned_len, versioned_keys) = run(&ds, &session, &versioned).await?;
	eprintln!("CURRENT keys_read={current_keys} VERSIONED keys_read={versioned_keys}");

	assert_eq!(current_len as u64, K, "filtered KNN should return K matches");
	assert_eq!(versioned_len as u64, K, "versioned filtered KNN should return K matches");
	assert!(
		versioned_keys <= current_keys + 2 * K,
		"a versioned filtered KNN should read each candidate record once: \
		 versioned keys_read={versioned_keys}, current keys_read={current_keys}"
	);
	ds.shutdown().await?;
	Ok(())
}

/// A filtered KNN over a compacted DiskANN graph finds the admitted documents
/// nearest the query even when the filter admits 1% of the table: the walk
/// evaluates the filter on every element it reaches rather than on the `l`
/// nearest overall, which here hold one of the five. Both truthy filters take
/// that walk — the table's per-record SELECT permission for a record user's
/// bare KNN, and a residual condition for the owner — and both match the
/// brute-force answer. A search list too narrow for one walk to admit K is
/// widened until it does.
#[tokio::test(flavor = "multi_thread")]
async fn diskann_selective_filter_finds_admitted_neighbours_past_the_search_list() -> Result<()> {
	// Compaction happens only where this test calls it, so the queries below
	// search the compacted graph rather than the pending queue.
	let ds = Datastore::builder().without_maintenance_tasks().build_with_path("memory").await?;
	{
		let tx = ds.transaction(TransactionType::Write).await?;
		tx.ensure_ns_db(None, "test", "test").await?;
		tx.commit().await?;
	}
	let owner = Session::owner()
		.with_ns("test")
		.with_db("test")
		.new_planner_strategy(NewPlannerStrategy::AllReadOnlyStatements);
	let record = Session::for_record(
		"test",
		"test",
		"user",
		crate::types::PublicValue::String("user:1".to_owned()),
	)
	.new_planner_strategy(NewPlannerStrategy::AllReadOnlyStatements);

	// 1000 pseudo-random 8-d points from a fixed-seed generator, spread like
	// real embeddings rather than along the few lines a modular lattice traces;
	// the permission and the residual condition both admit `category = 7`, one
	// row in a hundred.
	let mut setup = String::from(
		"DEFINE TABLE pts SCHEMALESS PERMISSIONS FOR select WHERE category = 7 \
		 FOR create, update, delete NONE;\n\
		 DEFINE INDEX emb ON pts FIELDS vec DISKANN DIMENSION 8 DIST EUCLIDEAN TYPE F32;\n",
	);
	let mut state = 0x9E37_79B9_7F4A_7C15u64;
	for i in 0..1000u32 {
		let mut v = String::new();
		for j in 0..8u32 {
			if j > 0 {
				v.push_str(", ");
			}
			// xorshift64*
			state ^= state >> 12;
			state ^= state << 25;
			state ^= state >> 27;
			let f = (state.wrapping_mul(0x2545_F491_4F6C_DD1D) >> 40) as f32 / (1u64 << 24) as f32;
			v.push_str(&format!("{f}f"));
		}
		setup.push_str(&format!("CREATE pts:{i} SET vec = [{v}], category = {};\n", i % 100));
	}
	for response in ds.execute(&setup, &owner, None).await? {
		response.result?;
	}
	Datastore::index_compaction(
		Arc::clone(&ds),
		std::time::Duration::from_secs(1),
		tokio_util::sync::CancellationToken::new(),
	)
	.await?;

	async fn ids(ds: &Datastore, session: &Session, query: &str) -> Result<Vec<String>> {
		let mut response = ds.execute(query, session, None).await?;
		let surrealdb_types::Value::Array(rows) = response.remove(0).result? else {
			panic!("Expected array result");
		};
		let mut ids: Vec<String> = rows.into_iter().map(|v| format!("{v:?}")).collect();
		ids.sort();
		Ok(ids)
	}

	let q = "[0.5f,0.5f,0.5f,0.5f,0.5f,0.5f,0.5f,0.5f]";
	let truth = ids(
		&ds,
		&owner,
		&format!("SELECT VALUE id FROM pts WHERE vec <|5,EUCLIDEAN|> {q} AND category = 7;"),
	)
	.await?;
	assert_eq!(truth.len(), 5, "brute force should find K admitted rows");

	// A search list too narrow for one walk to reach k visible rows widens and
	// walks again until it holds k.
	let narrow =
		ids(&ds, &record, &format!("SELECT VALUE id FROM pts WHERE vec <|5,10|> {q};")).await?;
	assert_eq!(narrow.len(), 5, "a narrow search list should still return K visible rows");

	let permitted =
		ids(&ds, &record, &format!("SELECT VALUE id FROM pts WHERE vec <|5,100|> {q};")).await?;
	assert_eq!(
		permitted, truth,
		"a record user's bare KNN should return its K nearest visible rows"
	);

	let conditioned = ids(
		&ds,
		&owner,
		&format!("SELECT VALUE id FROM pts WHERE vec <|5,100|> {q} AND category = 7;"),
	)
	.await?;
	assert_eq!(conditioned, truth, "a residual condition should return its K nearest matches");
	Ok(())
}