This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knot2 / crates / knot-pack / src / lib.rs
27 kB 891 lines
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}