#![allow(clippy::unwrap_used)]
use std::sync::Arc;
use std::time::Duration;
use surrealdb_kvs::TransactionType::Read;
use crate::key::KVKey;
use crate::key::schema::{BootstrapKey, NodePrefix, VersionKey};
use crate::kvs::testing::{NonRetryableErrorSite, inject_non_retryable_errors};
use crate::kvs::version::MajorVersion;
use crate::kvs::{Datastore, DatastoreError};
async fn ds() -> Arc<Datastore> {
Datastore::builder().without_maintenance_tasks().build_with_path("memory").await.unwrap()
}
async fn write_node_row(ds: &Datastore) {
ds.update_node().await.unwrap();
}
async fn stored_major(ds: &Datastore) -> Option<MajorVersion> {
let txn = ds.transaction(Read).await.unwrap();
let version = txn.get_key(&VersionKey {}, None).await.unwrap();
txn.cancel().await.unwrap();
version
}
async fn stored_sentinel(ds: &Datastore) -> Option<MajorVersion> {
let txn = ds.transaction(Read).await.unwrap();
let sentinel = txn.get_key(&BootstrapKey {}, None).await.unwrap();
txn.cancel().await.unwrap();
sentinel
}
async fn count_node_rows(ds: &Datastore) -> usize {
let txn = ds.transaction(Read).await.unwrap();
let count = txn.count(NodePrefix {}.range().unwrap(), None).await.unwrap();
txn.cancel().await.unwrap();
count
}
#[test]
fn the_sentinel_sorts_below_the_emptiness_probe() {
let sentinel = BootstrapKey {}.encode_key().unwrap();
let version = VersionKey {}.encode_key().unwrap();
assert!(
sentinel < version,
"the bootstrap sentinel {sentinel:?} does not sort before the version key {version:?}, \
so the emptiness probe would count it as data"
);
}
#[tokio::test]
async fn an_interrupted_bootstrap_is_completed_rather_than_refused() {
let ds = ds().await;
let guard =
inject_non_retryable_errors(NonRetryableErrorSite::VersionBootstrapStamp, ds.id(), 1);
ds.get_version().await.unwrap_err();
drop(guard);
assert_eq!(
stored_sentinel(&ds).await,
Some(MajorVersion::latest()),
"the interrupted bootstrap left behind no sentinel recording the version it targeted"
);
assert_eq!(stored_major(&ds).await, None, "the interrupted bootstrap stamped a version");
write_node_row(&ds).await;
let (version, is_new) = ds.check_version().await.unwrap();
assert_eq!(version, MajorVersion::latest());
assert!(is_new, "the start that completed the initialisation did not report creating it");
assert_eq!(stored_major(&ds).await, Some(MajorVersion::latest()));
assert_eq!(stored_sentinel(&ds).await, None, "the sentinel outlived the stamp that retires it");
}
#[tokio::test]
async fn completing_an_interrupted_bootstrap_only_creates_the_datastore_once() {
let ds = ds().await;
let guard =
inject_non_retryable_errors(NonRetryableErrorSite::VersionBootstrapStamp, ds.id(), 1);
ds.get_version().await.unwrap_err();
drop(guard);
assert!(ds.get_version().await.unwrap().1, "the resuming start did not report creating it");
assert_eq!(
ds.get_version().await.unwrap(),
(MajorVersion::latest(), false),
"a start against a settled datastore reported creating it"
);
}
#[tokio::test]
async fn storage_with_no_sentinel_is_still_refused_as_out_of_date() {
let ds = ds().await;
write_node_row(&ds).await;
let err = ds.check_version().await.unwrap_err();
let Some(DatastoreError::OutdatedStorageVersion {
expected,
actual,
}) = err.downcast_ref::<DatastoreError>()
else {
panic!("a datastore holding pre-versioning data was not refused: {err}");
};
assert_eq!(*expected, u16::from(MajorVersion::latest()));
assert_eq!(*actual, u16::from(MajorVersion::v1()));
}
#[tokio::test(start_paused = true)]
async fn the_bootstrap_waits_out_a_storage_layer_that_will_not_serve_it() {
let ds = ds().await;
let _guard =
inject_non_retryable_errors(NonRetryableErrorSite::VersionBootstrapBegin, ds.id(), 6);
let (version, is_new) = ds.check_version().await.unwrap();
assert_eq!(version, MajorVersion::latest());
assert!(is_new);
assert_eq!(stored_major(&ds).await, Some(MajorVersion::latest()));
}
#[tokio::test(start_paused = true)]
async fn the_bootstrap_gives_up_once_its_budget_is_spent() {
let ds = ds().await;
let _guard =
inject_non_retryable_errors(NonRetryableErrorSite::VersionBootstrapBegin, ds.id(), 100_000);
ds.check_version().await.unwrap_err();
}
#[tokio::test]
async fn held_back_maintenance_tasks_write_nothing_until_they_are_started() {
let options = crate::options::EngineOptions {
node_membership_refresh_interval: Duration::from_millis(1),
..Default::default()
};
let ds = Datastore::builder()
.without_maintenance_tasks()
.with_engine_options(options)
.build_with_path("memory")
.await
.unwrap();
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(count_node_rows(&ds).await, 0, "a held-back maintenance task wrote to storage");
ds.start_maintenance_tasks();
ds.start_maintenance_tasks();
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
while count_node_rows(&ds).await == 0 && tokio::time::Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert!(count_node_rows(&ds).await > 0, "the started maintenance tasks never wrote");
ds.shutdown().await.unwrap();
}
#[tokio::test]
async fn a_settled_datastore_carries_no_sentinel() {
let ds = ds().await;
ds.check_version().await.unwrap();
assert_eq!(stored_sentinel(&ds).await, None);
for _ in 0..3 {
ds.check_version().await.unwrap();
assert_eq!(stored_sentinel(&ds).await, None, "a settled start rewrote the sentinel");
}
assert_eq!(stored_major(&ds).await, Some(MajorVersion::latest()));
}
#[tokio::test(start_paused = true)]
async fn a_malformed_marker_is_reported_rather_than_waited_on() {
let ds = ds().await;
let txn = ds.transaction(surrealdb_kvs::TransactionType::Write).await.unwrap();
txn.set(VersionKey {}.encode_key().unwrap(), vec![0xff]).await.unwrap();
txn.commit().await.unwrap();
let started = tokio::time::Instant::now();
let err = ds.check_version().await.unwrap_err();
let waited = started.elapsed();
assert!(
matches!(err.downcast_ref::<DatastoreError>(), Some(DatastoreError::InvalidStorageVersion)),
"a malformed marker was not reported as one: {err}"
);
assert!(waited < Duration::from_secs(1), "a malformed marker was retried for {waited:?}");
}
#[tokio::test]
async fn a_bootstrap_begun_by_a_newer_build_is_left_for_that_build() {
let ds = ds().await;
let newer = MajorVersion::from(u16::from(MajorVersion::latest()) + 1);
let txn = ds.transaction(surrealdb_kvs::TransactionType::Write).await.unwrap();
txn.set_key(&BootstrapKey {}, &newer).await.unwrap();
txn.commit().await.unwrap();
let err = ds.check_version().await.unwrap_err();
let Some(DatastoreError::OutdatedStorageVersion {
expected,
actual,
}) = err.downcast_ref::<DatastoreError>()
else {
panic!("a bootstrap begun by a newer build was not refused: {err}");
};
assert_eq!(*expected, u16::from(MajorVersion::latest()));
assert_eq!(*actual, u16::from(newer));
assert_eq!(stored_major(&ds).await, None, "this build stamped a datastore it cannot complete");
assert_eq!(
stored_sentinel(&ds).await,
Some(newer),
"the sentinel the newer build needs to resume was retired by this one"
);
}
#[tokio::test]
async fn a_bootstrap_another_node_finished_is_not_reported_as_created_here() {
let ds = ds().await;
let txn = ds.transaction(surrealdb_kvs::TransactionType::Write).await.unwrap();
txn.set_key(&VersionKey {}, &MajorVersion::latest()).await.unwrap();
txn.commit().await.unwrap();
assert_eq!(
ds.finish_bootstrap(MajorVersion::latest()).await.unwrap(),
(MajorVersion::latest(), false),
"a node reported creating a datastore another node had already stamped"
);
}
#[tokio::test]
async fn a_default_datastore_is_stamped_before_its_writers_start() {
let options = crate::options::EngineOptions {
node_membership_refresh_interval: Duration::from_millis(1),
..Default::default()
};
let ds =
Datastore::builder().with_engine_options(options).build_with_path("memory").await.unwrap();
assert_eq!(
stored_major(&ds).await,
Some(MajorVersion::latest()),
"a default-built datastore was handed over with no version marker"
);
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(count_node_rows(&ds).await > 0, "the maintenance tasks never wrote");
assert_eq!(ds.check_version().await.unwrap().0, MajorVersion::latest());
ds.shutdown().await.unwrap();
}
#[tokio::test]
async fn a_default_datastore_still_reports_that_this_process_created_it() {
let ds = Datastore::builder().build_with_path("memory").await.unwrap();
assert!(
ds.check_version().await.unwrap().1,
"the process that created the datastore did not report creating it"
);
ds.shutdown().await.unwrap();
}