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, ClonePath, Handle, KnotHostname, OwnerDid, OwnerRef, ParseError, RepoDid,
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 limit: knot_git::ArchiveLimit,
223 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>,
224) -> Result<(), PackError> {
225 archive::stream(repo, request, limit, sink)
226}
227
228pub fn upload_archive(
229 repo: &Repo,
230 request: &[u8],
231 limit: knot_git::ArchiveLimit,
232) -> Result<Vec<u8>, PackError> {
233 let mut buf = Vec::new();
234 upload_archive_streamed(repo, request, limit, &mut |chunk| {
235 buf.extend_from_slice(chunk);
236 Ok(())
237 })?;
238 Ok(buf)
239}
240
241pub fn meter_pack(
242 pack: &[u8],
243 limits: &PackLimits,
244 kind: gix::hash::Kind,
245) -> Result<(), PackError> {
246 meter::meter(pack, limits, kind)
247}
248
249pub fn ingest_pack(
250 objects_dir: &std::path::Path,
251 pack: &[u8],
252 limits: &PackLimits,
253 kind: gix::hash::Kind,
254) -> Result<(), PackError> {
255 objects::index_pack(objects_dir, pack, limits, kind)
256}
257
258#[doc(hidden)]
259pub mod fuzz {
260 pub fn pkt(data: &[u8]) {
261 let _ = crate::pkt::data_payloads(data);
262 let _ = crate::pkt::data_payloads_all(data);
263 let _ = crate::pkt::split_receive(data);
264 }
265
266 pub fn pack(data: &[u8]) {
267 let _ = crate::meter::meter(data, &crate::PackLimits::default(), gix::hash::Kind::Sha1);
268 }
269
270 pub fn receive_commands(data: &[u8]) {
271 crate::receive::fuzz(data);
272 }
273
274 pub fn upload_args(data: &[u8]) {
275 crate::upload::fuzz(data);
276 }
277}
278
279const MAX_REQUEST_BYTES: usize = 16 * 1024 * 1024;
280
281#[derive(Debug, Clone, PartialEq, Eq)]
282pub enum RepoTarget {
283 Did(RepoDid),
284 OwnerPath(OwnerDid, ClonePath),
285}
286
287#[derive(Debug, Clone, PartialEq, Eq)]
288pub enum RepoLookup {
289 Hosted(RepoDid),
290 Unhosted,
291 Unavailable,
292}
293
294impl RepoLookup {
295 pub fn from_resolved<T>(
296 resolved: knot_index::Resolved<Option<T>>,
297 found: impl FnOnce(T) -> RepoDid,
298 ) -> RepoLookup {
299 match resolved {
300 knot_index::Resolved::Ready(Some(value)) => RepoLookup::Hosted(found(value)),
301 knot_index::Resolved::Ready(None) => RepoLookup::Unhosted,
302 knot_index::Resolved::Warming => RepoLookup::Unavailable,
303 }
304 }
305}
306
307pub trait RepoResolver: Send + Sync + 'static {
308 fn resolve(&self, target: &RepoTarget) -> RepoLookup;
309}
310
311pub trait HandleResolver: Send + Sync + 'static {
312 fn resolve(
313 &self,
314 handle: Handle,
315 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Option<AccountDid>> + Send + '_>>;
316}
317
318pub use knot_edge::SocketPeer;
319
320pub trait ReceiveAdvertiser: Send + Sync + 'static {
321 fn advertise(
322 &self,
323 repo: RepoDid,
324 peer: knot_edge::SocketPeer,
325 headers: HeaderMap,
326 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Response> + Send + '_>>;
327}
328
329impl<F> RepoResolver for F
330where
331 F: Fn(&RepoTarget) -> RepoLookup + Send + Sync + 'static,
332{
333 fn resolve(&self, target: &RepoTarget) -> RepoLookup {
334 self(target)
335 }
336}
337
338#[derive(Clone)]
339struct PackState {
340 layout: Layout,
341 resolver: Arc<dyn RepoResolver>,
342 receive: Option<Arc<dyn ReceiveAdvertiser>>,
343 handle_resolver: Option<Arc<dyn HandleResolver>>,
344 pack_slots: PackSlots,
345 cache: Arc<cache::PackCache>,
346 catalog: Arc<Catalog>,
347 hostname: KnotHostname,
348 archive_limit: knot_git::ArchiveLimit,
349}
350
351pub struct EdgeConfig {
352 pub layout: Layout,
353 pub resolver: Arc<dyn RepoResolver>,
354 pub receive: Option<Arc<dyn ReceiveAdvertiser>>,
355 pub handle_resolver: Option<Arc<dyn HandleResolver>>,
356 pub pack_slots: PackSlots,
357 pub cache: CacheConfig,
358 pub catalog: Arc<Catalog>,
359 pub hostname: KnotHostname,
360 pub clock: Arc<dyn Clock>,
361 pub archive_limit: knot_git::ArchiveLimit,
362}
363
364impl EdgeConfig {
365 pub fn serving(layout: Layout, resolver: Arc<dyn RepoResolver>, clock: Arc<dyn Clock>) -> Self {
366 Self {
367 layout,
368 resolver,
369 receive: None,
370 handle_resolver: None,
371 pack_slots: PackSlots::new(knot_resource::threads().get()),
372 cache: CacheConfig::default(),
373 catalog: Arc::new(Catalog::defaults()),
374 hostname: default_hostname().clone(),
375 clock,
376 archive_limit: knot_git::ArchiveLimit::default(),
377 }
378 }
379
380 pub fn with_pack_slots(self, pack_slots: PackSlots) -> Self {
381 Self { pack_slots, ..self }
382 }
383}
384
385pub fn router(layout: Layout, resolver: Arc<dyn RepoResolver>, clock: Arc<dyn Clock>) -> Router {
386 serving_router(EdgeConfig::serving(layout, resolver, clock))
387}
388
389pub fn router_with_pack_slots(
390 layout: Layout,
391 resolver: Arc<dyn RepoResolver>,
392 pack_slots: PackSlots,
393 clock: Arc<dyn Clock>,
394) -> Router {
395 serving_router(EdgeConfig::serving(layout, resolver, clock).with_pack_slots(pack_slots))
396}
397
398fn serving_router(config: EdgeConfig) -> Router {
399 let state = pack_state(config);
400 write_routes(state.clone()).merge(advertisement_routes(state).into_router())
401}
402
403pub fn edge_routes(config: EdgeConfig) -> (Router, knot_edge::ZeroRttRoutes) {
404 let state = pack_state(config);
405 (write_routes(state.clone()), advertisement_routes(state))
406}
407
408fn pack_state(config: EdgeConfig) -> PackState {
409 let EdgeConfig {
410 layout,
411 resolver,
412 receive,
413 handle_resolver,
414 pack_slots,
415 cache,
416 catalog,
417 hostname,
418 clock,
419 archive_limit,
420 } = config;
421 PackState {
422 layout,
423 resolver,
424 receive,
425 handle_resolver,
426 pack_slots,
427 cache: cache::PackCache::new(cache, clock),
428 catalog,
429 hostname,
430 archive_limit,
431 }
432}
433
434fn write_routes(state: PackState) -> Router {
435 Router::new()
436 .route("/{did}/{name}/git-upload-pack", post(upload_named))
437 .route(
438 "/{did}/{name}/git-upload-archive",
439 post(upload_archive_named),
440 )
441 .route("/{did}/git-upload-pack", post(upload_did))
442 .route("/{did}/git-upload-archive", post(upload_archive_did))
443 .layer(DefaultBodyLimit::max(MAX_REQUEST_BYTES))
444 .with_state(state)
445}
446
447fn advertisement_routes(state: PackState) -> knot_edge::ZeroRttRoutes {
448 let named_state = state.clone();
449 let did_state = state;
450 knot_edge::ZeroRttRoutes::new()
451 .get(
452 "/{did}/{name}/info/refs",
453 knot_edge::ZeroRttSafe::new(
454 move |Path((did, name)): Path<(String, String)>,
455 Query(query): Query<HashMap<String, String>>,
456 peer: knot_edge::SocketPeer,
457 headers: HeaderMap| {
458 let state = named_state.clone();
459 async move {
460 let owner = resolve_owner(&state, &did).await?;
461 let repo_did = resolve_named_did(&state, &owner, &name)?;
462 read_advertisement(
463 &state,
464 repo_did,
465 query.get("service").map(String::as_str),
466 peer,
467 &headers,
468 )
469 .await
470 }
471 },
472 ),
473 )
474 .get(
475 "/{did}/info/refs",
476 knot_edge::ZeroRttSafe::new(
477 move |Path(did): Path<String>,
478 Query(query): Query<HashMap<String, String>>,
479 peer: knot_edge::SocketPeer,
480 headers: HeaderMap| {
481 let state = did_state.clone();
482 async move {
483 let repo_did = resolve_did_did(&state, &did)?;
484 read_advertisement(
485 &state,
486 repo_did,
487 query.get("service").map(String::as_str),
488 peer,
489 &headers,
490 )
491 .await
492 }
493 },
494 ),
495 )
496}
497
498fn lookup_did(lookup: RepoLookup) -> Result<RepoDid, PackError> {
499 match lookup {
500 RepoLookup::Hosted(did) => Ok(did),
501 RepoLookup::Unhosted => Err(PackError::NotFound),
502 RepoLookup::Unavailable => Err(PackError::Unavailable),
503 }
504}
505
506fn bad_path(error: ParseError) -> PackError {
507 PackError::BadPath(error.to_string())
508}
509
510async fn resolve_owner(state: &PackState, owner: &str) -> Result<OwnerDid, PackError> {
511 match OwnerRef::parse(owner).ok_or(PackError::NotFound)? {
512 OwnerRef::Did(did) => Ok(did),
513 OwnerRef::Handle(handle) => {
514 let resolver = state.handle_resolver.as_ref().ok_or(PackError::NotFound)?;
515 let did = resolver.resolve(handle).await.ok_or(PackError::NotFound)?;
516 Ok(did.into())
517 }
518 }
519}
520
521fn resolve_named_did(
522 state: &PackState,
523 owner: &OwnerDid,
524 name: &str,
525) -> Result<RepoDid, PackError> {
526 let path = ClonePath::parse(name).ok_or(PackError::NotFound)?;
527 lookup_did(
528 state
529 .resolver
530 .resolve(&RepoTarget::OwnerPath(owner.clone(), path)),
531 )
532}
533
534fn resolve_did_did(state: &PackState, did: &str) -> Result<RepoDid, PackError> {
535 let did = RepoDid::new(did).map_err(bad_path)?;
536 lookup_did(state.resolver.resolve(&RepoTarget::Did(did)))
537}
538
539fn open_named(state: &PackState, owner: &OwnerDid, name: &str) -> Result<Repo, PackError> {
540 let did = resolve_named_did(state, owner, name)?;
541 state.layout.open(&did).map_err(|_| PackError::NotFound)
542}
543
544fn open_did(state: &PackState, did: &str) -> Result<Repo, PackError> {
545 let did = resolve_did_did(state, did)?;
546 state.layout.open(&did).map_err(|_| PackError::NotFound)
547}
548
549async fn read_advertisement(
550 state: &PackState,
551 repo_did: RepoDid,
552 service: Option<&str>,
553 peer: knot_edge::SocketPeer,
554 headers: &HeaderMap,
555) -> Result<Response, PackError> {
556 match service {
557 Some("git-upload-pack") => {
558 let repo = state
559 .layout
560 .open(&repo_did)
561 .map_err(|_| PackError::NotFound)?;
562 let body = if wants_v2(headers) {
563 advertise_upload(&repo)?
564 } else {
565 upload::advertise_v0(&repo)?
566 };
567 Ok(git_response(
568 "application/x-git-upload-pack-advertisement",
569 body,
570 ))
571 }
572 Some("git-receive-pack") => match &state.receive {
573 Some(advertiser) => Ok(advertiser.advertise(repo_did, peer, headers.clone()).await),
574 None => Err(PackError::PushOverSsh),
575 },
576 _ => Err(PackError::UnsupportedService),
577 }
578}
579
580fn decode_request(headers: &HeaderMap, body: Bytes) -> Result<Vec<u8>, PackError> {
581 let codings: Vec<&str> = headers
582 .get(header::CONTENT_ENCODING)
583 .and_then(|value| value.to_str().ok())
584 .map(|value| {
585 value
586 .split(',')
587 .map(str::trim)
588 .filter(|token| !token.is_empty() && !token.eq_ignore_ascii_case("identity"))
589 .collect()
590 })
591 .unwrap_or_default();
592 match codings.as_slice() {
593 [] => Ok(body.to_vec()),
594 [token] if token.eq_ignore_ascii_case("gzip") => {
595 let mut out = Vec::new();
596 flate2::read::GzDecoder::new(body.as_ref())
597 .take(MAX_REQUEST_BYTES as u64 + 1)
598 .read_to_end(&mut out)
599 .map_err(|error| PackError::Protocol(format!("gzip request body: {error}")))?;
600 match out.len() > MAX_REQUEST_BYTES {
601 true => Err(PackError::LimitExceeded(PackLimit::TotalBytes)),
602 false => Ok(out),
603 }
604 }
605 unsupported => Err(PackError::UnsupportedEncoding(unsupported.join(", "))),
606 }
607}
608
609fn wants_v2(headers: &HeaderMap) -> bool {
610 headers
611 .get("git-protocol")
612 .and_then(|value| value.to_str().ok())
613 .is_some_and(|value| value.split(':').any(|token| token.trim() == "version=2"))
614}
615
616fn refs_digest(repo: &Repo) -> Option<ids::RefsDigest> {
617 let refs = repo
618 .advertised_refs_for(knot_git::AdvertScope::Upload)
619 .ok()?;
620 let mut hasher = gix::hash::hasher(gix::hash::Kind::Sha256);
621 refs.iter().for_each(|record| {
622 hasher.update(record.name.as_str().as_bytes());
623 hasher.update(b"\0");
624 hasher.update(record.target.to_string().as_bytes());
625 hasher.update(b"\n");
626 });
627 let id = hasher.try_finalize().ok()?;
628 let mut token = [0u8; 32];
629 token.copy_from_slice(id.as_slice());
630 Some(ids::RefsDigest::new(token))
631}
632
633const CACHE_RETRY_BUDGET: usize = 4;
634
635async fn upload_dispatch(
636 state: &PackState,
637 repo: Repo,
638 body: &[u8],
639) -> Result<Response, PackError> {
640 upload_dispatch_within(state, repo, body, CACHE_RETRY_BUDGET).await
641}
642
643async fn upload_dispatch_within(
644 state: &PackState,
645 repo: Repo,
646 body: &[u8],
647 retries: usize,
648) -> Result<Response, PackError> {
649 let Some(token) = refs_digest(&repo) else {
650 return dispatch_plan(state, repo, body, None).await;
651 };
652 let key = cache::RequestKey::new(&repo.objects_dir(), &token, body);
653 match state.cache.decide(key) {
654 cache::Decision::Serve(bytes) => Ok(cached_response(bytes)),
655 cache::Decision::Await(receiver) => match cache::wait(receiver).await {
656 cache::Resolved::Bytes(bytes) => Ok(cached_response(bytes)),
657 cache::Resolved::Retry => match retries {
658 0 => dispatch_plan(state, repo, body, None).await,
659 _ => Box::pin(upload_dispatch_within(state, repo, body, retries - 1)).await,
660 },
661 cache::Resolved::Regenerate => dispatch_plan(state, repo, body, None).await,
662 },
663 cache::Decision::Lead(lease) => dispatch_plan(state, repo, body, Some(lease)).await,
664 cache::Decision::Stream | cache::Decision::Off => {
665 dispatch_plan(state, repo, body, None).await
666 }
667 }
668}
669
670async fn dispatch_plan(
671 state: &PackState,
672 repo: Repo,
673 body: &[u8],
674 lease: Option<cache::Lease>,
675) -> Result<Response, PackError> {
676 let permit = state.pack_slots.acquire().await;
677 let owned_body = body.to_vec();
678 let planned = tokio::task::spawn_blocking(move || {
679 let outcome = upload::plan(&repo, &owned_body);
680 (repo, outcome)
681 })
682 .await;
683 let (repo, plan) = match planned {
684 Ok((repo, Ok(plan))) => (repo, plan),
685 Ok((_repo, Err(error))) => {
686 if let Some(lease) = lease {
687 lease.regenerate();
688 }
689 return Err(error);
690 }
691 Err(_join) => {
692 if let Some(lease) = lease {
693 lease.regenerate();
694 }
695 return Err(PackError::Pack("upload planning task panicked".to_string()));
696 }
697 };
698 match plan {
699 UploadOutcome::Buffered(bytes) => {
700 drop(permit);
701 if let Some(lease) = lease {
702 lease.regenerate();
703 }
704 Ok(git_response("application/x-git-upload-pack-result", bytes))
705 }
706 UploadOutcome::Streaming {
707 preamble,
708 wants,
709 haves,
710 opts,
711 } => Ok(stream_response(
712 repo,
713 preamble,
714 wants,
715 haves,
716 opts,
717 permit,
718 lease,
719 Arc::clone(&state.catalog),
720 state.hostname.clone(),
721 )),
722 }
723}
724
725fn cached_response(bytes: Bytes) -> Response {
726 nocache(
727 Response::builder().header(header::CONTENT_TYPE, "application/x-git-upload-pack-result"),
728 )
729 .body(Body::from(bytes))
730 .expect("valid response")
731}
732
733#[allow(clippy::too_many_arguments)]
734fn stream_response(
735 repo: Repo,
736 preamble: Vec<u8>,
737 wants: WantOids,
738 haves: HaveOids,
739 opts: upload::StreamOpts,
740 permit: SlotPermit,
741 lease: Option<cache::Lease>,
742 catalog: Arc<Catalog>,
743 hostname: KnotHostname,
744) -> Response {
745 let side_band = opts.side_band;
746 let limit = lease.as_ref().map(cache::Lease::max_entry_bytes);
747 let (tx, rx) = mpsc::channel::<Result<Bytes, io::Error>>(16);
748 tokio::task::spawn_blocking(move || {
749 let _permit = permit;
750 let capture = std::cell::RefCell::new(cache::Capture::new(limit));
751 let mut emit = |chunk: &[u8]| -> io::Result<()> {
752 capture.borrow_mut().record(chunk);
753 tx.blocking_send(Ok(Bytes::copy_from_slice(chunk)))
754 .map_err(|_| io::Error::other("client disconnected"))
755 };
756 if emit(&preamble).is_err() {
757 if let Some(lease) = lease {
758 lease.retry();
759 }
760 return;
761 }
762 let result = upload::stream_pack(
763 &repo,
764 &wants,
765 &haves,
766 &opts,
767 &catalog.fetch,
768 &hostname,
769 &mut emit,
770 );
771 match result {
772 Ok(()) => {
773 if side_band {
774 let mut flush = Vec::new();
775 if pkt::write_flush(&mut flush).is_ok() {
776 let _ = emit(&flush);
777 }
778 }
779 if let Some(lease) = lease {
780 match capture.into_inner().into_bytes() {
781 Some(bytes) => lease.ready(Bytes::from(bytes)),
782 None => lease.too_large(),
783 }
784 }
785 }
786 Err(error) if side_band => {
787 let mut tail = Vec::new();
788 let line = catalog
789 .fetch
790 .fatal
791 .line(|ErrorKey::Error| error.to_string().replace('\n', " "));
792 let message = format!("{line}\n");
793 if pkt::write_band_error(&mut tail, message.as_bytes()).is_ok() {
794 let _ = pkt::write_flush(&mut tail);
795 let _ = emit(&tail);
796 }
797 if let Some(lease) = lease {
798 lease.regenerate();
799 }
800 }
801 Err(error) => {
802 let _ = tx.blocking_send(Err(io::Error::other(error.to_string())));
803 if let Some(lease) = lease {
804 lease.regenerate();
805 }
806 }
807 }
808 });
809 nocache(
810 Response::builder().header(header::CONTENT_TYPE, "application/x-git-upload-pack-result"),
811 )
812 .body(Body::from_stream(ReceiverStream::new(rx)))
813 .expect("valid response")
814}
815
816async fn upload_named(
817 State(state): State<PackState>,
818 Path((did, name)): Path<(String, String)>,
819 headers: HeaderMap,
820 body: Bytes,
821) -> Result<Response, PackError> {
822 let owner = resolve_owner(&state, &did).await?;
823 let repo = open_named(&state, &owner, &name)?;
824 upload_dispatch(&state, repo, &decode_request(&headers, body)?).await
825}
826
827async fn upload_did(
828 State(state): State<PackState>,
829 Path(did): Path<String>,
830 headers: HeaderMap,
831 body: Bytes,
832) -> Result<Response, PackError> {
833 let repo = open_did(&state, &did)?;
834 upload_dispatch(&state, repo, &decode_request(&headers, body)?).await
835}
836
837async fn archive_dispatch(
838 state: &PackState,
839 repo: Repo,
840 body: Vec<u8>,
841) -> Result<Response, PackError> {
842 let permit = state.pack_slots.acquire().await;
843 Ok(archive_response(repo, body, state.archive_limit, permit))
844}
845
846fn archive_response(
847 repo: Repo,
848 body: Vec<u8>,
849 limit: knot_git::ArchiveLimit,
850 permit: SlotPermit,
851) -> Response {
852 let (tx, rx) = mpsc::channel::<Result<Bytes, io::Error>>(16);
853 tokio::task::spawn_blocking(move || {
854 let _permit = permit;
855 let mut sink = |chunk: &[u8]| -> io::Result<()> {
856 tx.blocking_send(Ok(Bytes::copy_from_slice(chunk)))
857 .map_err(|_| io::Error::other("client disconnected"))
858 };
859 if let Err(error) = upload_archive_streamed(&repo, &body, limit, &mut sink) {
860 let _ = tx.blocking_send(Err(io::Error::other(error.to_string())));
861 }
862 });
863 nocache(Response::builder().header(
864 header::CONTENT_TYPE,
865 "application/x-git-upload-archive-result",
866 ))
867 .body(Body::from_stream(ReceiverStream::new(rx)))
868 .expect("valid response")
869}
870
871async fn upload_archive_named(
872 State(state): State<PackState>,
873 Path((did, name)): Path<(String, String)>,
874 headers: HeaderMap,
875 body: Bytes,
876) -> Result<Response, PackError> {
877 let owner = resolve_owner(&state, &did).await?;
878 let repo = open_named(&state, &owner, &name)?;
879 archive_dispatch(&state, repo, decode_request(&headers, body)?).await
880}
881
882async fn upload_archive_did(
883 State(state): State<PackState>,
884 Path(did): Path<String>,
885 headers: HeaderMap,
886 body: Bytes,
887) -> Result<Response, PackError> {
888 let repo = open_did(&state, &did)?;
889 archive_dispatch(&state, repo, decode_request(&headers, body)?).await
890}
891
892fn nocache(builder: axum::http::response::Builder) -> axum::http::response::Builder {
893 builder
894 .header(header::EXPIRES, "Fri, 01 Jan 1980 00:00:00 GMT")
895 .header(header::PRAGMA, "no-cache")
896 .header(
897 header::CACHE_CONTROL,
898 "no-cache, max-age=0, must-revalidate",
899 )
900}
901
902fn git_response(content_type: &'static str, body: Vec<u8>) -> Response {
903 nocache(Response::builder().header(header::CONTENT_TYPE, content_type))
904 .body(Body::from(body))
905 .expect("valid response")
906}