use crate::error::StoreFailureClass;
use crate::namespace::basis::{read_head_and_metadata_basis, MetadataBasis};
use loonfs_api::wire::control::{
decode_control_object, ControlObjectKind, HeadState, HeadStateEnvelope, MetadataRootEnvelope,
MetadataRootState, WalFloorEnvelope,
};
use loonfs_api::wire::envelope::EnvelopeCodecError;
use loonfs_api::NamespaceId;
use loonfs_objectstore::keys::{metadata_root, wal_floor, wal_head};
use loonfs_objectstore::ObjectStoreError;
use loonfs_objectstore::{ObjectMetadata, ObjectStore};
use serde::{Deserialize, Serialize};
use thiserror::Error;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct LoadedHeadObject {
pub(crate) object_key: String,
pub(crate) metadata: ObjectMetadata,
pub(crate) envelope: HeadStateEnvelope,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct LoadedMetadataRootObject {
pub(crate) object_key: String,
pub(crate) metadata: ObjectMetadata,
pub(crate) envelope: MetadataRootEnvelope,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct LoadedWalFloorObject {
pub(crate) object_key: String,
pub(crate) metadata: ObjectMetadata,
pub(crate) envelope: WalFloorEnvelope,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ControlObjectIdentity {
pub etag: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct LoadedMetadataRootControl {
pub object_key: String,
pub identity: ControlObjectIdentity,
pub state: MetadataRootState,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct LoadedWalFloorControl {
pub object_key: String,
pub identity: ControlObjectIdentity,
pub state: loonfs_api::wire::control::WalFloorState,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct LoadedHeadControl {
pub object_key: String,
pub identity: ControlObjectIdentity,
pub state: HeadState,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Error)]
pub enum ControlObjectLoadError {
#[error("missing control object `{object_key}`")]
MissingObject { object_key: String },
#[error(
"metadata root references seq `{root_manifest_head_seq}` beyond the reloaded head seq `{head_seq}`"
)]
RootAheadOfHead {
root_manifest_head_seq: loonfs_api::ChangeSeq,
head_seq: loonfs_api::ChangeSeq,
},
#[error(
"control object namespace mismatch for `{object_key}`: expected `{expected}`, actual `{actual}`"
)]
NamespaceMismatch {
object_key: String,
expected: NamespaceId,
actual: NamespaceId,
},
#[error(
"control object checksum mismatch for `{object_key}`: expected `{expected}`, actual `{actual}`"
)]
ChecksumMismatch {
object_key: String,
expected: String,
actual: String,
},
#[error("control object codec error for `{object_key}`: {message}")]
Codec { object_key: String, message: String },
#[error("control object store error for `{object_key}`: {message}")]
Store {
object_key: String,
message: String,
class: StoreFailureClass,
},
}
pub(crate) async fn read_wal_floor_object<S: ObjectStore + ?Sized>(
store: &S,
expected_namespace_id: &NamespaceId,
) -> Result<LoadedWalFloorObject, ControlObjectLoadError> {
let object_key = wal_floor(expected_namespace_id.as_str());
let (metadata, encoded_bytes) = read_control_object_bytes(store, &object_key).await?;
let envelope: WalFloorEnvelope =
decode_control_object(&encoded_bytes, ControlObjectKind::WalFloor)
.map_err(|err| map_control_codec_error(&object_key, err))?;
validate_expected_namespace(
&object_key,
expected_namespace_id,
&envelope.state.namespace_id,
)?;
Ok(LoadedWalFloorObject {
object_key,
metadata,
envelope,
})
}
pub(crate) async fn read_metadata_root_object<S: ObjectStore + ?Sized>(
store: &S,
expected_namespace_id: &NamespaceId,
) -> Result<LoadedMetadataRootObject, ControlObjectLoadError> {
let object_key = metadata_root(expected_namespace_id.as_str());
let (metadata, encoded_bytes) = read_control_object_bytes(store, &object_key).await?;
let envelope: MetadataRootEnvelope =
decode_control_object(&encoded_bytes, ControlObjectKind::MetadataRoot)
.map_err(|err| map_control_codec_error(&object_key, err))?;
validate_expected_namespace(
&object_key,
expected_namespace_id,
&envelope.state.namespace_id,
)?;
Ok(LoadedMetadataRootObject {
object_key,
metadata,
envelope,
})
}
pub(crate) async fn read_metadata_root_object_if_present<S: ObjectStore + ?Sized>(
store: &S,
expected_namespace_id: &NamespaceId,
) -> Result<Option<LoadedMetadataRootObject>, ControlObjectLoadError> {
match read_metadata_root_object(store, expected_namespace_id).await {
Ok(loaded) => Ok(Some(loaded)),
Err(ControlObjectLoadError::MissingObject { .. }) => Ok(None),
Err(error) => Err(error),
}
}
pub(crate) async fn read_head_and_metadata_root_if_present<S: ObjectStore + ?Sized>(
store: &S,
expected_namespace_id: &NamespaceId,
) -> Result<(LoadedHeadObject, Option<LoadedMetadataRootObject>), ControlObjectLoadError> {
const ROOT_AHEAD_HEAD_RELOADS: usize = 3;
let (head, root) = futures::join!(
read_head_object(store, expected_namespace_id),
read_metadata_root_object_if_present(store, expected_namespace_id)
);
let mut head = head?;
let Some(root) = root? else {
return Ok((head, None));
};
for _reload in 0..=ROOT_AHEAD_HEAD_RELOADS {
if root.envelope.state.manifest_head_seq <= head.envelope.state.seq {
return Ok((head, Some(root)));
}
head = read_head_object(store, expected_namespace_id).await?;
}
Err(ControlObjectLoadError::RootAheadOfHead {
root_manifest_head_seq: root.envelope.state.manifest_head_seq,
head_seq: head.envelope.state.seq,
})
}
pub(crate) async fn read_head_object<S: ObjectStore + ?Sized>(
store: &S,
expected_namespace_id: &NamespaceId,
) -> Result<LoadedHeadObject, ControlObjectLoadError> {
let object_key = wal_head(expected_namespace_id.as_str());
let (metadata, encoded_bytes) = read_control_object_bytes(store, &object_key).await?;
let envelope: HeadStateEnvelope =
decode_control_object(&encoded_bytes, ControlObjectKind::WalHead)
.map_err(|err| map_control_codec_error(&object_key, err))?;
validate_expected_namespace(
&object_key,
expected_namespace_id,
&envelope.state.namespace_id,
)?;
Ok(LoadedHeadObject {
object_key,
metadata,
envelope,
})
}
pub async fn load_namespace_wal_floor_control<S: ObjectStore + ?Sized>(
store: &S,
expected_namespace_id: &NamespaceId,
) -> Result<LoadedWalFloorControl, ControlObjectLoadError> {
let loaded = read_wal_floor_object(store, expected_namespace_id).await?;
let identity = control_identity(&loaded.object_key, &loaded.metadata)?;
Ok(LoadedWalFloorControl {
object_key: loaded.object_key,
identity,
state: loaded.envelope.state,
})
}
pub async fn load_namespace_checkpoint_record_control<S: ObjectStore + ?Sized>(
store: &S,
expected_namespace_id: &NamespaceId,
checkpoint_id: &loonfs_api::CheckpointId,
) -> Result<Option<loonfs_api::wire::control::CheckpointRecordState>, crate::error::CoreError> {
Ok(
crate::checkpoint::read_checkpoint_record(store, expected_namespace_id, checkpoint_id)
.await?
.map(|loaded| loaded.state),
)
}
pub async fn load_namespace_metadata_root_control<S: ObjectStore + ?Sized>(
store: &S,
expected_namespace_id: &NamespaceId,
) -> Result<LoadedMetadataRootControl, ControlObjectLoadError> {
let loaded = read_metadata_root_object(store, expected_namespace_id).await?;
let identity = control_identity(&loaded.object_key, &loaded.metadata)?;
Ok(LoadedMetadataRootControl {
object_key: loaded.object_key,
identity,
state: loaded.envelope.state,
})
}
pub async fn load_namespace_read_anchor<S: ObjectStore + ?Sized>(
store: &S,
expected_namespace_id: &NamespaceId,
) -> Result<(LoadedHeadControl, MetadataBasis), ControlObjectLoadError> {
let loaded = read_head_and_metadata_basis(store, expected_namespace_id).await?;
let head_identity = control_identity(&loaded.head.object_key, &loaded.head.metadata)?;
Ok((
LoadedHeadControl {
object_key: loaded.head.object_key,
identity: head_identity,
state: loaded.head.envelope.state,
},
loaded.basis,
))
}
pub async fn load_namespace_head_control<S: ObjectStore + ?Sized>(
store: &S,
expected_namespace_id: &NamespaceId,
) -> Result<LoadedHeadControl, ControlObjectLoadError> {
let loaded = read_head_object(store, expected_namespace_id).await?;
let identity = control_identity(&loaded.object_key, &loaded.metadata)?;
Ok(LoadedHeadControl {
object_key: loaded.object_key,
identity,
state: loaded.envelope.state,
})
}
fn control_identity(
object_key: &str,
metadata: &ObjectMetadata,
) -> Result<ControlObjectIdentity, ControlObjectLoadError> {
let etag = metadata
.etag
.clone()
.ok_or_else(|| ControlObjectLoadError::Store {
object_key: object_key.to_owned(),
message: "missing control object etag".to_owned(),
class: StoreFailureClass::Other,
})?;
Ok(ControlObjectIdentity { etag })
}
async fn read_control_object_bytes<S: ObjectStore + ?Sized>(
store: &S,
object_key: &str,
) -> Result<(ObjectMetadata, Vec<u8>), ControlObjectLoadError> {
let body = store
.get_with_metadata(object_key)
.await
.map_err(|err| map_store_load_error(object_key, err))?
.ok_or_else(|| ControlObjectLoadError::MissingObject {
object_key: object_key.to_owned(),
})?;
Ok((body.metadata, body.bytes))
}
fn validate_expected_namespace(
object_key: &str,
expected: &NamespaceId,
actual: &NamespaceId,
) -> Result<(), ControlObjectLoadError> {
if actual != expected {
return Err(ControlObjectLoadError::NamespaceMismatch {
object_key: object_key.to_owned(),
expected: expected.clone(),
actual: actual.clone(),
});
}
Ok(())
}
pub(crate) fn map_control_codec_error(
object_key: &str,
err: EnvelopeCodecError,
) -> ControlObjectLoadError {
match err {
EnvelopeCodecError::ChecksumMismatch { expected, actual } => {
ControlObjectLoadError::ChecksumMismatch {
object_key: object_key.to_owned(),
expected,
actual,
}
}
other => ControlObjectLoadError::Codec {
object_key: object_key.to_owned(),
message: other.to_string(),
},
}
}
fn map_store_load_error(object_key: &str, err: ObjectStoreError) -> ControlObjectLoadError {
ControlObjectLoadError::Store {
object_key: object_key.to_owned(),
message: err.message(),
class: StoreFailureClass::of(&err),
}
}