surrealdb-core 3.2.5

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
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
//! Tests for the changefeed garbage collector's bounds.
//!
//! The collector deletes every changefeed entry behind each database's retention
//! watermark. That backlog is sized by write throughput rather than by the
//! catalog, so it is deleted in bounded pages, each committed before the next is
//! scanned, and one pass spends a fixed key budget shared across the databases it
//! visits. These tests assert those bounds hold: the backlog is charged as writes
//! page by page rather than accumulating in one transaction, no database can
//! spend the whole pass on its own backlog, and an interrupted pass leaves what
//! it deleted deleted.
//!
//! Ported from the module #1023 added upstream. The ranges here are raw
//! `Range<Vec<u8>>` rather than a typed subspace, and the scans use `keys` —
//! ascending — because this branch has no typed key ranges and its `keysr` is
//! the reverse-order scan.

use std::ops::Range;
use std::sync::Arc;

use tokio_util::sync::CancellationToken;
use web_time::Duration;

use crate::catalog::providers::{DatabaseProvider, NamespaceProvider};
use crate::catalog::{DatabaseId, NamespaceId};
use crate::cf::gc::retention_still_reaches;
use crate::dbs::{Capabilities, Session};
use crate::kvs::LockType::Optimistic;
use crate::kvs::TransactionType::{Read, Write};
use crate::kvs::tasklease::{LeaseHandler, TaskLeaseType};
use crate::kvs::test_support::{WriteTxSizeObserver, cleanup_sizes};
use crate::kvs::{CHANGEFEED_GC_BATCH_SIZE, Datastore, KVKey};
use crate::observe::ExecutionObserver;

/// A retention of one nanosecond puts every entry written before the pass behind
/// the watermark, so a whole backlog is collectable without waiting for it to
/// age.
const STALE: &str = "1ns";

/// A memory datastore whose every write transaction is recorded. This branch
/// starts no maintenance task from the builder — the server drives them — so the
/// only transactions observed are the ones the test drives.
async fn observed_ds(observer: Arc<WriteTxSizeObserver>) -> Arc<Datastore> {
	Arc::new(
		Datastore::builder()
			.with_capabilities(Capabilities::all())
			.with_observer(observer as Arc<dyn ExecutionObserver>)
			.build_with_path("memory")
			.await
			.unwrap(),
	)
}

/// The whole changefeed keyspace of one database, whatever its entries'
/// timestamps.
async fn feed(ds: &Datastore, ns: &str, db: &str) -> Range<Vec<u8>> {
	let (ns_id, db_id) = db_ids(ds, ns, db).await;
	let beg = crate::key::change::prefix(ns_id, db_id).encode_key().unwrap();
	let end = crate::key::change::suffix(ns_id, db_id).encode_key().unwrap();
	beg..end
}

async fn count_range(ds: &Datastore, range: Range<Vec<u8>>) -> usize {
	let txn = ds.transaction(Read, Optimistic).await.unwrap();
	let keys = txn.keys(range, u32::MAX, 0, None).await.unwrap();
	let _ = txn.cancel().await;
	keys.len()
}

/// Grow a database's backlog to `count` entries beyond the one its setup wrote.
///
/// The entries are derived from a real changefeed key, so each carries a
/// timestamp the watermark covers — the subspace's own lower bound sorts below
/// the encoded earliest timestamp and so falls outside the collected range. One
/// statement produces one entry however many records it touches, a changefeed key
/// being per transaction, so a backlog of this size is built here rather than by
/// churning records.
async fn grow_backlog(ds: &Datastore, feed: &Range<Vec<u8>>, count: usize) {
	let txn = ds.transaction(Read, Optimistic).await.unwrap();
	let base = txn.keys(feed.clone(), 1, 0, None).await.unwrap().remove(0);
	let _ = txn.cancel().await;
	let tx = ds.transaction(Write, Optimistic).await.unwrap();
	for i in 0..count {
		let mut key = base.clone();
		key.extend_from_slice(&(i as u64).to_be_bytes());
		tx.set(&key, &vec![1u8]).await.unwrap();
	}
	tx.commit().await.unwrap();
}

/// Define a database whose changefeed already holds one entry.
async fn with_changefeed(ds: &Datastore, ns: &str, db: &str) {
	let ses = Session::owner().with_ns(ns).with_db(db);
	for res in ds
		.execute(
			&format!(
				"DEFINE DATABASE {db};
				 DEFINE TABLE thing CHANGEFEED {STALE};
				 CREATE thing:1;"
			),
			&ses,
			None,
		)
		.await
		.unwrap()
	{
		res.result.unwrap();
	}
}

/// The encoded watermark a pass would build right now for `expiry`.
async fn watermark_now(ds: &Datastore, expiry: std::time::Duration) -> Vec<u8> {
	let txn = ds.transaction(Read, Optimistic).await.unwrap();
	let ts_impl = txn.timestamp_impl();
	let ts = txn.timestamp().await.unwrap();
	let _ = txn.cancel().await;
	let watermark = ts.sub_checked(expiry).unwrap_or_else(|| ts_impl.earliest());
	let mut buf = [0u8; 32];
	watermark.encode(&mut buf).to_vec()
}

/// The ids of a database, for the revalidation call.
async fn db_ids(ds: &Datastore, ns: &str, db: &str) -> (NamespaceId, DatabaseId) {
	let txn = ds.transaction(Read, Optimistic).await.unwrap();
	let ns_def = txn.get_ns_by_name(ns, None).await.unwrap().unwrap();
	let db_def = txn.get_db_by_name(ns, db, None).await.unwrap().unwrap();
	let _ = txn.cancel().await;
	(ns_def.namespace_id, db_def.database_id)
}

/// This database's changefeed-retention fence, or `None` while no retention has
/// ever been written.
async fn fence(ds: &Datastore, ns: NamespaceId, db: DatabaseId) -> Option<u64> {
	let txn = ds.transaction(Read, Optimistic).await.unwrap();
	let v = txn.get(&crate::key::database::cf::new(ns, db), None).await.unwrap();
	let _ = txn.cancel().await;
	v
}

/// The whole backlog is deleted, in transactions each bounded to one page.
///
/// The backlog grows with write throughput and has no upper bound, so deleting it
/// in one transaction makes the collector's memory a function of how much was
/// written since the last pass. Paging it also charges every deleted entry to a
/// transaction's write set, which a whole-range delete does not report at all.
#[tokio::test]
async fn changefeed_gc_deletes_the_backlog_in_bounded_committed_pages() {
	let page = CHANGEFEED_GC_BATCH_SIZE as usize;
	let observer = Arc::new(WriteTxSizeObserver::default());
	let ds = observed_ds(Arc::clone(&observer)).await;
	ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
	with_changefeed(&ds, "test", "tenant").await;

	let feed = feed(&ds, "test", "tenant").await;
	// More than two pages of backlog, so the collector has to commit several
	// times.
	grow_backlog(&ds, &feed, page * 2 + page / 2).await;
	let total = count_range(&ds, feed.clone()).await;
	assert!(total > 2 * page, "the backlog must hold more than two pages, holds {total}");

	observer.clear();
	ds.changefeed_process(&Duration::from_secs(1), &CancellationToken::new()).await.unwrap();

	assert_eq!(count_range(&ds, feed).await, 0, "the collector must empty the backlog");
	let sizes = observer.sizes();
	// One page of deletes plus one: arming the retention fence writes the fence
	// key inside the same transaction. Upstream arms with a locked read, which
	// registers the conflict without writing, so its pages carry exactly one
	// page; this branch has no locked-read primitive and pays one write per page
	// for the same guarantee.
	assert!(
		sizes.iter().all(|&n| n <= page as u64 + 1),
		"a collection write transaction carried more than one page plus its fence arm: {sizes:?}"
	);
	assert!(
		sizes.iter().sum::<u64>() >= total as u64,
		"every collected entry must be charged as a write: {sizes:?}"
	);
	assert!(
		cleanup_sizes(&sizes).len() >= total / page,
		"the backlog must be spread across a transaction per page: {sizes:?}"
	);
}

/// A pass's key budget bounds the pass, not each database it visits.
///
/// Databases are visited in catalog order, so one that spent the whole budget on
/// its own backlog would leave every database behind it uncollected until it
/// finished — however many passes that takes.
#[tokio::test]
async fn changefeed_gc_shares_its_budget_across_databases() {
	let page = CHANGEFEED_GC_BATCH_SIZE as usize;
	let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
	ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
	// `aaa` is visited first and carries a backlog larger than its share of the
	// pass, so a pass that spent the budget in visit order would never reach
	// `zzz`.
	with_changefeed(&ds, "test", "aaa").await;
	with_changefeed(&ds, "test", "zzz").await;

	let aaa = feed(&ds, "test", "aaa").await;
	let zzz = feed(&ds, "test", "zzz").await;
	grow_backlog(&ds, &aaa, page * 4).await;
	grow_backlog(&ds, &zzz, page / 2).await;
	assert!(count_range(&ds, zzz.clone()).await > 0, "the trailing database needs a backlog");

	// A budget one database's backlog alone would exhaust.
	ds.changefeed_gc_with_budget(
		&Duration::from_secs(1),
		(page * 2) as u64,
		&CancellationToken::new(),
	)
	.await
	.unwrap();

	assert!(
		count_range(&ds, aaa).await > 0,
		"a budget smaller than the leading backlog must leave part of it behind"
	);
	assert_eq!(
		count_range(&ds, zzz).await,
		0,
		"the trailing database must be collected in the same pass"
	);
}

/// An interrupted pass leaves what it deleted deleted, so the next pass resumes
/// at the head of what survives instead of redoing the work.
///
/// The collected range's lower bound is fixed at the storage engine's earliest
/// timestamp, which is what lets a pass resume without a durable cursor:
/// everything below the first surviving entry is already gone.
#[tokio::test]
async fn changefeed_gc_makes_durable_progress_across_passes() {
	let page = CHANGEFEED_GC_BATCH_SIZE as usize;
	let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
	ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
	with_changefeed(&ds, "test", "tenant").await;

	let feed = feed(&ds, "test", "tenant").await;
	grow_backlog(&ds, &feed, page * 3).await;
	let mut left = count_range(&ds, feed.clone()).await;

	// One page's worth of budget per pass: each must retire a page and no more,
	// and the backlog must shrink monotonically to nothing.
	for _ in 0..10 {
		ds.changefeed_gc_with_budget(
			&Duration::from_secs(1),
			page as u64,
			&CancellationToken::new(),
		)
		.await
		.unwrap();
		let now = count_range(&ds, feed.clone()).await;
		assert!(now < left, "a pass must make progress: {left} then {now}");
		left = now;
		if left == 0 {
			break;
		}
	}
	assert_eq!(left, 0, "repeated bounded passes must drain the backlog");
}

/// A stale range's upper bound only holds while the retention that produced it
/// does.
///
/// The watermark is computed once per pass from the catalog, but the deletes it
/// authorises are committed in separate transactions over many pages. A
/// `DEFINE TABLE` or `DEFINE DATABASE` extending retention in that window moves
/// the watermark earlier, making entries the range still names live again — and a
/// deleted changefeed entry cannot be recovered. The bound is therefore
/// revalidated inside every page transaction, which is what this asserts.
#[tokio::test]
async fn an_extended_retention_withdraws_the_stale_bound() {
	let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
	ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
	with_changefeed(&ds, "test", "tenant").await;
	let (ns, db) = db_ids(&ds, "test", "tenant").await;

	// The bound a pass would start from under the one-nanosecond retention.
	let bound = watermark_now(&ds, std::time::Duration::from_nanos(1)).await;

	// Unchanged policy: the watermark a page recomputes still reaches the bound,
	// so the delete may proceed.
	let txn = ds.transaction(Write, Optimistic).await.unwrap();
	assert!(
		retention_still_reaches(&txn, ns, db, &bound).await.unwrap(),
		"an unchanged retention must keep authorising the range it made stale"
	);
	let _ = txn.cancel().await;

	// Extending retention to an hour moves the watermark an hour earlier, so
	// everything the old bound named is inside the window again.
	let ses = Session::owner().with_ns("test").with_db("tenant");
	ds.execute("DEFINE TABLE OVERWRITE thing CHANGEFEED 1h;", &ses, None).await.unwrap()[0]
		.result
		.as_ref()
		.unwrap();
	let txn = ds.transaction(Write, Optimistic).await.unwrap();
	assert!(
		!retention_still_reaches(&txn, ns, db, &bound).await.unwrap(),
		"an extended retention must withdraw the bound a pass started from"
	);
	let _ = txn.cancel().await;
}

/// A database that loses its retention, or goes away entirely, withdraws the
/// bound too: the first has nothing left to collect, and the second is reclaimed
/// by its own tombstone rather than page by page here.
#[tokio::test]
async fn a_removed_retention_withdraws_the_stale_bound() {
	let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
	ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
	with_changefeed(&ds, "test", "tenant").await;
	let (ns, db) = db_ids(&ds, "test", "tenant").await;
	let bound = watermark_now(&ds, std::time::Duration::from_nanos(1)).await;

	let ses = Session::owner().with_ns("test").with_db("tenant");
	ds.execute("DEFINE TABLE OVERWRITE thing;", &ses, None).await.unwrap()[0]
		.result
		.as_ref()
		.unwrap();
	let txn = ds.transaction(Write, Optimistic).await.unwrap();
	assert!(
		!retention_still_reaches(&txn, ns, db, &bound).await.unwrap(),
		"a database with no retention left has nothing this collector may delete"
	);
	let _ = txn.cancel().await;

	ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap()[0].result.as_ref().unwrap();
	let txn = ds.transaction(Write, Optimistic).await.unwrap();
	assert!(
		!retention_still_reaches(&txn, ns, db, &bound).await.unwrap(),
		"a removed database is reclaimed by its tombstone, not collected here"
	);
	let _ = txn.cancel().await;
}

/// A pass that loses the lease reports it, so the caller does not start the
/// live-query collector on the same lease.
#[tokio::test]
async fn a_pass_that_loses_the_lease_reports_it() {
	let page = CHANGEFEED_GC_BATCH_SIZE as usize;
	let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
	ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
	with_changefeed(&ds, "test", "tenant").await;
	let feed = feed(&ds, "test", "tenant").await;
	grow_backlog(&ds, &feed, page).await;
	let backlog = count_range(&ds, feed.clone()).await;
	assert!(backlog > 0, "there must be a backlog for the pass to page through");

	// Another node takes the changefeed-cleanup lease.
	let held = LeaseHandler::new(
		ds.sequences().clone(),
		uuid::Uuid::now_v7(),
		ds.transaction_factory().clone(),
		TaskLeaseType::ChangeFeedCleanup,
		std::time::Duration::from_secs(60),
	)
	.unwrap();
	assert!(held.has_lease().await.unwrap(), "the other node must take the lease");

	// This node's pass therefore has no lease to maintain.
	let ours = LeaseHandler::new(
		ds.sequences().clone(),
		uuid::Uuid::now_v7(),
		ds.transaction_factory().clone(),
		TaskLeaseType::ChangeFeedCleanup,
		std::time::Duration::from_secs(60),
	)
	.unwrap();
	let mut budget = u64::MAX;
	let held_after =
		crate::cf::gc::gc_all_at(&ds, &ours, &CancellationToken::new(), &mut budget).await.unwrap();

	assert!(!held_after, "a pass that lost the lease must report the loss to its caller");
	assert_eq!(
		count_range(&ds, feed).await,
		backlog,
		"a pass without the lease must not collect anything"
	);
}

/// A pass that keeps the lease reports that too, so the caller goes on to the
/// live-query collector.
#[tokio::test]
async fn a_completed_pass_reports_the_lease_still_held() {
	let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
	ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
	with_changefeed(&ds, "test", "tenant").await;
	let feed = feed(&ds, "test", "tenant").await;
	grow_backlog(&ds, &feed, CHANGEFEED_GC_BATCH_SIZE as usize).await;

	let ours = LeaseHandler::new(
		ds.sequences().clone(),
		uuid::Uuid::now_v7(),
		ds.transaction_factory().clone(),
		TaskLeaseType::ChangeFeedCleanup,
		std::time::Duration::from_secs(60),
	)
	.unwrap();
	assert!(ours.has_lease().await.unwrap(), "this node must hold the lease");
	let mut budget = u64::MAX;
	assert!(
		crate::cf::gc::gc_all_at(&ds, &ours, &CancellationToken::new(), &mut budget).await.unwrap(),
		"a pass that kept the lease must report it, so the caller runs the next collector"
	);
	assert_eq!(count_range(&ds, feed).await, 0, "the pass must collect the backlog");
}

/// Every statement that writes a changefeed retention bumps the fence.
///
/// The fence is the only key a collection page and a retention change both
/// touch, so a path that sets a retention without bumping it leaves that page
/// free to go on deleting under the watermark the old policy produced. Dropping a
/// retention is not bumped: it can only lower the maximum, which strands nothing.
#[tokio::test]
async fn every_retention_write_bumps_the_fence() {
	let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
	let ses = Session::owner().with_ns("test").with_db("tenant");
	ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();

	// DEFINE DATABASE with a retention.
	ds.execute("DEFINE DATABASE tenant CHANGEFEED 1h;", &ses, None).await.unwrap()[0]
		.result
		.as_ref()
		.unwrap();
	let (ns, db) = db_ids(&ds, "test", "tenant").await;
	let after_db = fence(&ds, ns, db).await.expect("DEFINE DATABASE must bump the fence");

	// DEFINE TABLE with a retention.
	ds.execute("DEFINE TABLE thing CHANGEFEED 1h;", &ses, None).await.unwrap()[0]
		.result
		.as_ref()
		.unwrap();
	let after_tb = fence(&ds, ns, db).await.unwrap();
	assert!(after_tb > after_db, "DEFINE TABLE must bump the fence: {after_db} -> {after_tb}");

	// ALTER TABLE setting a retention, on a table that already had one.
	ds.execute("ALTER TABLE thing CHANGEFEED 2h;", &ses, None).await.unwrap()[0]
		.result
		.as_ref()
		.unwrap();
	let after_alter = fence(&ds, ns, db).await.unwrap();
	assert!(after_alter > after_tb, "ALTER TABLE must bump the fence: {after_tb} -> {after_alter}");

	// ALTER TABLE setting a retention on a table that had none — raises the
	// database's maximum just as much, and is a different branch from the replace
	// above.
	ds.execute("DEFINE TABLE other;", &ses, None).await.unwrap()[0].result.as_ref().unwrap();
	let before_fresh = fence(&ds, ns, db).await.unwrap();
	ds.execute("ALTER TABLE other CHANGEFEED 3h;", &ses, None).await.unwrap()[0]
		.result
		.as_ref()
		.unwrap();
	let after_fresh = fence(&ds, ns, db).await.unwrap();
	assert!(
		after_fresh > before_fresh,
		"setting a retention on a table that had none must bump the fence: \
		 {before_fresh} -> {after_fresh}"
	);

	// Dropping one does not: it can only lower the maximum.
	ds.execute("ALTER TABLE other DROP CHANGEFEED;", &ses, None).await.unwrap()[0]
		.result
		.as_ref()
		.unwrap();
	assert_eq!(
		fence(&ds, ns, db).await.unwrap(),
		after_fresh,
		"dropping a retention needs no fence bump: it strands nothing"
	);
}

/// A collection page arms against the fence, so a retention change committed
/// while the page is open rejects the page instead of being ignored.
///
/// This is what a plain catalog read cannot do. The page reads catalog keys and
/// writes changefeed keys, so on a last-writer-wins backend the two transactions
/// share no key and the page commits regardless; arming the fence gives them one.
#[tokio::test]
async fn a_retention_change_rejects_an_open_collection_page() {
	let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
	ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
	with_changefeed(&ds, "test", "tenant").await;
	let (ns, db) = db_ids(&ds, "test", "tenant").await;
	let bound = watermark_now(&ds, std::time::Duration::from_nanos(1)).await;

	// A page transaction that has revalidated the bound and so holds the fence.
	let page = ds.transaction(Write, Optimistic).await.unwrap();
	assert!(
		retention_still_reaches(&page, ns, db, &bound).await.unwrap(),
		"the retention must still authorise the range before the change lands"
	);

	// The retention changes while that page is open.
	let ses = Session::owner().with_ns("test").with_db("tenant");
	ds.execute("ALTER TABLE thing CHANGEFEED 1h;", &ses, None).await.unwrap()[0]
		.result
		.as_ref()
		.unwrap();

	// The page must not be allowed to land: it would delete entries the committed
	// policy retains.
	page.set(&b"unused".to_vec(), &vec![1u8]).await.unwrap();
	page.commit()
		.await
		.expect_err("a page open across a retention change must be rejected at commit");
}