The journal.
6fn device(name: &str) -> DeviceName { 7 DeviceName::new(name).expect("a plain device name") 8} 9 10fn entry(n: u64) -> EntryLine { 11 EntryLine::parse(&format!(r#"{{"n":{n}}}"#)).expect("an entry") 12} 13 14fn entries(range: std::ops::Range<u64>) -> Vec<EntryLine> { 15 range.map(entry).collect() 16} 17 18fn cursor(n: u64) -> DeviceCursor { 19 DeviceCursor::new(n) 20}
Runs every journal law.
23pub fn journal<J: Journal>(fresh: &impl Fn() -> J, restart: &impl Fn(J) -> J) { 24 info!("conformance: journal"); 25 journal_keeps_arrival_order_across_devices(fresh()); 26 journal_takes_lines_only_where_the_device_left_off(fresh()); 27 journal_appending_nothing_is_a_question(fresh()); 28 journal_two_appends_at_one_cursor_take_one(fresh()); 29 journal_pull_since_and_the_rebuild_rule(fresh()); 30 journal_a_pull_is_bounded(fresh()); 31 journal_stores_the_same_bytes_everywhere(fresh()); 32 journal_survives_a_restart(fresh, restart); 33}
35pub fn journal_keeps_arrival_order_across_devices<J: Journal>(j: J) { 36 let (p1, p2) = (device("phone1"), device("phone2")); 37 assert_eq!(run_ready(j.append(&p1, cursor(0), &entries(0..2))).unwrap(), Appended::Accepted { held: cursor(2) }); 38 assert_eq!(run_ready(j.append(&p2, cursor(0), &entries(10..11))).unwrap(), Appended::Accepted { held: cursor(1) }); 39 assert_eq!(run_ready(j.append(&p1, cursor(2), &entries(2..3))).unwrap(), Appended::Accepted { held: cursor(3) }); 40 assert_eq!((run_ready(j.held(&p1)).unwrap(), run_ready(j.held(&p2)).unwrap()), (cursor(3), cursor(1)), "each device is counted on its own"); 41 assert_eq!(run_ready(j.held(&device("nobody"))).unwrap(), cursor(0)); 42 let page = run_ready(j.pull(LogCursor::START)).unwrap(); 43 let order: Vec<&str> = page.lines.iter().map(StoredLine::as_str).collect(); 44 assert_eq!( 45 order, 46 [ 47 r#"{"device":"phone1","entry":{"n":0}}"#, 48 r#"{"device":"phone1","entry":{"n":1}}"#, 49 r#"{"device":"phone2","entry":{"n":10}}"#, 50 r#"{"device":"phone1","entry":{"n":2}}"#, 51 ], 52 "one log, in the order the hub received the lines" 53 ); 54 assert_eq!((page.from.get(), page.total.get()), (0, 4)); 55} 56 57pub fn journal_takes_lines_only_where_the_device_left_off<J: Journal>(j: J) { 58 let p = device("phone1"); 59 run_ready(j.append(&p, cursor(0), &entries(0..2))).unwrap(); 60 // Behind (it did not hear the reply), ahead (the hub was restored from a backup): nothing is 61 // taken, and the device is told the truth so it can send from the right place. 62 assert_eq!(run_ready(j.append(&p, cursor(1), &entries(5..6))).unwrap(), Appended::Misaligned { held: cursor(2) }); 63 assert_eq!(run_ready(j.append(&p, cursor(9), &entries(5..6))).unwrap(), Appended::Misaligned { held: cursor(2) }); 64 assert_eq!(run_ready(j.pull(LogCursor::START)).unwrap().total.get(), 2, "and nothing was stored"); 65 assert_eq!(run_ready(j.append(&p, cursor(2), &entries(2..3))).unwrap().held(), cursor(3)); 66} 67 68pub fn journal_appending_nothing_is_a_question<J: Journal>(j: J) { 69 let p = device("phone1"); 70 assert_eq!(run_ready(j.append(&p, cursor(0), &[])).unwrap(), Appended::Accepted { held: cursor(0) }); 71 run_ready(j.append(&p, cursor(0), &entries(0..3))).unwrap(); 72 for at in [0, 3, 50] { 73 assert_eq!(run_ready(j.append(&p, cursor(at), &[])).unwrap(), Appended::Accepted { held: cursor(3) }, "asked at {at}"); 74 } 75} 76 77pub fn journal_two_appends_at_one_cursor_take_one<J: Journal>(j: J) { 78 let p = device("phone1"); 79 let first = run_ready(j.append(&p, cursor(0), &entries(0..2))).unwrap(); 80 let second = run_ready(j.append(&p, cursor(0), &entries(0..2))).unwrap(); 81 assert_eq!(first, Appended::Accepted { held: cursor(2) }); 82 assert_eq!(second, Appended::Misaligned { held: cursor(2) }, "the same lines are never taken twice"); 83 assert_eq!(run_ready(j.pull(LogCursor::START)).unwrap().total.get(), 2); 84} 85 86pub fn journal_pull_since_and_the_rebuild_rule<J: Journal>(j: J) { 87 let p = device("phone1"); 88 run_ready(j.append(&p, cursor(0), &entries(0..3))).unwrap(); 89 let page = run_ready(j.pull(LogCursor::new(2))).unwrap(); 90 assert_eq!((page.from.get(), page.total.get(), page.lines.len()), (2, 3, 1)); 91 assert!(!page.caller_is_ahead(LogCursor::new(2))); 92 let page = run_ready(j.pull(LogCursor::new(3))).unwrap(); 93 assert_eq!((page.from.get(), page.lines.len()), (3, 0), "at the end there is nothing more"); 94 // A caller whose copy is longer than the log (a damaged copy that duplicated lines) is told 95 // where the log really ends, and sees from the page that its copy is the wrong one. 96 let page = run_ready(j.pull(LogCursor::new(99))).unwrap(); 97 assert_eq!((page.from.get(), page.total.get(), page.lines.len()), (3, 3, 0), "clamped to the end of the log"); 98 assert!(page.caller_is_ahead(LogCursor::new(99)), "the caller rebuilds its copy"); 99} 100 101pub fn journal_a_pull_is_bounded<J: Journal>(j: J) { 102 let p = device("phone1"); 103 let n = (PULL_LIMIT + 5) as u64; 104 assert_eq!(run_ready(j.append(&p, cursor(0), &entries(0..n))).unwrap().held(), cursor(n), "a big append is one append"); 105 let page = run_ready(j.pull(LogCursor::START)).unwrap(); 106 assert_eq!((page.lines.len(), page.total.get()), (PULL_LIMIT, n)); 107 let rest = run_ready(j.pull(LogCursor::new(PULL_LIMIT as u64))).unwrap(); 108 assert_eq!((rest.from.get(), rest.lines.len()), (PULL_LIMIT as u64, 5), "the caller asks again for the rest"); 109} 110 111pub fn journal_stores_the_same_bytes_everywhere<J: Journal>(j: J) { 112 let p = device("phone1"); 113 let e = EntryLine::parse(r#"{ "b": 1, "a": { "z": 1, "y": 2 } }"#).unwrap(); 114 run_ready(j.append(&p, cursor(0), std::slice::from_ref(&e))).unwrap(); 115 let page = run_ready(j.pull(LogCursor::START)).unwrap(); 116 assert_eq!(page.lines[0].as_str(), r#"{"device":"phone1","entry":{"b":1,"a":{"z":1,"y":2}}}"#, "compact, keys in the order the device wrote them"); 117 assert_eq!(page.lines[0], StoredLine::of(&p, &e)); 118 assert_eq!(run_ready(j.pull(LogCursor::START)).unwrap(), page, "a line is never rewritten"); 119} 120 121pub fn journal_survives_a_restart<J: Journal>(fresh: &impl Fn() -> J, restart: &impl Fn(J) -> J) { 122 let j = fresh(); 123 let p = device("phone1"); 124 run_ready(j.append(&p, cursor(0), &entries(0..3))).unwrap(); 125 let before = run_ready(j.pull(LogCursor::START)).unwrap(); 126 let j = restart(j); 127 assert_eq!(run_ready(j.held(&p)).unwrap(), cursor(3), "the device is still counted"); 128 assert_eq!(run_ready(j.pull(LogCursor::START)).unwrap(), before); 129 assert_eq!(run_ready(j.append(&p, cursor(3), &entries(3..4))).unwrap().held(), cursor(4), "and carries on from where it was"); 130}