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};
17pub 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
24pub struct ObjectMetadata {
25 pub kind: ObjectType,
27 pub canonical_len: u64,
29 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}
42pub const OBJECT_READER_BATCH: usize = 16;
44pub const OBJECT_READER_CALLS: u32 = 8_500;
46#[derive(Debug)]
48pub enum ReaderView<'a> {
49 Public,
51 Owner(&'a RequestMeta<'a>),
53}
54#[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 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 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 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 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 #[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 #[allow(clippy::too_many_lines)] 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 !writer,
394 max_bytes,
395 )
396 .await
397 }
398 #[allow(clippy::too_many_lines, clippy::too_many_arguments)] 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 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 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 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 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}