| 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 | |
| 5 | use crate::model_ops::{self, AUTHOR}; |
| 6 | use crate::ops; |
| 7 | use crate::server::{Fault, Server, remote_snapshot, snapshot}; |
| 8 | use notebook::{EditStatus, Remote, Replica}; |
| 9 | use onestore::{ |
| 10 | CommitError, ExGuid, RevisionIndex, Stamp, Store, Transaction, |
| 11 | document::Document, |
| 12 | page::{Outline, Page, PageObject}, |
| 13 | }; |
| 14 | use std::{io, sync::LazyLock}; |
| 15 | |
| 16 | static 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 | |
| 35 | fn 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. |
| 46 | fn 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. |
| 62 | fn 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. |
| 76 | pub(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. |
| 132 | struct 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 | |
| 139 | impl 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. |
| 161 | fn local(cache: &Replica) -> Vec<String> { |
| 162 | shape(&cache.page(SOURCE.1).unwrap()) |
| 163 | } |
| 164 | |
| 165 | fn 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. |
| 176 | fn 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 | |
| 187 | pub 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 | } |