use std::ops::Range;
use std::time::Duration;
use anyhow::Result;
use tokio_util::sync::CancellationToken;
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider};
use crate::key::lqe;
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};
#[instrument(level = "trace", target = "surrealdb::core::lq", skip_all)]
pub async fn gc_all_at(
ds: &Datastore,
lh: &LeaseHandler,
retention: Duration,
canceller: &CancellationToken,
budget: &mut u64,
) -> Result<()> {
if retention.is_zero() {
return Ok(());
}
let stale = stale_ranges(ds, retention).await?;
let mut remaining = stale.len() as u64;
for range in stale {
Datastore::ensure_not_cancelled(canceller)?;
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: range,
companion: PageCompanion::None,
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 || *budget == 0 {
break;
}
yield_now!();
}
Ok(())
}
async fn stale_ranges(ds: &Datastore, retention: Duration) -> Result<Vec<Range<Key>>> {
let txn = ds.transaction(Read, Optimistic).await?;
let res = async {
let ts_impl = txn.timestamp_impl();
let ts = txn.timestamp().await?;
let watermark = ts.sub_checked(retention).unwrap_or_else(|| ts_impl.earliest());
let mut buf = [0u8; _];
let beg = ts_impl.earliest().encode(&mut buf).to_vec();
let mut buf = [0u8; _];
let end = watermark.encode(&mut buf);
let mut ranges = 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 from =
lqe::prefix_ts(db.namespace_id, db.database_id, beg.as_slice()).encode_key()?;
let to = lqe::prefix_ts(db.namespace_id, db.database_id, end).encode_key()?;
ranges.push(from..to);
}
}
Ok(ranges)
}
.await;
let _ = txn.cancel().await;
res
}