1use mkit_core::hash::{Hash, to_hex};
10use mkit_core::object::ObjectType;
11
12use tracing::Instrument as _;
13
14use super::{
15 HookSet, MultipartBlobStore, NamespaceStore, OpKind, Operation, Pipeline, Principal, Procedure,
16 ms,
17};
18use crate::http_objects::range::{self, Selection};
19use crate::http_objects::reach::{self, Reach};
20use crate::http_objects::resolve::{self, Budget, Env};
21use crate::http_objects::seams::{AdmitDecision, AdmitRequest, TakedownVerdict};
22use crate::http_objects::{
23 Fail, HttpBody, HttpObjectRequest, HttpObjectResponse, HttpSeams, METRIC_HTTP_REACH_CAPPED,
24 ParsedUrl, Target, route,
25};
26use crate::repo::{NamespaceKey, RepoId, RepoName};
27use crate::store::read;
28use crate::{Code, ServerError};
29
30const ALLOW: &str = "GET, HEAD, OPTIONS";
31const PACKMAP_PREFIX: &str = "refs/mkit/packmap/";
33
34impl<B: MultipartBlobStore, N: NamespaceStore + Clone + 'static, H: HookSet> Pipeline<B, N, H> {
35 #[must_use]
37 pub fn http_objects_enabled(&self) -> bool {
38 self.cfg.indexed.is_some() && self.cfg.http_objects.is_some()
39 }
40
41 #[must_use]
43 pub fn url_token_config(&self) -> Option<&crate::url_token::UrlTokenConfig> {
44 self.cfg.url_tokens.as_ref()
45 }
46
47 #[must_use]
52 pub fn with_http_seams(mut self, edit: impl FnOnce(HttpSeams) -> HttpSeams) -> Self {
53 self.http_seams = self.http_seams.take().map(edit);
54 self
55 }
56
57 pub async fn serve_http_object(&self, req: &HttpObjectRequest<'_>) -> HttpObjectResponse {
62 self.serve_http_with_runtime(req, None, None).await
63 }
64
65 pub async fn serve_http_object_with_runtime(
67 &self,
68 req: &HttpObjectRequest<'_>,
69 runtime: crate::http_objects::HttpReadRuntime,
70 ) -> HttpObjectResponse {
71 self.serve_http_with_runtime(req, Some(runtime), None).await
72 }
73
74 pub async fn serve_http_object_with_proofs(
76 &self,
77 req: &HttpObjectRequest<'_>,
78 runtime: crate::http_objects::HttpReadRuntime,
79 proofs: std::sync::Arc<dyn crate::http_objects::ProofServer>,
80 ) -> HttpObjectResponse {
81 self.serve_http_with_runtime(req, Some(runtime), Some(proofs))
82 .await
83 }
84
85 async fn serve_http_with_runtime(
86 &self,
87 req: &HttpObjectRequest<'_>,
88 runtime: Option<crate::http_objects::HttpReadRuntime>,
89 proofs: Option<std::sync::Arc<dyn crate::http_objects::ProofServer>>,
90 ) -> HttpObjectResponse {
91 let mut response = self.serve_http_inner(req, runtime, proofs).await;
92 if req.method == "HEAD" {
93 response.body = HttpBody::Empty;
95 }
96 response
97 }
98
99 async fn serve_http_inner(
100 &self,
101 req: &HttpObjectRequest<'_>,
102 runtime: Option<crate::http_objects::HttpReadRuntime>,
103 proofs: Option<std::sync::Arc<dyn crate::http_objects::ProofServer>>,
104 ) -> HttpObjectResponse {
105 let (Some(_), Some(seams)) = (&self.cfg.http_objects, &self.http_seams) else {
106 return HttpObjectResponse::not_found();
107 };
108 let mut seams = seams.clone();
109 if let Some(proofs) = proofs {
110 seams.proofs = proofs;
111 }
112 if let Some(runtime) = runtime {
113 seams.read_runtime = Some(runtime);
114 }
115 match req.method {
116 "OPTIONS" => return HttpObjectResponse::new(204).with_header("Allow", ALLOW),
117 "GET" | "HEAD" => {}
118 _ => return HttpObjectResponse::error(405).with_header("Allow", ALLOW),
119 }
120 let Ok(parsed) = route::parse_request(req) else {
121 return HttpObjectResponse::error(400);
122 };
123 let token = parsed
126 .query
127 .token
128 .as_ref()
129 .map(|token| seams.tokens.precheck(token, self.clock.now_ms()));
130 let Some(repo) = repo_id(&parsed) else {
131 return HttpObjectResponse::not_found();
132 };
133 let ref_path = matches!(parsed.target, Target::Ref { .. });
134 let procedure = if ref_path {
135 Procedure::HttpGetRefPath
136 } else {
137 Procedure::HttpGetObject
138 };
139 let identity = parsed
140 .repository
141 .as_ref()
142 .map(ToString::to_string)
143 .unwrap_or_default();
144 let mut outcome = self.outcome_for(procedure, "anonymous", &identity);
145 let served = async {
146 self.serve_repository(req, &seams, &parsed, &repo, token)
147 .await
148 }
149 .instrument(outcome.span.clone())
150 .await;
151 match served {
152 Ok(response) => {
153 let code = match response.status {
154 ..400 => None,
155 402 | 403 | 451 => Some(Code::PermissionDenied),
156 416 => Some(Code::OutOfRange),
157 404 => Some(Code::NotFound),
158 _ => Some(Code::Unknown),
159 };
160 match code {
161 None => outcome.record(Ok(())),
162 Some(code) => outcome.record(Err(&ServerError::new(code, "http"))),
163 }
164 response
165 }
166 Err(fail) => {
167 outcome.record(Err(&ServerError::new(fail.code(), "http")));
168 fail.into_response()
169 }
170 }
171 }
172
173 #[allow(clippy::too_many_lines)] async fn serve_repository(
175 &self,
176 req: &HttpObjectRequest<'_>,
177 seams: &HttpSeams,
178 parsed: &ParsedUrl,
179 repo: &RepoId,
180 token: Option<Result<crate::url_token::Prechecked, crate::url_token::TokenRejected>>,
181 ) -> Result<HttpObjectResponse, Fail> {
182 let head = req.method == "HEAD";
183 let (Some(cfg), Some(indexed)) = (&self.cfg.http_objects, &self.cfg.indexed) else {
184 return Err(Fail::Unavailable);
185 };
186 let ref_name = match &parsed.target {
188 Target::Ref { name, .. } => Some(name.clone()),
189 Target::Object(_) => None,
190 };
191 let op = Operation::new(
192 repo.clone(),
193 Principal::Anonymous,
194 None,
195 OpKind::HttpGet { ref_name },
196 );
197 let expiry = self
198 .authorize_http_read(&op, &parsed.target, seams, token)
199 .await?;
200
201 let denial_budget = crate::indexed::budget::SliceBudget::new(9000);
202 let proof_meta = crate::indexed::budget::Budgeted::new(&self.meta, &denial_budget);
203 let proof_blobs = crate::indexed::budget::Budgeted::new(&self.blobs, &denial_budget);
204 let view = crate::store::view::ViewStore {
205 store: &proof_meta,
206 repo,
207 writer: false,
208 policy: self.publication_policy.as_deref(),
209 };
210 let env = Env {
211 no_reads: &std::collections::BTreeSet::new(),
212 blobs: &proof_blobs,
213 meta: &view,
214 shards: self.shards.as_ref(),
215 repo,
216 indexed,
217 cfg,
218 metrics: self.metrics.as_ref(),
219 };
220 let mut budget = Budget(cfg.http_decode_budget);
221 let now = ms(self.clock.now_ms());
222
223 let (leaf_id, commit, located) = match &parsed.target {
225 Target::Ref { name, path } => {
226 let shard = self.shards.ref_shard(repo, name);
227 let tip = read::read_ref(&view, &shard, &repo.name, name)
228 .await
229 .map_err(|error| {
230 tracing::warn!(detail = %error, "ref read failed");
231 Fail::Unavailable
232 })?
233 .ok_or(Fail::NotFound)?;
234 let resolved = resolve::resolve_ref(&env, tip, path, &mut budget).await?;
235 let located = resolve::locate(&env, resolved.leaf).await?;
236 for id in [&resolved.commit, &resolved.leaf] {
238 seams.reachability.record(repo, id, now);
239 }
240 (resolved.leaf, Some(resolved.commit), located)
241 }
242 Target::Object(id) => {
243 let located = resolve::locate(&env, *id).await?;
244 if self.cfg.takedown_denial
245 || !seams
246 .reachability
247 .known_reachable(repo, id, now)
248 .await
249 .map_err(|error| Fail::from_server_error(&error))?
250 {
251 match self.prove_reachable(&env, seams, id, &mut budget).await {
252 Err(Fail::Unavailable) if denial_budget.remaining() == 0 => {
254 return Err(Fail::NotFound);
255 }
256 other => other?,
257 }
258 seams.reachability.record(repo, id, now);
259 }
260 (*id, None, located)
261 }
262 };
263
264 for id in [&leaf_id, &located.pack] {
265 crate::takedown::denial::require_clear(&proof_meta, id)
266 .await
267 .map_err(|e| {
268 if e.public_message() == "object blocked" {
269 Fail::NotFound
270 } else {
271 Fail::Unavailable
272 }
273 })?;
274 }
275 if self.cfg.takedown_denial {
276 crate::takedown::denial::require_object_clear(
277 &self.meta,
278 self.shards.as_ref(),
279 repo,
280 &leaf_id,
281 indexed,
282 &denial_budget,
283 )
284 .await
285 .map_err(|e| {
286 if e.public_message() == "object blocked" {
287 Fail::NotFound
288 } else {
289 Fail::Unavailable
290 }
291 })?;
292 }
293 let takedown = seams
296 .takedown
297 .check(repo, &leaf_id)
298 .await
299 .map_err(|error| Fail::from_server_error(&error))?;
300 if matches!(takedown, TakedownVerdict::NotFound) {
301 return Err(Fail::NotFound);
302 }
303
304 let context = if parsed.query.proof {
307 let (proof_commit, path) = match &parsed.target {
308 Target::Ref { path, .. } => (commit.ok_or(Fail::NotFound)?, path.as_slice()),
309 Target::Object(_) => (
310 parsed.query.commit.ok_or(Fail::NotFound)?,
311 parsed.query.path.as_deref().ok_or(Fail::NotFound)?,
312 ),
313 };
314 if !ref_path_target(&parsed.target) {
315 self.prove_reachable(&env, seams, &proof_commit, &mut budget)
316 .await?;
317 }
318 Some(
319 crate::http_objects::proof::context(
320 &env,
321 seams.takedown.as_ref(),
322 proof_commit,
323 path,
324 leaf_id,
325 &mut budget,
326 )
327 .await?,
328 )
329 } else {
330 None
331 };
332 if let TakedownVerdict::Respond(response) = takedown {
334 return Ok(response);
335 }
336 let metadata = if let Some(context) = &context {
337 Some(context.metadata(&env, &mut budget).await?)
338 } else {
339 None
340 };
341 let ref_path = matches!(parsed.target, Target::Ref { .. });
342 let mut inline = Budget(cfg.max_inline_object_bytes);
343 let leaf = if context.is_none() {
346 Some(resolve::open_leaf(&env, leaf_id, located, &mut inline).await?)
347 } else {
348 None
349 };
350 let proof = if let (Some(context), Some(metadata)) = (&context, &metadata) {
354 match context
355 .select(&env, seams.takedown.as_ref(), metadata, parsed.query.range)
356 .await
357 {
358 Ok(proof) => Some(Ok(proof)),
359 Err(Fail::ProofRange) => Some(Err(Fail::ProofRange)),
360 Err(other) => return Err(other),
361 }
362 } else {
363 None
364 };
365 let etag = context.as_ref().map_or_else(
366 || format!("\"{}\"", to_hex(&leaf_id)),
367 |c| c.etag(parsed.query.range),
368 );
369 let ty = metadata
370 .as_ref()
371 .map(|c| c.ty)
372 .or_else(|| leaf.as_ref().map(|l| l.ty))
373 .ok_or(Fail::Unavailable)?;
374 let metadata_commit = context.as_ref().map(|c| c.commit).or(commit);
375 let values = |name: &str| (req.headers)(name);
376 let mut success = Vec::with_capacity(8);
377 success.push(("ETag", etag.clone()));
378 success.push(("X-Mkit-Object", to_hex(&leaf_id)));
379 success.push(("X-Mkit-Object-Type", ty.name().to_owned()));
380 if let Some(commit) = &metadata_commit {
381 success.push(("X-Mkit-Commit", to_hex(commit)));
382 }
383
384 let private_policy = cfg.admit_reads || seams.admission.is_configured();
387
388 if range::if_none_match(&values("if-none-match"), &etag) {
390 let mut response = HttpObjectResponse::new(304).with_header(
391 "Cache-Control",
392 super::http_tokens::cache(ref_path, private_policy, expiry, self.clock.now_ms()),
393 );
394 response.headers.extend(success);
395 return Ok(response);
396 }
397
398 let proof = proof.transpose()?;
400 if proof.is_some() && !seams.proofs.is_supported() {
401 return Err(Fail::ProofRange);
402 }
403 let window = if let Some(leaf) = &leaf {
404 let single = |name: &str| {
405 let mut all = values(name);
406 (all.len() == 1).then(|| all.remove(0))
407 };
408 let if_range = match values("if-range").len() {
409 0 => None,
410 1 => single("if-range"),
411 _ => Some(String::new()),
412 };
413 match range::select(
414 single("range").as_deref(),
415 if_range.as_deref(),
416 &etag,
417 leaf.len,
418 ) {
419 Selection::Full => None,
420 Selection::Partial { start, end } => Some((start, end)),
421 Selection::Unsatisfiable => {
422 return Ok(HttpObjectResponse::error(416)
423 .with_header("Content-Range", format!("bytes */{}", leaf.len)));
424 }
425 }
426 } else {
427 None
428 };
429 let selected_len = proof
430 .as_ref()
431 .map(|p| p.encoded_len)
432 .or_else(|| {
433 leaf.as_ref()
434 .map(|l| window.map_or(l.len, |(a, b)| b - a + 1))
435 })
436 .ok_or(Fail::Unavailable)?;
437
438 let first = |name: &str| (req.headers)(name).into_iter().next();
440 let credentials = {
441 let meta = super::RequestMeta {
442 procedure: op.procedure(),
443 header: &first,
444 header_values: Some(req.headers),
445 unary_body: None,
446 transport_principal: None,
447 };
448 if private_policy {
449 super::admission::validate_credentials(
450 &super::admission::capture_http_credentials(
451 &meta,
452 &self.cfg.admission_credential_headers,
453 req.header_names,
454 )
455 .map_err(|e| Fail::from_server_error(&e))?,
456 )
457 .map_err(|e| Fail::from_server_error(&e))?
458 } else {
459 Vec::new()
460 }
461 };
462 let request = AdmitRequest {
463 repo,
464 procedure: op.procedure(),
465 head,
466 ref_path,
467 declared_bytes: selected_len,
468 credential_headers: &credentials,
469 };
470 let (admitted, mut finalizer) = if cfg.admit_reads {
471 match self.admit_http_read(seams, &op, &request, leaf_id).await {
472 Ok(result) => result,
473 Err(error) if error.http_status() == Some(402) => {
474 return Ok(super::http_admission::challenge_response(&error, head));
475 }
476 Err(error) => return Err(Fail::from_server_error(&error)),
477 }
478 } else {
479 match seams
480 .admission
481 .admit(&request)
482 .await
483 .map_err(|e| Fail::from_server_error(&e))?
484 {
485 AdmitDecision::Allow(admitted) => (admitted, None),
486 AdmitDecision::Respond(response) => return Ok(response),
487 }
488 };
489
490 if cfg.redirect_public_refs
493 && ref_path
494 && proof.is_none()
495 && expiry.is_none()
496 && !private_policy
497 {
498 let prefix = req
499 .raw_path
500 .split_once("/-/")
501 .map_or("", |(prefix, _)| prefix);
502 return Ok(HttpObjectResponse::new(302)
503 .with_header(
504 "Location",
505 format!("{prefix}/-/objects/{}", to_hex(&leaf_id)),
506 )
507 .with_header("Cache-Control", "no-cache")
508 .with_header("Content-Length", "0"));
509 }
510
511 let mut hook = admitted.on_end;
513 let body = if head || selected_len == 0 {
514 if let Some(hook) = hook.take() {
515 hook(0, Ok(()));
516 }
517 if let Some(finalizer) = finalizer.take() {
518 finalizer.complete(true).await;
519 }
520 HttpBody::Empty
521 } else {
522 let opened = if let Some(proof) = &proof {
523 let mut source = crate::http_objects::proof::RepositorySource {
524 env: Env {
525 no_reads: &std::collections::BTreeSet::new(),
526 blobs: &proof_blobs,
527 meta: &view,
528 shards: self.shards.as_ref(),
529 repo,
530 indexed,
531 cfg,
532 metrics: self.metrics.as_ref(),
533 },
534 budget: Budget(cfg.http_decode_budget),
535 gate: seams.takedown.as_ref(),
536 };
537 match seams.proofs.build(proof, &mut source).await {
538 Ok(bytes) if bytes.len() as u64 == selected_len => {
539 Ok(crate::http_objects::body_with_hook(
540 HttpBody::Bytes(bytes.into()),
541 hook.take(),
542 ))
543 }
544 _ => Err(resolve::Miss::Unavailable),
545 }
546 } else if let Some(leaf) = &leaf {
547 resolve::open_body(&proof_blobs, leaf, window, &mut hook).await
548 } else {
549 Err(resolve::Miss::Unavailable)
550 };
551 match opened {
552 Ok(body) => body,
553 Err(miss) => {
554 if let Some(hook) = hook.take() {
555 hook(0, Err(&ServerError::unavailable("object unavailable")));
556 }
557 if let Some(finalizer) = finalizer.take() {
558 finalizer.complete(false).await;
559 }
560 return Err(miss.into());
561 }
562 }
563 };
564 let filename = match (&parsed.target, ty, proof.is_none()) {
565 (Target::Ref { path, .. }, ObjectType::Blob | ObjectType::ChunkedBlob, true) => {
566 path.last().map(Vec::as_slice)
567 }
568 _ => None,
569 };
570 let file_headers = filename.map(crate::http_objects::content_headers::media_type);
571 let mut response = HttpObjectResponse::new(if window.is_some() { 206 } else { 200 })
572 .with_header(
573 "Accept-Ranges",
574 if proof.is_some() { "none" } else { "bytes" },
575 )
576 .with_header(
577 "Cache-Control",
578 super::http_tokens::cache(
579 ref_path,
580 admitted.private || private_policy,
581 expiry,
582 self.clock.now_ms(),
583 ),
584 )
585 .with_header("Content-Length", selected_len.to_string())
586 .with_header(
587 "Content-Type",
588 match &proof {
589 Some(p) if p.span => "application/vnd.mkit.disclosure-span",
590 Some(_) => "application/vnd.mkit.disclosure",
591 None => match ty {
592 ObjectType::Blob | ObjectType::ChunkedBlob => {
593 file_headers.map_or("application/octet-stream", |(media, _)| media)
594 }
595 _ => "application/vnd.mkit.object",
596 },
597 },
598 );
599 if let (Some((start, end)), Some(leaf)) = (window, &leaf) {
600 response =
601 response.with_header("Content-Range", format!("bytes {start}-{end}/{}", leaf.len));
602 }
603 if let (Some(name), Some((_, kind))) = (filename, file_headers) {
604 response = response.with_header(
605 "Content-Disposition",
606 crate::http_objects::content_headers::disposition(name, kind),
607 );
608 }
609 response.headers.extend(success);
610 response.headers.extend(admitted.headers);
611 response.body = match finalizer {
612 Some(f) => f.wrap(body),
613 None => body,
614 };
615 Ok(response)
616 }
617
618 async fn prove_reachable(
620 &self,
621 env: &Env<'_, impl crate::BlobStore, impl NamespaceStore>,
622 seams: &HttpSeams,
623 id: &Hash,
624 budget: &mut Budget,
625 ) -> Result<(), Fail> {
626 let (tips, truncated) = self
627 .reader_tips(env.meta, env.repo, env.cfg.max_walk_objects, false)
628 .await?;
629 let reached = if truncated {
630 Reach::Capped
631 } else {
632 reach::walk(env, seams.takedown.as_ref(), &tips, *id, budget).await?
633 };
634 match reached {
635 Reach::Reachable => Ok(()),
636 Reach::Unreachable => Err(Fail::NotFound),
637 Reach::Capped => {
638 tracing::warn!("reachability walk hit a cap");
639 self.metrics.incr(METRIC_HTTP_REACH_CAPPED, &[], 1);
640 Err(Fail::NotFound)
641 }
642 }
643 }
644
645 pub(super) async fn reader_tips(
647 &self,
648 store: &impl NamespaceStore,
649 repo: &RepoId,
650 cap: usize,
651 writer: bool,
652 ) -> Result<(Vec<Hash>, bool), Fail> {
653 let view = crate::store::view::ViewStore {
654 store,
655 repo,
656 writer,
657 policy: self.publication_policy.as_deref(),
658 };
659 let scan = crate::refs::list_scan_prefix("refs/");
660 let partitions = self.shards.ref_index_partitions(repo);
661 let (mut tips, mut last) = (Vec::new(), None::<String>);
662 let mut rows_left = cap;
663 let mut pages_left = cap.div_ceil(self.cfg.list_page_limit as usize);
664 loop {
665 let limit = self
669 .cfg
670 .list_page_limit
671 .min(u32::try_from(rows_left / partitions.len()).unwrap_or(u32::MAX));
672 if limit == 0 || pages_left == 0 {
673 return Ok((tips, true));
674 }
675 rows_left -= limit as usize * partitions.len();
676 pages_left -= 1;
677 let page = if partitions.len() == 1 {
678 let bucket = super::list::RefBucket {
679 store: &view,
680 partition: &partitions[0],
681 };
682 super::list::page(
683 &[bucket],
684 repo,
685 &scan,
686 last.as_deref(),
687 limit,
688 super::list::MAX_RESPONSE_BYTES,
689 )
690 .await
691 } else {
692 let buckets: Vec<_> = partitions
693 .iter()
694 .map(|partition| super::list::IndexBucket {
695 store: &view,
696 partition,
697 })
698 .collect();
699 super::list::page(
700 &buckets,
701 repo,
702 &scan,
703 last.as_deref(),
704 limit,
705 super::list::MAX_RESPONSE_BYTES,
706 )
707 .await
708 }
709 .map_err(|error| {
710 tracing::warn!(detail = %error, "ref listing scan failed");
711 Fail::Unavailable
712 })?;
713 for entry in page.refs {
714 if entry.name.starts_with(PACKMAP_PREFIX)
715 || !crate::refs::is_served_ref_name(&entry.name)
716 {
717 continue;
718 }
719 tips.push(entry.id);
720 }
721 let Some(next) = page.next else {
722 return Ok((tips, false));
723 };
724 last = Some(super::list::decode_token(repo, &scan, &next).ok_or(Fail::Unavailable)?);
725 }
726 }
727}
728
729fn ref_path_target(target: &Target) -> bool {
730 matches!(target, Target::Ref { .. })
731}
732
733fn repo_id(parsed: &ParsedUrl) -> Option<RepoId> {
735 let identity = parsed.repository.as_ref()?;
736 Some(RepoId {
737 namespace: NamespaceKey::from_namespace(identity.namespace()?),
738 name: RepoName::new(identity.name()).ok()?,
739 })
740}