#![cfg(any(feature = "kv-mem", feature = "kv-rocksdb", feature = "kv-surrealkv"))]
use std::sync::Arc;
use surrealdb_kvs::TransactionType::{Read, Write};
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
use web_time::Duration;
use crate::catalog::providers::DatabaseProvider;
use crate::dbs::{Capabilities, Session};
use crate::key::schema::{DbRoot, ReclaimPrefix};
use crate::key::{AnyRange, Key, RawRange};
use crate::kvs::{Datastore, RECLAIM_BATCH_SIZE};
async fn open(path: &str) -> Arc<Datastore> {
let node_id = Uuid::parse_str("6f1c9d3e-3a12-4a7b-9c4e-5b8d2f0a7e11").unwrap();
Datastore::builder()
.with_id(node_id)
.with_capabilities(Capabilities::all())
.without_maintenance_tasks()
.build_with_path(path)
.await
.unwrap()
}
async fn close(ds: Arc<Datastore>) {
ds.shutdown().await.unwrap();
drop(ds);
}
async fn reclaim_queue_len(ds: &Datastore) -> usize {
let range = ReclaimPrefix {}.range().unwrap();
let tx = ds.transaction(Read).await.unwrap();
let count = tx.count(range, None).await.unwrap();
let _ = tx.cancel().await;
count
}
async fn count_range(ds: &Datastore, range: impl AnyRange) -> usize {
let tx = ds.transaction(Read).await.unwrap();
let count = tx.count(range, None).await.unwrap();
let _ = tx.cancel().await;
count
}
async fn create_tenant(ds: &Datastore, ses: &Session) -> RawRange {
ds.execute("DEFINE NAMESPACE test; DEFINE DATABASE tenant;", ses, None).await.unwrap();
let tx = ds.transaction(Read).await.unwrap();
let def = tx.get_db_by_name("test", "tenant", None).await.unwrap().unwrap();
let _ = tx.cancel().await;
DbRoot {
ns: def.namespace_id,
db: def.database_id,
}
.range()
.unwrap()
}
async fn seed_range(ds: &Datastore, range: &RawRange, count: usize) {
let base = range.start().as_ref().to_vec();
for chunk in 0..count.div_ceil(500) {
let tx = ds.transaction(Write).await.unwrap();
for i in (chunk * 500)..(chunk * 500 + 500).min(count) {
let mut key = base.clone();
key.extend_from_slice(&(i as u64).to_be_bytes());
tx.set(Key::from(key), vec![1u8]).await.unwrap();
}
tx.commit().await.unwrap();
}
}
async fn reclaim_destroys_prefix(path: &str, reopen: bool) {
let ds = open(path).await;
let ses = Session::owner().with_ns("test").with_db("tenant");
let range = create_tenant(&ds, &ses).await;
seed_range(&ds, &range, RECLAIM_BATCH_SIZE as usize * 2).await;
assert!(count_range(&ds, range.clone()).await > 0);
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap();
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
let assert_empty = |ds: Arc<Datastore>, range: RawRange| async move {
let tx = ds.transaction(Read).await.unwrap();
assert_eq!(tx.count(range.clone(), None).await.unwrap(), 0, "the prefix must be empty");
assert!(
tx.getr_raw(range, None).await.unwrap().is_empty(),
"no key of the prefix may remain in the raw keyspace"
);
let _ = tx.cancel().await;
};
assert_empty(Arc::clone(&ds), range.clone()).await;
assert_eq!(reclaim_queue_len(&ds).await, 0, "the reclaim queue must be drained");
if reopen {
close(ds).await;
let ds = open(path).await;
assert_empty(Arc::clone(&ds), range).await;
assert_eq!(reclaim_queue_len(&ds).await, 0, "the drained queue must stay drained");
close(ds).await;
} else {
close(ds).await;
}
}
#[cfg(feature = "kv-mem")]
#[tokio::test]
async fn mem_reclaim_destroys_prefix() {
reclaim_destroys_prefix("memory", false).await;
}
#[cfg(feature = "kv-rocksdb")]
#[tokio::test]
async fn rocksdb_reclaim_destroys_prefix() {
use temp_dir::TempDir;
let dir = TempDir::new().unwrap();
let path = format!("rocksdb:{}", dir.path().to_string_lossy());
reclaim_destroys_prefix(&path, true).await;
}
#[cfg(feature = "kv-surrealkv")]
#[tokio::test]
async fn surrealkv_reclaim_destroys_prefix() {
use temp_dir::TempDir;
let dir = TempDir::new().unwrap();
let path = format!("surrealkv:{}", dir.path().to_string_lossy());
reclaim_destroys_prefix(&path, true).await;
}
#[cfg(any(feature = "kv-rocksdb", feature = "kv-surrealkv"))]
async fn reclaim_resumes_across_a_restart(path: &str) {
use crate::key::KVKeyDecode;
use crate::key::schema::ReclaimKey;
let page = RECLAIM_BATCH_SIZE as u64;
let ds = open(path).await;
let ses = Session::owner().with_ns("test").with_db("tenant");
let range = create_tenant(&ds, &ses).await;
let baseline = count_range(&ds, range.clone()).await;
seed_range(&ds, &range, page as usize * 3 - baseline).await;
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap();
let total = count_range(&ds, range.clone()).await as u64;
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
page,
CancellationToken::new(),
)
.await
.unwrap();
let after_first = count_range(&ds, range.clone()).await as u64;
assert_eq!(after_first, total - page, "the first pass must delete exactly one page");
close(ds).await;
let ds = open(path).await;
assert_eq!(reclaim_queue_len(&ds).await, 1, "the queue entry must survive the restart");
let cursor = {
let tx = ds.transaction(Read).await.unwrap();
let items = tx.getr(ReclaimPrefix {}.range().unwrap(), None).await.unwrap();
let _ = tx.cancel().await;
ReclaimKey::decode_key(&items[0].0).unwrap();
items[0].1.cursor.clone()
};
assert!(cursor.is_some(), "the resume cursor must survive the restart");
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
page,
CancellationToken::new(),
)
.await
.unwrap();
let after_second = count_range(&ds, range.clone()).await as u64;
assert_eq!(after_second, after_first - page, "the second pass must resume, not restart");
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(count_range(&ds, range).await, 0, "the prefix must end up empty");
assert_eq!(reclaim_queue_len(&ds).await, 0, "an emptied prefix must retire its queue entry");
close(ds).await;
}
#[cfg(feature = "kv-rocksdb")]
#[tokio::test]
async fn rocksdb_reclaim_resumes_across_a_restart() {
use temp_dir::TempDir;
let dir = TempDir::new().unwrap();
let path = format!("rocksdb:{}", dir.path().to_string_lossy());
reclaim_resumes_across_a_restart(&path).await;
}
#[cfg(feature = "kv-surrealkv")]
#[tokio::test]
async fn surrealkv_reclaim_resumes_across_a_restart() {
use temp_dir::TempDir;
let dir = TempDir::new().unwrap();
let path = format!("surrealkv:{}", dir.path().to_string_lossy());
reclaim_resumes_across_a_restart(&path).await;
}