| 1 | use std::{ |
| 2 | collections::HashSet, |
| 3 | path::Path, |
| 4 | }; |
| 5 | |
| 6 | use bytes::Bytes; |
| 7 | use sha1::{ |
| 8 | Digest, |
| 9 | Sha1, |
| 10 | }; |
| 11 | use tokio::io::AsyncRead; |
| 12 | use tokio_util::io::StreamReader; |
| 13 | |
| 14 | use crate::{ |
| 15 | error::{ |
| 16 | Error, |
| 17 | Result, |
| 18 | }, |
| 19 | pktline, |
| 20 | }; |
| 21 | |
| 22 | #[derive(Clone, Debug, Default)] |
| 23 | pub struct UploadPackCapabilities { |
| 24 | pub ofs_delta: bool, |
| 25 | pub multi_ack: bool, |
| 26 | pub multi_ack_detailed: bool, |
| 27 | } |
| 28 | |
| 29 | #[derive(Clone, Debug, Default)] |
| 30 | pub struct ShallowRequest { |
| 31 | pub depth: Option<usize>, |
| 32 | pub client_shallows: Vec<gix::ObjectId>, |
| 33 | pub deepen_relative: bool, |
| 34 | } |
| 35 | |
| 36 | /// A parsed upload-pack request from a Git client. |
| 37 | pub struct UploadPackRequest { |
| 38 | pub wants: Vec<gix::ObjectId>, |
| 39 | pub haves: Vec<gix::ObjectId>, |
| 40 | pub done: bool, |
| 41 | pub capabilities: UploadPackCapabilities, |
| 42 | pub shallow: ShallowRequest, |
| 43 | pub object_ids: Option<Vec<gix::ObjectId>>, |
| 44 | } |
| 45 | |
| 46 | impl UploadPackRequest { |
| 47 | /// Parse a pkt-line encoded upload-pack request body. |
| 48 | /// |
| 49 | /// The body contains: |
| 50 | /// - "want <oid> [capabilities]\n" lines |
| 51 | /// - flush packet "0000" |
| 52 | /// - "have <oid>\n" lines (optional) |
| 53 | /// - "done\n" |
| 54 | pub fn parse(body: &[u8]) -> Result<Self> { |
| 55 | let mut wants = Vec::new(); |
| 56 | let mut haves = Vec::new(); |
| 57 | let mut done = false; |
| 58 | let mut capabilities = UploadPackCapabilities::default(); |
| 59 | let mut shallow = ShallowRequest::default(); |
| 60 | let mut pos = 0; |
| 61 | |
| 62 | while pos < body.len() { |
| 63 | // Check for flush packet |
| 64 | if body[pos..].starts_with(b"0000") { |
| 65 | pos += 4; |
| 66 | continue; |
| 67 | } |
| 68 | |
| 69 | // Read 4-byte hex length prefix |
| 70 | if pos + 4 > body.len() { |
| 71 | break; |
| 72 | } |
| 73 | let len_str = std::str::from_utf8(&body[pos..pos + 4]) |
| 74 | .map_err(|_| Error::Protocol("invalid pkt-line length prefix".into()))?; |
| 75 | let len = usize::from_str_radix(len_str, 16) |
| 76 | .map_err(|_| Error::Protocol("invalid pkt-line length".into()))?; |
| 77 | |
| 78 | if len == 0 { |
| 79 | // flush packet already handled above, but just in case |
| 80 | pos += 4; |
| 81 | continue; |
| 82 | } |
| 83 | |
| 84 | if len < 4 || pos + len > body.len() { |
| 85 | break; |
| 86 | } |
| 87 | |
| 88 | let payload = &body[pos + 4..pos + len]; |
| 89 | let line = std::str::from_utf8(payload) |
| 90 | .map_err(|_| Error::Protocol("invalid UTF-8 in pkt-line".into()))?; |
| 91 | let line = line.trim_end_matches('\n'); |
| 92 | |
| 93 | if line == "done" { |
| 94 | done = true; |
| 95 | } else if let Some(rest) = line.strip_prefix("deepen ") { |
| 96 | let depth = rest |
| 97 | .parse::<usize>() |
| 98 | .map_err(|_| Error::Protocol(format!("invalid deepen value: {rest}")))?; |
| 99 | shallow.depth = Some(depth); |
| 100 | } else if line == "deepen-relative" { |
| 101 | shallow.deepen_relative = true; |
| 102 | } else if let Some(rest) = line.strip_prefix("shallow ") { |
| 103 | let oid = gix::ObjectId::from_hex(rest.as_bytes()) |
| 104 | .map_err(|_| Error::Protocol(format!("invalid OID in shallow: {rest}")))?; |
| 105 | shallow.client_shallows.push(oid); |
| 106 | } else if let Some(rest) = line.strip_prefix("want ") { |
| 107 | let mut parts = rest.split_ascii_whitespace(); |
| 108 | let oid_hex = parts |
| 109 | .next() |
| 110 | .ok_or_else(|| Error::Protocol("missing OID in want".into()))?; |
| 111 | let oid = gix::ObjectId::from_hex(oid_hex.as_bytes()) |
| 112 | .map_err(|_| Error::Protocol(format!("invalid OID in want: {oid_hex}")))?; |
| 113 | if wants.is_empty() { |
| 114 | for capability in parts { |
| 115 | if capability == "ofs-delta" { |
| 116 | capabilities.ofs_delta = true; |
| 117 | } else if capability == "multi_ack" { |
| 118 | capabilities.multi_ack = true; |
| 119 | } else if capability == "multi_ack_detailed" { |
| 120 | capabilities.multi_ack = true; |
| 121 | capabilities.multi_ack_detailed = true; |
| 122 | } |
| 123 | } |
| 124 | } |
| 125 | wants.push(oid); |
| 126 | } else if let Some(rest) = line.strip_prefix("have ") { |
| 127 | let oid_hex = rest |
| 128 | .split_ascii_whitespace() |
| 129 | .next() |
| 130 | .ok_or_else(|| Error::Protocol("missing OID in have".into()))?; |
| 131 | let oid = gix::ObjectId::from_hex(oid_hex.as_bytes()) |
| 132 | .map_err(|_| Error::Protocol(format!("invalid OID in have: {oid_hex}")))?; |
| 133 | haves.push(oid); |
| 134 | } |
| 135 | |
| 136 | pos += len; |
| 137 | } |
| 138 | |
| 139 | Ok(Self { |
| 140 | wants, |
| 141 | haves, |
| 142 | done, |
| 143 | capabilities, |
| 144 | shallow, |
| 145 | object_ids: None, |
| 146 | }) |
| 147 | } |
| 148 | } |
| 149 | |
| 150 | /// Encode the variable-length pack object header. |
| 151 | /// |
| 152 | /// Format: first byte = MSB continuation + 3-bit type + 4-bit size |
| 153 | /// Subsequent bytes: 7-bit size chunks with MSB continuation |
| 154 | fn encode_pack_object_header(obj_type: u8, size: usize) -> Vec<u8> { |
| 155 | let mut header = Vec::new(); |
| 156 | let mut byte = (obj_type << 4) | (size as u8 & 0x0f); |
| 157 | let mut remaining = size >> 4; |
| 158 | |
| 159 | if remaining > 0 { |
| 160 | byte |= 0x80; // set continuation bit |
| 161 | header.push(byte); |
| 162 | while remaining > 0 { |
| 163 | byte = remaining as u8 & 0x7f; |
| 164 | remaining >>= 7; |
| 165 | if remaining > 0 { |
| 166 | byte |= 0x80; |
| 167 | } |
| 168 | header.push(byte); |
| 169 | } |
| 170 | } else { |
| 171 | header.push(byte); |
| 172 | } |
| 173 | |
| 174 | header |
| 175 | } |
| 176 | |
| 177 | fn encode_ofs_delta_base_distance(mut distance: u64) -> Vec<u8> { |
| 178 | debug_assert!(distance > 0, "offset deltas must point backwards"); |
| 179 | |
| 180 | let mut buf = [0u8; 10]; |
| 181 | let mut bytes_written = 1; |
| 182 | buf[buf.len() - 1] = distance as u8 & 0x7f; |
| 183 | |
| 184 | for out in buf.iter_mut().rev().skip(1) { |
| 185 | distance >>= 7; |
| 186 | if distance == 0 { |
| 187 | break; |
| 188 | } |
| 189 | distance -= 1; |
| 190 | *out = 0x80 | (distance as u8 & 0x7f); |
| 191 | bytes_written += 1; |
| 192 | } |
| 193 | |
| 194 | buf[buf.len() - bytes_written..].to_vec() |
| 195 | } |
| 196 | |
| 197 | fn encode_delta_size(mut size: usize, out: &mut Vec<u8>) { |
| 198 | loop { |
| 199 | let mut byte = (size & 0x7f) as u8; |
| 200 | size >>= 7; |
| 201 | if size > 0 { |
| 202 | byte |= 0x80; |
| 203 | } |
| 204 | out.push(byte); |
| 205 | if size == 0 { |
| 206 | break; |
| 207 | } |
| 208 | } |
| 209 | } |
| 210 | |
| 211 | fn encode_delta_copy_instruction(out: &mut Vec<u8>, offset: usize, size: usize) { |
| 212 | debug_assert!(size > 0 && size <= 0x10000); |
| 213 | |
| 214 | let command_pos = out.len(); |
| 215 | out.push(0x80); |
| 216 | let mut command = 0x80; |
| 217 | |
| 218 | if offset & 0xff != 0 { |
| 219 | command |= 0x01; |
| 220 | out.push(offset as u8); |
| 221 | } |
| 222 | if (offset >> 8) & 0xff != 0 { |
| 223 | command |= 0x02; |
| 224 | out.push((offset >> 8) as u8); |
| 225 | } |
| 226 | if (offset >> 16) & 0xff != 0 { |
| 227 | command |= 0x04; |
| 228 | out.push((offset >> 16) as u8); |
| 229 | } |
| 230 | if (offset >> 24) & 0xff != 0 { |
| 231 | command |= 0x08; |
| 232 | out.push((offset >> 24) as u8); |
| 233 | } |
| 234 | |
| 235 | if size != 0x10000 { |
| 236 | if size & 0xff != 0 { |
| 237 | command |= 0x10; |
| 238 | out.push(size as u8); |
| 239 | } |
| 240 | if (size >> 8) & 0xff != 0 { |
| 241 | command |= 0x20; |
| 242 | out.push((size >> 8) as u8); |
| 243 | } |
| 244 | if (size >> 16) & 0xff != 0 { |
| 245 | command |= 0x40; |
| 246 | out.push((size >> 16) as u8); |
| 247 | } |
| 248 | } |
| 249 | |
| 250 | out[command_pos] = command; |
| 251 | } |
| 252 | |
| 253 | fn encode_delta_copy(out: &mut Vec<u8>, mut offset: usize, mut size: usize) { |
| 254 | while size > 0 { |
| 255 | let chunk = size.min(0x10000); |
| 256 | encode_delta_copy_instruction(out, offset, chunk); |
| 257 | offset += chunk; |
| 258 | size -= chunk; |
| 259 | } |
| 260 | } |
| 261 | |
| 262 | fn encode_delta_insert(out: &mut Vec<u8>, data: &[u8]) { |
| 263 | for chunk in data.chunks(0x7f) { |
| 264 | out.push(chunk.len() as u8); |
| 265 | out.extend_from_slice(chunk); |
| 266 | } |
| 267 | } |
| 268 | |
| 269 | fn encode_blob_delta(base: &[u8], target: &[u8]) -> Option<Vec<u8>> { |
| 270 | let mut prefix = 0; |
| 271 | let max_prefix = base.len().min(target.len()); |
| 272 | while prefix < max_prefix && base[prefix] == target[prefix] { |
| 273 | prefix += 1; |
| 274 | } |
| 275 | |
| 276 | let max_suffix = base |
| 277 | .len() |
| 278 | .saturating_sub(prefix) |
| 279 | .min(target.len().saturating_sub(prefix)); |
| 280 | let mut suffix = 0; |
| 281 | while suffix < max_suffix && base[base.len() - 1 - suffix] == target[target.len() - 1 - suffix] |
| 282 | { |
| 283 | suffix += 1; |
| 284 | } |
| 285 | |
| 286 | if prefix == 0 && suffix == 0 { |
| 287 | return None; |
| 288 | } |
| 289 | |
| 290 | let mut delta = Vec::new(); |
| 291 | encode_delta_size(base.len(), &mut delta); |
| 292 | encode_delta_size(target.len(), &mut delta); |
| 293 | |
| 294 | if prefix > 0 { |
| 295 | encode_delta_copy(&mut delta, 0, prefix); |
| 296 | } |
| 297 | |
| 298 | let insert_start = prefix; |
| 299 | let insert_end = target.len() - suffix; |
| 300 | encode_delta_insert(&mut delta, &target[insert_start..insert_end]); |
| 301 | |
| 302 | if suffix > 0 { |
| 303 | encode_delta_copy(&mut delta, base.len() - suffix, suffix); |
| 304 | } |
| 305 | |
| 306 | Some(delta) |
| 307 | } |
| 308 | |
| 309 | fn build_base_entry(kind: gix::object::Kind, data: &[u8]) -> Vec<u8> { |
| 310 | let type_num = object_type_number(kind); |
| 311 | let obj_header = encode_pack_object_header(type_num, data.len()); |
| 312 | let compressed = miniz_oxide::deflate::compress_to_vec_zlib(data, 6); |
| 313 | |
| 314 | let mut entry = Vec::with_capacity(obj_header.len() + compressed.len()); |
| 315 | entry.extend_from_slice(&obj_header); |
| 316 | entry.extend_from_slice(&compressed); |
| 317 | entry |
| 318 | } |
| 319 | |
| 320 | fn build_ofs_delta_entry( |
| 321 | pack_offset: u64, |
| 322 | base_pack_offset: u64, |
| 323 | base_data: &[u8], |
| 324 | target_data: &[u8], |
| 325 | ) -> Option<Vec<u8>> { |
| 326 | let delta = encode_blob_delta(base_data, target_data)?; |
| 327 | let obj_header = encode_pack_object_header(6, delta.len()); |
| 328 | let base_distance = encode_ofs_delta_base_distance(pack_offset - base_pack_offset); |
| 329 | let compressed = miniz_oxide::deflate::compress_to_vec_zlib(&delta, 6); |
| 330 | |
| 331 | let mut entry = Vec::with_capacity(obj_header.len() + base_distance.len() + compressed.len()); |
| 332 | entry.extend_from_slice(&obj_header); |
| 333 | entry.extend_from_slice(&base_distance); |
| 334 | entry.extend_from_slice(&compressed); |
| 335 | Some(entry) |
| 336 | } |
| 337 | |
| 338 | struct BlobDeltaBase { |
| 339 | pack_offset: u64, |
| 340 | data: Vec<u8>, |
| 341 | } |
| 342 | |
| 343 | /// Map gix object kind to pack type number. |
| 344 | fn object_type_number(kind: gix::object::Kind) -> u8 { |
| 345 | match kind { |
| 346 | gix::object::Kind::Commit => 1, |
| 347 | gix::object::Kind::Tree => 2, |
| 348 | gix::object::Kind::Blob => 3, |
| 349 | gix::object::Kind::Tag => 4, |
| 350 | } |
| 351 | } |
| 352 | |
| 353 | /// Send raw bytes through the channel. |
| 354 | fn send( |
| 355 | tx: &tokio::sync::mpsc::Sender<std::result::Result<Bytes, std::io::Error>>, |
| 356 | data: &[u8], |
| 357 | ) -> std::result::Result<(), Box<dyn std::error::Error + Send + Sync>> { |
| 358 | tx.blocking_send(Ok(Bytes::copy_from_slice(data))) |
| 359 | .map_err(|_| "receiver dropped".into()) |
| 360 | } |
| 361 | |
| 362 | /// Send pack data through the channel wrapped in side-band-64k framing |
| 363 | /// (band 1 = pack data). |
| 364 | /// |
| 365 | /// Respects LARGE_PACKET_MAX: each pkt-line frame carries at most |
| 366 | /// 65520 - 4 (prefix) - 1 (band byte) = 65515 bytes of payload. |
| 367 | fn send_sideband( |
| 368 | tx: &tokio::sync::mpsc::Sender<std::result::Result<Bytes, std::io::Error>>, |
| 369 | data: &[u8], |
| 370 | ) -> std::result::Result<(), Box<dyn std::error::Error + Send + Sync>> { |
| 371 | const MAX_DATA_PER_FRAME: usize = 65515; |
| 372 | |
| 373 | for chunk in data.chunks(MAX_DATA_PER_FRAME) { |
| 374 | let pkt_len = 4 + 1 + chunk.len(); |
| 375 | let mut frame = Vec::with_capacity(pkt_len); |
| 376 | frame.extend_from_slice(format!("{pkt_len:04x}").as_bytes()); |
| 377 | frame.push(0x01); // band 1 = pack data |
| 378 | frame.extend_from_slice(chunk); |
| 379 | send(tx, &frame)?; |
| 380 | } |
| 381 | |
| 382 | Ok(()) |
| 383 | } |
| 384 | |
| 385 | fn encode_ack_line(oid: gix::ObjectId, suffix: Option<&str>) -> Vec<u8> { |
| 386 | let mut line = format!("ACK {oid}"); |
| 387 | if let Some(suffix) = suffix { |
| 388 | line.push(' '); |
| 389 | line.push_str(suffix); |
| 390 | } |
| 391 | line.push('\n'); |
| 392 | pktline::encode(line.as_bytes()) |
| 393 | } |
| 394 | |
| 395 | /// Recursively collect tree and blob OIDs reachable from `tree_oid`. |
| 396 | /// |
| 397 | /// Uses a single `find_object` call per object and parses raw tree |
| 398 | /// bytes via `TreeRefIter` to avoid a second ODB lookup. |
| 399 | fn collect_tree_oids( |
| 400 | repo: &gix::Repository, |
| 401 | tree_oid: gix::ObjectId, |
| 402 | seen: &mut HashSet<gix::ObjectId>, |
| 403 | oids: &mut Vec<gix::ObjectId>, |
| 404 | ) -> std::result::Result<(), Box<dyn std::error::Error + Send + Sync>> { |
| 405 | if !seen.insert(tree_oid) { |
| 406 | return Ok(()); |
| 407 | } |
| 408 | |
| 409 | let tree_obj = repo.find_object(tree_oid)?; |
| 410 | let tree_data = tree_obj.data.to_vec(); |
| 411 | oids.push(tree_oid); |
| 412 | |
| 413 | for entry_result in gix::objs::TreeRefIter::from_bytes(&tree_data, gix::hash::Kind::Sha1) { |
| 414 | let entry = entry_result?; |
| 415 | let entry_oid = entry.oid.to_owned(); |
| 416 | let entry_mode = entry.mode; |
| 417 | |
| 418 | if entry_mode.is_tree() { |
| 419 | collect_tree_oids(repo, entry_oid, seen, oids)?; |
| 420 | } else if seen.insert(entry_oid) && !entry_mode.is_commit() { |
| 421 | oids.push(entry_oid); |
| 422 | } |
| 423 | } |
| 424 | |
| 425 | Ok(()) |
| 426 | } |
| 427 | |
| 428 | /// Peel `wants` for a pack walk: tag wants contribute the tag object itself |
| 429 | /// (plus the chain down to a non-tag), and the peeled target joins the walk — |
| 430 | /// commits as walk tips, trees/blobs (rare tag targets) collected directly. |
| 431 | /// `git fetch`/`push` send tag *object* ids whenever annotated tags are |
| 432 | /// involved, so every pack path needs this. |
| 433 | pub(crate) fn peel_wants( |
| 434 | repo: &gix::Repository, |
| 435 | wants: &[gix::ObjectId], |
| 436 | ) -> std::result::Result< |
| 437 | (Vec<gix::ObjectId>, Vec<gix::ObjectId>), |
| 438 | Box<dyn std::error::Error + Send + Sync>, |
| 439 | > { |
| 440 | let mut commit_wants = Vec::new(); |
| 441 | let mut extras = Vec::new(); |
| 442 | let mut seen = HashSet::new(); |
| 443 | for want in wants { |
| 444 | let mut oid = *want; |
| 445 | loop { |
| 446 | let obj = repo.find_object(oid)?; |
| 447 | match obj.kind { |
| 448 | gix::object::Kind::Tag => { |
| 449 | if seen.insert(oid) { |
| 450 | extras.push(oid); |
| 451 | } |
| 452 | oid = gix::objs::TagRefIter::from_bytes(&obj.data, gix::hash::Kind::Sha1) |
| 453 | .target_id()?; |
| 454 | } |
| 455 | gix::object::Kind::Commit => { |
| 456 | commit_wants.push(oid); |
| 457 | break; |
| 458 | } |
| 459 | gix::object::Kind::Tree => { |
| 460 | collect_tree_oids(repo, oid, &mut seen, &mut extras)?; |
| 461 | break; |
| 462 | } |
| 463 | gix::object::Kind::Blob => { |
| 464 | if seen.insert(oid) { |
| 465 | extras.push(oid); |
| 466 | } |
| 467 | break; |
| 468 | } |
| 469 | } |
| 470 | } |
| 471 | } |
| 472 | Ok((commit_wants, extras)) |
| 473 | } |
| 474 | |
| 475 | /// Walk commits from `wants` (excluding `haves`) and collect all |
| 476 | /// reachable ObjectIds (commits, trees, blobs). Tag wants are peeled, with |
| 477 | /// the tag objects themselves included in the result. |
| 478 | /// |
| 479 | /// Pass 1 of the two-pass streaming approach: only OIDs are stored, |
| 480 | /// not object data. |
| 481 | fn collect_all_oids( |
| 482 | repo: &gix::Repository, |
| 483 | wants: &[gix::ObjectId], |
| 484 | haves: &[gix::ObjectId], |
| 485 | ) -> std::result::Result<Vec<gix::ObjectId>, Box<dyn std::error::Error + Send + Sync>> { |
| 486 | let (commit_wants, extras) = peel_wants(repo, wants)?; |
| 487 | let have_set: HashSet<gix::ObjectId> = haves.iter().copied().collect(); |
| 488 | let mut seen = HashSet::new(); |
| 489 | let mut oids = Vec::new(); |
| 490 | |
| 491 | // Mark have objects as already seen so we skip them |
| 492 | for have in haves { |
| 493 | seen.insert(*have); |
| 494 | } |
| 495 | |
| 496 | let walk = repo |
| 497 | .rev_walk(commit_wants.iter().copied()) |
| 498 | .with_hidden(haves.iter().copied()) |
| 499 | .all()?; |
| 500 | |
| 501 | for info_result in walk { |
| 502 | let info = info_result?; |
| 503 | let commit_oid = info.id; |
| 504 | |
| 505 | if have_set.contains(&commit_oid) || !seen.insert(commit_oid) { |
| 506 | continue; |
| 507 | } |
| 508 | |
| 509 | // Extract tree OID from raw commit bytes (single ODB read) |
| 510 | let commit_obj = repo.find_object(commit_oid)?; |
| 511 | let tree_oid = |
| 512 | gix::objs::CommitRefIter::from_bytes(&commit_obj.data, gix::hash::Kind::Sha1) |
| 513 | .tree_id()?; |
| 514 | |
| 515 | oids.push(commit_oid); |
| 516 | |
| 517 | collect_tree_oids(repo, tree_oid, &mut seen, &mut oids)?; |
| 518 | } |
| 519 | |
| 520 | // Tag objects (and direct tree/blob tag targets) ride along after the |
| 521 | // walk; anything already collected or known to the remote is skipped. |
| 522 | oids.extend( |
| 523 | extras |
| 524 | .into_iter() |
| 525 | .filter(|oid| !seen.contains(oid) && !have_set.contains(oid)), |
| 526 | ); |
| 527 | |
| 528 | Ok(oids) |
| 529 | } |
| 530 | |
| 531 | fn common_haves( |
| 532 | repo: &gix::Repository, |
| 533 | wants: &[gix::ObjectId], |
| 534 | haves: &[gix::ObjectId], |
| 535 | ) -> std::result::Result<Vec<gix::ObjectId>, Box<dyn std::error::Error + Send + Sync>> { |
| 536 | let want_set: HashSet<gix::ObjectId> = |
| 537 | collect_all_oids(repo, wants, &[])?.into_iter().collect(); |
| 538 | |
| 539 | Ok(haves |
| 540 | .iter() |
| 541 | .copied() |
| 542 | .filter(|oid| want_set.contains(oid)) |
| 543 | .collect()) |
| 544 | } |
| 545 | |
| 546 | /// Build a complete, self-contained pack (header + entries + SHA-1 trailer) |
| 547 | /// of everything reachable from `wants` but not from `haves`, returned as one |
| 548 | /// buffer. Entries are undeltified base objects — used by the *push* |
| 549 | /// (send-pack) client, where mirror pushes are infrequent enough that pack |
| 550 | /// size doesn't justify delta search. |
| 551 | /// |
| 552 | /// `haves` must exist in the local object database (callers filter the remote |
| 553 | /// advertisement accordingly). |
| 554 | pub fn build_raw_pack( |
| 555 | repo_path: &Path, |
| 556 | wants: &[gix::ObjectId], |
| 557 | haves: &[gix::ObjectId], |
| 558 | ) -> Result<Vec<u8>> { |
| 559 | let repo = gix::open(repo_path)?; |
| 560 | // Tag wants (annotated tags) are peeled inside `collect_all_oids`, with |
| 561 | // the tag objects themselves included in the result. |
| 562 | let oids = collect_all_oids(&repo, wants, haves) |
| 563 | .map_err(|e| Error::Protocol(format!("collecting pack objects: {e}")))?; |
| 564 | |
| 565 | let mut out = Vec::new(); |
| 566 | out.extend_from_slice(b"PACK"); |
| 567 | out.extend_from_slice(&2u32.to_be_bytes()); |
| 568 | out.extend_from_slice(&(oids.len() as u32).to_be_bytes()); |
| 569 | for oid in &oids { |
| 570 | let obj = repo |
| 571 | .find_object(*oid) |
| 572 | .map_err(|e| Error::Protocol(format!("reading {oid}: {e}")))?; |
| 573 | out.extend_from_slice(&build_base_entry(obj.kind, &obj.data)); |
| 574 | } |
| 575 | let checksum = Sha1::digest(&out); |
| 576 | out.extend_from_slice(&checksum); |
| 577 | Ok(out) |
| 578 | } |
| 579 | |
| 580 | /// Generate the complete pack response for a Git upload-pack request. |
| 581 | /// |
| 582 | /// Returns an `AsyncRead` producing the side-band-64k framed response that |
| 583 | /// can be streamed as the HTTP response body. |
| 584 | pub fn generate_pack( |
| 585 | repo_path: &Path, |
| 586 | request: &UploadPackRequest, |
| 587 | ) -> Result<impl AsyncRead + Send + Unpin + use<>> { |
| 588 | let repo_path = repo_path.to_path_buf(); |
| 589 | let wants: Vec<gix::ObjectId> = request.wants.clone(); |
| 590 | let haves: Vec<gix::ObjectId> = request.haves.clone(); |
| 591 | let object_ids = request.object_ids.clone(); |
| 592 | let done = request.done; |
| 593 | let ofs_delta = request.capabilities.ofs_delta; |
| 594 | let multi_ack = request.capabilities.multi_ack; |
| 595 | let multi_ack_detailed = request.capabilities.multi_ack_detailed; |
| 596 | |
| 597 | let (tx, rx) = tokio::sync::mpsc::channel::<std::result::Result<Bytes, std::io::Error>>(64); |
| 598 | |
| 599 | let handle = tokio::task::spawn_blocking(move || { |
| 600 | if let Err(e) = generate_pack_sync( |
| 601 | &repo_path, |
| 602 | &wants, |
| 603 | &haves, |
| 604 | object_ids, |
| 605 | GeneratePackOptions { |
| 606 | done, |
| 607 | ofs_delta, |
| 608 | multi_ack, |
| 609 | multi_ack_detailed, |
| 610 | }, |
| 611 | &tx, |
| 612 | ) { |
| 613 | let _ = tx.blocking_send(Err(std::io::Error::other(e.to_string()))); |
| 614 | } |
| 615 | }); |
| 616 | |
| 617 | // Log panics from the blocking task without blocking the stream |
| 618 | tokio::spawn(async move { |
| 619 | if let Err(e) = handle.await { |
| 620 | tracing::error!("pack generation task panicked: {e}"); |
| 621 | } |
| 622 | }); |
| 623 | |
| 624 | let stream = tokio_stream::wrappers::ReceiverStream::new(rx); |
| 625 | Ok(StreamReader::new(stream)) |
| 626 | } |
| 627 | |
| 628 | /// Synchronous two-pass streaming pack generator. |
| 629 | /// |
| 630 | /// Pass 1: collect OIDs only (lightweight -- no object data retained). |
| 631 | /// Pass 2: re-read each object, compress, and stream it through `tx`. |
| 632 | struct GeneratePackOptions { |
| 633 | done: bool, |
| 634 | ofs_delta: bool, |
| 635 | multi_ack: bool, |
| 636 | multi_ack_detailed: bool, |
| 637 | } |
| 638 | |
| 639 | fn generate_pack_sync( |
| 640 | repo_path: &Path, |
| 641 | wants: &[gix::ObjectId], |
| 642 | haves: &[gix::ObjectId], |
| 643 | object_ids: Option<Vec<gix::ObjectId>>, |
| 644 | options: GeneratePackOptions, |
| 645 | tx: &tokio::sync::mpsc::Sender<std::result::Result<Bytes, std::io::Error>>, |
| 646 | ) -> std::result::Result<(), Box<dyn std::error::Error + Send + Sync>> { |
| 647 | const MAX_DELTA_BASES: usize = 8; |
| 648 | const MIN_DELTA_BLOB_SIZE: usize = 1024; |
| 649 | |
| 650 | let repo = gix::open(repo_path)?; |
| 651 | |
| 652 | let common = if !haves.is_empty() { |
| 653 | common_haves(&repo, wants, haves)? |
| 654 | } else { |
| 655 | Vec::new() |
| 656 | }; |
| 657 | |
| 658 | if options.multi_ack && !haves.is_empty() && !options.done { |
| 659 | for oid in &common { |
| 660 | let suffix = if options.multi_ack_detailed { |
| 661 | "common" |
| 662 | } else { |
| 663 | "continue" |
| 664 | }; |
| 665 | send(tx, &encode_ack_line(*oid, Some(suffix)))?; |
| 666 | } |
| 667 | send(tx, &pktline::encode(b"NAK\n"))?; |
| 668 | return Ok(()); |
| 669 | } |
| 670 | |
| 671 | if options.multi_ack && !common.is_empty() { |
| 672 | send(tx, &encode_ack_line(*common.last().unwrap(), None))?; |
| 673 | } else { |
| 674 | // NAK line |
| 675 | send(tx, &pktline::encode(b"NAK\n"))?; |
| 676 | } |
| 677 | |
| 678 | // Pass 1: collect OIDs only |
| 679 | let oids = match object_ids { |
| 680 | Some(oids) => oids, |
| 681 | None => collect_all_oids(&repo, wants, haves)?, |
| 682 | }; |
| 683 | |
| 684 | // Pass 2: stream each object |
| 685 | let mut hasher = Sha1::new(); |
| 686 | |
| 687 | // Pack header |
| 688 | let mut header = Vec::with_capacity(12); |
| 689 | header.extend_from_slice(b"PACK"); |
| 690 | header.extend_from_slice(&2u32.to_be_bytes()); |
| 691 | header.extend_from_slice(&(oids.len() as u32).to_be_bytes()); |
| 692 | hasher.update(&header); |
| 693 | send_sideband(tx, &header)?; |
| 694 | |
| 695 | let mut pack_offset = header.len() as u64; |
| 696 | let mut recent_blob_bases = Vec::<BlobDeltaBase>::new(); |
| 697 | |
| 698 | // Each object: read, compress, frame, send |
| 699 | for oid in &oids { |
| 700 | let obj = repo.find_object(*oid)?; |
| 701 | let full_entry = build_base_entry(obj.kind, &obj.data); |
| 702 | let mut used_delta = false; |
| 703 | let entry = if options.ofs_delta |
| 704 | && obj.kind == gix::object::Kind::Blob |
| 705 | && obj.data.len() >= MIN_DELTA_BLOB_SIZE |
| 706 | { |
| 707 | recent_blob_bases |
| 708 | .iter() |
| 709 | .filter(|base| base.data.len() >= MIN_DELTA_BLOB_SIZE) |
| 710 | .filter_map(|base| { |
| 711 | build_ofs_delta_entry(pack_offset, base.pack_offset, &base.data, &obj.data) |
| 712 | }) |
| 713 | .min_by_key(Vec::len) |
| 714 | .filter(|delta_entry| delta_entry.len() < full_entry.len()) |
| 715 | .inspect(|_| { |
| 716 | used_delta = true; |
| 717 | }) |
| 718 | .unwrap_or(full_entry) |
| 719 | } else { |
| 720 | full_entry |
| 721 | }; |
| 722 | |
| 723 | hasher.update(&entry); |
| 724 | send_sideband(tx, &entry)?; |
| 725 | |
| 726 | if obj.kind == gix::object::Kind::Blob |
| 727 | && !used_delta |
| 728 | && obj.data.len() >= MIN_DELTA_BLOB_SIZE |
| 729 | { |
| 730 | recent_blob_bases.push(BlobDeltaBase { |
| 731 | pack_offset, |
| 732 | data: obj.data.to_vec(), |
| 733 | }); |
| 734 | if recent_blob_bases.len() > MAX_DELTA_BASES { |
| 735 | recent_blob_bases.remove(0); |
| 736 | } |
| 737 | } |
| 738 | |
| 739 | pack_offset += entry.len() as u64; |
| 740 | } |
| 741 | |
| 742 | // SHA-1 checksum over raw pack bytes |
| 743 | let checksum = hasher.finalize(); |
| 744 | send_sideband(tx, &checksum)?; |
| 745 | |
| 746 | // Flush |
| 747 | send(tx, b"0000")?; |
| 748 | |
| 749 | Ok(()) |
| 750 | } |
| 751 | |
| 752 | #[cfg(test)] |
| 753 | mod tests { |
| 754 | use std::{ |
| 755 | path::{ |
| 756 | Path, |
| 757 | PathBuf, |
| 758 | }, |
| 759 | process::Command, |
| 760 | }; |
| 761 | |
| 762 | use tempfile::TempDir; |
| 763 | use tokio::io::AsyncReadExt; |
| 764 | |
| 765 | use super::*; |
| 766 | |
| 767 | fn make_pktline(data: &str) -> Vec<u8> { |
| 768 | let len = data.len() + 4; |
| 769 | format!("{len:04x}{data}").into_bytes() |
| 770 | } |
| 771 | |
| 772 | /// Create a bare repo with a single commit on the `main` branch. |
| 773 | fn create_repo_with_commit(root: &Path) -> PathBuf { |
| 774 | let bare_path = root.join("test.git"); |
| 775 | let clone_path = root.join("workdir"); |
| 776 | |
| 777 | let out = Command::new("git") |
| 778 | .args(["init", "--bare", bare_path.to_str().unwrap()]) |
| 779 | .output() |
| 780 | .expect("git init --bare failed"); |
| 781 | assert!(out.status.success(), "git init --bare failed: {:?}", out); |
| 782 | |
| 783 | let out = Command::new("git") |
| 784 | .args(["symbolic-ref", "HEAD", "refs/heads/main"]) |
| 785 | .current_dir(&bare_path) |
| 786 | .output() |
| 787 | .expect("git symbolic-ref failed"); |
| 788 | assert!(out.status.success()); |
| 789 | |
| 790 | let out = Command::new("git") |
| 791 | .args([ |
| 792 | "clone", |
| 793 | bare_path.to_str().unwrap(), |
| 794 | clone_path.to_str().unwrap(), |
| 795 | ]) |
| 796 | .output() |
| 797 | .expect("git clone failed"); |
| 798 | assert!(out.status.success(), "git clone failed: {:?}", out); |
| 799 | |
| 800 | for (key, val) in [("user.name", "Test User"), ("user.email", "test@test.com")] { |
| 801 | Command::new("git") |
| 802 | .args(["config", key, val]) |
| 803 | .current_dir(&clone_path) |
| 804 | .output() |
| 805 | .expect("git config failed"); |
| 806 | } |
| 807 | |
| 808 | // Create a file and commit |
| 809 | std::fs::write(clone_path.join("README.md"), "# Test\n").unwrap(); |
| 810 | |
| 811 | Command::new("git") |
| 812 | .args(["add", "README.md"]) |
| 813 | .current_dir(&clone_path) |
| 814 | .output() |
| 815 | .expect("git add failed"); |
| 816 | |
| 817 | let out = Command::new("git") |
| 818 | .args(["commit", "-m", "initial commit"]) |
| 819 | .current_dir(&clone_path) |
| 820 | .env("GIT_AUTHOR_NAME", "Test User") |
| 821 | .env("GIT_AUTHOR_EMAIL", "test@test.com") |
| 822 | .env("GIT_COMMITTER_NAME", "Test User") |
| 823 | .env("GIT_COMMITTER_EMAIL", "test@test.com") |
| 824 | .output() |
| 825 | .expect("git commit failed"); |
| 826 | assert!(out.status.success(), "git commit failed: {:?}", out); |
| 827 | |
| 828 | let out = Command::new("git") |
| 829 | .args(["push", "origin", "main"]) |
| 830 | .current_dir(&clone_path) |
| 831 | .output() |
| 832 | .expect("git push failed"); |
| 833 | assert!(out.status.success(), "git push failed: {:?}", out); |
| 834 | |
| 835 | bare_path |
| 836 | } |
| 837 | |
| 838 | #[test] |
| 839 | fn parse_simple_want() { |
| 840 | let hash = "0000000000000000000000000000000000000001"; |
| 841 | let mut body = make_pktline(&format!("want {hash}\n")); |
| 842 | body.extend_from_slice(b"00000009done\n"); |
| 843 | let req = UploadPackRequest::parse(&body).unwrap(); |
| 844 | assert_eq!(req.wants.len(), 1); |
| 845 | assert!(req.haves.is_empty()); |
| 846 | assert!(req.done); |
| 847 | assert!(!req.capabilities.ofs_delta); |
| 848 | assert_eq!(req.shallow.depth, None); |
| 849 | } |
| 850 | |
| 851 | #[test] |
| 852 | fn parse_wants_and_haves() { |
| 853 | let want = "0000000000000000000000000000000000000001"; |
| 854 | let have = "0000000000000000000000000000000000000002"; |
| 855 | let mut body = make_pktline(&format!("want {want}\n")); |
| 856 | body.extend_from_slice(b"0000"); |
| 857 | body.extend_from_slice(&make_pktline(&format!("have {have}\n"))); |
| 858 | body.extend_from_slice(b"0009done\n"); |
| 859 | let req = UploadPackRequest::parse(&body).unwrap(); |
| 860 | assert_eq!(req.wants.len(), 1); |
| 861 | assert_eq!(req.haves.len(), 1); |
| 862 | assert!(req.done); |
| 863 | assert!(!req.capabilities.ofs_delta); |
| 864 | assert!(req.shallow.client_shallows.is_empty()); |
| 865 | } |
| 866 | |
| 867 | #[test] |
| 868 | fn parse_ofs_delta_capability() { |
| 869 | let hash = "0000000000000000000000000000000000000001"; |
| 870 | let mut body = make_pktline(&format!("want {hash} side-band-64k ofs-delta\n")); |
| 871 | body.extend_from_slice(b"0009done\n"); |
| 872 | let req = UploadPackRequest::parse(&body).unwrap(); |
| 873 | assert!(req.capabilities.ofs_delta); |
| 874 | } |
| 875 | |
| 876 | #[test] |
| 877 | fn parse_multi_ack_capability() { |
| 878 | let hash = "0000000000000000000000000000000000000001"; |
| 879 | let mut body = make_pktline(&format!("want {hash} multi_ack side-band-64k\n")); |
| 880 | body.extend_from_slice(b"0009done\n"); |
| 881 | let req = UploadPackRequest::parse(&body).unwrap(); |
| 882 | assert!(req.capabilities.multi_ack); |
| 883 | } |
| 884 | |
| 885 | #[test] |
| 886 | fn parse_multi_ack_detailed_capability() { |
| 887 | let hash = "0000000000000000000000000000000000000001"; |
| 888 | let mut body = make_pktline(&format!("want {hash} multi_ack_detailed side-band-64k\n")); |
| 889 | body.extend_from_slice(b"0009done\n"); |
| 890 | let req = UploadPackRequest::parse(&body).unwrap(); |
| 891 | assert!(req.capabilities.multi_ack); |
| 892 | assert!(req.capabilities.multi_ack_detailed); |
| 893 | } |
| 894 | |
| 895 | #[test] |
| 896 | fn parse_shallow_request() { |
| 897 | let hash = "0000000000000000000000000000000000000001"; |
| 898 | let mut body = make_pktline(&format!("want {hash}\n")); |
| 899 | body.extend_from_slice(&make_pktline("deepen 2\n")); |
| 900 | body.extend_from_slice(&make_pktline(&format!("shallow {hash}\n"))); |
| 901 | body.extend_from_slice(&make_pktline("deepen-relative\n")); |
| 902 | body.extend_from_slice(b"0009done\n"); |
| 903 | let req = UploadPackRequest::parse(&body).unwrap(); |
| 904 | assert_eq!(req.shallow.depth, Some(2)); |
| 905 | assert_eq!( |
| 906 | req.shallow.client_shallows, |
| 907 | vec![gix::ObjectId::from_hex(hash.as_bytes()).unwrap()] |
| 908 | ); |
| 909 | assert!(req.shallow.deepen_relative); |
| 910 | } |
| 911 | |
| 912 | #[tokio::test] |
| 913 | async fn generate_pack_for_clone() { |
| 914 | let dir = TempDir::new().unwrap(); |
| 915 | let repo_path = create_repo_with_commit(dir.path()); |
| 916 | |
| 917 | // Get HEAD OID |
| 918 | let repo = gix::open(&repo_path).unwrap(); |
| 919 | let head_oid = repo.head_id().unwrap().detach(); |
| 920 | drop(repo); |
| 921 | |
| 922 | let request = UploadPackRequest { |
| 923 | wants: vec![head_oid], |
| 924 | haves: vec![], |
| 925 | done: true, |
| 926 | capabilities: UploadPackCapabilities::default(), |
| 927 | shallow: ShallowRequest::default(), |
| 928 | object_ids: None, |
| 929 | }; |
| 930 | |
| 931 | let mut reader = generate_pack(&repo_path, &request).unwrap(); |
| 932 | let mut buf = Vec::new(); |
| 933 | reader.read_to_end(&mut buf).await.unwrap(); |
| 934 | |
| 935 | let response = String::from_utf8_lossy(&buf); |
| 936 | assert!( |
| 937 | response.contains("NAK"), |
| 938 | "response should contain NAK: {response:?}" |
| 939 | ); |
| 940 | |
| 941 | // Find PACK signature in the binary response |
| 942 | let pack_found = buf.windows(4).any(|window| window == b"PACK"); |
| 943 | assert!(pack_found, "response should contain PACK signature"); |
| 944 | } |
| 945 | |
| 946 | /// Regression: a want may be an annotated tag *object* (protocol-v2 |
| 947 | /// `ls-refs` advertises unpeeled ids). The walk used to choke on it |
| 948 | /// mid-stream, killing `git clone` of any repo with an annotated tag; |
| 949 | /// now the tag is peeled and the tag object itself joins the pack. |
| 950 | #[tokio::test] |
| 951 | async fn tag_object_wants_are_peeled_and_packed() { |
| 952 | let dir = TempDir::new().unwrap(); |
| 953 | let repo_path = create_repo_with_commit(dir.path()); |
| 954 | |
| 955 | let clone_path = dir.path().join("workdir"); |
| 956 | let out = Command::new("git") |
| 957 | .args(["tag", "-a", "-m", "annotated", "v1"]) |
| 958 | .current_dir(&clone_path) |
| 959 | .env("GIT_AUTHOR_NAME", "t") |
| 960 | .env("GIT_AUTHOR_EMAIL", "t@t") |
| 961 | .env("GIT_COMMITTER_NAME", "t") |
| 962 | .env("GIT_COMMITTER_EMAIL", "t@t") |
| 963 | .output() |
| 964 | .unwrap(); |
| 965 | assert!(out.status.success()); |
| 966 | let out = Command::new("git") |
| 967 | .args(["push", "-q", "origin", "v1"]) |
| 968 | .current_dir(&clone_path) |
| 969 | .output() |
| 970 | .unwrap(); |
| 971 | assert!(out.status.success()); |
| 972 | |
| 973 | let repo = gix::open(&repo_path).unwrap(); |
| 974 | let tag_oid = repo |
| 975 | .find_reference("refs/tags/v1") |
| 976 | .unwrap() |
| 977 | .try_id() |
| 978 | .unwrap() |
| 979 | .detach(); |
| 980 | assert_eq!( |
| 981 | repo.find_object(tag_oid).unwrap().kind, |
| 982 | gix::object::Kind::Tag, |
| 983 | "fixture must produce a real tag object" |
| 984 | ); |
| 985 | drop(repo); |
| 986 | |
| 987 | let pack = build_raw_pack(&repo_path, &[tag_oid], &[]).unwrap(); |
| 988 | assert!(pack.starts_with(b"PACK")); |
| 989 | let count = u32::from_be_bytes(pack[8..12].try_into().unwrap()); |
| 990 | // tag + commit + tree + blob |
| 991 | assert_eq!(count, 4, "tag object and full closure packed"); |
| 992 | |
| 993 | // The streaming (fetch-serving) path accepts the same want. |
| 994 | let request = UploadPackRequest { |
| 995 | wants: vec![tag_oid], |
| 996 | haves: vec![], |
| 997 | done: true, |
| 998 | capabilities: UploadPackCapabilities::default(), |
| 999 | shallow: ShallowRequest::default(), |
| 1000 | object_ids: None, |
| 1001 | }; |
| 1002 | let mut reader = generate_pack(&repo_path, &request).unwrap(); |
| 1003 | let mut buf = Vec::new(); |
| 1004 | reader.read_to_end(&mut buf).await.unwrap(); |
| 1005 | assert!( |
| 1006 | buf.windows(4).any(|w| w == b"PACK"), |
| 1007 | "fetch response should contain a pack, got: {:?}", |
| 1008 | String::from_utf8_lossy(&buf[..buf.len().min(200)]) |
| 1009 | ); |
| 1010 | } |
| 1011 | } |