use std::ops::Range;
use std::time::Duration;
use anyhow::Result;
use tokio_util::sync::CancellationToken;
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider, TableProvider};
use crate::catalog::{DatabaseDefinition, DatabaseId, NamespaceId, TableDefinition};
use crate::key::change;
use crate::kvs::LockType::Optimistic;
use crate::kvs::TransactionType::Read;
use crate::kvs::paging::{PageCompanion, PagedDelete, PagedOutcome};
use crate::kvs::tasklease::LeaseHandler;
use crate::kvs::{CHANGEFEED_GC_BATCH_SIZE, Datastore, KVKey, Key, Transaction};
struct Target<R> {
ns: NamespaceId,
db: DatabaseId,
stale: R,
watermark: Vec<u8>,
}
#[instrument(level = "trace", target = "surrealdb::core::cfs", skip_all)]
pub async fn gc_all_at(
ds: &Datastore,
lh: &LeaseHandler,
canceller: &CancellationToken,
budget: &mut u64,
) -> Result<bool> {
let targets = targets(ds).await?;
let mut remaining = targets.len() as u64;
for target in targets {
Datastore::ensure_not_cancelled(canceller)?;
trace!("Performing changefeed garbage collection on {}:{}", target.ns, target.db);
let share = (*budget / remaining.max(1)).max(1).min(*budget);
remaining = remaining.saturating_sub(1);
if share == 0 {
break;
}
let mut left = share;
let outcome = PagedDelete {
window: target.stale,
companion: PageCompanion::ChangefeedRetention {
ns: target.ns,
db: target.db,
watermark: &target.watermark,
},
expunge: false,
page: CHANGEFEED_GC_BATCH_SIZE,
}
.run(ds, lh, canceller, &mut left)
.await?;
*budget = budget.saturating_sub(share.saturating_sub(left));
if outcome == PagedOutcome::LeaseLost {
return Ok(false);
}
if *budget == 0 {
break;
}
yield_now!();
}
Ok(true)
}
async fn targets(ds: &Datastore) -> Result<Vec<Target<Range<Key>>>> {
let txn = ds.transaction(Read, Optimistic).await?;
let res = async {
let ts_impl = txn.timestamp_impl();
let ts = txn.timestamp().await?;
let mut buf = [0u8; _];
let earliest = ts_impl.earliest().encode(&mut buf).to_vec();
let mut targets = Vec::new();
for ns in txn.all_ns(None).await?.as_ref() {
for db in txn.all_db(ns.namespace_id, None).await?.as_ref() {
let tbs = txn.all_tb(db.namespace_id, db.database_id, None).await?;
let expiry = retention_of(db, tbs.as_ref());
if expiry.is_zero() {
continue;
}
let watermark = ts.sub_checked(expiry).unwrap_or_else(|| ts_impl.earliest());
let mut buf = [0u8; _];
let end = watermark.encode(&mut buf);
let beg = change::prefix_ts(db.namespace_id, db.database_id, earliest.as_slice())
.encode_key()?;
let fin = change::prefix_ts(db.namespace_id, db.database_id, end).encode_key()?;
let stale = beg..fin;
targets.push(Target {
ns: db.namespace_id,
db: db.database_id,
stale,
watermark: end.to_vec(),
});
}
}
Ok(targets)
}
.await;
let _ = txn.cancel().await;
res
}
fn retention_of(db: &DatabaseDefinition, tbs: &[TableDefinition]) -> Duration {
let db_expiry = db.changefeed.map(|v| v.expiry).unwrap_or_default();
let tb_expiry = tbs
.iter()
.filter_map(|tb| tb.changefeed.as_ref())
.map(|cf| cf.expiry)
.filter(|&dur| !dur.is_zero())
.max()
.unwrap_or(Duration::ZERO);
db_expiry.max(tb_expiry)
}
pub(crate) async fn retention_still_reaches(
txn: &Transaction,
ns: NamespaceId,
db: DatabaseId,
watermark: &[u8],
) -> Result<bool> {
txn.arm_changefeed_retention_fence(ns, db).await?;
let dbs = txn.all_db(ns, None).await?;
let Some(def) = dbs.as_ref().iter().find(|d| d.database_id == db) else {
return Ok(false);
};
let tbs = txn.all_tb(ns, db, None).await?;
let expiry = retention_of(def, tbs.as_ref());
if expiry.is_zero() {
return Ok(false);
}
let ts_impl = txn.timestamp_impl();
let ts = txn.timestamp().await?;
let current = ts.sub_checked(expiry).unwrap_or_else(|| ts_impl.earliest());
let mut buf = [0u8; _];
Ok(current.encode(&mut buf) >= watermark)
}