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 906 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 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}