1use crate::indexed::{
3 geometry,
4 resolve::{self, MemberCache, MemberSourceLimits},
5};
6use crate::pipeline::ShardMap;
7use crate::{BlobStore, Metrics, NamespaceStore, RepoId, ServerError};
8use mkit_core::hash::Hash;
9use std::{collections::BTreeSet, sync::Arc};
10
11#[derive(Debug, Clone, Copy)]
13pub struct Profile {
14 pub(super) limits: MemberSourceLimits,
15 pub(super) chain_depth: u32,
16 retained: u64,
17 retain_latest: bool,
18 resident: u64,
19 slice_calls: u32,
20}
21impl Profile {
22 pub fn inline(decode_budget: u64, chain_depth: u32) -> Result<Self, ServerError> {
27 if decode_budget < 8 || chain_depth == 0 || chain_depth > u32::from(u16::MAX) {
28 return Err(ServerError::invalid_argument("invalid acquisition limits"));
29 }
30 let resident = decode_budget
31 .checked_mul(8)
32 .and_then(|n| n.checked_add(128 << 20))
33 .ok_or_else(|| ServerError::invalid_argument("acquisition resident bound overflow"))?;
34 Ok(Self {
35 limits: MemberSourceLimits {
36 max_frame_bytes: decode_budget.min(geometry::FRAME_BYTES),
37 max_decoded_bytes: decode_budget.min(geometry::CANONICAL_BYTES),
38 },
39 chain_depth,
40 retained: decode_budget,
41 retain_latest: false,
42 resident,
43 slice_calls: ((chain_depth + 1) * 8 + 256).max(700),
44 })
45 }
46 #[must_use]
56 pub const fn scheduled() -> Self {
57 Self {
58 limits: MemberSourceLimits {
59 max_frame_bytes: geometry::FRAME_BYTES,
60 max_decoded_bytes: geometry::CANONICAL_BYTES,
61 },
62 chain_depth: 50,
63 retained: geometry::CANONICAL_BYTES,
64 retain_latest: true,
65 resident: geometry::RESIDENT_BYTES,
66 slice_calls: 700,
67 }
68 }
69 #[must_use]
71 pub const fn resident_upper_bound(&self) -> u64 {
72 self.resident
73 }
74 pub(super) const fn slice_calls(&self) -> u32 {
75 self.slice_calls
76 }
77}
78
79#[derive(Debug)]
81pub struct Verified {
82 pub id: Hash,
83 pub kind: u8,
84 pub canonical: Arc<[u8]>,
85}
86
87#[allow(clippy::too_many_arguments)]
93pub async fn resolve<B: BlobStore, S: NamespaceStore>(
94 blobs: &B,
95 store: &S,
96 shards: &dyn ShardMap,
97 repo: &RepoId,
98 id: Hash,
99 profile: &Profile,
100 metrics: &dyn Metrics,
101) -> Result<Verified, ServerError> {
102 let located = resolve::locate_split(store, shards, repo, &[id], metrics)
103 .await?
104 .remove(&id)
105 .and_then(Result::ok)
106 .flatten()
107 .ok_or_else(|| ServerError::unavailable("canonical member source unavailable"))?;
108 let mut memo = MemberCache::with_work_budget(profile.chain_depth + 1);
109 decode(
110 blobs, store, shards, repo, id, located, profile, metrics, &mut memo,
111 )
112 .await
113}
114
115#[allow(clippy::too_many_arguments)]
116pub(crate) async fn resolve_selected<B: BlobStore, S: NamespaceStore>(
117 blobs: &B,
118 store: &S,
119 shards: &dyn ShardMap,
120 repo: &RepoId,
121 id: Hash,
122 profile: &Profile,
123 root: &crate::Partition,
124 prefix: &crate::Key,
125) -> Result<Verified, ServerError> {
126 let raw = store
127 .get(
128 root,
129 &crate::Key::new([prefix.as_bytes(), &0u32.to_be_bytes()].concat()),
130 )
131 .await
132 .map_err(|_| ServerError::unavailable("selected source unavailable"))?
133 .ok_or_else(|| ServerError::unavailable("selected source unavailable"))?;
134 let (found, located) = super::source::decode_frame(&raw)
135 .map_err(|_| ServerError::unavailable("selected source unavailable"))?;
136 if found != id {
137 return Err(ServerError::unavailable("selected source unavailable"));
138 }
139 let mut memo =
140 MemberCache::with_selection(profile.chain_depth + 1, root.clone(), prefix.clone());
141 decode(
142 blobs,
143 store,
144 shards,
145 repo,
146 id,
147 located,
148 profile,
149 &crate::NoopMetrics,
150 &mut memo,
151 )
152 .await
153}
154
155#[allow(clippy::too_many_arguments)]
156async fn decode<B: BlobStore, S: NamespaceStore>(
157 blobs: &B,
158 store: &S,
159 shards: &dyn ShardMap,
160 repo: &RepoId,
161 id: Hash,
162 located: crate::store::index::LocatedObject,
163 profile: &Profile,
164 metrics: &dyn Metrics,
165 memo: &mut MemberCache,
166) -> Result<Verified, ServerError> {
167 if profile.retain_latest {
168 memo.retain_latest();
169 }
170 let (canonical, _) = resolve::member_object_for_preservation_bounded(
171 blobs,
172 store,
173 shards,
174 repo,
175 id,
176 located,
177 profile.chain_depth,
178 profile.retained,
179 memo,
180 &mut BTreeSet::new(),
181 metrics,
182 profile.limits,
183 )
184 .await
185 .map_err(|failure| match failure {
186 resolve::ResolveFailure::Corrupt(_) => {
187 ServerError::new(crate::Code::DataLoss, "canonical member source corrupt")
188 }
189 _ => ServerError::unavailable("canonical member source unavailable"),
190 })?;
191 let kind = canonical.first().copied().unwrap_or(0);
192 if !matches!(kind, 1 | 2 | 3 | 4 | 5 | 7) {
193 return Err(ServerError::invalid_argument(
194 "preservation target must be a canonical storable object",
195 ));
196 }
197 Ok(Verified {
198 id,
199 kind,
200 canonical,
201 })
202}
203
204#[cfg(test)]
205#[path = "acquisition_tests.rs"]
206mod tests;