diff --git a/crates/notebook/README.md b/crates/notebook/README.md index 0371f3ecdb671ae0ad48befe308b0b7d50ac5882..90b37cc4d67c08f711224afc9273858290938e88 100644 --- a/crates/notebook/README.md +++ b/crates/notebook/README.md @@ -2,8 +2,9 @@ The application-facing crate: notebook discovery, a durable local replica with reconnect reconciliation, external-asset caching, recovery export and optional -embedded SMB access. The replica stores a complete working image of one section -and the editor's intents in a local SQLite database. `sync_once` provides a +embedded SMB access. The replica stores the last observed remote image of one section, +the working image as its difference from that, and the editor's intents in a +local SQLite database; a save writes the revision it appended, not the section. `sync_once` provides a reconciliation step and `start_sync` owns automatic polling and reconnects. Local success does not acknowledge publication to a shared notebook. diff --git a/crates/notebook/src/assets.rs b/crates/notebook/src/assets.rs index 8cc13e69d213b64b4127e43ab78ec1ba2fe115ca..d176734210a5156d5d41c39184272841ab989c3f 100644 --- a/crates/notebook/src/assets.rs +++ b/crates/notebook/src/assets.rs @@ -79,12 +79,8 @@ pub(crate) fn key(filename: &str) -> Result { } fn referenced(connection: &Connection, key: &str) -> Result { - for column in ["working", "base"] { - let image: Vec = connection.query_row( - &format!("SELECT {column} FROM replica WHERE id=1"), - [], - |row| row.get(0), - )?; + let (base, working) = crate::images::both(connection)?; + for image in [working, base] { let store = Store::parse(&image)?; let index = RevisionIndex::parse(&store)?; let document = Document::parse(&index)?; diff --git a/crates/notebook/src/images.rs b/crates/notebook/src/images.rs new file mode 100644 index 0000000000000000000000000000000000000000..87b2c6111b2eefb5acf39dc8b8fee69074866ca5 --- /dev/null +++ b/crates/notebook/src/images.rs @@ -0,0 +1,151 @@ +//! The cache's two section images. `base` is stored whole; `working` is stored as the byte +//! ranges where it differs from `base`, so saving an edit writes the revision it appended +//! instead of the section. + +use crate::Result; +use rusqlite::{Connection, params}; +use std::io; + +const BLOCK: usize = 4096; + +/// `image` as its length and the `(offset, bytes)` runs of blocks that differ from `base`; +/// empty when the images are equal. +fn difference(base: &[u8], image: &[u8]) -> Vec { + if base == image { + return Vec::new(); + } + let mut patch = (image.len() as u64).to_le_bytes().to_vec(); + let mut run: Option = None; + let close = |patch: &mut Vec, start: usize, end: usize| { + patch.extend_from_slice(&(start as u64).to_le_bytes()); + patch.extend_from_slice(&((end - start) as u64).to_le_bytes()); + patch.extend_from_slice(&image[start..end]); + }; + for start in (0..image.len()).step_by(BLOCK) { + let end = (start + BLOCK).min(image.len()); + if base.get(start..end) == Some(&image[start..end]) { + if let Some(from) = run.take() { + close(&mut patch, from, start); + } + } else { + run.get_or_insert(start); + } + } + if let Some(from) = run { + close(&mut patch, from, image.len()); + } + patch +} + +fn restore(mut base: Vec, patch: &[u8]) -> Result> { + let damaged = || io::Error::new(io::ErrorKind::InvalidData, "Damaged working image"); + let Some((length, mut runs)) = patch.split_first_chunk::<8>() else { + return if patch.is_empty() { + Ok(base) + } else { + Err(damaged().into()) + }; + }; + base.resize( + usize::try_from(u64::from_le_bytes(*length)).map_err(|_| damaged())?, + 0, + ); + while let Some((header, rest)) = runs.split_first_chunk::<16>() { + let offset = usize::try_from(u64::from_le_bytes(header[..8].try_into().unwrap())) + .map_err(|_| damaged())?; + let size = usize::try_from(u64::from_le_bytes(header[8..].try_into().unwrap())) + .map_err(|_| damaged())?; + let (bytes, rest) = rest.split_at_checked(size).ok_or_else(damaged)?; + base.get_mut(offset..) + .and_then(|target| target.get_mut(..size)) + .ok_or_else(damaged)? + .copy_from_slice(bytes); + runs = rest; + } + if runs.is_empty() { + Ok(base) + } else { + Err(damaged().into()) + } +} + +pub(crate) fn base(connection: &Connection) -> Result> { + Ok(connection.query_row("SELECT base FROM replica WHERE id=1", [], |row| row.get(0))?) +} + +/// `(base, working)`. +pub(crate) fn both(connection: &Connection) -> Result<(Vec, Vec)> { + let (base, patch): (Vec, Vec) = + connection.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { + Ok((row.get(0)?, row.get(1)?)) + })?; + let working = restore(base.clone(), &patch)?; + Ok((base, working)) +} + +pub(crate) fn working(connection: &Connection) -> Result> { + let (base, patch): (Vec, Vec) = + connection.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { + Ok((row.get(0)?, row.get(1)?)) + })?; + restore(base, &patch) +} + +/// Replaces the working image; `base` is the stored base image. +pub(crate) fn set_working(connection: &Connection, base: &[u8], working: &[u8]) -> Result<()> { + connection.execute( + "UPDATE replica SET working=?1 WHERE id=1", + [difference(base, working)], + )?; + Ok(()) +} + +/// Replaces the base image, keeping `working` as the working image. An unchanged base +/// leaves its stored bytes alone. +pub(crate) fn set_base( + connection: &Connection, + old: &[u8], + base: &[u8], + working: &[u8], +) -> Result<()> { + if old == base { + return set_working(connection, base, working); + } + connection.execute( + "UPDATE replica SET base=?1, working=?2 WHERE id=1", + params![base, difference(base, working)], + )?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn a_working_image_survives_as_its_difference_from_any_base() { + let base: Vec = (0..3 * BLOCK + 17).map(|i| i as u8).collect(); + let mut edited = base.clone(); + edited[5] ^= 1; + edited[2 * BLOCK + 1] ^= 1; + edited.extend_from_slice(b"appended revision"); + for image in [ + base.clone(), + edited.clone(), + base[..BLOCK + 3].to_vec(), + Vec::new(), + vec![7; 5 * BLOCK], + ] { + let patch = difference(&base, &image); + assert_eq!(restore(base.clone(), &patch).unwrap(), image); + } + assert!(difference(&base, &base).is_empty()); + assert!(difference(&base, &edited).len() < 3 * BLOCK); + let patch = difference(&base, &edited); + for cut in 1..patch.len() { + if let Ok(image) = restore(base.clone(), &patch[..cut]) { + assert_ne!(image, edited); + } + } + } +} diff --git a/crates/notebook/src/lib.rs b/crates/notebook/src/lib.rs index e35e6a7e4c8c52a22a195225f4bb0480bfec3ab6..c8dde05afef867990772847468788668854917d9 100644 --- a/crates/notebook/src/lib.rs +++ b/crates/notebook/src/lib.rs @@ -20,6 +20,7 @@ use std::{ }; mod assets; +mod images; mod merge; mod pages; mod rebase; @@ -59,7 +60,7 @@ pub enum Error { type Result = std::result::Result; const APPLICATION_ID: u32 = 0x4f4e454f; -const SCHEMA_VERSION: u32 = 13; +const SCHEMA_VERSION: u32 = 14; /// An edited page model together with the stored model it was edited from. /// `before` is the precondition reconciliation checks against the remote page. @@ -217,7 +218,7 @@ impl Replica { ", )?; schema::create(&transaction)?; - transaction.execute("INSERT INTO replica VALUES (1, ?1, ?1)", [source])?; + transaction.execute("INSERT INTO replica VALUES (1, ?1, x'')", [source])?; } else { if application != APPLICATION_ID { return Err( @@ -251,11 +252,7 @@ impl Replica { .connection .lock() .map_err(|_| io::Error::other("Cache owner panicked"))?; - Ok( - connection.query_row("SELECT working FROM replica WHERE id=1", [], |row| { - row.get(0) - })?, - ) + images::working(&connection) } pub fn pending(&self) -> Result> { @@ -294,10 +291,7 @@ impl Replica { .map_err(|_| io::Error::other("Cache owner panicked"))?; let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?; - let current: Vec = - transaction.query_row("SELECT working FROM replica WHERE id=1", [], |row| { - row.get(0) - })?; + let (base, current) = images::both(&transaction)?; if current != source { return Err(io::Error::new( io::ErrorKind::ResourceBusy, @@ -334,10 +328,7 @@ impl Replica { id ], )?; - transaction.execute( - "UPDATE replica SET working=?1 WHERE id=1", - [prepared.as_bytes()], - )?; + images::set_working(&transaction, &base, prepared.as_bytes())?; transaction.commit()?; drop(connection); self.wake_sync(); @@ -388,10 +379,7 @@ impl Replica { .lock() .map_err(|_| io::Error::other("Cache owner panicked"))?; let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?; - let current: Vec = - transaction.query_row("SELECT working FROM replica WHERE id=1", [], |row| { - row.get(0) - })?; + let (base, current) = images::both(&transaction)?; if current != source { return Err(io::Error::new( io::ErrorKind::ResourceBusy, @@ -410,10 +398,7 @@ impl Replica { ], )?; let id = u64::try_from(transaction.last_insert_rowid()).map_err(io::Error::other)?; - transaction.execute( - "UPDATE replica SET working=?1 WHERE id=1", - [edit.as_bytes()], - )?; + images::set_working(&transaction, &base, edit.as_bytes())?; transaction.commit()?; drop(connection); self.wake_sync(); @@ -514,10 +499,7 @@ fn validate_images(connection: &Connection) -> Result<()> { io::Error::new(io::ErrorKind::InvalidData, "Cache integrity check failed").into(), ); } - let (base, working): (Vec, Vec) = - connection.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { - Ok((row.get(0)?, row.get(1)?)) - })?; + let (base, working) = images::both(connection)?; if validate(&base)? != validate(&working)? { return Err(io::Error::new( io::ErrorKind::InvalidData, diff --git a/crates/notebook/src/recovery.rs b/crates/notebook/src/recovery.rs index 48a5dd884aaae2d251dfaa59e84f30ac9821c326..499aefd668821636be875d7b95e793795bcf2663 100644 --- a/crates/notebook/src/recovery.rs +++ b/crates/notebook/src/recovery.rs @@ -58,17 +58,11 @@ impl Recovery { } pub fn snapshot(&self) -> Result> { - Ok(self - .connection - .query_row("SELECT working FROM replica WHERE id=1", [], |row| { - row.get(0) - })?) + crate::images::working(&self.connection) } pub fn remote_snapshot(&self) -> Result> { - Ok(self - .connection - .query_row("SELECT base FROM replica WHERE id=1", [], |row| row.get(0))?) + crate::images::base(&self.connection) } pub fn pending(&self) -> Result> { @@ -150,10 +144,11 @@ fn summary(connection: &Connection) -> Result { [], |row| Ok((unsigned(row, 0)?, unsigned(row, 1)?)), )?; + let working_bytes = crate::images::working(connection)?.len() as u64; Ok(connection.query_row( "SELECT (SELECT count(*) FROM edits), (SELECT count(*) FROM conflicts), (SELECT count(*) FROM attempt), (SELECT count(*) FROM receipts), - length(working), length(base) FROM replica WHERE id=1", + length(base) FROM replica WHERE id=1", [], |row| { Ok(RecoverySummary { @@ -161,8 +156,8 @@ fn summary(connection: &Connection) -> Result { conflicts: unsigned(row, 1)?, uncertain_edits: unsigned(row, 2)?, published_receipts: unsigned(row, 3)?, - working_bytes: unsigned(row, 4)?, - remote_bytes: unsigned(row, 5)?, + working_bytes, + remote_bytes: unsigned(row, 4)?, cached_assets, cached_asset_bytes, }) diff --git a/crates/notebook/src/sync.rs b/crates/notebook/src/sync.rs index 5093fbe14797971599f6262ef65685d9bc65ca27..037b91b75c9fde1e9dec8d339d98d864bd911b37 100644 --- a/crates/notebook/src/sync.rs +++ b/crates/notebook/src/sync.rs @@ -63,7 +63,7 @@ impl Replica { .connection .lock() .map_err(|_| io::Error::other("Cache owner panicked"))?; - Ok(connection.query_row("SELECT base FROM replica WHERE id=1", [], |row| row.get(0))?) + images::base(&connection) } /// Reconciles one pending edit, or refreshes the working image when the queue is empty. @@ -79,7 +79,12 @@ impl Replica { } let _step = Step(self); let snapshot = remote.read().map_err(Error::RemoteIo)?; - let identity = validate(&snapshot)?; + // An unchanged remote is the base image, validated when it was stored. + let identity = if snapshot == self.remote_snapshot()? { + None + } else { + Some(validate(&snapshot)?) + }; let (intent, attempted) = { let mut connection = self .connection @@ -87,11 +92,10 @@ impl Replica { .map_err(|_| io::Error::other("Cache owner panicked"))?; let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?; - let base: Vec = - transaction - .query_row("SELECT base FROM replica WHERE id=1", [], |row| row.get(0))?; - let base_store = Store::parse(&base)?; - if RevisionIndex::parse(&base_store)?.root != identity { + let (base, working) = images::both(&transaction)?; + if let Some(identity) = identity + && RevisionIndex::parse(&Store::parse(&base)?)?.root != identity + { return Err(io::Error::new( io::ErrorKind::InvalidInput, "Remote snapshot belongs to another document", @@ -105,10 +109,7 @@ impl Replica { rows.next()?.map(pending_edit).transpose()? }; let Some(intent) = intent else { - transaction.execute( - "UPDATE replica SET base=?1, working=?1 WHERE id=1", - [&snapshot], - )?; + images::set_base(&transaction, &base, &snapshot, &snapshot)?; transaction.commit()?; return Ok(None); }; @@ -123,7 +124,7 @@ impl Replica { |row| row.get::<_, String>(0), ) .optional()?; - transaction.execute("UPDATE replica SET base=?1 WHERE id=1", [&snapshot])?; + images::set_base(&transaction, &base, &snapshot, &working)?; transaction.commit()?; (intent, attempted) }; @@ -415,10 +416,7 @@ impl Replica { .lock() .map_err(|_| io::Error::other("Cache owner panicked"))?; let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?; - let (base, working): (Vec, Vec) = - transaction.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { - Ok((row.get(0)?, row.get(1)?)) - })?; + let (base, working) = images::both(&transaction)?; if working != local || base != remote { return Err(io::Error::new( io::ErrorKind::ResourceBusy, @@ -464,7 +462,7 @@ impl Replica { [archive.to_string_lossy().into_owned()], )?; transaction.execute("DELETE FROM edits", [])?; - transaction.execute("UPDATE replica SET working=base WHERE id=1", [])?; + transaction.execute("UPDATE replica SET working=x'' WHERE id=1", [])?; } } transaction.commit()?; @@ -499,10 +497,7 @@ impl Replica { .lock() .map_err(|_| io::Error::other("Cache owner panicked"))?; let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?; - let (base, working): (Vec, Vec) = - transaction.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { - Ok((row.get(0)?, row.get(1)?)) - })?; + let (base, working) = images::both(&transaction)?; if working != local || base != remote { return Err(io::Error::new( io::ErrorKind::ResourceBusy, @@ -561,7 +556,15 @@ impl Replica { params![id, revision.to_string()], )?; transaction.execute("DELETE FROM edits WHERE id=?1", [id])?; - transaction.execute("UPDATE replica SET base=?1, working=CASE WHEN EXISTS(SELECT 1 FROM edits) THEN working ELSE ?1 END WHERE id=1", [snapshot])?; + let (base, working) = images::both(&transaction)?; + let pending: bool = + transaction.query_row("SELECT EXISTS(SELECT 1 FROM edits)", [], |row| row.get(0))?; + images::set_base( + &transaction, + &base, + snapshot, + if pending { &working } else { snapshot }, + )?; transaction.commit()?; Ok(()) } diff --git a/crates/onestore/src/edit.rs b/crates/onestore/src/edit.rs index f28739636606a7f56878110da26fa6e417cc7ed2..32ecf570a446a1bce57570dc99ec74f75158e23b 100644 --- a/crates/onestore/src/edit.rs +++ b/crates/onestore/src/edit.rs @@ -237,7 +237,7 @@ pub fn replace_text( let Some((page, automatic, title_text)) = page_title(revision, &pages, Some((object, &changed)))? else { - return crate::write::replace_objects(source, space, &edits); + return crate::write::replace_objects(&index, space, &edits); }; let Kind::Page { alternate_title, .. @@ -296,7 +296,7 @@ pub fn replace_text( &[] }, }); - crate::write::replace_objects(source, space, &edits) + crate::write::replace_objects(&index, space, &edits) } pub(crate) fn editable_parents( diff --git a/crates/onestore/src/formatting.rs b/crates/onestore/src/formatting.rs index a1bae1fb1b91785ba316abdf8bb6b9767d8da146..428c22dbb23081b59fbd4e449b85441c985533f4 100644 --- a/crates/onestore/src/formatting.rs +++ b/crates/onestore/src/formatting.rs @@ -3,7 +3,7 @@ use crate::{ create::{current_timestamps, properties, string}, document::{Document, Kind}, edit::editable_parents, - write::{PropertyObject, fresh_guid, write_revision}, + write::{PropertyObject, fresh_guid, write_revision_on}, }; use serde::{Deserialize, Serialize}; use std::{ @@ -172,7 +172,7 @@ pub(crate) fn format_text( } } let modified = current_timestamps()?.0.to_le_bytes(); - write_revision(source, space, |raw| { + write_revision_on(&index, space, |raw| { let mut target = PropertyObject::from_object(&raw.objects[&object])?; let fields = PropertySets::parse(&target.bytes)?; if fields.sets[0].iter().any(|p| p.id == 0x24003458) { diff --git a/crates/onestore/src/insertion.rs b/crates/onestore/src/insertion.rs index 7a9989e2c9cb8718a199b5803c159468985b3ff2..b9a8081f7388f6fbffbfe423217a0d4754abbfa9 100644 --- a/crates/onestore/src/insertion.rs +++ b/crates/onestore/src/insertion.rs @@ -3,7 +3,7 @@ use crate::{ create::{current_timestamps, default_text_style, properties, string}, document::{Document, Element, Kind}, edit::{editable_parents, page_title}, - write::{PropertyObject, fresh_guid, write_revision}, + write::{PropertyObject, fresh_guid, write_revision_on}, }; use serde::{Deserialize, Serialize}; use std::{ @@ -394,7 +394,7 @@ impl Insertion { drop(view); title }; - write_revision(source, space, |raw| { + write_revision_on(&index, space, |raw| { let mut changed = new; for id in &ancestors { let mut object = PropertyObject::from_object(&raw.objects[id])?; diff --git a/crates/onestore/src/objects.rs b/crates/onestore/src/objects.rs index 3460a0f2c37583ea1c27b9aa57c1640115cc8031..c2546e7960f9aa8b97dc71323ffa500a00e5c531 100644 --- a/crates/onestore/src/objects.rs +++ b/crates/onestore/src/objects.rs @@ -277,6 +277,8 @@ impl<'a> RevisionIndex<'a> { let mut pending: Vec<&Node<'a>> = revision.nodes.iter().rev().collect(); let mut groups = BTreeSet::new(); let mut defining_table = false; + // Entries of the table being defined; the map is built once, at its end. + let mut defining: Vec<(u32, [u8; 16])> = Vec::new(); let initial_crc = if self.store.header.file_type == crate::FileType::Section { u32::MAX } else { @@ -333,6 +335,7 @@ impl<'a> RevisionIndex<'a> { }); } table = Arc::new(GlobalIds::new()); + defining.clear(); defining_table = true; } 0x24..=0x26 => { @@ -343,8 +346,9 @@ impl<'a> RevisionIndex<'a> { }); } let first = u32::from_le_bytes(c.read()?); - let entries: Vec<_> = match node.id { - 0x24 => vec![(first, c.read()?)], + let start = defining.len(); + match node.id { + 0x24 => defining.push((first, c.read()?)), 0x25 | 0x26 => { let count = if node.id == 0x26 { u32::from_le_bytes(c.read()?) @@ -364,30 +368,28 @@ impl<'a> RevisionIndex<'a> { message: "Global ID import range exceeds its table", }); } - let entries: Vec<_> = dependency_table - .range(first..end) - .map(|(from, guid)| (to + from - first, *guid)) - .collect(); - if entries.len() as u64 != u64::from(count) { + defining.extend( + dependency_table + .range(first..end) + .map(|(from, guid)| (to + from - first, *guid)), + ); + if (defining.len() - start) as u64 != u64::from(count) { return Err(Error { offset: node.offset, message: "Global ID import refers to a missing entry", }); } - entries } _ => unreachable!(), - }; - for (index, guid) in entries { - if index >= 0xffffff - || guid == [0; 16] - || Arc::make_mut(&mut table).insert(index, guid).is_some() - { - return Err(Error { - offset: node.offset, - message: "Invalid or repeated global ID entry", - }); - } + } + if defining[start..] + .iter() + .any(|(index, guid)| *index >= 0xffffff || *guid == [0; 16]) + { + return Err(Error { + offset: node.offset, + message: "Invalid or repeated global ID entry", + }); } } 0x28 => { @@ -398,13 +400,22 @@ impl<'a> RevisionIndex<'a> { }); } defining_table = false; - let unique: BTreeSet<_> = table.values().collect(); - if unique.len() != table.len() { + defining.sort_unstable(); + if defining.windows(2).any(|pair| pair[0].0 == pair[1].0) { + return Err(Error { + offset: node.offset, + message: "Invalid or repeated global ID entry", + }); + } + let mut guids: Vec<_> = defining.iter().map(|(_, guid)| *guid).collect(); + guids.sort_unstable(); + if guids.windows(2).any(|pair| pair[0] == pair[1]) { return Err(Error { offset: node.offset, message: "Global ID table repeats a GUID", }); } + table = Arc::new(defining.drain(..).collect()); } 0x59 | 0x5a => { let id = if node.id == 0x5a { diff --git a/crates/onestore/src/outline.rs b/crates/onestore/src/outline.rs index b17661e4de6ebea601e6bf94ee36083d4df3b283..f7207d209526aed9fb263ced2542ea7de1819a86 100644 --- a/crates/onestore/src/outline.rs +++ b/crates/onestore/src/outline.rs @@ -3,7 +3,7 @@ use crate::{ create::current_timestamps, document::{Document, Kind}, edit::{editable_parents, update_title}, - write::{PropertyObject, write_revision}, + write::{PropertyObject, write_revision_on}, }; use serde::{Deserialize, Serialize}; use std::collections::{BTreeMap, BTreeSet}; @@ -105,7 +105,7 @@ impl OutlineEdit { pending.extend(parents.get(&id).into_iter().flatten().copied()); } let modified = current_timestamps()?.0.to_le_bytes(); - write_revision(source, space, |raw| { + write_revision_on(&index, space, |raw| { let mut target = PropertyObject::from_object(&raw.objects[&object])?; target.set( &values @@ -139,6 +139,7 @@ impl OutlineEdit { #[cfg(test)] mod tests { use super::*; + use crate::write::write_revision; #[test] fn protection_on_the_target_or_ancestor_prevents_layout_edits() { diff --git a/crates/onestore/src/page/write.rs b/crates/onestore/src/page/write.rs index 9b9a829421cd23890116a4c139a2c2de7d6e3b61..d1ae945263914f74833e831c410b1a93d66656a3 100644 --- a/crates/onestore/src/page/write.rs +++ b/crates/onestore/src/page/write.rs @@ -339,25 +339,34 @@ pub(crate) fn write_page( if author.contains('\0') { return Err(invalid("Choose an author name without NUL")); } - let store = Store::parse(source)?; - let index = RevisionIndex::parse(&store)?; - index.validate_current()?; - let document = Document::parse(&index)?; - let pages = document.pages_in(space)?; - let [page] = pages.as_slice() else { - return Err(invalid("Choose an object space containing one active page")); + // The parsed source is released before the writers parse their own images. + let (page, before, existing) = { + let store = Store::parse(source)?; + let index = RevisionIndex::parse(&store)?; + index.validate_current()?; + let document = Document::parse(&index)?; + let pages = document.pages_in(space)?; + let [page] = pages.as_slice() else { + return Err(invalid("Choose an object space containing one active page")); + }; + let existing = index + .resolve_active(space)? + .objects + .keys() + .copied() + .collect(); + (*page, Page::from_space(&document, space)?, existing) }; - let before = Page::from_space(&document, space)?; - let raw = index.resolve_active(space)?; let mut lowering = Lowering { image: source.to_vec(), + current: Some(before.clone()), space, - page: *page, + page, author, alias: BTreeMap::new(), built: BTreeSet::new(), }; - lowering.run(&before, after, &raw.objects.keys().copied().collect())?; + lowering.run(&before, after, &existing)?; if lowering.image == source { return Ok(lowering.image); } @@ -453,6 +462,8 @@ impl<'a> View<'a> { struct Lowering<'a> { image: Vec, + /// The page as `image` stores it, until a writer changes the image. + current: Option, space: ExGuid, page: ExGuid, author: &'a str, @@ -478,15 +489,22 @@ impl Lowering<'_> { } fn apply(&mut self, edit: impl FnOnce(&[u8]) -> Result, Error>) -> Result<(), Error> { - self.image = edit(&self.image)?; + let image = edit(&self.image)?; + if image != self.image { + self.image = image; + self.current = None; + } Ok(()) } - fn current(&self) -> Result { - let store = Store::parse(&self.image)?; - let index = RevisionIndex::parse(&store)?; - let document = Document::parse(&index)?; - Page::from_space(&document, self.space) + fn current(&mut self) -> Result { + if self.current.is_none() { + let store = Store::parse(&self.image)?; + let index = RevisionIndex::parse(&store)?; + let document = Document::parse(&index)?; + self.current = Some(Page::from_space(&document, self.space)?); + } + Ok(self.current.clone().unwrap()) } fn run( @@ -2921,10 +2939,6 @@ pub(crate) fn squash( .iter() .map(|(model, image)| (*image, *model)) .collect(); - let applied_store = Store::parse(applied)?; - let applied_index = RevisionIndex::parse(&applied_store)?; - // Payloads the typed edits embedded travel into the squashed transaction as well. - let source_store = Store::parse(source)?; let declared = |store: &Store<'_>| -> Vec<[u8; 16]> { store .lists @@ -2934,7 +2948,10 @@ pub(crate) fn squash( .filter_map(|node| node.payload.get(..16).and_then(|g| g.try_into().ok())) .collect() }; - let existing = declared(&source_store); + let existing = declared(&Store::parse(source)?); + let applied_store = Store::parse(applied)?; + let applied_index = RevisionIndex::parse(&applied_store)?; + // Payloads the typed edits embedded travel into the squashed transaction as well. let mut payloads = Vec::new(); for guid in declared(&applied_store) { if !existing.contains(&guid) { @@ -3018,10 +3035,13 @@ pub(crate) fn squash( } Ok(changes) }; - match protection { - Some(_) => crate::write::append_revisions(source, &payloads, protection, edit), - None => crate::write::write_revisions_with_payloads(source, &payloads, edit), - } + let validate = protection.is_none(); + let output = crate::write::build(source, &payloads, protection, validate, edit)?; + // The parsed images are released before the result is parsed. + drop(applied_index); + drop(applied_store); + crate::write::check(&output, validate)?; + Ok(output) } fn remap(object: &mut PropertyObject, rename: &BTreeMap) -> Result<(), Error> { diff --git a/crates/onestore/src/paragraph.rs b/crates/onestore/src/paragraph.rs index 86ffcdb169df8f55723c3b43e9ad00a2a63a5856..66105442433e16ef48dd085dfc4745a72d0de671 100644 --- a/crates/onestore/src/paragraph.rs +++ b/crates/onestore/src/paragraph.rs @@ -3,7 +3,7 @@ use crate::{ create::{current_timestamps, properties, string}, document::{Document, Element, Kind}, edit::{editable_parents, update_title}, - write::{PropertyObject, fresh_guid, write_revision}, + write::{PropertyObject, fresh_guid, write_revision_on}, }; use serde::{Deserialize, Serialize}; use std::{ @@ -252,7 +252,7 @@ impl ParagraphSplit { object.set(&[(0x14001d7a, &modified)])?; } update_title(&store, &raw, view, &pages, &mut changed)?; - write_revision(source, space, |_| Ok(changed)) + write_revision_on(&index, space, |_| Ok(changed)) } } @@ -682,7 +682,7 @@ impl ParagraphJoin { pending.extend(parents.get(&id).into_iter().flatten().copied()); } update_title(&store, &raw, view, &pages, &mut changed)?; - write_revision(source, space, |_| Ok(changed)) + write_revision_on(&index, space, |_| Ok(changed)) } } diff --git a/crates/onestore/src/store.rs b/crates/onestore/src/store.rs index 18c57e2aad6216fe68e30f290daf8353717d60f4..d97a32017fd30c110bb607a9c539aeb17bb22505 100644 --- a/crates/onestore/src/store.rs +++ b/crates/onestore/src/store.rs @@ -505,6 +505,7 @@ impl<'a> Store<'a> { pending.push(reference); } } + nodes.shrink_to_fit(); lists.insert(list_id.unwrap(), NodeList { fragments, nodes }); } Ok(Self { diff --git a/crates/onestore/src/tree.rs b/crates/onestore/src/tree.rs index e3d4d8d8e576d34af0506289cb86299d095834ba..b717409fb3f5fa5a7fb01a9d8b47167b451aa69b 100644 --- a/crates/onestore/src/tree.rs +++ b/crates/onestore/src/tree.rs @@ -3,7 +3,7 @@ use crate::{ create::{current_timestamps, properties, string}, document::{Document, Kind, Revision}, edit::{editable_parents, update_title}, - write::{PropertyObject, fresh_guid, write_revision}, + write::{PropertyObject, fresh_guid, write_revision_on}, }; use serde::{Deserialize, Serialize}; use std::{ @@ -371,7 +371,7 @@ impl TreeEdit { changed.insert(self.object, object); } update_title(&store, &raw, view, &pages, &mut changed)?; - write_revision(source, space, |_| Ok(changed)) + write_revision_on(&index, space, |_| Ok(changed)) } } @@ -416,6 +416,7 @@ fn checked_path( #[cfg(test)] mod tests { use super::*; + use crate::write::write_revision; use crate::{Insertion, PreparedEdit}; #[test] diff --git a/crates/onestore/src/write.rs b/crates/onestore/src/write.rs index 06f8621b819c07207ba254e024cbdc2de916a261..ee0bbd413e10291cbe6f01babb4e634b03c04237 100644 --- a/crates/onestore/src/write.rs +++ b/crates/onestore/src/write.rs @@ -305,8 +305,11 @@ pub fn replace_property_bytes( message: "Property does not contain scalar bytes", }); } + let store = Store::parse(source)?; + let index = RevisionIndex::parse(&store)?; + index.validate_current()?; replace_objects( - source, + &index, space, &[ObjectEdit { object: object_id, @@ -316,12 +319,13 @@ pub fn replace_property_bytes( ) } +/// Patches objects of a source the caller has parsed and validated. pub(crate) fn replace_objects( - source: &[u8], + index: &RevisionIndex<'_>, space: ExGuid, edits: &[ObjectEdit<'_>], ) -> Result> { - write_revision(source, space, |revision| { + write_revision_on(index, space, |revision| { let mut changed = BTreeMap::new(); for edit in edits { if changed.contains_key(&edit.object) { @@ -788,23 +792,86 @@ pub(crate) fn write_revisions_with_payloads( payloads: &[([u8; 16], &[u8])], edit: impl FnOnce(&RevisionIndex<'_>) -> Result>, ) -> Result> { - let store = Store::parse(source)?; - RevisionIndex::parse(&store)?.validate_current()?; - let output = append_revisions(source, payloads, None, edit)?; - let store = Store::parse(&output)?; - RevisionIndex::parse(&store)?.validate_current()?; - Ok(output) + publish(source, payloads, None, true, edit) } /// `write_revisions_with_payloads` without validating that current revisions are /// complete: a protected section is validated by unlocking it, through `protection`. +#[cfg(feature = "protected")] pub(crate) fn append_revisions( source: &[u8], payloads: &[([u8; 16], &[u8])], protection: Option<&dyn Protection>, edit: impl FnOnce(&RevisionIndex<'_>) -> Result>, +) -> Result> { + publish(source, payloads, protection, false, edit) +} + +fn publish( + source: &[u8], + payloads: &[([u8; 16], &[u8])], + protection: Option<&dyn Protection>, + validate: bool, + edit: impl FnOnce(&RevisionIndex<'_>) -> Result>, +) -> Result> { + // The parsed source is released before the result is parsed. + let output = build(source, payloads, protection, validate, edit)?; + check(&output, validate)?; + Ok(output) +} + +/// Parses a written image, and with `validate` requires its current revisions complete. +pub(crate) fn check(output: &[u8], validate: bool) -> Result<()> { + let store = Store::parse(output)?; + let index = RevisionIndex::parse(&store)?; + if validate { + index.validate_current()?; + } + Ok(()) +} + +/// The written image, unchecked: `check` follows once the caller has released whatever +/// `edit` borrowed. +pub(crate) fn build( + source: &[u8], + payloads: &[([u8; 16], &[u8])], + protection: Option<&dyn Protection>, + validate: bool, + edit: impl FnOnce(&RevisionIndex<'_>) -> Result>, ) -> Result> { let store = Store::parse(source)?; + let index = RevisionIndex::parse(&store)?; + if validate { + index.validate_current()?; + } + build_on(&index, payloads, protection, edit) +} + +/// `write_revision` on a source the caller has parsed and validated. +pub(crate) fn write_revision_on( + index: &RevisionIndex<'_>, + space: ExGuid, + edit: impl FnOnce(&crate::ResolvedRevision<'_>) -> Result>, +) -> Result> { + let output = build_on(index, &[], None, |index| { + let revision = index.resolve(space, index.active(space)?)?; + Ok(BTreeMap::from([( + space, + RevisionEdit::Update(edit(&revision)?), + )])) + })?; + check(&output, true)?; + Ok(output) +} + +fn build_on( + index: &RevisionIndex<'_>, + payloads: &[([u8; 16], &[u8])], + protection: Option<&dyn Protection>, + edit: impl FnOnce(&RevisionIndex<'_>) -> Result>, +) -> Result> { + let store = index.store; + let source = store.data; let is_section = store.header.file_type == FileType::Section; if !store.checksum_mismatches.is_empty() { return Err(Error { @@ -812,13 +879,14 @@ pub(crate) fn append_revisions( message: "Cannot write a file with transaction checksum damage", }); } - let index = RevisionIndex::parse(&store)?; let resolve = |space, rid| match protection { Some(protection) => protection.resolve(space, rid), None => index.resolve(space, rid), }; - let changes = edit(&index)?; - let mut output = source.to_vec(); + let changes = edit(index)?; + // Room for the revision, so appending does not double the image. + let mut output = Vec::with_capacity(source.len() + source.len() / 16 + (1 << 16)); + output.extend_from_slice(source); // Native files reserve 1 KiB per transaction-log fragment; a fragment that ends the file // keeps that room before the first new chunk. let tail = store.transaction_fragments.last().unwrap().chunk; @@ -1348,12 +1416,12 @@ pub(crate) fn append_revisions( let space_node = root .nodes .iter() - .find(|node| node.id == 8 && node.fields(&store).exguid() == Ok(space)) + .find(|node| node.id == 8 && node.fields(store).exguid() == Ok(space)) .ok_or(Error { offset: 0, message: "Object space is absent from the root list", })?; - let space_list = space_node.referenced_list(&store)?; + let space_list = space_node.referenced_list(store)?; let revision_node = space_list.iter().rfind(|node| node.id == 0x10).unwrap(); let Some(Reference::NodeList(manifest_reference)) = revision_node.reference else { unreachable!() @@ -1522,6 +1590,5 @@ pub(crate) fn append_revisions( message: "File generation counter is exhausted", })?; output[228..236].copy_from_slice(&generation.to_le_bytes()); - RevisionIndex::parse(&Store::parse(&output)?)?; Ok(output) } diff --git a/tools/cache_images.py b/tools/cache_images.py new file mode 100644 index 0000000000000000000000000000000000000000..6b241891b9f8f36bc438d2b23151af21f4fcc179 --- /dev/null +++ b/tools/cache_images.py @@ -0,0 +1,17 @@ +"""The cache's working image, stored as the byte ranges where it differs from the base.""" +import struct + + +def working(connection): + base, patch = connection.execute('SELECT base, working FROM replica WHERE id=1').fetchone() + if not patch: return base + length, = struct.unpack_from('