| 1 | use notebook::{EditStatus, Remote, Replica}; |
| 2 | use onestore::{ |
| 3 | Arena, CommitError, ExGuid, PageCreation, Section, Stamp, Transaction, |
| 4 | op::{Edit, Op, PageOp, SectionOp}, |
| 5 | page::PageObject, |
| 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 stamp(&mut self) -> io::Result<Stamp> { |
| 16 | Stamp::of(&self.read()?).map_err(io::Error::other) |
| 17 | } |
| 18 | fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> { |
| 19 | transaction.commit_file(&self.0) |
| 20 | } |
| 21 | fn confirm(&mut self, base: &Stamp) -> Result<(), CommitError> { |
| 22 | onestore::confirm_file(&self.0, base) |
| 23 | } |
| 24 | } |
| 25 | |
| 26 | fn main() -> Result<(), Box<dyn std::error::Error>> { |
| 27 | let args: Vec<_> = std::env::args_os().skip(1).collect(); |
| 28 | assert_eq!(args.len(), 2, "queue_scale NEW_DIRECTORY EDIT_COUNT"); |
| 29 | let directory = PathBuf::from(&args[0]); |
| 30 | let count: usize = args[1].to_str().unwrap().parse()?; |
| 31 | assert!(count >= 2); |
| 32 | std::fs::create_dir(&directory)?; |
| 33 | let source = onestore::create_section("queue.one", "First", "Fixture")?; |
| 34 | let creation = PageCreation::new(None, Some("Second"), "Fixture")?; |
| 35 | let (source, pages) = { |
| 36 | let arena = Arena::default(); |
| 37 | let mut section = Section::open(&arena, source)?; |
| 38 | let ops = vec![Op::Section(SectionOp::Create(creation))]; |
| 39 | section.apply( |
| 40 | "Fixture", |
| 41 | &Edit { |
| 42 | at: 133_000_000_000_000_000, |
| 43 | ops, |
| 44 | }, |
| 45 | )?; |
| 46 | section.seal()?; |
| 47 | let pages: Vec<ExGuid> = section |
| 48 | .pages()? |
| 49 | .into_iter() |
| 50 | .map(|(space, ..)| space) |
| 51 | .collect(); |
| 52 | (section.image(), pages) |
| 53 | }; |
| 54 | assert_eq!(pages.len(), 2); |
| 55 | let path = directory.join("cache.sqlite"); |
| 56 | let mut remote = FileRemote(directory.join("queue.one")); |
| 57 | std::fs::write(&remote.0, &source)?; |
| 58 | let cache = Replica::create(&path, &source)?; |
| 59 | let mut ids = Vec::new(); |
| 60 | let mut expected = [String::new(), String::new()]; |
| 61 | let start = Instant::now(); |
| 62 | for n in 0..count { |
| 63 | let edit_start = Instant::now(); |
| 64 | let slot = n % 2; |
| 65 | let space = pages[slot]; |
| 66 | let page = cache.page(space)?; |
| 67 | let text = page |
| 68 | .objects |
| 69 | .iter() |
| 70 | .find_map(|object| match object { |
| 71 | PageObject::Title(title) => title |
| 72 | .outlines |
| 73 | .iter() |
| 74 | .flat_map(|outline| &outline.paragraphs) |
| 75 | .find_map(|p| p.text().filter(|text| text.date_field.is_none())), |
| 76 | _ => None, |
| 77 | }) |
| 78 | .unwrap(); |
| 79 | expected[slot] = format!("Edit {n} 🦀 e\u{301}"); |
| 80 | let end = u32::try_from(text.text.text().encode_utf16().count())?; |
| 81 | let op = PageOp::Text { |
| 82 | text: text.id, |
| 83 | range: 0..end, |
| 84 | with: expected[slot].clone(), |
| 85 | }; |
| 86 | let id = cache.apply( |
| 87 | "Fixture", |
| 88 | Edit { |
| 89 | at: 133_000_000_000_000_000, |
| 90 | ops: vec![Op::Page { space, op }], |
| 91 | }, |
| 92 | )?; |
| 93 | assert!(ids.last().is_none_or(|previous| *previous < id)); |
| 94 | ids.push(id); |
| 95 | println!( |
| 96 | "{}", |
| 97 | serde_json::json!({"phase":"ack","n":n,"id":id, |
| 98 | "ms":edit_start.elapsed().as_secs_f64()*1000.0}) |
| 99 | ); |
| 100 | } |
| 101 | println!( |
| 102 | "{}", |
| 103 | serde_json::json!({"phase":"queued","count":count, |
| 104 | "seconds":start.elapsed().as_secs_f64(),"cache_bytes":std::fs::metadata(&path)?.len()}) |
| 105 | ); |
| 106 | drop(cache); |
| 107 | let start = Instant::now(); |
| 108 | let cache = Replica::open(&path)?; |
| 109 | assert_eq!( |
| 110 | cache |
| 111 | .pending()? |
| 112 | .iter() |
| 113 | .map(|edit| edit.id) |
| 114 | .collect::<Vec<_>>(), |
| 115 | ids |
| 116 | ); |
| 117 | verify(&cache, &pages, &expected)?; |
| 118 | println!( |
| 119 | "{}", |
| 120 | serde_json::json!({"phase":"reopen","seconds":start.elapsed().as_secs_f64()}) |
| 121 | ); |
| 122 | let start = Instant::now(); |
| 123 | cache.export_recovery(directory.join("recovery.sqlite"))?; |
| 124 | println!( |
| 125 | "{}", |
| 126 | serde_json::json!({"phase":"export","seconds":start.elapsed().as_secs_f64()}) |
| 127 | ); |
| 128 | let start = Instant::now(); |
| 129 | // The queue publishes as one batch. |
| 130 | let step = Instant::now(); |
| 131 | assert!(matches!( |
| 132 | cache.sync_once(&mut remote)?.edit, |
| 133 | Some((actual, EditStatus::Published { .. })) if Some(&actual) == ids.last() |
| 134 | )); |
| 135 | println!( |
| 136 | "{}", |
| 137 | serde_json::json!({"phase":"publish","count":ids.len(),"ms":step.elapsed().as_secs_f64()*1000.0}) |
| 138 | ); |
| 139 | assert!(cache.pending()?.is_empty()); |
| 140 | let published = onestore::read_file(&remote.0)?; |
| 141 | verify(&cache, &pages, &expected)?; |
| 142 | let arena = Arena::default(); |
| 143 | let section = Section::open(&arena, published.clone())?; |
| 144 | for space in &pages { |
| 145 | assert_eq!(section.page(*space)?, cache.page(*space)?); |
| 146 | } |
| 147 | drop(cache); |
| 148 | let cache = Replica::open(&path)?; |
| 149 | for id in &ids { |
| 150 | assert!(matches!( |
| 151 | cache.status(*id)?, |
| 152 | Some(EditStatus::Published { .. }) |
| 153 | )); |
| 154 | } |
| 155 | let archive = notebook::Recovery::open(directory.join("recovery.sqlite"))?; |
| 156 | assert_eq!( |
| 157 | archive |
| 158 | .pending()? |
| 159 | .iter() |
| 160 | .map(|edit| edit.id) |
| 161 | .collect::<Vec<_>>(), |
| 162 | ids |
| 163 | ); |
| 164 | let arena = Arena::default(); |
| 165 | let section = Section::open(&arena, archive.snapshot()?)?; |
| 166 | for (slot, space) in pages.iter().enumerate() { |
| 167 | assert_eq!(section.page(*space)?.title, expected[slot]); |
| 168 | } |
| 169 | println!( |
| 170 | "{}", |
| 171 | serde_json::json!({"phase":"complete","count":count, |
| 172 | "publish_seconds":start.elapsed().as_secs_f64(),"cache_bytes":std::fs::metadata(&path)?.len(), |
| 173 | "remote_bytes":published.len()}) |
| 174 | ); |
| 175 | Ok(()) |
| 176 | } |
| 177 | |
| 178 | fn verify( |
| 179 | cache: &Replica, |
| 180 | pages: &[ExGuid], |
| 181 | expected: &[String; 2], |
| 182 | ) -> Result<(), Box<dyn std::error::Error>> { |
| 183 | for (slot, space) in pages.iter().enumerate() { |
| 184 | assert_eq!(cache.page(*space)?.title, expected[slot]); |
| 185 | } |
| 186 | Ok(()) |
| 187 | } |