#![forbid(unsafe_code)]
use lunaris_core::retention::{RetentionPolicy, retention_policy_key};
use lunaris_core::{Hlc, LunarisError, Scope, StoragePort, WriteOp};
use crate::forget::{ForgetReceipt, ForgetTarget};
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub struct RetentionReceipt {
pub policy: Option<RetentionPolicy>,
pub cutoff: Option<Hlc>,
pub forget: Option<ForgetReceipt>,
}
impl RetentionReceipt {
pub fn rows_swept(&self) -> u64 {
match &self.forget {
Some(r) => r.rows_written + r.rows_deleted,
None => 0,
}
}
}
pub async fn read_policy(
storage: &std::sync::Arc<dyn StoragePort>,
scope: &Scope,
clock: &lunaris_core::HlcClock,
) -> Result<Option<RetentionPolicy>, LunarisError> {
let key = retention_policy_key(scope);
let row = storage.read_as_of(scope, &key, clock.tick()).await.map_err(LunarisError::Storage)?;
let Some(row) = row else { return Ok(None) };
match serde_json::from_slice::<RetentionPolicy>(&row.value) {
Ok(p) => Ok(Some(p)),
Err(e) => Err(LunarisError::Storage(lunaris_core::error::StorageError::Backend(format!(
"retention policy for scope `{}` did not parse: {e}",
scope.as_str()
)))),
}
}
pub async fn write_policy(
storage: &std::sync::Arc<dyn StoragePort>,
scope: &Scope,
policy: RetentionPolicy,
) -> Result<(), LunarisError> {
let value = serde_json::to_vec(&policy).map_err(|e| {
LunarisError::Storage(lunaris_core::error::StorageError::Backend(format!(
"retention policy serialize: {e}"
)))
})?;
storage
.atomic_write(scope, &[WriteOp::KvPut { key: retention_policy_key(scope), value }])
.await
.map_err(LunarisError::Storage)?;
Ok(())
}
pub async fn enforce_at(
engine: &crate::handle::Lunaris,
scope: &Scope,
now_ms: u64,
) -> Result<RetentionReceipt, LunarisError> {
run(engine, scope, now_ms, Mode::Commit).await
}
pub async fn preview_at(
engine: &crate::handle::Lunaris,
scope: &Scope,
now_ms: u64,
) -> Result<RetentionReceipt, LunarisError> {
run(engine, scope, now_ms, Mode::Preview).await
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum Mode {
Commit,
Preview,
}
async fn run(
engine: &crate::handle::Lunaris,
scope: &Scope,
now_ms: u64,
mode: Mode,
) -> Result<RetentionReceipt, LunarisError> {
let storage = engine.storage();
let Some(policy) = read_policy(&storage, scope, &engine.clock()).await? else {
return Ok(RetentionReceipt { policy: None, cutoff: None, forget: None });
};
let cutoff = Hlc { wall_ms: now_ms.saturating_sub(policy.max_age_ms), counter: 0, node_id: 0 };
let scoped = engine.scoped(scope.clone());
let forget = if mode == Mode::Preview {
scoped.forget(ForgetTarget::Before(cutoff).dry_run()).await?
} else if policy.hard {
let preview = scoped.forget(ForgetTarget::Before(cutoff).dry_run()).await?;
let token = engine.confirm_hard_forget(preview).await?;
scoped.forget(ForgetTarget::Before(cutoff).hard().with_token(token)).await?
} else {
scoped.forget(ForgetTarget::Before(cutoff)).await?
};
Ok(RetentionReceipt { policy: Some(policy), cutoff: Some(cutoff), forget: Some(forget) })
}