1use std::collections::{BTreeMap, BTreeSet};
51use std::path::{Path, PathBuf};
52use std::sync::{Arc, Mutex};
53
54use rudb_common::{Error, Result};
55
56use crate::submit::{Completion, Request, Response};
57use crate::{File, Filesystem, OpenMode};
58
59#[derive(Debug, Clone, PartialEq, Eq)]
64pub enum Op {
65 Open {
67 path: PathBuf,
69 mode: OpenMode,
71 },
72 Write {
74 path: PathBuf,
76 offset: u64,
78 len: usize,
80 },
81 Sync {
83 path: PathBuf,
85 },
86 Truncate {
88 path: PathBuf,
90 len: u64,
92 },
93 Rename {
95 from: PathBuf,
97 to: PathBuf,
99 },
100 Remove {
102 path: PathBuf,
104 },
105 CreateDir {
107 path: PathBuf,
109 },
110 SyncDir {
112 path: PathBuf,
114 },
115}
116
117impl Op {
118 #[must_use]
123 pub fn is_durability_point(&self) -> bool {
124 matches!(self, Self::Sync { .. } | Self::SyncDir { .. })
125 }
126}
127
128#[derive(Debug, Clone, PartialEq, Eq)]
130pub enum Crash {
131 LosingUnsynced,
133 KeepingEverything,
136 Keeping(Vec<u64>),
141}
142
143impl Crash {
144 fn keeps(&self, seq: u64) -> bool {
145 match self {
146 Self::LosingUnsynced => false,
147 Self::KeepingEverything => true,
148 Self::Keeping(kept) => kept.contains(&seq),
149 }
150 }
151}
152
153#[derive(Debug, Clone)]
155struct Pending {
156 seq: u64,
157 change: Change,
158}
159
160#[derive(Debug, Clone)]
161enum Change {
162 Write { offset: u64, data: Vec<u8> },
163 Truncate { len: u64 },
164}
165
166#[derive(Debug, Clone, Default)]
167struct SimFile {
168 durable: Vec<u8>,
170 pending: Vec<Pending>,
172}
173
174impl SimFile {
175 fn visible(&self) -> Vec<u8> {
177 let mut bytes = self.durable.clone();
178 for entry in &self.pending {
179 apply(&mut bytes, &entry.change);
180 }
181 bytes
182 }
183}
184
185fn apply(bytes: &mut Vec<u8>, change: &Change) {
186 match change {
187 Change::Write { offset, data } => {
188 let end = *offset as usize + data.len();
189 if bytes.len() < end {
190 bytes.resize(end, 0);
191 }
192 bytes[*offset as usize..end].copy_from_slice(data);
193 }
194 Change::Truncate { len } => bytes.resize(*len as usize, 0),
195 }
196}
197
198#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
204pub enum Completions {
205 #[default]
207 InOrder,
208 Reversed,
211 Shuffled(u64),
213}
214
215#[derive(Debug, Default)]
217struct ReadFaults {
218 served: u64,
222 short: BTreeMap<u64, usize>,
224 failing: BTreeSet<u64>,
226 order: Completions,
228}
229
230#[derive(Debug, Default)]
231struct Inner {
232 files: BTreeMap<PathBuf, SimFile>,
233 dirs: BTreeSet<PathBuf>,
234 log: Vec<Op>,
235 next_seq: u64,
236 fail_at: Option<usize>,
238 reads: ReadFaults,
239}
240
241impl Inner {
242 fn record(&mut self, op: Op) -> Result<()> {
247 let index = self.log.len();
248 self.log.push(op);
249 if self.fail_at == Some(index) {
250 self.fail_at = None;
251 return Err(Error::io(format!("injected failure at operation {index}")));
252 }
253 Ok(())
254 }
255
256 fn serve_read(&mut self, path: &Path, offset: u64, buf: &mut [u8]) -> Result<usize> {
261 let number = self.reads.served;
262 self.reads.served += 1;
263 let failing = self.reads.failing.remove(&number);
264 let short = self.reads.short.remove(&number);
265 if failing {
266 return Err(Error::io(format!("injected read failure on read {number}")));
267 }
268 let file = self
269 .files
270 .get(path)
271 .ok_or_else(|| Error::io(format!("{} was removed while open", path.display())))?;
272 let bytes = file.visible();
273 let start = offset as usize;
274 if start >= bytes.len() {
275 return Ok(0);
276 }
277 let mut n = buf.len().min(bytes.len() - start);
278 if let Some(cap) = short {
279 n = n.min(cap);
280 }
281 buf[..n].copy_from_slice(&bytes[start..start + n]);
282 Ok(n)
283 }
284}
285
286fn reorder<T>(outcomes: &mut [T], order: Completions) {
292 match order {
293 Completions::InOrder => {}
294 Completions::Reversed => outcomes.reverse(),
295 Completions::Shuffled(seed) => {
296 let mut state = seed;
300 let mut next = move || {
301 state = state.wrapping_add(0x9e37_79b9_7f4a_7c15);
302 let mut z = state;
303 z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
304 z = (z ^ (z >> 27)).wrapping_mul(0x94d0_49bb_1331_11eb);
305 z ^ (z >> 31)
306 };
307 for i in (1..outcomes.len()).rev() {
308 let j = (next() % (i as u64 + 1)) as usize;
309 outcomes.swap(i, j);
310 }
311 }
312 }
313}
314
315#[derive(Debug, Clone, Default)]
320pub struct SimFilesystem {
321 inner: Arc<Mutex<Inner>>,
322}
323
324impl SimFilesystem {
325 #[must_use]
327 pub fn new() -> Self {
328 Self::default()
329 }
330
331 fn lock(&self) -> std::sync::MutexGuard<'_, Inner> {
332 self.inner.lock().unwrap_or_else(std::sync::PoisonError::into_inner)
336 }
337
338 #[must_use]
340 pub fn ops(&self) -> Vec<Op> {
341 self.lock().log.clone()
342 }
343
344 #[must_use]
349 pub fn op_count(&self) -> usize {
350 self.lock().log.len()
351 }
352
353 pub fn clear_log(&self) {
358 self.lock().log.clear();
359 }
360
361 pub fn fail_at(&self, index: usize) {
368 self.lock().fail_at = Some(index);
369 }
370
371 pub fn clear_failure(&self) {
373 self.lock().fail_at = None;
374 }
375
376 #[must_use]
382 pub fn reads_served(&self) -> u64 {
383 self.lock().reads.served
384 }
385
386 pub fn short_read_at(&self, read: u64, len: usize) {
392 self.lock().reads.short.insert(read, len);
393 }
394
395 pub fn fail_read_at(&self, read: u64) {
400 self.lock().reads.failing.insert(read);
401 }
402
403 pub fn complete(&self, order: Completions) {
405 self.lock().reads.order = order;
406 }
407
408 pub fn clear_read_faults(&self) {
410 let mut inner = self.lock();
411 inner.reads.short.clear();
412 inner.reads.failing.clear();
413 inner.reads.order = Completions::InOrder;
414 }
415
416 #[must_use]
422 pub fn pending(&self) -> Vec<(u64, PathBuf)> {
423 let inner = self.lock();
424 let mut out: Vec<(u64, PathBuf)> = inner
425 .files
426 .iter()
427 .flat_map(|(path, file)| file.pending.iter().map(|p| (p.seq, path.clone())))
428 .collect();
429 out.sort_by_key(|(seq, _)| *seq);
430 out
431 }
432
433 #[must_use]
439 pub fn crash(&self, crash: &Crash) -> Self {
440 let inner = self.lock();
441 let mut files = BTreeMap::new();
442 for (path, file) in &inner.files {
443 let mut bytes = file.durable.clone();
444 for entry in &file.pending {
445 if crash.keeps(entry.seq) {
446 apply(&mut bytes, &entry.change);
447 }
448 }
449 files.insert(path.clone(), SimFile { durable: bytes, pending: Vec::new() });
450 }
451 Self {
452 inner: Arc::new(Mutex::new(Inner {
453 files,
454 dirs: inner.dirs.clone(),
455 log: Vec::new(),
456 next_seq: 0,
457 fail_at: None,
458 reads: ReadFaults::default(),
459 })),
460 }
461 }
462
463 #[must_use]
467 pub fn durable_contents(&self, path: &Path) -> Option<Vec<u8>> {
468 self.lock().files.get(path).map(|file| file.durable.clone())
469 }
470
471 #[must_use]
473 pub fn contents(&self, path: &Path) -> Option<Vec<u8>> {
474 self.lock().files.get(path).map(SimFile::visible)
475 }
476}
477
478impl Filesystem for SimFilesystem {
479 fn open(&self, path: &Path, mode: OpenMode) -> Result<Box<dyn File>> {
480 let mut inner = self.lock();
481 let exists = inner.files.contains_key(path);
482 match mode {
483 OpenMode::Read | OpenMode::ReadWrite if !exists => {
484 inner.record(Op::Open { path: path.to_path_buf(), mode })?;
487 return Err(Error::io(format!("{} does not exist", path.display())));
488 }
489 OpenMode::CreateNew if exists => {
490 inner.record(Op::Open { path: path.to_path_buf(), mode })?;
491 return Err(Error::io(format!("{} already exists", path.display())));
492 }
493 _ => {}
494 }
495 inner.record(Op::Open { path: path.to_path_buf(), mode })?;
496 inner.files.entry(path.to_path_buf()).or_default();
497 Ok(Box::new(SimHandle {
498 fs: self.clone(),
499 path: path.to_path_buf(),
500 writable: mode.writable(),
501 }))
502 }
503
504 fn exists(&self, path: &Path) -> bool {
505 let inner = self.lock();
506 inner.files.contains_key(path) || inner.dirs.contains(path)
507 }
508
509 fn is_dir(&self, path: &Path) -> bool {
510 let inner = self.lock();
511 inner.dirs.contains(path)
512 }
513
514 fn read_dir(&self, path: &Path) -> Result<Vec<PathBuf>> {
515 let inner = self.lock();
516 if !inner.dirs.contains(path) {
517 return Err(Error::io(format!("{} is not a directory", path.display())));
518 }
519 let mut found: Vec<PathBuf> = inner
522 .files
523 .keys()
524 .chain(inner.dirs.iter())
525 .filter(|entry| entry.parent() == Some(path))
526 .cloned()
527 .collect();
528 found.sort();
529 found.dedup();
530 Ok(found)
531 }
532
533 fn remove(&self, path: &Path) -> Result<()> {
534 let mut inner = self.lock();
535 inner.record(Op::Remove { path: path.to_path_buf() })?;
536 if inner.files.remove(path).is_none() {
537 return Err(Error::io(format!("{} does not exist", path.display())));
538 }
539 Ok(())
540 }
541
542 fn rename(&self, from: &Path, to: &Path) -> Result<()> {
543 let mut inner = self.lock();
544 inner.record(Op::Rename { from: from.to_path_buf(), to: to.to_path_buf() })?;
545 let Some(file) = inner.files.remove(from) else {
546 return Err(Error::io(format!("{} does not exist", from.display())));
547 };
548 inner.files.insert(to.to_path_buf(), file);
549 Ok(())
550 }
551
552 fn create_dir_all(&self, path: &Path) -> Result<()> {
553 let mut inner = self.lock();
554 inner.record(Op::CreateDir { path: path.to_path_buf() })?;
555 let mut current = PathBuf::new();
556 for part in path {
557 current.push(part);
558 inner.dirs.insert(current.clone());
559 }
560 Ok(())
561 }
562
563 fn sync_dir(&self, path: &Path) -> Result<()> {
564 let mut inner = self.lock();
565 inner.record(Op::SyncDir { path: path.to_path_buf() })
566 }
567}
568
569#[derive(Debug)]
571struct SimHandle {
572 fs: SimFilesystem,
573 path: PathBuf,
574 writable: bool,
575}
576
577impl SimHandle {
578 fn missing(&self) -> Error {
579 Error::io(format!("{} was removed while open", self.path.display()))
580 }
581}
582
583impl File for SimHandle {
584 fn read_at(&self, offset: u64, buf: &mut [u8]) -> Result<usize> {
585 self.fs.lock().serve_read(&self.path, offset, buf)
586 }
587
588 fn submit(&self, requests: Vec<Request>) -> Completion {
589 let (completion, filler) = Completion::pending(requests.len());
590 let mut outcomes = Vec::with_capacity(requests.len());
591 for (index, request) in requests.into_iter().enumerate() {
592 let offset = request.offset();
593 let mut buf = request.into_buffer();
594 let outcome =
595 self.read_at(offset, &mut buf).map(|read| Response::new(index, offset, read, buf));
596 outcomes.push((index, outcome));
597 }
598 let order = self.fs.lock().reads.order;
599 reorder(&mut outcomes, order);
600 for (index, outcome) in outcomes {
601 filler.finish(index, outcome);
602 }
603 completion
604 }
605
606 fn write_at(&self, offset: u64, data: &[u8]) -> Result<()> {
607 if !self.writable {
608 return Err(Error::io("this file was opened for reading"));
609 }
610 let mut inner = self.fs.lock();
611 inner.record(Op::Write { path: self.path.clone(), offset, len: data.len() })?;
612 let seq = inner.next_seq;
613 inner.next_seq += 1;
614 let file = inner.files.get_mut(&self.path).ok_or_else(|| self.missing())?;
615 file.pending.push(Pending { seq, change: Change::Write { offset, data: data.to_vec() } });
616 Ok(())
617 }
618
619 fn sync(&self) -> Result<()> {
620 let mut inner = self.fs.lock();
621 inner.record(Op::Sync { path: self.path.clone() })?;
622 let file = inner.files.get_mut(&self.path).ok_or_else(|| self.missing())?;
623 let pending = std::mem::take(&mut file.pending);
624 let mut durable = std::mem::take(&mut file.durable);
625 for entry in &pending {
626 apply(&mut durable, &entry.change);
627 }
628 file.durable = durable;
629 Ok(())
630 }
631
632 fn truncate(&self, len: u64) -> Result<()> {
633 if !self.writable {
634 return Err(Error::io("this file was opened for reading"));
635 }
636 let mut inner = self.fs.lock();
637 inner.record(Op::Truncate { path: self.path.clone(), len })?;
638 let seq = inner.next_seq;
639 inner.next_seq += 1;
640 let file = inner.files.get_mut(&self.path).ok_or_else(|| self.missing())?;
641 file.pending.push(Pending { seq, change: Change::Truncate { len } });
642 Ok(())
643 }
644
645 fn len(&self) -> Result<u64> {
646 let inner = self.fs.lock();
647 let file = inner.files.get(&self.path).ok_or_else(|| self.missing())?;
648 Ok(file.visible().len() as u64)
649 }
650}
651
652#[cfg(test)]
653mod tests {
654 use std::path::Path;
655
656 use super::{Completions, Crash, Op, SimFilesystem};
657 use crate::submit::{Request, Response};
658 use crate::{Filesystem, OpenMode};
659
660 fn write_two_unsynced(fs: &SimFilesystem) {
661 let file = fs.open(Path::new("/db"), OpenMode::Create).unwrap();
662 file.write_at(0, b"AAAA").unwrap();
663 file.sync().unwrap();
664 file.write_at(0, b"BBBB").unwrap();
665 file.write_at(4, b"CCCC").unwrap();
666 }
667
668 #[test]
669 fn a_reader_sees_a_write_before_it_is_durable() {
670 let fs = SimFilesystem::new();
674 let file = fs.open(Path::new("/db"), OpenMode::Create).unwrap();
675 file.write_at(0, b"hello").unwrap();
676 let mut buf = [0u8; 5];
677 file.read_exact_at(0, &mut buf).unwrap();
678 assert_eq!(&buf, b"hello");
679 assert_eq!(fs.durable_contents(Path::new("/db")).unwrap(), Vec::<u8>::new());
680 }
681
682 #[test]
683 fn a_crash_can_leave_either_both_or_neither() {
684 let mut seen = Vec::new();
688 for kept in [vec![], vec![1], vec![2], vec![1, 2]] {
689 let fs = SimFilesystem::new();
690 write_two_unsynced(&fs);
691 let pending = fs.pending();
692 assert_eq!(pending.len(), 2, "both writes are unsynced");
693 let after = fs.crash(&Crash::Keeping(kept.clone()));
694 seen.push(after.durable_contents(Path::new("/db")).unwrap());
695 }
696 assert_eq!(seen[0], b"AAAA".to_vec(), "neither landed");
697 assert_eq!(seen[1], b"BBBB".to_vec(), "the first landed");
698 assert_eq!(seen[2], b"AAAACCCC".to_vec(), "the second landed and left a hole of zeroes");
699 assert_eq!(seen[3], b"BBBBCCCC".to_vec(), "both landed");
700 }
701
702 #[test]
703 fn everything_before_a_sync_survives_a_crash() {
704 let fs = SimFilesystem::new();
705 write_two_unsynced(&fs);
706 let after = fs.crash(&Crash::LosingUnsynced);
707 assert_eq!(after.durable_contents(Path::new("/db")).unwrap(), b"AAAA".to_vec());
708 assert!(after.pending().is_empty(), "a crashed filesystem has nothing in flight");
709 assert_eq!(after.op_count(), 0, "the log belonged to the process that died");
710 }
711
712 #[test]
713 fn keeping_everything_is_what_a_clean_shutdown_looks_like() {
714 let fs = SimFilesystem::new();
715 write_two_unsynced(&fs);
716 let after = fs.crash(&Crash::KeepingEverything);
717 assert_eq!(after.durable_contents(Path::new("/db")).unwrap(), b"BBBBCCCC".to_vec());
718 }
719
720 #[test]
721 fn every_operation_is_recorded_in_order() {
722 let fs = SimFilesystem::new();
723 fs.create_dir_all(Path::new("/data")).unwrap();
724 let file = fs.open(Path::new("/data/db"), OpenMode::CreateNew).unwrap();
725 file.write_at(0, b"xyz").unwrap();
726 file.truncate(2).unwrap();
727 file.sync().unwrap();
728 fs.rename(Path::new("/data/db"), Path::new("/data/live")).unwrap();
729 fs.sync_dir(Path::new("/data")).unwrap();
730 fs.remove(Path::new("/data/live")).unwrap();
731
732 let ops = fs.ops();
733 assert!(matches!(ops[0], Op::CreateDir { .. }));
734 assert!(matches!(ops[1], Op::Open { .. }));
735 assert_eq!(ops[2], Op::Write { path: "/data/db".into(), offset: 0, len: 3 });
736 assert_eq!(ops[3], Op::Truncate { path: "/data/db".into(), len: 2 });
737 assert!(ops[4].is_durability_point());
738 assert!(matches!(ops[5], Op::Rename { .. }));
739 assert!(ops[6].is_durability_point());
740 assert!(matches!(ops[7], Op::Remove { .. }));
741 assert_eq!(fs.op_count(), 8);
742 }
743
744 #[test]
745 fn an_injected_failure_hits_the_operation_it_was_aimed_at() {
746 let fs = SimFilesystem::new();
747 let file = fs.open(Path::new("/db"), OpenMode::Create).unwrap();
748 fs.fail_at(2);
750 assert!(file.write_at(0, b"a").is_ok());
751 assert!(file.write_at(1, b"b").is_err());
752 assert!(file.write_at(1, b"c").is_ok(), "one shot, not a permanently broken disk");
753 assert_eq!(fs.contents(Path::new("/db")).unwrap(), b"ac".to_vec());
756 assert_eq!(fs.op_count(), 4);
757 }
758
759 #[test]
760 fn a_failed_sync_leaves_the_writes_unsynced() {
761 let fs = SimFilesystem::new();
764 let file = fs.open(Path::new("/db"), OpenMode::Create).unwrap();
765 file.write_at(0, b"data").unwrap();
766 fs.fail_at(2);
767 assert!(file.sync().is_err());
768 assert_eq!(fs.durable_contents(Path::new("/db")).unwrap(), Vec::<u8>::new());
769 assert_eq!(fs.pending().len(), 1);
770 }
771
772 #[test]
773 fn the_enumeration_the_crash_tests_will_run_is_expressible_today() {
774 let fs = SimFilesystem::new();
778 let file = fs.open(Path::new("/db"), OpenMode::Create).unwrap();
779 file.write_at(0, b"1").unwrap();
780 file.write_at(1, b"2").unwrap();
781 file.write_at(2, b"3").unwrap();
782 let seqs: Vec<u64> = fs.pending().into_iter().map(|(seq, _)| seq).collect();
783 assert_eq!(seqs.len(), 3);
784
785 let mut outcomes = std::collections::BTreeSet::new();
786 for mask in 0u32..(1 << seqs.len()) {
787 let kept: Vec<u64> = seqs
788 .iter()
789 .enumerate()
790 .filter(|(bit, _)| mask & (1 << bit) != 0)
791 .map(|(_, seq)| *seq)
792 .collect();
793 let after = fs.crash(&Crash::Keeping(kept));
794 outcomes.insert(after.durable_contents(Path::new("/db")).unwrap());
795 }
796 assert_eq!(outcomes.len(), 8);
800 }
801
802 #[test]
803 fn a_handle_on_a_removed_file_reports_that_rather_than_pretending() {
804 let fs = SimFilesystem::new();
805 let file = fs.open(Path::new("/db"), OpenMode::Create).unwrap();
806 file.write_at(0, b"x").unwrap();
807 fs.remove(Path::new("/db")).unwrap();
808 assert!(file.len().is_err());
809 assert!(file.write_at(0, b"y").is_err());
810 }
811
812 #[test]
813 fn opening_a_file_that_is_not_there_fails_and_creating_one_that_is_fails_too() {
814 let fs = SimFilesystem::new();
815 assert!(fs.open(Path::new("/nope"), OpenMode::Read).is_err());
816 assert!(fs.open(Path::new("/nope"), OpenMode::ReadWrite).is_err());
817 fs.open(Path::new("/db"), OpenMode::CreateNew).unwrap();
818 assert!(fs.open(Path::new("/db"), OpenMode::CreateNew).is_err());
819 assert!(fs.open(Path::new("/db"), OpenMode::Create).is_ok());
820 }
821
822 #[test]
823 fn clearing_the_log_keeps_the_contents() {
824 let fs = SimFilesystem::new();
825 let file = fs.open(Path::new("/db"), OpenMode::Create).unwrap();
826 file.write_at(0, b"kept").unwrap();
827 file.sync().unwrap();
828 fs.clear_log();
829 assert_eq!(fs.op_count(), 0);
830 assert_eq!(fs.durable_contents(Path::new("/db")).unwrap(), b"kept".to_vec());
831 }
832
833 fn alphabet(fs: &SimFilesystem) -> Box<dyn crate::File> {
835 let file = fs.open(Path::new("/data"), OpenMode::Create).unwrap();
836 file.write_at(0, b"abcdefghijklmnop").unwrap();
837 file.sync().unwrap();
838 file
839 }
840
841 #[test]
842 fn a_submitted_batch_comes_back_whole_and_in_submission_order() {
843 let fs = SimFilesystem::new();
844 let file = alphabet(&fs);
845 let responses = file
846 .submit(vec![Request::new(0, 4), Request::new(8, 4), Request::new(4, 4)])
847 .wait()
848 .unwrap();
849 let bytes: Vec<&[u8]> = responses.iter().map(Response::bytes).collect();
850 assert_eq!(bytes, [b"abcd", b"ijkl", b"efgh"]);
851 assert_eq!(fs.reads_served(), 3);
852 }
853
854 #[test]
855 fn reordered_completions_still_say_which_request_they_answer() {
856 let fs = SimFilesystem::new();
859 let file = alphabet(&fs);
860 fs.complete(Completions::Reversed);
861 let mut completion = file.submit(vec![Request::new(0, 4), Request::new(4, 4)]);
862 let first = completion.take().unwrap().unwrap();
863 assert_eq!(first.index(), 1);
864 assert_eq!(first.bytes(), b"efgh");
865 assert_eq!(completion.take().unwrap().unwrap().index(), 0);
866 assert!(completion.take().is_none());
867 }
868
869 #[test]
870 fn a_shuffle_is_the_same_shuffle_every_time_for_a_seed() {
871 let order = |seed| {
872 let fs = SimFilesystem::new();
873 let file = alphabet(&fs);
874 fs.complete(Completions::Shuffled(seed));
875 let mut completion =
876 file.submit((0..8).map(|i| Request::new(i * 2, 2)).collect::<Vec<_>>());
877 let mut seen = Vec::new();
878 while let Some(response) = completion.take() {
879 seen.push(response.unwrap().index());
880 }
881 seen
882 };
883 assert_eq!(order(7), order(7), "the same seed is the same run");
884 assert_ne!(order(7), order(8), "and a different one is a different run");
885 let mut sorted = order(7);
886 sorted.sort_unstable();
887 assert_eq!(sorted, (0..8).collect::<Vec<_>>(), "every request is answered exactly once");
888 }
889
890 #[test]
891 fn an_injected_short_read_is_short_and_is_not_an_error() {
892 let fs = SimFilesystem::new();
893 let file = alphabet(&fs);
894 fs.short_read_at(1, 2);
895 let responses = file.submit(vec![Request::new(0, 4), Request::new(4, 4)]).wait().unwrap();
896 assert!(!responses[0].is_short());
897 assert!(responses[1].is_short());
898 assert_eq!(responses[1].bytes(), b"ef");
899 assert_eq!(file.read_at(4, &mut [0u8; 4]).unwrap(), 4);
902 }
903
904 #[test]
905 fn an_error_part_way_through_a_batch_leaves_the_rest_of_the_batch_alone() {
906 let fs = SimFilesystem::new();
907 let file = alphabet(&fs);
908 fs.fail_read_at(1);
909 let mut completion =
910 file.submit(vec![Request::new(0, 4), Request::new(4, 4), Request::new(8, 4)]);
911 let mut answered = 0;
912 let mut failed = 0;
913 while let Some(outcome) = completion.take() {
914 match outcome {
915 Ok(_) => answered += 1,
916 Err(_) => failed += 1,
917 }
918 }
919 assert_eq!((answered, failed), (2, 1));
920 }
921
922 #[test]
923 fn a_failed_read_does_not_shift_which_read_the_next_fault_lands_on() {
924 let fs = SimFilesystem::new();
925 let file = alphabet(&fs);
926 fs.fail_read_at(0);
927 fs.short_read_at(1, 1);
928 let responses = file.submit(vec![Request::new(0, 4), Request::new(4, 4)]);
929 let mut outcomes = responses;
930 let first = outcomes.take().unwrap();
931 let second = outcomes.take().unwrap();
932 assert!(first.is_err());
933 assert_eq!(second.unwrap().bytes(), b"e");
934 }
935
936 #[test]
937 fn read_at_and_a_batch_of_one_are_the_same_read() {
938 let fs = SimFilesystem::new();
939 let file = alphabet(&fs);
940 let mut buf = [0u8; 5];
941 file.read_exact_at(3, &mut buf).unwrap();
942 let batched = file.submit(vec![Request::new(3, 5)]).wait().unwrap();
943 assert_eq!(batched[0].bytes(), &buf);
944 }
945}