use super::reap::lease_expired;
use crate::checkpoint::record::{encode_checkpoint_record, release_checkpoint_record};
use crate::context::MutationContext;
use crate::error::{CoreError, Result};
use crate::namespace::control::{read_head_object, ControlObjectLoadError};
use loonfs_api::wire::control::{
decode_control_object, CheckpointOwner, CheckpointRecordLifecycle, CheckpointRecordState,
ControlObjectKind, NamespaceState,
};
use loonfs_api::NamespaceId;
use loonfs_objectstore::keys::metadata_manifest_object;
use loonfs_objectstore::ObjectStore;
pub(super) enum ForkCheckpointSweep {
Released,
Retained,
NotAnActiveFork,
}
pub(super) async fn release_missing_basis_checkpoint<S: ObjectStore + ?Sized>(
store: &S,
namespace_id: &NamespaceId,
key: &str,
grace_window_ms: u64,
context: &MutationContext,
) -> Result<bool> {
let Some(body) = store
.get_with_metadata(key)
.await
.map_err(|error| CoreError::store(key, &error))?
else {
return Ok(false);
};
let Ok(envelope) = decode_control_object::<CheckpointRecordState>(
&body.bytes,
ControlObjectKind::CheckpointRecord,
) else {
return Ok(false);
};
let record = envelope.state;
if record.state != (CheckpointRecordLifecycle::Active {}) {
return Ok(false);
}
if context.now_ms.saturating_sub(record.created_at_ms) < grace_window_ms {
return Ok(false);
}
let manifest_key = metadata_manifest_object(namespace_id.as_str(), &record.manifest_object_id);
if store
.head(&manifest_key)
.await
.map_err(|error| CoreError::store(&manifest_key, &error))?
.is_some()
{
return Ok(false);
}
release_checkpoint_record(store, namespace_id, &record.checkpoint_id, context.now_ms).await?;
Ok(true)
}
pub(super) async fn maybe_release_fork_checkpoint<S: ObjectStore + ?Sized>(
store: &S,
key: &str,
context: &MutationContext,
) -> Result<ForkCheckpointSweep> {
let Some(body) = store
.get_with_metadata(key)
.await
.map_err(|error| CoreError::store(key, &error))?
else {
return Ok(ForkCheckpointSweep::NotAnActiveFork);
};
let Ok(envelope) = decode_control_object::<CheckpointRecordState>(
&body.bytes,
ControlObjectKind::CheckpointRecord,
) else {
return Ok(ForkCheckpointSweep::Retained);
};
let record = envelope.state;
if record.state != (CheckpointRecordLifecycle::Active {}) {
return Ok(ForkCheckpointSweep::NotAnActiveFork);
}
let CheckpointOwner::Fork {
target_namespace_id,
} = &record.owner
else {
return Ok(ForkCheckpointSweep::NotAnActiveFork);
};
if !fork_target_proven_gone(store, target_namespace_id, &record, context).await? {
return Ok(ForkCheckpointSweep::Retained);
}
let Some(etag) = body.metadata.etag.as_deref() else {
return Ok(ForkCheckpointSweep::Retained);
};
let mut released = record;
released.state = CheckpointRecordLifecycle::Released {
released_at_ms: context.now_ms,
};
let encoded = encode_checkpoint_record(&released)?;
match store.compare_and_swap(key, etag, encoded).await {
Ok(_) => Ok(ForkCheckpointSweep::Released),
Err(loonfs_objectstore::ObjectStoreError::PreconditionFailed { .. }) => {
Ok(ForkCheckpointSweep::Retained)
}
Err(error) => Err(CoreError::store(key, &error)),
}
}
pub(super) async fn fork_target_proven_gone<S: ObjectStore + ?Sized>(
store: &S,
target_namespace_id: &NamespaceId,
record: &CheckpointRecordState,
context: &MutationContext,
) -> Result<bool> {
match read_head_object(store, target_namespace_id).await {
Ok(loaded) => Ok(loaded.envelope.state.state == NamespaceState::Deleted),
Err(ControlObjectLoadError::MissingObject { .. }) => {
Ok(lease_expired(record, context.now_ms))
}
Err(_) => Ok(false),
}
}