The journal.

3use ::log::info;
4use whiskers_ports::{Appended, DeviceCursor, DeviceName, EntryLine, Journal, LogCursor, PULL_LIMIT, StoredLine, run_ready};
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}