use std::ops::Range;
use std::sync::Arc;
use tokio_util::sync::CancellationToken;
use web_time::Duration;
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider};
use crate::catalog::{DatabaseId, NamespaceId};
use crate::cf::gc::retention_still_reaches;
use crate::dbs::{Capabilities, Session};
use crate::kvs::LockType::Optimistic;
use crate::kvs::TransactionType::{Read, Write};
use crate::kvs::tasklease::{LeaseHandler, TaskLeaseType};
use crate::kvs::test_support::{WriteTxSizeObserver, cleanup_sizes};
use crate::kvs::{CHANGEFEED_GC_BATCH_SIZE, Datastore, KVKey};
use crate::observe::ExecutionObserver;
const STALE: &str = "1ns";
async fn observed_ds(observer: Arc<WriteTxSizeObserver>) -> Arc<Datastore> {
Arc::new(
Datastore::builder()
.with_capabilities(Capabilities::all())
.with_observer(observer as Arc<dyn ExecutionObserver>)
.build_with_path("memory")
.await
.unwrap(),
)
}
async fn feed(ds: &Datastore, ns: &str, db: &str) -> Range<Vec<u8>> {
let (ns_id, db_id) = db_ids(ds, ns, db).await;
let beg = crate::key::change::prefix(ns_id, db_id).encode_key().unwrap();
let end = crate::key::change::suffix(ns_id, db_id).encode_key().unwrap();
beg..end
}
async fn count_range(ds: &Datastore, range: Range<Vec<u8>>) -> usize {
let txn = ds.transaction(Read, Optimistic).await.unwrap();
let keys = txn.keys(range, u32::MAX, 0, None).await.unwrap();
let _ = txn.cancel().await;
keys.len()
}
async fn grow_backlog(ds: &Datastore, feed: &Range<Vec<u8>>, count: usize) {
let txn = ds.transaction(Read, Optimistic).await.unwrap();
let base = txn.keys(feed.clone(), 1, 0, None).await.unwrap().remove(0);
let _ = txn.cancel().await;
let tx = ds.transaction(Write, Optimistic).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, &vec![1u8]).await.unwrap();
}
tx.commit().await.unwrap();
}
async fn with_changefeed(ds: &Datastore, ns: &str, db: &str) {
let ses = Session::owner().with_ns(ns).with_db(db);
for res in ds
.execute(
&format!(
"DEFINE DATABASE {db};
DEFINE TABLE thing CHANGEFEED {STALE};
CREATE thing:1;"
),
&ses,
None,
)
.await
.unwrap()
{
res.result.unwrap();
}
}
async fn watermark_now(ds: &Datastore, expiry: std::time::Duration) -> Vec<u8> {
let txn = ds.transaction(Read, Optimistic).await.unwrap();
let ts_impl = txn.timestamp_impl();
let ts = txn.timestamp().await.unwrap();
let _ = txn.cancel().await;
let watermark = ts.sub_checked(expiry).unwrap_or_else(|| ts_impl.earliest());
let mut buf = [0u8; 32];
watermark.encode(&mut buf).to_vec()
}
async fn db_ids(ds: &Datastore, ns: &str, db: &str) -> (NamespaceId, DatabaseId) {
let txn = ds.transaction(Read, Optimistic).await.unwrap();
let ns_def = txn.get_ns_by_name(ns, None).await.unwrap().unwrap();
let db_def = txn.get_db_by_name(ns, db, None).await.unwrap().unwrap();
let _ = txn.cancel().await;
(ns_def.namespace_id, db_def.database_id)
}
async fn fence(ds: &Datastore, ns: NamespaceId, db: DatabaseId) -> Option<u64> {
let txn = ds.transaction(Read, Optimistic).await.unwrap();
let v = txn.get(&crate::key::database::cf::new(ns, db), None).await.unwrap();
let _ = txn.cancel().await;
v
}
#[tokio::test]
async fn changefeed_gc_deletes_the_backlog_in_bounded_committed_pages() {
let page = CHANGEFEED_GC_BATCH_SIZE as usize;
let observer = Arc::new(WriteTxSizeObserver::default());
let ds = observed_ds(Arc::clone(&observer)).await;
ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
with_changefeed(&ds, "test", "tenant").await;
let feed = feed(&ds, "test", "tenant").await;
grow_backlog(&ds, &feed, page * 2 + page / 2).await;
let total = count_range(&ds, feed.clone()).await;
assert!(total > 2 * page, "the backlog must hold more than two pages, holds {total}");
observer.clear();
ds.changefeed_process(&Duration::from_secs(1), &CancellationToken::new()).await.unwrap();
assert_eq!(count_range(&ds, feed).await, 0, "the collector must empty the backlog");
let sizes = observer.sizes();
assert!(
sizes.iter().all(|&n| n <= page as u64 + 1),
"a collection write transaction carried more than one page plus its fence arm: {sizes:?}"
);
assert!(
sizes.iter().sum::<u64>() >= total as u64,
"every collected entry must be charged as a write: {sizes:?}"
);
assert!(
cleanup_sizes(&sizes).len() >= total / page,
"the backlog must be spread across a transaction per page: {sizes:?}"
);
}
#[tokio::test]
async fn changefeed_gc_shares_its_budget_across_databases() {
let page = CHANGEFEED_GC_BATCH_SIZE as usize;
let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
with_changefeed(&ds, "test", "aaa").await;
with_changefeed(&ds, "test", "zzz").await;
let aaa = feed(&ds, "test", "aaa").await;
let zzz = feed(&ds, "test", "zzz").await;
grow_backlog(&ds, &aaa, page * 4).await;
grow_backlog(&ds, &zzz, page / 2).await;
assert!(count_range(&ds, zzz.clone()).await > 0, "the trailing database needs a backlog");
ds.changefeed_gc_with_budget(
&Duration::from_secs(1),
(page * 2) as u64,
&CancellationToken::new(),
)
.await
.unwrap();
assert!(
count_range(&ds, aaa).await > 0,
"a budget smaller than the leading backlog must leave part of it behind"
);
assert_eq!(
count_range(&ds, zzz).await,
0,
"the trailing database must be collected in the same pass"
);
}
#[tokio::test]
async fn changefeed_gc_makes_durable_progress_across_passes() {
let page = CHANGEFEED_GC_BATCH_SIZE as usize;
let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
with_changefeed(&ds, "test", "tenant").await;
let feed = feed(&ds, "test", "tenant").await;
grow_backlog(&ds, &feed, page * 3).await;
let mut left = count_range(&ds, feed.clone()).await;
for _ in 0..10 {
ds.changefeed_gc_with_budget(
&Duration::from_secs(1),
page as u64,
&CancellationToken::new(),
)
.await
.unwrap();
let now = count_range(&ds, feed.clone()).await;
assert!(now < left, "a pass must make progress: {left} then {now}");
left = now;
if left == 0 {
break;
}
}
assert_eq!(left, 0, "repeated bounded passes must drain the backlog");
}
#[tokio::test]
async fn an_extended_retention_withdraws_the_stale_bound() {
let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
with_changefeed(&ds, "test", "tenant").await;
let (ns, db) = db_ids(&ds, "test", "tenant").await;
let bound = watermark_now(&ds, std::time::Duration::from_nanos(1)).await;
let txn = ds.transaction(Write, Optimistic).await.unwrap();
assert!(
retention_still_reaches(&txn, ns, db, &bound).await.unwrap(),
"an unchanged retention must keep authorising the range it made stale"
);
let _ = txn.cancel().await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("DEFINE TABLE OVERWRITE thing CHANGEFEED 1h;", &ses, None).await.unwrap()[0]
.result
.as_ref()
.unwrap();
let txn = ds.transaction(Write, Optimistic).await.unwrap();
assert!(
!retention_still_reaches(&txn, ns, db, &bound).await.unwrap(),
"an extended retention must withdraw the bound a pass started from"
);
let _ = txn.cancel().await;
}
#[tokio::test]
async fn a_removed_retention_withdraws_the_stale_bound() {
let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
with_changefeed(&ds, "test", "tenant").await;
let (ns, db) = db_ids(&ds, "test", "tenant").await;
let bound = watermark_now(&ds, std::time::Duration::from_nanos(1)).await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("DEFINE TABLE OVERWRITE thing;", &ses, None).await.unwrap()[0]
.result
.as_ref()
.unwrap();
let txn = ds.transaction(Write, Optimistic).await.unwrap();
assert!(
!retention_still_reaches(&txn, ns, db, &bound).await.unwrap(),
"a database with no retention left has nothing this collector may delete"
);
let _ = txn.cancel().await;
ds.execute("REMOVE DATABASE tenant;", &ses, None).await.unwrap()[0].result.as_ref().unwrap();
let txn = ds.transaction(Write, Optimistic).await.unwrap();
assert!(
!retention_still_reaches(&txn, ns, db, &bound).await.unwrap(),
"a removed database is reclaimed by its tombstone, not collected here"
);
let _ = txn.cancel().await;
}
#[tokio::test]
async fn a_pass_that_loses_the_lease_reports_it() {
let page = CHANGEFEED_GC_BATCH_SIZE as usize;
let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
with_changefeed(&ds, "test", "tenant").await;
let feed = feed(&ds, "test", "tenant").await;
grow_backlog(&ds, &feed, page).await;
let backlog = count_range(&ds, feed.clone()).await;
assert!(backlog > 0, "there must be a backlog for the pass to page through");
let held = LeaseHandler::new(
ds.sequences().clone(),
uuid::Uuid::now_v7(),
ds.transaction_factory().clone(),
TaskLeaseType::ChangeFeedCleanup,
std::time::Duration::from_secs(60),
)
.unwrap();
assert!(held.has_lease().await.unwrap(), "the other node must take the lease");
let ours = LeaseHandler::new(
ds.sequences().clone(),
uuid::Uuid::now_v7(),
ds.transaction_factory().clone(),
TaskLeaseType::ChangeFeedCleanup,
std::time::Duration::from_secs(60),
)
.unwrap();
let mut budget = u64::MAX;
let held_after =
crate::cf::gc::gc_all_at(&ds, &ours, &CancellationToken::new(), &mut budget).await.unwrap();
assert!(!held_after, "a pass that lost the lease must report the loss to its caller");
assert_eq!(
count_range(&ds, feed).await,
backlog,
"a pass without the lease must not collect anything"
);
}
#[tokio::test]
async fn a_completed_pass_reports_the_lease_still_held() {
let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
with_changefeed(&ds, "test", "tenant").await;
let feed = feed(&ds, "test", "tenant").await;
grow_backlog(&ds, &feed, CHANGEFEED_GC_BATCH_SIZE as usize).await;
let ours = LeaseHandler::new(
ds.sequences().clone(),
uuid::Uuid::now_v7(),
ds.transaction_factory().clone(),
TaskLeaseType::ChangeFeedCleanup,
std::time::Duration::from_secs(60),
)
.unwrap();
assert!(ours.has_lease().await.unwrap(), "this node must hold the lease");
let mut budget = u64::MAX;
assert!(
crate::cf::gc::gc_all_at(&ds, &ours, &CancellationToken::new(), &mut budget).await.unwrap(),
"a pass that kept the lease must report it, so the caller runs the next collector"
);
assert_eq!(count_range(&ds, feed).await, 0, "the pass must collect the backlog");
}
#[tokio::test]
async fn every_retention_write_bumps_the_fence() {
let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
ds.execute("DEFINE DATABASE tenant CHANGEFEED 1h;", &ses, None).await.unwrap()[0]
.result
.as_ref()
.unwrap();
let (ns, db) = db_ids(&ds, "test", "tenant").await;
let after_db = fence(&ds, ns, db).await.expect("DEFINE DATABASE must bump the fence");
ds.execute("DEFINE TABLE thing CHANGEFEED 1h;", &ses, None).await.unwrap()[0]
.result
.as_ref()
.unwrap();
let after_tb = fence(&ds, ns, db).await.unwrap();
assert!(after_tb > after_db, "DEFINE TABLE must bump the fence: {after_db} -> {after_tb}");
ds.execute("ALTER TABLE thing CHANGEFEED 2h;", &ses, None).await.unwrap()[0]
.result
.as_ref()
.unwrap();
let after_alter = fence(&ds, ns, db).await.unwrap();
assert!(after_alter > after_tb, "ALTER TABLE must bump the fence: {after_tb} -> {after_alter}");
ds.execute("DEFINE TABLE other;", &ses, None).await.unwrap()[0].result.as_ref().unwrap();
let before_fresh = fence(&ds, ns, db).await.unwrap();
ds.execute("ALTER TABLE other CHANGEFEED 3h;", &ses, None).await.unwrap()[0]
.result
.as_ref()
.unwrap();
let after_fresh = fence(&ds, ns, db).await.unwrap();
assert!(
after_fresh > before_fresh,
"setting a retention on a table that had none must bump the fence: \
{before_fresh} -> {after_fresh}"
);
ds.execute("ALTER TABLE other DROP CHANGEFEED;", &ses, None).await.unwrap()[0]
.result
.as_ref()
.unwrap();
assert_eq!(
fence(&ds, ns, db).await.unwrap(),
after_fresh,
"dropping a retention needs no fence bump: it strands nothing"
);
}
#[tokio::test]
async fn a_retention_change_rejects_an_open_collection_page() {
let ds = observed_ds(Arc::new(WriteTxSizeObserver::default())).await;
ds.execute("DEFINE NAMESPACE test;", &Session::owner(), None).await.unwrap();
with_changefeed(&ds, "test", "tenant").await;
let (ns, db) = db_ids(&ds, "test", "tenant").await;
let bound = watermark_now(&ds, std::time::Duration::from_nanos(1)).await;
let page = ds.transaction(Write, Optimistic).await.unwrap();
assert!(
retention_still_reaches(&page, ns, db, &bound).await.unwrap(),
"the retention must still authorise the range before the change lands"
);
let ses = Session::owner().with_ns("test").with_db("tenant");
ds.execute("ALTER TABLE thing CHANGEFEED 1h;", &ses, None).await.unwrap()[0]
.result
.as_ref()
.unwrap();
page.set(&b"unused".to_vec(), &vec![1u8]).await.unwrap();
page.commit()
.await
.expect_err("a page open across a retention change must be rejected at commit");
}