This repository has no description
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}