Skip to main content

mkit_server/indexed/
resolve.rs

1//! Repository-isolated external delta resolution from member index rows.
2
3use 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
22/// The public message of a decode-budget overrun.
23pub(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/// A membership miss needs the consuming ticket's lag window; a cap is
30/// permanent even inside that window.
31#[derive(Debug)]
32pub enum ResolveFailure {
33    Missing,
34    Capped,
35    /// A selected member frame failed decoding or canonical id verification.
36    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/// Whether a ticket still lies inside the ยง9.4 repository-membership window.
60#[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/// Exact error for an unresolved base, independent of global blob existence.
66#[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
83/// Split retryable capped batch lookups to one id, preserving only
84/// repository-scoped answers. A final capped id is returned for the caller's
85/// own error mapping.
86pub 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
96/// Speculatively locate syntactic bases in 256-id batches. The caller reports
97/// a cap only if the decoder later requests that id as an external base.
98pub 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
151/// Read one bounded range into memory, preserving the blob contract's
152/// streaming cap and refusing short or overlong backend responses.
153pub(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
219/// Canonical member bytes and their total delta depth.
220pub type ResolvedMember = (Arc<[u8]>, u32);
221
222/// Canonical member bytes retained during one verification. Each location is
223/// charged once, even when multiple deltas reuse it as an external base.
224#[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    /// A linear reconstruction only needs its newest canonical base. Keep
253    /// retention separate from decode admission so eviction cannot exhaust
254    /// the next frame's entry allowance.
255    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    /// The retained locations, with their canonical bytes and total depth.
294    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            // The caller still owns the current base until this decode ends;
308            // older intermediates have no remaining consumer in this chain.
309            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/// Resolve a member object's canonical bytes and total depth. Memoization
326/// shares repeated external bases across the consuming advance.
327#[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/// Restricted canonical acquisition; never used by serving or byte reuse.
347#[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/// Independent limits for every encoded frame and decoded chain member.
367#[derive(Debug, Clone, Copy)]
368pub struct MemberSourceLimits {
369    pub max_frame_bytes: u64,
370    pub max_decoded_bytes: u64,
371}
372
373/// Restricted acquisition with admission-derived limits on every chain source.
374#[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
437/// Select exactly the base location used by canonical reconstruction, without bytes.
438async 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/// Prove the reconstruction chain clear using metadata only, with the same
483/// location-cycle and delta-hop limits as `member_object`.
484#[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        // An ancestry frontier is already charged by the walk; recursive,
579        // uncached delta bases share its work budget. Cache hits are free.
580        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                // A raw terminal base may be one node beyond the hop cap;
599                // another delta may not. Stop before a long chain recurses.
600                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                // A selected member's verified metadata is immutable. A changed
640                // object claim is corruption, even when it now exceeds the budget.
641                // Raw deltas carry their reconstructed size in the delta header;
642                // compressed delta outer claims instead describe stream size.
643                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                    // The outer claim is a stream size. Inspect just the inner
652                    // delta header before a corrupted result can hit the budget.
653                    // An unsuccessful peek proves no mismatch: let normal decode
654                    // retain its existing corruption/resource error separation.
655                    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}