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