1use crate::pipeline::ShardMap;
4use crate::repo::RepoId;
5use crate::store::{
6 codec,
7 index::{self, IndexValue, LocatedObject, LookupError, ObjectLookup},
8 keys,
9};
10use crate::telemetry::{METRIC_INDEX_LOOKUP_CAPPED, Metrics};
11use crate::{BlobBody, BlobKey, BlobStore, BoxFuture, ByteRange, NamespaceStore, ServerError};
12use futures::StreamExt as _;
13use mkit_core::hash::Hash;
14use mkit_core::pack::{DeltaBaseSource, PackError, decode_frame_with, peek_delta_header};
15use std::collections::{BTreeMap, BTreeSet};
16use std::sync::Arc;
17
18fn unavailable() -> ServerError {
19 ServerError::unavailable("object storage request failed")
20}
21
22pub(crate) const DECODE_BUDGET_MESSAGE: &str = "pack exceeds indexed decode budget";
24
25fn budget_exceeded() -> ServerError {
26 ServerError::invalid_argument(DECODE_BUDGET_MESSAGE)
27}
28
29#[derive(Debug)]
32pub enum ResolveFailure {
33 Missing,
34 Capped,
35 Corrupt(ServerError),
37 Other(ServerError),
38}
39
40impl From<ServerError> for ResolveFailure {
41 fn from(error: ServerError) -> Self {
42 Self::Other(error)
43 }
44}
45
46impl ResolveFailure {
47 #[must_use]
48 pub fn public_error(self, now: u64, created: u64, bound: u64) -> ServerError {
49 match self {
50 Self::Missing => missing_base(now, created, bound),
51 Self::Capped => {
52 ServerError::failed_precondition("delta base not available in this repository")
53 }
54 Self::Other(error) | Self::Corrupt(error) => error,
55 }
56 }
57}
58
59#[must_use]
61pub fn lagged(now_ms: u64, created_at_ms: u64, bound_ms: u64) -> bool {
62 now_ms.saturating_sub(created_at_ms) < bound_ms
63}
64
65#[must_use]
67pub fn missing_base(now_ms: u64, created_at_ms: u64, bound_ms: u64) -> ServerError {
68 if lagged(now_ms, created_at_ms, bound_ms) {
69 ServerError::unavailable("repository membership not yet visible")
70 } else {
71 ServerError::failed_precondition("delta base not available in this repository")
72 }
73}
74
75fn cap_reason(cap: LookupError) -> &'static str {
76 match cap {
77 LookupError::TooManyRows => "rows",
78 LookupError::TooManyPages => "pages",
79 LookupError::TooManyMembershipReads => "membership_reads",
80 }
81}
82
83pub async fn locate_split<S: NamespaceStore>(
87 store: &S,
88 shards: &dyn ShardMap,
89 repo: &RepoId,
90 ids: &[Hash],
91 metrics: &dyn Metrics,
92) -> Result<BTreeMap<Hash, ObjectLookup>, ServerError> {
93 locate_split_inner(store, shards, repo, ids, Some(metrics)).await
94}
95
96pub async fn locate_split_quiet<S: NamespaceStore>(
99 store: &S,
100 shards: &dyn ShardMap,
101 repo: &RepoId,
102 ids: &[Hash],
103) -> Result<BTreeMap<Hash, ObjectLookup>, ServerError> {
104 locate_split_inner(store, shards, repo, ids, None).await
105}
106
107async fn locate_split_inner<S: NamespaceStore>(
108 store: &S,
109 shards: &dyn ShardMap,
110 repo: &RepoId,
111 ids: &[Hash],
112 metrics: Option<&dyn Metrics>,
113) -> Result<BTreeMap<Hash, ObjectLookup>, ServerError> {
114 let mut todo: Vec<Vec<Hash>> = ids
115 .chunks(index::MAX_LOOKUP_IDS)
116 .map(<[Hash]>::to_vec)
117 .collect();
118 let mut found = BTreeMap::new();
119 while let Some(chunk) = todo.pop() {
120 let answers = index::locate_many(store, shards, repo, &chunk)
121 .await
122 .map_err(|_| unavailable())?;
123 let split = chunk.len() > 1
124 && answers.iter().any(|answer| {
125 matches!(
126 answer,
127 Err(LookupError::TooManyPages | LookupError::TooManyMembershipReads)
128 )
129 });
130 if split {
131 let mid = chunk.len() / 2;
132 todo.push(chunk[mid..].to_vec());
133 todo.push(chunk[..mid].to_vec());
134 continue;
135 }
136 for (id, answer) in chunk.into_iter().zip(answers) {
137 if let (Err(cap), Some(metrics)) = (answer, metrics) {
138 tracing::error!(reason = cap_reason(cap), "object index lookup capped");
139 metrics.incr(
140 METRIC_INDEX_LOOKUP_CAPPED,
141 &[("reason", cap_reason(cap))],
142 1,
143 );
144 }
145 found.insert(id, answer);
146 }
147 }
148 Ok(found)
149}
150
151pub(super) async fn frame_bytes<B: BlobStore>(
154 blobs: &B,
155 pack: Hash,
156 offset: u64,
157 length: u64,
158 budget: u64,
159) -> Result<Vec<u8>, ServerError> {
160 if length == 0 {
161 return Err(ServerError::invalid_argument("object hash mismatch"));
162 }
163 if length > budget {
164 return Err(budget_exceeded());
165 }
166 let end = offset.checked_add(length - 1).ok_or_else(unavailable)?;
167 let body = blobs
168 .get(
169 &BlobKey::pack(pack),
170 Some(ByteRange {
171 start: offset,
172 end_inclusive: end,
173 }),
174 )
175 .await
176 .map_err(|_| unavailable())?
177 .ok_or_else(unavailable)?;
178 let mut bytes = Vec::new();
179 bytes
180 .try_reserve_exact(usize::try_from(length).map_err(|_| budget_exceeded())?)
181 .map_err(|_| budget_exceeded())?;
182 match body {
183 BlobBody::Bytes(value) => {
184 if value.len() as u64 != length {
185 return Err(unavailable());
186 }
187 bytes.extend_from_slice(&value);
188 }
189 BlobBody::Stream { mut stream, .. } => {
190 while let Some(chunk) = stream.next().await {
191 let chunk = chunk.map_err(|_| unavailable())?;
192 if (bytes.len() as u64).saturating_add(chunk.len() as u64) > length {
193 return Err(unavailable());
194 }
195 bytes.extend_from_slice(&chunk);
196 }
197 }
198 }
199 if bytes.len() as u64 != length {
200 return Err(unavailable());
201 }
202 Ok(bytes)
203}
204
205struct CachedBase(Option<(Hash, Arc<[u8]>)>);
206impl DeltaBaseSource for CachedBase {
207 const VERIFIED: bool = false;
208 fn base(&mut self, id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
209 Ok(self
210 .0
211 .as_ref()
212 .filter(|(base, _)| base == id)
213 .map(|(_, bytes)| bytes.to_vec()))
214 }
215}
216
217type Location = (Hash, Hash, u64);
218
219pub type ResolvedMember = (Arc<[u8]>, u32);
221
222#[derive(Debug, Clone, Default, PartialEq, Eq)]
225pub struct MemberCache {
226 no_reads: BTreeSet<Hash>,
227 rows: BTreeMap<Location, ResolvedMember>,
228 retained_bytes: u64,
229 retain_latest: bool,
230 remaining_work: Option<u32>,
231 selection: Option<(crate::Partition, crate::Key)>,
232}
233
234impl MemberCache {
235 #[cfg(feature = "http-objects")]
236 pub(crate) fn forbid_reads(&mut self, ids: &BTreeSet<Hash>) {
237 self.no_reads.clone_from(ids);
238 }
239 pub(crate) fn with_selection(limit: u32, root: crate::Partition, prefix: crate::Key) -> Self {
240 Self {
241 selection: Some((root, prefix)),
242 ..Self::with_work_budget(limit)
243 }
244 }
245 pub(crate) fn with_work_budget(limit: u32) -> Self {
246 Self {
247 remaining_work: Some(limit),
248 ..Self::default()
249 }
250 }
251
252 pub(crate) fn retain_latest(&mut self) {
256 self.retain_latest = true;
257 }
258
259 fn available(&self, budget: u64) -> Result<u64, ServerError> {
260 if self.retain_latest {
261 Ok(budget)
262 } else {
263 budget
264 .checked_sub(self.retained_bytes)
265 .ok_or_else(budget_exceeded)
266 }
267 }
268
269 pub(crate) fn charge_work(&mut self, amount: u32) -> Result<(), ResolveFailure> {
270 if let Some(remaining) = &mut self.remaining_work {
271 *remaining = remaining
272 .checked_sub(amount)
273 .ok_or(ResolveFailure::Capped)?;
274 }
275 Ok(())
276 }
277
278 #[must_use]
279 pub fn len(&self) -> usize {
280 self.rows.len()
281 }
282
283 #[must_use]
284 pub fn is_empty(&self) -> bool {
285 self.rows.is_empty()
286 }
287
288 #[must_use]
289 pub fn retained_bytes(&self) -> u64 {
290 self.retained_bytes
291 }
292
293 pub(crate) fn rows(&self) -> impl Iterator<Item = (&Location, (&Arc<[u8]>, &u32))> {
295 self.rows
296 .iter()
297 .map(|(location, (bytes, depth))| (location, (bytes, depth)))
298 }
299
300 fn insert(
301 &mut self,
302 location: Location,
303 value: ResolvedMember,
304 budget: u64,
305 ) -> Result<(), ResolveFailure> {
306 if self.retain_latest {
307 self.rows.clear();
310 self.retained_bytes = 0;
311 }
312 let used = self
313 .retained_bytes
314 .checked_add(value.0.len() as u64)
315 .ok_or_else(budget_exceeded)?;
316 if used > budget {
317 return Err(budget_exceeded().into());
318 }
319 self.rows.insert(location, value);
320 self.retained_bytes = used;
321 Ok(())
322 }
323}
324
325#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
328pub fn member_object<'a, B: BlobStore, S: NamespaceStore>(
329 blobs: &'a B,
330 store: &'a S,
331 shards: &'a dyn ShardMap,
332 repo: &'a RepoId,
333 id: Hash,
334 located: LocatedObject,
335 cap: u32,
336 budget: u64,
337 memo: &'a mut MemberCache,
338 visiting: &'a mut BTreeSet<Location>,
339 metrics: &'a dyn Metrics,
340) -> BoxFuture<'a, Result<ResolvedMember, ResolveFailure>> {
341 member_object_inner(
342 blobs, store, shards, repo, id, located, cap, budget, memo, visiting, metrics, true, None,
343 )
344}
345
346#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
348pub fn member_object_for_preservation<'a, B: BlobStore, S: NamespaceStore>(
349 blobs: &'a B,
350 store: &'a S,
351 shards: &'a dyn ShardMap,
352 repo: &'a RepoId,
353 id: Hash,
354 located: LocatedObject,
355 cap: u32,
356 budget: u64,
357 memo: &'a mut MemberCache,
358 visiting: &'a mut BTreeSet<Location>,
359 metrics: &'a dyn Metrics,
360) -> BoxFuture<'a, Result<ResolvedMember, ResolveFailure>> {
361 member_object_inner(
362 blobs, store, shards, repo, id, located, cap, budget, memo, visiting, metrics, false, None,
363 )
364}
365
366#[derive(Debug, Clone, Copy)]
368pub struct MemberSourceLimits {
369 pub max_frame_bytes: u64,
370 pub max_decoded_bytes: u64,
371}
372
373#[allow(clippy::too_many_arguments)]
375pub fn member_object_for_preservation_bounded<'a, B: BlobStore, S: NamespaceStore>(
376 blobs: &'a B,
377 store: &'a S,
378 shards: &'a dyn ShardMap,
379 repo: &'a RepoId,
380 id: Hash,
381 located: LocatedObject,
382 cap: u32,
383 budget: u64,
384 memo: &'a mut MemberCache,
385 visiting: &'a mut BTreeSet<Location>,
386 metrics: &'a dyn Metrics,
387 limits: MemberSourceLimits,
388) -> BoxFuture<'a, Result<ResolvedMember, ResolveFailure>> {
389 member_object_inner(
390 blobs,
391 store,
392 shards,
393 repo,
394 id,
395 located,
396 cap,
397 budget,
398 memo,
399 visiting,
400 metrics,
401 false,
402 Some(limits),
403 )
404}
405
406async fn selected_frame<S: NamespaceStore>(
407 store: &S,
408 selection: Option<&(crate::Partition, crate::Key)>,
409 level: usize,
410 id: Hash,
411) -> Result<Option<LocatedObject>, ServerError> {
412 let Some((root, prefix)) = selection else {
413 return Ok(None);
414 };
415 let key = crate::Key::new(
416 [
417 prefix.as_bytes(),
418 &u32::try_from(level)
419 .map_err(|_| unavailable())?
420 .to_be_bytes(),
421 ]
422 .concat(),
423 );
424 let raw = store
425 .get(root, &key)
426 .await
427 .map_err(|_| unavailable())?
428 .ok_or_else(unavailable)?;
429 let (found, selected) =
430 crate::takedown::source::decode_frame(&raw).map_err(|_| unavailable())?;
431 if found != id {
432 return Err(unavailable());
433 }
434 Ok(Some(selected))
435}
436
437async fn member_base<S: NamespaceStore>(
439 store: &S,
440 shards: &dyn ShardMap,
441 repo: &RepoId,
442 base: Hash,
443 located: LocatedObject,
444 metrics: &dyn Metrics,
445 selected: Option<LocatedObject>,
446) -> Result<LocatedObject, ResolveFailure> {
447 if let Some(selected) = selected {
448 return Ok(selected);
449 }
450 let partition = shards.object_index(repo, &base);
451 let key = keys::object_index(&repo.name, &base, &located.pack);
452 let same = store
453 .get_many(&partition, &[key])
454 .await
455 .map_err(|_| unavailable())?;
456 if same.len() != 1 {
457 return Err(unavailable().into());
458 }
459 let same = same.into_iter().next().flatten();
460 let next = if let Some(value) = same {
461 let value = codec::decode_object_index(&base, &value).map_err(|_| unavailable())?;
462 (value.frame_offset < located.value.frame_offset).then_some(LocatedObject {
463 pack: located.pack,
464 value,
465 })
466 } else {
467 None
468 };
469 Ok(match next {
470 Some(next) => next,
471 None => match locate_split(store, shards, repo, &[base], metrics)
472 .await?
473 .remove(&base)
474 {
475 Some(Ok(Some(next))) => next,
476 Some(Err(_)) => return Err(ResolveFailure::Capped),
477 _ => return Err(ResolveFailure::Missing),
478 },
479 })
480}
481
482#[cfg(feature = "http-objects")]
485pub(crate) async fn member_dependencies_clear<S: NamespaceStore>(
486 store: &S,
487 shards: &dyn ShardMap,
488 repo: &RepoId,
489 mut id: Hash,
490 mut located: LocatedObject,
491 cap: u32,
492 metrics: &dyn Metrics,
493) -> Result<bool, ServerError> {
494 let mut visiting = BTreeSet::new();
495 loop {
496 if crate::takedown::denial::denied(store, &id).await?
497 || crate::takedown::denial::denied(store, &located.pack).await?
498 || !visiting.insert((id, located.pack, located.value.frame_offset))
499 {
500 return Ok(false);
501 }
502 let Some(base) = located.value.delta_base else {
503 return Ok(true);
504 };
505 if visiting.len() > usize::try_from(cap).unwrap_or(usize::MAX) {
506 return Ok(false);
507 }
508 located = match member_base(store, shards, repo, base, located, metrics, None).await {
509 Ok(next) => next,
510 Err(ResolveFailure::Missing | ResolveFailure::Capped) => return Ok(false),
511 Err(ResolveFailure::Other(error) | ResolveFailure::Corrupt(error)) => {
512 return Err(error);
513 }
514 };
515 id = base;
516 }
517}
518
519#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
520fn member_object_inner<'a, B: BlobStore, S: NamespaceStore>(
521 blobs: &'a B,
522 store: &'a S,
523 shards: &'a dyn ShardMap,
524 repo: &'a RepoId,
525 id: Hash,
526 located: LocatedObject,
527 cap: u32,
528 budget: u64,
529 memo: &'a mut MemberCache,
530 visiting: &'a mut BTreeSet<Location>,
531 metrics: &'a dyn Metrics,
532 enforce_denial: bool,
533 source_limits: Option<MemberSourceLimits>,
534) -> BoxFuture<'a, Result<ResolvedMember, ResolveFailure>> {
535 Box::pin(async move {
536 if memo.no_reads.contains(&id) {
537 return Err(budget_exceeded().into());
538 }
539 if enforce_denial {
540 crate::takedown::denial::require_clear(store, &id).await?;
541 crate::takedown::denial::require_clear(store, &located.pack).await?;
542 }
543 let source_limits = Some(source_limits.unwrap_or(MemberSourceLimits {
544 max_frame_bytes: super::geometry::FRAME_BYTES,
545 max_decoded_bytes: super::geometry::CANONICAL_BYTES,
546 }));
547 if source_limits.is_some_and(|limits| {
548 located.value.frame_length > limits.max_frame_bytes
549 || located.value.decoded_size > limits.max_decoded_bytes
550 }) {
551 return Err(budget_exceeded().into());
552 }
553 let location = (id, located.pack, located.value.frame_offset);
554 if let Some(selected) =
555 selected_frame(store, memo.selection.as_ref(), visiting.len(), id).await?
556 {
557 if selected != located {
558 return Err(unavailable().into());
559 }
560 if !store
561 .has(
562 &shards.membership(repo, &BlobKey::pack(located.pack)),
563 &keys::membership(&repo.name, &located.pack),
564 )
565 .await
566 .map_err(|_| unavailable())?
567 {
568 return Err(ResolveFailure::Missing);
569 }
570 }
571 let available = memo.available(budget)?;
572 if let Some(value) = memo.rows.get(&location) {
573 if value.1 > cap {
574 return Err(ServerError::invalid_argument("delta chain too deep").into());
575 }
576 return Ok(value.clone());
577 }
578 if !visiting.is_empty() {
581 memo.charge_work(1)?;
582 }
583 if !visiting.insert(location) {
584 return Err(ServerError::invalid_argument("delta chain too deep").into());
585 }
586 let result: Result<ResolvedMember, ResolveFailure> = async {
587 let IndexValue {
588 frame_offset,
589 frame_length,
590 delta_base,
591 ..
592 } = located.value;
593 let prefix = frame_bytes(blobs, located.pack, 0, 8, available).await?;
594 let version = u32::from_le_bytes(prefix[4..8].try_into().map_err(|_| unavailable())?);
595 let mut depth = 0;
596 let mut base_bytes = None;
597 if let Some(base) = delta_base {
598 if visiting.len() > usize::try_from(cap).unwrap_or(usize::MAX) {
601 return Err(ServerError::invalid_argument("delta chain too deep").into());
602 }
603 let selected =
604 selected_frame(store, memo.selection.as_ref(), visiting.len(), base).await?;
605 let next =
606 member_base(store, shards, repo, base, located, metrics, selected).await?;
607 let (canonical, base_depth) = member_object_inner(
608 blobs,
609 store,
610 shards,
611 repo,
612 base,
613 next,
614 cap,
615 budget,
616 memo,
617 visiting,
618 metrics,
619 enforce_denial,
620 source_limits,
621 )
622 .await?;
623 base_bytes = Some((base, canonical));
624 depth = base_depth.saturating_add(1);
625 if depth > cap {
626 return Err(ServerError::invalid_argument("delta chain too deep").into());
627 }
628 }
629 let available = memo.available(budget)?;
630 let frame = frame_bytes(
631 blobs,
632 located.pack,
633 frame_offset,
634 frame_length,
635 source_limits.map_or(available, |limits| limits.max_frame_bytes),
636 )
637 .await?;
638 if source_limits.is_some() {
639 let claim = match frame.first() {
644 Some(0x00) => Some(frame.len().saturating_sub(5) as u64),
645 Some(0x03) => frame.get(5..9).and_then(|bytes| {
646 bytes.try_into().ok().map(u32::from_le_bytes).map(u64::from)
647 }),
648 Some(0x02) => frame.get(42..46).and_then(|bytes| {
649 bytes.try_into().ok().map(u32::from_le_bytes).map(u64::from)
650 }),
651 Some(0x04) => frame
656 .get(41..)
657 .and_then(|bytes| peek_delta_header(bytes).ok())
658 .map(|(_, result)| u64::from(result)),
659 _ => None,
660 };
661 if frame.first().copied() != Some(located.value.wire_type)
662 || claim.is_some_and(|size| size != located.value.decoded_size)
663 {
664 return Err(ResolveFailure::Corrupt(ServerError::invalid_argument(
665 "verified source frame metadata mismatch",
666 )));
667 }
668 }
669 let mut source = CachedBase(base_bytes);
670 let (actual, bytes) = decode_frame_with(
671 &frame,
672 version,
673 &mut source,
674 super::geometry::entry_limits(
675 source_limits
676 .map_or(available, |limits| available.min(limits.max_decoded_bytes)),
677 ),
678 )
679 .map_err(|error| {
680 if matches!(error, PackError::PackfileTooLarge) {
681 ResolveFailure::Other(budget_exceeded())
682 } else {
683 ResolveFailure::Corrupt(ServerError::invalid_argument("object hash mismatch"))
684 }
685 })?;
686 if actual != id {
687 return Err(ResolveFailure::Corrupt(ServerError::invalid_argument(
688 "object hash mismatch",
689 )));
690 }
691 Ok((Arc::from(bytes), depth))
692 }
693 .await;
694 visiting.remove(&location);
695 let value = result?;
696 memo.insert(location, value.clone(), budget)?;
697 Ok(value)
698 })
699}