1use 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
26const BATCH: usize = 256;
28const MANIFEST_FIXED: u64 = 22;
31const MAX_MANIFEST_CHUNKS: u64 = 1_000_000;
34const SWEEP_MS: u64 = 1_000;
36
37pub trait Reachability: MaybeSend + MaybeSync {
39 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 fn record(&self, repo: &RepoId, id: &Hash, now_ms: u64);
53
54 fn invalidate(&self, repo: &RepoId);
57}
58
59#[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 swept_ms: u64,
76}
77
78impl TtlReachability {
79 #[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 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
137pub(crate) enum Reach {
138 Reachable,
139 Unreachable,
140 Capped,
143}
144
145#[derive(Clone, Copy, PartialEq, Eq)]
146enum Kind {
147 Node,
149 Tree,
150 File,
152}
153
154fn 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
163struct Frontier {
166 cap: usize,
167 seen: BTreeSet<Hash>,
168 nodes: VecDeque<Hash>,
170 work: VecDeque<(Hash, Kind)>,
172 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 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
215fn 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
257pub(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
277pub(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 cache.record(&a, &[3; 32], 6_000);
396 assert!(!known(&cache, &a, 3, 6_000));
397 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}