mkit_server/pipeline/
info.rs1use mkit_core::upload_parts::MIN_PART_SIZE;
4
5use super::{Admission, HookSet, Pipeline, PipelineConfig};
6use crate::ServerError;
7use crate::repo::Addressing;
8use crate::store::{BlobStore, INDEX_FANOUT, MultipartBlobStore, NamespaceStore};
9
10#[derive(Debug, Clone, PartialEq, Eq)]
12#[non_exhaustive]
13pub struct ServerInfo {
14 pub protocol: &'static str,
16 pub spec_version: u32,
18 pub max_pack_bytes: u64,
20 pub part_size: u64,
22 pub max_parts: u32,
24 pub max_list_refs_page_size: u32,
26 pub begin_upload_threshold_bytes: u64,
28 pub atomic_advance: bool,
30 pub indexed_mode: bool,
32 pub admission: bool,
34 pub receipt_public_key: Vec<u8>,
36 pub receipt_key_id: String,
38 pub grant_schemes: Vec<String>,
40 pub namespace_policy: &'static str,
42 pub index_fanout: u32,
44 pub max_delta_chain_depth: u32,
46 pub inspection_max_objects: Option<u32>,
48}
49
50impl PipelineConfig {
51 pub(super) fn validate_server_info_limits(&self) -> Result<(), ServerError> {
52 let refusal = if !self.part_size.is_power_of_two()
53 || !(MIN_PART_SIZE..=32 * 1024 * 1024).contains(&self.part_size)
54 {
55 "part_size must be a power of two in 8..=32 MiB"
56 } else if self.max_parts == 0 {
57 "max_parts must be at least 1"
58 } else if self.part_size * u64::from(self.max_parts) < self.upload_limits.max_total_bytes {
59 "max_pack_bytes is unreachable with part_size × max_parts"
61 } else if !(1..=10_000).contains(&self.max_list_refs_page_size) {
62 "max_list_refs_page_size must be in 1..=10_000"
63 } else {
64 return Ok(());
65 };
66 Err(ServerError::invalid_argument(refusal))
67 }
68}
69
70impl<B: BlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
71 pub(super) fn effective_threshold(&self) -> u64 {
73 if !self.hooks.admission().is_default()
74 || matches!(self.cfg.addressing, Addressing::Multi(_))
75 {
76 0
77 } else {
78 self.cfg.begin_upload_threshold_bytes
79 }
80 }
81}
82
83impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
84 #[must_use]
87 pub fn server_info(&self) -> ServerInfo {
88 let admission = !self.hooks.admission().is_default();
89 ServerInfo {
90 protocol: "mkit.transport.v1",
91 spec_version: 2,
92 max_pack_bytes: if self.blobs.supports_multipart() {
93 self.cfg.upload_limits.max_total_bytes
94 } else {
95 self.cfg
96 .upload_limits
97 .max_total_bytes
98 .min(self.cfg.part_size)
99 },
100 part_size: self.cfg.part_size,
101 max_parts: self.cfg.max_parts,
102 max_list_refs_page_size: self.cfg.max_list_refs_page_size,
103 begin_upload_threshold_bytes: self.effective_threshold(),
104 atomic_advance: self.capabilities().atomic_advance,
105 indexed_mode: self.cfg.indexed_mode(),
106 admission,
107 receipt_public_key: self
108 .cfg
109 .receipt_publication
110 .as_ref()
111 .map_or_else(Vec::new, |keys| keys.public_key.to_vec()),
112 receipt_key_id: self
113 .cfg
114 .receipt_publication
115 .as_ref()
116 .map_or_else(String::new, |keys| keys.key_id.clone()),
117 grant_schemes: self.cfg.grants.as_ref().map_or_else(Vec::new, |grants| {
118 grants.schemes().tokens().map(str::to_owned).collect()
119 }),
120 namespace_policy: self.cfg.advertised_namespace_policy(),
121 index_fanout: u32::from(INDEX_FANOUT),
122 max_delta_chain_depth: self
123 .cfg
124 .indexed
125 .as_ref()
126 .map_or(0, |cfg| cfg.max_delta_chain_depth),
127 inspection_max_objects: self.inspection_max_objects(),
128 }
129 }
130}