This repository has no description
1mod archive;
2mod cache;
3mod error;
4mod fetch;
5mod frame;
6mod guard;
7mod ids;
8mod idxwrite;
9mod meter;
10mod objects;
11mod oids;
12mod pkt;
13mod quarantine;
14mod receive;
15mod receiver;
16mod resolve;
17mod upload;
18
19use std::collections::HashMap;
20use std::io::{self, Read};
21
22use axum::Router;
23use axum::body::{Body, Bytes};
24use axum::extract::{DefaultBodyLimit, Path, Query, State};
25use axum::http::{HeaderMap, header};
26use axum::response::Response;
27use axum::routing::post;
28use knot_git::{Layout, Repo};
29use knot_messages::{Catalog, ErrorKey, FetchMessages};
30use knot_resource::{PackSlots, SlotPermit};
31use knot_runtime::Clock;
32use knot_types::{
33 AccountDid, Handle, KnotHostname, OwnerDid, OwnerRef, ParseError, RepoDid, RepoRkey,
34};
35use std::sync::Arc;
36use tokio::sync::mpsc;
37use tokio_stream::wrappers::ReceiverStream;
38
39pub use cache::{CacheConfig, MaxCacheBytes, MaxEntryBytes};
40pub use error::{PackError, PackLimit};
41pub use fetch::{FetchError, UpstreamRefs, local_pack, local_refs, remote_pack, remote_refs};
42pub use frame::{
43 ReceiveFramer, UploadFramer, archive_request_complete, receive_request_complete, upload_v0_nak,
44};
45pub use guard::PushGuard;
46pub use ids::{DeltaDepth, MaxObjectBytes, MaxTotalBytes, MaxWireBytes};
47pub use meter::PackLimits;
48pub use objects::{ExpandedPack, count_expanded, write_expanded, write_pack};
49pub use oids::{HaveOids, WantOids};
50pub use pkt::frame_report;
51pub use quarantine::sweep_incoming;
52pub use receive::{Preflight, ReceiveCommand, ReceiveGuard, ReceiveOutcome, RefDecision};
53pub use receiver::{PackReceiver, ReceiveReadError, ReceivedPack};
54pub use upload::{SelectionLimits, init_selection_limits, selection_budget};
55
56use upload::UploadOutcome;
57
58pub use knot_messages::default_catalog;
59
60pub fn default_hostname() -> &'static KnotHostname {
61 static HOSTNAME: std::sync::LazyLock<KnotHostname> = std::sync::LazyLock::new(|| {
62 KnotHostname::new("knot.invalid").expect("static hostname parses")
63 });
64 &HOSTNAME
65}
66
67pub fn advertise_upload(repo: &Repo) -> Result<Vec<u8>, PackError> {
68 upload::advertise(repo)
69}
70
71pub fn advertise_upload_v0(repo: &Repo) -> Result<Vec<u8>, PackError> {
72 upload::advertise_v0(repo)
73}
74
75pub fn upload_pack(repo: &Repo, request: &[u8]) -> Result<Vec<u8>, PackError> {
76 upload::buffered(repo, request, &default_catalog().fetch, default_hostname())
77}
78
79pub fn upload_pack_streamed(
80 repo: &Repo,
81 request: &[u8],
82 messages: &FetchMessages,
83 knot: &KnotHostname,
84 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>,
85) -> Result<(), PackError> {
86 upload::streamed(repo, request, messages, knot, sink)
87}
88
89pub fn advertise_receive(repo: &Repo) -> Result<Vec<u8>, PackError> {
90 receive::advertise(repo)
91}
92
93pub fn advertise_upload_ssh(repo: &Repo) -> Result<Vec<u8>, PackError> {
94 upload::advertise_ssh(repo)
95}
96
97pub fn advertise_upload_v0_ssh(repo: &Repo) -> Result<Vec<u8>, PackError> {
98 upload::advertise_v0_ssh(repo)
99}
100
101pub fn advertise_receive_ssh(repo: &Repo) -> Result<Vec<u8>, PackError> {
102 receive::advertise_ssh(repo)
103}
104
105#[doc(hidden)]
106pub fn receive_pack(repo: &Repo, request: &[u8]) -> Result<Vec<u8>, PackError> {
107 receive::handle_bytes(repo, request, &PackLimits::default())
108}
109
110#[doc(hidden)]
111pub fn receive_pack_with_limits(
112 repo: &Repo,
113 request: &[u8],
114 limits: &PackLimits,
115) -> Result<Vec<u8>, PackError> {
116 receive::handle_bytes(repo, request, limits)
117}
118
119#[doc(hidden)]
120pub fn bench_ingest_fresh(
121 objects_dir: &std::path::Path,
122 pack_path: &std::path::Path,
123 kind: gix::hash::Kind,
124) -> Result<bool, PackError> {
125 let pack = gix_pack::data::File::at(pack_path, kind)
126 .map_err(|error| PackError::Pack(error.to_string()))?;
127 objects::ingest_and_close(
128 objects_dir,
129 &pack,
130 &PackLimits::default(),
131 kind,
132 knot_resource::ingest_base_budget(),
133 false,
134 )
135 .map(|closure| closure.is_some())
136}
137
138#[doc(hidden)]
139pub fn bench_ingest_external(
140 objects_dir: &std::path::Path,
141 pack_path: &std::path::Path,
142 kind: gix::hash::Kind,
143) -> Result<Option<bool>, PackError> {
144 let pack = gix_pack::data::File::at(pack_path, kind)
145 .map_err(|error| PackError::Pack(error.to_string()))?;
146 objects::admit_ingest(&pack, kind)?;
147 objects::ingest_and_close(
148 objects_dir,
149 &pack,
150 &PackLimits::default(),
151 kind,
152 knot_resource::ingest_base_budget(),
153 true,
154 )
155 .map(|closure| closure.map(|closure| closure.self_contained))
156}
157
158#[doc(hidden)]
159pub fn bench_admit_and_ingest(
160 objects_dir: &std::path::Path,
161 pack_path: &std::path::Path,
162 kind: gix::hash::Kind,
163) -> Result<bool, PackError> {
164 bench_ingest_with_base_budget(
165 objects_dir,
166 pack_path,
167 kind,
168 knot_resource::ingest_base_budget(),
169 )
170}
171
172#[doc(hidden)]
173pub fn bench_ingest_with_base_budget(
174 objects_dir: &std::path::Path,
175 pack_path: &std::path::Path,
176 kind: gix::hash::Kind,
177 base_budget: Option<usize>,
178) -> Result<bool, PackError> {
179 let pack = gix_pack::data::File::at(pack_path, kind)
180 .map_err(|error| PackError::Pack(error.to_string()))?;
181 objects::admit_ingest(&pack, kind)?;
182 objects::ingest_and_close(
183 objects_dir,
184 &pack,
185 &PackLimits::default(),
186 kind,
187 base_budget,
188 false,
189 )
190 .map(|closure| closure.is_some())
191}
192
193pub fn receive_pack_guarded(
194 repo: &Repo,
195 request: &[u8],
196 limits: &PackLimits,
197 guard: &dyn ReceiveGuard,
198 seal: &dyn Fn(&[knot_git::RefUpdate]),
199 messages: &knot_messages::RejectMessages,
200) -> Result<ReceiveOutcome, PackError> {
201 receive::handle_guarded_bytes(repo, request, limits, guard, seal, messages)
202}
203
204pub fn receive_pack_guarded_streamed(
205 repo: &Repo,
206 received: &ReceivedPack,
207 limits: &PackLimits,
208 guard: &dyn ReceiveGuard,
209 seal: &dyn Fn(&[knot_git::RefUpdate]),
210 messages: &knot_messages::RejectMessages,
211) -> Result<ReceiveOutcome, PackError> {
212 receive::handle_guarded_streamed(repo, received, limits, guard, seal, messages)
213}
214
215pub fn receive_preflight(request: &[u8]) -> Preflight {
216 receive::preflight(request)
217}
218
219pub fn upload_archive_streamed(
220 repo: &Repo,
221 request: &[u8],
222 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>,
223) -> Result<(), PackError> {
224 archive::stream(repo, request, sink)
225}
226
227pub fn upload_archive(repo: &Repo, request: &[u8]) -> Result<Vec<u8>, PackError> {
228 let mut buf = Vec::new();
229 upload_archive_streamed(repo, request, &mut |chunk| {
230 buf.extend_from_slice(chunk);
231 Ok(())
232 })?;
233 Ok(buf)
234}
235
236pub fn meter_pack(
237 pack: &[u8],
238 limits: &PackLimits,
239 kind: gix::hash::Kind,
240) -> Result<(), PackError> {
241 meter::meter(pack, limits, kind)
242}
243
244pub fn ingest_pack(
245 objects_dir: &std::path::Path,
246 pack: &[u8],
247 limits: &PackLimits,
248 kind: gix::hash::Kind,
249) -> Result<(), PackError> {
250 objects::index_pack(objects_dir, pack, limits, kind)
251}
252
253#[doc(hidden)]
254pub mod fuzz {
255 pub fn pkt(data: &[u8]) {
256 let _ = crate::pkt::data_payloads(data);
257 let _ = crate::pkt::data_payloads_all(data);
258 let _ = crate::pkt::split_receive(data);
259 }
260
261 pub fn pack(data: &[u8]) {
262 let _ = crate::meter::meter(data, &crate::PackLimits::default(), gix::hash::Kind::Sha1);
263 }
264
265 pub fn receive_commands(data: &[u8]) {
266 crate::receive::fuzz(data);
267 }
268
269 pub fn upload_args(data: &[u8]) {
270 crate::upload::fuzz(data);
271 }
272}
273
274const MAX_REQUEST_BYTES: usize = 16 * 1024 * 1024;
275
276#[derive(Debug, Clone, PartialEq, Eq)]
277pub enum RepoTarget {
278 Did(RepoDid),
279 OwnerRkey(OwnerDid, RepoRkey),
280}
281
282#[derive(Debug, Clone, PartialEq, Eq)]
283pub enum RepoLookup {
284 Hosted(RepoDid),
285 Unhosted,
286 Unavailable,
287}
288
289impl RepoLookup {
290 pub fn or_else(self, next: impl FnOnce() -> RepoLookup) -> RepoLookup {
291 match self {
292 RepoLookup::Unhosted => next(),
293 decided => decided,
294 }
295 }
296
297 pub fn first(
298 candidates: impl IntoIterator<Item = RepoRkey>,
299 resolve: impl Fn(RepoRkey) -> RepoLookup,
300 ) -> RepoLookup {
301 candidates
302 .into_iter()
303 .fold(RepoLookup::Unhosted, |acc, rkey| {
304 acc.or_else(|| resolve(rkey))
305 })
306 }
307
308 pub fn from_resolved<T>(
309 resolved: knot_index::Resolved<Option<T>>,
310 found: impl FnOnce(T) -> RepoDid,
311 ) -> RepoLookup {
312 match resolved {
313 knot_index::Resolved::Ready(Some(value)) => RepoLookup::Hosted(found(value)),
314 knot_index::Resolved::Ready(None) => RepoLookup::Unhosted,
315 knot_index::Resolved::Warming => RepoLookup::Unavailable,
316 }
317 }
318}
319
320pub trait RepoResolver: Send + Sync + 'static {
321 fn resolve(&self, target: &RepoTarget) -> RepoLookup;
322}
323
324pub trait HandleResolver: Send + Sync + 'static {
325 fn resolve(
326 &self,
327 handle: Handle,
328 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Option<AccountDid>> + Send + '_>>;
329}
330
331pub use knot_edge::SocketPeer;
332
333pub trait ReceiveAdvertiser: Send + Sync + 'static {
334 fn advertise(
335 &self,
336 repo: RepoDid,
337 peer: knot_edge::SocketPeer,
338 headers: HeaderMap,
339 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Response> + Send + '_>>;
340}
341
342impl<F> RepoResolver for F
343where
344 F: Fn(&RepoTarget) -> RepoLookup + Send + Sync + 'static,
345{
346 fn resolve(&self, target: &RepoTarget) -> RepoLookup {
347 self(target)
348 }
349}
350
351#[derive(Clone)]
352struct PackState {
353 layout: Layout,
354 resolver: Arc<dyn RepoResolver>,
355 receive: Option<Arc<dyn ReceiveAdvertiser>>,
356 handle_resolver: Option<Arc<dyn HandleResolver>>,
357 pack_slots: PackSlots,
358 cache: Arc<cache::PackCache>,
359 catalog: Arc<Catalog>,
360 hostname: KnotHostname,
361}
362
363pub fn router(layout: Layout, resolver: Arc<dyn RepoResolver>, clock: Arc<dyn Clock>) -> Router {
364 router_with_pack_slots(
365 layout,
366 resolver,
367 PackSlots::new(knot_resource::threads().get()),
368 clock,
369 )
370}
371
372pub fn router_with_pack_slots(
373 layout: Layout,
374 resolver: Arc<dyn RepoResolver>,
375 pack_slots: PackSlots,
376 clock: Arc<dyn Clock>,
377) -> Router {
378 let state = pack_state(
379 layout,
380 resolver,
381 None,
382 None,
383 pack_slots,
384 CacheConfig::default(),
385 Arc::new(Catalog::defaults()),
386 default_hostname().clone(),
387 clock,
388 );
389 write_routes(state.clone()).merge(advertisement_routes(state).into_router())
390}
391
392#[allow(clippy::too_many_arguments)]
393pub fn edge_routes(
394 layout: Layout,
395 resolver: Arc<dyn RepoResolver>,
396 receive: Option<Arc<dyn ReceiveAdvertiser>>,
397 handle_resolver: Option<Arc<dyn HandleResolver>>,
398 pack_slots: PackSlots,
399 cache: CacheConfig,
400 catalog: Arc<Catalog>,
401 hostname: KnotHostname,
402 clock: Arc<dyn Clock>,
403) -> (Router, knot_edge::ZeroRttRoutes) {
404 let state = pack_state(
405 layout,
406 resolver,
407 receive,
408 handle_resolver,
409 pack_slots,
410 cache,
411 catalog,
412 hostname,
413 clock,
414 );
415 (write_routes(state.clone()), advertisement_routes(state))
416}
417
418#[allow(clippy::too_many_arguments)]
419fn pack_state(
420 layout: Layout,
421 resolver: Arc<dyn RepoResolver>,
422 receive: Option<Arc<dyn ReceiveAdvertiser>>,
423 handle_resolver: Option<Arc<dyn HandleResolver>>,
424 pack_slots: PackSlots,
425 cache: CacheConfig,
426 catalog: Arc<Catalog>,
427 hostname: KnotHostname,
428 clock: Arc<dyn Clock>,
429) -> PackState {
430 PackState {
431 layout,
432 resolver,
433 receive,
434 handle_resolver,
435 pack_slots,
436 cache: cache::PackCache::new(cache, clock),
437 catalog,
438 hostname,
439 }
440}
441
442fn write_routes(state: PackState) -> Router {
443 Router::new()
444 .route("/{did}/{name}/git-upload-pack", post(upload_named))
445 .route(
446 "/{did}/{name}/git-upload-archive",
447 post(upload_archive_named),
448 )
449 .route("/{did}/git-upload-pack", post(upload_did))
450 .route("/{did}/git-upload-archive", post(upload_archive_did))
451 .layer(DefaultBodyLimit::max(MAX_REQUEST_BYTES))
452 .with_state(state)
453}
454
455fn advertisement_routes(state: PackState) -> knot_edge::ZeroRttRoutes {
456 let named_state = state.clone();
457 let did_state = state;
458 knot_edge::ZeroRttRoutes::new()
459 .get(
460 "/{did}/{name}/info/refs",
461 knot_edge::ZeroRttSafe::new(
462 move |Path((did, name)): Path<(String, String)>,
463 Query(query): Query<HashMap<String, String>>,
464 peer: knot_edge::SocketPeer,
465 headers: HeaderMap| {
466 let state = named_state.clone();
467 async move {
468 let owner = resolve_owner(&state, &did).await?;
469 let repo_did = resolve_named_did(&state, &owner, &name)?;
470 read_advertisement(
471 &state,
472 repo_did,
473 query.get("service").map(String::as_str),
474 peer,
475 &headers,
476 )
477 .await
478 }
479 },
480 ),
481 )
482 .get(
483 "/{did}/info/refs",
484 knot_edge::ZeroRttSafe::new(
485 move |Path(did): Path<String>,
486 Query(query): Query<HashMap<String, String>>,
487 peer: knot_edge::SocketPeer,
488 headers: HeaderMap| {
489 let state = did_state.clone();
490 async move {
491 let repo_did = resolve_did_did(&state, &did)?;
492 read_advertisement(
493 &state,
494 repo_did,
495 query.get("service").map(String::as_str),
496 peer,
497 &headers,
498 )
499 .await
500 }
501 },
502 ),
503 )
504}
505
506fn lookup_did(lookup: RepoLookup) -> Result<RepoDid, PackError> {
507 match lookup {
508 RepoLookup::Hosted(did) => Ok(did),
509 RepoLookup::Unhosted => Err(PackError::NotFound),
510 RepoLookup::Unavailable => Err(PackError::Unavailable),
511 }
512}
513
514fn bad_path(error: ParseError) -> PackError {
515 PackError::BadPath(error.to_string())
516}
517
518async fn resolve_owner(state: &PackState, owner: &str) -> Result<OwnerDid, PackError> {
519 match OwnerRef::parse(owner).ok_or(PackError::NotFound)? {
520 OwnerRef::Did(did) => Ok(did),
521 OwnerRef::Handle(handle) => {
522 let resolver = state.handle_resolver.as_ref().ok_or(PackError::NotFound)?;
523 let did = resolver.resolve(handle).await.ok_or(PackError::NotFound)?;
524 Ok(did.into())
525 }
526 }
527}
528
529fn resolve_named_did(
530 state: &PackState,
531 owner: &OwnerDid,
532 name: &str,
533) -> Result<RepoDid, PackError> {
534 lookup_did(RepoLookup::first(
535 RepoRkey::clone_path_candidates(name),
536 |rkey| {
537 state
538 .resolver
539 .resolve(&RepoTarget::OwnerRkey(owner.clone(), rkey))
540 },
541 ))
542}
543
544fn resolve_did_did(state: &PackState, did: &str) -> Result<RepoDid, PackError> {
545 let did = RepoDid::new(did).map_err(bad_path)?;
546 lookup_did(state.resolver.resolve(&RepoTarget::Did(did)))
547}
548
549fn open_named(state: &PackState, owner: &OwnerDid, name: &str) -> Result<Repo, PackError> {
550 let did = resolve_named_did(state, owner, name)?;
551 state.layout.open(&did).map_err(|_| PackError::NotFound)
552}
553
554fn open_did(state: &PackState, did: &str) -> Result<Repo, PackError> {
555 let did = resolve_did_did(state, did)?;
556 state.layout.open(&did).map_err(|_| PackError::NotFound)
557}
558
559async fn read_advertisement(
560 state: &PackState,
561 repo_did: RepoDid,
562 service: Option<&str>,
563 peer: knot_edge::SocketPeer,
564 headers: &HeaderMap,
565) -> Result<Response, PackError> {
566 match service {
567 Some("git-upload-pack") => {
568 let repo = state
569 .layout
570 .open(&repo_did)
571 .map_err(|_| PackError::NotFound)?;
572 let body = if wants_v2(headers) {
573 advertise_upload(&repo)?
574 } else {
575 upload::advertise_v0(&repo)?
576 };
577 Ok(git_response(
578 "application/x-git-upload-pack-advertisement",
579 body,
580 ))
581 }
582 Some("git-receive-pack") => match &state.receive {
583 Some(advertiser) => Ok(advertiser.advertise(repo_did, peer, headers.clone()).await),
584 None => Err(PackError::PushOverSsh),
585 },
586 _ => Err(PackError::UnsupportedService),
587 }
588}
589
590fn decode_request(headers: &HeaderMap, body: Bytes) -> Result<Vec<u8>, PackError> {
591 let codings: Vec<&str> = headers
592 .get(header::CONTENT_ENCODING)
593 .and_then(|value| value.to_str().ok())
594 .map(|value| {
595 value
596 .split(',')
597 .map(str::trim)
598 .filter(|token| !token.is_empty() && !token.eq_ignore_ascii_case("identity"))
599 .collect()
600 })
601 .unwrap_or_default();
602 match codings.as_slice() {
603 [] => Ok(body.to_vec()),
604 [token] if token.eq_ignore_ascii_case("gzip") => {
605 let mut out = Vec::new();
606 flate2::read::GzDecoder::new(body.as_ref())
607 .take(MAX_REQUEST_BYTES as u64 + 1)
608 .read_to_end(&mut out)
609 .map_err(|error| PackError::Protocol(format!("gzip request body: {error}")))?;
610 match out.len() > MAX_REQUEST_BYTES {
611 true => Err(PackError::LimitExceeded(PackLimit::TotalBytes)),
612 false => Ok(out),
613 }
614 }
615 unsupported => Err(PackError::UnsupportedEncoding(unsupported.join(", "))),
616 }
617}
618
619fn wants_v2(headers: &HeaderMap) -> bool {
620 headers
621 .get("git-protocol")
622 .and_then(|value| value.to_str().ok())
623 .is_some_and(|value| value.split(':').any(|token| token.trim() == "version=2"))
624}
625
626fn refs_digest(repo: &Repo) -> Option<ids::RefsDigest> {
627 let refs = repo
628 .advertised_refs_for(knot_git::AdvertScope::Upload)
629 .ok()?;
630 let mut hasher = gix::hash::hasher(gix::hash::Kind::Sha256);
631 refs.iter().for_each(|record| {
632 hasher.update(record.name.as_str().as_bytes());
633 hasher.update(b"\0");
634 hasher.update(record.target.to_string().as_bytes());
635 hasher.update(b"\n");
636 });
637 let id = hasher.try_finalize().ok()?;
638 let mut token = [0u8; 32];
639 token.copy_from_slice(id.as_slice());
640 Some(ids::RefsDigest::new(token))
641}
642
643const CACHE_RETRY_BUDGET: usize = 4;
644
645async fn upload_dispatch(
646 state: &PackState,
647 repo: Repo,
648 body: &[u8],
649) -> Result<Response, PackError> {
650 upload_dispatch_within(state, repo, body, CACHE_RETRY_BUDGET).await
651}
652
653async fn upload_dispatch_within(
654 state: &PackState,
655 repo: Repo,
656 body: &[u8],
657 retries: usize,
658) -> Result<Response, PackError> {
659 let Some(token) = refs_digest(&repo) else {
660 return dispatch_plan(state, repo, body, None).await;
661 };
662 let key = cache::RequestKey::new(&repo.objects_dir(), &token, body);
663 match state.cache.decide(key) {
664 cache::Decision::Serve(bytes) => Ok(cached_response(bytes)),
665 cache::Decision::Await(receiver) => match cache::wait(receiver).await {
666 cache::Resolved::Bytes(bytes) => Ok(cached_response(bytes)),
667 cache::Resolved::Retry => match retries {
668 0 => dispatch_plan(state, repo, body, None).await,
669 _ => Box::pin(upload_dispatch_within(state, repo, body, retries - 1)).await,
670 },
671 cache::Resolved::Regenerate => dispatch_plan(state, repo, body, None).await,
672 },
673 cache::Decision::Lead(lease) => dispatch_plan(state, repo, body, Some(lease)).await,
674 cache::Decision::Stream | cache::Decision::Off => {
675 dispatch_plan(state, repo, body, None).await
676 }
677 }
678}
679
680async fn dispatch_plan(
681 state: &PackState,
682 repo: Repo,
683 body: &[u8],
684 lease: Option<cache::Lease>,
685) -> Result<Response, PackError> {
686 let permit = state.pack_slots.acquire().await;
687 let owned_body = body.to_vec();
688 let planned = tokio::task::spawn_blocking(move || {
689 let outcome = upload::plan(&repo, &owned_body);
690 (repo, outcome)
691 })
692 .await;
693 let (repo, plan) = match planned {
694 Ok((repo, Ok(plan))) => (repo, plan),
695 Ok((_repo, Err(error))) => {
696 if let Some(lease) = lease {
697 lease.regenerate();
698 }
699 return Err(error);
700 }
701 Err(_join) => {
702 if let Some(lease) = lease {
703 lease.regenerate();
704 }
705 return Err(PackError::Pack("upload planning task panicked".to_string()));
706 }
707 };
708 match plan {
709 UploadOutcome::Buffered(bytes) => {
710 drop(permit);
711 if let Some(lease) = lease {
712 lease.regenerate();
713 }
714 Ok(git_response("application/x-git-upload-pack-result", bytes))
715 }
716 UploadOutcome::Streaming {
717 preamble,
718 wants,
719 haves,
720 opts,
721 } => Ok(stream_response(
722 repo,
723 preamble,
724 wants,
725 haves,
726 opts,
727 permit,
728 lease,
729 Arc::clone(&state.catalog),
730 state.hostname.clone(),
731 )),
732 }
733}
734
735fn cached_response(bytes: Bytes) -> Response {
736 nocache(
737 Response::builder().header(header::CONTENT_TYPE, "application/x-git-upload-pack-result"),
738 )
739 .body(Body::from(bytes))
740 .expect("valid response")
741}
742
743#[allow(clippy::too_many_arguments)]
744fn stream_response(
745 repo: Repo,
746 preamble: Vec<u8>,
747 wants: WantOids,
748 haves: HaveOids,
749 opts: upload::StreamOpts,
750 permit: SlotPermit,
751 lease: Option<cache::Lease>,
752 catalog: Arc<Catalog>,
753 hostname: KnotHostname,
754) -> Response {
755 let side_band = opts.side_band;
756 let limit = lease.as_ref().map(cache::Lease::max_entry_bytes);
757 let (tx, rx) = mpsc::channel::<Result<Bytes, io::Error>>(16);
758 tokio::task::spawn_blocking(move || {
759 let _permit = permit;
760 let capture = std::cell::RefCell::new(cache::Capture::new(limit));
761 let mut emit = |chunk: &[u8]| -> io::Result<()> {
762 capture.borrow_mut().record(chunk);
763 tx.blocking_send(Ok(Bytes::copy_from_slice(chunk)))
764 .map_err(|_| io::Error::other("client disconnected"))
765 };
766 if emit(&preamble).is_err() {
767 if let Some(lease) = lease {
768 lease.retry();
769 }
770 return;
771 }
772 let result = upload::stream_pack(
773 &repo,
774 &wants,
775 &haves,
776 &opts,
777 &catalog.fetch,
778 &hostname,
779 &mut emit,
780 );
781 match result {
782 Ok(()) => {
783 if side_band {
784 let mut flush = Vec::new();
785 if pkt::write_flush(&mut flush).is_ok() {
786 let _ = emit(&flush);
787 }
788 }
789 if let Some(lease) = lease {
790 match capture.into_inner().into_bytes() {
791 Some(bytes) => lease.ready(Bytes::from(bytes)),
792 None => lease.too_large(),
793 }
794 }
795 }
796 Err(error) if side_band => {
797 let mut tail = Vec::new();
798 let line = catalog
799 .fetch
800 .fatal
801 .line(|ErrorKey::Error| error.to_string().replace('\n', " "));
802 let message = format!("{line}\n");
803 if pkt::write_band_error(&mut tail, message.as_bytes()).is_ok() {
804 let _ = pkt::write_flush(&mut tail);
805 let _ = emit(&tail);
806 }
807 if let Some(lease) = lease {
808 lease.regenerate();
809 }
810 }
811 Err(error) => {
812 let _ = tx.blocking_send(Err(io::Error::other(error.to_string())));
813 if let Some(lease) = lease {
814 lease.regenerate();
815 }
816 }
817 }
818 });
819 nocache(
820 Response::builder().header(header::CONTENT_TYPE, "application/x-git-upload-pack-result"),
821 )
822 .body(Body::from_stream(ReceiverStream::new(rx)))
823 .expect("valid response")
824}
825
826async fn upload_named(
827 State(state): State<PackState>,
828 Path((did, name)): Path<(String, String)>,
829 headers: HeaderMap,
830 body: Bytes,
831) -> Result<Response, PackError> {
832 let owner = resolve_owner(&state, &did).await?;
833 let repo = open_named(&state, &owner, &name)?;
834 upload_dispatch(&state, repo, &decode_request(&headers, body)?).await
835}
836
837async fn upload_did(
838 State(state): State<PackState>,
839 Path(did): Path<String>,
840 headers: HeaderMap,
841 body: Bytes,
842) -> Result<Response, PackError> {
843 let repo = open_did(&state, &did)?;
844 upload_dispatch(&state, repo, &decode_request(&headers, body)?).await
845}
846
847async fn archive_dispatch(
848 state: &PackState,
849 repo: Repo,
850 body: Vec<u8>,
851) -> Result<Response, PackError> {
852 let permit = state.pack_slots.acquire().await;
853 Ok(archive_response(repo, body, permit))
854}
855
856fn archive_response(repo: Repo, body: Vec<u8>, permit: SlotPermit) -> Response {
857 let (tx, rx) = mpsc::channel::<Result<Bytes, io::Error>>(16);
858 tokio::task::spawn_blocking(move || {
859 let _permit = permit;
860 let mut sink = |chunk: &[u8]| -> io::Result<()> {
861 tx.blocking_send(Ok(Bytes::copy_from_slice(chunk)))
862 .map_err(|_| io::Error::other("client disconnected"))
863 };
864 if let Err(error) = upload_archive_streamed(&repo, &body, &mut sink) {
865 let _ = tx.blocking_send(Err(io::Error::other(error.to_string())));
866 }
867 });
868 nocache(Response::builder().header(
869 header::CONTENT_TYPE,
870 "application/x-git-upload-archive-result",
871 ))
872 .body(Body::from_stream(ReceiverStream::new(rx)))
873 .expect("valid response")
874}
875
876async fn upload_archive_named(
877 State(state): State<PackState>,
878 Path((did, name)): Path<(String, String)>,
879 headers: HeaderMap,
880 body: Bytes,
881) -> Result<Response, PackError> {
882 let owner = resolve_owner(&state, &did).await?;
883 let repo = open_named(&state, &owner, &name)?;
884 archive_dispatch(&state, repo, decode_request(&headers, body)?).await
885}
886
887async fn upload_archive_did(
888 State(state): State<PackState>,
889 Path(did): Path<String>,
890 headers: HeaderMap,
891 body: Bytes,
892) -> Result<Response, PackError> {
893 let repo = open_did(&state, &did)?;
894 archive_dispatch(&state, repo, decode_request(&headers, body)?).await
895}
896
897fn nocache(builder: axum::http::response::Builder) -> axum::http::response::Builder {
898 builder
899 .header(header::EXPIRES, "Fri, 01 Jan 1980 00:00:00 GMT")
900 .header(header::PRAGMA, "no-cache")
901 .header(
902 header::CACHE_CONTROL,
903 "no-cache, max-age=0, must-revalidate",
904 )
905}
906
907fn git_response(content_type: &'static str, body: Vec<u8>) -> Response {
908 nocache(Response::builder().header(header::CONTENT_TYPE, content_type))
909 .body(Body::from(body))
910 .expect("valid response")
911}