This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knot2 / crates / knot-lfs / src / transfer.rs
41 kB 1175 lines
1use std::io::{Read, Write}; 2 3use gix_packetline::PacketLineRef; 4use gix_packetline::blocking_io::{StreamingPeekableIter, encode}; 5use knot_messages::{ 6 AlgorithmKey, CommandKey, DeclaredComputedKey, DeclaredLimitKey, DeclaredReceivedKey, 7 DetailKey, FreeFloorKey, LfsMessages, OidKey, ValueKey, VersionKey, WhatLimitKey, 8}; 9use knot_types::{HttpStatus, RepoDid}; 10 11use crate::store::for_each_chunk; 12use crate::{ 13 BatchObject, ClaimedSize, LfsError, LfsOid, LfsSize, LfsStore, MAX_BATCH_OBJECTS, 14 UploadAdmission, 15}; 16 17pub const CAPABILITY_VERSION: &str = "version=1"; 18const PKT_DATA_MAX: usize = 65516; 19const MAX_MESSAGE_ARGS: usize = 64; 20 21#[derive(Debug, Clone, Copy, PartialEq, Eq)] 22pub enum TransferOp { 23 Upload, 24 Download, 25} 26 27impl TransferOp { 28 pub fn parse(token: &str) -> Option<Self> { 29 match token { 30 "upload" => Some(Self::Upload), 31 "download" => Some(Self::Download), 32 _ => None, 33 } 34 } 35} 36 37#[allow(clippy::too_many_arguments)] 38pub fn serve_transfer( 39 store: &dyn LfsStore, 40 admission: &dyn UploadAdmission, 41 repo: &RepoDid, 42 op: TransferOp, 43 messages: &LfsMessages, 44 input: impl Read, 45 mut output: impl Write, 46) -> Result<(), LfsError> { 47 write_text(&mut output, CAPABILITY_VERSION)?; 48 write_flush(&mut output)?; 49 let mut session = Session { 50 store, 51 admission, 52 repo, 53 op, 54 messages, 55 lines: StreamingPeekableIter::new(input, &[], false), 56 out: output, 57 }; 58 let mut done = false; 59 std::iter::from_fn(|| { 60 (!done).then(|| { 61 session.step().map(|flow| match flow { 62 Flow::Continue => (), 63 Flow::Quit => done = true, 64 }) 65 }) 66 }) 67 .try_for_each(std::convert::identity) 68} 69 70enum Flow { 71 Continue, 72 Quit, 73} 74 75enum Pkt { 76 Data(Vec<u8>), 77 Flush, 78 Delim, 79 Eof, 80} 81 82enum Ended { 83 Flush, 84 Delim, 85} 86 87fn next_pkt<R: Read>(lines: &mut StreamingPeekableIter<R>) -> Result<Pkt, LfsError> { 88 match lines.read_line() { 89 None => Ok(Pkt::Eof), 90 Some(Err(source)) if source.kind() == std::io::ErrorKind::UnexpectedEof => Ok(Pkt::Eof), 91 Some(Err(source)) => Err(LfsError::Channel { source }), 92 Some(Ok(Err(fault))) => Err(LfsError::Framing { 93 detail: fault.to_string(), 94 }), 95 Some(Ok(Ok(PacketLineRef::Data(payload)))) => Ok(Pkt::Data(payload.to_vec())), 96 Some(Ok(Ok(PacketLineRef::Flush))) => Ok(Pkt::Flush), 97 Some(Ok(Ok(PacketLineRef::Delimiter))) => Ok(Pkt::Delim), 98 Some(Ok(Ok(PacketLineRef::ResponseEnd))) => Err(LfsError::Framing { 99 detail: "unexpected response-end packet".to_string(), 100 }), 101 } 102} 103 104fn text_of(payload: Vec<u8>) -> Result<String, LfsError> { 105 String::from_utf8(payload) 106 .map(|line| line.trim_end_matches('\n').to_string()) 107 .map_err(|_| LfsError::Framing { 108 detail: "non-utf8 text packet".to_string(), 109 }) 110} 111 112fn fault_text(fault: &LfsError, messages: &LfsMessages) -> String { 113 match fault { 114 LfsError::InvalidOid { value } => messages 115 .invalid_oid 116 .line(|ValueKey::Value| format!("{value:?}")), 117 LfsError::HashMismatch { declared, computed } => { 118 messages.hash_mismatch.line(|key| match key { 119 DeclaredComputedKey::Declared => declared.to_string(), 120 DeclaredComputedKey::Computed => computed.to_string(), 121 }) 122 } 123 LfsError::SizeMismatch { declared, received } => { 124 messages.size_mismatch.line(|key| match key { 125 DeclaredReceivedKey::Declared => declared.to_string(), 126 DeclaredReceivedKey::Received => received.to_string(), 127 }) 128 } 129 LfsError::SizeLimitExceeded { declared, limit } => { 130 messages.size_limit_exceeded.line(|key| match key { 131 DeclaredLimitKey::Declared => declared.to_string(), 132 DeclaredLimitKey::Limit => limit.to_string(), 133 }) 134 } 135 LfsError::FreeSpaceDenied { free, floor } => { 136 messages.free_space_denied.line(|key| match key { 137 FreeFloorKey::Free => free.to_string(), 138 FreeFloorKey::Floor => floor.to_string(), 139 }) 140 } 141 LfsError::NotFound { oid } => messages.not_found.line(|OidKey::Oid| oid.to_string()), 142 LfsError::Framing { detail } => messages.framing.line(|DetailKey::Detail| detail.clone()), 143 LfsError::TooMany { what, limit } => messages.too_many.line(|key| match key { 144 WhatLimitKey::What => what.to_string(), 145 WhatLimitKey::Limit => limit.to_string(), 146 }), 147 other => other.to_string(), 148 } 149} 150 151fn status_of(fault: &LfsError) -> HttpStatus { 152 HttpStatus::new(match fault { 153 LfsError::NotFound { .. } => 404, 154 LfsError::InvalidOid { .. } 155 | LfsError::HashMismatch { .. } 156 | LfsError::SizeMismatch { .. } 157 | LfsError::Framing { .. } => 400, 158 LfsError::SizeLimitExceeded { .. } => 413, 159 LfsError::FreeSpaceDenied { .. } => 429, 160 _ => 500, 161 }) 162} 163 164fn write_text(out: &mut impl Write, line: &str) -> Result<(), LfsError> { 165 encode::data_to_write(format!("{line}\n").as_bytes(), &mut *out) 166 .map(|_| ()) 167 .map_err(|source| LfsError::Channel { source }) 168} 169 170fn write_flush(out: &mut impl Write) -> Result<(), LfsError> { 171 encode::flush_to_write(&mut *out) 172 .and_then(|_| out.flush().map(|()| 0)) 173 .map(|_| ()) 174 .map_err(|source| LfsError::Channel { source }) 175} 176 177fn write_delim(out: &mut impl Write) -> Result<(), LfsError> { 178 encode::delim_to_write(&mut *out) 179 .map(|_| ()) 180 .map_err(|source| LfsError::Channel { source }) 181} 182 183fn arg_value<'a>(args: &'a [String], key: &str) -> Option<&'a str> { 184 args.iter().find_map(|arg| { 185 arg.strip_prefix(key) 186 .and_then(|rest| rest.strip_prefix('=')) 187 }) 188} 189 190struct Session<'a, R: Read, W: Write> { 191 store: &'a dyn LfsStore, 192 admission: &'a dyn UploadAdmission, 193 repo: &'a RepoDid, 194 op: TransferOp, 195 messages: &'a LfsMessages, 196 lines: StreamingPeekableIter<R>, 197 out: W, 198} 199 200impl<R: Read, W: Write> Session<'_, R, W> { 201 fn step(&mut self) -> Result<Flow, LfsError> { 202 match next_pkt(&mut self.lines)? { 203 Pkt::Eof => Ok(Flow::Quit), 204 Pkt::Flush | Pkt::Delim => Err(LfsError::Framing { 205 detail: "expected a command packet".to_string(), 206 }), 207 Pkt::Data(payload) => self.dispatch(text_of(payload)?), 208 } 209 } 210 211 fn dispatch(&mut self, command: String) -> Result<Flow, LfsError> { 212 let (verb, rest) = command 213 .split_once(' ') 214 .map_or((command.as_str(), ""), |(verb, rest)| (verb, rest)); 215 match verb { 216 "version" => self.handle_version(rest), 217 "batch" => self.handle_batch(), 218 "put-object" => self.handle_put(rest), 219 "verify-object" => self.handle_verify(rest), 220 "get-object" => self.handle_get(rest), 221 "quit" => { 222 self.drain_message()?; 223 self.respond_ok(&[])?; 224 Ok(Flow::Quit) 225 } 226 _ => { 227 self.drain_message()?; 228 let message = self 229 .messages 230 .unknown_command 231 .line(|CommandKey::Command| format!("{verb:?}")); 232 self.respond_error(HttpStatus::new(400), &message)?; 233 Ok(Flow::Continue) 234 } 235 } 236 } 237 238 fn read_args(&mut self) -> Result<(Vec<String>, Ended), LfsError> { 239 let mut args = Vec::new(); 240 std::iter::from_fn(|| Some(next_pkt(&mut self.lines))) 241 .find_map(|pkt| match pkt { 242 Err(fault) => Some(Err(fault)), 243 Ok(Pkt::Data(payload)) => match text_of(payload) { 244 Ok(_) if args.len() >= MAX_MESSAGE_ARGS => Some(Err(LfsError::TooMany { 245 what: "arguments", 246 limit: MAX_MESSAGE_ARGS, 247 })), 248 Ok(line) => { 249 args.push(line); 250 None 251 } 252 Err(fault) => Some(Err(fault)), 253 }, 254 Ok(Pkt::Flush) => Some(Ok(Ended::Flush)), 255 Ok(Pkt::Delim) => Some(Ok(Ended::Delim)), 256 Ok(Pkt::Eof) => Some(Err(LfsError::Framing { 257 detail: "message truncated before flush".to_string(), 258 })), 259 }) 260 .expect("an endless packet iterator always yields a terminator") 261 .map(|ended| (args, ended)) 262 } 263 264 fn drain_budget(&self) -> LfsSize { 265 self.admission.max_object() 266 } 267 268 fn drain_to_flush(&mut self, limit: LfsSize) -> Result<(), LfsError> { 269 let mut discarded = LfsSize::new(0); 270 std::iter::from_fn(|| Some(next_pkt(&mut self.lines))) 271 .find_map(|pkt| match pkt { 272 Err(fault) => Some(Err(fault)), 273 Ok(Pkt::Flush) => Some(Ok(())), 274 Ok(Pkt::Eof) => Some(Err(LfsError::Framing { 275 detail: "message truncated before flush".to_string(), 276 })), 277 Ok(Pkt::Data(payload)) => { 278 discarded = discarded.saturating_add(LfsSize::new(payload.len() as u64)); 279 match discarded > limit { 280 true => Some(Err(LfsError::Framing { 281 detail: "message body exceeds the drain bound".to_string(), 282 })), 283 false => None, 284 } 285 } 286 Ok(Pkt::Delim) => None, 287 }) 288 .expect("an endless packet iterator always yields a terminator") 289 } 290 291 fn drain_message(&mut self) -> Result<(), LfsError> { 292 match self.read_args()?.1 { 293 Ended::Flush => Ok(()), 294 Ended::Delim => self.drain_to_flush(self.drain_budget()), 295 } 296 } 297 298 fn respond_ok(&mut self, args: &[String]) -> Result<(), LfsError> { 299 write_text(&mut self.out, "status 200")?; 300 args.iter() 301 .try_for_each(|arg| write_text(&mut self.out, arg))?; 302 write_flush(&mut self.out) 303 } 304 305 fn respond_error(&mut self, code: HttpStatus, message: &str) -> Result<(), LfsError> { 306 write_text(&mut self.out, &format!("status {:03}", code.get()))?; 307 write_delim(&mut self.out)?; 308 write_text(&mut self.out, &format!("error: {message}"))?; 309 write_flush(&mut self.out) 310 } 311 312 fn respond_fault(&mut self, fault: &LfsError) -> Result<(), LfsError> { 313 let message = match fault { 314 LfsError::Io { .. } => { 315 tracing::warn!(repo = self.repo.as_str(), %fault, "lfs store fault"); 316 "internal storage fault".to_string() 317 } 318 other => fault_text(other, self.messages), 319 }; 320 self.respond_error(status_of(fault), &message) 321 } 322 323 fn handle_version(&mut self, rest: &str) -> Result<Flow, LfsError> { 324 self.drain_message()?; 325 match rest.trim() { 326 "1" => self.respond_ok(&[])?, 327 other => { 328 let message = self 329 .messages 330 .unsupported_version 331 .line(|VersionKey::Version| format!("{other:?}")); 332 self.respond_error(HttpStatus::new(400), &message)?; 333 } 334 } 335 Ok(Flow::Continue) 336 } 337 338 fn handle_batch(&mut self) -> Result<Flow, LfsError> { 339 let (args, ended) = self.read_args()?; 340 if let Some(algo) = arg_value(&args, "hash-algo") 341 && algo != crate::HASH_ALGO 342 { 343 if matches!(ended, Ended::Delim) { 344 self.drain_to_flush(self.drain_budget())?; 345 } 346 let message = self 347 .messages 348 .unsupported_hash 349 .line(|AlgorithmKey::Algorithm| format!("{algo:?}")); 350 self.respond_error(HttpStatus::new(400), &message)?; 351 return Ok(Flow::Continue); 352 } 353 let items = match ended { 354 Ended::Flush => Ok(Vec::new()), 355 Ended::Delim => self.read_batch_items(), 356 }; 357 let items = match items { 358 Ok(items) => items, 359 Err(fault @ (LfsError::InvalidOid { .. } | LfsError::Framing { .. })) => { 360 self.drain_to_flush(self.drain_budget())?; 361 self.respond_fault(&fault)?; 362 return Ok(Flow::Continue); 363 } 364 Err(fault) => return Err(fault), 365 }; 366 let lines: Result<Vec<String>, LfsError> = 367 items.iter().map(|item| self.batch_line(item)).collect(); 368 match lines { 369 Ok(lines) => { 370 write_text(&mut self.out, "status 200")?; 371 write_text(&mut self.out, &format!("hash-algo={}", crate::HASH_ALGO))?; 372 write_delim(&mut self.out)?; 373 lines 374 .iter() 375 .try_for_each(|line| write_text(&mut self.out, line))?; 376 write_flush(&mut self.out)?; 377 } 378 Err(fault) => self.respond_fault(&fault)?, 379 } 380 Ok(Flow::Continue) 381 } 382 383 fn read_batch_items(&mut self) -> Result<Vec<BatchObject>, LfsError> { 384 let mut items = Vec::new(); 385 std::iter::from_fn(|| Some(next_pkt(&mut self.lines))) 386 .find_map(|pkt| match pkt { 387 Err(fault) => Some(Err(fault)), 388 Ok(Pkt::Data(payload)) => match text_of(payload).and_then(parse_batch_item) { 389 Ok(_) if items.len() >= MAX_BATCH_OBJECTS => Some(Err(LfsError::TooMany { 390 what: "batch items", 391 limit: MAX_BATCH_OBJECTS, 392 })), 393 Ok(item) => { 394 items.push(item); 395 None 396 } 397 Err(fault) => Some(Err(fault)), 398 }, 399 Ok(Pkt::Flush) => Some(Ok(())), 400 Ok(Pkt::Delim | Pkt::Eof) => Some(Err(LfsError::Framing { 401 detail: "batch items truncated before flush".to_string(), 402 })), 403 }) 404 .expect("an endless packet iterator always yields a terminator") 405 .map(|()| items) 406 } 407 408 fn batch_line(&self, item: &BatchObject) -> Result<String, LfsError> { 409 let stored = match self.op { 410 TransferOp::Upload => self.store.touch(self.repo, &item.oid)?, 411 TransferOp::Download => self.store.probe(self.repo, &item.oid)?, 412 }; 413 let (size, action) = match (self.op, stored) { 414 (TransferOp::Upload, Some(_)) => (item.size.get(), "noop"), 415 (TransferOp::Upload, None) => (item.size.get(), "upload"), 416 (TransferOp::Download, Some(actual)) => (actual.get(), "download"), 417 (TransferOp::Download, None) => (item.size.get(), "download"), 418 }; 419 Ok(format!("{} {} {}", item.oid, size, action)) 420 } 421 422 fn handle_put(&mut self, rest: &str) -> Result<Flow, LfsError> { 423 if self.op != TransferOp::Upload { 424 self.drain_message()?; 425 self.respond_error(HttpStatus::new(403), &self.messages.put_on_download.text())?; 426 return Ok(Flow::Continue); 427 } 428 let (args, ended) = self.read_args()?; 429 let checked = LfsOid::new(rest).and_then(|oid| { 430 let declared = arg_value(&args, "size") 431 .and_then(|value| value.parse().ok()) 432 .map(ClaimedSize::new) 433 .ok_or_else(|| LfsError::Framing { 434 detail: "put-object requires a size argument".to_string(), 435 })?; 436 let permit = self.admission.admit(declared)?; 437 Ok((oid, declared, permit)) 438 }); 439 let (oid, declared, permit) = match checked { 440 Ok(admitted) => admitted, 441 Err(fault) => { 442 if matches!(ended, Ended::Delim) { 443 self.drain_to_flush(self.drain_budget())?; 444 } 445 self.respond_fault(&fault)?; 446 return Ok(Flow::Continue); 447 } 448 }; 449 if matches!(ended, Ended::Flush) { 450 self.respond_error(HttpStatus::new(400), &self.messages.put_no_body.text())?; 451 return Ok(Flow::Continue); 452 } 453 let mut body = PktBody { 454 lines: &mut self.lines, 455 buffer: Vec::new(), 456 offset: 0, 457 done: false, 458 }; 459 let stored = self.store.put(self.repo, &oid, declared, &mut body); 460 drop(permit); 461 let synced = body.done; 462 if !synced { 463 self.drain_to_flush(self.drain_budget())?; 464 } 465 match stored { 466 Ok(()) => { 467 tracing::info!( 468 repo = self.repo.as_str(), 469 oid = oid.as_str(), 470 size = declared.get(), 471 "lfs object received over ssh" 472 ); 473 self.respond_ok(&[])? 474 } 475 Err(fault) => self.respond_fault(&fault)?, 476 } 477 Ok(Flow::Continue) 478 } 479 480 fn handle_verify(&mut self, rest: &str) -> Result<Flow, LfsError> { 481 if self.op != TransferOp::Upload { 482 self.drain_message()?; 483 self.respond_error( 484 HttpStatus::new(403), 485 &self.messages.verify_on_download.text(), 486 )?; 487 return Ok(Flow::Continue); 488 } 489 let (args, ended) = self.read_args()?; 490 if matches!(ended, Ended::Delim) { 491 self.drain_to_flush(self.drain_budget())?; 492 } 493 let declared = arg_value(&args, "size") 494 .and_then(|value| value.parse().ok()) 495 .map(ClaimedSize::new); 496 let verdict = LfsOid::new(rest).and_then(|oid| { 497 match (self.store.touch(self.repo, &oid)?, declared) { 498 (None, _) => Err(LfsError::NotFound { oid }), 499 (Some(actual), Some(declared)) if !declared.matches(actual) => { 500 Err(LfsError::SizeMismatch { 501 declared, 502 received: actual, 503 }) 504 } 505 (Some(_), _) => Ok(()), 506 } 507 }); 508 match verdict { 509 Ok(()) => self.respond_ok(&[])?, 510 Err(fault) => self.respond_fault(&fault)?, 511 } 512 Ok(Flow::Continue) 513 } 514 515 fn handle_get(&mut self, rest: &str) -> Result<Flow, LfsError> { 516 if self.op != TransferOp::Download { 517 self.drain_message()?; 518 self.respond_error(HttpStatus::new(403), &self.messages.get_on_upload.text())?; 519 return Ok(Flow::Continue); 520 } 521 self.drain_message()?; 522 let opened = LfsOid::new(rest).and_then(|oid| { 523 let size = self 524 .store 525 .probe(self.repo, &oid)? 526 .ok_or(LfsError::NotFound { oid: oid.clone() })?; 527 let body = self.store.read(self.repo, &oid)?; 528 Ok((size, body)) 529 }); 530 let (size, mut body) = match opened { 531 Ok(found) => found, 532 Err(fault) => { 533 self.respond_fault(&fault)?; 534 return Ok(Flow::Continue); 535 } 536 }; 537 write_text(&mut self.out, "status 200")?; 538 write_text(&mut self.out, &format!("size={size}"))?; 539 write_delim(&mut self.out)?; 540 let out = &mut self.out; 541 for_each_chunk(&mut body, |chunk| { 542 chunk.chunks(PKT_DATA_MAX).try_for_each(|piece| { 543 encode::data_to_write(piece, &mut *out) 544 .map(|_| ()) 545 .map_err(|source| LfsError::Channel { source }) 546 }) 547 })?; 548 write_flush(&mut self.out)?; 549 tracing::info!( 550 repo = self.repo.as_str(), 551 oid = rest, 552 size = size.get(), 553 "lfs object served over ssh" 554 ); 555 Ok(Flow::Continue) 556 } 557} 558 559fn parse_batch_item(line: String) -> Result<BatchObject, LfsError> { 560 let mut tokens = line.split(' '); 561 let oid = LfsOid::new(tokens.next().unwrap_or_default())?; 562 let size = tokens 563 .next() 564 .and_then(|token| token.parse().ok()) 565 .map(ClaimedSize::new) 566 .ok_or_else(|| LfsError::Framing { 567 detail: format!("malformed batch item {line:?}"), 568 })?; 569 Ok(BatchObject { oid, size }) 570} 571 572struct PktBody<'a, R: Read> { 573 lines: &'a mut StreamingPeekableIter<R>, 574 buffer: Vec<u8>, 575 offset: usize, 576 done: bool, 577} 578 579impl<R: Read> Read for PktBody<'_, R> { 580 fn read(&mut self, out: &mut [u8]) -> std::io::Result<usize> { 581 if self.done { 582 return Ok(0); 583 } 584 if self.offset >= self.buffer.len() { 585 match next_pkt(self.lines).map_err(std::io::Error::other)? { 586 Pkt::Data(payload) => { 587 self.buffer = payload; 588 self.offset = 0; 589 } 590 Pkt::Flush => { 591 self.done = true; 592 return Ok(0); 593 } 594 Pkt::Delim => { 595 return Err(std::io::Error::other("unexpected delimiter in object body")); 596 } 597 Pkt::Eof => return Err(std::io::Error::other("object body truncated")), 598 } 599 } 600 let take = out.len().min(self.buffer.len() - self.offset); 601 out[..take].copy_from_slice(&self.buffer[self.offset..self.offset + take]); 602 self.offset += take; 603 Ok(take) 604 } 605} 606 607#[cfg(test)] 608mod tests { 609 use sha2::{Digest, Sha256}; 610 611 use super::*; 612 use crate::admission::Unbounded; 613 use crate::{FreeSpaceFloor, MemoryStore}; 614 615 fn oid_of(bytes: &[u8]) -> LfsOid { 616 LfsOid::from_digest(Sha256::digest(bytes).into()) 617 } 618 619 fn repo() -> RepoDid { 620 RepoDid::new("did:plc:squid").unwrap() 621 } 622 623 fn put_text(buf: &mut Vec<u8>, line: &str) { 624 encode::data_to_write(format!("{line}\n").as_bytes(), &mut *buf).unwrap(); 625 } 626 627 fn msg(buf: &mut Vec<u8>, command: &str, args: &[&str]) { 628 put_text(buf, command); 629 args.iter().for_each(|arg| put_text(buf, arg)); 630 encode::flush_to_write(&mut *buf).unwrap(); 631 } 632 633 fn msg_lines(buf: &mut Vec<u8>, command: &str, args: &[&str], lines: &[String]) { 634 put_text(buf, command); 635 args.iter().for_each(|arg| put_text(buf, arg)); 636 encode::delim_to_write(&mut *buf).unwrap(); 637 lines.iter().for_each(|line| put_text(buf, line)); 638 encode::flush_to_write(&mut *buf).unwrap(); 639 } 640 641 fn msg_data(buf: &mut Vec<u8>, command: &str, args: &[&str], data: &[u8]) { 642 put_text(buf, command); 643 args.iter().for_each(|arg| put_text(buf, arg)); 644 encode::delim_to_write(&mut *buf).unwrap(); 645 data.chunks(PKT_DATA_MAX).for_each(|chunk| { 646 encode::data_to_write(chunk, &mut *buf).unwrap(); 647 }); 648 encode::flush_to_write(&mut *buf).unwrap(); 649 } 650 651 #[derive(Debug, PartialEq, Eq, Clone)] 652 enum Out { 653 Line(String), 654 Bin(Vec<u8>), 655 Delim, 656 Flush, 657 } 658 659 fn parse_out(bytes: &[u8]) -> Vec<Out> { 660 let mut lines = StreamingPeekableIter::new(bytes, &[], false); 661 std::iter::from_fn(|| { 662 lines.read_line().and_then(|pkt| { 663 if matches!(&pkt, Err(fault) if fault.kind() == std::io::ErrorKind::UnexpectedEof) { 664 return None; 665 } 666 Some( 667 match pkt.expect("readable output").expect("well-formed output") { 668 PacketLineRef::Data(payload) => std::str::from_utf8(payload) 669 .ok() 670 .filter(|text| text.ends_with('\n')) 671 .map(|text| Out::Line(text.trim_end_matches('\n').to_string())) 672 .unwrap_or_else(|| Out::Bin(payload.to_vec())), 673 PacketLineRef::Flush => Out::Flush, 674 PacketLineRef::Delimiter => Out::Delim, 675 PacketLineRef::ResponseEnd => panic!("server never writes response-end"), 676 }, 677 ) 678 }) 679 }) 680 .collect() 681 } 682 683 fn line(text: &str) -> Out { 684 Out::Line(text.to_string()) 685 } 686 687 fn run(op: TransferOp, store: &dyn LfsStore, script: &[u8]) -> Vec<Out> { 688 let mut output = Vec::new(); 689 serve_transfer( 690 store, 691 &Unbounded, 692 &repo(), 693 op, 694 &knot_messages::default_catalog().lfs, 695 script, 696 &mut output, 697 ) 698 .unwrap(); 699 parse_out(&output) 700 } 701 702 fn responses(out: &[Out]) -> Vec<Vec<Out>> { 703 out.split_inclusive(|item| *item == Out::Flush) 704 .map(<[Out]>::to_vec) 705 .collect() 706 } 707 708 #[test] 709 fn an_upload_session_lands_and_verifies_an_object() { 710 let store = MemoryStore::new(); 711 let seeded: &[u8] = b"already on the server"; 712 let seeded_oid = oid_of(seeded); 713 store 714 .put( 715 &repo(), 716 &seeded_oid, 717 ClaimedSize::new(seeded.len() as u64), 718 &mut &seeded[..], 719 ) 720 .unwrap(); 721 722 let fresh: &[u8] = b"\xffnew binary media\x00with raw bytes"; 723 let fresh_oid = oid_of(fresh); 724 let fresh_len = fresh.len(); 725 726 let mut script = Vec::new(); 727 msg(&mut script, "version 1", &[]); 728 msg_lines( 729 &mut script, 730 "batch", 731 &[ 732 "transfer=ssh", 733 "hash-algo=sha256", 734 "refname=refs/heads/main", 735 ], 736 &[ 737 format!("{seeded_oid} {} extra=ignored", seeded.len()), 738 format!("{fresh_oid} {fresh_len}"), 739 ], 740 ); 741 msg_data( 742 &mut script, 743 &format!("put-object {fresh_oid}"), 744 &[&format!("size={fresh_len}")], 745 fresh, 746 ); 747 msg( 748 &mut script, 749 &format!("verify-object {fresh_oid}"), 750 &[&format!("size={fresh_len}")], 751 ); 752 msg(&mut script, "quit", &[]); 753 754 let out = run(TransferOp::Upload, &store, &script); 755 let turns = responses(&out); 756 assert_eq!(turns[0], vec![line(CAPABILITY_VERSION), Out::Flush]); 757 assert_eq!(turns[1], vec![line("status 200"), Out::Flush]); 758 assert_eq!( 759 turns[2], 760 vec![ 761 line("status 200"), 762 line("hash-algo=sha256"), 763 Out::Delim, 764 line(&format!("{seeded_oid} {} noop", seeded.len())), 765 line(&format!("{fresh_oid} {fresh_len} upload")), 766 Out::Flush, 767 ] 768 ); 769 assert_eq!(turns[3], vec![line("status 200"), Out::Flush]); 770 assert_eq!(turns[4], vec![line("status 200"), Out::Flush]); 771 assert_eq!(turns[5], vec![line("status 200"), Out::Flush]); 772 assert_eq!( 773 store.probe(&repo(), &fresh_oid).unwrap(), 774 Some(LfsSize::new(fresh_len as u64)) 775 ); 776 } 777 778 #[test] 779 fn a_download_session_streams_bytes_and_404s_the_missing() { 780 let store = MemoryStore::new(); 781 let media: &[u8] = b"\xff\x00streamable media"; 782 let media_oid = oid_of(media); 783 let absent = oid_of(b"never uploaded"); 784 store 785 .put( 786 &repo(), 787 &media_oid, 788 ClaimedSize::new(media.len() as u64), 789 &mut &media[..], 790 ) 791 .unwrap(); 792 793 let mut script = Vec::new(); 794 msg_lines( 795 &mut script, 796 "batch", 797 &["transfer=ssh", "hash-algo=sha256"], 798 &[format!("{media_oid} 1"), format!("{absent} 9")], 799 ); 800 msg(&mut script, &format!("get-object {media_oid}"), &[]); 801 msg(&mut script, &format!("get-object {absent}"), &[]); 802 msg(&mut script, "quit", &[]); 803 804 let out = run(TransferOp::Download, &store, &script); 805 let turns = responses(&out); 806 assert_eq!( 807 turns[1], 808 vec![ 809 line("status 200"), 810 line("hash-algo=sha256"), 811 Out::Delim, 812 line(&format!("{media_oid} {} download", media.len())), 813 line(&format!("{absent} 9 download")), 814 Out::Flush, 815 ] 816 ); 817 assert_eq!( 818 turns[2], 819 vec![ 820 line("status 200"), 821 line(&format!("size={}", media.len())), 822 Out::Delim, 823 Out::Bin(media.to_vec()), 824 Out::Flush, 825 ] 826 ); 827 assert_eq!(turns[3][0], line("status 404")); 828 assert_eq!(turns[4], vec![line("status 200"), Out::Flush]); 829 } 830 831 #[test] 832 fn the_channel_mode_gates_every_write_and_read_verb() { 833 let store = MemoryStore::new(); 834 let oid = oid_of(b"whatever"); 835 836 let mut download_script = Vec::new(); 837 msg_data( 838 &mut download_script, 839 &format!("put-object {oid}"), 840 &["size=3"], 841 b"abc", 842 ); 843 msg( 844 &mut download_script, 845 &format!("verify-object {oid}"), 846 &["size=3"], 847 ); 848 let out = run(TransferOp::Download, &store, &download_script); 849 let turns = responses(&out); 850 assert_eq!(turns[1][0], line("status 403")); 851 assert_eq!(turns[2][0], line("status 403")); 852 assert_eq!(store.probe(&repo(), &oid).unwrap(), None); 853 854 let mut upload_script = Vec::new(); 855 msg(&mut upload_script, &format!("get-object {oid}"), &[]); 856 let out = run(TransferOp::Upload, &store, &upload_script); 857 assert_eq!(responses(&out)[1][0], line("status 403")); 858 } 859 860 #[test] 861 fn admission_rejections_surface_as_typed_statuses() { 862 struct Deny(LfsSize); 863 impl UploadAdmission for Deny { 864 fn admit(&self, declared: ClaimedSize) -> Result<crate::UploadPermit, LfsError> { 865 match declared.get() > self.0.get() { 866 true => Err(LfsError::SizeLimitExceeded { 867 declared, 868 limit: self.0, 869 }), 870 false => Err(LfsError::FreeSpaceDenied { 871 free: LfsSize::new(1), 872 floor: FreeSpaceFloor::new(2), 873 }), 874 } 875 } 876 877 fn max_object(&self) -> LfsSize { 878 self.0 879 } 880 } 881 let store = MemoryStore::new(); 882 let body: &[u8] = b"denied"; 883 let oid = oid_of(body); 884 let mut script = Vec::new(); 885 msg_data( 886 &mut script, 887 &format!("put-object {oid}"), 888 &["size=999"], 889 body, 890 ); 891 msg_data(&mut script, &format!("put-object {oid}"), &["size=6"], body); 892 let mut output = Vec::new(); 893 serve_transfer( 894 &store, 895 &Deny(LfsSize::new(10)), 896 &repo(), 897 TransferOp::Upload, 898 &knot_messages::default_catalog().lfs, 899 &script[..], 900 &mut output, 901 ) 902 .unwrap(); 903 let turns = responses(&parse_out(&output)); 904 assert_eq!(turns[1][0], line("status 413")); 905 assert_eq!(turns[2][0], line("status 429")); 906 assert_eq!(store.probe(&repo(), &oid).unwrap(), None); 907 } 908 909 #[test] 910 fn tampered_and_truncated_uploads_fail_closed_and_the_session_survives() { 911 let store = MemoryStore::new(); 912 let body: &[u8] = b"the true bytes"; 913 let liar = oid_of(b"some other bytes"); 914 let valid = oid_of(body); 915 916 let mut script = Vec::new(); 917 msg_data( 918 &mut script, 919 &format!("put-object {liar}"), 920 &[&format!("size={}", body.len())], 921 body, 922 ); 923 msg_data( 924 &mut script, 925 &format!("put-object {valid}"), 926 &[&format!("size={}", body.len() + 5)], 927 body, 928 ); 929 msg( 930 &mut script, 931 &format!("verify-object {valid}"), 932 &[&format!("size={}", body.len())], 933 ); 934 msg_data( 935 &mut script, 936 &format!("put-object {valid}"), 937 &[&format!("size={}", body.len())], 938 body, 939 ); 940 msg(&mut script, "quit", &[]); 941 942 let out = run(TransferOp::Upload, &store, &script); 943 let turns = responses(&out); 944 assert_eq!(turns[1][0], line("status 400"), "hash mismatch"); 945 assert_eq!(turns[2][0], line("status 400"), "size mismatch"); 946 assert_eq!(turns[3][0], line("status 404"), "nothing landed"); 947 assert_eq!( 948 turns[4][0], 949 line("status 200"), 950 "retry with matching bytes succeeds" 951 ); 952 assert_eq!( 953 store.probe(&repo(), &valid).unwrap(), 954 Some(LfsSize::new(body.len() as u64)) 955 ); 956 assert_eq!(store.probe(&repo(), &liar).unwrap(), None); 957 } 958 959 #[test] 960 fn hostile_commands_get_clean_rejections_and_never_a_panic() { 961 let store = MemoryStore::new(); 962 let mut script = Vec::new(); 963 msg(&mut script, "version 9", &[]); 964 msg(&mut script, "steal-the-objects now", &[]); 965 msg_lines( 966 &mut script, 967 "batch", 968 &["hash-algo=sha256"], 969 &["../../../etc/passwd0000000000000000000000000000000000000000 5".to_string()], 970 ); 971 msg_lines( 972 &mut script, 973 "batch", 974 &["hash-algo=sha1"], 975 &[format!("{} 5", oid_of(b"x"))], 976 ); 977 msg(&mut script, "quit", &[]); 978 979 let out = run(TransferOp::Upload, &store, &script); 980 let turns = responses(&out); 981 assert_eq!(turns[1][0], line("status 400"), "unsupported version"); 982 assert_eq!(turns[2][0], line("status 400"), "unknown command"); 983 assert_eq!(turns[3][0], line("status 400"), "traversal oid"); 984 assert_eq!(turns[4][0], line("status 400"), "foreign hash algo"); 985 assert_eq!(turns[5][0], line("status 200"), "quit still answers"); 986 } 987 988 #[test] 989 fn floods_kill_the_session_instead_of_accumulating() { 990 let store = MemoryStore::new(); 991 let flooded_batch: Vec<String> = (0..MAX_BATCH_OBJECTS + 1) 992 .map(|index| format!("{} 1", oid_of(index.to_string().as_bytes()))) 993 .collect(); 994 let mut script = Vec::new(); 995 msg_lines(&mut script, "batch", &["hash-algo=sha256"], &flooded_batch); 996 let mut output = Vec::new(); 997 let verdict = serve_transfer( 998 &store, 999 &Unbounded, 1000 &repo(), 1001 TransferOp::Upload, 1002 &knot_messages::default_catalog().lfs, 1003 &script[..], 1004 &mut output, 1005 ); 1006 assert!(matches!( 1007 verdict, 1008 Err(LfsError::TooMany { 1009 what: "batch items", 1010 .. 1011 }) 1012 )); 1013 1014 let flooded_args: Vec<&str> = std::iter::repeat_n("size=1", MAX_MESSAGE_ARGS + 1).collect(); 1015 let mut script = Vec::new(); 1016 msg( 1017 &mut script, 1018 &format!("put-object {}", oid_of(b"flooded")), 1019 &flooded_args, 1020 ); 1021 let mut output = Vec::new(); 1022 let verdict = serve_transfer( 1023 &store, 1024 &Unbounded, 1025 &repo(), 1026 TransferOp::Upload, 1027 &knot_messages::default_catalog().lfs, 1028 &script[..], 1029 &mut output, 1030 ); 1031 assert!(matches!( 1032 verdict, 1033 Err(LfsError::TooMany { 1034 what: "arguments", 1035 .. 1036 }) 1037 )); 1038 } 1039 1040 #[test] 1041 fn malformed_and_oversized_put_streams_end_the_session() { 1042 let store = MemoryStore::new(); 1043 1044 let mut sink = Vec::new(); 1045 let garbage = serve_transfer( 1046 &store, 1047 &Unbounded, 1048 &repo(), 1049 TransferOp::Upload, 1050 &knot_messages::default_catalog().lfs, 1051 &b"zzzz not pkt-line at all"[..], 1052 &mut sink, 1053 ); 1054 assert!( 1055 matches!(garbage, Err(LfsError::Framing { .. })), 1056 "raw garbage on the wire is a framing fault, not a hang" 1057 ); 1058 1059 let admission = crate::StoreAdmission::new( 1060 crate::LfsStorePath::new("/"), 1061 LfsSize::new(10), 1062 FreeSpaceFloor::new(0), 1063 ); 1064 let oid = oid_of(b"whatever"); 1065 let mut script = Vec::new(); 1066 put_text(&mut script, &format!("put-object {oid}")); 1067 put_text(&mut script, "size=1"); 1068 encode::delim_to_write(&mut script).unwrap(); 1069 (0..3).for_each(|_| { 1070 encode::data_to_write(&[0u8; 100][..], &mut script).unwrap(); 1071 }); 1072 encode::flush_to_write(&mut script).unwrap(); 1073 let mut output = Vec::new(); 1074 let overshoot = serve_transfer( 1075 &store, 1076 &admission, 1077 &repo(), 1078 TransferOp::Upload, 1079 &knot_messages::default_catalog().lfs, 1080 &script[..], 1081 &mut output, 1082 ); 1083 assert!( 1084 matches!(overshoot, Err(LfsError::Framing { .. })), 1085 "a body overshooting the object size limit ends the session instead of draining unbounded bytes" 1086 ); 1087 assert_eq!(store.probe(&repo(), &oid).unwrap(), None); 1088 } 1089 1090 #[test] 1091 fn store_faults_reach_the_client_without_the_path() { 1092 struct BrokenStore; 1093 impl LfsStore for BrokenStore { 1094 fn put( 1095 &self, 1096 _repo: &RepoDid, 1097 _oid: &LfsOid, 1098 _size: ClaimedSize, 1099 _body: &mut dyn std::io::Read, 1100 ) -> Result<(), LfsError> { 1101 // for now!!!! 1102 unreachable!("this session never puts") 1103 } 1104 1105 fn read( 1106 &self, 1107 _repo: &RepoDid, 1108 _oid: &LfsOid, 1109 ) -> Result<Box<dyn std::io::Read + Send>, LfsError> { 1110 unreachable!("this session never reads") 1111 } 1112 1113 fn probe(&self, _repo: &RepoDid, _oid: &LfsOid) -> Result<Option<LfsSize>, LfsError> { 1114 Err(LfsError::Io { 1115 op: "stat", 1116 path: "/srv/secret-lfs-root/plc/sq/uid".into(), 1117 source: std::io::Error::other("disk fell off"), 1118 }) 1119 } 1120 1121 fn touch(&self, repo: &RepoDid, oid: &LfsOid) -> Result<Option<LfsSize>, LfsError> { 1122 self.probe(repo, oid) 1123 } 1124 } 1125 1126 let mut script = Vec::new(); 1127 msg_lines( 1128 &mut script, 1129 "batch", 1130 &["hash-algo=sha256"], 1131 &[format!("{} 5", oid_of(b"whatever"))], 1132 ); 1133 let out = run(TransferOp::Upload, &BrokenStore, &script); 1134 let turns = responses(&out); 1135 assert_eq!(turns[1][0], line("status 500")); 1136 assert!(turns[1].contains(&line("error: internal storage fault"))); 1137 turns[1].iter().for_each(|item| { 1138 if let Out::Line(text) = item { 1139 assert!(!text.contains("secret-lfs-root"), "leaked path in {text:?}"); 1140 } 1141 }); 1142 } 1143 1144 #[test] 1145 fn a_large_object_round_trips_across_many_packets() { 1146 let store = MemoryStore::new(); 1147 let media: Vec<u8> = (0..500_000u32).map(|n| (n % 251) as u8).collect(); 1148 let oid = oid_of(&media); 1149 1150 let mut script = Vec::new(); 1151 msg_data( 1152 &mut script, 1153 &format!("put-object {oid}"), 1154 &[&format!("size={}", media.len())], 1155 &media, 1156 ); 1157 let out = run(TransferOp::Upload, &store, &script); 1158 assert_eq!(responses(&out)[1][0], line("status 200")); 1159 1160 let mut fetch = Vec::new(); 1161 msg(&mut fetch, &format!("get-object {oid}"), &[]); 1162 let out = run(TransferOp::Download, &store, &fetch); 1163 let body: Vec<u8> = out 1164 .iter() 1165 .skip_while(|item| **item != Out::Delim) 1166 .filter_map(|item| match item { 1167 Out::Bin(chunk) => Some(chunk.clone()), 1168 Out::Line(text) => Some(format!("{text}\n").into_bytes()), 1169 _ => None, 1170 }) 1171 .flatten() 1172 .collect(); 1173 assert_eq!(body, media); 1174 } 1175}