| 1 | //! `open_probe SECTION | --paragraphs N`: times opening a section through a fresh and a |
| 2 | //! reopened cache up to its page list and first page, then how long page reads wait while |
| 3 | //! the sync thread rebases the queue onto a changed remote. |
| 4 | |
| 5 | use notebook::{Remote, Replica, session::Section}; |
| 6 | use onestore::{ |
| 7 | CommitError, ExGuid, Stamp, Transaction, |
| 8 | document::{Format, Layout}, |
| 9 | op::{Edit, Op, PageOp, SectionOp}, |
| 10 | page::{ |
| 11 | Outline, Page, PageObject, PageParagraph, Paragraph, ParagraphContent, TextObject, |
| 12 | text::new_id, |
| 13 | }, |
| 14 | }; |
| 15 | use std::{ |
| 16 | io, |
| 17 | time::{Duration, Instant}, |
| 18 | }; |
| 19 | |
| 20 | /// A section holding one page of `count` paragraphs. |
| 21 | fn generated(count: usize) -> Result<Vec<u8>, Box<dyn std::error::Error>> { |
| 22 | let paragraph = |i: usize| PageParagraph { |
| 23 | id: new_id().unwrap(), |
| 24 | parent: None, |
| 25 | level: 1, |
| 26 | style: None, |
| 27 | format: Format::default(), |
| 28 | content: ParagraphContent::Text(TextObject { |
| 29 | id: new_id().unwrap(), |
| 30 | date_field: None, |
| 31 | text: Paragraph::new(format!("Paragraph {i} of the probe"), Format::default()), |
| 32 | tags: Vec::new(), |
| 33 | }), |
| 34 | lists: Vec::new(), |
| 35 | tags: Vec::new(), |
| 36 | media: Default::default(), |
| 37 | collapsed: false, |
| 38 | }; |
| 39 | let arena = onestore::Arena::default(); |
| 40 | let mut section = |
| 41 | onestore::Section::open(&arena, onestore::create_section("probe.one", "", "Probe")?)?; |
| 42 | let page = Page { |
| 43 | title: "Probe".into(), |
| 44 | identity: None, |
| 45 | created: None, |
| 46 | margin_origin: [0.0; 2], |
| 47 | rtl: false, |
| 48 | color: None, |
| 49 | rule_lines: None, |
| 50 | objects: vec![PageObject::Outline(Outline { |
| 51 | id: new_id()?, |
| 52 | title: false, |
| 53 | min_width: None, |
| 54 | layout: Layout { |
| 55 | x: Some(36.0), |
| 56 | y: Some(86.4), |
| 57 | ..Layout::default() |
| 58 | }, |
| 59 | indents: Vec::new(), |
| 60 | paragraphs: (0..count).map(paragraph).collect(), |
| 61 | unsupported: Vec::new(), |
| 62 | })], |
| 63 | definitions: Default::default(), |
| 64 | }; |
| 65 | let creation = onestore::PageCreation::new(None, Some("Probe"), "Probe")?; |
| 66 | section.apply( |
| 67 | "Probe", |
| 68 | &Edit { |
| 69 | at: 133_000_000_000_000_000, |
| 70 | ops: vec![Op::Section(SectionOp::Import { creation, page })], |
| 71 | }, |
| 72 | )?; |
| 73 | section.seal()?; |
| 74 | Ok(section.image()) |
| 75 | } |
| 76 | |
| 77 | /// The first text of the page with the most outline paragraphs, and that page's space. |
| 78 | fn text(page: &Page) -> Option<ExGuid> { |
| 79 | page.objects.iter().find_map(|object| match object { |
| 80 | PageObject::Outline(outline) => outline |
| 81 | .paragraphs |
| 82 | .iter() |
| 83 | .find_map(|p| p.text().map(|t| t.id)), |
| 84 | _ => None, |
| 85 | }) |
| 86 | } |
| 87 | |
| 88 | struct Memory(Vec<u8>); |
| 89 | |
| 90 | impl onestore::CommitIo for Memory { |
| 91 | fn read_at(&mut self, offset: u64, output: &mut [u8]) -> io::Result<usize> { |
| 92 | let offset = offset as usize; |
| 93 | let size = output.len().min(self.0.len().saturating_sub(offset)); |
| 94 | output[..size].copy_from_slice(&self.0[offset..offset + size]); |
| 95 | Ok(size) |
| 96 | } |
| 97 | fn write_at(&mut self, offset: u64, bytes: &[u8]) -> io::Result<usize> { |
| 98 | let offset = offset as usize; |
| 99 | self.0.resize(self.0.len().max(offset + bytes.len()), 0); |
| 100 | self.0[offset..offset + bytes.len()].copy_from_slice(bytes); |
| 101 | Ok(bytes.len()) |
| 102 | } |
| 103 | fn flush(&mut self) -> io::Result<()> { |
| 104 | Ok(()) |
| 105 | } |
| 106 | } |
| 107 | |
| 108 | impl Remote for Memory { |
| 109 | fn read(&mut self) -> io::Result<Vec<u8>> { |
| 110 | Ok(self.0.clone()) |
| 111 | } |
| 112 | fn stamp(&mut self) -> io::Result<Stamp> { |
| 113 | Stamp::of(&self.0).map_err(io::Error::other) |
| 114 | } |
| 115 | fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> { |
| 116 | transaction.commit(self) |
| 117 | } |
| 118 | fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> { |
| 119 | onestore::confirm(self, base) |
| 120 | } |
| 121 | } |
| 122 | |
| 123 | fn main() -> Result<(), Box<dyn std::error::Error>> { |
| 124 | let mut args = std::env::args().skip(1); |
| 125 | let source = match (args.next().ok_or("SECTION | --paragraphs N")?, args.next()) { |
| 126 | (flag, Some(count)) if flag == "--paragraphs" => generated(count.parse()?)?, |
| 127 | (path, _) => std::fs::read(path)?, |
| 128 | }; |
| 129 | let directory = tempfile::tempdir()?; |
| 130 | let file = directory.path().join("section.one"); |
| 131 | std::fs::write(&file, &source)?; |
| 132 | let cache = directory.path().join("cache"); |
| 133 | println!("section {} bytes", source.len()); |
| 134 | for label in ["fresh cache", "reopened cache"] { |
| 135 | let start = Instant::now(); |
| 136 | let section = Section::open(&file, &cache, || {})?; |
| 137 | let opened = start.elapsed(); |
| 138 | let pages = section.pages()?; |
| 139 | let listed = start.elapsed(); |
| 140 | let page = section.page(pages[0].0)?; |
| 141 | let read = start.elapsed(); |
| 142 | println!( |
| 143 | "{label}: open {opened:.2?}, page list ({} pages) at {listed:.2?}, first page ({} objects) at {read:.2?}", |
| 144 | pages.len(), |
| 145 | page.objects.len() |
| 146 | ); |
| 147 | section.close()?; |
| 148 | } |
| 149 | |
| 150 | // A local edit queued against a remote another writer changed elsewhere: the next |
| 151 | // step rebases, while this thread keeps reading the page list and a page. |
| 152 | let replica = std::sync::Arc::new(Replica::create( |
| 153 | directory.path().join("rebase.sqlite"), |
| 154 | &source, |
| 155 | )?); |
| 156 | let (space, ..) = replica |
| 157 | .pages()? |
| 158 | .into_iter() |
| 159 | .max_by_key(|(space, ..)| { |
| 160 | replica |
| 161 | .page(*space) |
| 162 | .map(|page| page.objects.len()) |
| 163 | .unwrap_or(0) |
| 164 | }) |
| 165 | .ok_or("a page")?; |
| 166 | let target = text(&replica.page(space)?).ok_or("a text")?; |
| 167 | let typed = |with: &str| Edit { |
| 168 | at: 133_000_000_000_000_000, |
| 169 | ops: vec![Op::Page { |
| 170 | space, |
| 171 | op: PageOp::Text { |
| 172 | text: target, |
| 173 | range: 0..0, |
| 174 | with: with.into(), |
| 175 | }, |
| 176 | }], |
| 177 | }; |
| 178 | replica.apply("Probe", typed("Local "))?; |
| 179 | let remote = { |
| 180 | let arena = onestore::Arena::default(); |
| 181 | let mut section = onestore::Section::open(&arena, source.clone())?; |
| 182 | let root = section.root(); |
| 183 | let other = section |
| 184 | .pages()? |
| 185 | .into_iter() |
| 186 | .map(|(space, ..)| space) |
| 187 | .find(|other| *other != space); |
| 188 | match other.and_then(|other| Some((other, text(&section.page(other).ok()?)?))) { |
| 189 | Some((other, text)) => section.apply( |
| 190 | "Native", |
| 191 | &Edit { |
| 192 | at: 133_000_000_000_000_000, |
| 193 | ops: vec![Op::Page { |
| 194 | space: other, |
| 195 | op: PageOp::Text { |
| 196 | text, |
| 197 | range: 0..0, |
| 198 | with: "Remote ".into(), |
| 199 | }, |
| 200 | }], |
| 201 | }, |
| 202 | )?, |
| 203 | None => section.apply( |
| 204 | "Native", |
| 205 | &Edit { |
| 206 | at: 133_000_000_000_000_000, |
| 207 | ops: vec![Op::Section(SectionOp::Create(onestore::PageCreation::new( |
| 208 | None, |
| 209 | Some("Remote"), |
| 210 | "Native", |
| 211 | )?))], |
| 212 | }, |
| 213 | )?, |
| 214 | } |
| 215 | let _ = root; |
| 216 | section.seal()?; |
| 217 | section.image() |
| 218 | }; |
| 219 | let syncing = std::sync::Arc::clone(&replica); |
| 220 | let sync = std::thread::spawn(move || { |
| 221 | let start = Instant::now(); |
| 222 | let result = syncing.sync_once(&mut Memory(remote)); |
| 223 | (start.elapsed(), result.map(|synced| synced.edit.is_some())) |
| 224 | }); |
| 225 | let (mut slowest, mut reads) = (Duration::ZERO, 0); |
| 226 | while !sync.is_finished() { |
| 227 | let start = Instant::now(); |
| 228 | replica.pages()?; |
| 229 | replica.page(space)?; |
| 230 | slowest = slowest.max(start.elapsed()); |
| 231 | reads += 1; |
| 232 | } |
| 233 | let (rebase, published) = sync.join().map_err(|_| "the sync step panicked")?; |
| 234 | println!( |
| 235 | "rebase and publish {rebase:.2?} (published: {:?}); {reads} page-list + page reads meanwhile, slowest {slowest:.2?}", |
| 236 | published? |
| 237 | ); |
| 238 | Ok(()) |
| 239 | } |