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