Skip to main content

mkit_server/pipeline/
object_reader.rs

1use super::{
2    AuthMode, CallerView, HookSet, OpKind, Operation, Pipeline, Principal, RequestMeta, ms,
3    read_policy,
4};
5use crate::http_objects::{
6    Fail, TakedownVerdict, Target, reach,
7    resolve::{self, Budget, Env},
8};
9use crate::indexed::budget::{Budgeted, SliceBudget};
10use crate::store::{MultipartBlobStore, NamespaceStore, view::ViewStore};
11use crate::takedown::{
12    denial::{denied, object_denials},
13    inventory,
14};
15use crate::url_token::UrlTarget;
16use crate::{Code, RepoId, ServerError};
17/// A signed token and expiry; its credential is redacted from Debug.
18pub type IssuedUrl = crate::url_token::MintedToken;
19use mkit_core::{hash::Hash, object::ObjectType, repo_identity::Namespace};
20use std::collections::{BTreeMap, BTreeSet};
21type Prefetched = (BTreeMap<Hash, Vec<u8>>, BTreeMap<Hash, ObjectMetadata>);
22/// Verified lengths describe canonical objects separately from logical files.
23#[derive(Debug, Clone, Copy, PartialEq, Eq)]
24pub struct ObjectMetadata {
25    /// Reconstructed canonical object kind.
26    pub kind: ObjectType,
27    /// Length of the complete canonical serialization.
28    pub canonical_len: u64,
29    /// `Blob` payload or `ChunkedBlob::total_size`; absent for non-file objects.
30    pub logical_len: Option<u64>,
31}
32fn exhausted() -> ServerError {
33    ServerError::resource_exhausted("object reader byte limit exceeded")
34}
35fn resolution_failure(miss: resolve::Miss, limited: bool) -> ServerError {
36    if limited && miss == resolve::Miss::Capped {
37        exhausted()
38    } else {
39        failure(miss)
40    }
41}
42/// Maximum IDs per call; duplicates preserve input order and share proof work.
43pub const OBJECT_READER_BATCH: usize = 16;
44/// Core call cap inside the Worker invocation allowance.
45pub const OBJECT_READER_CALLS: u32 = 8_500;
46/// Repository view, with envelope-verified owner authority.
47#[derive(Debug)]
48pub enum ReaderView<'a> {
49    /// Anonymous public-repository reads through published membership and refs.
50    Public,
51    /// Signed `ListRefs` envelope; authority and epoch are rechecked per batch.
52    Owner(&'a RequestMeta<'a>),
53}
54/// Repository prefetch for `mkit_core::store::MemorySource` and core builders.
55#[derive(Debug)]
56pub struct ObjectReader<'a, B, N, H> {
57    pipe: &'a Pipeline<B, N, H>,
58    repo: RepoId,
59    view: ReaderView<'a>,
60    cfg: &'a crate::http_objects::HttpObjectsConfig,
61    indexed: &'a crate::indexed::IndexedConfig,
62    seams: &'a crate::http_objects::HttpSeams,
63}
64fn failure<E>(_: E) -> ServerError {
65    ServerError::unavailable("object reader unavailable")
66}
67fn http_failure(fail: &Fail) -> ServerError {
68    ServerError::new(fail.code(), "object reader unavailable")
69}
70impl<B: MultipartBlobStore, N: NamespaceStore + Clone + 'static, H: HookSet> Pipeline<B, N, H> {
71    /// Construct a repository reader using verified owner authority.
72    /// # Errors
73    /// Missing indexed/HTTP configuration, or invalid/non-owner credentials.
74    pub async fn object_reader<'a>(
75        &'a self,
76        repo: RepoId,
77        view: ReaderView<'a>,
78    ) -> Result<ObjectReader<'a, B, N, H>, ServerError> {
79        let (Some(cfg), Some(indexed), Some(seams)) =
80            (&self.cfg.http_objects, &self.cfg.indexed, &self.http_seams)
81        else {
82            return Err(ServerError::new(Code::Unimplemented, "indexed HTTP config"));
83        };
84        let reader = ObjectReader {
85            pipe: self,
86            repo,
87            view,
88            cfg,
89            indexed,
90            seams,
91        };
92        if matches!(reader.view, ReaderView::Owner(_)) {
93            let calls = SliceBudget::new(OBJECT_READER_CALLS);
94            reader.authorize(&calls).await?;
95        }
96        Ok(reader)
97    }
98}
99impl<B: MultipartBlobStore, N: NamespaceStore + Clone + 'static, H: HookSet>
100    ObjectReader<'_, B, N, H>
101{
102    async fn authorize(&self, budget: &SliceBudget) -> Result<bool, ServerError> {
103        for _ in 0..2 {
104            budget.charge().map_err(failure)?;
105        }
106        match &self.view {
107            ReaderView::Public => {
108                let op = Operation::new(
109                    self.repo.clone(),
110                    Principal::Anonymous,
111                    None,
112                    OpKind::HttpGet { ref_name: None },
113                );
114                self.pipe
115                    .authorize_http_read(&op, &Target::Object([0; 32]), self.seams, None)
116                    .await
117                    .map_err(|e| http_failure(&e))?;
118                Ok(false)
119            }
120            ReaderView::Owner(meta) => {
121                if !matches!(self.pipe.cfg.auth, AuthMode::AuthV2(_)) {
122                    return Err(ServerError::unauthenticated("auth v2 required"));
123                }
124                let a = self.pipe.authenticate(meta)?;
125                if a.auth.is_none() || a.repo().repo != self.repo {
126                    return Err(ServerError::unauthenticated("envelope mismatch"));
127                }
128                let op = self.pipe.identify(
129                    &a,
130                    OpKind::ListRefs {
131                        prefix: "refs/".into(),
132                    },
133                )?;
134                let auth = self.pipe.authorize_read(&op).await?;
135                let owner = Namespace::parse(self.repo.namespace.as_str()).is_ok_and(
136                    |n| matches!(n, Namespace::Ed25519(key) if a.principal.ed25519() == Some(&key)),
137                );
138                let grant = self.pipe.visibility_gates_reads()
139                    && auth.facts.grant.is_some()
140                    && op
141                        .write_grant
142                        .as_ref()
143                        .zip(self.pipe.cfg.grants.as_ref())
144                        .and_then(|(h, c)| read_policy::check_grant(c, h.expose(), &op))
145                        .is_some_and(|g| g.write);
146                if auth.facts.caller_view != CallerView::Writer || !(owner || grant) {
147                    return Err(ServerError::permission_denied("writer authority required"));
148                }
149                Ok(true)
150            }
151        }
152    }
153    /// Prefetch canonical bytes; inaccessible IDs are uniformly absent.
154    /// # Errors
155    /// Oversized batches, invalid authority, exhausted budgets or store failures.
156    pub async fn read_canonical(&self, ids: &[Hash]) -> Result<Vec<Option<Vec<u8>>>, ServerError> {
157        self.read_canonical_with_limit(ids, self.cfg.http_decode_budget)
158            .await
159    }
160    /// Bound canonical decode work (including ancestors/bases) and output bytes
161    /// for this call. Duplicate outputs count individually. Inaccessible ids
162    /// remain absent. Length checks precede requested-object allocation.
163    /// # Errors
164    /// `ResourceExhausted` for the caller cap; other errors as `read_canonical`.
165    pub async fn read_canonical_with_limit(
166        &self,
167        ids: &[Hash],
168        max_bytes: u64,
169    ) -> Result<Vec<Option<Vec<u8>>>, ServerError> {
170        let (bytes, _) = self.batch_limited(ids, false, Some(max_bytes)).await?;
171        Ok(ids.iter().map(|id| bytes.get(id).cloned()).collect())
172    }
173    /// Verified object metadata without fetching requested canonical bytes.
174    /// # Errors
175    /// Invalid authority, incomplete proof, corrupt facts or storage failure.
176    pub async fn object_metadata(
177        &self,
178        ids: &[Hash],
179    ) -> Result<Vec<Option<ObjectMetadata>>, ServerError> {
180        let (_, metadata) = self.batch(ids, true).await?;
181        Ok(ids.iter().map(|id| metadata.get(id).copied()).collect())
182    }
183    /// Historical mixed sizes: `Blob` payload, other kinds' canonical length.
184    /// # Errors
185    /// As `object_metadata`; incomplete proofs remain unavailable.
186    #[deprecated(note = "use object_metadata for kind, canonical_len and logical_len")]
187    pub async fn object_sizes(&self, ids: &[Hash]) -> Result<Vec<Option<u64>>, ServerError> {
188        Ok(self
189            .object_metadata(ids)
190            .await?
191            .into_iter()
192            .map(|row| {
193                row.map(|m| {
194                    if m.kind == ObjectType::Blob {
195                        m.logical_len.unwrap_or(0)
196                    } else {
197                        m.canonical_len
198                    }
199                })
200            })
201            .collect())
202    }
203    /// Issue at most 16 URL tokens, preserving order and duplicates.
204    /// Requires configured URL-token keys. Inaccessible targets are uniformly
205    /// absent. Tokens bind unresolved targets exactly as `IssueObjectUrl` does,
206    /// after a bounded published-view preflight, including for Owner readers.
207    /// # Errors
208    /// As [`Self::read_canonical`], plus `unimplemented` without URL-token keys.
209    /// Targets whose reachability cannot be proved within decode/walk limits are absent.
210    #[allow(clippy::too_many_lines)] // One shared-budget preflight and credential issuance pass.
211    pub async fn issue_urls(
212        &self,
213        targets: &[UrlTarget],
214        ttl_s: u32,
215    ) -> Result<Vec<Option<IssuedUrl>>, ServerError> {
216        if targets.len() > OBJECT_READER_BATCH {
217            return Err(ServerError::invalid_argument("batch exceeds 16 targets"));
218        }
219        if self.pipe.cfg.url_tokens.is_none() {
220            return Err(ServerError::unimplemented("URL tokens not configured"));
221        }
222        let calls = SliceBudget::new(OBJECT_READER_CALLS);
223        match self.authorize(&calls).await {
224            Err(e) if e.code() == Code::NotFound && matches!(self.view, ReaderView::Public) => {
225                return Ok(vec![None; targets.len()]);
226            }
227            other => {
228                other?;
229            }
230        }
231        let mut op = match &self.view {
232            ReaderView::Owner(meta) => {
233                let a = self.pipe.authenticate(meta)?;
234                self.pipe.identify(
235                    &a,
236                    OpKind::ListRefs {
237                        prefix: "refs/".into(),
238                    },
239                )?
240            }
241            ReaderView::Public => Operation::new(
242                self.repo.clone(),
243                Principal::Anonymous,
244                None,
245                OpKind::ListRefs {
246                    prefix: "refs/".into(),
247                },
248            ),
249        };
250        let repository = if self.repo.namespace == crate::NamespaceKey::deployment_default() {
251            self.repo.name.as_str().to_owned()
252        } else {
253            format!(
254                "{}/{}",
255                self.repo.namespace.as_str(),
256                self.repo.name.as_str()
257            )
258        };
259        let now = self.pipe.clock.now_ms();
260        let mut issued = Vec::with_capacity(targets.len());
261        for target in targets {
262            calls.charge().map_err(failure)?;
263            calls.charge().map_err(failure)?;
264            op.kind = OpKind::IssueObjectUrl {
265                target: target.clone(),
266                ttl_seconds: ttl_s,
267            };
268            issued.push(
269                match self
270                    .pipe
271                    .issue_url(&op, &repository, target, ttl_s, now)
272                    .await
273                {
274                    Ok(token) => Some(token),
275                    Err(e)
276                        if matches!(
277                            e.code(),
278                            Code::NotFound | Code::PermissionDenied | Code::Unauthenticated
279                        ) =>
280                    {
281                        None
282                    }
283                    Err(e) => return Err(e),
284                },
285            );
286        }
287        let meta = Budgeted::new(&self.pipe.meta, &calls);
288        let blobs = Budgeted::new(&self.pipe.blobs, &calls);
289        let view = ViewStore {
290            store: &meta,
291            repo: &self.repo,
292            writer: false,
293            policy: self.pipe.publication_policy.as_deref(),
294        };
295        let env = Env {
296            no_reads: &BTreeSet::new(),
297            blobs: &blobs,
298            meta: &view,
299            shards: self.pipe.shards.as_ref(),
300            repo: &self.repo,
301            indexed: self.indexed,
302            cfg: self.cfg,
303            metrics: self.pipe.metrics.as_ref(),
304        };
305        let mut decode = Budget(self.cfg.http_decode_budget);
306        let mut ids = Vec::with_capacity(targets.len());
307        for (target, token) in targets.iter().zip(&issued) {
308            if token.is_none() {
309                ids.push(None);
310                continue;
311            }
312            let id = match target {
313                UrlTarget::Object(id) => Some(*id),
314                UrlTarget::Path { reference, path } => {
315                    let shard = self.pipe.shards.ref_shard(&self.repo, reference);
316                    let tip =
317                        crate::store::read::read_ref(&view, &shard, &self.repo.name, reference)
318                            .await
319                            .map_err(failure)?;
320                    if let Some(tip) = tip {
321                        let path = if path.is_empty() {
322                            Vec::new()
323                        } else {
324                            path.split('/').map(|p| p.as_bytes().to_vec()).collect()
325                        };
326                        match resolve::resolve_ref(&env, tip, &path, &mut decode).await {
327                            Ok(resolved) => Some(resolved.leaf),
328                            Err(resolve::Miss::NotFound | resolve::Miss::Capped) => None,
329                            Err(miss) => return Err(failure(miss)),
330                        }
331                    } else {
332                        None
333                    }
334                }
335            };
336            ids.push(id);
337        }
338        let leaves = ids.iter().flatten().copied().collect::<Vec<_>>();
339        let (_, sizes) = self
340            .batch_with_budget(
341                &leaves,
342                true,
343                &calls,
344                false,
345                &BTreeSet::new(),
346                &mut decode,
347                true,
348                None,
349            )
350            .await?;
351        Ok(ids
352            .into_iter()
353            .zip(issued)
354            .map(|(id, token)| id.filter(|id| sizes.contains_key(id)).and(token))
355            .collect())
356    }
357    async fn batch(&self, ids: &[Hash], sizes_only: bool) -> Result<Prefetched, ServerError> {
358        self.batch_limited(ids, sizes_only, None).await
359    }
360    async fn batch_limited(
361        &self,
362        ids: &[Hash],
363        sizes_only: bool,
364        max_bytes: Option<u64>,
365    ) -> Result<Prefetched, ServerError> {
366        if ids.len() > OBJECT_READER_BATCH {
367            return Err(ServerError::invalid_argument("batch exceeds 16 ids"));
368        }
369        let calls = SliceBudget::new(OBJECT_READER_CALLS);
370        let writer = match self.authorize(&calls).await {
371            Err(e) if e.code() == Code::NotFound && matches!(self.view, ReaderView::Public) => {
372                return Ok((BTreeMap::new(), BTreeMap::new()));
373            }
374            other => other?,
375        };
376        self.batch_with_budget(
377            ids,
378            sizes_only,
379            &calls,
380            writer,
381            &if sizes_only {
382                ids.iter().copied().collect()
383            } else {
384                BTreeSet::new()
385            },
386            &mut Budget(
387                max_bytes
388                    .unwrap_or(self.cfg.http_decode_budget)
389                    .min(self.cfg.http_decode_budget),
390            ),
391            // Non-writer views never distinguish a stored but unprovable id
392            // from an unknown one.
393            !writer,
394            max_bytes,
395        )
396        .await
397    }
398    #[allow(clippy::too_many_lines, clippy::too_many_arguments)] // One shared-budget authorization/resolution pass, in precedence order.
399    async fn batch_with_budget(
400        &self,
401        ids: &[Hash],
402        sizes_only: bool,
403        calls: &SliceBudget,
404        writer: bool,
405        forbidden: &BTreeSet<Hash>,
406        decode: &mut Budget,
407        capped_as_absent: bool,
408        max_bytes: Option<u64>,
409    ) -> Result<Prefetched, ServerError> {
410        let mut output_left = max_bytes
411            .unwrap_or(self.cfg.http_decode_budget)
412            .min(self.cfg.http_decode_budget);
413        let limited = max_bytes.is_some() && !capped_as_absent;
414        let pipe = self.pipe;
415        let (cfg, indexed, seams) = (self.cfg, self.indexed, self.seams);
416        let meta = Budgeted::new(&pipe.meta, calls);
417        let blobs = Budgeted::new(&pipe.blobs, calls);
418        let view = ViewStore {
419            store: &meta,
420            repo: &self.repo,
421            writer,
422            policy: pipe.publication_policy.as_deref(),
423        };
424        let env = Env {
425            no_reads: forbidden,
426            blobs: &blobs,
427            meta: &view,
428            shards: pipe.shards.as_ref(),
429            repo: &self.repo,
430            indexed,
431            cfg,
432            metrics: pipe.metrics.as_ref(),
433        };
434        let mut located = if capped_as_absent {
435            Vec::new()
436        } else {
437            resolve::locate_many(&env, ids).await.map_err(failure)?
438        };
439        // Directly denied members are absent even when their reachability
440        // cannot be proved within the caller's byte or walk budget.
441        let mut clear = Vec::with_capacity(located.len());
442        for (id, location) in located {
443            if !denied(&meta, &id).await.map_err(failure)?
444                && !denied(&meta, &location.pack).await.map_err(failure)?
445            {
446                clear.push((id, location));
447            }
448        }
449        located = clear;
450        // Issuance proves missing IDs too: proof cost must not expose membership.
451        let mut targets = if capped_as_absent {
452            ids.iter().copied().collect::<BTreeSet<_>>()
453        } else {
454            located.iter().map(|(id, _)| *id).collect::<BTreeSet<_>>()
455        };
456        let mut reached = BTreeSet::new();
457        let mut fresh = BTreeSet::new();
458        if !writer && !pipe.cfg.takedown_denial {
459            for id in &targets {
460                if let Err(e) = calls.charge() {
461                    if capped_as_absent {
462                        break;
463                    }
464                    return Err(failure(e));
465                }
466                if seams
467                    .reachability
468                    .known_reachable(&self.repo, id, ms(pipe.clock.now_ms()))
469                    .await?
470                {
471                    reached.insert(*id);
472                }
473            }
474        }
475        targets.retain(|id| !reached.contains(id));
476        if !targets.is_empty() {
477            let tips = pipe
478                .reader_tips(&meta, &self.repo, cfg.max_walk_objects, writer)
479                .await;
480            // A spent call budget is an unprovable proof, not a store fault.
481            let (tips, truncated) = match tips {
482                Err(_) if capped_as_absent && calls.remaining() == 0 => (Vec::new(), true),
483                other => other.map_err(|e| http_failure(&e))?,
484            };
485            if truncated && sizes_only && !capped_as_absent {
486                return Err(failure(resolve::Miss::Capped));
487            }
488            if !truncated {
489                let walked =
490                    reach::walk_many(&env, seams.takedown.as_ref(), &tips, &targets, decode).await;
491                let (found, incomplete) = match walked {
492                    Err(_) if capped_as_absent && calls.remaining() == 0 => {
493                        (BTreeSet::new(), Some(resolve::Miss::Capped))
494                    }
495                    other => other.map_err(|e| resolution_failure(e, limited))?,
496                };
497                if limited && incomplete == Some(resolve::Miss::Capped) {
498                    return Err(exhausted());
499                }
500                if sizes_only && !capped_as_absent && incomplete == Some(resolve::Miss::Capped) {
501                    return Err(failure(resolve::Miss::Capped));
502                }
503                fresh.extend(found.iter().copied());
504                reached.extend(found);
505            }
506        }
507        if capped_as_absent {
508            // Do not locate inaccessible targets: membership-dependent work can
509            // distinguish a stored orphan from a missing ID near the call cap.
510            let mut accessible = Vec::new();
511            for id in &reached {
512                if !denied(&meta, id).await.map_err(failure)? {
513                    accessible.push(*id);
514                }
515            }
516            reached = accessible.iter().copied().collect();
517            located = resolve::locate_many(&env, &accessible)
518                .await
519                .map_err(failure)?;
520        }
521        let blocked = if pipe.cfg.takedown_denial && !reached.is_empty() {
522            object_denials(&view, pipe.shards.as_ref(), &self.repo, &reached, indexed).await?
523        } else {
524            BTreeSet::new()
525        };
526        let (mut bytes, mut sizes) = (BTreeMap::new(), BTreeMap::new());
527        for (id, located) in located {
528            if !reached.contains(&id) || blocked.contains(&id) {
529                continue;
530            }
531            if denied(&meta, &id).await.map_err(failure)?
532                || denied(&meta, &located.pack).await.map_err(failure)?
533            {
534                continue;
535            }
536            calls.charge().map_err(failure)?;
537            if !matches!(
538                seams.takedown.check(&self.repo, &id).await?,
539                TakedownVerdict::Clear
540            ) {
541                continue;
542            }
543            if !writer && fresh.contains(&id) {
544                seams
545                    .reachability
546                    .record(&self.repo, &id, ms(pipe.clock.now_ms()));
547            }
548            if sizes_only {
549                if !crate::indexed::resolve::member_dependencies_clear(
550                    &view,
551                    pipe.shards.as_ref(),
552                    &self.repo,
553                    id,
554                    located,
555                    indexed.max_delta_chain_depth,
556                    pipe.metrics.as_ref(),
557                )
558                .await?
559                {
560                    continue;
561                }
562                let row = inventory::entry(&meta, &located.pack, &id)
563                    .await
564                    .map_err(failure)?
565                    .ok_or_else(|| failure(resolve::Miss::Unavailable))?;
566                if row.kind == ObjectType::Delta as u8 {
567                    continue;
568                }
569                let kind = match row.kind {
570                    1 => ObjectType::Blob,
571                    2 => ObjectType::Tree,
572                    3 => ObjectType::Commit,
573                    4 => ObjectType::Remix,
574                    5 => ObjectType::ChunkedBlob,
575                    7 => ObjectType::Tag,
576                    _ => return Err(failure(resolve::Miss::Unavailable)),
577                };
578                if row.canonical_len != located.value.decoded_size {
579                    return Err(failure(resolve::Miss::Unavailable));
580                }
581                sizes.insert(
582                    id,
583                    ObjectMetadata {
584                        kind,
585                        canonical_len: row.canonical_len,
586                        logical_len: row.logical_len,
587                    },
588                );
589            } else {
590                let output = located
591                    .value
592                    .decoded_size
593                    .checked_mul(ids.iter().filter(|requested| **requested == id).count() as u64)
594                    .filter(|n| *n <= output_left)
595                    .ok_or_else(exhausted)?;
596                match resolve::load(&env, id, located, decode).await {
597                    Ok(canonical) if resolve::type_of(&canonical) != Some(ObjectType::Delta) => {
598                        if canonical.len() as u64 != located.value.decoded_size {
599                            return Err(failure(resolve::Miss::Unavailable));
600                        }
601                        output_left -= output;
602                        bytes.insert(id, canonical.to_vec());
603                    }
604                    Ok(_) | Err(resolve::Miss::NotFound) => {}
605                    Err(miss) => return Err(resolution_failure(miss, limited)),
606                }
607            }
608        }
609        Ok((bytes, sizes))
610    }
611}