| 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 | |
| 7 | use crate::{Result, base, queue, schema}; |
| 8 | use onestore::{ |
| 9 | Arena, ExGuid, PageCreation, PageEdit, RevisionIndex, Section, Store, |
| 10 | document::Document, |
| 11 | op::{Edit, Op, SectionOp}, |
| 12 | page::Page, |
| 13 | }; |
| 14 | use rusqlite::{Connection, OptionalExtension, TransactionBehavior, params}; |
| 15 | use std::{collections::BTreeMap, io, path::Path}; |
| 16 | |
| 17 | pub(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)] |
| 22 | enum 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)] |
| 31 | struct 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)] |
| 40 | struct PageEdits { |
| 41 | edits: Vec<PageEdit>, |
| 42 | #[allow(dead_code)] |
| 43 | observed: Vec<(ExGuid, u32)>, |
| 44 | } |
| 45 | |
| 46 | fn 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. |
| 52 | fn 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 | |
| 85 | struct 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 | |
| 93 | pub(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. |
| 112 | fn 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 | |
| 164 | fn 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. |
| 334 | fn 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 | } |