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