use mkit_core::upload_parts::MIN_PART_SIZE;
use super::{Admission, HookSet, Pipeline, PipelineConfig};
use crate::ServerError;
use crate::repo::Addressing;
use crate::store::{BlobStore, INDEX_FANOUT, MultipartBlobStore, NamespaceStore};
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct ServerInfo {
pub protocol: &'static str,
pub spec_version: u32,
pub max_pack_bytes: u64,
pub part_size: u64,
pub max_parts: u32,
pub max_list_refs_page_size: u32,
pub begin_upload_threshold_bytes: u64,
pub atomic_advance: bool,
pub indexed_mode: bool,
pub admission: bool,
pub receipt_public_key: Vec<u8>,
pub receipt_key_id: String,
pub grant_schemes: Vec<String>,
pub namespace_policy: &'static str,
pub index_fanout: u32,
pub max_delta_chain_depth: u32,
pub inspection_max_objects: Option<u32>,
}
impl PipelineConfig {
pub(super) fn validate_server_info_limits(&self) -> Result<(), ServerError> {
let refusal = if !self.part_size.is_power_of_two()
|| !(MIN_PART_SIZE..=32 * 1024 * 1024).contains(&self.part_size)
{
"part_size must be a power of two in 8..=32 MiB"
} else if self.max_parts == 0 {
"max_parts must be at least 1"
} else if self.part_size * u64::from(self.max_parts) < self.upload_limits.max_total_bytes {
"max_pack_bytes is unreachable with part_size × max_parts"
} else if !(1..=10_000).contains(&self.max_list_refs_page_size) {
"max_list_refs_page_size must be in 1..=10_000"
} else {
return Ok(());
};
Err(ServerError::invalid_argument(refusal))
}
}
impl<B: BlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
pub(super) fn effective_threshold(&self) -> u64 {
if !self.hooks.admission().is_default()
|| matches!(self.cfg.addressing, Addressing::Multi(_))
{
0
} else {
self.cfg.begin_upload_threshold_bytes
}
}
}
impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
#[must_use]
pub fn server_info(&self) -> ServerInfo {
let admission = !self.hooks.admission().is_default();
ServerInfo {
protocol: "mkit.transport.v1",
spec_version: 2,
max_pack_bytes: if self.blobs.supports_multipart() {
self.cfg.upload_limits.max_total_bytes
} else {
self.cfg
.upload_limits
.max_total_bytes
.min(self.cfg.part_size)
},
part_size: self.cfg.part_size,
max_parts: self.cfg.max_parts,
max_list_refs_page_size: self.cfg.max_list_refs_page_size,
begin_upload_threshold_bytes: self.effective_threshold(),
atomic_advance: self.capabilities().atomic_advance,
indexed_mode: self.cfg.indexed_mode(),
admission,
receipt_public_key: self
.cfg
.receipt_publication
.as_ref()
.map_or_else(Vec::new, |keys| keys.public_key.to_vec()),
receipt_key_id: self
.cfg
.receipt_publication
.as_ref()
.map_or_else(String::new, |keys| keys.key_id.clone()),
grant_schemes: self.cfg.grants.as_ref().map_or_else(Vec::new, |grants| {
grants.schemes().tokens().map(str::to_owned).collect()
}),
namespace_policy: self.cfg.advertised_namespace_policy(),
index_fanout: u32::from(INDEX_FANOUT),
max_delta_chain_depth: self
.cfg
.indexed
.as_ref()
.map_or(0, |cfg| cfg.max_delta_chain_depth),
inspection_max_objects: self.inspection_max_objects(),
}
}
}