1//! Converts a schema-14 cache, whose queue holds whole-page edits over a working image,
2//! to this schema's op queue. The cache is exported first; each queued page is lowered to
3//! ops against the page the conversion has so far, and every converted page must equal
4//! the page in the old working image, or nothing changes. A page edit schema 14 held in
5//! conflict becomes the conflict page a merge makes now, the remote's page staying.
6
7use crate::{Result, base, queue, schema};
8use onestore::{
9 Arena, ExGuid, PageCreation, PageEdit, RevisionIndex, Section, Store,
10 document::Document,
11 op::{Edit, Op, SectionOp},
12 page::Page,
13};
14use rusqlite::{Connection, OptionalExtension, TransactionBehavior, params};
15use std::{collections::BTreeMap, io, path::Path};
16
17pub(crate) const VERSION: u32 = 14;
18
19/// Schema 14's queued intents, as it serialized them.
20#[derive(serde::Deserialize)]
21#[serde(deny_unknown_fields)]
22enum Operation {
23 Page(PageIntent),
24 CreatePage(PageCreation),
25 Pages(PageEdits),
26 DeletePages(Vec<ExGuid>),
27}
28
29#[derive(serde::Deserialize)]
30#[serde(deny_unknown_fields)]
31struct PageIntent {
32 #[allow(dead_code)]
33 before: Page,
34 after: Page,
35 author: String,
36}
37
38#[derive(serde::Deserialize)]
39#[serde(deny_unknown_fields)]
40struct PageEdits {
41 edits: Vec<PageEdit>,
42 #[allow(dead_code)]
43 observed: Vec<(ExGuid, u32)>,
44}
45
46fn invalid(message: String) -> crate::Error {
47 io::Error::new(io::ErrorKind::InvalidData, message).into()
48}
49
50/// Schema 14 stored the working image as its length and the runs where it differs from
51/// the base.
52fn working(mut base: Vec<u8>, patch: &[u8]) -> Result<Vec<u8>> {
53 let damaged = || invalid("Damaged working image".into());
54 let Some((length, mut runs)) = patch.split_first_chunk::<8>() else {
55 return if patch.is_empty() {
56 Ok(base)
57 } else {
58 Err(damaged())
59 };
60 };
61 let length = usize::try_from(u64::from_le_bytes(*length))
62 .ok()
63 .filter(|length| *length <= base.len() + patch.len())
64 .ok_or_else(damaged)?;
65 base.resize(length, 0);
66 while let Some((header, rest)) = runs.split_first_chunk::<16>() {
67 let offset = usize::try_from(u64::from_le_bytes(header[..8].try_into().unwrap()))
68 .map_err(|_| damaged())?;
69 let size = usize::try_from(u64::from_le_bytes(header[8..].try_into().unwrap()))
70 .map_err(|_| damaged())?;
71 let (bytes, rest) = rest.split_at_checked(size).ok_or_else(damaged)?;
72 base.get_mut(offset..)
73 .and_then(|target| target.get_mut(..size))
74 .ok_or_else(damaged)?
75 .copy_from_slice(bytes);
76 runs = rest;
77 }
78 if runs.is_empty() {
79 Ok(base)
80 } else {
81 Err(damaged())
82 }
83}
84
85struct Converted {
86 id: i64,
87 author: String,
88 edit: Edit,
89 /// Schema 14's publication evidence, when this edit was attempted.
90 attempted: Option<String>,
91}
92
93pub(crate) fn migrate(connection: &mut Connection, path: &Path) -> Result<()> {
94 let mut archive = path.as_os_str().to_owned();
95 archive.push(".v14-recovery");
96 let archive = std::path::PathBuf::from(archive);
97 crate::recovery::export(connection, &archive, true)?;
98 let failed = |error: crate::Error| {
99 invalid(format!(
100 "The schema-14 cache could not be converted ({error}); it is unchanged and archived at {}",
101 archive.display()
102 ))
103 };
104 let transaction = connection.transaction_with_behavior(TransactionBehavior::Exclusive)?;
105 convert(&transaction).map_err(failed)?;
106 transaction.commit()?;
107 Ok(())
108}
109
110/// The conflict page keeping `local`, the version of page `space` its merge with `remote`
111/// did not take, marking the texts `remote` does not hold as they are.
112fn conflict(space: ExGuid, remote: &Page, local: &Page, author: &str) -> Result<Op> {
113 let texts = |page: &Page| -> BTreeMap<ExGuid, String> {
114 let mut texts = BTreeMap::new();
115 let mut pending: Vec<&onestore::page::PageParagraph> = Vec::new();
116 for object in &page.objects {
117 match object {
118 onestore::page::PageObject::Outline(outline) => pending.extend(&outline.paragraphs),
119 onestore::page::PageObject::Title(title) => pending.extend(
120 title
121 .outlines
122 .iter()
123 .flat_map(|outline| &outline.paragraphs),
124 ),
125 _ => {}
126 }
127 }
128 while let Some(paragraph) = pending.pop() {
129 match &paragraph.content {
130 onestore::page::ParagraphContent::Text(text) => {
131 texts.insert(text.id, format!("{:?}", text.text));
132 }
133 onestore::page::ParagraphContent::Table(table) => pending.extend(
134 table
135 .rows
136 .iter()
137 .flat_map(|row| &row.cells)
138 .flat_map(|cell| &cell.paragraphs),
139 ),
140 _ => {}
141 }
142 }
143 texts
144 };
145 let kept = texts(remote);
146 let mut objects: Vec<ExGuid> = texts(local)
147 .into_iter()
148 .filter(|(id, text)| kept.get(id) != Some(text))
149 .map(|(id, _)| id)
150 .collect();
151 let page = local.copy_with(&mut objects)?;
152 let titled = local
153 .objects
154 .iter()
155 .any(|object| matches!(object, onestore::page::PageObject::Title(_)));
156 Ok(Op::Section(SectionOp::Conflict {
157 of: space,
158 creation: PageCreation::new(None, titled.then_some(local.title.as_str()), author)?,
159 page,
160 objects,
161 }))
162}
163
164fn convert(transaction: &rusqlite::Transaction<'_>) -> Result<()> {
165 let (base, patch): (Vec<u8>, Vec<u8>) =
166 transaction.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| {
167 Ok((row.get(0)?, row.get(1)?))
168 })?;
169 let expected = working(base.clone(), &patch)?;
170 let attempt: Option<(i64, String)> = transaction
171 .query_row("SELECT edit_id, revisions FROM attempt", [], |row| {
172 Ok((row.get(0)?, row.get(1)?))
173 })
174 .optional()?;
175 let conflicts: std::collections::BTreeSet<i64> = transaction
176 .prepare("SELECT edit_id FROM conflicts")?
177 .query_map([], |row| row.get(0))?
178 .collect::<rusqlite::Result<_>>()?;
179 let queued: Vec<(i64, String, String)> = transaction
180 .prepare("SELECT id, space, operation FROM edits ORDER BY id")?
181 .query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))?
182 .collect::<rusqlite::Result<_>>()?;
183 let sequence: i64 = transaction.query_row(
184 "SELECT max(coalesce((SELECT seq FROM sqlite_sequence WHERE name='edits'), 0),
185 coalesce((SELECT max(edit_id) FROM receipts), 0),
186 coalesce((SELECT max(edit_id) FROM archived), 0))",
187 [],
188 |row| row.get(0),
189 )?;
190
191 let arena = Arena::default();
192 let mut section = Section::open(&arena, base.clone())?;
193 let at = crate::now();
194 let mut converted = Vec::new();
195 let mut sealed = None;
196 let mut edited = BTreeMap::new();
197 let (mut created, mut deleted) = (Vec::new(), Vec::new());
198 for (id, space, operation) in queued {
199 let space: ExGuid = space.parse()?;
200 let operation: Operation = serde_json::from_str(&operation)
201 .map_err(|error| invalid(format!("Queued edit {id} is unreadable: {error}")))?;
202 let (author, ops) = match operation {
203 Operation::Page(intent) if conflicts.contains(&id) => {
204 let current = section.page(space)?;
205 let op = conflict(space, &current, &intent.after, &intent.author)?;
206 (intent.author, vec![op])
207 }
208 Operation::Page(intent) => {
209 let current = section.page(space)?;
210 let ops = onestore::op::lower_page(&current, &intent.after)?;
211 edited.insert(space, id);
212 (
213 intent.author,
214 ops.into_iter().map(|op| Op::Page { space, op }).collect(),
215 )
216 }
217 Operation::CreatePage(creation) => {
218 created.push(creation.space());
219 (
220 String::new(),
221 vec![Op::Section(SectionOp::Create(creation))],
222 )
223 }
224 Operation::Pages(batch) => (
225 String::new(),
226 vec![Op::Section(SectionOp::Pages(batch.edits))],
227 ),
228 Operation::DeletePages(pages) => {
229 deleted.extend(&pages);
230 (String::new(), vec![Op::Section(SectionOp::Delete(pages))])
231 }
232 };
233 let edit = Edit { at, ops };
234 section.apply(&author, &edit)?;
235 let attempted = attempt
236 .as_ref()
237 .filter(|(edit, _)| *edit == id)
238 .map(|(_, revisions)| revisions.clone());
239 if attempted.is_some() {
240 if !converted.is_empty() {
241 return Err(invalid(
242 "Only the oldest queued edit can hold an attempt".into(),
243 ));
244 }
245 sealed = Some(section.seal()?);
246 }
247 // A schema-14 conflict is queued work: the next rebase keeps both versions.
248 converted.push(Converted {
249 id,
250 author,
251 edit,
252 attempted,
253 });
254 }
255 let listed = |image: &[u8]| -> Result<Vec<ExGuid>> {
256 let store = Store::parse(image)?;
257 let index = RevisionIndex::parse(&store)?;
258 Ok(Document::parse(&index)?
259 .pages()?
260 .into_iter()
261 .map(|(space, _)| space)
262 .collect())
263 };
264 let based = listed(&base)?;
265 let listing: Vec<ExGuid> = section
266 .pages()?
267 .into_iter()
268 .map(|(space, ..)| space)
269 .collect();
270 let cached = listed(&expected)?;
271 // A page the remote added or removed while edits waited is in one list and not the
272 // other; one the queue created or deleted must agree.
273 let agree = cached
274 .iter()
275 .filter(|space| !listing.contains(space))
276 .all(|space| !based.contains(space) && !created.contains(space))
277 && listing
278 .iter()
279 .filter(|space| !cached.contains(space))
280 .all(|space| based.contains(space) && !deleted.contains(space));
281 if !agree {
282 return Err(invalid(
283 "The converted page list differs from the cached one".into(),
284 ));
285 }
286 verify(&mut section, &expected, &edited)?;
287
288 transaction.execute_batch(
289 "DROP TABLE replica; DROP TABLE attempt; DROP TABLE conflicts; DROP TABLE edits;",
290 )?;
291 transaction.execute_batch(schema::QUEUE)?;
292 base::write(transaction, base::Image::Base, &base)?;
293 let mut batch = None;
294 for edit in converted {
295 let current = match (batch, &edit.attempted) {
296 (Some(current), None) => current,
297 _ => {
298 transaction.execute("INSERT INTO batches DEFAULT VALUES", [])?;
299 transaction.last_insert_rowid()
300 }
301 };
302 if let Some(evidence) = &edit.attempted {
303 transaction.execute(
304 "UPDATE batches SET sealed=?1, revisions=?2, attempted=1 WHERE id=?3",
305 params![
306 serde_json::to_string(sealed.as_ref().unwrap()).map_err(io::Error::other)?,
307 evidence,
308 current
309 ],
310 )?;
311 batch = None;
312 } else {
313 batch = Some(current);
314 }
315 queue::insert(
316 transaction,
317 None,
318 Some(crate::unsigned(edit.id)?),
319 current,
320 &edit.author,
321 &edit.edit,
322 )?;
323 }
324 transaction.execute("DELETE FROM sqlite_sequence WHERE name='edits'", [])?;
325 transaction.execute(
326 "INSERT INTO sqlite_sequence(name, seq) VALUES ('edits', ?1)",
327 [sequence],
328 )?;
329 transaction.pragma_update(None, "user_version", schema::VERSION)?;
330 Ok(())
331}
332
333/// Requires each page a queued page edit changed to read as the old working image has it.
334fn verify(
335 section: &mut Section<'_>,
336 expected: &[u8],
337 edited: &BTreeMap<ExGuid, i64>,
338) -> Result<()> {
339 let store = Store::parse(expected)?;
340 let index = RevisionIndex::parse(&store)?;
341 let document = Document::parse(&index)?;
342 for (space, id) in edited {
343 let converted = section.page(*space).ok();
344 let stored = Page::from_space(&document, *space).ok();
345 if converted != stored {
346 return Err(invalid(format!(
347 "Queued edit {id} converts to a page that differs from the cached page"
348 )));
349 }
350 }
351 Ok(())
352}