| ... | @@ -0,0 +1,172 @@ |
| 1 | use notebook::{EditStatus, Remote, Replica}; |
| 2 | use onestore::{ |
| 3 | CommitError, PageCreation, PreparedEdit, RevisionIndex, Store, |
| 4 | document::Document, |
| 5 | page::{Page, PageObject, Paragraph, text::Edit}, |
| 6 | }; |
| 7 | use std::{io, path::PathBuf, time::Instant}; |
| 8 | |
| 9 | struct FileRemote(PathBuf); |
| 10 | |
| 11 | impl Remote for FileRemote { |
| 12 | fn read(&mut self) -> io::Result<Vec<u8>> { |
| 13 | onestore::read_file(&self.0) |
| 14 | } |
| 15 | fn publish(&mut self, edit: &PreparedEdit<'_>) -> Result<(), CommitError> { |
| 16 | edit.commit_file(&self.0) |
| 17 | } |
| 18 | fn confirm(&mut self, snapshot: &[u8]) -> Result<(), CommitError> { |
| 19 | onestore::confirm_file_snapshot(&self.0, snapshot) |
| 20 | } |
| 21 | } |
| 22 | |
| 23 | fn main() -> Result<(), Box<dyn std::error::Error>> { |
| 24 | let args: Vec<_> = std::env::args_os().skip(1).collect(); |
| 25 | assert_eq!(args.len(), 2, "queue_scale NEW_DIRECTORY EDIT_COUNT"); |
| 26 | let directory = PathBuf::from(&args[0]); |
| 27 | let count: usize = args[1].to_str().unwrap().parse()?; |
| 28 | assert!(count >= 2); |
| 29 | std::fs::create_dir(&directory)?; |
| 30 | let source = onestore::create_section("queue.one", "First", "Fixture")?; |
| 31 | let creation = PageCreation::new(None, Some("Second"), "Fixture")?; |
| 32 | let source = PreparedEdit::create_page(&source, &creation)? |
| 33 | .as_bytes() |
| 34 | .to_vec(); |
| 35 | let pages = { |
| 36 | let store = Store::parse(&source)?; |
| 37 | let index = RevisionIndex::parse(&store)?; |
| 38 | Document::parse(&index)?.pages()? |
| 39 | }; |
| 40 | assert_eq!(pages.len(), 2); |
| 41 | let path = directory.join("cache.sqlite"); |
| 42 | let mut remote = FileRemote(directory.join("queue.one")); |
| 43 | std::fs::write(&remote.0, &source)?; |
| 44 | let cache = Replica::create(&path, &source)?; |
| 45 | let mut ids = Vec::new(); |
| 46 | let mut expected = [String::new(), String::new()]; |
| 47 | let start = Instant::now(); |
| 48 | for n in 0..count { |
| 49 | let edit_start = Instant::now(); |
| 50 | let source = cache.snapshot()?; |
| 51 | let store = Store::parse(&source)?; |
| 52 | let index = RevisionIndex::parse(&store)?; |
| 53 | let document = Document::parse(&index)?; |
| 54 | let slot = n % 2; |
| 55 | let space = pages[slot].0; |
| 56 | let mut page = Page::from_space(&document, space)?; |
| 57 | let text = page |
| 58 | .objects |
| 59 | .iter_mut() |
| 60 | .find_map(|object| match object { |
| 61 | PageObject::Outline(outline) => { |
| 62 | outline.paragraphs.iter_mut().find_map(|p| p.text_mut()) |
| 63 | } |
| 64 | PageObject::Title(title) => title |
| 65 | .outlines |
| 66 | .iter_mut() |
| 67 | .flat_map(|outline| &mut outline.paragraphs) |
| 68 | .find_map(|p| p.text_mut().filter(|text| text.date_field.is_none())), |
| 69 | _ => None, |
| 70 | }) |
| 71 | .unwrap(); |
| 72 | expected[slot] = format!("Edit {n} 🦀 e\u{301}"); |
| 73 | let end = u32::try_from(text.text.text().encode_utf16().count())?; |
| 74 | let format = text.text.format_at(0)?.clone(); |
| 75 | text.text.apply(Edit { |
| 76 | range: 0..end, |
| 77 | replacement: Paragraph::new(expected[slot].clone(), format), |
| 78 | })?; |
| 79 | let id = cache.save(&source, space, &page, "Fixture")?.unwrap(); |
| 80 | assert!(ids.last().is_none_or(|previous| *previous < id)); |
| 81 | ids.push(id); |
| 82 | println!( |
| 83 | "{}", |
| 84 | serde_json::json!({"phase":"ack","n":n,"id":id, |
| 85 | "ms":edit_start.elapsed().as_secs_f64()*1000.0,"working_bytes":source.len()}) |
| 86 | ); |
| 87 | } |
| 88 | println!( |
| 89 | "{}", |
| 90 | serde_json::json!({"phase":"queued","count":count, |
| 91 | "seconds":start.elapsed().as_secs_f64(),"cache_bytes":std::fs::metadata(&path)?.len()}) |
| 92 | ); |
| 93 | drop(cache); |
| 94 | let start = Instant::now(); |
| 95 | let cache = Replica::open(&path)?; |
| 96 | assert_eq!( |
| 97 | cache |
| 98 | .pending()? |
| 99 | .iter() |
| 100 | .map(|edit| edit.id) |
| 101 | .collect::<Vec<_>>(), |
| 102 | ids |
| 103 | ); |
| 104 | let source = cache.snapshot()?; |
| 105 | verify(&source, &pages, &expected)?; |
| 106 | println!( |
| 107 | "{}", |
| 108 | serde_json::json!({"phase":"reopen","seconds":start.elapsed().as_secs_f64()}) |
| 109 | ); |
| 110 | let start = Instant::now(); |
| 111 | cache.export_recovery(directory.join("recovery.sqlite"))?; |
| 112 | println!( |
| 113 | "{}", |
| 114 | serde_json::json!({"phase":"export","seconds":start.elapsed().as_secs_f64()}) |
| 115 | ); |
| 116 | let start = Instant::now(); |
| 117 | for (n, id) in ids.iter().enumerate() { |
| 118 | let step = Instant::now(); |
| 119 | assert!( |
| 120 | matches!(cache.sync_once(&mut remote)?, Some((actual, EditStatus::Published { .. })) if actual == *id) |
| 121 | ); |
| 122 | println!( |
| 123 | "{}", |
| 124 | serde_json::json!({"phase":"publish","n":n,"id":id,"ms":step.elapsed().as_secs_f64()*1000.0}) |
| 125 | ); |
| 126 | } |
| 127 | assert!(cache.pending()?.is_empty()); |
| 128 | let published = cache.snapshot()?; |
| 129 | assert_eq!(onestore::read_file(&remote.0)?, published); |
| 130 | verify(&published, &pages, &expected)?; |
| 131 | drop(cache); |
| 132 | let cache = Replica::open(&path)?; |
| 133 | for id in &ids { |
| 134 | assert!(matches!( |
| 135 | cache.status(*id)?, |
| 136 | Some(EditStatus::Published { .. }) |
| 137 | )); |
| 138 | } |
| 139 | let archive = notebook::Recovery::open(directory.join("recovery.sqlite"))?; |
| 140 | assert_eq!( |
| 141 | archive |
| 142 | .pending()? |
| 143 | .iter() |
| 144 | .map(|edit| edit.id) |
| 145 | .collect::<Vec<_>>(), |
| 146 | ids |
| 147 | ); |
| 148 | assert_eq!(archive.snapshot()?, source); |
| 149 | println!( |
| 150 | "{}", |
| 151 | serde_json::json!({"phase":"complete","count":count, |
| 152 | "publish_seconds":start.elapsed().as_secs_f64(),"cache_bytes":std::fs::metadata(&path)?.len(), |
| 153 | "remote_bytes":published.len()}) |
| 154 | ); |
| 155 | Ok(()) |
| 156 | } |
| 157 | |
| 158 | fn verify( |
| 159 | bytes: &[u8], |
| 160 | pages: &[(onestore::ExGuid, onestore::ExGuid)], |
| 161 | expected: &[String; 2], |
| 162 | ) -> Result<(), Box<dyn std::error::Error>> { |
| 163 | let store = Store::parse(bytes)?; |
| 164 | assert!(store.checksum_mismatches.is_empty()); |
| 165 | let index = RevisionIndex::parse(&store)?; |
| 166 | index.validate_current()?; |
| 167 | let document = Document::parse(&index)?; |
| 168 | for (slot, (space, _)) in pages.iter().enumerate() { |
| 169 | assert_eq!(Page::from_space(&document, *space)?.title, expected[slot]); |
| 170 | } |
| 171 | Ok(()) |
| 172 | } |