use anyhow::Result;
use surrealdb_kvs::TransactionType::{Read, Write};
use tokio_util::sync::CancellationToken;
use crate::catalog::{DatabaseId, NamespaceId};
use crate::key::reclaim::ReclaimState;
use crate::key::schema::ReclaimKey;
use crate::key::{AnyRange, Key, Resumable};
use crate::kvs::tasklease::LeaseHandler;
use crate::kvs::{Datastore, Direction};
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum PagedOutcome {
Complete,
Incomplete,
LeaseLost,
Cancelled,
}
pub(crate) enum PageCompanion<'a> {
None,
ReclaimEntry {
rc: &'a ReclaimKey<'a>,
state: ReclaimState,
},
ChangefeedRetention {
ns: NamespaceId,
db: DatabaseId,
watermark: &'a [u8],
},
}
pub(crate) struct PagedDelete<'a, R> {
pub(crate) window: R,
pub(crate) companion: PageCompanion<'a>,
pub(crate) expunge: bool,
pub(crate) page: u32,
}
impl<R: AnyRange + Resumable + Clone> PagedDelete<'_, R> {
pub(crate) async fn run(
mut self,
ds: &Datastore,
lh: &LeaseHandler,
canceller: &CancellationToken,
budget: &mut u64,
) -> Result<PagedOutcome> {
loop {
Datastore::ensure_not_cancelled(canceller)?;
if *budget == 0 {
let txn = ds.transaction(Read).await?;
let rest = txn.keys_raw(self.window.clone(), 1, 0, None).await;
let _ = txn.cancel().await;
return Ok(match rest?.is_empty() {
true => PagedOutcome::Complete,
false => PagedOutcome::Incomplete,
});
}
if !lh.try_maintain_lease().await? {
return Ok(PagedOutcome::LeaseLost);
}
let limit = (*budget).min(self.page.max(1) as u64) as u32;
let txn = ds.transaction(Write).await?;
if let PageCompanion::ChangefeedRetention {
ns,
db,
watermark,
} = &self.companion
{
let covered = catch!(
txn,
crate::cf::gc::retention_still_reaches(&txn, *ns, *db, watermark).await
);
if !covered {
let _ = txn.cancel().await;
return Ok(PagedOutcome::Cancelled);
}
}
let keys = catch!(txn, txn.keys_raw(self.window.clone(), limit, 0, None).await);
let exhausted = keys.len() < limit as usize;
let Some(last) = keys.last().cloned() else {
let _ = txn.cancel().await;
return Ok(PagedOutcome::Complete);
};
for key in &keys {
if self.expunge {
catch!(txn, txn.clr(Key::from(key)).await);
} else {
catch!(txn, txn.del(Key::from(key)).await);
}
}
let advanced = if let PageCompanion::ReclaimEntry {
rc,
state,
} = &self.companion
{
let advanced = ReclaimState {
observed_ms: state.observed_ms,
cursor: Some(last.clone()),
};
match txn.put_compare_key(*rc, &advanced, Some(state)).await {
Ok(()) => Some(advanced),
Err(e) if super::is_conditional_write_conflict(&e) => {
let _ = txn.cancel().await;
return Ok(PagedOutcome::Cancelled);
}
Err(e) => {
let _ = txn.cancel().await;
return Err(e);
}
}
} else {
None
};
match txn.commit().await {
Ok(()) => {}
Err(e) if super::is_conditional_write_conflict(&e) => {
let _ = txn.cancel().await;
return Ok(PagedOutcome::Cancelled);
}
Err(e) => {
let _ = txn.cancel().await;
return Err(e);
}
}
if let (
Some(advanced),
PageCompanion::ReclaimEntry {
state,
..
},
) = (advanced, &mut self.companion)
{
*state = advanced;
}
*budget = budget.saturating_sub(keys.len() as u64);
if exhausted {
return Ok(PagedOutcome::Complete);
}
self.window = self.window.resume_after(&last, Direction::Forward);
yield_now!();
}
}
}