use crate::error::CoreError;
use crate::namespace::control::read_head_object;
use crate::options::DeleteNamespaceOptions;
use bytes::Bytes;
use loonfs_api::wire::control::{
encode_control_object, AcquiredWriter, ControlObjectKind, HeadState, HeadStateEnvelope,
NamespaceState,
};
use loonfs_api::{DeleteNamespaceResponse, NamespaceId};
use loonfs_objectstore::{ObjectStore, ObjectStoreError, PutMode};
const MAX_DELETE_CAS_ATTEMPTS: usize = 8;
pub(crate) async fn delete_namespace<S: ObjectStore + ?Sized>(
store: &S,
namespace_id: &NamespaceId,
options: DeleteNamespaceOptions,
acquired_writer: AcquiredWriter,
) -> Result<DeleteNamespaceResponse, CoreError> {
let mut attempted_swap = false;
for _attempt in 0..MAX_DELETE_CAS_ATTEMPTS {
let loaded = read_head_object(store, namespace_id)
.await
.map_err(|error| CoreError::MetadataProjection(error.into()))?;
let head = loaded.envelope.state.clone();
if head.state == NamespaceState::Deleted {
if attempted_swap {
return Ok(DeleteNamespaceResponse {
namespace_id: namespace_id.clone(),
head_seq: head.seq,
});
}
return Err(CoreError::NamespaceDeleted {
namespace_id: namespace_id.clone(),
});
}
if head.writer_epoch != acquired_writer.writer_epoch {
return Err(CoreError::WriterFenced(crate::error::WriterFence {
fenced_epoch: acquired_writer.writer_epoch,
active_epoch: head.writer_epoch,
active_writer: head.writer.as_ref().map(|writer| writer.writer_id.clone()),
active_acquired_at_ms: head.writer.as_ref().map(|writer| writer.acquired_at_ms),
}));
}
if let Some(expected) = options.expected_head_seq {
if head.seq != expected {
return Err(CoreError::StaleHeadPrecondition {
expected,
actual: head.seq,
});
}
}
let head_etag = loaded.metadata.etag.clone().ok_or_else(|| {
CoreError::NamespaceCorrupt(format!("missing head etag for `{}`", loaded.object_key))
})?;
let deleted_head = HeadState {
state: NamespaceState::Deleted,
..head.clone()
};
let envelope = HeadStateEnvelope::from_state(ControlObjectKind::WalHead, deleted_head)
.map_err(|err| CoreError::Internal(format!("failed to build head envelope: {err}")))?;
let encoded = encode_control_object(&envelope)
.map_err(|err| CoreError::Internal(format!("failed to encode head object: {err}")))?;
let swap = store
.put(
&loaded.object_key,
Bytes::from(encoded),
PutMode::CompareAndSwap {
expected_etag: head_etag,
},
)
.await;
match swap {
Ok(_) => {
return Ok(DeleteNamespaceResponse {
namespace_id: namespace_id.clone(),
head_seq: head.seq,
});
}
Err(ObjectStoreError::PreconditionFailed { .. }) => continue,
Err(ObjectStoreError::Transport { .. }) => {
attempted_swap = true;
continue;
}
Err(other) => return Err(CoreError::store(&loaded.object_key, &other)),
}
}
Err(CoreError::HeadPublish(
crate::commit::CommitHeadPublishError::StaleHead,
))
}