use crate::checkpoint::record::encode_checkpoint_record;
use crate::context::MutationContext;
use crate::error::{CoreError, Result};
use loonfs_api::wire::control::{
decode_control_object, CheckpointOwner, CheckpointRecordLifecycle, CheckpointRecordState,
ControlObjectKind,
};
use loonfs_api::{GeneratedIdValidationError, ManifestObjectId, NamespaceId, RetainedReason};
use loonfs_objectstore::{ObjectStore, ObjectStoreError};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum CheckpointSweep {
Delete,
Released,
Retain,
}
pub(super) fn lease_expired(record: &CheckpointRecordState, now_ms: u64) -> bool {
record
.expires_at_ms
.is_some_and(|expires_at_ms| expires_at_ms <= now_ms)
}
pub(super) async fn sweep_checkpoint_record<S: ObjectStore + ?Sized>(
store: &S,
namespace_id: &NamespaceId,
key: &str,
grace_window_ms: u64,
namespace_deleted: bool,
context: &MutationContext,
) -> Result<CheckpointSweep> {
let Some(body) = store
.get_with_metadata(key)
.await
.map_err(|error| CoreError::store(key, &error))?
else {
return Ok(CheckpointSweep::Retain);
};
let Ok(envelope) = decode_control_object::<CheckpointRecordState>(
&body.bytes,
ControlObjectKind::CheckpointRecord,
) else {
return Ok(CheckpointSweep::Retain);
};
let record = envelope.state;
if record.namespace_id != *namespace_id {
return Ok(CheckpointSweep::Retain);
}
if let CheckpointRecordLifecycle::Released { released_at_ms } = record.state {
let aged = context.now_ms.saturating_sub(released_at_ms) >= grace_window_ms;
return Ok(if aged {
CheckpointSweep::Delete
} else {
CheckpointSweep::Retain
});
}
if matches!(record.owner, CheckpointOwner::Fork { .. }) {
return Ok(CheckpointSweep::Retain);
}
let releasable = lease_expired(&record, context.now_ms)
|| (namespace_deleted
&& context.now_ms.saturating_sub(record.created_at_ms) >= grace_window_ms);
if !releasable {
return Ok(CheckpointSweep::Retain);
}
let Some(etag) = body.metadata.etag.as_deref() else {
return Ok(CheckpointSweep::Retain);
};
let mut released = record;
released.state = CheckpointRecordLifecycle::Released {
released_at_ms: context.now_ms,
};
let bytes = encode_checkpoint_record(&released)?;
match store.compare_and_swap(key, etag, bytes).await {
Ok(_) => Ok(CheckpointSweep::Released),
Err(ObjectStoreError::PreconditionFailed { .. }) => {
tracing::debug!(
namespace_id = %namespace_id,
object_key = key,
"checkpoint release lost its inspected etag; retaining"
);
Ok(CheckpointSweep::Retain)
}
Err(error) => Err(CoreError::store(key, &error)),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AgedSweep {
Deleted,
AlreadyGone,
RetainedInGraceWindow,
RetainedWithoutTimestamp,
}
impl AgedSweep {
pub fn deleted(self) -> bool {
self == Self::Deleted
}
pub fn retained_reason(self) -> Option<RetainedReason> {
match self {
Self::Deleted | Self::AlreadyGone => None,
Self::RetainedInGraceWindow => Some(RetainedReason::GraceWindow),
Self::RetainedWithoutTimestamp => Some(RetainedReason::NoProviderTimestamp),
}
}
}
pub async fn delete_if_aged<S: ObjectStore + ?Sized>(
store: &S,
key: &str,
grace_window_ms: u64,
now_ms: u64,
) -> std::result::Result<AgedSweep, ObjectStoreError> {
let Some(metadata) = store.head(key).await? else {
return Ok(AgedSweep::AlreadyGone);
};
let Some(last_modified_ms) = metadata.last_modified_ms else {
return Ok(AgedSweep::RetainedWithoutTimestamp);
};
if now_ms.saturating_sub(last_modified_ms) < grace_window_ms {
return Ok(AgedSweep::RetainedInGraceWindow);
}
store.delete(key).await?;
Ok(AgedSweep::Deleted)
}
pub(super) fn manifest_object_id_of(
key: &str,
) -> Option<std::result::Result<ManifestObjectId, GeneratedIdValidationError>> {
let name = key.rsplit('/').next()?;
let object_id = name.strip_suffix(".manifest.json")?;
Some(ManifestObjectId::parse(object_id))
}