Skip to main content

mkit_server/http_objects/
reach.rs

1//! Reachability for id URLs (SPEC-HTTP-OBJECTS ยง4, D2): a member id is
2//! served only when a published ref reaches it. The default is a bounded
3//! breadth-first walk from the published refs with a positive-result cache;
4//! [`Reachability`] is the seam a maintained reachable set (WP-5.3a) will
5//! replace it behind.
6//!
7//! The walk follows commit and remix parents and trees, tree entries,
8//! manifest chunks and tag targets. It never follows remix `sources` or
9//! delta bases (the same edges as `mkit_core::ops::graph::children` in
10//! history mode), and it stops at a tombstoned or blocked object
11//! ([`TakedownGate::stops_descent`]). A cap never aborts the walk: the
12//! object or subtree it hides is skipped and the walk reports
13//! `Capped` only if the target was not found anywhere else.
14
15use std::collections::{BTreeMap, BTreeSet, VecDeque};
16use std::sync::Mutex;
17
18use mkit_core::hash::Hash;
19use mkit_core::object::{EntryMode, Object, ObjectType};
20
21use super::TakedownGate;
22use super::resolve::{self, Budget, Env, Miss};
23use crate::repo::RepoId;
24use crate::{BlobStore, BoxFuture, MaybeSend, MaybeSync, NamespaceStore, ServerError};
25
26/// Ids located per index lookup (`store::index::MAX_LOOKUP_IDS`).
27const BATCH: usize = 256;
28/// Canonical manifest bytes without chunk hashes: prologue 6, `total_size`
29/// 8, `chunk_size` 4, `chunk_count` 4.
30const MANIFEST_FIXED: u64 = 22;
31/// The decode-side chunk-count cap of `mkit-core`'s deserializer
32/// (`MAX_CHUNKS`, `serialize.rs`): a larger manifest cannot exist.
33const MAX_MANIFEST_CHUNKS: u64 = 1_000_000;
34/// Least time between sweeps of a full [`TtlReachability`] table.
35const SWEEP_MS: u64 = 1_000;
36
37/// A source of positive reachability answers, consulted before the walk.
38pub trait Reachability: MaybeSend + MaybeSync {
39    /// Whether `id` is already known to be reachable from a published ref of
40    /// `repo` as of `now_ms`. `false` means unknown, never unreachable. A
41    /// positive answer skips the walk, so it also skips the walk's takedown
42    /// stop predicate: [`TakedownGate::check`] must therefore refuse a leaf
43    /// that only a blocked or tombstoned manifest reaches.
44    fn known_reachable<'a>(
45        &'a self,
46        repo: &'a RepoId,
47        id: &'a Hash,
48        now_ms: u64,
49    ) -> BoxFuture<'a, Result<bool, ServerError>>;
50
51    /// Record a proof: a walk found `id`, or a ref-path serve resolved it.
52    fn record(&self, repo: &RepoId, id: &Hash, now_ms: u64);
53
54    /// Forget every proof of `repo`: a takedown, suspension or visibility
55    /// change (WP-5.9a) that cannot wait out the lag calls it.
56    fn invalidate(&self, repo: &RepoId);
57}
58
59/// The default: proofs live for `lag_ms`, in a table of at most
60/// `max_entries` rows. A rewind or ref deletion is therefore visible within
61/// the configured `reachability_lag`; only [`Reachability::invalidate`]
62/// forgets one sooner.
63#[derive(Debug)]
64pub struct TtlReachability {
65    lag_ms: u64,
66    max_entries: usize,
67    table: Mutex<Table>,
68}
69
70#[derive(Debug, Default)]
71struct Table {
72    rows: BTreeMap<(RepoId, Hash), u64>,
73    /// When expired rows were last swept: a full table of live rows is not
74    /// re-scanned on every record.
75    swept_ms: u64,
76}
77
78impl TtlReachability {
79    /// A cache whose rows expire `lag_ms` after they are recorded.
80    #[must_use]
81    pub fn new(lag_ms: u64, max_entries: usize) -> Self {
82        Self {
83            lag_ms,
84            max_entries,
85            table: Mutex::default(),
86        }
87    }
88
89    fn table(&self) -> std::sync::MutexGuard<'_, Table> {
90        self.table
91            .lock()
92            .unwrap_or_else(std::sync::PoisonError::into_inner)
93    }
94}
95
96impl Reachability for TtlReachability {
97    fn known_reachable<'a>(
98        &'a self,
99        repo: &'a RepoId,
100        id: &'a Hash,
101        now_ms: u64,
102    ) -> BoxFuture<'a, Result<bool, ServerError>> {
103        Box::pin(async move {
104            Ok(self
105                .table()
106                .rows
107                .get(&(repo.clone(), *id))
108                .is_some_and(|expires| *expires > now_ms))
109        })
110    }
111
112    fn record(&self, repo: &RepoId, id: &Hash, now_ms: u64) {
113        let mut table = self.table();
114        if table.rows.len() >= self.max_entries && now_ms >= table.swept_ms.saturating_add(SWEEP_MS)
115        {
116            table.rows.retain(|_, expires| *expires > now_ms);
117            table.swept_ms = now_ms;
118        }
119        // A full table of live rows drops the new proof: the next request
120        // walks again. Memory stays bounded.
121        if table.rows.len() < self.max_entries {
122            table
123                .rows
124                .insert((repo.clone(), *id), now_ms.saturating_add(self.lag_ms));
125        }
126    }
127
128    fn invalidate(&self, repo: &RepoId) {
129        self.table()
130            .rows
131            .retain(|(row_repo, _), _| row_repo != repo);
132    }
133}
134
135/// The walk's answer.
136#[derive(Debug, Clone, Copy, PartialEq, Eq)]
137pub(crate) enum Reach {
138    Reachable,
139    Unreachable,
140    /// A cap hid part of the graph and the target was not found elsewhere:
141    /// the uniform 404 and a metric.
142    Capped,
143}
144
145#[derive(Clone, Copy, PartialEq, Eq)]
146enum Kind {
147    /// A commit, remix, tag or ref tip: type unknown until decoded.
148    Node,
149    Tree,
150    /// A tree entry that is a Blob or a `ChunkedBlob`.
151    File,
152}
153
154/// Whether a `File` of `size` canonical bytes can be a manifest: an exact
155/// necessary condition (`22 + 32n`, `n` at most the decode cap), so a plain
156/// Blob is decoded only about one time in 32, and never when it is too big
157/// to be one.
158fn manifest_sized(size: u64) -> bool {
159    (MANIFEST_FIXED..=MANIFEST_FIXED + 32 * MAX_MANIFEST_CHUNKS).contains(&size)
160        && (size - MANIFEST_FIXED).is_multiple_of(32)
161}
162
163/// The queues one walk drains, bounded by `cap` distinct objects so a wide
164/// tree cannot grow them without limit.
165struct Frontier {
166    cap: usize,
167    seen: BTreeSet<Hash>,
168    /// Commits, remixes and tags: history, walked after the trees.
169    nodes: VecDeque<Hash>,
170    /// Trees and files: tip trees are the likeliest targets.
171    work: VecDeque<(Hash, Kind)>,
172    /// A reference was not queued because of `cap`, or an object was skipped
173    /// because of the decode budget: the walk is incomplete.
174    incomplete: Option<Miss>,
175}
176
177impl Frontier {
178    fn push(&mut self, id: Hash, kind: Kind) {
179        if self.seen.contains(&id) {
180            return;
181        }
182        if self.seen.len() >= self.cap {
183            self.incomplete = Some(Miss::Capped);
184            return;
185        }
186        self.seen.insert(id);
187        match kind {
188            Kind::Node => self.nodes.push_back(id),
189            Kind::Tree | Kind::File => self.work.push_back((id, kind)),
190        }
191    }
192
193    /// The next batch: trees and files first, then history.
194    fn batch(&mut self) -> Vec<(Hash, Kind)> {
195        let mut batch = Vec::new();
196        if self.work.is_empty() {
197            while batch.len() < BATCH {
198                let Some(id) = self.nodes.pop_front() else {
199                    break;
200                };
201                batch.push((id, Kind::Node));
202            }
203        } else {
204            while batch.len() < BATCH {
205                let Some(item) = self.work.pop_front() else {
206                    break;
207                };
208                batch.push(item);
209            }
210        }
211        batch
212    }
213}
214
215/// Queue the references of `object`; whether `target` is among them. The
216/// comparison happens for every reference, queued or not.
217fn expand(object: &Object, frontier: &mut Frontier, mut reached: impl FnMut(Hash)) {
218    let mut visit = |child: Hash, kind: Option<Kind>| {
219        reached(child);
220        if let Some(kind) = kind {
221            frontier.push(child, kind);
222        }
223    };
224    match object {
225        Object::Commit(c) => {
226            visit(c.tree_hash, Some(Kind::Tree));
227            c.parents.iter().for_each(|p| visit(*p, Some(Kind::Node)));
228        }
229        Object::Remix(r) => {
230            visit(r.tree_hash, Some(Kind::Tree));
231            r.parents.iter().for_each(|p| visit(*p, Some(Kind::Node)));
232        }
233        Object::Tree(t) => {
234            for entry in &t.entries {
235                let kind = if entry.mode == EntryMode::Tree {
236                    Kind::Tree
237                } else {
238                    Kind::File
239                };
240                visit(entry.object_hash, Some(kind));
241            }
242        }
243        Object::ChunkedBlob(cb) => cb.chunks.iter().for_each(|c| visit(*c, None)),
244        Object::Tag(t) => visit(
245            t.target,
246            match t.target_type {
247                ObjectType::Tree => Some(Kind::Tree),
248                ObjectType::Blob | ObjectType::ChunkedBlob => Some(Kind::File),
249                ObjectType::Delta => None,
250                _ => Some(Kind::Node),
251            },
252        ),
253        Object::Blob(_) | Object::Delta(_) => {}
254    }
255}
256
257/// Search the published `tips` for `target`, queueing at most
258/// `max_walk_objects` objects and charging every decode to `budget`.
259pub(crate) async fn walk<B: BlobStore, N: NamespaceStore>(
260    env: &Env<'_, B, N>,
261    takedown: &dyn TakedownGate,
262    tips: &[Hash],
263    target: Hash,
264    budget: &mut Budget,
265) -> Result<Reach, Miss> {
266    let targets = BTreeSet::from([target]);
267    let (reached, incomplete) = walk_many(env, takedown, tips, &targets, budget).await?;
268    Ok(if reached.contains(&target) {
269        Reach::Reachable
270    } else if incomplete == Some(Miss::Capped) {
271        Reach::Capped
272    } else {
273        Reach::Unreachable
274    })
275}
276
277/// One walk for a batch. `Env::no_reads` forbids loading selected objects, including
278/// as ancestors; unresolved descendants then report an incomplete proof.
279/// The residual miss distinguishes a budget cap from blocked ancestry.
280pub(crate) async fn walk_many<B: BlobStore, N: NamespaceStore>(
281    env: &Env<'_, B, N>,
282    takedown: &dyn TakedownGate,
283    tips: &[Hash],
284    targets: &BTreeSet<Hash>,
285    budget: &mut Budget,
286) -> Result<(BTreeSet<Hash>, Option<Miss>), Miss> {
287    let mut reached = BTreeSet::new();
288    let mut frontier = Frontier {
289        cap: env.cfg.max_walk_objects,
290        seen: BTreeSet::new(),
291        nodes: VecDeque::new(),
292        work: VecDeque::new(),
293        incomplete: None,
294    };
295    let mut record = |id| {
296        if targets.contains(&id) {
297            reached.insert(id);
298        }
299    };
300    for tip in tips {
301        record(*tip);
302        frontier.push(*tip, Kind::Node);
303    }
304    loop {
305        if reached.len() == targets.len() {
306            return Ok((reached, None));
307        }
308        let batch = frontier.batch();
309        if batch.is_empty() {
310            return Ok((reached, frontier.incomplete));
311        }
312        let ids: Vec<Hash> = batch.iter().map(|(id, _)| *id).collect();
313        let kinds: BTreeMap<Hash, Kind> = batch.into_iter().collect();
314        for (id, located) in resolve::locate_many(env, &ids).await? {
315            if crate::takedown::denial::denied(env.meta, &id)
316                .await
317                .map_err(|_| Miss::Unavailable)?
318                || takedown.stops_descent(env.repo, &id)
319            {
320                continue;
321            }
322            if kinds[&id] == Kind::File && !manifest_sized(located.value.decoded_size) {
323                continue;
324            }
325            if env.no_reads.contains(&id) {
326                frontier.incomplete = Some(Miss::Capped);
327                continue;
328            }
329            let bytes = match resolve::load(env, id, located, budget).await {
330                Ok(bytes) => bytes,
331                Err(Miss::Capped) => {
332                    frontier.incomplete = Some(Miss::Capped);
333                    continue;
334                }
335                Err(Miss::NotFound) => {
336                    frontier.incomplete = Some(Miss::NotFound);
337                    return Ok((reached, frontier.incomplete));
338                }
339                Err(other) => return Err(other),
340            };
341            let object =
342                mkit_core::serialize::deserialize(&bytes).map_err(|_| Miss::Unavailable)?;
343            expand(&object, &mut frontier, |child| {
344                if targets.contains(&child) {
345                    reached.insert(child);
346                }
347            });
348            if reached.len() == targets.len() {
349                return Ok((reached, None));
350            }
351        }
352    }
353}
354
355#[cfg(test)]
356mod tests {
357    use futures_executor::block_on;
358
359    use super::*;
360    use crate::repo::{NamespaceKey, RepoName};
361
362    fn repo(name: &str) -> RepoId {
363        RepoId {
364            namespace: NamespaceKey::deployment_default(),
365            name: RepoName::new(name).unwrap(),
366        }
367    }
368
369    fn known(cache: &TtlReachability, repo: &RepoId, id: u8, now: u64) -> bool {
370        block_on(cache.known_reachable(repo, &[id; 32], now)).unwrap()
371    }
372
373    #[test]
374    fn a_proof_expires_after_the_lag_and_is_per_repository() {
375        let cache = TtlReachability::new(60_000, 8);
376        let (a, b) = (repo("a"), repo("b"));
377        cache.record(&a, &[1; 32], 1_000);
378        assert!(known(&cache, &a, 1, 60_999));
379        assert!(!known(&cache, &a, 1, 61_000));
380        assert!(!known(&cache, &b, 1, 1_000), "another repository's proof");
381        assert!(!known(&cache, &a, 2, 1_000));
382        cache.record(&a, &[2; 32], 1_000);
383        cache.record(&b, &[2; 32], 1_000);
384        cache.invalidate(&a);
385        assert!(!known(&cache, &a, 2, 1_001) && known(&cache, &b, 2, 1_001));
386    }
387
388    #[test]
389    fn a_full_table_is_bounded_and_sweeps_at_most_once_a_second() {
390        let cache = TtlReachability::new(10_000, 2);
391        let a = repo("a");
392        cache.record(&a, &[1; 32], 5_000);
393        cache.record(&a, &[2; 32], 5_000);
394        // Full of live rows: the new proof is dropped, never stored.
395        cache.record(&a, &[3; 32], 6_000);
396        assert!(!known(&cache, &a, 3, 6_000));
397        // Once the rows expire, the next record sweeps and stores.
398        cache.record(&a, &[3; 32], 16_000);
399        assert!(known(&cache, &a, 3, 16_000));
400        assert!(!known(&cache, &a, 1, 16_000));
401    }
402
403    #[test]
404    fn only_a_manifest_shaped_file_is_ever_decoded() {
405        for (size, want) in [
406            (0, false),
407            (21, false),
408            (22, true),
409            (23, false),
410            (54, true),
411            (22 + 32 * 1_000_000, true),
412            (22 + 32 * 1_000_001, false),
413            (10 + 300 * 1024 * 1024, false),
414        ] {
415            assert_eq!(manifest_sized(size), want, "{size}");
416        }
417    }
418}