use std::time::Duration;
use tracing::debug;
use crate::error::RuntimeError;
use crate::runtime::Runtime;
use crate::state::EffectStatus;
use crate::store::{EffectStore, PruneQuery};
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct RetentionPolicy {
committed: Option<Duration>,
failed: Option<Duration>,
rejected: Option<Duration>,
compensated: Option<Duration>,
}
impl RetentionPolicy {
pub const KEEP_ALL: Self = Self {
committed: None,
failed: None,
rejected: None,
compensated: None,
};
pub const fn settled(older_than: Duration) -> Self {
Self {
committed: Some(older_than),
failed: Some(older_than),
rejected: Some(older_than),
compensated: Some(older_than),
}
}
#[must_use]
pub const fn committed(mut self, older_than: Duration) -> Self {
self.committed = Some(older_than);
self
}
#[must_use]
pub const fn failed(mut self, older_than: Duration) -> Self {
self.failed = Some(older_than);
self
}
#[must_use]
pub const fn rejected(mut self, older_than: Duration) -> Self {
self.rejected = Some(older_than);
self
}
#[must_use]
pub const fn compensated(mut self, older_than: Duration) -> Self {
self.compensated = Some(older_than);
self
}
pub const fn keeps_all(&self) -> bool {
self.committed.is_none()
&& self.failed.is_none()
&& self.rejected.is_none()
&& self.compensated.is_none()
}
pub fn rules(&self) -> impl Iterator<Item = (EffectStatus, Duration)> {
[
(EffectStatus::Committed, self.committed),
(EffectStatus::Failed, self.failed),
(EffectStatus::Rejected, self.rejected),
(EffectStatus::Compensated, self.compensated),
]
.into_iter()
.filter_map(|(status, age)| Some((status, age?)))
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
#[non_exhaustive]
pub struct PruneReport {
pub committed: u64,
pub failed: u64,
pub rejected: u64,
pub compensated: u64,
}
impl PruneReport {
pub const fn total(&self) -> u64 {
self.committed + self.failed + self.rejected + self.compensated
}
fn add(&mut self, status: EffectStatus, deleted: u64) {
match status {
EffectStatus::Committed => self.committed += deleted,
EffectStatus::Failed => self.failed += deleted,
EffectStatus::Rejected => self.rejected += deleted,
EffectStatus::Compensated => self.compensated += deleted,
_ => {}
}
}
}
impl<S: EffectStore> Runtime<S> {
pub async fn prune(&self) -> Result<PruneReport, RuntimeError> {
let mut report = PruneReport::default();
let now = self.now();
for (status, older_than) in self.retention().rules() {
let query = PruneQuery::new(status, older_than, now);
loop {
let deleted = self.store().prune(query.clone()).await?;
report.add(status, deleted);
if deleted < query.limit as u64 {
break;
}
}
}
if report.total() > 0 {
debug!(?report, "pruned settled effects");
}
Ok(report)
}
}