1use notebook::{EditStatus, Remote, Replica};
2use onestore::{
3 Arena, CommitError, ExGuid, PageCreation, Section, Stamp, Transaction,
4 op::{Edit, Op, PageOp, SectionOp},
5 page::PageObject,
6};
7use std::{io, path::PathBuf, time::Instant};
8
9struct FileRemote(PathBuf);
10
11impl 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
26fn 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
178fn 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}