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
28 kB 911 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, 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}