1//! A deterministic multi-actor schedule over page-model saves: several replicas edit one
2//! page offline, publish through a fault-injecting remote, reopen, and delete the conflict
3//! pages the merges leave.
4
5use crate::model_ops::{self, AUTHOR};
6use crate::ops;
7use crate::server::{Fault, Server, remote_snapshot, snapshot};
8use notebook::{EditStatus, Remote, Replica};
9use onestore::{
10 CommitError, ExGuid, RevisionIndex, Stamp, Store, Transaction,
11 document::Document,
12 page::{Outline, Page, PageObject},
13};
14use std::{io, sync::LazyLock};
15
16static SOURCE: LazyLock<(Vec<u8>, ExGuid)> = LazyLock::new(|| {
17 let source =
18 onestore::create_section("model-schedule.one", "Original 🦀 é 0", "Author").unwrap();
19 let store = Store::parse(&source).unwrap();
20 let index = RevisionIndex::parse(&store).unwrap();
21 let document = Document::parse(&index).unwrap();
22 let space = document.pages().unwrap()[0].0;
23 let mut page = Page::from_space(&document, space).unwrap();
24 let mut anchor = body(&page).paragraphs[0].text().unwrap().id;
25 for at in 1..8 {
26 anchor = model_ops::insert_after(&mut page, anchor, &format!("Original 🦀 é {at}")).1;
27 }
28 let source = ops::saved(&source, space, &page)
29 .unwrap()
30 .as_slice()
31 .to_vec();
32 (source, space)
33});
34
35fn body(page: &Page) -> &Outline {
36 page.objects
37 .iter()
38 .find_map(|object| match object {
39 PageObject::Outline(outline) => Some(outline),
40 _ => None,
41 })
42 .expect("the page keeps one body outline")
43}
44
45/// Everything the schedule asserts about a page, free of the identities the writers mint.
46fn shape(page: &Page) -> Vec<String> {
47 let outline = body(page);
48 let mut out = vec![format!("{:?}", outline.layout)];
49 out.extend(outline.paragraphs.iter().map(|paragraph| {
50 format!(
51 "{} {} {} {:?}",
52 paragraph.level,
53 u8::from(paragraph.collapsed),
54 u8::from(paragraph.parent.is_some()),
55 paragraph.text().map(|text| text.text.text()),
56 )
57 }));
58 out
59}
60
61/// Paragraph and text identities of the body outline with their text, in model order.
62fn rows(bytes: &[u8]) -> Vec<(ExGuid, ExGuid, String)> {
63 let page = model_ops::page_of(bytes, SOURCE.1);
64 body(&page)
65 .paragraphs
66 .iter()
67 .filter_map(|paragraph| {
68 let text = paragraph.text()?;
69 Some((paragraph.id, text.id, text.text.text().to_owned()))
70 })
71 .collect()
72}
73
74/// Moves a paragraph and its descendants immediately before `anchor`, or last among
75/// `parent`'s children; `anchor` supplies its own parent and level.
76pub(crate) fn move_subtree(
77 page: &mut Page,
78 paragraph: ExGuid,
79 parent: Option<ExGuid>,
80 anchor: Option<ExGuid>,
81) {
82 for outline in model_ops::outlines_mut(page) {
83 let Some(at) = outline.paragraphs.iter().position(|p| p.id == paragraph) else {
84 continue;
85 };
86 let level = outline.paragraphs[at].level;
87 let mut end = at + 1;
88 while end < outline.paragraphs.len() && outline.paragraphs[end].level > level {
89 end += 1;
90 }
91 if parent
92 .into_iter()
93 .chain(anchor)
94 .any(|id| outline.paragraphs[at..end].iter().any(|p| p.id == id))
95 {
96 return;
97 }
98 let mut subtree: Vec<_> = outline.paragraphs.drain(at..end).collect();
99 let (destination, parent, depth) = match anchor {
100 Some(anchor) => {
101 let Some(at) = outline.paragraphs.iter().position(|p| p.id == anchor) else {
102 return;
103 };
104 let target = &outline.paragraphs[at];
105 (at, target.parent, target.level)
106 }
107 None => match parent {
108 Some(parent) => {
109 let Some(at) = outline.paragraphs.iter().position(|p| p.id == parent) else {
110 return;
111 };
112 let level = outline.paragraphs[at].level;
113 let mut end = at + 1;
114 while end < outline.paragraphs.len() && outline.paragraphs[end].level > level {
115 end += 1;
116 }
117 (end, Some(parent), level + 1)
118 }
119 None => (outline.paragraphs.len(), None, 1),
120 },
121 };
122 subtree[0].parent = parent;
123 for descendant in &mut subtree {
124 descendant.level = descendant.level + depth - level;
125 }
126 outline.paragraphs.splice(destination..destination, subtree);
127 return;
128 }
129}
130
131/// Asserts the publication invariants, then delegates.
132struct Session<'a> {
133 server: &'a mut Server,
134 /// The head batch was already attempted, so its publication must never be repeated.
135 retired: bool,
136 publications: usize,
137}
138
139impl Remote for Session<'_> {
140 fn read(&mut self) -> io::Result<Vec<u8>> {
141 self.server.read()
142 }
143
144 fn stamp(&mut self) -> io::Result<Stamp> {
145 self.server.stamp()
146 }
147
148 fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> {
149 assert!(!self.retired, "a retired attempt was replayed");
150 self.publications += 1;
151 assert_eq!(self.publications, 1);
152 self.server.publish(transaction)
153 }
154
155 fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> {
156 self.server.confirm(base)
157 }
158}
159
160/// The page's shape as the replica's queue leaves it.
161fn local(cache: &Replica) -> Vec<String> {
162 shape(&cache.page(SOURCE.1).unwrap())
163}
164
165fn statuses(cache: &Replica) -> Vec<(u64, Option<EditStatus>)> {
166 cache
167 .pending()
168 .unwrap()
169 .iter()
170 .map(|edit| (edit.id, cache.status(edit.id).unwrap()))
171 .collect()
172}
173
174/// Fails unless the durable remote holds `revision`: in the edited page's space, or, for
175/// the edit adding a conflict page, in the space it changed first.
176fn durable(server: &Server, revision: ExGuid) {
177 let store = Store::parse(&server.durable).unwrap();
178 let index = RevisionIndex::parse(&store).unwrap();
179 assert!(
180 index
181 .spaces
182 .values()
183 .any(|space| space.revisions.contains_key(&revision))
184 );
185}
186
187pub fn run(input: &[u8]) {
188 let (source, space) = &*SOURCE;
189 let directory = tempfile::tempdir().unwrap();
190 let mut replicas: [Option<Replica>; 12] = std::array::from_fn(|_| None);
191 let mut server = Server::new(source);
192 for step in input.chunks_exact(8).take(48) {
193 let actor = usize::from(step[0]) % replicas.len();
194 let path = directory.path().join(format!("{actor}.sqlite"));
195 let cache = replicas[actor].get_or_insert_with(|| Replica::create(&path, source).unwrap());
196 let image = snapshot(cache);
197 let before = local(cache);
198 let pending = cache.pending().unwrap();
199 let rows_now = rows(&image);
200 match step[1] % 8 {
201 0..=2 => {
202 let row = rows_now[usize::from(step[2]) % rows_now.len()].clone();
203 let other = rows_now[usize::from(step[4]) % rows_now.len()].clone();
204 let change: Box<dyn Fn(&mut Page)> = match step[3] % 6 {
205 0 if !row.2.is_empty() => Box::new(move |page| {
206 model_ops::replace_text(
207 page,
208 row.1,
209 0..1,
210 &char::from(b'A' + step[5] % 26).to_string(),
211 );
212 }),
213 1 => Box::new(move |page| {
214 model_ops::insert_after(page, row.1, "Inserted 🦋 é");
215 }),
216 2 => Box::new(move |page| {
217 move_subtree(page, row.0, (step[5] & 1 != 0).then_some(other.0), None);
218 }),
219 3 => Box::new(move |page| model_ops::delete_paragraph(page, row.1)),
220 4 => Box::new(move |page| {
221 let outline = model_ops::outlines_mut(page).swap_remove(0);
222 if step[5] & 1 == 0 {
223 outline.layout.x = Some(f32::from(step[5]) + 36.0);
224 outline.layout.y = Some(f32::from(step[6]) + 36.0);
225 } else {
226 outline.layout.max_width = Some(f32::from(step[6]) + 72.0);
227 outline.layout.width_set_by_user = Some(step[6] & 2 == 0);
228 }
229 }),
230 5 if !row.2.is_empty() => Box::new(move |page| {
231 model_ops::restyle(page, row.1, 0..1, |format| {
232 format.font_size = Some(f32::from(step[5] % 20) + 12.0);
233 });
234 }),
235 _ => Box::new(|_| {}),
236 };
237 let mut after = cache.page(*space).unwrap();
238 change(&mut after);
239 match model_ops::save_as(cache, *space, &after, AUTHOR) {
240 Ok(id) => {
241 let next = cache.pending().unwrap();
242 let ids = |edits: &[notebook::PendingEdit]| -> Vec<u64> {
243 edits.iter().map(|edit| edit.id).collect()
244 };
245 assert_eq!(ids(&next)[..pending.len()], ids(&pending)[..]);
246 assert_eq!(next.len(), pending.len() + usize::from(id.is_some()));
247 assert_eq!(local(cache), shape(&after));
248 }
249 Err(_) => {
250 assert_eq!(local(cache), before);
251 assert_eq!(cache.pending().unwrap(), pending);
252 }
253 }
254 }
255 3 | 4 => {
256 server.fault = match step[2] % 8 {
257 0 => Fault::Before,
258 1 => Fault::UnknownBefore,
259 2 => Fault::UnknownAfter,
260 3 => Fault::Committed,
261 4 => Fault::Confirm,
262 5 => Fault::ConfirmCommitted,
263 _ => Fault::None,
264 };
265 let head = pending.first().map(|edit| edit.id);
266 let result = {
267 let mut session = Session {
268 retired: head.is_some_and(|id| {
269 matches!(
270 cache.status(id).unwrap(),
271 Some(EditStatus::AwaitingConfirmation { .. })
272 )
273 }),
274 server: &mut server,
275 publications: 0,
276 };
277 cache.sync_once(&mut session)
278 };
279 server.fault = Fault::None;
280 let next = cache.pending().unwrap();
281 // Edits leave the queue oldest first, each with a durable receipt; a rebase
282 // queues the conflict pages it makes after them.
283 let queued =
284 |edit: &notebook::PendingEdit| pending.iter().any(|kept| kept.id == edit.id);
285 let left = pending.len() - next.iter().filter(|edit| queued(edit)).count();
286 assert!(next.iter().all(|edit| {
287 pending[left..].iter().any(|kept| kept.id == edit.id)
288 || matches!(
289 &edit.edit.ops[..],
290 [onestore::op::Op::Section(
291 onestore::op::SectionOp::Conflict { .. }
292 )]
293 )
294 }));
295 for edit in &pending[..left] {
296 let Some(EditStatus::Published { revision }) = cache.status(edit.id).unwrap()
297 else {
298 panic!("an edit leaves the queue only with its receipt")
299 };
300 durable(&server, revision);
301 }
302 if next.is_empty() && result.is_ok() {
303 assert_eq!(
304 local(cache),
305 shape(&model_ops::page_of(&remote_snapshot(cache), *space))
306 );
307 }
308 }
309 5 => {
310 let remote = rows(&server.visible);
311 let row = remote[usize::from(step[2]) % remote.len()].clone();
312 let letter = char::from(b'a' + step[3] % 26);
313 // Retyping the letter already there changes nothing, so stores no revision.
314 if row.2.chars().next().is_some_and(|first| first != letter) {
315 let mut page = model_ops::page_of(&server.visible, *space);
316 model_ops::replace_text(&mut page, row.1, 0..1, &letter.to_string());
317 let visible = server.visible.clone();
318 ops::save(&visible, *space, &page)
319 .unwrap()
320 .commit(&mut server)
321 .unwrap();
322 }
323 }
324 6 => {
325 // The writer merged a conflict page's version by hand and deletes it.
326 let listed = cache.conflicts().unwrap();
327 if let Some(conflict) = listed
328 .iter()
329 .find(|(page, _)| page == space)
330 .map(|(_, pages)| pages[usize::from(step[2]) % pages.len()].space)
331 {
332 let delete =
333 onestore::op::Op::Section(onestore::op::SectionOp::Delete(vec![conflict]));
334 cache
335 .apply(
336 AUTHOR,
337 onestore::op::Edit {
338 at: model_ops::now(),
339 ops: vec![delete],
340 },
341 )
342 .unwrap();
343 assert!(
344 cache
345 .conflicts()
346 .unwrap()
347 .iter()
348 .flat_map(|(_, pages)| pages)
349 .all(|page| page.space != conflict)
350 );
351 assert_eq!(local(cache), before);
352 let (result, publications) = {
353 let mut session = Session {
354 retired: false,
355 server: &mut server,
356 publications: 0,
357 };
358 (cache.sync_once(&mut session), session.publications)
359 };
360 assert!(publications <= 1);
361 if let Ok(notebook::Synced {
362 edit: Some((id, EditStatus::Published { revision })),
363 ..
364 }) = result
365 {
366 durable(&server, revision);
367 assert!(cache.pending().unwrap().iter().all(|e| e.id != id));
368 }
369 }
370 }
371 _ => {
372 let recorded = statuses(cache);
373 replicas[actor] = None;
374 let cache = Replica::open(&path).unwrap();
375 assert_eq!(local(&cache), before);
376 assert_eq!(cache.pending().unwrap(), pending);
377 assert_eq!(statuses(&cache), recorded);
378 replicas[actor] = Some(cache);
379 }
380 }
381 rows(&server.durable);
382 rows(&snapshot(replicas[actor].as_ref().unwrap()));
383 }
384}