use std::future::Future;
use std::pin::Pin;
use std::time::Duration;
use taquba::{Clock, Expired, ExpiryIndex, Queue, SettlementEffects};
use tokio_util::sync::CancellationToken;
use tracing::warn;
use crate::error::Result;
use crate::keys::RunId;
type ClearError = Box<dyn std::error::Error + Send + Sync>;
type ClearFuture<'a> =
Pin<Box<dyn Future<Output = std::result::Result<Vec<Vec<u8>>, ClearError>> + Send + 'a>>;
pub(crate) trait Clearable: Send + Sync + 'static {
type Error: Into<ClearError>;
fn clear(
&self,
id: &RunId,
) -> impl Future<Output = std::result::Result<Vec<Vec<u8>>, Self::Error>> + Send;
}
trait DynClearable: Send + Sync {
fn clear_dyn<'a>(&'a self, id: &'a RunId) -> ClearFuture<'a>;
}
impl<C: Clearable> DynClearable for C {
fn clear_dyn<'a>(&'a self, id: &'a RunId) -> ClearFuture<'a> {
Box::pin(async move { self.clear(id).await.map_err(Into::into) })
}
}
pub(crate) struct Sweep {
index: ExpiryIndex,
retention: Duration,
store: Box<dyn DynClearable>,
}
impl Sweep {
pub(crate) fn new(prefix: &'static [u8], retention: Duration, store: impl Clearable) -> Self {
Self {
index: ExpiryIndex::new(prefix),
retention,
store: Box::new(store),
}
}
#[cfg(test)]
pub(crate) fn marker_key(&self, id: &RunId, at_ms: u64) -> Vec<u8> {
self.index.entry_key(at_ms, id.as_str().as_bytes())
}
pub(crate) fn mark(
&self,
effects: SettlementEffects,
id: &RunId,
at_ms: u64,
) -> SettlementEffects {
effects.expiry_entry(&self.index, at_ms, id.as_str().as_bytes())
}
pub(crate) async fn run(
&self,
queue: &Queue,
clock: &dyn Clock,
interval: Duration,
stop: CancellationToken,
) {
run_periodically(interval, &stop, (), |()| async move {
if let Err(err) = self.pass(queue, clock).await {
warn!("retention sweep failed: {err}");
}
})
.await;
}
pub(crate) async fn pass(&self, queue: &Queue, clock: &dyn Clock) -> Result<usize> {
let now_ms = clock.now_ms();
let store = &self.store;
let removed = self
.index
.pass(queue, now_ms, self.retention, |_, suffix| async move {
let id = std::str::from_utf8(&suffix)
.ok()
.and_then(|id| RunId::new(id).ok());
let Some(id) = id else {
warn!(
suffix = %String::from_utf8_lossy(&suffix),
"marker without a run id; deleting without clearing",
);
return Expired::Delete(SettlementEffects::default());
};
match store.clear_dyn(&id).await {
Ok(kv_deletes) => {
Expired::Delete(SettlementEffects::default().kv_deletes(kv_deletes))
}
Err(err) => {
warn!(id = %id, "clear failed during sweep: {err}");
Expired::Keep
}
}
})
.await?;
Ok(removed)
}
}
pub(crate) async fn run_periodically<S, Fut>(
interval: Duration,
stop: &CancellationToken,
mut state: S,
mut pass: impl FnMut(S) -> Fut,
) where
Fut: Future<Output = S>,
{
loop {
state = pass(state).await;
tokio::select! {
_ = stop.cancelled() => return,
_ = tokio::time::sleep(interval) => {}
}
}
}