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 / upload.rs
31 kB 979 lines
1use std::collections::HashSet; 2use std::io::{self, Read, Write}; 3use std::sync::OnceLock; 4use std::time::Duration; 5 6use knot_git::{ 7 CommitDepth, Deepen, Filter, Haves, PackBudget, PackfileUri, Repo, ShallowCommits, Wants, 8}; 9use knot_messages::{CountKey, FetchMessages, KnotKey}; 10use knot_types::{KnotHostname, ObjectCount, ObjectFormat, Oid, UnixSeconds}; 11 12use crate::error::PackError; 13use crate::objects; 14use crate::pkt; 15use crate::{HaveOids, WantOids}; 16 17const AGENT: &[u8] = b"agent=knot/0\n"; 18const SELECTION_MAX_OBJECTS: ObjectCount = ObjectCount::new(16_000_000); 19const SELECTION_TIME_BUDGET: Duration = Duration::from_secs(600); 20 21#[derive(Debug, Clone, Copy)] 22pub struct SelectionLimits { 23 pub max_objects: ObjectCount, 24 pub time_budget: Duration, 25} 26 27impl Default for SelectionLimits { 28 fn default() -> Self { 29 Self { 30 max_objects: SELECTION_MAX_OBJECTS, 31 time_budget: SELECTION_TIME_BUDGET, 32 } 33 } 34} 35 36static SELECTION_LIMITS: OnceLock<SelectionLimits> = OnceLock::new(); 37 38pub fn init_selection_limits(limits: SelectionLimits) { 39 if SELECTION_LIMITS.set(limits).is_err() { 40 debug_assert!(false, "selection limits initialized more than once"); 41 } 42} 43 44fn selection_limits() -> SelectionLimits { 45 SELECTION_LIMITS.get().copied().unwrap_or_default() 46} 47const V0_CAPS_BASE: &str = "multi_ack_detailed no-done side-band-64k ofs-delta shallow deepen-since deepen-not filter agent=knot/0"; 48 49fn v0_caps(format: ObjectFormat) -> String { 50 format!("{V0_CAPS_BASE} object-format={}", format.capability()) 51} 52 53pub struct StreamOpts { 54 pub side_band: bool, 55 pub sideband_all: bool, 56 pub no_progress: bool, 57 pub thin: bool, 58 pub filter: Filter, 59 pub shallow_commits: Option<Vec<Oid>>, 60 pub packfile_uris: Vec<PackfileUri>, 61 pub emit_packfile_header: bool, 62} 63 64pub enum UploadOutcome { 65 Buffered(Vec<u8>), 66 Streaming { 67 preamble: Vec<u8>, 68 wants: WantOids, 69 haves: HaveOids, 70 opts: StreamOpts, 71 }, 72} 73 74fn write_v2_caps(buf: &mut Vec<u8>, format: ObjectFormat) -> Result<(), PackError> { 75 pkt::write_data(buf, b"version 2\n")?; 76 pkt::write_data(buf, AGENT)?; 77 pkt::write_data(buf, b"ls-refs\n")?; 78 pkt::write_data( 79 buf, 80 b"fetch=shallow filter wait-for-done packfile-uris sideband-all\n", 81 )?; 82 pkt::write_data(buf, b"server-option\n")?; 83 pkt::write_data( 84 buf, 85 format!("object-format={}\n", format.capability()).as_bytes(), 86 )?; 87 pkt::write_flush(buf)?; 88 Ok(()) 89} 90 91pub fn advertise(repo: &Repo) -> Result<Vec<u8>, PackError> { 92 let mut buf = Vec::new(); 93 pkt::write_data(&mut buf, b"# service=git-upload-pack\n")?; 94 pkt::write_flush(&mut buf)?; 95 write_v2_caps(&mut buf, repo.object_format())?; 96 Ok(buf) 97} 98 99pub fn advertise_ssh(repo: &Repo) -> Result<Vec<u8>, PackError> { 100 let mut buf = Vec::new(); 101 write_v2_caps(&mut buf, repo.object_format())?; 102 Ok(buf) 103} 104 105pub fn advertise_v0(repo: &Repo) -> Result<Vec<u8>, PackError> { 106 let mut buf = Vec::new(); 107 pkt::write_data(&mut buf, b"# service=git-upload-pack\n")?; 108 pkt::write_flush(&mut buf)?; 109 write_v0_advert(&mut buf, repo)?; 110 Ok(buf) 111} 112 113pub fn advertise_v0_ssh(repo: &Repo) -> Result<Vec<u8>, PackError> { 114 let mut buf = Vec::new(); 115 write_v0_advert(&mut buf, repo)?; 116 Ok(buf) 117} 118 119fn write_v0_advert(buf: &mut Vec<u8>, repo: &Repo) -> Result<(), PackError> { 120 let format = repo.object_format(); 121 let caps = v0_caps(format); 122 let refs = repo.advertised_refs_for(knot_git::AdvertScope::Upload)?; 123 match repo.head() { 124 Some(head) => { 125 pkt::write_data( 126 buf, 127 format!("{} HEAD\0{caps} symref=HEAD:{}\n", head.target, head.name).as_bytes(), 128 )?; 129 write_plain_refs(buf, &refs)?; 130 } 131 None => match refs.split_first() { 132 Some((first, rest)) => { 133 pkt::write_data( 134 buf, 135 format!("{} {}\0{caps}\n", first.target, first.name).as_bytes(), 136 )?; 137 write_plain_refs(buf, rest)?; 138 } 139 None => { 140 pkt::write_data( 141 buf, 142 format!("{} capabilities^{{}}\0{caps}\n", format.null_oid()).as_bytes(), 143 )?; 144 } 145 }, 146 } 147 pkt::write_flush(buf)?; 148 Ok(()) 149} 150 151fn write_plain_refs(buf: &mut Vec<u8>, refs: &[knot_git::RefRecord]) -> Result<(), PackError> { 152 refs.iter().try_for_each(|record| { 153 pkt::write_data( 154 buf, 155 format!("{} {}\n", record.target, record.name).as_bytes(), 156 ) 157 .map_err(PackError::from) 158 }) 159} 160 161pub(crate) fn fuzz(body: &[u8]) { 162 if let Ok(lines) = pkt::data_payloads_all(body) { 163 let _ = parse_wants(&lines); 164 let _ = parse_oids(&lines, b"want "); 165 let _ = parse_oids(&lines, b"have "); 166 let _ = first_caps(&lines); 167 let _ = parse_ls_refs_args(&lines); 168 } 169} 170 171pub fn plan(repo: &Repo, body: &[u8]) -> Result<UploadOutcome, PackError> { 172 let peek = pkt::data_payloads(body)?; 173 match peek.first() { 174 Some(line) if line.starts_with(b"command=") => plan_v2(repo, &peek), 175 _ => plan_v0(repo, body), 176 } 177} 178 179pub fn buffered( 180 repo: &Repo, 181 body: &[u8], 182 messages: &FetchMessages, 183 knot: &KnotHostname, 184) -> Result<Vec<u8>, PackError> { 185 let mut out = Vec::new(); 186 streamed(repo, body, messages, knot, &mut |chunk| { 187 out.extend_from_slice(chunk); 188 Ok(()) 189 })?; 190 Ok(out) 191} 192 193pub fn streamed( 194 repo: &Repo, 195 body: &[u8], 196 messages: &FetchMessages, 197 knot: &KnotHostname, 198 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>, 199) -> Result<(), PackError> { 200 match plan(repo, body)? { 201 UploadOutcome::Buffered(bytes) => sink(&bytes).map_err(PackError::from), 202 UploadOutcome::Streaming { 203 preamble, 204 wants, 205 haves, 206 opts, 207 } => { 208 sink(&preamble)?; 209 stream_pack(repo, &wants, &haves, &opts, messages, knot, sink)?; 210 if opts.side_band { 211 let mut flush = Vec::new(); 212 pkt::write_flush(&mut flush)?; 213 sink(&flush)?; 214 } 215 Ok(()) 216 } 217 } 218} 219 220fn plan_v2(repo: &Repo, lines: &[&[u8]]) -> Result<UploadOutcome, PackError> { 221 match lines.first().copied().unwrap_or_default() { 222 command if command.starts_with(b"command=ls-refs") => { 223 Ok(UploadOutcome::Buffered(ls_refs(repo, lines)?)) 224 } 225 command if command.starts_with(b"command=fetch") => plan_v2_fetch(repo, lines), 226 _ => Err(PackError::Protocol( 227 "unsupported protocol v2 command".to_string(), 228 )), 229 } 230} 231 232struct LsRefsArgs { 233 symrefs: bool, 234 peel: bool, 235 prefixes: Vec<String>, 236} 237 238fn parse_ls_refs_args(lines: &[&[u8]]) -> LsRefsArgs { 239 LsRefsArgs { 240 symrefs: lines.iter().any(|line| line.starts_with(b"symrefs")), 241 peel: lines.iter().any(|line| line.starts_with(b"peel")), 242 prefixes: lines 243 .iter() 244 .filter_map(|line| { 245 std::str::from_utf8(line) 246 .ok()? 247 .trim_end() 248 .strip_prefix("ref-prefix ") 249 .map(str::to_string) 250 }) 251 .collect(), 252 } 253} 254 255pub(crate) fn matches_prefix<S: AsRef<str>>(name: &str, prefixes: &[S]) -> bool { 256 prefixes.is_empty() 257 || prefixes 258 .iter() 259 .any(|prefix| name.starts_with(prefix.as_ref())) 260} 261 262fn ls_refs(repo: &Repo, lines: &[&[u8]]) -> Result<Vec<u8>, PackError> { 263 let args = parse_ls_refs_args(lines); 264 let mut buf = Vec::new(); 265 if let Some(head) = repo.head() 266 && matches_prefix("HEAD", &args.prefixes) 267 { 268 let mut line = format!("{} HEAD", head.target); 269 if args.symrefs { 270 line.push_str(&format!(" symref-target:{}", head.name)); 271 } 272 line.push('\n'); 273 pkt::write_data(&mut buf, line.as_bytes())?; 274 } 275 repo.advertised_refs_for(knot_git::AdvertScope::Upload)? 276 .iter() 277 .filter(|record| matches_prefix(record.name.as_str(), &args.prefixes)) 278 .try_fold(&mut buf, |buf, record| { 279 let mut line = format!("{} {}", record.target, record.name); 280 if args.peel 281 && let Some(peeled) = repo.peeled_target(record.target)? 282 { 283 line.push_str(&format!(" peeled:{peeled}")); 284 } 285 line.push('\n'); 286 pkt::write_data(buf, line.as_bytes())?; 287 Ok::<_, PackError>(buf) 288 })?; 289 pkt::write_flush(&mut buf)?; 290 Ok(buf) 291} 292 293fn plan_v2_fetch(repo: &Repo, lines: &[&[u8]]) -> Result<UploadOutcome, PackError> { 294 let wants = parse_wants(lines)?; 295 ensure_wanted(repo, &wants)?; 296 let haves = HaveOids::new(parse_oids(lines, b"have ")); 297 let done = lines.iter().any(|line| line.starts_with(b"done")); 298 let wait_for_done = lines.iter().any(|line| line.starts_with(b"wait-for-done")); 299 let sideband_all = lines.iter().any(|line| line.starts_with(b"sideband-all")); 300 let no_progress = lines.iter().any(|line| line.starts_with(b"no-progress")); 301 let thin = lines.iter().any(|line| line.starts_with(b"thin-pack")); 302 let filter = parse_filter(lines)?; 303 let deepen = parse_deepen(repo, lines)?; 304 let client_shallow = parse_oids(lines, b"shallow "); 305 let common: HaveOids = haves 306 .iter() 307 .copied() 308 .filter(|oid| repo.contains(*oid)) 309 .collect(); 310 311 let mut preamble = Vec::new(); 312 if !haves.is_empty() && !done { 313 seg(&mut preamble, sideband_all, b"acknowledgments\n")?; 314 if common.is_empty() { 315 seg(&mut preamble, sideband_all, b"NAK\n")?; 316 pkt::write_flush(&mut preamble)?; 317 return Ok(UploadOutcome::Buffered(preamble)); 318 } 319 common.iter().try_for_each(|oid| { 320 seg( 321 &mut preamble, 322 sideband_all, 323 format!("ACK {oid}\n").as_bytes(), 324 ) 325 })?; 326 if wait_for_done { 327 pkt::write_flush(&mut preamble)?; 328 return Ok(UploadOutcome::Buffered(preamble)); 329 } 330 seg(&mut preamble, sideband_all, b"ready\n")?; 331 pkt::write_delim(&mut preamble)?; 332 } 333 334 let shallow_commits = if deepen.is_shallow_request() || repo.is_shallow() { 335 let plan = 336 repo.shallow_walk(wants.wants(), &deepen, ShallowCommits::new(&client_shallow))?; 337 seg(&mut preamble, sideband_all, b"shallow-info\n")?; 338 plan.shallow.iter().try_for_each(|oid| { 339 seg( 340 &mut preamble, 341 sideband_all, 342 format!("shallow {oid}\n").as_bytes(), 343 ) 344 })?; 345 plan.unshallow.iter().try_for_each(|oid| { 346 seg( 347 &mut preamble, 348 sideband_all, 349 format!("unshallow {oid}\n").as_bytes(), 350 ) 351 })?; 352 pkt::write_delim(&mut preamble)?; 353 Some(plan.commits) 354 } else { 355 None 356 }; 357 358 Ok(UploadOutcome::Streaming { 359 preamble, 360 wants, 361 haves: common, 362 opts: StreamOpts { 363 side_band: true, 364 sideband_all, 365 no_progress, 366 thin, 367 filter, 368 shallow_commits, 369 packfile_uris: packfile_uri_candidates(repo, lines), 370 emit_packfile_header: true, 371 }, 372 }) 373} 374 375fn seg(buf: &mut Vec<u8>, sideband_all: bool, content: &[u8]) -> io::Result<()> { 376 if sideband_all { 377 pkt::write_band(buf, content) 378 } else { 379 pkt::write_data(buf, content) 380 } 381} 382 383fn packfile_uri_candidates(repo: &Repo, lines: &[&[u8]]) -> Vec<PackfileUri> { 384 let Some(protocols) = lines.iter().find_map(|line| { 385 std::str::from_utf8(line) 386 .ok()? 387 .trim_end() 388 .strip_prefix("packfile-uris ") 389 .map(str::to_string) 390 }) else { 391 return Vec::new(); 392 }; 393 let allowed: Vec<&str> = protocols.split(',').map(str::trim).collect(); 394 repo.blob_packfile_uris() 395 .into_iter() 396 .filter(|candidate| { 397 candidate 398 .uri 399 .as_str() 400 .split_once("://") 401 .is_some_and(|(scheme, _)| allowed.contains(&scheme)) 402 }) 403 .collect() 404} 405 406fn parse_filter(lines: &[&[u8]]) -> Result<Filter, PackError> { 407 let spec = lines.iter().find_map(|line| { 408 std::str::from_utf8(line) 409 .ok()? 410 .trim_end() 411 .strip_prefix("filter ") 412 }); 413 match spec { 414 None => Ok(Filter::None), 415 Some("blob:none") => Ok(Filter::BlobNone), 416 Some(rest) if rest.starts_with("blob:limit=") => parse_size(&rest["blob:limit=".len()..]) 417 .map(Filter::BlobLimit) 418 .ok_or_else(|| PackError::Protocol(format!("bad blob:limit filter: {rest}"))), 419 Some(rest) if rest.starts_with("tree:") => rest["tree:".len()..] 420 .parse::<u32>() 421 .map(|depth| Filter::TreeDepth(knot_git::TreeDepth::new(depth))) 422 .map_err(|_| PackError::Protocol(format!("bad tree filter: {rest}"))), 423 Some(other) => Err(PackError::Protocol(format!("unsupported filter: {other}"))), 424 } 425} 426 427fn parse_size(text: &str) -> Option<u64> { 428 let (digits, scale) = match text.chars().last() { 429 Some('k') | Some('K') => (&text[..text.len() - 1], 1024), 430 Some('m') | Some('M') => (&text[..text.len() - 1], 1024 * 1024), 431 Some('g') | Some('G') => (&text[..text.len() - 1], 1024 * 1024 * 1024), 432 _ => (text, 1), 433 }; 434 digits 435 .parse::<u64>() 436 .ok() 437 .and_then(|value| value.checked_mul(scale)) 438} 439 440fn parse_deepen(repo: &Repo, lines: &[&[u8]]) -> Result<Deepen, PackError> { 441 let value = |prefix: &str| -> Option<&str> { 442 lines.iter().find_map(|line| { 443 std::str::from_utf8(line) 444 .ok()? 445 .trim_end() 446 .strip_prefix(prefix) 447 }) 448 }; 449 let depth = value("deepen ") 450 .and_then(|text| text.parse::<u32>().ok()) 451 .map(CommitDepth::new); 452 let since = value("deepen-since ") 453 .and_then(|text| text.parse::<i64>().ok()) 454 .map(UnixSeconds::new); 455 let relative = lines 456 .iter() 457 .any(|line| line.starts_with(b"deepen-relative")); 458 let not = lines 459 .iter() 460 .filter_map(|line| { 461 std::str::from_utf8(line) 462 .ok()? 463 .trim_end() 464 .strip_prefix("deepen-not ") 465 }) 466 .map(|spec| resolve_commitish(repo, spec)) 467 .collect::<Result<Vec<_>, _>>()?; 468 Ok(Deepen { 469 depth, 470 since, 471 not, 472 relative, 473 }) 474} 475 476fn resolve_commitish(repo: &Repo, spec: &str) -> Result<Oid, PackError> { 477 if let Ok(oid) = Oid::from_hex(spec) 478 && repo.contains(oid) 479 { 480 return Ok(oid); 481 } 482 let candidates = [ 483 spec.to_string(), 484 format!("refs/{spec}"), 485 format!("refs/tags/{spec}"), 486 format!("refs/heads/{spec}"), 487 ]; 488 repo.advertised_refs_for(knot_git::AdvertScope::Upload)? 489 .iter() 490 .find(|record| candidates.iter().any(|name| record.name.as_str() == name)) 491 .map(|record| record.target) 492 .ok_or_else(|| PackError::Protocol(format!("deepen-not {spec}: unknown ref"))) 493} 494 495fn plan_v0(repo: &Repo, body: &[u8]) -> Result<UploadOutcome, PackError> { 496 let lines = pkt::data_payloads_all(body)?; 497 let wants = parse_wants(&lines)?; 498 ensure_wanted(repo, &wants)?; 499 let haves = HaveOids::new(parse_oids(&lines, b"have ")); 500 let done = lines.iter().any(|line| line.starts_with(b"done")); 501 let caps = first_caps(&lines); 502 let side_band = caps 503 .map(|caps| { 504 caps.split(' ') 505 .any(|cap| cap == "side-band-64k" || cap == "side-band") 506 }) 507 .unwrap_or(false); 508 let no_progress = caps 509 .map(|caps| caps.split(' ').any(|cap| cap == "no-progress")) 510 .unwrap_or(false); 511 let thin = caps 512 .map(|caps| caps.split(' ').any(|cap| cap == "thin-pack")) 513 .unwrap_or(false); 514 let multi_ack_detailed = caps 515 .map(|caps| caps.split(' ').any(|cap| cap == "multi_ack_detailed")) 516 .unwrap_or(false); 517 let no_done = caps 518 .map(|caps| caps.split(' ').any(|cap| cap == "no-done")) 519 .unwrap_or(false); 520 let filter = parse_filter(&lines)?; 521 let deepen = parse_deepen(repo, &lines)?; 522 let client_shallow = parse_oids(&lines, b"shallow "); 523 let common: HaveOids = haves 524 .iter() 525 .copied() 526 .filter(|oid| repo.contains(*oid)) 527 .collect(); 528 529 let mut preamble = Vec::new(); 530 let shallow_commits = if deepen.is_shallow_request() || repo.is_shallow() { 531 let plan = 532 repo.shallow_walk(wants.wants(), &deepen, ShallowCommits::new(&client_shallow))?; 533 plan.shallow.iter().try_for_each(|oid| { 534 pkt::write_data(&mut preamble, format!("shallow {oid}\n").as_bytes()) 535 })?; 536 plan.unshallow.iter().try_for_each(|oid| { 537 pkt::write_data(&mut preamble, format!("unshallow {oid}\n").as_bytes()) 538 })?; 539 pkt::write_flush(&mut preamble)?; 540 Some(plan.commits) 541 } else { 542 None 543 }; 544 545 if haves.is_empty() && !done { 546 if !(deepen.is_shallow_request() || repo.is_shallow()) { 547 pkt::write_data(&mut preamble, b"NAK\n")?; 548 } 549 return Ok(UploadOutcome::Buffered(preamble)); 550 } 551 552 let ready = multi_ack_detailed 553 && !done 554 && !common.is_empty() 555 && common.len() == haves.len() 556 && repo.wants_satisfied_by(wants.wants(), common.haves())?; 557 558 let stream = move |preamble: Vec<u8>, common: HaveOids| UploadOutcome::Streaming { 559 preamble, 560 wants, 561 haves: common, 562 opts: StreamOpts { 563 side_band, 564 sideband_all: false, 565 no_progress, 566 thin, 567 filter, 568 shallow_commits, 569 packfile_uris: Vec::new(), 570 emit_packfile_header: false, 571 }, 572 }; 573 574 if multi_ack_detailed { 575 common.iter().try_for_each(|oid| { 576 pkt::write_data(&mut preamble, format!("ACK {oid} common\n").as_bytes()) 577 })?; 578 let last = common.as_slice().last().copied(); 579 if done { 580 match last { 581 Some(oid) => { 582 pkt::write_data(&mut preamble, format!("ACK {oid}\n").as_bytes())?; 583 } 584 None => pkt::write_data(&mut preamble, b"NAK\n")?, 585 } 586 return Ok(stream(preamble, common)); 587 } 588 match (ready, last) { 589 (true, Some(oid)) => { 590 pkt::write_data(&mut preamble, format!("ACK {oid} ready\n").as_bytes())?; 591 pkt::write_data(&mut preamble, b"NAK\n")?; 592 if no_done { 593 pkt::write_data(&mut preamble, format!("ACK {oid}\n").as_bytes())?; 594 return Ok(stream(preamble, common)); 595 } 596 } 597 _ => pkt::write_data(&mut preamble, b"NAK\n")?, 598 } 599 return Ok(UploadOutcome::Buffered(preamble)); 600 } 601 602 match common.as_slice().first() { 603 Some(oid) => pkt::write_data(&mut preamble, format!("ACK {oid}\n").as_bytes())?, 604 None => pkt::write_data(&mut preamble, b"NAK\n")?, 605 } 606 match done { 607 true => Ok(stream(preamble, common)), 608 false => Ok(UploadOutcome::Buffered(preamble)), 609 } 610} 611 612pub fn stream_pack( 613 repo: &Repo, 614 wants: &WantOids, 615 haves: &HaveOids, 616 opts: &StreamOpts, 617 messages: &FetchMessages, 618 knot: &KnotHostname, 619 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>, 620) -> Result<(), PackError> { 621 let progress = opts.side_band && !opts.no_progress; 622 if opts.shallow_commits.is_none() 623 && haves.is_empty() 624 && opts.filter == Filter::None 625 && opts.packfile_uris.is_empty() 626 { 627 if let Ok(Some(pack)) = knot_git::verbatim_clone_pack(repo, wants.wants()) { 628 write_packfile_header(opts, sink)?; 629 return stream_verbatim_pack(pack, opts, progress, messages, knot, sink); 630 } 631 if let Ok(Some(oids)) = knot_git::reachable_via_bitmap(repo, wants.wants(), Haves::new(&[])) 632 { 633 write_packfile_header(opts, sink)?; 634 return stream_object_set(repo, oids, opts, progress, messages, knot, sink); 635 } 636 write_packfile_header(opts, sink)?; 637 return stream_full_clone(repo, wants.as_slice(), opts, progress, messages, knot, sink); 638 } 639 let budget = selection_budget(); 640 let mut selection = match &opts.shallow_commits { 641 Some(commits) => repo.select_shallow_objects( 642 wants.wants(), 643 ShallowCommits::new(commits), 644 haves.haves(), 645 opts.filter, 646 budget, 647 )?, 648 None => { 649 repo.select_pack_objects_filtered(wants.wants(), haves.haves(), opts.filter, budget)? 650 } 651 }; 652 if !opts.packfile_uris.is_empty() { 653 offload_packfile_uris( 654 &mut selection.send, 655 &opts.packfile_uris, 656 opts.sideband_all, 657 sink, 658 )?; 659 } 660 write_packfile_header(opts, sink)?; 661 let count = ObjectCount::new(selection.send.len()); 662 emit_preamble(progress, messages, knot, count, sink)?; 663 { 664 let mut pack_sink = PackSink { 665 side_band: opts.side_band, 666 sink: &mut *sink, 667 }; 668 let thin_bases = opts.thin.then_some(&selection.client_has); 669 objects::write_pack( 670 &repo.objects_dir(), 671 selection.send, 672 thin_bases, 673 &mut pack_sink, 674 repo.object_format().kind(), 675 )?; 676 } 677 emit_total(progress, messages, count, sink) 678} 679 680fn emit_progress( 681 progress: bool, 682 lines: impl FnOnce() -> Vec<String>, 683 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>, 684) -> Result<(), PackError> { 685 if !progress { 686 return Ok(()); 687 } 688 let rendered = lines(); 689 if rendered.is_empty() { 690 return Ok(()); 691 } 692 let mut buf = Vec::new(); 693 rendered.iter().try_for_each(|line| { 694 format!("{line}\n") 695 .into_bytes() 696 .chunks(pkt::MAX_BAND) 697 .try_for_each(|chunk| pkt::write_band_progress(&mut buf, chunk)) 698 })?; 699 sink(&buf).map_err(PackError::from) 700} 701 702fn emit_preamble( 703 progress: bool, 704 messages: &FetchMessages, 705 knot: &KnotHostname, 706 count: ObjectCount, 707 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>, 708) -> Result<(), PackError> { 709 emit_progress( 710 progress, 711 || { 712 messages 713 .motd 714 .lines(|KnotKey::Knot| knot.as_str().to_string()) 715 .into_iter() 716 .chain( 717 messages 718 .enumerating 719 .lines(|CountKey::Count| count.to_string()), 720 ) 721 .collect() 722 }, 723 sink, 724 ) 725} 726 727fn emit_total( 728 progress: bool, 729 messages: &FetchMessages, 730 count: ObjectCount, 731 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>, 732) -> Result<(), PackError> { 733 emit_progress( 734 progress, 735 || messages.total.lines(|CountKey::Count| count.to_string()), 736 sink, 737 ) 738} 739 740fn write_packfile_header( 741 opts: &StreamOpts, 742 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>, 743) -> Result<(), PackError> { 744 if !opts.emit_packfile_header { 745 return Ok(()); 746 } 747 let mut buf = Vec::new(); 748 seg(&mut buf, opts.sideband_all, b"packfile\n")?; 749 sink(&buf)?; 750 Ok(()) 751} 752 753fn offload_packfile_uris( 754 send: &mut Vec<Oid>, 755 candidates: &[PackfileUri], 756 sideband_all: bool, 757 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>, 758) -> Result<(), PackError> { 759 let present: std::collections::HashSet<Oid> = send.iter().copied().collect(); 760 let kept: Vec<&PackfileUri> = candidates 761 .iter() 762 .filter(|candidate| present.contains(&candidate.oid)) 763 .collect(); 764 if kept.is_empty() { 765 return Ok(()); 766 } 767 let excluded: std::collections::HashSet<Oid> = 768 kept.iter().map(|candidate| candidate.oid).collect(); 769 send.retain(|oid| !excluded.contains(oid)); 770 let mut buf = Vec::new(); 771 seg(&mut buf, sideband_all, b"packfile-uris\n")?; 772 kept.iter().try_for_each(|candidate| { 773 seg( 774 &mut buf, 775 sideband_all, 776 format!( 777 "{} {}\n", 778 candidate.pack_hash.as_str(), 779 candidate.uri.as_str() 780 ) 781 .as_bytes(), 782 ) 783 })?; 784 pkt::write_delim(&mut buf)?; 785 sink(&buf)?; 786 Ok(()) 787} 788 789fn stream_verbatim_pack( 790 mut file: std::fs::File, 791 opts: &StreamOpts, 792 progress: bool, 793 messages: &FetchMessages, 794 knot: &KnotHostname, 795 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>, 796) -> Result<(), PackError> { 797 let mut header = [0u8; 12]; 798 file.read_exact(&mut header)?; 799 if &header[..4] != b"PACK" { 800 return Err(PackError::Pack( 801 "reused pack is missing its PACK signature".to_string(), 802 )); 803 } 804 let count = ObjectCount::from(u32::from_be_bytes([ 805 header[8], header[9], header[10], header[11], 806 ])); 807 emit_preamble(progress, messages, knot, count, sink)?; 808 { 809 let mut pack_sink = PackSink { 810 side_band: opts.side_band, 811 sink: &mut *sink, 812 }; 813 pack_sink.write_all(&header)?; 814 io::copy(&mut file, &mut pack_sink)?; 815 } 816 emit_total(progress, messages, count, sink) 817} 818 819fn stream_object_set( 820 repo: &Repo, 821 oids: Vec<Oid>, 822 opts: &StreamOpts, 823 progress: bool, 824 messages: &FetchMessages, 825 knot: &KnotHostname, 826 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>, 827) -> Result<(), PackError> { 828 let count = ObjectCount::new(oids.len()); 829 emit_preamble(progress, messages, knot, count, sink)?; 830 { 831 let mut pack_sink = PackSink { 832 side_band: opts.side_band, 833 sink: &mut *sink, 834 }; 835 objects::write_pack( 836 &repo.objects_dir(), 837 oids, 838 None, 839 &mut pack_sink, 840 repo.object_format().kind(), 841 )?; 842 } 843 emit_total(progress, messages, count, sink) 844} 845 846fn stream_full_clone( 847 repo: &Repo, 848 wants: &[Oid], 849 opts: &StreamOpts, 850 progress: bool, 851 messages: &FetchMessages, 852 knot: &KnotHostname, 853 sink: &mut dyn FnMut(&[u8]) -> io::Result<()>, 854) -> Result<(), PackError> { 855 let limits = selection_limits(); 856 let stall = limits.time_budget; 857 let roots = repo.clone_roots(wants, PackBudget::new(limits.max_objects, stall))?; 858 let pack = objects::count_expanded( 859 &repo.objects_dir(), 860 roots, 861 limits.max_objects, 862 stall, 863 repo.object_format().kind(), 864 )?; 865 let count = ObjectCount::new(pack.len()); 866 emit_preamble(progress, messages, knot, count, sink)?; 867 { 868 let mut pack_sink = PackSink { 869 side_band: opts.side_band, 870 sink: &mut *sink, 871 }; 872 objects::write_expanded(pack, &mut pack_sink)?; 873 } 874 emit_total(progress, messages, count, sink) 875} 876 877struct PackSink<'a> { 878 side_band: bool, 879 sink: &'a mut dyn FnMut(&[u8]) -> io::Result<()>, 880} 881 882impl Write for PackSink<'_> { 883 fn write(&mut self, data: &[u8]) -> io::Result<usize> { 884 if self.side_band { 885 data.chunks(pkt::MAX_BAND).try_for_each(|chunk| { 886 let mut framed = Vec::with_capacity(chunk.len() + 5); 887 pkt::write_band(&mut framed, chunk)?; 888 (self.sink)(&framed) 889 })?; 890 } else { 891 (self.sink)(data)?; 892 } 893 Ok(data.len()) 894 } 895 896 fn flush(&mut self) -> io::Result<()> { 897 Ok(()) 898 } 899} 900 901pub fn selection_budget() -> PackBudget { 902 let limits = selection_limits(); 903 PackBudget::new(limits.max_objects, limits.time_budget) 904} 905 906fn ensure_wanted(repo: &Repo, wants: &WantOids) -> Result<(), PackError> { 907 let tips: Vec<Oid> = repo 908 .advertised_refs_for(knot_git::AdvertScope::Upload)? 909 .iter() 910 .map(|record| record.target) 911 .collect(); 912 let advertised: HashSet<Oid> = tips.iter().copied().collect(); 913 if wants.iter().all(|want| advertised.contains(want)) { 914 return Ok(()); 915 } 916 let commit_closure = repo.reachable_commits(&tips, selection_budget())?; 917 let unresolved: Vec<Oid> = wants 918 .iter() 919 .copied() 920 .filter(|want| !advertised.contains(want) && !commit_closure.contains(want)) 921 .collect(); 922 if unresolved.is_empty() { 923 return Ok(()); 924 } 925 let reachable: HashSet<Oid> = repo 926 .select_pack_objects_filtered( 927 Wants::new(&tips), 928 Haves::new(&[]), 929 Filter::None, 930 selection_budget(), 931 )? 932 .send 933 .into_iter() 934 .collect(); 935 match unresolved.iter().find(|want| !reachable.contains(want)) { 936 Some(hidden) => Err(PackError::Protocol(format!( 937 "want {hidden} isn't reachable from public ref" 938 ))), 939 None => Ok(()), 940 } 941} 942 943fn first_caps<'a>(lines: &[&'a [u8]]) -> Option<&'a str> { 944 let line = lines.iter().find(|line| line.starts_with(b"want "))?; 945 let text = std::str::from_utf8(line).ok()?.trim_end(); 946 text.strip_prefix("want ")? 947 .split_once(' ') 948 .map(|(_oid, caps)| caps) 949} 950 951fn parse_wants(lines: &[&[u8]]) -> Result<WantOids, PackError> { 952 lines 953 .iter() 954 .filter_map(|line| line.strip_prefix(b"want ")) 955 .map(|rest| { 956 let hex = rest 957 .split(|byte| *byte == b' ' || *byte == b'\n') 958 .next() 959 .unwrap_or_default(); 960 std::str::from_utf8(hex) 961 .ok() 962 .and_then(|text| Oid::from_hex(text).ok()) 963 .ok_or_else(|| PackError::Protocol("malformed want line".to_string())) 964 }) 965 .collect::<Result<Vec<_>, _>>() 966 .map(WantOids::new) 967} 968 969fn parse_oids(lines: &[&[u8]], prefix: &[u8]) -> Vec<Oid> { 970 lines 971 .iter() 972 .filter_map(|line| line.strip_prefix(prefix)) 973 .filter_map(|rest| { 974 let hex = rest.split(|byte| *byte == b' ' || *byte == b'\n').next()?; 975 let hex = std::str::from_utf8(hex).ok()?; 976 Oid::from_hex(hex).ok() 977 }) 978 .collect() 979}