1use crate::{
2 disk, model_ops, ops,
3 server::{remote_snapshot, snapshot},
4};
5use notebook::{Remote, Replica};
6use onestore::{
7 CommitError, ExGuid, PageEdit, RevisionIndex, Stamp, Store, Transaction,
8 document::{Document, Kind},
9 op::SectionOp,
10 page::{Page, PageObject},
11};
12use std::{io, sync::LazyLock};
13
14const BODY: &str = "Body 🦋 é";
15
16#[path = "../../../onestore/tests/support/current.rs"]
17mod current;
18
19static SOURCE: LazyLock<Vec<u8>> = LazyLock::new(|| {
20 let mut source = onestore::create_section("page-schedule.one", "Original", "Author").unwrap();
21 for _ in 0..3 {
22 let page = onestore::PageCreation::new(None, Some("Same title"), "Author").unwrap();
23 source = ops::section_op(&source, SectionOp::Create(page.clone()))
24 .unwrap()
25 .as_bytes()
26 .to_vec();
27 }
28 source
29});
30
31fn pages(source: &[u8]) -> Vec<(ExGuid, ExGuid, u32)> {
32 let store = Store::parse(source).unwrap();
33 let index = RevisionIndex::parse(&store).unwrap();
34 index.validate_current().unwrap();
35 let document = Document::parse(&index).unwrap();
36 document
37 .pages()
38 .unwrap()
39 .into_iter()
40 .map(|(sid, page)| {
41 let space = &document.spaces[&sid];
42 let view = &space.revisions[&space.contexts[&ExGuid::default()]];
43 let Kind::Metadata { level, .. } = view.nodes[&view.roots[&2]].kind else {
44 panic!()
45 };
46 (sid, page, level.unwrap_or(1))
47 })
48 .collect()
49}
50
51/// Adds a body outline holding one plain paragraph, returning its text identity. A page the
52/// section just created has no body text for `model_ops::insert_outline` to copy formatting from.
53pub fn body_outline(page: &mut Page, text: &str) -> ExGuid {
54 let paragraph = ops::paragraph(text);
55 let id = paragraph.text().unwrap().id;
56 let at = page
57 .objects
58 .iter()
59 .position(|object| matches!(object, PageObject::Title(_)))
60 .unwrap_or(page.objects.len());
61 page.objects
62 .insert(at, ops::outline(36.0, 36.0, vec![paragraph]));
63 id
64}
65
66/// Publishes through the crash-injecting disk, asserting each commit is atomic.
67struct Session<'a> {
68 disk: &'a mut disk::Disk,
69}
70
71impl Remote for Session<'_> {
72 fn read(&mut self) -> io::Result<Vec<u8>> {
73 Ok(self.disk.visible.clone())
74 }
75
76 fn stamp(&mut self) -> io::Result<Stamp> {
77 Stamp::of(&self.disk.visible).map_err(io::Error::other)
78 }
79
80 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
81 let mut image = self.disk.visible.clone();
82 transaction.apply(&mut image).unwrap();
83 let old = current::current(&self.disk.durable);
84 let new = current::current(&image);
85 let result = transaction.commit(self.disk);
86 let observed = current::current(&self.disk.durable);
87 assert!(observed == old || observed == new);
88 if result.is_ok() {
89 assert_eq!(observed, new);
90 }
91 result
92 }
93
94 fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> {
95 onestore::confirm(self.disk, base)
96 }
97}
98
99fn section_op(cache: &Replica, op: onestore::op::SectionOp) -> Result<u64, notebook::Error> {
100 cache.apply(
101 "Author",
102 onestore::op::Edit {
103 at: 133_000_000_000_000_000,
104 ops: vec![onestore::op::Op::Section(op)],
105 },
106 )
107}
108
109pub fn run(input: &[u8]) {
110 let directory = tempfile::tempdir().unwrap();
111 let mut replicas: [Option<Replica>; 12] = std::array::from_fn(|_| None);
112 let mut owned: [Vec<(ExGuid, ExGuid)>; 12] = std::array::from_fn(|_| Vec::new());
113 let mut disk = disk::Disk {
114 visible: SOURCE.clone(),
115 durable: SOURCE.clone(),
116 operation: 0,
117 fail_at: None,
118 write_limit: 4096,
119 random: 1,
120 };
121 for step in input.chunks_exact(8).take(48) {
122 let actor = usize::from(step[0]) % replicas.len();
123 let path = directory.path().join(format!("{actor}.sqlite"));
124 let cache = replicas[actor].get_or_insert_with(|| Replica::create(&path, &SOURCE).unwrap());
125 let image = snapshot(cache);
126 let listed = pages(&image);
127 let pending = cache.pending().unwrap();
128 match step[1] % 8 {
129 0 | 1 => {
130 let title = (step[2] & 1 != 0).then_some("Same 🦋 é");
131 let page = onestore::PageCreation::new(None, title, "Author").unwrap();
132 section_op(cache, onestore::op::SectionOp::Create(page.clone())).unwrap();
133 owned[actor].push((page.space(), page.object()));
134 if step[1] % 8 == 1 {
135 let mut model = cache.page(page.space()).unwrap();
136 body_outline(&mut model, BODY);
137 model_ops::save_as(cache, page.space(), &model, "Author")
138 .unwrap()
139 .unwrap();
140 let local = cache.page(page.space()).unwrap();
141 assert!(
142 local
143 .objects
144 .iter()
145 .any(|object| matches!(object, PageObject::Outline(o) if o.paragraphs[0].text().is_some_and(|t| t.text.text() == BODY)))
146 );
147 }
148 }
149 2 => {
150 let statuses: Vec<_> = pending
151 .iter()
152 .map(|edit| cache.status(edit.id).unwrap())
153 .collect();
154 replicas[actor] = None;
155 let cache = Replica::open(&path).unwrap();
156 assert_eq!(pages(&snapshot(&cache)), listed);
157 assert_eq!(cache.pending().unwrap(), pending);
158 assert_eq!(
159 pending
160 .iter()
161 .map(|edit| cache.status(edit.id).unwrap())
162 .collect::<Vec<_>>(),
163 statuses
164 );
165 replicas[actor] = Some(cache);
166 }
167 3 | 4 => {
168 disk.operation = 0;
169 disk.write_limit = if step[3] & 1 == 0 { 17 } else { 4096 };
170 disk.fail_at =
171 (step[2] != 0).then_some(usize::from(u16::from_le_bytes([step[2], step[3]])));
172 disk.random = u64::from(step[4]) + 1;
173 let result = cache.sync_once(&mut Session { disk: &mut disk });
174 if result.is_ok() && cache.pending().unwrap().is_empty() {
175 assert_eq!(pages(&snapshot(cache)), pages(&remote_snapshot(cache)));
176 }
177 }
178 5 => {
179 disk.visible.clone_from(&disk.durable);
180 disk.fail_at = None;
181 }
182 6 => {
183 let count = 1 + usize::from(step[3] & 1 != 0);
184 let selected: Vec<_> = (0..count)
185 .map(|i| listed[(usize::from(step[2]) + i) % listed.len()].0)
186 .collect();
187 let anchor = (step[4] & 1 != 0)
188 .then_some(listed[usize::from(step[5]) % listed.len()].0)
189 .filter(|sid| !selected.contains(sid));
190 let edits: Vec<_> = selected
191 .iter()
192 .enumerate()
193 .map(|(i, sid)| {
194 let level = u32::from(step[6 + i]) % 3 + 1;
195 if step[4] & 2 == 0 {
196 PageEdit::set_level(*sid, level).unwrap()
197 } else {
198 PageEdit::move_to(*sid, anchor, level).unwrap()
199 }
200 })
201 .collect();
202 let expected = ops::section_op(&image, SectionOp::Pages(edits.to_vec()));
203 match section_op(cache, onestore::op::SectionOp::Pages(edits)) {
204 Ok(_) => {
205 assert_eq!(pages(&snapshot(cache)), pages(expected.unwrap().as_bytes()))
206 }
207 Err(_) => {
208 assert!(expected.is_err());
209 assert_eq!(pages(&snapshot(cache)), listed);
210 assert_eq!(cache.pending().unwrap(), pending);
211 }
212 }
213 }
214 7 => {
215 // The actor merged a conflict page's version by hand and deletes it.
216 if let Some(conflict) = cache
217 .conflicts()
218 .unwrap()
219 .first()
220 .map(|(_, pages)| pages[usize::from(step[2]) % pages.len()].space)
221 {
222 section_op(cache, onestore::op::SectionOp::Delete(vec![conflict])).unwrap();
223 assert_eq!(pages(&snapshot(cache)), listed);
224 }
225 }
226 _ => unreachable!(),
227 }
228 let local = pages(&snapshot(replicas[actor].as_ref().unwrap()));
229 for page in &owned[actor] {
230 assert!(local.iter().any(|(sid, oid, _)| (*sid, *oid) == *page));
231 }
232 }
233}