use std::any::{Any, TypeId};
use std::borrow::Cow;
use std::sync::{Arc, Mutex};
use common::future::BoxFut;
use surrealdb_cnf::ConfigMap;
use surrealdb_kvs::TransactionType::{Read, Write};
use surrealdb_kvs::{
DestroyRange, DestroyRangeHandle, KeyRange, Metrics, Result as KvsResult, Transactable,
TransactionBuilder, TransactionType,
};
use surrealdb_strand::TableName;
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
use web_time::{Duration, SystemTime};
use crate::CommunityComposer;
use crate::catalog::providers::{DatabaseProvider, TableProvider};
use crate::catalog::{DatabaseId, IndexId, NamespaceId};
use crate::dbs::{Capabilities, Session};
use crate::key::reclaim::{ALL_DOC_ID_KINDS, Expunge, ReclaimKind, ReclaimState};
use crate::key::schema::{
DbRoot, DocKeyPrefix, DocLookupPrefix, DocPendingPrefix, IdxRoot, NsRoot, ReclaimKey,
ReclaimPrefix, ReclaimTbPrefix,
};
use crate::key::{AnyRange, KVKeyDecode, KVSubspace, Key, RawRange};
use crate::kvs::test_support::{WriteTxSizeObserver, cleanup_sizes};
use crate::kvs::tx::DocIdsReclaimClaim;
use crate::kvs::{
Datastore, RECLAIM_BATCH_SIZE, TransactionBuilderFactory, TransactionBuilderParts,
};
use crate::observe::ExecutionObserver;
async fn mem_ds() -> Arc<Datastore> {
Datastore::builder()
.with_capabilities(Capabilities::all())
.build_with_path("memory")
.await
.unwrap()
}
async fn observed_ds(observer: Arc<WriteTxSizeObserver>) -> Arc<Datastore> {
Datastore::builder()
.with_capabilities(Capabilities::all())
.without_maintenance_tasks()
.with_observer(observer as Arc<dyn ExecutionObserver>)
.build_with_path("memory")
.await
.unwrap()
}
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 reclaim_queue_len(ds: &Datastore) -> usize {
let range = ReclaimPrefix {}.range().unwrap();
count_range(ds, range).await
}
async fn reclaim_queue_kinds(ds: &Datastore) -> Vec<ReclaimKind> {
let range = ReclaimPrefix {}.range().unwrap();
keys_in_range(ds, range)
.await
.iter()
.map(|k| ReclaimKey::decode_key(k).expect("a queued entry must decode").kind)
.collect()
}
fn count_kind(kinds: &[ReclaimKind], kind: ReclaimKind) -> usize {
kinds.iter().filter(|k| **k == kind).count()
}
fn count_doc_kinds(kinds: &[ReclaimKind]) -> usize {
kinds.iter().filter(|k| k.is_doc_id()).count()
}
async fn reclaim_entry(ds: &Datastore) -> (Vec<u8>, ReclaimState) {
let range = ReclaimPrefix {}.range().unwrap();
let tx = ds.transaction(Read).await.unwrap();
let items = tx.getr(range, None).await.unwrap();
let _ = tx.cancel().await;
assert_eq!(items.len(), 1, "expected exactly one reclaim queue entry");
items[0].clone()
}
async fn reclaim_entry_observed_ms(ds: &Datastore) -> u64 {
reclaim_entry(ds).await.1.observed_ms
}
fn now_ms() -> u64 {
SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap().as_millis() as u64
}
async fn ids(ds: &Datastore, ns: &str, db: &str) -> (NamespaceId, DatabaseId) {
let tx = ds.transaction(Read).await.unwrap();
let def = tx.get_db_by_name(ns, db, None).await.unwrap().unwrap();
let _ = tx.cancel().await;
(def.namespace_id, def.database_id)
}
async fn index_id(ds: &Datastore, ns: NamespaceId, db: DatabaseId, tb: &str, ix: &str) -> IndexId {
let table = TableName::from(tb);
let tx = ds.transaction(Read).await.unwrap();
let def = tx.get_tb_index(ns, db, &table, ix, None).await.unwrap().unwrap();
let _ = tx.cancel().await;
def.index_id
}
async fn seed_range(ds: &Datastore, range: &RawRange, count: usize) {
let base = range.start().as_ref().to_vec();
let tx = ds.transaction(Write).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::from(key), vec![1u8]).await.unwrap();
}
tx.commit().await.unwrap();
}
async fn keys_in_range(ds: &Datastore, range: impl AnyRange) -> Vec<Vec<u8>> {
let tx = ds.transaction(Read).await.unwrap();
let keys = tx.keys_raw(range, u32::MAX, 0, None).await.unwrap();
let _ = tx.cancel().await;
keys
}
fn expected_page_sizes(total: usize) -> Vec<u64> {
let page = RECLAIM_BATCH_SIZE as usize;
let mut sizes = Vec::new();
let mut left = total;
while left > 0 {
let deletes = left.min(page);
sizes.push(deletes as u64 + 1);
left -= deletes;
}
sizes
}
#[tokio::test]
async fn remove_database_defers_data_reclaim() {
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test; DEFINE DATABASE tenant; CREATE thing:1 SET v = 1; CREATE thing:2 SET v = 2;",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = {
let tx = ds.transaction(Read).await.unwrap();
let db = tx.get_db_by_name("test", "tenant", None).await.unwrap().unwrap();
let _ = tx.cancel().await;
(db.namespace_id, db.database_id)
};
let db_prefix = DbRoot {
ns: ns_id,
db: db_id,
};
let range = db_prefix.range().unwrap();
assert!(
count_range(&ds, range.clone()).await > 0,
"the database prefix should hold data before reclaim"
);
assert_eq!(reclaim_queue_len(&ds).await, 0, "no reclaim jobs before removal");
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap();
{
let tx = ds.transaction(Read).await.unwrap();
let gone = tx.get_db_by_name("test", "tenant", None).await.unwrap();
let _ = tx.cancel().await;
assert!(gone.is_none(), "database must be invisible immediately after REMOVE");
}
assert_eq!(reclaim_queue_len(&ds).await, 1, "REMOVE DATABASE must enqueue one reclaim job");
assert!(
count_range(&ds, range.clone()).await > 0,
"data must still be present before the reclaim task runs"
);
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
count_range(&ds, range.clone()).await,
0,
"the reclaim task must destroy the database data prefix"
);
assert_eq!(reclaim_queue_len(&ds).await, 0, "the reclaim task must drain the reclaim queue");
ds.execute("DEFINE DATABASE tenant;", &ses, None).await.unwrap();
let new_db_id = {
let tx = ds.transaction(Read).await.unwrap();
let db = tx.get_db_by_name("test", "tenant", None).await.unwrap().unwrap();
let _ = tx.cancel().await;
db.database_id
};
assert_ne!(new_db_id, db_id, "recreated database must get a fresh, never-reused id");
let new_prefix = DbRoot {
ns: ns_id,
db: new_db_id,
};
let new_range = new_prefix.range().unwrap();
assert_eq!(
count_range(&ds, new_range.clone()).await,
0,
"recreated database must start empty with no data leaked from the removed one"
);
}
#[tokio::test]
async fn reclaim_is_idempotent_and_safe_when_empty() {
let ds = mem_ds().await;
let (iters, errors) = Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(errors, 0);
assert_eq!(iters, 0, "an empty queue performs no reclaim iterations");
}
#[tokio::test]
async fn reclaim_respects_grace_period() {
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test; DEFINE DATABASE tenant; CREATE thing:1 SET v = 1; CREATE thing:2 SET v = 2;",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = {
let tx = ds.transaction(Read).await.unwrap();
let db = tx.get_db_by_name("test", "tenant", None).await.unwrap().unwrap();
let _ = tx.cancel().await;
(db.namespace_id, db.database_id)
};
let db_prefix = DbRoot {
ns: ns_id,
db: db_id,
};
let range = db_prefix.range().unwrap();
let reader = ds.transaction(Read).await.unwrap();
let reader_seen_before = reader.getr_raw(range.clone(), None).await.unwrap().len();
assert!(reader_seen_before > 0, "reader should see the data before removal");
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap();
assert_eq!(
reclaim_entry_observed_ms(&ds).await,
0,
"a freshly-enqueued reclaim entry must start unobserved"
);
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::from_secs(3600),
CancellationToken::new(),
)
.await
.unwrap();
assert!(
count_range(&ds, range.clone()).await > 0,
"data must NOT be destroyed while inside the grace window"
);
assert_eq!(reclaim_queue_len(&ds).await, 1, "the reclaim job must remain queued during grace");
assert_ne!(
reclaim_entry_observed_ms(&ds).await,
0,
"the first pass must stamp an observation time instead of reclaiming"
);
assert_eq!(
reader.getr_raw(range.clone(), None).await.unwrap().len(),
reader_seen_before,
"an in-flight reader opened before REMOVE must still see its data"
);
let _ = reader.cancel().await;
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
count_range(&ds, range.clone()).await,
0,
"data must be reclaimed once past the grace window"
);
assert_eq!(reclaim_queue_len(&ds).await, 0, "the reclaim job must be drained after reclaim");
}
#[tokio::test]
async fn reclaim_runs_once_observation_ages_past_grace() {
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test; DEFINE DATABASE tenant; CREATE thing:1 SET v = 1;",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = {
let tx = ds.transaction(Read).await.unwrap();
let db = tx.get_db_by_name("test", "tenant", None).await.unwrap().unwrap();
let _ = tx.cancel().await;
(db.namespace_id, db.database_id)
};
let range = DbRoot {
ns: ns_id,
db: db_id,
}
.range()
.unwrap();
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap();
{
let range = ReclaimPrefix {}.range().unwrap();
let tx = ds.transaction(Write).await.unwrap();
let items = tx.getr(range, None).await.unwrap();
assert_eq!(items.len(), 1);
let rc = ReclaimKey::decode_key(&items[0].0).unwrap();
tx.set_key(
&rc,
&ReclaimState {
observed_ms: now_ms().saturating_sub(3_600_000),
cursor: None,
},
)
.await
.unwrap();
tx.commit().await.unwrap();
}
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::from_secs(60),
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
count_range(&ds, range.clone()).await,
0,
"an entry observed longer than the grace ago must be reclaimed"
);
assert_eq!(reclaim_queue_len(&ds).await, 0, "the reclaim job must be drained");
}
#[tokio::test]
async fn reclaim_deletes_in_bounded_committed_batches() {
let page = RECLAIM_BATCH_SIZE as usize;
let observer = Arc::new(WriteTxSizeObserver::default());
let ds = observed_ds(Arc::clone(&observer)).await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("DEFINE NAMESPACE test; DEFINE DATABASE tenant;", &ses, None).await.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let range = DbRoot {
ns: ns_id,
db: db_id,
}
.range()
.unwrap();
let baseline = count_range(&ds, range.clone()).await;
seed_range(&ds, &range, page * 2 + page / 2 - baseline).await;
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap();
let total = count_range(&ds, range.clone()).await;
assert_eq!(total, page * 2 + page / 2, "the prefix must be sized to 2.5 pages");
observer.clear();
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 reclaim must empty the prefix");
assert_eq!(reclaim_queue_len(&ds).await, 0, "an emptied prefix must retire its queue entry");
let sizes = observer.sizes();
assert!(
sizes.iter().all(|&n| n <= page as u64 + 1),
"a reclaim write transaction carried more than one page: {sizes:?}"
);
assert!(
sizes.iter().sum::<u64>() >= total as u64,
"every reclaimed key must be charged as a write: {sizes:?}"
);
assert_eq!(
cleanup_sizes(&sizes),
expected_page_sizes(total),
"unexpected reclaim page sequence"
);
}
#[tokio::test]
async fn reclaim_keeps_tombstone_queued_until_prefix_is_empty() {
let page = RECLAIM_BATCH_SIZE as usize;
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("DEFINE NAMESPACE test; DEFINE DATABASE tenant;", &ses, None).await.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let range = DbRoot {
ns: ns_id,
db: db_id,
}
.range()
.unwrap();
let baseline = count_range(&ds, range.clone()).await;
seed_range(&ds, &range, page * 2 + page / 2 - baseline).await;
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap();
let total = count_range(&ds, range.clone()).await;
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
page as u64,
CancellationToken::new(),
)
.await
.unwrap();
let remaining = count_range(&ds, range.clone()).await;
assert_eq!(remaining, total - page, "the pass must delete exactly its budget");
assert_eq!(
reclaim_queue_len(&ds).await,
1,
"a prefix that is not yet empty must keep its queue entry"
);
assert!(
reclaim_entry(&ds).await.1.cursor.is_some(),
"a partial pass must leave a durable resume cursor"
);
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(count_range(&ds, range).await, 0);
assert_eq!(reclaim_queue_len(&ds).await, 0, "an emptied prefix must retire its queue entry");
}
#[tokio::test]
async fn reclaim_resumes_after_the_committed_cursor() {
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("DEFINE NAMESPACE test; DEFINE DATABASE tenant;", &ses, None).await.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let range = DbRoot {
ns: ns_id,
db: db_id,
}
.range()
.unwrap();
seed_range(&ds, &range, 16).await;
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap();
let keys = keys_in_range(&ds, range.clone()).await;
assert!(keys.len() > 8, "the prefix must hold enough keys to split");
let already_done = 5;
{
let (k, state) = reclaim_entry(&ds).await;
let rc = ReclaimKey::decode_key(&k).unwrap();
let tx = ds.transaction(Write).await.unwrap();
tx.set_key(
&rc,
&ReclaimState {
observed_ms: state.observed_ms,
cursor: Some(keys[already_done - 1].clone()),
},
)
.await
.unwrap();
tx.commit().await.unwrap();
}
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
keys_in_range(&ds, range).await,
keys[..already_done],
"the reclaim must delete only the keys ordering after the cursor"
);
assert_eq!(reclaim_queue_len(&ds).await, 0, "the entry is retired once the scan finds nothing");
}
#[tokio::test]
async fn reclaim_pages_are_invisible_to_an_in_flight_reader() {
let page = RECLAIM_BATCH_SIZE as usize;
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("DEFINE NAMESPACE test; DEFINE DATABASE tenant;", &ses, None).await.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let range = DbRoot {
ns: ns_id,
db: db_id,
}
.range()
.unwrap();
let baseline = count_range(&ds, range.clone()).await;
seed_range(&ds, &range, page * 2 - baseline).await;
let reader = ds.transaction(Read).await.unwrap();
let seen_before = reader.count(range.clone(), None).await.unwrap();
assert_eq!(seen_before, page * 2);
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();
assert_eq!(
reader.count(range.clone(), None).await.unwrap(),
seen_before,
"a reader opened before the reclaim must not see its pages commit"
);
let _ = reader.cancel().await;
assert_eq!(count_range(&ds, range).await, 0, "a fresh reader sees the prefix emptied");
}
enum Removal {
Namespace {
expunge: bool,
},
Database {
expunge: bool,
},
Index,
IndexExpunged,
}
#[tokio::test]
async fn reclaim_bounds_every_kind_and_expunge_mode() {
for removal in [
Removal::Namespace {
expunge: false,
},
Removal::Namespace {
expunge: true,
},
Removal::Database {
expunge: false,
},
Removal::Database {
expunge: true,
},
Removal::Index,
Removal::IndexExpunged,
] {
assert_removal_reclaims_in_bounded_pages(removal).await;
}
}
async fn assert_removal_reclaims_in_bounded_pages(removal: Removal) {
let page = RECLAIM_BATCH_SIZE as usize;
let observer = Arc::new(WriteTxSizeObserver::default());
let ds = observed_ds(Arc::clone(&observer)).await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE DATABASE other;
DEFINE TABLE thing SCHEMALESS;
DEFINE INDEX idx ON thing FIELDS val;",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let (_, other_id) = ids(&ds, "test", "other").await;
let ix_id = index_id(&ds, ns_id, db_id, "thing", "idx").await;
let target = match removal {
Removal::Namespace {
..
} => NsRoot {
ns: ns_id,
}
.range()
.unwrap(),
Removal::Database {
..
} => DbRoot {
ns: ns_id,
db: db_id,
}
.range()
.unwrap(),
Removal::Index | Removal::IndexExpunged => IdxRoot {
ns: ns_id,
db: db_id,
tb: Cow::Owned(TableName::from("thing")),
ix: ix_id,
}
.range()
.unwrap(),
};
let sized = page * 2 + page / 2;
let baseline = count_range(&ds, target.clone()).await;
let to_seed = sized - baseline;
match removal {
Removal::Namespace {
..
} => {
let tenant = DbRoot {
ns: ns_id,
db: db_id,
}
.range()
.unwrap();
let other = DbRoot {
ns: ns_id,
db: other_id,
}
.range()
.unwrap();
seed_range(&ds, &tenant, to_seed / 2).await;
seed_range(&ds, &other, to_seed - to_seed / 2).await;
}
_ => seed_range(&ds, &target, to_seed).await,
}
let statement = match removal {
Removal::Namespace {
expunge: false,
} => Some("REMOVE NAMESPACE test;"),
Removal::Namespace {
expunge: true,
} => Some("REMOVE NAMESPACE AND EXPUNGE test;"),
Removal::Database {
expunge: false,
} => Some("REMOVE DATABASE tenant;"),
Removal::Database {
expunge: true,
} => Some("REMOVE DATABASE AND EXPUNGE tenant;"),
Removal::Index => Some("REMOVE INDEX idx ON thing;"),
Removal::IndexExpunged => None,
};
match statement {
Some(statement) => {
ds.execute(statement, &ses, None).await.unwrap();
}
None => {
let rc = ReclaimKey {
kind: ReclaimKind::Index,
ns: ns_id,
db: db_id,
tb: Cow::Owned(TableName::from("thing")),
ix: ix_id,
expunge: Expunge::Expunge,
uid: Uuid::now_v7(),
};
let tx = ds.transaction(Write).await.unwrap();
tx.set_key(&rc, &ReclaimState::enqueued()).await.unwrap();
tx.commit().await.unwrap();
}
}
let kinds = reclaim_queue_kinds(&ds).await;
let own_kind = match removal {
Removal::Namespace {
..
} => ReclaimKind::Namespace,
Removal::Database {
..
} => ReclaimKind::Database,
Removal::Index | Removal::IndexExpunged => ReclaimKind::Index,
};
assert_eq!(
count_kind(&kinds, own_kind),
1,
"unexpected number of {own_kind:?} entries queued for the removal"
);
let expected_doc_entries = match removal {
Removal::Index => ALL_DOC_ID_KINDS.len(),
_ => 0,
};
assert_eq!(
count_doc_kinds(&kinds),
expected_doc_entries,
"unexpected number of shared doc-ID entries queued for the removal"
);
let total = count_range(&ds, target.clone()).await;
assert!(total > 2 * page, "the target must hold more than two pages, holds {total}");
observer.clear();
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(count_range(&ds, target).await, 0, "the reclaim must empty the target prefix");
assert_eq!(reclaim_queue_len(&ds).await, 0, "an emptied prefix must retire its queue entry");
let sizes = observer.sizes();
assert!(
sizes.iter().all(|&n| n <= page as u64 + 1),
"a reclaim write transaction carried more than one page: {sizes:?}"
);
assert!(
sizes.iter().sum::<u64>() >= total as u64,
"every reclaimed key must be charged as a write: {sizes:?}"
);
let mut expected = expected_page_sizes(total);
if kinds.len() > 1 {
expected.push(kinds.len() as u64);
}
assert_eq!(cleanup_sizes(&sizes), expected, "unexpected reclaim page sequence");
}
#[tokio::test]
async fn reclaim_shares_its_budget_across_queued_entries() {
let page = RECLAIM_BATCH_SIZE as usize;
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("one");
ds.execute("DEFINE NAMESPACE test; DEFINE DATABASE one; DEFINE DATABASE two;", &ses, None)
.await
.unwrap();
let (ns_id, one_id) = ids(&ds, "test", "one").await;
let (_, two_id) = ids(&ds, "test", "two").await;
let one = DbRoot {
ns: ns_id,
db: one_id,
}
.range()
.unwrap();
let two = DbRoot {
ns: ns_id,
db: two_id,
}
.range()
.unwrap();
for range in [&one, &two] {
let baseline = count_range(&ds, range.clone()).await;
seed_range(&ds, range, page * 2 - baseline).await;
}
ds.execute("REMOVE DATABASE one; REMOVE DATABASE two;", &ses, None).await.unwrap();
assert_eq!(reclaim_queue_len(&ds).await, 2, "both removals must enqueue an entry");
let (first, second) = match one.start().as_ref() < two.start().as_ref() {
true => (one, two),
false => (two, one),
};
let before_first = count_range(&ds, first.clone()).await;
let before_second = count_range(&ds, second.clone()).await;
let budget = page + page / 2;
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
budget as u64,
CancellationToken::new(),
)
.await
.unwrap();
let deleted_first = before_first - count_range(&ds, first.clone()).await;
let deleted_second = before_second - count_range(&ds, second.clone()).await;
assert_eq!(
deleted_first + deleted_second,
budget,
"the pass must spend its whole budget and no more, across both entries"
);
assert_eq!(
deleted_first,
budget / 2,
"the leading entry takes its share, not the whole budget"
);
assert_eq!(deleted_second, budget - budget / 2, "the entry behind it must get what is left");
assert_eq!(reclaim_queue_len(&ds).await, 2, "neither prefix is empty yet, so both stay queued");
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(count_range(&ds, first).await, 0, "the leading prefix must end up empty");
assert_eq!(count_range(&ds, second).await, 0, "the trailing prefix must end up empty");
assert_eq!(reclaim_queue_len(&ds).await, 0, "both entries retire once their prefixes are gone");
}
#[tokio::test]
async fn reclaim_queue_updates_are_bounded_by_the_batch() {
let batch = RECLAIM_BATCH_SIZE as usize;
let observer = Arc::new(WriteTxSizeObserver::default());
let ds = observed_ds(Arc::clone(&observer)).await;
let queued = batch + batch / 2;
{
let tx = ds.transaction(Write).await.unwrap();
for i in 0..queued {
let rc = ReclaimKey::database(
NamespaceId(1),
DatabaseId(i as u32 + 1),
false,
Uuid::now_v7(),
);
tx.set_key(&rc, &ReclaimState::enqueued()).await.unwrap();
}
tx.commit().await.unwrap();
}
assert_eq!(reclaim_queue_len(&ds).await, queued);
observer.clear();
const MAX_PASSES: usize = 8;
let mut passes = 0;
while reclaim_queue_len(&ds).await > 0 {
passes += 1;
assert!(passes <= MAX_PASSES, "the queue must drain within {MAX_PASSES} passes");
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
}
let sizes = observer.sizes();
assert!(
sizes.iter().all(|&n| n <= batch as u64),
"a reclaim write transaction carried more than one batch: {sizes:?}"
);
let updates = cleanup_sizes(&sizes);
assert_eq!(
updates,
vec![batch as u64, (queued - batch) as u64],
"unexpected queue update sequence"
);
assert_eq!(
updates.iter().sum::<u64>(),
queued as u64,
"every queued entry must be accounted for by a queue update"
);
}
#[tokio::test]
async fn unsafe_destroy_range_reports_unsupported_rather_than_succeeding() {
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("DEFINE NAMESPACE test; DEFINE DATABASE tenant;", &ses, None).await.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let range = DbRoot {
ns: ns_id,
db: db_id,
}
.range()
.unwrap();
seed_range(&ds, &range, 8).await;
let seeded = count_range(&ds, range.clone()).await;
assert!(seeded > 0, "the range must hold keys for the assertion below to mean anything");
let err = ds.unsafe_destroy_range(range.clone().into_key_range()).await.unwrap_err();
assert!(
matches!(
err.downcast_ref::<crate::kvs::Error>(),
Some(crate::kvs::Error::RangeDestroyNotSupported)
),
"a backend offering no range destroy must report it: {err}"
);
assert_eq!(
count_range(&ds, range).await,
seeded,
"a destroy reported as unsupported must leave every key in place"
);
}
type DestroyedRange = (Vec<u8>, Vec<u8>);
#[derive(Clone, Default)]
struct DestroyedRanges(Arc<Mutex<Vec<DestroyedRange>>>);
impl DestroyedRanges {
fn ranges(&self) -> Vec<DestroyedRange> {
self.0.lock().unwrap().clone()
}
}
struct RecordingDestroyRange {
inner: Arc<dyn TransactionBuilder>,
destroyed: DestroyedRanges,
}
impl DestroyRange for RecordingDestroyRange {
fn mechanism(&self) -> &'static str {
"recording"
}
fn destroy_range(&self, range: KeyRange<'static>) -> BoxFut<'_, KvsResult<()>> {
Box::pin(async move {
self.destroyed
.0
.lock()
.unwrap()
.push((range.start.as_ref().to_vec(), range.end.as_ref().to_vec()));
let (txn, _) = self.inner.new_transaction(TransactionType::Write).await?;
txn.delr(range).await?;
txn.commit().await
})
}
}
struct DestroyingBuilder {
inner: Arc<dyn TransactionBuilder>,
handle: Arc<DestroyRangeHandle>,
}
impl TransactionBuilder for DestroyingBuilder {
fn name(&self) -> &'static str {
self.inner.name()
}
fn new_transaction(
&self,
write: TransactionType,
) -> BoxFut<'_, KvsResult<(Box<dyn Transactable>, bool)>> {
self.inner.new_transaction(write)
}
fn shutdown(&self) -> BoxFut<'_, KvsResult<()>> {
self.inner.shutdown()
}
fn register_metrics(&self) -> Option<Metrics> {
self.inner.register_metrics()
}
fn collect_u64_metric(&self, metric: &str) -> Option<u64> {
self.inner.collect_u64_metric(metric)
}
fn wait_until_serve_ready(&self) -> BoxFut<'_, KvsResult<()>> {
self.inner.wait_until_serve_ready()
}
fn extension(&self, id: TypeId) -> Option<Arc<dyn Any + Send + Sync>> {
match id == TypeId::of::<DestroyRangeHandle>() {
true => Some(Arc::clone(&self.handle) as Arc<dyn Any + Send + Sync>),
false => self.inner.extension(id),
}
}
}
struct DestroyRangeComposer(DestroyedRanges);
impl TransactionBuilderFactory for DestroyRangeComposer {
type RouterState = ();
async fn new_transaction_builder(
&self,
path: &str,
canceller: CancellationToken,
config: ConfigMap,
) -> anyhow::Result<TransactionBuilderParts<Self::RouterState>> {
let parts = CommunityComposer().new_transaction_builder(path, canceller, config).await?;
let inner: Arc<dyn TransactionBuilder> = Arc::from(parts.builder);
let handle = Arc::new(DestroyRangeHandle(Arc::new(RecordingDestroyRange {
inner: Arc::clone(&inner),
destroyed: self.0.clone(),
})));
Ok(TransactionBuilderParts::without_router_state(Box::new(DestroyingBuilder {
inner,
handle,
})))
}
fn path_valid(&self, v: &str) -> anyhow::Result<String> {
CommunityComposer().path_valid(v)
}
}
#[tokio::test]
async fn reclaim_destroys_the_prefix_in_one_call_where_the_backend_offers_it() {
let page = RECLAIM_BATCH_SIZE as usize;
let destroyed = DestroyedRanges::default();
let ds = Datastore::builder()
.with_capabilities(Capabilities::all())
.without_maintenance_tasks()
.build_with_factory_path("memory", DestroyRangeComposer(destroyed.clone()))
.await
.unwrap();
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("DEFINE NAMESPACE test; DEFINE DATABASE tenant;", &ses, None).await.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let range = DbRoot {
ns: ns_id,
db: db_id,
}
.range()
.unwrap();
let baseline = count_range(&ds, range.clone()).await;
seed_range(&ds, &range, page * 2 - baseline).await;
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap();
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
1,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
count_range(&ds, range.clone()).await,
0,
"the whole prefix must go, however small the pass budget"
);
assert_eq!(
reclaim_queue_len(&ds).await,
0,
"a prefix destroyed in one call retires its entry in the same pass"
);
assert_eq!(
destroyed.ranges(),
vec![(range.start().as_ref().to_vec(), range.end().as_ref().to_vec())],
"the reclaim must hand over exactly the entry's prefix, once"
);
}
#[tokio::test]
async fn a_doc_id_reclaim_never_takes_the_out_of_transaction_destroy() {
let page = RECLAIM_BATCH_SIZE as usize;
let destroyed = DestroyedRanges::default();
let ds = Datastore::builder()
.with_capabilities(Capabilities::all())
.without_maintenance_tasks()
.build_with_factory_path("memory", DestroyRangeComposer(destroyed.clone()))
.await
.unwrap();
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE TABLE thing SCHEMALESS;
DEFINE INDEX idx ON thing FIELDS val;",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let tb = TableName::from("thing");
let dd = DocKeyPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let di = DocLookupPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
seed_range(&ds, &dd, page).await;
seed_range(&ds, &di, page).await;
ds.execute("REMOVE INDEX idx ON thing;", &ses, None).await.unwrap();
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
1,
CancellationToken::new(),
)
.await
.unwrap();
for (name, range) in [("!dd", &dd), ("!di", &di)] {
let bounds = (range.start().as_ref().to_vec(), range.end().as_ref().to_vec());
assert!(
!destroyed.ranges().contains(&bounds),
"{name} must not be handed to the out-of-transaction destroy"
);
assert!(
count_range(&ds, range.clone()).await > 0,
"{name} must be deleted in bounded pages, not taken whole"
);
}
let left = count_range(&ds, dd.clone()).await + count_range(&ds, di.clone()).await;
let mut response =
ds.execute("DEFINE INDEX idx ON thing FIELDS val;", &ses, None).await.unwrap();
let error = response.remove(0).result.unwrap_err().to_string();
assert!(error.contains("still being reclaimed"));
assert_eq!(
count_range(&ds, dd.clone()).await + count_range(&ds, di.clone()).await,
left,
"the rejected DDL must not synchronously purge the remainder"
);
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(count_range(&ds, dd).await + count_range(&ds, di).await, 0);
let mut response =
ds.execute("DEFINE INDEX idx ON thing FIELDS val;", &ses, None).await.unwrap();
response.remove(0).result.unwrap();
}
#[tokio::test]
async fn reclaim_retires_a_prefix_the_budget_ran_out_on_emptying() {
let page = RECLAIM_BATCH_SIZE as usize;
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("DEFINE NAMESPACE test; DEFINE DATABASE tenant;", &ses, None).await.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let range = DbRoot {
ns: ns_id,
db: db_id,
}
.range()
.unwrap();
let baseline = count_range(&ds, range.clone()).await;
seed_range(&ds, &range, page - baseline).await;
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap();
let total = count_range(&ds, range.clone()).await;
assert_eq!(total, page, "the prefix must hold exactly one page");
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
total as u64,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(count_range(&ds, range).await, 0, "the pass must empty the prefix");
assert_eq!(
reclaim_queue_len(&ds).await,
0,
"a prefix emptied as the budget ran out must still retire its queue entry"
);
}
#[tokio::test]
async fn reclaim_restarts_an_entry_whose_state_cannot_be_read() {
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("DEFINE NAMESPACE test; DEFINE DATABASE tenant;", &ses, None).await.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let range = DbRoot {
ns: ns_id,
db: db_id,
}
.range()
.unwrap();
let baseline = count_range(&ds, range.clone()).await;
seed_range(&ds, &range, 32 - baseline).await;
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap();
assert!(count_range(&ds, range.clone()).await > 0, "the prefix must hold data to reclaim");
{
let (k, _) = reclaim_entry(&ds).await;
let tx = ds.transaction(Write).await.unwrap();
tx.set(Key::from(&k), vec![0xde, 0xad, 0xbe, 0xef]).await.unwrap();
tx.commit().await.unwrap();
}
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
count_range(&ds, range).await,
0,
"an entry whose state cannot be read must still have its prefix reclaimed"
);
assert_eq!(reclaim_queue_len(&ds).await, 0, "and must still be retired");
}
#[tokio::test]
async fn reclaim_shares_its_budget_only_with_entries_past_the_grace() {
let page = RECLAIM_BATCH_SIZE as usize;
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("big");
ds.execute(
"DEFINE NAMESPACE test; DEFINE DATABASE big;
DEFINE DATABASE f1; DEFINE DATABASE f2; DEFINE DATABASE f3; DEFINE DATABASE f4;",
&ses,
None,
)
.await
.unwrap();
let (ns_id, big_id) = ids(&ds, "test", "big").await;
let big = DbRoot {
ns: ns_id,
db: big_id,
}
.range()
.unwrap();
let baseline = count_range(&ds, big.clone()).await;
seed_range(&ds, &big, page * 3 - baseline).await;
ds.execute(
"REMOVE DATABASE big; REMOVE DATABASE f1; REMOVE DATABASE f2;
REMOVE DATABASE f3; REMOVE DATABASE f4;",
&ses,
None,
)
.await
.unwrap();
assert_eq!(reclaim_queue_len(&ds).await, 5, "every removal must enqueue an entry");
{
let tx = ds.transaction(Write).await.unwrap();
let items = tx.getr(ReclaimPrefix {}.range().unwrap(), None).await.unwrap();
for (k, state) in &items {
let rc = ReclaimKey::decode_key(k).unwrap();
if rc.db == big_id {
tx.set_key(
&rc,
&ReclaimState {
observed_ms: now_ms().saturating_sub(3_600_000),
cursor: state.cursor.clone(),
},
)
.await
.unwrap();
}
}
tx.commit().await.unwrap();
}
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::from_secs(60),
(page * 3) as u64,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
count_range(&ds, big).await,
0,
"the entry past its grace must receive the budget the ageing entries cannot spend"
);
}
#[tokio::test]
async fn remove_index_defers_the_shared_doc_id_reclaim() {
let page = RECLAIM_BATCH_SIZE as usize;
let observer = Arc::new(WriteTxSizeObserver::default());
let ds = observed_ds(Arc::clone(&observer)).await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE TABLE thing SCHEMALESS;
DEFINE INDEX idx ON thing FIELDS val;",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let tb = TableName::from("thing");
let di = DocLookupPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let dd = DocKeyPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let dp = DocPendingPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let sized = page * 2 + page / 2;
seed_range(&ds, &di, sized).await;
seed_range(&ds, &dd, sized).await;
let seeded = count_range(&ds, di.clone()).await + count_range(&ds, dd.clone()).await;
assert!(seeded > 4 * page, "the shared space must hold more than four pages, holds {seeded}");
observer.clear();
ds.execute("REMOVE INDEX idx ON thing;", &ses, None).await.unwrap();
let sizes = observer.sizes();
assert!(
sizes.iter().all(|&n| n <= page as u64 + 1),
"a REMOVE INDEX write transaction carried more than one page: {sizes:?}"
);
assert_eq!(
count_range(&ds, di.clone()).await + count_range(&ds, dd.clone()).await,
seeded,
"REMOVE INDEX must defer the shared doc-ID space, not delete it inline"
);
let kinds = reclaim_queue_kinds(&ds).await;
assert_eq!(
count_kind(&kinds, ReclaimKind::Index),
1,
"the removal must enqueue the index prefix"
);
for kind in ALL_DOC_ID_KINDS {
assert_eq!(
count_kind(&kinds, kind),
1,
"the removal must enqueue the {kind:?} shared doc-ID prefix"
);
}
observer.clear();
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
for (name, range) in [("!di", di), ("!dd", dd), ("!dp", dp)] {
assert_eq!(count_range(&ds, range).await, 0, "the reclaim must empty {name}");
}
assert_eq!(reclaim_queue_len(&ds).await, 0, "an emptied prefix must retire its queue entry");
let sizes = observer.sizes();
assert!(
sizes.iter().all(|&n| n <= page as u64 + 1),
"a reclaim write transaction carried more than one page: {sizes:?}"
);
assert!(
sizes.iter().sum::<u64>() >= seeded as u64,
"every reclaimed key must be charged as a write: {sizes:?}"
);
}
#[tokio::test]
async fn defining_a_doc_id_index_cancels_a_queued_doc_id_reclaim() {
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE TABLE thing SCHEMALESS;
DEFINE INDEX idx ON thing FIELDS val;
CREATE thing:1 SET val = 'a';
CREATE thing:2 SET val = 'b';",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let tb = TableName::from("thing");
let di = DocLookupPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let dd = DocKeyPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let mapped = keys_in_range(&ds, di.clone()).await;
assert!(!mapped.is_empty(), "the index must have mapped the records it covers");
ds.execute("REMOVE INDEX idx ON thing;", &ses, None).await.unwrap();
assert_eq!(
keys_in_range(&ds, di.clone()).await,
mapped,
"the drop must leave the shared space for the background reclaim"
);
ds.execute("DEFINE INDEX idx ON thing FIELDS val;", &ses, None).await.unwrap();
let kinds = reclaim_queue_kinds(&ds).await;
assert_eq!(
count_doc_kinds(&kinds),
0,
"defining a consumer must cancel every queued shared doc-ID reclaim"
);
assert_eq!(
count_kind(&kinds, ReclaimKind::Index),
kinds.len(),
"only the dropped index's own prefix may stay queued"
);
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
keys_in_range(&ds, di).await,
mapped,
"the reclaim must not delete mappings the redefined index has adopted"
);
assert!(
!keys_in_range(&ds, dd).await.is_empty(),
"the reverse direction must survive with the forward one"
);
let mut res = ds.execute("SELECT val FROM thing WHERE val = 'a';", &ses, None).await.unwrap();
assert_eq!(
format!("{:?}", res.remove(0).result.unwrap()),
r#"Array(Array([Object(Object({"val": String("a")}))]))"#,
"the redefined index must resolve the records it adopted"
);
}
#[tokio::test]
async fn defining_a_doc_id_index_waits_for_a_partly_reclaimed_doc_id_space() {
let page = RECLAIM_BATCH_SIZE as usize;
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE TABLE thing SCHEMALESS;
DEFINE INDEX idx ON thing FIELDS val;",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let tb = TableName::from("thing");
let di = DocLookupPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let dd = DocKeyPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
seed_range(&ds, &di, page * 2).await;
seed_range(&ds, &dd, page * 2).await;
ds.execute("REMOVE INDEX idx ON thing;", &ses, None).await.unwrap();
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
page as u64,
CancellationToken::new(),
)
.await
.unwrap();
let left = count_range(&ds, di.clone()).await + count_range(&ds, dd.clone()).await;
assert!(
left > 0 && left < page * 4,
"the pass must have torn the space, not finished or skipped it: {left} left"
);
let queue_before = reclaim_queue_len(&ds).await;
let mut response =
ds.execute("DEFINE INDEX idx ON thing FIELDS val;", &ses, None).await.unwrap();
let error = response.remove(0).result.unwrap_err().to_string();
assert!(
error.contains("still being reclaimed"),
"the DDL must ask the caller to retry while bounded cleanup continues: {error}"
);
assert_eq!(
count_range(&ds, di.clone()).await + count_range(&ds, dd.clone()).await,
left,
"the rejected DDL must not synchronously purge the torn space"
);
assert_eq!(
reclaim_queue_len(&ds).await,
queue_before,
"the rejected DDL must leave the bounded reclaim resumable"
);
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
count_range(&ds, di).await + count_range(&ds, dd).await,
0,
"the background reaper must finish the torn space"
);
assert_eq!(count_doc_kinds(&reclaim_queue_kinds(&ds).await), 0);
let mut response =
ds.execute("DEFINE INDEX idx ON thing FIELDS val;", &ses, None).await.unwrap();
response.remove(0).result.unwrap();
let _ = index_id(&ds, ns_id, db_id, "thing", "idx").await;
}
#[tokio::test]
async fn defining_a_doc_id_index_waits_when_a_sibling_entry_has_retired() {
let page = RECLAIM_BATCH_SIZE as usize;
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE TABLE thing SCHEMALESS;
DEFINE INDEX idx ON thing FIELDS val;",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let tb = TableName::from("thing");
let di = DocLookupPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let dd = DocKeyPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
seed_range(&ds, &dd, page).await;
seed_range(&ds, &di, page).await;
ds.execute("REMOVE INDEX idx ON thing;", &ses, None).await.unwrap();
let kinds = reclaim_queue_kinds(&ds).await;
assert_eq!(
count_kind(&kinds, ReclaimKind::Index),
1,
"the removal must enqueue the index's own prefix"
);
assert_eq!(
count_doc_kinds(&kinds),
ALL_DOC_ID_KINDS.len(),
"the removal must enqueue every shared doc-ID prefix"
);
let drained = keys_in_range(&ds, dd.clone()).await;
let retired = keys_in_range(
&ds,
ReclaimTbPrefix {
kind: ReclaimKind::DocKey,
ns: ns_id,
db: db_id,
tb: Cow::Borrowed(&tb),
}
.range()
.unwrap(),
)
.await;
assert_eq!(retired.len(), 1, "`!dd` must have exactly one queue entry to retire");
let txn = ds.transaction(Write).await.unwrap();
for key in drained.into_iter().chain(retired) {
txn.del(Key::from(key)).await.unwrap();
}
txn.commit().await.unwrap();
assert_eq!(count_range(&ds, dd).await, 0, "`!dd` must be drained");
assert_eq!(count_range(&ds, di.clone()).await, page, "`!di` must be untouched");
let kinds = reclaim_queue_kinds(&ds).await;
assert_eq!(count_kind(&kinds, ReclaimKind::DocKey), 0, "`!dd`'s entry must be retired");
assert_eq!(
count_doc_kinds(&kinds),
ALL_DOC_ID_KINDS.len() - 1,
"`!dd`'s siblings must be left queued, never having been reached"
);
let queue_before = reclaim_queue_len(&ds).await;
let mut response =
ds.execute("DEFINE INDEX idx ON thing FIELDS val;", &ses, None).await.unwrap();
let error = response.remove(0).result.unwrap_err().to_string();
assert!(
error.contains("still being reclaimed"),
"the missing sibling must make the DDL wait for cleanup: {error}"
);
assert_eq!(
count_range(&ds, di.clone()).await,
page,
"the rejected DDL must not synchronously purge surviving mappings"
);
assert_eq!(reclaim_queue_len(&ds).await, queue_before);
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(count_range(&ds, di).await, 0);
let kinds = reclaim_queue_kinds(&ds).await;
assert_eq!(count_doc_kinds(&kinds), 0, "the bounded reclaim must finish");
assert_eq!(
count_kind(&kinds, ReclaimKind::Index),
kinds.len(),
"only unrelated index-prefix work may remain"
);
let mut response =
ds.execute("DEFINE INDEX idx ON thing FIELDS val;", &ses, None).await.unwrap();
response.remove(0).result.unwrap();
let _ = index_id(&ds, ns_id, db_id, "thing", "idx").await;
}
#[tokio::test]
async fn a_partly_reclaimed_doc_id_space_is_left_whole_when_a_consumer_reappears() {
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE TABLE thing SCHEMALESS;
DEFINE INDEX idx ON thing FIELDS val;
CREATE thing:1 SET val = 'a';",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let tb = TableName::from("thing");
let di = DocLookupPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let mapped = keys_in_range(&ds, di.clone()).await;
assert!(!mapped.is_empty(), "the index must have mapped the record it covers");
let tx = ds.transaction(Write).await.unwrap();
for kind in ALL_DOC_ID_KINDS {
let rc = ReclaimKey {
kind,
ns: ns_id,
db: db_id,
tb: Cow::Owned(tb.clone()),
ix: IndexId(0),
expunge: Expunge::Keep,
uid: Uuid::now_v7(),
};
let torn = ReclaimState {
observed_ms: 1,
cursor: Some(mapped[0].clone()),
};
tx.set_key(&rc, &torn).await.unwrap();
}
tx.commit().await.unwrap();
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
keys_in_range(&ds, di).await,
mapped,
"the reclaim must leave a space its table still consumes whole"
);
assert_eq!(
reclaim_queue_len(&ds).await,
0,
"a lapsed obligation must be retired rather than left to requeue forever"
);
}
#[tokio::test]
async fn a_started_doc_id_reclaim_cannot_be_cancelled() {
let page = RECLAIM_BATCH_SIZE as usize;
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE TABLE thing SCHEMALESS;",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let tb = TableName::from("thing");
let di = DocLookupPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
seed_range(&ds, &di, page * 2).await;
let rc = ReclaimKey {
kind: ReclaimKind::DocLookup,
ns: ns_id,
db: db_id,
tb: Cow::Owned(tb.clone()),
ix: IndexId(0),
expunge: Expunge::Keep,
uid: Uuid::now_v7(),
};
let tx = ds.transaction(Write).await.unwrap();
tx.set_key(&rc, &ReclaimState::enqueued()).await.unwrap();
tx.commit().await.unwrap();
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
page as u64,
CancellationToken::new(),
)
.await
.unwrap();
let after_first_pass = count_range(&ds, di.clone()).await;
assert!(
after_first_pass > 0 && after_first_pass < page * 2,
"the first pass must stop mid-prefix: {after_first_pass} left"
);
let tx = ds.transaction(Write).await.unwrap();
let claim = tx.claim_tb_doc_ids_reclaim(ns_id, db_id, &tb).await.unwrap();
tx.commit().await.unwrap();
assert_eq!(claim, DocIdsReclaimClaim::InProgress);
assert_eq!(reclaim_queue_len(&ds).await, 1, "the started entry must remain queued");
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
count_range(&ds, di).await,
0,
"the bounded reclaim must finish after the claim attempt is rejected"
);
assert_eq!(reclaim_queue_len(&ds).await, 0, "the completed entry must retire");
}
#[tokio::test]
async fn a_queued_doc_id_reclaim_is_retired_once_the_table_has_a_consumer_again() {
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE TABLE thing SCHEMALESS;
DEFINE INDEX idx ON thing FIELDS val;
CREATE thing:1 SET val = 'a';",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let tb = TableName::from("thing");
let di = DocLookupPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let mapped = keys_in_range(&ds, di.clone()).await;
assert!(!mapped.is_empty(), "the index must have mapped the record it covers");
let tx = ds.transaction(Write).await.unwrap();
for kind in ALL_DOC_ID_KINDS {
let rc = ReclaimKey {
kind,
ns: ns_id,
db: db_id,
tb: Cow::Owned(tb.clone()),
ix: IndexId(0),
expunge: Expunge::Keep,
uid: Uuid::now_v7(),
};
tx.set_key(&rc, &ReclaimState::enqueued()).await.unwrap();
}
tx.commit().await.unwrap();
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(
keys_in_range(&ds, di).await,
mapped,
"the reclaim must leave a space its table still consumes whole"
);
assert_eq!(reclaim_queue_len(&ds).await, 0, "a lapsed obligation must be retired");
}
#[tokio::test]
async fn an_orphaned_namespace_prefix_can_be_requeued_and_reclaimed() {
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE TABLE thing SCHEMALESS;
CREATE thing:1 SET val = 'a';
CREATE thing:2 SET val = 'b';",
&ses,
None,
)
.await
.unwrap();
let (ns_id, _) = ids(&ds, "test", "tenant").await;
let ns_root = NsRoot::new(ns_id).raw_range().unwrap();
assert!(count_range(&ds, ns_root.clone()).await > 0, "the namespace must hold data");
ds.execute("REMOVE NAMESPACE test;", &ses, None).await.unwrap();
let tx = ds.transaction(Write).await.unwrap();
for k in keys_in_range(&ds, ReclaimPrefix {}.range().unwrap()).await {
tx.del(Key::from(k)).await.unwrap();
}
tx.commit().await.unwrap();
assert_eq!(reclaim_queue_len(&ds).await, 0, "nothing may name the prefix");
assert!(count_range(&ds, ns_root.clone()).await > 0, "the data must outlive its catalog");
let tx = ds.transaction(Write).await.unwrap();
tx.set_key(&ReclaimKey::namespace(ns_id, false, Uuid::now_v7()), &ReclaimState::enqueued())
.await
.unwrap();
tx.commit().await.unwrap();
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(count_range(&ds, ns_root).await, 0, "the requeued prefix must be destroyed");
assert_eq!(reclaim_queue_len(&ds).await, 0, "and its entry retired");
}
async fn downgrade_index_format(ds: &Datastore, tb: &TableName, ix: &str, format_version: u16) {
let txn = ds.transaction(Write).await.unwrap();
let (ns, db) = ids(ds, "test", "tenant").await;
let def = txn.get_tb_index(ns, db, tb, ix, None).await.unwrap().expect("index should exist");
let mut old = (*def).clone();
old.format_version = format_version;
txn.put_tb_index(ns, db, tb, &old).await.unwrap();
let tb_def = txn.expect_tb(ns, db, tb).await.unwrap();
txn.put_tb("test", "tenant", &tb_def).await.unwrap();
txn.commit().await.unwrap();
}
#[tokio::test]
async fn rebuilding_a_legacy_index_cancels_a_queued_doc_id_reclaim() {
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE TABLE thing SCHEMALESS;
DEFINE INDEX legacy ON thing FIELDS a;",
&ses,
None,
)
.await
.unwrap();
let tb = TableName::from("thing");
downgrade_index_format(&ds, &tb, "legacy", 0).await;
ds.execute(
"DEFINE INDEX current ON thing FIELDS b;
CREATE thing:1 SET a = 'x', b = 'y';
CREATE thing:2 SET a = 'p', b = 'q';",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let di = DocLookupPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let mapped = keys_in_range(&ds, di.clone()).await;
assert!(!mapped.is_empty(), "the current-format index must have mapped its records");
ds.execute("REMOVE INDEX current ON thing;", &ses, None).await.unwrap();
let kinds = reclaim_queue_kinds(&ds).await;
assert_eq!(
count_doc_kinds(&kinds),
ALL_DOC_ID_KINDS.len(),
"dropping the last consumer must queue every shared doc-ID prefix"
);
ds.execute("REBUILD INDEX legacy ON thing;", &ses, None).await.unwrap();
let kinds = reclaim_queue_kinds(&ds).await;
assert_eq!(
count_doc_kinds(&kinds),
0,
"the rebuild must withdraw every queued reclaim of the space it now consumes"
);
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert!(
!keys_in_range(&ds, di).await.is_empty(),
"the reclaim must not delete mappings the rebuilt index resolves through"
);
}
#[tokio::test]
async fn rebuilding_a_legacy_index_waits_for_a_partly_reclaimed_doc_id_space() {
let page = RECLAIM_BATCH_SIZE as usize;
let ds = mem_ds().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute(
"DEFINE NAMESPACE test;
DEFINE DATABASE tenant;
DEFINE TABLE thing SCHEMALESS;
DEFINE INDEX legacy ON thing FIELDS a;",
&ses,
None,
)
.await
.unwrap();
let tb = TableName::from("thing");
downgrade_index_format(&ds, &tb, "legacy", 0).await;
ds.execute(
"DEFINE INDEX current ON thing FIELDS b;
CREATE thing:1 SET a = 'x', b = 'y';
CREATE thing:2 SET a = 'p', b = 'q';",
&ses,
None,
)
.await
.unwrap();
let (ns_id, db_id) = ids(&ds, "test", "tenant").await;
let di = DocLookupPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
let dd = DocKeyPrefix::new(ns_id, db_id, Cow::Borrowed(&tb)).raw_range().unwrap();
seed_range(&ds, &di, page * 2).await;
seed_range(&ds, &dd, page * 2).await;
ds.execute("REMOVE INDEX current ON thing;", &ses, None).await.unwrap();
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
page as u64,
CancellationToken::new(),
)
.await
.unwrap();
let left = count_range(&ds, di.clone()).await + count_range(&ds, dd.clone()).await;
let mut response = ds.execute("REBUILD INDEX legacy ON thing;", &ses, None).await.unwrap();
let error = response.remove(0).result.unwrap_err().to_string();
assert!(error.contains("still being reclaimed"));
assert_eq!(
count_range(&ds, di.clone()).await + count_range(&ds, dd.clone()).await,
left,
"the rejected rebuild must not synchronously purge the remainder"
);
let txn = ds.transaction(Read).await.unwrap();
let legacy = txn.get_tb_index(ns_id, db_id, &tb, "legacy", None).await.unwrap().unwrap();
assert!(!legacy.uses_shared_doc_ids(), "the rejected rebuild must not stamp the new format");
let _ = txn.cancel().await;
Datastore::reclaim_tombstones(
Arc::clone(&ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(count_range(&ds, di).await + count_range(&ds, dd).await, 0);
let mut response = ds.execute("REBUILD INDEX legacy ON thing;", &ses, None).await.unwrap();
response.remove(0).result.unwrap();
let txn = ds.transaction(Read).await.unwrap();
let legacy = txn.get_tb_index(ns_id, db_id, &tb, "legacy", None).await.unwrap().unwrap();
assert!(legacy.uses_shared_doc_ids(), "the retry must stamp the current format");
let _ = txn.cancel().await;
}
async fn queue_database_entry(ds: &Datastore, db: u32, observed_ms: u64, keys: usize) -> RawRange {
let range = DbRoot {
ns: NamespaceId(1),
db: DatabaseId(db),
}
.range()
.unwrap();
if keys > 0 {
seed_range(ds, &range, keys).await;
}
let rc = ReclaimKey::database(NamespaceId(1), DatabaseId(db), false, Uuid::now_v7());
let tx = ds.transaction(Write).await.unwrap();
tx.set_key(
&rc,
&ReclaimState {
observed_ms,
cursor: None,
},
)
.await
.unwrap();
tx.commit().await.unwrap();
range
}
async fn queue_observations(ds: &Datastore) -> Vec<u64> {
let txn = ds.transaction(Read).await.unwrap();
let items: Vec<(Vec<u8>, ReclaimState)> =
txn.getr(ReclaimPrefix {}.range().unwrap(), None).await.unwrap();
let _ = txn.cancel().await;
items.iter().map(|(_, state)| state.observed_ms).collect()
}
#[tokio::test]
async fn the_walk_stamps_entries_the_budget_cannot_reach() {
let page = RECLAIM_BATCH_SIZE as usize;
let ds = mem_ds().await;
let grace = Duration::from_secs(600);
let aged = now_ms().saturating_sub(2 * grace.as_millis() as u64);
for db in 1..=2u32 {
queue_database_entry(&ds, db, aged, page * 2).await;
}
for db in 10..=12u32 {
queue_database_entry(&ds, db, 0, page).await;
}
assert_eq!(reclaim_queue_len(&ds).await, 5);
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
grace,
page as u64,
CancellationToken::new(),
)
.await
.unwrap();
let observed = queue_observations(&ds).await;
assert_eq!(observed.len(), 5, "no entry may be retired: none of the prefixes is empty");
assert!(
observed.iter().all(|&ms| ms != 0),
"every entry the pass walked past must be stamped, budget or not: {observed:?}"
);
}
#[tokio::test]
async fn the_budget_goes_to_the_longest_waiting_entry() {
let page = RECLAIM_BATCH_SIZE as usize;
let ds = mem_ds().await;
let grace = Duration::from_secs(600);
let base = now_ms().saturating_sub(2 * grace.as_millis() as u64);
let newest = queue_database_entry(&ds, 1, base + 60_000, page * 2).await;
let middle = queue_database_entry(&ds, 2, base + 30_000, page * 2).await;
let oldest = queue_database_entry(&ds, 3, base, page * 2).await;
let before = [
count_range(&ds, newest.clone()).await,
count_range(&ds, middle.clone()).await,
count_range(&ds, oldest.clone()).await,
];
Datastore::reclaim_tombstones_with_budget(
Arc::clone(&ds),
Duration::from_secs(1),
grace,
1,
CancellationToken::new(),
)
.await
.unwrap();
let deleted = [
before[0] - count_range(&ds, newest).await,
before[1] - count_range(&ds, middle).await,
before[2] - count_range(&ds, oldest).await,
];
assert_eq!(
deleted.iter().sum::<usize>(),
1,
"the pass must spend its whole budget and no more: {deleted:?}"
);
assert_eq!(
deleted,
[0, 0, 1],
"the budget must go to the oldest observation, not the lowest key: {deleted:?}"
);
}