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
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
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
1163
1164
1165
use core::future::Future;
use std::pin::Pin;
use std::sync::{Arc, Weak};
use std::time::Duration;

use common::time::{Instant, MissedTickBehavior, sleep};
use futures::StreamExt;
use surrealdb_datastore::triggers::CommitTriggers;
use surrealdb_types::Error;
#[cfg(not(target_family = "wasm"))]
use tokio::spawn;
use tokio_util::sync::CancellationToken;
#[cfg(target_family = "wasm")]
use wasm_bindgen_futures::spawn_local as spawn;

use crate::err::{is_query_cancelled, is_query_timedout};
use crate::kvs::{Datastore, LiveQueryEngine};
use crate::observe::process::RefreshClaim;
use crate::options::EngineOptions;

mod interval;

use self::interval::IntervalStream;

#[cfg(not(target_family = "wasm"))]
type Task = Pin<Box<dyn Future<Output = Result<(), tokio::task::JoinError>> + Send + 'static>>;

#[cfg(target_family = "wasm")]
type Task = Pin<Box<()>>;

/// Spawns `fut` on the ambient runtime and returns a handle the caller can await
/// to observe the task finishing.
///
/// Awaiting the returned [`Task`] is only meaningful off wasm. `spawn_local`
/// yields no join handle, so on wasm the task is detached and the returned
/// `Task` completes immediately without waiting for `fut`.
#[cfg(not(target_family = "wasm"))]
fn into_task<F>(fut: F) -> Task
where
	F: Future<Output = ()> + Send + 'static,
{
	Box::pin(spawn(fut))
}

#[cfg(target_family = "wasm")]
fn into_task<F>(fut: F) -> Task
where
	F: Future<Output = ()> + 'static,
{
	spawn(fut);
	Box::pin(())
}

const NODE_MEMBERSHIP_UPDATE_TIMEOUT: Duration = Duration::from_secs(60);

/// How long a trigger-driven index-compaction pass waits before running.
///
/// The compaction interval is the floor on how *stale* the queue may get; this
/// is the floor on how *often* a write may start a pass. It bounds the cost of
/// the wake-up path — most importantly the lease check on nodes that do not
/// hold the compaction lease, which returns almost immediately and would
/// otherwise spin for as long as writes keep arriving — and batches the commits
/// that land during the wait into a single pass.
const INDEX_COMPACTION_TRIGGER_DEBOUNCE: Duration = Duration::from_millis(250);

/// How long the graph fold task waits after a commit-time wake-up before
/// running, for the same batching reason as the compaction debounce above.
const GRAPH_FOLD_TRIGGER_DEBOUNCE: Duration = Duration::from_millis(250);

enum NodeMembershipUpdateResult {
	Updated,
	Cancelled,
	TimedOut,
	Failed(anyhow::Error),
}

pub struct Tasks(#[cfg_attr(target_family = "wasm", expect(dead_code))] Vec<Task>);

impl Tasks {
	#[cfg(target_family = "wasm")]
	pub async fn resolve(self) -> Result<(), Error> {
		Ok(())
	}
	#[cfg(not(target_family = "wasm"))]
	pub async fn resolve(self) -> Result<(), Error> {
		for task in self.0 {
			// Surface a task that panicked or was aborted. The maintenance
			// scheduler carries several jobs, so losing it silently would stop
			// all of them while shutdown still reported success.
			if let Err(e) = task.await {
				error!("Background task did not shut down cleanly: {e}");
			}
		}
		Ok(())
	}
}

/// Starts this node's background tasks and returns handles to await at shutdown.
///
/// Work is split across tasks by what its cadence has to guarantee, not one task
/// per job:
///
/// - The **cluster heartbeat** runs alone. A stalled heartbeat gets this node archived by another
///   member's expiry scan and fails its readiness probe, so nothing may ever share its task.
/// - **Async event processing** and **index compaction** each keep their own task because both are
///   driven by the write path and drain their queue to empty, so under sustained load they do not
///   return between ticks.
/// - The remaining jobs run on two [`spawn_task_scheduler`] instances, split by whether one pass
///   has a bound. [`maintenance_slots`] carries the jobs whose cost is bounded by catalog size, so
///   none can delay a peer by more than one short pass. [`sweep_slots`] carries the two whose queue
///   is nearly always empty but whose individual entries are not bounded: reclaiming one tombstone
///   destroys an entire namespace or database prefix, and a session purge pages the whole session
///   keyspace opening a write transaction per expired entry. Keeping the groups apart means a long
///   sweep cannot stall the short jobs — in particular it cannot leave the cached process metrics
///   stale for its duration.
/// - The **live-query router** is spawned only under [`LiveQueryEngine::Router`]; under the default
///   inline engine it has nothing to deliver.
///
/// Every task holds a [`Weak`] reference to the datastore, never a strong one.
/// The datastore owns these handles, so a strong reference would close a cycle
/// through them and the datastore could never be dropped. Each pass upgrades
/// for the duration of that pass and exits the task once the upgrade fails,
/// which is what lets an embedder that simply drops its datastore — rather than
/// calling [`Datastore::shutdown`] — still wind the tasks down.
///
/// The datastore starts these itself, so the only reason to call this directly
/// is to start them against a datastore built with
/// [`Builder::without_maintenance_tasks`](crate::kvs::ds::builder::Builder::without_maintenance_tasks).
pub fn init(dbs: &Arc<Datastore>, canceller: CancellationToken, opts: &EngineOptions) -> Tasks {
	let weak = Arc::downgrade(dbs);
	// The triggers are shared state in their own right, so a task holding them
	// does not keep the datastore alive.
	let triggers = dbs.commit_triggers();
	let mut tasks = Vec::with_capacity(6);
	tasks.push(spawn_task_node_membership_refresh(Weak::clone(&weak), canceller.clone(), opts));
	tasks.push(spawn_task_event_processing(
		Weak::clone(&weak),
		Arc::clone(triggers),
		canceller.clone(),
		opts,
	));
	tasks.push(spawn_task_index_compaction(
		Weak::clone(&weak),
		Arc::clone(triggers),
		canceller.clone(),
		opts,
	));
	tasks.push(spawn_task_graph_fold(
		Weak::clone(&weak),
		Arc::clone(triggers),
		canceller.clone(),
		opts,
	));
	for (group, slots) in [("maintenance", maintenance_slots(opts)), ("sweep", sweep_slots(opts))] {
		// Every job in a group can be disabled by interval, so a group can end up
		// empty; spawning a task that immediately exits would serve no purpose.
		if !slots.is_empty() {
			tasks.push(spawn_task_scheduler(
				group,
				slots,
				Weak::clone(&weak),
				canceller.clone(),
				opts,
			));
		}
	}
	if dbs.live_query_engine() == LiveQueryEngine::Router {
		tasks.push(spawn_task_live_query_router(weak, canceller, opts));
	}
	Tasks(tasks)
}

/// Spawns the per-node live-query router task.
///
/// Tails the dedicated `lqe` keyspace and delivers notifications off the write
/// path. The cadence bounds steady-state delivery latency, so it ticks
/// frequently and keeps its own task rather than queueing behind maintenance
/// work. Only spawned under [`LiveQueryEngine::Router`].
fn spawn_task_live_query_router(
	dbs: Weak<Datastore>,
	canceller: CancellationToken,
	opts: &EngineOptions,
) -> Task {
	let interval = opts.live_query_router_interval;
	into_task(async move {
		trace!("Running the live-query router every {interval:?}");
		let mut ticker = interval_ticker(interval).await;
		loop {
			tokio::select! {
				biased;
				_ = canceller.cancelled() => break,
				Some(_) = ticker.next() => {
					let Some(dbs) = dbs.upgrade() else { break };
					if let Err(e) = dbs.live_query_router_process().await {
						error!("Error running the live-query router: {e}");
					}
				}
			}
		}
		trace!("Background task exited: Running the live-query router");
	})
}

fn spawn_task_node_membership_refresh(
	dbs: Weak<Datastore>,
	canceller: CancellationToken,
	opts: &EngineOptions,
) -> Task {
	// Get the delay interval from the config
	let interval = opts.node_membership_refresh_interval;
	// Spawn a future
	into_task(async move {
		// Log the interval frequency
		trace!("Updating node registration information every {interval:?}");
		// Create a new time-based interval ticket
		let mut ticker = interval_ticker(interval).await;
		// Loop continuously until the task is cancelled
		loop {
			tokio::select! {
				biased;
				// Check if this has shutdown
				_ = canceller.cancelled() => break,
				// Receive a notification on the channel
				Some(_) = ticker.next() => {
					let Some(dbs) = dbs.upgrade() else { break };
					if !run_node_membership_update(
						NODE_MEMBERSHIP_UPDATE_TIMEOUT,
						update_node_membership(
							&dbs,
							&canceller,
							NODE_MEMBERSHIP_UPDATE_TIMEOUT,
						),
					).await {
						break;
					}
				}
			}
		}
		trace!("Background task exited: Updating node registration information");
	})
}

/// Spawns a background task for index compaction
///
/// This function creates a background task that periodically runs the index
/// compaction process. The compaction process optimizes indexes (particularly
/// full-text indexes) by consolidating changes and removing unnecessary data,
/// which helps maintain query performance over time.
///
/// The task runs at the interval specified by `opts.index_compaction_interval`.
/// It keeps its own task rather than joining the maintenance scheduler because
/// the queue is fed by the write path and each pass drains it to empty, so under
/// sustained indexed writes a pass does not return between ticks.
///
/// # Arguments
///
/// * `dbs` - The datastore instance
/// * `triggers` - The commit wake-ups, so a write can start a pass early
/// * `canceller` - Token used to cancel the task when the engine is shutting down
/// * `opts` - Engine options containing the compaction interval
///
/// # Returns
///
/// * A pinned task that can be awaited
fn spawn_task_index_compaction(
	dbs: Weak<Datastore>,
	triggers: Arc<CommitTriggers>,
	canceller: CancellationToken,
	opts: &EngineOptions,
) -> Task {
	// Get the delay interval from the config
	let interval = opts.index_compaction_interval;
	// Spawn a future
	into_task(async move {
		// Log the interval frequency
		trace!("Running index compaction every {interval:?}");
		// Create a new time-based interval ticket
		let mut ticker = interval_ticker(interval).await;
		// Loop continuously until the task is cancelled
		loop {
			tokio::select! {
				biased;
				// Check if this has shutdown
				_ = canceller.cancelled() => break,
				// Wake early when a commit queues compaction work, so the queue
				// drains at the write rate rather than at this task's cadence.
				// Nothing otherwise relates the two, and the gap between them is
				// what lets a count index's delta log — which every read of that
				// index sums — grow without bound.
				_ = triggers.index_compaction.notified() => {
					// Debounce before running. A node that does not hold the
					// compaction lease returns from a pass almost immediately,
					// so honouring every notification would spin on the lease
					// check for as long as writes keep arriving. The wait also
					// batches the commits landing in the meantime into one pass.
					tokio::select! {
						biased;
						_ = canceller.cancelled() => break,
						_ = sleep(INDEX_COMPACTION_TRIGGER_DEBOUNCE) => {}
					}
					let Some(dbs) = dbs.upgrade() else { break };
					if let Err(e) =
						Datastore::index_compaction(dbs, interval, canceller.clone()).await
					{
						if canceller.is_cancelled() {
							break;
						}
						error!("Error running index compaction: {e}");
					}
				}
				// Receive a notification on the channel
				Some(_) = ticker.next() => {
					let Some(dbs) = dbs.upgrade() else { break };
					if let Err(e) =
						Datastore::index_compaction(dbs, interval, canceller.clone()).await
					{
						if canceller.is_cancelled() {
							break;
						}
						error!("Error running index compaction: {e}");
					}
				}
			}
		}
		trace!("Background task exited: Running index compaction");
	})
}

/// Spawns the graph adjacency fold task.
///
/// Runs at `opts.graph_fold_interval` and wakes early when a commit queues
/// fold work or a read observes a hot scope. Keeps its own task rather than
/// joining the maintenance scheduler for the same reason as index
/// compaction: the queue is fed by the write and read paths, and each pass
/// drains it to empty. The pass itself is a no-op while the fold threshold
/// is 0.
fn spawn_task_graph_fold(
	dbs: Weak<Datastore>,
	triggers: Arc<CommitTriggers>,
	canceller: CancellationToken,
	opts: &EngineOptions,
) -> Task {
	let interval = opts.graph_fold_interval;
	into_task(async move {
		trace!("Running graph fold every {interval:?}");
		let mut ticker = interval_ticker(interval).await;
		loop {
			tokio::select! {
				biased;
				_ = canceller.cancelled() => break,
				_ = triggers.graph_fold.notified() => {
					tokio::select! {
						biased;
						_ = canceller.cancelled() => break,
						_ = sleep(GRAPH_FOLD_TRIGGER_DEBOUNCE) => {}
					}
					let Some(dbs) = dbs.upgrade() else { break };
					if let Err(e) = Datastore::graph_fold(dbs, interval, canceller.clone()).await {
						if canceller.is_cancelled() {
							break;
						}
						error!("Error running graph fold: {e}");
					}
				}
				Some(_) = ticker.next() => {
					let Some(dbs) = dbs.upgrade() else { break };
					if let Err(e) = Datastore::graph_fold(dbs, interval, canceller.clone()).await {
						if canceller.is_cancelled() {
							break;
						}
						error!("Error running graph fold: {e}");
					}
				}
			}
		}
		trace!("Background task exited: Running graph fold");
	})
}

/// Spawns the async event processing task.
///
/// Keeps its own task rather than joining the maintenance scheduler for two
/// reasons: the queue is fed by the write path and each pass drains it to empty,
/// and the events themselves run user-defined SurrealQL of unbounded duration.
fn spawn_task_event_processing(
	dbs: Weak<Datastore>,
	triggers: Arc<CommitTriggers>,
	canceller: CancellationToken,
	opts: &EngineOptions,
) -> Task {
	// Get the delay interval from the config
	let interval = opts.event_processing_interval;
	// Spawn a future
	into_task(async move {
		// Log the interval frequency
		trace!("Running event processing every {interval:?}");
		// Create a new time-based interval ticket
		let mut ticker = interval_ticker(interval).await;
		// Reports whether the datastore is still alive; a `false` ends the task.
		let process_events = async || {
			let Some(dbs) = dbs.upgrade() else {
				return false;
			};
			// The pass stops at the next batch boundary once cancelled, so a
			// shutdown does not wait out a queue the write path keeps refilling.
			if let Err(e) = dbs.event_processing(interval, &canceller).await
				&& !canceller.is_cancelled()
			{
				error!("Error running event processing: {e}");
			}
			true
		};
		// Loop continuously until the task is cancelled
		loop {
			tokio::select! {
				biased;
				// Check if this has shutdown
				_ = canceller.cancelled() => break,
				// Wake early when new async events are committed.
				_ = triggers.async_event.notified() => if !process_events().await { break },
				// Receive a notification on the channel
				Some(_) = ticker.next() => if !process_events().await { break }
			}
		}
		trace!("Background task exited: Running event processing");
	})
}

// --------------------------------------------------
// Maintenance scheduler
// --------------------------------------------------

/// A periodic maintenance job multiplexed onto a shared scheduler task.
///
/// Every variant runs on a cadence of seconds to minutes. They are split across
/// two schedulers by whether a single pass is bounded — [`maintenance_slots`]
/// versus [`sweep_slots`] — so an unbounded sweep cannot stall the short jobs.
/// Work that cannot share a task at all — the cluster heartbeat, async event
/// processing, index compaction and the live-query router — keeps its own; see
/// [`init`].
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum MaintenanceJob {
	/// Archive cluster members whose heartbeat has gone stale.
	NodeExpire,
	/// Delete archived members and garbage-collect their live queries.
	NodeCleanup,
	/// Garbage-collect expired changefeed data.
	ChangefeedGc,
	/// Destroy data left behind by `REMOVE NAMESPACE/DATABASE/INDEX`.
	ReclaimTombstones,
	/// Adopt index builds stranded by a crashed or expired owner node.
	ResumeIndexBuilds,
	/// Purge expired durable RPC sessions.
	RpcSessionGc,
	/// Advance the TiKV MVCC garbage-collection safepoint.
	TikvGc,
	/// Resolve stale TiKV transactional locks.
	TikvLockCleanup,
	/// Refresh the cached process/system utilisation metrics.
	SystemMetricsRefresh,
}

impl MaintenanceJob {
	/// Every variant, used to assert the two schedules partition the job set.
	/// Must list all of them; a new variant belongs here and in exactly one of
	/// [`maintenance_slots`] / [`sweep_slots`].
	#[cfg(test)]
	const ALL: [Self; 9] = [
		Self::NodeExpire,
		Self::NodeCleanup,
		Self::ChangefeedGc,
		Self::ReclaimTombstones,
		Self::ResumeIndexBuilds,
		Self::RpcSessionGc,
		Self::TikvGc,
		Self::TikvLockCleanup,
		Self::SystemMetricsRefresh,
	];

	/// Description used in the scheduler's log lines.
	fn label(self) -> &'static str {
		match self {
			Self::NodeExpire => "inactive node expiry",
			Self::NodeCleanup => "archived node cleanup",
			Self::ChangefeedGc => "changefeed garbage collection",
			Self::ReclaimTombstones => "tombstone reclaim",
			Self::ResumeIndexBuilds => "stalled index build recovery",
			Self::RpcSessionGc => "expired RPC session purge",
			Self::TikvGc => "TiKV MVCC GC",
			Self::TikvLockCleanup => "TiKV lock cleanup",
			Self::SystemMetricsRefresh => "system metrics refresh",
		}
	}
}

/// One job's place in the schedule.
struct Slot {
	job: MaintenanceJob,
	interval: Duration,
	next_due: Instant,
}

/// The jobs whose cost per pass is bounded up front.
///
/// Each of these either reads a fixed number of catalog rows, or bounds its own
/// pass explicitly: `changefeed_process` deletes at most one key budget per pass,
/// shared across the databases it visits, however far the retention backlog has
/// run ahead of it. So they can share one task without delaying a peer by more
/// than one short pass, and a job with more work than one pass allows continues
/// on the next tick rather than holding the task.
fn maintenance_slots(opts: &EngineOptions) -> Vec<Slot> {
	slots(&[
		(MaintenanceJob::SystemMetricsRefresh, opts.system_metrics_refresh_interval),
		(MaintenanceJob::NodeExpire, opts.node_membership_check_interval),
		(MaintenanceJob::NodeCleanup, opts.node_membership_cleanup_interval),
		(MaintenanceJob::ChangefeedGc, opts.changefeed_gc_interval),
		(MaintenanceJob::ResumeIndexBuilds, opts.index_build_resume_interval),
		(MaintenanceJob::TikvGc, opts.tikv_gc_interval),
		(MaintenanceJob::TikvLockCleanup, opts.tikv_lock_cleanup_interval),
	])
}

/// The jobs whose queue is nearly always empty but whose cost per pass is not
/// bounded by catalog size.
///
/// `purge_expired_rpc_sessions` pages the entire session keyspace, opening a
/// write transaction for every expired entry, so one pass is arbitrarily long on
/// a busy durable-session deployment. `reclaim_tombstones` bounds its own pass
/// by a key budget, but the queue it drains is filled by `REMOVE` statements
/// rather than by catalog size, and a large removal takes many passes to finish.
/// Both are therefore a no-op scan in the steady state and neither belongs
/// alongside the jobs whose pass length follows the catalog; they share a task
/// with each other instead.
fn sweep_slots(opts: &EngineOptions) -> Vec<Slot> {
	slots(&[
		(MaintenanceJob::ReclaimTombstones, opts.reclaim_interval),
		(MaintenanceJob::RpcSessionGc, opts.rpc_session_gc_interval),
	])
}

/// Builds a schedule from the given job/interval pairs.
///
/// A zero interval leaves the job unregistered, which is how the documented
/// "set to zero to disable" options are honoured. Zero is not a valid tick
/// period for any job, so treating it uniformly as "disabled" also avoids
/// turning a misconfigured interval into a busy loop.
fn slots(jobs: &[(MaintenanceJob, Duration)]) -> Vec<Slot> {
	// One `now` for the whole schedule, so jobs sharing an interval share a
	// deadline exactly and registration order — not nanosecond skew between
	// clock reads — decides which of them runs first.
	let now = Instant::now();
	jobs.iter()
		.copied()
		.filter(|(_, interval)| !interval.is_zero())
		.map(|(job, interval)| Slot {
			job,
			interval,
			// The metrics refresh is due immediately so a cold start reports real
			// utilisation rather than zeroes. Every other job waits out one full
			// interval, which keeps startup from firing several lease acquisitions
			// at the same moment.
			next_due: match job {
				MaintenanceJob::SystemMetricsRefresh => now,
				_ => now + interval,
			},
		})
		.collect()
}

/// Index of the slot to run next: the earliest deadline, and among equal
/// deadlines the one registered first (`min_by_key` yields the first minimum).
///
/// An overdue job sorts ahead of one whose deadline was just set past the
/// current instant, so a saturated schedule drains oldest-first and a job that
/// has just run cannot immediately run again while a peer is waiting.
fn next_slot(slots: &[Slot]) -> Option<usize> {
	slots.iter().enumerate().min_by_key(|(_, s)| s.next_due).map(|(i, _)| i)
}

/// Runs one pass of `job`, logging and swallowing its error so that a failing
/// job cannot stop the others. Cancellation is not an error.
async fn run_maintenance_job(
	job: MaintenanceJob,
	dbs: &Arc<Datastore>,
	canceller: &CancellationToken,
	opts: &EngineOptions,
	reclaim_grace: Duration,
	metrics_claim: &RefreshClaim,
) {
	let res = match job {
		MaintenanceJob::NodeExpire => dbs.expire_nodes().await,
		MaintenanceJob::NodeCleanup => dbs.remove_nodes().await,
		MaintenanceJob::ChangefeedGc => {
			dbs.changefeed_process(&opts.changefeed_gc_interval, canceller).await
		}
		MaintenanceJob::ReclaimTombstones => Datastore::reclaim_tombstones(
			Arc::clone(dbs),
			opts.reclaim_interval,
			reclaim_grace,
			canceller.clone(),
		)
		.await
		.map(|_| ()),
		MaintenanceJob::ResumeIndexBuilds => dbs
			.resume_stalled_index_builds(opts.index_build_resume_interval, canceller.clone())
			.await
			.map(|_| ()),
		MaintenanceJob::RpcSessionGc => {
			dbs.purge_expired_rpc_sessions(&opts.rpc_session_gc_interval).await
		}
		MaintenanceJob::TikvGc => dbs.run_mvcc_gc(opts.tikv_gc_lifetime).await,
		MaintenanceJob::TikvLockCleanup => dbs.run_lock_cleanup(opts.tikv_gc_lifetime).await,
		MaintenanceJob::SystemMetricsRefresh => {
			// The snapshot is one cache per process, and its CPU percentage is a
			// delta since the previous refresh of it, so only the datastore
			// holding the process-wide claim refreshes; the others let their
			// pass go by. Retried every pass, so the refresh moves to a
			// surviving datastore when the holder's task ends.
			if metrics_claim.take() {
				crate::observe::refresh_process_snapshot().await;
			}
			Ok(())
		}
	};
	if let Err(e) = res
		&& !canceller.is_cancelled()
	{
		error!("Error running {}: {e}", job.label());
	}
}

/// Drives the schedule until cancelled, running one job per iteration.
///
/// Dispatch is injected so the scheduling behaviour can be exercised without a
/// datastore.
async fn maintenance_loop<F, Fut>(mut slots: Vec<Slot>, canceller: CancellationToken, run: F)
where
	F: Fn(MaintenanceJob) -> Fut,
	Fut: Future<Output = ()>,
{
	while let Some(i) = next_slot(&slots) {
		let delay = slots[i].next_due.saturating_duration_since(Instant::now());
		tokio::select! {
			biased;
			_ = canceller.cancelled() => break,
			_ = sleep(delay) => {}
		}
		run(slots[i].job).await;
		// Measure the next deadline from completion, not from the deadline just
		// met. A pass that overruns its interval therefore still rests for a
		// full interval instead of coming due the moment it returns, and its
		// deadline cannot drift permanently into the past.
		slots[i].next_due = Instant::now() + slots[i].interval;
	}
}

/// Spawns one task that runs the given group of jobs on a shared schedule.
///
/// `group` names the group in the task's log lines so an operator can tell the
/// two schedulers apart.
fn spawn_task_scheduler(
	group: &'static str,
	slots: Vec<Slot>,
	dbs: Weak<Datastore>,
	canceller: CancellationToken,
	opts: &EngineOptions,
) -> Task {
	let opts = *opts;
	// Clamp the grace up to at least the TiKV GC lifetime: the reclaim's
	// out-of-transaction range destroy bypasses MVCC, so data must not be
	// reclaimed while a snapshot older than the GC safepoint
	// (`now - tikv_gc_lifetime`) could still read it. Deriving the effective
	// grace here means a longer `--tikv-gc-lifetime` can never be undercut by
	// leaving `--reclaim-grace` at its default.
	let reclaim_grace = opts.reclaim_grace.max(opts.tikv_gc_lifetime);
	into_task(async move {
		trace!(
			"Running {} {group} jobs on a shared schedule: {}",
			slots.len(),
			slots.iter().map(|s| s.job.label()).collect::<Vec<_>>().join(", ")
		);
		let jobs_canceller = canceller.clone();
		// This task's holder of the process-wide metrics claim: taken by its
		// first refresh pass, and released when the task ends or is dropped so
		// another datastore can take over. Only the group carrying that job
		// ever takes it.
		let metrics_claim = RefreshClaim::process();
		maintenance_loop(slots, canceller, move |job| {
			let dbs = Weak::clone(&dbs);
			let canceller = jobs_canceller.clone();
			let metrics_claim = metrics_claim.clone();
			async move {
				// The datastore is gone, so there is nothing left to maintain.
				// Cancelling ends the schedule on its next iteration, and does the
				// same for every other task sharing this token.
				let Some(dbs) = dbs.upgrade() else {
					canceller.cancel();
					return;
				};
				run_maintenance_job(job, &dbs, &canceller, &opts, reclaim_grace, &metrics_claim)
					.await
			}
		})
		.await;
		trace!("Background task exited: Running {group} jobs");
	})
}

async fn update_node_membership(
	dbs: &Datastore,
	canceller: &CancellationToken,
	timeout_duration: Duration,
) -> NodeMembershipUpdateResult {
	match dbs.update_node_with_timeout(timeout_duration, canceller).await {
		Ok(()) => NodeMembershipUpdateResult::Updated,
		Err(e) if is_query_cancelled(&e) => NodeMembershipUpdateResult::Cancelled,
		Err(e) if is_query_timedout(&e) => NodeMembershipUpdateResult::TimedOut,
		Err(e) => NodeMembershipUpdateResult::Failed(e),
	}
}

async fn run_node_membership_update<Fut>(timeout_duration: Duration, update_node: Fut) -> bool
where
	Fut: Future<Output = NodeMembershipUpdateResult>,
{
	match update_node.await {
		NodeMembershipUpdateResult::Updated => true,
		NodeMembershipUpdateResult::Cancelled => false,
		NodeMembershipUpdateResult::TimedOut => {
			warn!("Timed out updating node registration information after {timeout_duration:?}");
			true
		}
		NodeMembershipUpdateResult::Failed(e) => {
			error!("Error updating node registration information: {e}");
			true
		}
	}
}

async fn interval_ticker(interval: Duration) -> IntervalStream {
	// Create a new interval timer
	let mut interval = common::time::interval(interval);
	// Don't bombard the database if we miss some ticks
	interval.set_missed_tick_behavior(MissedTickBehavior::Delay);
	interval.tick().await;
	IntervalStream::new(interval)
}

#[cfg(test)]
mod test {
	use std::sync::{Arc, Mutex};
	use std::time::Duration;

	use tokio::time::Instant;
	use tokio_util::sync::CancellationToken;

	// Used only by the tests that build a datastore, so it is gated with them:
	// the lib tests build under `-D warnings` once per kv backend.
	#[cfg(feature = "kv-mem")]
	use super::RefreshClaim;
	use super::{
		MaintenanceJob, Slot, maintenance_loop, maintenance_slots, next_slot, sweep_slots,
	};
	#[cfg(feature = "kv-mem")]
	use crate::kvs::Datastore;
	#[cfg(feature = "kv-mem")]
	use crate::kvs::tasks;
	use crate::options::EngineOptions;

	/// A slot due one interval from now, as the scheduler registers them.
	fn slot(job: MaintenanceJob, interval: Duration) -> Slot {
		Slot {
			job,
			interval,
			next_due: Instant::now() + interval,
		}
	}

	/// Drives `maintenance_loop` with a recorder that stops it after `limit`
	/// passes, then returns the jobs in the order they ran.
	async fn record_passes(slots: Vec<Slot>, limit: usize) -> Vec<MaintenanceJob> {
		record_passes_with_delay(slots, limit, |_| Duration::ZERO)
			.await
			.into_iter()
			.map(|(job, _)| job)
			.collect()
	}

	/// As [`record_passes`], but each pass takes the duration `cost` reports for
	/// its job, and each pass is paired with the instant it started, so a slow
	/// job's effect on the schedule can be observed.
	async fn record_passes_with_delay(
		slots: Vec<Slot>,
		limit: usize,
		cost: impl Fn(MaintenanceJob) -> Duration,
	) -> Vec<(MaintenanceJob, Instant)> {
		let canceller = CancellationToken::new();
		let log = Arc::new(Mutex::new(Vec::new()));
		let stop = canceller.clone();
		let sink = Arc::clone(&log);
		maintenance_loop(slots, canceller, move |job| {
			let sink = Arc::clone(&sink);
			let stop = stop.clone();
			let delay = cost(job);
			async move {
				{
					let mut passes = sink.lock().unwrap();
					passes.push((job, Instant::now()));
					if passes.len() >= limit {
						stop.cancel();
					}
				}
				if !delay.is_zero() {
					tokio::time::sleep(delay).await;
				}
			}
		})
		.await;
		Arc::into_inner(log).unwrap().into_inner().unwrap()
	}

	#[test]
	fn next_slot_is_none_when_nothing_is_registered() {
		assert!(next_slot(&[]).is_none());
	}

	#[test]
	fn next_slot_picks_the_earliest_deadline() {
		let slots = vec![
			slot(MaintenanceJob::NodeCleanup, Duration::from_secs(300)),
			slot(MaintenanceJob::NodeExpire, Duration::from_secs(15)),
			slot(MaintenanceJob::ChangefeedGc, Duration::from_secs(30)),
		];
		assert_eq!(next_slot(&slots), Some(1));
	}

	#[test]
	fn next_slot_breaks_deadline_ties_on_registration_order() {
		let due = Instant::now() + Duration::from_secs(60);
		let mut slots = vec![
			slot(MaintenanceJob::RpcSessionGc, Duration::from_secs(60)),
			slot(MaintenanceJob::ReclaimTombstones, Duration::from_secs(60)),
		];
		for s in &mut slots {
			s.next_due = due;
		}
		assert_eq!(next_slot(&slots), Some(0));
	}

	#[test]
	fn next_slot_prefers_an_overdue_job_over_one_just_run() {
		let mut slots = vec![
			// Just ran, so its deadline sits in the future.
			slot(MaintenanceJob::ChangefeedGc, Duration::from_secs(30)),
			// Overdue.
			slot(MaintenanceJob::NodeExpire, Duration::from_secs(15)),
		];
		slots[1].next_due = Instant::now() - Duration::from_secs(5);
		assert_eq!(next_slot(&slots), Some(1));
	}

	#[test]
	fn slots_skip_zero_intervals() {
		let opts = EngineOptions::default()
			.with_tikv_gc_interval(Duration::ZERO)
			.with_rpc_session_gc_interval(Duration::ZERO);
		let maintenance: Vec<_> = maintenance_slots(&opts).into_iter().map(|s| s.job).collect();
		let sweeps: Vec<_> = sweep_slots(&opts).into_iter().map(|s| s.job).collect();
		assert!(!maintenance.contains(&MaintenanceJob::TikvGc));
		assert!(!sweeps.contains(&MaintenanceJob::RpcSessionGc));
		// The rest of each schedule is unaffected.
		assert!(maintenance.contains(&MaintenanceJob::NodeExpire));
		assert!(maintenance.contains(&MaintenanceJob::TikvLockCleanup));
		assert!(sweeps.contains(&MaintenanceJob::ReclaimTombstones));
	}

	/// The two groups must partition the job set: a job in neither would never
	/// run, and a job in both would run twice per interval and contend with
	/// itself for its own lease.
	#[test]
	fn the_two_groups_partition_every_job() {
		let opts = EngineOptions::default();
		let mut scheduled: Vec<_> =
			maintenance_slots(&opts).into_iter().chain(sweep_slots(&opts)).map(|s| s.job).collect();
		let total = scheduled.len();
		scheduled.dedup();
		assert_eq!(total, scheduled.len(), "a job is registered in both groups");
		for job in MaintenanceJob::ALL {
			assert!(scheduled.contains(&job), "{} is in neither group", job.label());
		}
	}

	/// The unbounded sweeps must not share a task with the bounded jobs: a
	/// reclaim that destroys a whole database, or a purge that pages the session
	/// keyspace, would otherwise stall every short job for its duration.
	#[test]
	fn unbounded_sweeps_are_not_on_the_maintenance_schedule() {
		let opts = EngineOptions::default();
		let maintenance: Vec<_> = maintenance_slots(&opts).into_iter().map(|s| s.job).collect();
		assert!(!maintenance.contains(&MaintenanceJob::ReclaimTombstones));
		assert!(!maintenance.contains(&MaintenanceJob::RpcSessionGc));
		// And the metrics refresh must stay off the sweep schedule, so a long
		// sweep cannot leave `INFO FOR ROOT` and the process gauges stale.
		let sweeps: Vec<_> = sweep_slots(&opts).into_iter().map(|s| s.job).collect();
		assert!(!sweeps.contains(&MaintenanceJob::SystemMetricsRefresh));
	}

	#[test]
	fn maintenance_slots_registers_the_metrics_refresh_immediately() {
		let opts = EngineOptions::default();
		let slots = maintenance_slots(&opts);
		let metrics = slots
			.iter()
			.find(|s| s.job == MaintenanceJob::SystemMetricsRefresh)
			.expect("the metrics refresh is always registered");
		// Due now, so a cold start reports real utilisation rather than zeroes.
		assert!(metrics.next_due <= Instant::now());
		// Every other job waits out one full interval.
		let others = slots.iter().filter(|s| s.job != MaintenanceJob::SystemMetricsRefresh);
		for s in others {
			assert!(s.next_due > Instant::now(), "{} should not be due yet", s.job.label());
		}
	}

	#[test_log::test(tokio::test(start_paused = true))]
	async fn maintenance_loop_runs_every_registered_job() {
		let passes = record_passes(
			vec![
				slot(MaintenanceJob::NodeExpire, Duration::from_millis(5)),
				slot(MaintenanceJob::ChangefeedGc, Duration::from_millis(10)),
				slot(MaintenanceJob::NodeCleanup, Duration::from_millis(40)),
			],
			40,
		)
		.await;
		for job in
			[MaintenanceJob::NodeExpire, MaintenanceJob::ChangefeedGc, MaintenanceJob::NodeCleanup]
		{
			assert!(passes.contains(&job), "{} never ran: {passes:?}", job.label());
		}
		// Cadence is honoured, not just liveness: the 5ms job runs more often
		// than the 40ms one.
		let fast = passes.iter().filter(|j| **j == MaintenanceJob::NodeExpire).count();
		let slow = passes.iter().filter(|j| **j == MaintenanceJob::NodeCleanup).count();
		assert!(fast > slow, "5ms job ran {fast}x, 40ms job ran {slow}x");
	}

	#[test_log::test(tokio::test(start_paused = true))]
	async fn maintenance_loop_does_not_let_one_job_monopolise_the_schedule() {
		// A pass that takes far longer than every interval leaves both jobs
		// permanently overdue, which is exactly when a scheduler starves one.
		let passes = record_passes_with_delay(
			vec![
				slot(MaintenanceJob::ReclaimTombstones, Duration::from_millis(5)),
				slot(MaintenanceJob::NodeExpire, Duration::from_millis(5)),
			],
			20,
			|job| match job {
				MaintenanceJob::ReclaimTombstones => Duration::from_millis(100),
				_ => Duration::ZERO,
			},
		)
		.await;
		for pair in passes.windows(2) {
			assert_ne!(
				pair[0].0, pair[1].0,
				"the same job ran twice in a row while the other was overdue: {passes:?}"
			);
		}
	}

	/// A job's next deadline is measured from when its pass finished, so a pass
	/// that overruns its interval still rests a full interval afterwards rather
	/// than being immediately due again.
	#[test_log::test(tokio::test(start_paused = true))]
	async fn maintenance_loop_rests_a_full_interval_after_an_overrunning_pass() {
		let interval = Duration::from_millis(50);
		let cost = Duration::from_millis(100);
		let passes = record_passes_with_delay(
			vec![slot(MaintenanceJob::ReclaimTombstones, interval)],
			4,
			|_| cost,
		)
		.await;
		for pair in passes.windows(2) {
			let gap = pair[1].1.saturating_duration_since(pair[0].1);
			assert!(
				gap >= cost + interval,
				"passes started {gap:?} apart; the overrunning pass did not rest for {interval:?}"
			);
		}
	}

	#[test_log::test(tokio::test(start_paused = true))]
	async fn maintenance_loop_stops_dispatching_once_cancelled() {
		// The recorder cancels on its third pass; the loop must break rather
		// than start a fourth.
		let passes = record_passes(
			vec![
				slot(MaintenanceJob::NodeExpire, Duration::from_millis(5)),
				slot(MaintenanceJob::ChangefeedGc, Duration::from_millis(5)),
			],
			3,
		)
		.await;
		assert_eq!(passes.len(), 3, "dispatched after cancellation: {passes:?}");
	}

	#[test_log::test(tokio::test)]
	async fn node_membership_update_exits_when_cancelled() {
		let should_continue = super::run_node_membership_update(Duration::from_secs(60), async {
			super::NodeMembershipUpdateResult::Cancelled
		})
		.await;

		assert!(!should_continue);
	}

	#[test_log::test(tokio::test)]
	async fn node_membership_update_continues_after_timeout() {
		let should_continue = super::run_node_membership_update(Duration::from_secs(60), async {
			super::NodeMembershipUpdateResult::TimedOut
		})
		.await;

		assert!(should_continue);
	}

	#[test_log::test(tokio::test)]
	async fn node_membership_update_continues_after_success() {
		let should_continue = super::run_node_membership_update(Duration::from_secs(60), async {
			super::NodeMembershipUpdateResult::Updated
		})
		.await;

		assert!(should_continue);
	}

	#[test_log::test(tokio::test)]
	async fn node_membership_update_continues_after_error() {
		let should_continue = super::run_node_membership_update(Duration::from_secs(60), async {
			super::NodeMembershipUpdateResult::Failed(anyhow::anyhow!("update failed"))
		})
		.await;

		assert!(should_continue);
	}

	#[cfg(feature = "kv-mem")]
	#[test_log::test(tokio::test)]
	pub async fn tasks_complete() {
		let can = CancellationToken::new();
		let opt = EngineOptions::default();
		let dbs = Datastore::new("memory").await.unwrap();
		let tasks = tasks::init(&dbs, can.clone(), &opt);
		can.cancel();
		tasks.resolve().await.unwrap();
	}

	#[cfg(feature = "kv-mem")]
	#[test_log::test(tokio::test)]
	pub async fn tasks_complete_channel_closed() {
		let can = CancellationToken::new();
		// Tick fast enough that passes are genuinely in flight when the
		// cancellation lands, rather than cancelling before anything has run.
		let opt = EngineOptions::default()
			.with_node_membership_refresh_interval(Duration::from_millis(10))
			.with_node_membership_check_interval(Duration::from_millis(10))
			.with_index_compaction_interval(Duration::from_millis(10))
			.with_event_processing_interval(Duration::from_millis(10));
		let dbs = Datastore::new("memory").await.unwrap();
		let tasks = tasks::init(&dbs, can.clone(), &opt);
		tokio::time::sleep(Duration::from_millis(200)).await;
		can.cancel();
		tokio::time::timeout(Duration::from_secs(10), tasks.resolve())
			.await
			.map_err(|e| format!("Timed out after {e}"))
			.unwrap()
			.map_err(|e| format!("Resolution failed: {e}"))
			.unwrap();
	}

	/// The heartbeat must keep its cadence while the rest of the engine's
	/// background work is running, because a stale heartbeat gets this node
	/// archived by another member's expiry scan.
	#[cfg(feature = "kv-mem")]
	#[test_log::test(tokio::test)]
	pub async fn heartbeat_keeps_its_cadence_alongside_other_tasks() {
		let can = CancellationToken::new();
		let opt = EngineOptions::default()
			.with_node_membership_refresh_interval(Duration::from_millis(50));
		let dbs = Datastore::new("memory").await.unwrap();
		dbs.insert_node().await.unwrap();
		let tasks = tasks::init(&dbs, can.clone(), &opt);
		tokio::time::sleep(Duration::from_millis(500)).await;
		let age = dbs.node_heartbeat_age().await.unwrap().expect("the node was registered");
		can.cancel();
		tasks.resolve().await.unwrap();
		assert!(age < Duration::from_millis(500), "heartbeat was {age:?} stale");
	}

	/// However many datastores a process builds, one of them refreshes the
	/// process metrics: the snapshot is a single per-process cache whose CPU
	/// percentage is a delta since its own previous refresh, so a second
	/// refresher would cut that window short by an amount that depends on
	/// nothing but scheduling.
	///
	/// The claim is process-wide, so every datastore this test binary builds
	/// contends for it: the holder that ends the wait below may be one of those
	/// rather than one of these two. What is asserted is the property itself —
	/// that a scheduler holds the claim and a further holder is refused.
	#[cfg(feature = "kv-mem")]
	#[test_log::test(tokio::test)]
	pub async fn only_one_datastore_refreshes_the_process_metrics() {
		let opts = EngineOptions::default()
			.with_system_metrics_refresh_interval(Duration::from_millis(10));
		let _first =
			Datastore::builder().with_engine_options(opts).build_with_path("memory").await.unwrap();
		let _second =
			Datastore::builder().with_engine_options(opts).build_with_path("memory").await.unwrap();
		// A scheduler takes the claim on a pass of its own task, so wait for one
		// rather than assuming it has already run. A probe that wins the race
		// against a pass holds the claim only for the length of the check, and
		// the next pass takes it back.
		tokio::time::timeout(Duration::from_secs(30), async {
			while RefreshClaim::process().take() {
				tokio::time::sleep(Duration::from_millis(10)).await;
			}
		})
		.await
		.expect("no scheduler holds the metrics claim, so every datastore refreshes");
	}

	/// The live-query router is the only conditionally-spawned task: under the
	/// default inline engine it has nothing to deliver, so it is not spawned.
	#[cfg(feature = "kv-mem")]
	#[test_log::test(tokio::test)]
	pub async fn live_query_router_is_not_spawned_under_the_inline_engine() {
		let can = CancellationToken::new();
		let opt = EngineOptions::default();
		let dbs = Datastore::new("memory").await.unwrap();
		let tasks = tasks::init(&dbs, can.clone(), &opt);
		// heartbeat, event processing, index compaction, graph fold,
		// maintenance, sweeps.
		assert_eq!(tasks.0.len(), 6);
		can.cancel();
		tasks.resolve().await.unwrap();
	}

	/// A group whose every job is disabled must not leave an idle task behind.
	#[cfg(feature = "kv-mem")]
	#[test_log::test(tokio::test)]
	pub async fn a_fully_disabled_group_is_not_spawned() {
		let can = CancellationToken::new();
		// Both sweep jobs off; the maintenance group is untouched.
		let opt = EngineOptions::default()
			.with_reclaim_interval(Duration::ZERO)
			.with_rpc_session_gc_interval(Duration::ZERO);
		let dbs = Datastore::new("memory").await.unwrap();
		let tasks = tasks::init(&dbs, can.clone(), &opt);
		assert_eq!(tasks.0.len(), 5, "the empty sweep group should not be spawned");
		can.cancel();
		tasks.resolve().await.unwrap();
	}
}