use std::time::Duration;
use anyhow::Result;
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider};
use crate::catalog::{DatabaseId, NamespaceId};
use crate::key::lqe;
use crate::kvs::tasklease::LeaseHandler;
use crate::kvs::{BoxTimeStamp, BoxTimeStampImpl, KVKey, Transaction};
#[instrument(level = "trace", target = "surrealdb::core::lq", skip_all)]
pub async fn gc_all_at(lh: &LeaseHandler, tx: &Transaction, retention: Duration) -> Result<()> {
if retention.is_zero() {
return Ok(());
}
let ts_impl = tx.timestamp_impl();
let nss = tx.all_ns(None).await?;
for ns in nss.as_ref() {
let dbs = tx.all_db(ns.namespace_id, None).await?;
for db in dbs.as_ref() {
let ts = tx.timestamp().await?;
let watermark_ts = ts.sub_checked(retention).unwrap_or_else(|| ts_impl.earliest());
gc_range(tx, db.namespace_id, db.database_id, &watermark_ts, &ts_impl).await?;
lh.try_maintain_lease().await?;
yield_now!();
}
lh.try_maintain_lease().await?;
yield_now!();
}
Ok(())
}
async fn gc_range(
tx: &Transaction,
ns: NamespaceId,
db: DatabaseId,
ts: &BoxTimeStamp,
ts_impl: &BoxTimeStampImpl,
) -> Result<()> {
let mut buf = [0u8; _];
let beg_ts = ts_impl.earliest().encode(&mut buf);
let mut buf = [0u8; _];
let end_ts = ts.encode(&mut buf);
let beg = lqe::prefix_ts(ns, db, beg_ts).encode_key()?;
let end = lqe::prefix_ts(ns, db, end_ts).encode_key()?;
tx.delr(beg..end).await?;
Ok(())
}