1//! Converting schema-14 caches: whole-page edits become ops, verified page by page against
2//! the old working image, with a recovery archive first and nothing changed on a mismatch;
3//! and schema-15 caches, whose batches lose the review conflict.
4
5#[path = "../../onestore/tests/support/ops.rs"]
6mod ops;
7
8use notebook::{EditStatus, Replica};
9use onestore::{ExGuid, PageCreation, PageEdit, op::SectionOp, page::Page};
10use rusqlite::{Connection, params};
11use serde_json::json;
12use std::path::Path;
13
14#[path = "support/server.rs"]
15mod server;
16use server::*;
17#[path = "support/model_ops.rs"]
18mod model_ops;
19
20const SOURCE: &[u8] = include_bytes!("../../../corpus/outline-edit/before/notebook/synthetic.one");
21
22/// The body pages of the fixture with their first body text.
23fn body_pages(source: &[u8]) -> Vec<(ExGuid, ExGuid)> {
24 pages(source)
25 .into_iter()
26 .filter_map(|(space, page)| {
27 let text = page.objects.iter().find_map(|object| match object {
28 onestore::page::PageObject::Outline(outline) => {
29 outline.paragraphs.first()?.text().map(|text| text.id)
30 }
31 _ => None,
32 })?;
33 Some((space, text))
34 })
35 .collect()
36}
37
38/// A schema-14 queue entry.
39enum Queued {
40 Page(ExGuid, Page),
41 Create(PageCreation),
42 Pages(Vec<PageEdit>),
43 Delete(Vec<ExGuid>),
44}
45
46/// Writes a schema-14 cache holding `queue` on `base`, as that schema stored it; returns
47/// the working image.
48fn v14(path: &Path, base: &[u8], queue: &[Queued], working: Option<&[u8]>) -> Vec<u8> {
49 let mut image = base.to_vec();
50 let mut edits = Vec::new();
51 for entry in queue {
52 let (space, operation, next) = match entry {
53 Queued::Page(space, after) => {
54 let before = model_ops::page_of(&image, *space);
55 let next = ops::saved(&image, *space, after).unwrap();
56 (
57 *space,
58 json!({"Page": {"before": before, "after": after, "author": "Author"}}),
59 next.as_slice().to_vec(),
60 )
61 }
62 Queued::Create(creation) => (
63 creation.space(),
64 json!({"CreatePage": creation}),
65 ops::section_op(&image, SectionOp::Create(creation.clone()))
66 .unwrap()
67 .as_bytes()
68 .to_vec(),
69 ),
70 Queued::Pages(batch) => {
71 let observed: Vec<_> = pages(&image)
72 .into_iter()
73 .map(|(space, _)| (space, 1))
74 .collect();
75 (
76 root(&image),
77 json!({"Pages": {"edits": batch, "observed": observed}}),
78 ops::section_op(&image, SectionOp::Pages(batch.to_vec()))
79 .unwrap()
80 .as_bytes()
81 .to_vec(),
82 )
83 }
84 Queued::Delete(spaces) => (
85 root(&image),
86 json!({"DeletePages": spaces}),
87 ops::section_op(&image, SectionOp::Delete(spaces.to_vec()))
88 .unwrap()
89 .as_bytes()
90 .to_vec(),
91 ),
92 };
93 edits.push((space, operation));
94 image = next;
95 }
96 let working = working.map_or(image, <[u8]>::to_vec);
97 std::fs::File::create_new(path).unwrap();
98 let connection = Connection::open(path).unwrap();
99 connection
100 .execute_batch(
101 "PRAGMA application_id=1330529615; PRAGMA user_version=14;
102 CREATE TABLE replica (id INTEGER PRIMARY KEY CHECK(id=1), base BLOB NOT NULL, working BLOB NOT NULL) STRICT;
103 CREATE TABLE edits (id INTEGER PRIMARY KEY AUTOINCREMENT CHECK(id>0), space TEXT NOT NULL, operation TEXT NOT NULL) STRICT;
104 CREATE TABLE attempt (id INTEGER PRIMARY KEY CHECK(id=1), edit_id INTEGER NOT NULL UNIQUE REFERENCES edits(id) ON DELETE CASCADE, revisions TEXT NOT NULL) STRICT;
105 CREATE TABLE receipts (edit_id INTEGER PRIMARY KEY CHECK(edit_id>0), revision TEXT NOT NULL) STRICT;
106 CREATE TABLE archived (edit_id INTEGER PRIMARY KEY CHECK(edit_id>0), archive TEXT NOT NULL) STRICT;
107 CREATE TABLE conflicts (edit_id INTEGER PRIMARY KEY REFERENCES edits(id) ON DELETE CASCADE, kind INTEGER NOT NULL CHECK(kind BETWEEN 0 AND 3)) STRICT;
108 CREATE TABLE assets (name TEXT PRIMARY KEY NOT NULL, data BLOB NOT NULL, sha256 BLOB NOT NULL CHECK(length(sha256)=32)) STRICT;",
109 )
110 .unwrap();
111 // Schema 14 kept the working image as its length and the runs differing from the base.
112 let mut patch = (working.len() as u64).to_le_bytes().to_vec();
113 patch.extend_from_slice(&0_u64.to_le_bytes());
114 patch.extend_from_slice(&(working.len() as u64).to_le_bytes());
115 patch.extend_from_slice(&working);
116 connection
117 .execute(
118 "INSERT INTO replica VALUES (1, ?1, ?2)",
119 params![base, patch],
120 )
121 .unwrap();
122 connection
123 .execute(
124 "INSERT INTO receipts VALUES (1, ?1)",
125 [ExGuid {
126 guid: [7; 16],
127 n: 1,
128 }
129 .to_string()],
130 )
131 .unwrap();
132 connection
133 .execute(
134 "INSERT INTO sqlite_sequence(name, seq) VALUES ('edits', 1)",
135 [],
136 )
137 .unwrap();
138 for (space, operation) in edits {
139 connection
140 .execute(
141 "INSERT INTO edits(space, operation) VALUES (?1, ?2)",
142 params![space.to_string(), operation.to_string()],
143 )
144 .unwrap();
145 }
146 working
147}
148
149fn root(image: &[u8]) -> ExGuid {
150 onestore::RevisionIndex::parse(&onestore::Store::parse(image).unwrap())
151 .unwrap()
152 .root
153}
154
155fn appended(source: &[u8], space: ExGuid, text: ExGuid, suffix: &str) -> Page {
156 let mut page = model_ops::page_of(source, space);
157 let end = model_ops::paragraph_with(&page, text)
158 .unwrap()
159 .text()
160 .unwrap()
161 .text
162 .text()
163 .encode_utf16()
164 .count() as u32;
165 model_ops::replace_text(&mut page, text, end..end, suffix);
166 page
167}
168
169fn archive(path: &Path) -> std::path::PathBuf {
170 let mut archive = path.as_os_str().to_owned();
171 archive.push(".v14-recovery");
172 archive.into()
173}
174
175#[test]
176fn a_schema_14_queue_converts_to_ops_that_reach_its_working_pages_and_publish() {
177 let body = body_pages(SOURCE);
178 let ((a, a_text), (b, b_text)) = (body[0], body[1]);
179 let directory = tempfile::tempdir().unwrap();
180 let path = directory.path().join("cache.sqlite");
181 let creation = PageCreation::new(None, Some("Created offline"), "Author").unwrap();
182 let first = appended(SOURCE, a, a_text, " one");
183 let working = v14(
184 &path,
185 SOURCE,
186 &[
187 Queued::Page(a, first.clone()),
188 Queued::Page(b, appended(SOURCE, b, b_text, " other")),
189 Queued::Create(creation.clone()),
190 Queued::Pages(vec![PageEdit::move_to(b, Some(a), 1).unwrap()]),
191 Queued::Page(
192 a,
193 appended(
194 ops::saved(SOURCE, a, &first).unwrap().as_slice(),
195 a,
196 a_text,
197 " two",
198 ),
199 ),
200 Queued::Delete(vec![body[2].0]),
201 ],
202 None,
203 );
204 let cache = Replica::open(&path).unwrap();
205 assert!(archive(&path).exists());
206 let pending = cache.pending().unwrap();
207 assert_eq!(
208 pending.iter().map(|edit| edit.id).collect::<Vec<_>>(),
209 [2, 3, 4, 5, 6, 7]
210 );
211 assert!(pending.iter().all(|edit| !edit.edit.ops.is_empty()));
212 assert_eq!(
213 cache.status(1).unwrap(),
214 Some(EditStatus::Published {
215 revision: ExGuid {
216 guid: [7; 16],
217 n: 1
218 }
219 })
220 );
221 let local = pages(&snapshot(&cache));
222 assert_eq!(local, pages(&working));
223 drop(cache);
224 // A converted cache opens as it is.
225 let cache = Replica::open(&path).unwrap();
226 assert_eq!(pages(&snapshot(&cache)), local);
227 let next = cache
228 .apply(
229 "Author",
230 onestore::op::Edit {
231 at: 133_000_000_000_000_000,
232 ops: vec![onestore::op::Op::Page {
233 space: b,
234 op: onestore::op::PageOp::Text {
235 text: b_text,
236 range: 0..0,
237 with: " three".into(),
238 },
239 }],
240 },
241 )
242 .unwrap();
243 assert!(next > 7);
244 let mut server = Server::new(SOURCE);
245 assert!(matches!(
246 cache.sync_once(&mut server).unwrap().edit,
247 Some((id, EditStatus::Published { .. })) if id == next
248 ));
249 let published = pages(&server.durable);
250 assert_eq!(published, pages(&snapshot(&cache)));
251 assert!(
252 published
253 .iter()
254 .any(|(space, _)| *space == creation.space())
255 );
256 assert!(published.iter().all(|(space, _)| *space != body[2].0));
257}
258
259#[test]
260fn an_uncertain_attempt_keeps_its_state_and_a_conflict_becomes_a_conflict_page() {
261 let body = body_pages(SOURCE);
262 let (a, a_text) = body[0];
263 let directory = tempfile::tempdir().unwrap();
264
265 // The oldest edit was attempted: its schema-14 evidence names the revision it published.
266 let path = directory.path().join("attempted.sqlite");
267 let after = appended(SOURCE, a, a_text, " attempted");
268 let published = ops::saved(SOURCE, a, &after).unwrap();
269 let revision =
270 onestore::RevisionIndex::parse(&onestore::Store::parse(published.as_slice()).unwrap())
271 .unwrap()
272 .active(a)
273 .unwrap();
274 v14(&path, SOURCE, &[Queued::Page(a, after.clone())], None);
275 Connection::open(&path)
276 .unwrap()
277 .execute(
278 "INSERT INTO attempt VALUES (1, 2, ?1)",
279 [json!({a.to_string(): revision.to_string()}).to_string()],
280 )
281 .unwrap();
282 let cache = Replica::open(&path).unwrap();
283 let awaiting = EditStatus::AwaitingConfirmation { revision };
284 assert_eq!(cache.status(2).unwrap(), Some(awaiting.clone()));
285 let mut unchanged = Server::new(SOURCE);
286 assert_eq!(
287 cache.sync_once(&mut unchanged).unwrap().edit,
288 Some((2, awaiting))
289 );
290 assert_eq!(unchanged.publications, 0);
291 let mut landed = Server::new(published.as_slice());
292 assert_eq!(
293 cache.sync_once(&mut landed).unwrap().edit,
294 Some((2, EditStatus::Published { revision }))
295 );
296 assert_eq!(landed.publications, 0);
297 drop(cache);
298
299 // The oldest edit was in conflict: the conversion keeps the remote page the conflict was
300 // recorded against and queues the local version as its conflict page.
301 let path = directory.path().join("conflicted.sqlite");
302 let remote = ops::saved(SOURCE, a, &appended(SOURCE, a, a_text, " remote"))
303 .unwrap()
304 .as_slice()
305 .to_vec();
306 let local = appended(SOURCE, a, a_text, " local");
307 let working = ops::saved(SOURCE, a, &local).unwrap().as_slice().to_vec();
308 {
309 // Schema 14 stored the remote as the base and kept the local working image.
310 v14(&path, &remote, &[], Some(&working));
311 let connection = Connection::open(&path).unwrap();
312 connection
313 .execute(
314 "INSERT INTO edits(space, operation) VALUES (?1, ?2)",
315 params![
316 a.to_string(),
317 json!({"Page": {"before": model_ops::page_of(SOURCE, a), "after": local, "author": "Author"}}).to_string()
318 ],
319 )
320 .unwrap();
321 connection
322 .execute("INSERT INTO conflicts VALUES (2, 3)", [])
323 .unwrap();
324 }
325 let cache = Replica::open(&path).unwrap();
326 assert_eq!(cache.status(2).unwrap(), Some(EditStatus::Pending));
327 assert_eq!(
328 model_ops::page_of(&snapshot(&cache), a),
329 model_ops::page_of(&remote, a)
330 );
331 let mut server = Server::new(&remote);
332 assert!(matches!(
333 cache.sync_once(&mut server).unwrap().edit,
334 Some((2, EditStatus::Published { .. }))
335 ));
336 assert_eq!(
337 model_ops::page_of(&server.durable, a),
338 model_ops::page_of(&remote, a)
339 );
340 assert_eq!(
341 conflicts(&server.durable),
342 [(
343 a,
344 vec![(
345 "Author".to_owned(),
346 page_texts(&model_ops::page_of(&working, a))
347 )]
348 )]
349 );
350}
351
352#[test]
353fn a_page_the_conversion_cannot_reproduce_leaves_the_cache_untouched() {
354 let body = body_pages(SOURCE);
355 let (a, a_text) = body[0];
356 let directory = tempfile::tempdir().unwrap();
357 let path = directory.path().join("cache.sqlite");
358 // The working image holds a page the queued edit does not reach.
359 let stray = ops::saved(SOURCE, a, &appended(SOURCE, a, a_text, " stray"))
360 .unwrap()
361 .as_slice()
362 .to_vec();
363 v14(
364 &path,
365 SOURCE,
366 &[Queued::Page(a, appended(SOURCE, a, a_text, " queued"))],
367 Some(&stray),
368 );
369 let before = std::fs::read(&path).unwrap();
370 let Err(notebook::Error::Io(error)) = Replica::open(&path) else {
371 panic!("the conversion must fail")
372 };
373 assert_eq!(error.kind(), std::io::ErrorKind::InvalidData);
374 assert!(error.to_string().contains(".v14-recovery"), "{error}");
375 assert_eq!(std::fs::read(&path).unwrap(), before);
376 assert!(archive(&path).exists());
377 let version: u32 = Connection::open(archive(&path))
378 .unwrap()
379 .pragma_query_value(None, "user_version", |row| row.get(0))
380 .unwrap();
381 assert_eq!(version, 14);
382 // A second open tries again from the same untouched cache.
383 assert!(Replica::open(&path).is_err());
384 assert_eq!(std::fs::read(&path).unwrap(), before);
385}
386
387/// A schema-15 cache upgrades in place: its queue, receipts and recovery archives stay, and
388/// a conflict it recorded is queued work the next sync keeps as a conflict page.
389#[test]
390fn a_schema_15_cache_upgrades_in_place_and_its_conflict_becomes_a_conflict_page() {
391 let (space, text) = body_pages(SOURCE)[0];
392 let directory = tempfile::tempdir().unwrap();
393 let path = directory.path().join("cache.sqlite");
394 let cache = Replica::create(&path, SOURCE).unwrap();
395 let id = model_ops::save(&cache, text, |page| {
396 model_ops::replace_text(page, text, 0..0, "Local ")
397 })
398 .unwrap()
399 .unwrap();
400 let pending = cache.pending().unwrap();
401 let archive = directory.path().join("before.sqlite");
402 drop(cache);
403 {
404 // Schema 15 recorded the review conflict and its page on the batch.
405 let connection = Connection::open(&path).unwrap();
406 connection
407 .execute_batch(&format!(
408 "PRAGMA foreign_keys=OFF;
409 BEGIN;
410 CREATE TABLE old (
411 id INTEGER PRIMARY KEY AUTOINCREMENT CHECK(id>0),
412 sealed TEXT,
413 revisions TEXT,
414 attempted INTEGER NOT NULL DEFAULT 0 CHECK(attempted IN (0,1)),
415 conflict INTEGER CHECK(conflict BETWEEN 0 AND 3),
416 space TEXT,
417 CHECK((sealed IS NULL) = (revisions IS NULL)),
418 CHECK(attempted=0 OR sealed IS NOT NULL),
419 CHECK((conflict IS NULL) = (space IS NULL))
420 ) STRICT;
421 INSERT INTO old SELECT id, sealed, revisions, attempted, 3, '{space}' FROM batches;
422 DROP TABLE batches;
423 ALTER TABLE old RENAME TO batches;
424 PRAGMA user_version=15;
425 COMMIT;"
426 ))
427 .unwrap();
428 }
429 let cache = Replica::open(&path).unwrap();
430 assert_eq!(cache.pending().unwrap(), pending);
431 assert_eq!(cache.status(id).unwrap(), Some(EditStatus::Pending));
432 cache.export_recovery(&archive).unwrap();
433 drop(cache);
434 let connection = Connection::open(&path).unwrap();
435 let version: u32 = connection
436 .pragma_query_value(None, "user_version", |row| row.get(0))
437 .unwrap();
438 let columns: Vec<String> = connection
439 .prepare("SELECT name FROM pragma_table_info('batches')")
440 .unwrap()
441 .query_map([], |row| row.get(0))
442 .unwrap()
443 .collect::<Result<_, _>>()
444 .unwrap();
445 assert_eq!(version, 16);
446 assert_eq!(columns, ["id", "sealed", "revisions", "attempted"]);
447 drop(connection);
448 assert_eq!(
449 notebook::Recovery::open(&archive)
450 .unwrap()
451 .pending()
452 .unwrap(),
453 pending
454 );
455 let cache = Replica::open(&path).unwrap();
456 let remote = typed(SOURCE, space, text, 0..0, "Remote ");
457 let mut server = Server::new(&remote);
458 assert!(matches!(
459 cache.sync_once(&mut server).unwrap().edit,
460 Some((_, EditStatus::Published { .. }))
461 ));
462 assert!(conflicted(&server.durable, space));
463}