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<()>>;
#[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);
const INDEX_COMPACTION_TRIGGER_DEBOUNCE: Duration = Duration::from_millis(250);
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 {
if let Err(e) = task.await {
error!("Background task did not shut down cleanly: {e}");
}
}
Ok(())
}
}
pub fn init(dbs: &Arc<Datastore>, canceller: CancellationToken, opts: &EngineOptions) -> Tasks {
let weak = Arc::downgrade(dbs);
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))] {
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)
}
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 {
let interval = opts.node_membership_refresh_interval;
into_task(async move {
trace!("Updating node registration information 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 !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");
})
}
fn spawn_task_index_compaction(
dbs: Weak<Datastore>,
triggers: Arc<CommitTriggers>,
canceller: CancellationToken,
opts: &EngineOptions,
) -> Task {
let interval = opts.index_compaction_interval;
into_task(async move {
trace!("Running index compaction every {interval:?}");
let mut ticker = interval_ticker(interval).await;
loop {
tokio::select! {
biased;
_ = canceller.cancelled() => break,
_ = triggers.index_compaction.notified() => {
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}");
}
}
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");
})
}
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");
})
}
fn spawn_task_event_processing(
dbs: Weak<Datastore>,
triggers: Arc<CommitTriggers>,
canceller: CancellationToken,
opts: &EngineOptions,
) -> Task {
let interval = opts.event_processing_interval;
into_task(async move {
trace!("Running event processing every {interval:?}");
let mut ticker = interval_ticker(interval).await;
let process_events = async || {
let Some(dbs) = dbs.upgrade() else {
return false;
};
if let Err(e) = dbs.event_processing(interval, &canceller).await
&& !canceller.is_cancelled()
{
error!("Error running event processing: {e}");
}
true
};
loop {
tokio::select! {
biased;
_ = canceller.cancelled() => break,
_ = triggers.async_event.notified() => if !process_events().await { break },
Some(_) = ticker.next() => if !process_events().await { break }
}
}
trace!("Background task exited: Running event processing");
})
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum MaintenanceJob {
NodeExpire,
NodeCleanup,
ChangefeedGc,
ReclaimTombstones,
ResumeIndexBuilds,
RpcSessionGc,
TikvGc,
TikvLockCleanup,
SystemMetricsRefresh,
}
impl MaintenanceJob {
#[cfg(test)]
const ALL: [Self; 9] = [
Self::NodeExpire,
Self::NodeCleanup,
Self::ChangefeedGc,
Self::ReclaimTombstones,
Self::ResumeIndexBuilds,
Self::RpcSessionGc,
Self::TikvGc,
Self::TikvLockCleanup,
Self::SystemMetricsRefresh,
];
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",
}
}
}
struct Slot {
job: MaintenanceJob,
interval: Duration,
next_due: Instant,
}
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),
])
}
fn sweep_slots(opts: &EngineOptions) -> Vec<Slot> {
slots(&[
(MaintenanceJob::ReclaimTombstones, opts.reclaim_interval),
(MaintenanceJob::RpcSessionGc, opts.rpc_session_gc_interval),
])
}
fn slots(jobs: &[(MaintenanceJob, Duration)]) -> Vec<Slot> {
let now = Instant::now();
jobs.iter()
.copied()
.filter(|(_, interval)| !interval.is_zero())
.map(|(job, interval)| Slot {
job,
interval,
next_due: match job {
MaintenanceJob::SystemMetricsRefresh => now,
_ => now + interval,
},
})
.collect()
}
fn next_slot(slots: &[Slot]) -> Option<usize> {
slots.iter().enumerate().min_by_key(|(_, s)| s.next_due).map(|(i, _)| i)
}
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 => {
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());
}
}
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;
slots[i].next_due = Instant::now() + slots[i].interval;
}
}
fn spawn_task_scheduler(
group: &'static str,
slots: Vec<Slot>,
dbs: Weak<Datastore>,
canceller: CancellationToken,
opts: &EngineOptions,
) -> Task {
let opts = *opts;
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();
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 {
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 {
let mut interval = common::time::interval(interval);
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;
#[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;
fn slot(job: MaintenanceJob, interval: Duration) -> Slot {
Slot {
job,
interval,
next_due: Instant::now() + interval,
}
}
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()
}
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![
slot(MaintenanceJob::ChangefeedGc, Duration::from_secs(30)),
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));
assert!(maintenance.contains(&MaintenanceJob::NodeExpire));
assert!(maintenance.contains(&MaintenanceJob::TikvLockCleanup));
assert!(sweeps.contains(&MaintenanceJob::ReclaimTombstones));
}
#[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());
}
}
#[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));
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");
assert!(metrics.next_due <= Instant::now());
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());
}
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() {
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:?}"
);
}
}
#[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() {
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();
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();
}
#[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");
}
#[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();
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");
}
#[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);
assert_eq!(tasks.0.len(), 6);
can.cancel();
tasks.resolve().await.unwrap();
}
#[cfg(feature = "kv-mem")]
#[test_log::test(tokio::test)]
pub async fn a_fully_disabled_group_is_not_spawned() {
let can = CancellationToken::new();
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();
}
}