| author | |
| committer | |
| log | be6b82da9a9e19b3a47c24655949fd2d66db0800 |
| tree | afe4c81369cdf10faf1f1ddb8aff5d220f73bacb |
| parent | 3d0dd98f39482111b7e2e32734d71bc806d826c8 |
| signature | Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU |
A page save parsed the section about twenty times with up to four parsed copies
alive: the lowering re-read the unchanged page before every writer pass, the
write path parsed and validated the source its caller had just parsed and
validated, and squash kept both parsed images alive while checking the result.
The lowered page is now cached until a writer changes the image, typed writers
hand their validated index to the write layer, parsed sources are released
before results are checked, GUID tables are built once from a sorted list and
node lists drop their slack.
The cache stored base and working whole in one row, so every save rewrote twice
the section and every idle poll parsed, validated and rewrote both. working is
now the byte ranges where it differs from base (schema 14) and an unchanged
remote is recognised by comparison alone.
Owner 12 MiB section: 141 ms / +65 MiB to 65 ms / +29 MiB per save. 24 MiB soak
section: 780 ms / +506 MiB to 310 ms / +154 MiB; launch with two publications
7.3 s to 2.7 s; cache half the size.
Assisted-by: claude-fable-5.120 files changed, 394 insertions(+), 147 deletions(-)
crates/notebook/README.md+3-2| ... | ... | @@ -2,8 +2,9 @@ |
| 2 | 2 | |
| 3 | 3 | The application-facing crate: notebook discovery, a durable local replica with |
| 4 | 4 | reconnect reconciliation, external-asset caching, recovery export and optional |
| 5 | embedded SMB access. The replica stores a complete working image of one section | |
| 6 | and the editor's intents in a local SQLite database. `sync_once` provides a | |
| 5 | embedded SMB access. The replica stores the last observed remote image of one section, | |
| 6 | the working image as its difference from that, and the editor's intents in a | |
| 7 | local SQLite database; a save writes the revision it appended, not the section. `sync_once` provides a | |
| 7 | 8 | reconciliation step and `start_sync` owns automatic polling and reconnects. Local |
| 8 | 9 | success does not acknowledge publication to a shared notebook. |
| 9 | 10 |
crates/notebook/src/assets.rs+2-6| ... | ... | @@ -79,12 +79,8 @@ pub(crate) fn key(filename: &str) -> Result<String> { |
| 79 | 79 | } |
| 80 | 80 | |
| 81 | 81 | fn referenced(connection: &Connection, key: &str) -> Result<bool> { |
| 82 | for column in ["working", "base"] { | |
| 83 | let image: Vec<u8> = connection.query_row( | |
| 84 | &format!("SELECT {column} FROM replica WHERE id=1"), | |
| 85 | [], | |
| 86 | |row| row.get(0), | |
| 87 | )?; | |
| 82 | let (base, working) = crate::images::both(connection)?; | |
| 83 | for image in [working, base] { | |
| 88 | 84 | let store = Store::parse(&image)?; |
| 89 | 85 | let index = RevisionIndex::parse(&store)?; |
| 90 | 86 | let document = Document::parse(&index)?; |
crates/notebook/src/images.rs created+151| ... | ... | @@ -0,0 +1,151 @@ |
| 1 | //! The cache's two section images. `base` is stored whole; `working` is stored as the byte | |
| 2 | //! ranges where it differs from `base`, so saving an edit writes the revision it appended | |
| 3 | //! instead of the section. | |
| 4 | ||
| 5 | use crate::Result; | |
| 6 | use rusqlite::{Connection, params}; | |
| 7 | use std::io; | |
| 8 | ||
| 9 | const BLOCK: usize = 4096; | |
| 10 | ||
| 11 | /// `image` as its length and the `(offset, bytes)` runs of blocks that differ from `base`; | |
| 12 | /// empty when the images are equal. | |
| 13 | fn difference(base: &[u8], image: &[u8]) -> Vec<u8> { | |
| 14 | if base == image { | |
| 15 | return Vec::new(); | |
| 16 | } | |
| 17 | let mut patch = (image.len() as u64).to_le_bytes().to_vec(); | |
| 18 | let mut run: Option<usize> = None; | |
| 19 | let close = |patch: &mut Vec<u8>, start: usize, end: usize| { | |
| 20 | patch.extend_from_slice(&(start as u64).to_le_bytes()); | |
| 21 | patch.extend_from_slice(&((end - start) as u64).to_le_bytes()); | |
| 22 | patch.extend_from_slice(&image[start..end]); | |
| 23 | }; | |
| 24 | for start in (0..image.len()).step_by(BLOCK) { | |
| 25 | let end = (start + BLOCK).min(image.len()); | |
| 26 | if base.get(start..end) == Some(&image[start..end]) { | |
| 27 | if let Some(from) = run.take() { | |
| 28 | close(&mut patch, from, start); | |
| 29 | } | |
| 30 | } else { | |
| 31 | run.get_or_insert(start); | |
| 32 | } | |
| 33 | } | |
| 34 | if let Some(from) = run { | |
| 35 | close(&mut patch, from, image.len()); | |
| 36 | } | |
| 37 | patch | |
| 38 | } | |
| 39 | ||
| 40 | fn restore(mut base: Vec<u8>, patch: &[u8]) -> Result<Vec<u8>> { | |
| 41 | let damaged = || io::Error::new(io::ErrorKind::InvalidData, "Damaged working image"); | |
| 42 | let Some((length, mut runs)) = patch.split_first_chunk::<8>() else { | |
| 43 | return if patch.is_empty() { | |
| 44 | Ok(base) | |
| 45 | } else { | |
| 46 | Err(damaged().into()) | |
| 47 | }; | |
| 48 | }; | |
| 49 | base.resize( | |
| 50 | usize::try_from(u64::from_le_bytes(*length)).map_err(|_| damaged())?, | |
| 51 | 0, | |
| 52 | ); | |
| 53 | while let Some((header, rest)) = runs.split_first_chunk::<16>() { | |
| 54 | let offset = usize::try_from(u64::from_le_bytes(header[..8].try_into().unwrap())) | |
| 55 | .map_err(|_| damaged())?; | |
| 56 | let size = usize::try_from(u64::from_le_bytes(header[8..].try_into().unwrap())) | |
| 57 | .map_err(|_| damaged())?; | |
| 58 | let (bytes, rest) = rest.split_at_checked(size).ok_or_else(damaged)?; | |
| 59 | base.get_mut(offset..) | |
| 60 | .and_then(|target| target.get_mut(..size)) | |
| 61 | .ok_or_else(damaged)? | |
| 62 | .copy_from_slice(bytes); | |
| 63 | runs = rest; | |
| 64 | } | |
| 65 | if runs.is_empty() { | |
| 66 | Ok(base) | |
| 67 | } else { | |
| 68 | Err(damaged().into()) | |
| 69 | } | |
| 70 | } | |
| 71 | ||
| 72 | pub(crate) fn base(connection: &Connection) -> Result<Vec<u8>> { | |
| 73 | Ok(connection.query_row("SELECT base FROM replica WHERE id=1", [], |row| row.get(0))?) | |
| 74 | } | |
| 75 | ||
| 76 | /// `(base, working)`. | |
| 77 | pub(crate) fn both(connection: &Connection) -> Result<(Vec<u8>, Vec<u8>)> { | |
| 78 | let (base, patch): (Vec<u8>, Vec<u8>) = | |
| 79 | connection.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { | |
| 80 | Ok((row.get(0)?, row.get(1)?)) | |
| 81 | })?; | |
| 82 | let working = restore(base.clone(), &patch)?; | |
| 83 | Ok((base, working)) | |
| 84 | } | |
| 85 | ||
| 86 | pub(crate) fn working(connection: &Connection) -> Result<Vec<u8>> { | |
| 87 | let (base, patch): (Vec<u8>, Vec<u8>) = | |
| 88 | connection.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { | |
| 89 | Ok((row.get(0)?, row.get(1)?)) | |
| 90 | })?; | |
| 91 | restore(base, &patch) | |
| 92 | } | |
| 93 | ||
| 94 | /// Replaces the working image; `base` is the stored base image. | |
| 95 | pub(crate) fn set_working(connection: &Connection, base: &[u8], working: &[u8]) -> Result<()> { | |
| 96 | connection.execute( | |
| 97 | "UPDATE replica SET working=?1 WHERE id=1", | |
| 98 | [difference(base, working)], | |
| 99 | )?; | |
| 100 | Ok(()) | |
| 101 | } | |
| 102 | ||
| 103 | /// Replaces the base image, keeping `working` as the working image. An unchanged base | |
| 104 | /// leaves its stored bytes alone. | |
| 105 | pub(crate) fn set_base( | |
| 106 | connection: &Connection, | |
| 107 | old: &[u8], | |
| 108 | base: &[u8], | |
| 109 | working: &[u8], | |
| 110 | ) -> Result<()> { | |
| 111 | if old == base { | |
| 112 | return set_working(connection, base, working); | |
| 113 | } | |
| 114 | connection.execute( | |
| 115 | "UPDATE replica SET base=?1, working=?2 WHERE id=1", | |
| 116 | params![base, difference(base, working)], | |
| 117 | )?; | |
| 118 | Ok(()) | |
| 119 | } | |
| 120 | ||
| 121 | #[cfg(test)] | |
| 122 | mod tests { | |
| 123 | use super::*; | |
| 124 | ||
| 125 | #[test] | |
| 126 | fn a_working_image_survives_as_its_difference_from_any_base() { | |
| 127 | let base: Vec<u8> = (0..3 * BLOCK + 17).map(|i| i as u8).collect(); | |
| 128 | let mut edited = base.clone(); | |
| 129 | edited[5] ^= 1; | |
| 130 | edited[2 * BLOCK + 1] ^= 1; | |
| 131 | edited.extend_from_slice(b"appended revision"); | |
| 132 | for image in [ | |
| 133 | base.clone(), | |
| 134 | edited.clone(), | |
| 135 | base[..BLOCK + 3].to_vec(), | |
| 136 | Vec::new(), | |
| 137 | vec![7; 5 * BLOCK], | |
| 138 | ] { | |
| 139 | let patch = difference(&base, &image); | |
| 140 | assert_eq!(restore(base.clone(), &patch).unwrap(), image); | |
| 141 | } | |
| 142 | assert!(difference(&base, &base).is_empty()); | |
| 143 | assert!(difference(&base, &edited).len() < 3 * BLOCK); | |
| 144 | let patch = difference(&base, &edited); | |
| 145 | for cut in 1..patch.len() { | |
| 146 | if let Ok(image) = restore(base.clone(), &patch[..cut]) { | |
| 147 | assert_ne!(image, edited); | |
| 148 | } | |
| 149 | } | |
| 150 | } | |
| 151 | } |
crates/notebook/src/lib.rs+9-27| ... | ... | @@ -20,6 +20,7 @@ use std::{ |
| 20 | 20 | }; |
| 21 | 21 | |
| 22 | 22 | mod assets; |
| 23 | mod images; | |
| 23 | 24 | mod merge; |
| 24 | 25 | mod pages; |
| 25 | 26 | mod rebase; |
| ... | ... | @@ -59,7 +60,7 @@ pub enum Error { |
| 59 | 60 | type Result<T> = std::result::Result<T, Error>; |
| 60 | 61 | |
| 61 | 62 | const APPLICATION_ID: u32 = 0x4f4e454f; |
| 62 | const SCHEMA_VERSION: u32 = 13; | |
| 63 | const SCHEMA_VERSION: u32 = 14; | |
| 63 | 64 | |
| 64 | 65 | /// An edited page model together with the stored model it was edited from. |
| 65 | 66 | /// `before` is the precondition reconciliation checks against the remote page. |
| ... | ... | @@ -217,7 +218,7 @@ impl Replica { |
| 217 | 218 | ", |
| 218 | 219 | )?; |
| 219 | 220 | schema::create(&transaction)?; |
| 220 | transaction.execute("INSERT INTO replica VALUES (1, ?1, ?1)", [source])?; | |
| 221 | transaction.execute("INSERT INTO replica VALUES (1, ?1, x'')", [source])?; | |
| 221 | 222 | } else { |
| 222 | 223 | if application != APPLICATION_ID { |
| 223 | 224 | return Err( |
| ... | ... | @@ -251,11 +252,7 @@ impl Replica { |
| 251 | 252 | .connection |
| 252 | 253 | .lock() |
| 253 | 254 | .map_err(|_| io::Error::other("Cache owner panicked"))?; |
| 254 | Ok( | |
| 255 | connection.query_row("SELECT working FROM replica WHERE id=1", [], |row| { | |
| 256 | row.get(0) | |
| 257 | })?, | |
| 258 | ) | |
| 255 | images::working(&connection) | |
| 259 | 256 | } |
| 260 | 257 | |
| 261 | 258 | pub fn pending(&self) -> Result<Vec<PendingEdit>> { |
| ... | ... | @@ -294,10 +291,7 @@ impl Replica { |
| 294 | 291 | .map_err(|_| io::Error::other("Cache owner panicked"))?; |
| 295 | 292 | let transaction = |
| 296 | 293 | connection.transaction_with_behavior(TransactionBehavior::Immediate)?; |
| 297 | let current: Vec<u8> = | |
| 298 | transaction.query_row("SELECT working FROM replica WHERE id=1", [], |row| { | |
| 299 | row.get(0) | |
| 300 | })?; | |
| 294 | let (base, current) = images::both(&transaction)?; | |
| 301 | 295 | if current != source { |
| 302 | 296 | return Err(io::Error::new( |
| 303 | 297 | io::ErrorKind::ResourceBusy, |
| ... | ... | @@ -334,10 +328,7 @@ impl Replica { |
| 334 | 328 | id |
| 335 | 329 | ], |
| 336 | 330 | )?; |
| 337 | transaction.execute( | |
| 338 | "UPDATE replica SET working=?1 WHERE id=1", | |
| 339 | [prepared.as_bytes()], | |
| 340 | )?; | |
| 331 | images::set_working(&transaction, &base, prepared.as_bytes())?; | |
| 341 | 332 | transaction.commit()?; |
| 342 | 333 | drop(connection); |
| 343 | 334 | self.wake_sync(); |
| ... | ... | @@ -388,10 +379,7 @@ impl Replica { |
| 388 | 379 | .lock() |
| 389 | 380 | .map_err(|_| io::Error::other("Cache owner panicked"))?; |
| 390 | 381 | let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?; |
| 391 | let current: Vec<u8> = | |
| 392 | transaction.query_row("SELECT working FROM replica WHERE id=1", [], |row| { | |
| 393 | row.get(0) | |
| 394 | })?; | |
| 382 | let (base, current) = images::both(&transaction)?; | |
| 395 | 383 | if current != source { |
| 396 | 384 | return Err(io::Error::new( |
| 397 | 385 | io::ErrorKind::ResourceBusy, |
| ... | ... | @@ -410,10 +398,7 @@ impl Replica { |
| 410 | 398 | ], |
| 411 | 399 | )?; |
| 412 | 400 | let id = u64::try_from(transaction.last_insert_rowid()).map_err(io::Error::other)?; |
| 413 | transaction.execute( | |
| 414 | "UPDATE replica SET working=?1 WHERE id=1", | |
| 415 | [edit.as_bytes()], | |
| 416 | )?; | |
| 401 | images::set_working(&transaction, &base, edit.as_bytes())?; | |
| 417 | 402 | transaction.commit()?; |
| 418 | 403 | drop(connection); |
| 419 | 404 | self.wake_sync(); |
| ... | ... | @@ -514,10 +499,7 @@ fn validate_images(connection: &Connection) -> Result<()> { |
| 514 | 499 | io::Error::new(io::ErrorKind::InvalidData, "Cache integrity check failed").into(), |
| 515 | 500 | ); |
| 516 | 501 | } |
| 517 | let (base, working): (Vec<u8>, Vec<u8>) = | |
| 518 | connection.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { | |
| 519 | Ok((row.get(0)?, row.get(1)?)) | |
| 520 | })?; | |
| 502 | let (base, working) = images::both(connection)?; | |
| 521 | 503 | if validate(&base)? != validate(&working)? { |
| 522 | 504 | return Err(io::Error::new( |
| 523 | 505 | io::ErrorKind::InvalidData, |
crates/notebook/src/recovery.rs+6-11| ... | ... | @@ -58,17 +58,11 @@ impl Recovery { |
| 58 | 58 | } |
| 59 | 59 | |
| 60 | 60 | pub fn snapshot(&self) -> Result<Vec<u8>> { |
| 61 | Ok(self | |
| 62 | .connection | |
| 63 | .query_row("SELECT working FROM replica WHERE id=1", [], |row| { | |
| 64 | row.get(0) | |
| 65 | })?) | |
| 61 | crate::images::working(&self.connection) | |
| 66 | 62 | } |
| 67 | 63 | |
| 68 | 64 | pub fn remote_snapshot(&self) -> Result<Vec<u8>> { |
| 69 | Ok(self | |
| 70 | .connection | |
| 71 | .query_row("SELECT base FROM replica WHERE id=1", [], |row| row.get(0))?) | |
| 65 | crate::images::base(&self.connection) | |
| 72 | 66 | } |
| 73 | 67 | |
| 74 | 68 | pub fn pending(&self) -> Result<Vec<PendingEdit>> { |
| ... | ... | @@ -150,10 +144,11 @@ fn summary(connection: &Connection) -> Result<RecoverySummary> { |
| 150 | 144 | [], |
| 151 | 145 | |row| Ok((unsigned(row, 0)?, unsigned(row, 1)?)), |
| 152 | 146 | )?; |
| 147 | let working_bytes = crate::images::working(connection)?.len() as u64; | |
| 153 | 148 | Ok(connection.query_row( |
| 154 | 149 | "SELECT (SELECT count(*) FROM edits), (SELECT count(*) FROM conflicts), |
| 155 | 150 | (SELECT count(*) FROM attempt), (SELECT count(*) FROM receipts), |
| 156 | length(working), length(base) FROM replica WHERE id=1", | |
| 151 | length(base) FROM replica WHERE id=1", | |
| 157 | 152 | [], |
| 158 | 153 | |row| { |
| 159 | 154 | Ok(RecoverySummary { |
| ... | ... | @@ -161,8 +156,8 @@ fn summary(connection: &Connection) -> Result<RecoverySummary> { |
| 161 | 156 | conflicts: unsigned(row, 1)?, |
| 162 | 157 | uncertain_edits: unsigned(row, 2)?, |
| 163 | 158 | published_receipts: unsigned(row, 3)?, |
| 164 | working_bytes: unsigned(row, 4)?, | |
| 165 | remote_bytes: unsigned(row, 5)?, | |
| 159 | working_bytes, | |
| 160 | remote_bytes: unsigned(row, 4)?, | |
| 166 | 161 | cached_assets, |
| 167 | 162 | cached_asset_bytes, |
| 168 | 163 | }) |
crates/notebook/src/sync.rs+25-22| ... | ... | @@ -63,7 +63,7 @@ impl Replica { |
| 63 | 63 | .connection |
| 64 | 64 | .lock() |
| 65 | 65 | .map_err(|_| io::Error::other("Cache owner panicked"))?; |
| 66 | Ok(connection.query_row("SELECT base FROM replica WHERE id=1", [], |row| row.get(0))?) | |
| 66 | images::base(&connection) | |
| 67 | 67 | } |
| 68 | 68 | |
| 69 | 69 | /// Reconciles one pending edit, or refreshes the working image when the queue is empty. |
| ... | ... | @@ -79,7 +79,12 @@ impl Replica { |
| 79 | 79 | } |
| 80 | 80 | let _step = Step(self); |
| 81 | 81 | let snapshot = remote.read().map_err(Error::RemoteIo)?; |
| 82 | let identity = validate(&snapshot)?; | |
| 82 | // An unchanged remote is the base image, validated when it was stored. | |
| 83 | let identity = if snapshot == self.remote_snapshot()? { | |
| 84 | None | |
| 85 | } else { | |
| 86 | Some(validate(&snapshot)?) | |
| 87 | }; | |
| 83 | 88 | let (intent, attempted) = { |
| 84 | 89 | let mut connection = self |
| 85 | 90 | .connection |
| ... | ... | @@ -87,11 +92,10 @@ impl Replica { |
| 87 | 92 | .map_err(|_| io::Error::other("Cache owner panicked"))?; |
| 88 | 93 | let transaction = |
| 89 | 94 | connection.transaction_with_behavior(TransactionBehavior::Immediate)?; |
| 90 | let base: Vec<u8> = | |
| 91 | transaction | |
| 92 | .query_row("SELECT base FROM replica WHERE id=1", [], |row| row.get(0))?; | |
| 93 | let base_store = Store::parse(&base)?; | |
| 94 | if RevisionIndex::parse(&base_store)?.root != identity { | |
| 95 | let (base, working) = images::both(&transaction)?; | |
| 96 | if let Some(identity) = identity | |
| 97 | && RevisionIndex::parse(&Store::parse(&base)?)?.root != identity | |
| 98 | { | |
| 95 | 99 | return Err(io::Error::new( |
| 96 | 100 | io::ErrorKind::InvalidInput, |
| 97 | 101 | "Remote snapshot belongs to another document", |
| ... | ... | @@ -105,10 +109,7 @@ impl Replica { |
| 105 | 109 | rows.next()?.map(pending_edit).transpose()? |
| 106 | 110 | }; |
| 107 | 111 | let Some(intent) = intent else { |
| 108 | transaction.execute( | |
| 109 | "UPDATE replica SET base=?1, working=?1 WHERE id=1", | |
| 110 | [&snapshot], | |
| 111 | )?; | |
| 112 | images::set_base(&transaction, &base, &snapshot, &snapshot)?; | |
| 112 | 113 | transaction.commit()?; |
| 113 | 114 | return Ok(None); |
| 114 | 115 | }; |
| ... | ... | @@ -123,7 +124,7 @@ impl Replica { |
| 123 | 124 | |row| row.get::<_, String>(0), |
| 124 | 125 | ) |
| 125 | 126 | .optional()?; |
| 126 | transaction.execute("UPDATE replica SET base=?1 WHERE id=1", [&snapshot])?; | |
| 127 | images::set_base(&transaction, &base, &snapshot, &working)?; | |
| 127 | 128 | transaction.commit()?; |
| 128 | 129 | (intent, attempted) |
| 129 | 130 | }; |
| ... | ... | @@ -415,10 +416,7 @@ impl Replica { |
| 415 | 416 | .lock() |
| 416 | 417 | .map_err(|_| io::Error::other("Cache owner panicked"))?; |
| 417 | 418 | let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?; |
| 418 | let (base, working): (Vec<u8>, Vec<u8>) = | |
| 419 | transaction.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { | |
| 420 | Ok((row.get(0)?, row.get(1)?)) | |
| 421 | })?; | |
| 419 | let (base, working) = images::both(&transaction)?; | |
| 422 | 420 | if working != local || base != remote { |
| 423 | 421 | return Err(io::Error::new( |
| 424 | 422 | io::ErrorKind::ResourceBusy, |
| ... | ... | @@ -464,7 +462,7 @@ impl Replica { |
| 464 | 462 | [archive.to_string_lossy().into_owned()], |
| 465 | 463 | )?; |
| 466 | 464 | transaction.execute("DELETE FROM edits", [])?; |
| 467 | transaction.execute("UPDATE replica SET working=base WHERE id=1", [])?; | |
| 465 | transaction.execute("UPDATE replica SET working=x'' WHERE id=1", [])?; | |
| 468 | 466 | } |
| 469 | 467 | } |
| 470 | 468 | transaction.commit()?; |
| ... | ... | @@ -499,10 +497,7 @@ impl Replica { |
| 499 | 497 | .lock() |
| 500 | 498 | .map_err(|_| io::Error::other("Cache owner panicked"))?; |
| 501 | 499 | let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?; |
| 502 | let (base, working): (Vec<u8>, Vec<u8>) = | |
| 503 | transaction.query_row("SELECT base, working FROM replica WHERE id=1", [], |row| { | |
| 504 | Ok((row.get(0)?, row.get(1)?)) | |
| 505 | })?; | |
| 500 | let (base, working) = images::both(&transaction)?; | |
| 506 | 501 | if working != local || base != remote { |
| 507 | 502 | return Err(io::Error::new( |
| 508 | 503 | io::ErrorKind::ResourceBusy, |
| ... | ... | @@ -561,7 +556,15 @@ impl Replica { |
| 561 | 556 | params![id, revision.to_string()], |
| 562 | 557 | )?; |
| 563 | 558 | transaction.execute("DELETE FROM edits WHERE id=?1", [id])?; |
| 564 | transaction.execute("UPDATE replica SET base=?1, working=CASE WHEN EXISTS(SELECT 1 FROM edits) THEN working ELSE ?1 END WHERE id=1", [snapshot])?; | |
| 559 | let (base, working) = images::both(&transaction)?; | |
| 560 | let pending: bool = | |
| 561 | transaction.query_row("SELECT EXISTS(SELECT 1 FROM edits)", [], |row| row.get(0))?; | |
| 562 | images::set_base( | |
| 563 | &transaction, | |
| 564 | &base, | |
| 565 | snapshot, | |
| 566 | if pending { &working } else { snapshot }, | |
| 567 | )?; | |
| 565 | 568 | transaction.commit()?; |
| 566 | 569 | Ok(()) |
| 567 | 570 | } |
crates/onestore/src/edit.rs+2-2| ... | ... | @@ -237,7 +237,7 @@ pub fn replace_text( |
| 237 | 237 | let Some((page, automatic, title_text)) = |
| 238 | 238 | page_title(revision, &pages, Some((object, &changed)))? |
| 239 | 239 | else { |
| 240 | return crate::write::replace_objects(source, space, &edits); | |
| 240 | return crate::write::replace_objects(&index, space, &edits); | |
| 241 | 241 | }; |
| 242 | 242 | let Kind::Page { |
| 243 | 243 | alternate_title, .. |
| ... | ... | @@ -296,7 +296,7 @@ pub fn replace_text( |
| 296 | 296 | &[] |
| 297 | 297 | }, |
| 298 | 298 | }); |
| 299 | crate::write::replace_objects(source, space, &edits) | |
| 299 | crate::write::replace_objects(&index, space, &edits) | |
| 300 | 300 | } |
| 301 | 301 | |
| 302 | 302 | pub(crate) fn editable_parents( |
crates/onestore/src/formatting.rs+2-2| ... | ... | @@ -3,7 +3,7 @@ use crate::{ |
| 3 | 3 | create::{current_timestamps, properties, string}, |
| 4 | 4 | document::{Document, Kind}, |
| 5 | 5 | edit::editable_parents, |
| 6 | write::{PropertyObject, fresh_guid, write_revision}, | |
| 6 | write::{PropertyObject, fresh_guid, write_revision_on}, | |
| 7 | 7 | }; |
| 8 | 8 | use serde::{Deserialize, Serialize}; |
| 9 | 9 | use std::{ |
| ... | ... | @@ -172,7 +172,7 @@ pub(crate) fn format_text( |
| 172 | 172 | } |
| 173 | 173 | } |
| 174 | 174 | let modified = current_timestamps()?.0.to_le_bytes(); |
| 175 | write_revision(source, space, |raw| { | |
| 175 | write_revision_on(&index, space, |raw| { | |
| 176 | 176 | let mut target = PropertyObject::from_object(&raw.objects[&object])?; |
| 177 | 177 | let fields = PropertySets::parse(&target.bytes)?; |
| 178 | 178 | if fields.sets[0].iter().any(|p| p.id == 0x24003458) { |
crates/onestore/src/insertion.rs+2-2| ... | ... | @@ -3,7 +3,7 @@ use crate::{ |
| 3 | 3 | create::{current_timestamps, default_text_style, properties, string}, |
| 4 | 4 | document::{Document, Element, Kind}, |
| 5 | 5 | edit::{editable_parents, page_title}, |
| 6 | write::{PropertyObject, fresh_guid, write_revision}, | |
| 6 | write::{PropertyObject, fresh_guid, write_revision_on}, | |
| 7 | 7 | }; |
| 8 | 8 | use serde::{Deserialize, Serialize}; |
| 9 | 9 | use std::{ |
| ... | ... | @@ -394,7 +394,7 @@ impl Insertion { |
| 394 | 394 | drop(view); |
| 395 | 395 | title |
| 396 | 396 | }; |
| 397 | write_revision(source, space, |raw| { | |
| 397 | write_revision_on(&index, space, |raw| { | |
| 398 | 398 | let mut changed = new; |
| 399 | 399 | for id in &ancestors { |
| 400 | 400 | let mut object = PropertyObject::from_object(&raw.objects[id])?; |
crates/onestore/src/objects.rs+32-21| ... | ... | @@ -277,6 +277,8 @@ impl<'a> RevisionIndex<'a> { |
| 277 | 277 | let mut pending: Vec<&Node<'a>> = revision.nodes.iter().rev().collect(); |
| 278 | 278 | let mut groups = BTreeSet::new(); |
| 279 | 279 | let mut defining_table = false; |
| 280 | // Entries of the table being defined; the map is built once, at its end. | |
| 281 | let mut defining: Vec<(u32, [u8; 16])> = Vec::new(); | |
| 280 | 282 | let initial_crc = if self.store.header.file_type == crate::FileType::Section { |
| 281 | 283 | u32::MAX |
| 282 | 284 | } else { |
| ... | ... | @@ -333,6 +335,7 @@ impl<'a> RevisionIndex<'a> { |
| 333 | 335 | }); |
| 334 | 336 | } |
| 335 | 337 | table = Arc::new(GlobalIds::new()); |
| 338 | defining.clear(); | |
| 336 | 339 | defining_table = true; |
| 337 | 340 | } |
| 338 | 341 | 0x24..=0x26 => { |
| ... | ... | @@ -343,8 +346,9 @@ impl<'a> RevisionIndex<'a> { |
| 343 | 346 | }); |
| 344 | 347 | } |
| 345 | 348 | let first = u32::from_le_bytes(c.read()?); |
| 346 | let entries: Vec<_> = match node.id { | |
| 347 | 0x24 => vec![(first, c.read()?)], | |
| 349 | let start = defining.len(); | |
| 350 | match node.id { | |
| 351 | 0x24 => defining.push((first, c.read()?)), | |
| 348 | 352 | 0x25 | 0x26 => { |
| 349 | 353 | let count = if node.id == 0x26 { |
| 350 | 354 | u32::from_le_bytes(c.read()?) |
| ... | ... | @@ -364,30 +368,28 @@ impl<'a> RevisionIndex<'a> { |
| 364 | 368 | message: "Global ID import range exceeds its table", |
| 365 | 369 | }); |
| 366 | 370 | } |
| 367 | let entries: Vec<_> = dependency_table | |
| 368 | .range(first..end) | |
| 369 | .map(|(from, guid)| (to + from - first, *guid)) | |
| 370 | .collect(); | |
| 371 | if entries.len() as u64 != u64::from(count) { | |
| 371 | defining.extend( | |
| 372 | dependency_table | |
| 373 | .range(first..end) | |
| 374 | .map(|(from, guid)| (to + from - first, *guid)), | |
| 375 | ); | |
| 376 | if (defining.len() - start) as u64 != u64::from(count) { | |
| 372 | 377 | return Err(Error { |
| 373 | 378 | offset: node.offset, |
| 374 | 379 | message: "Global ID import refers to a missing entry", |
| 375 | 380 | }); |
| 376 | 381 | } |
| 377 | entries | |
| 378 | 382 | } |
| 379 | 383 | _ => unreachable!(), |
| 380 | }; | |
| 381 | for (index, guid) in entries { | |
| 382 | if index >= 0xffffff | |
| 383 | || guid == [0; 16] | |
| 384 | || Arc::make_mut(&mut table).insert(index, guid).is_some() | |
| 385 | { | |
| 386 | return Err(Error { | |
| 387 | offset: node.offset, | |
| 388 | message: "Invalid or repeated global ID entry", | |
| 389 | }); | |
| 390 | } | |
| 384 | } | |
| 385 | if defining[start..] | |
| 386 | .iter() | |
| 387 | .any(|(index, guid)| *index >= 0xffffff || *guid == [0; 16]) | |
| 388 | { | |
| 389 | return Err(Error { | |
| 390 | offset: node.offset, | |
| 391 | message: "Invalid or repeated global ID entry", | |
| 392 | }); | |
| 391 | 393 | } |
| 392 | 394 | } |
| 393 | 395 | 0x28 => { |
| ... | ... | @@ -398,13 +400,22 @@ impl<'a> RevisionIndex<'a> { |
| 398 | 400 | }); |
| 399 | 401 | } |
| 400 | 402 | defining_table = false; |
| 401 | let unique: BTreeSet<_> = table.values().collect(); | |
| 402 | if unique.len() != table.len() { | |
| 403 | defining.sort_unstable(); | |
| 404 | if defining.windows(2).any(|pair| pair[0].0 == pair[1].0) { | |
| 405 | return Err(Error { | |
| 406 | offset: node.offset, | |
| 407 | message: "Invalid or repeated global ID entry", | |
| 408 | }); | |
| 409 | } | |
| 410 | let mut guids: Vec<_> = defining.iter().map(|(_, guid)| *guid).collect(); | |
| 411 | guids.sort_unstable(); | |
| 412 | if guids.windows(2).any(|pair| pair[0] == pair[1]) { | |
| 403 | 413 | return Err(Error { |
| 404 | 414 | offset: node.offset, |
| 405 | 415 | message: "Global ID table repeats a GUID", |
| 406 | 416 | }); |
| 407 | 417 | } |
| 418 | table = Arc::new(defining.drain(..).collect()); | |
| 408 | 419 | } |
| 409 | 420 | 0x59 | 0x5a => { |
| 410 | 421 | let id = if node.id == 0x5a { |
crates/onestore/src/outline.rs+3-2| ... | ... | @@ -3,7 +3,7 @@ use crate::{ |
| 3 | 3 | create::current_timestamps, |
| 4 | 4 | document::{Document, Kind}, |
| 5 | 5 | edit::{editable_parents, update_title}, |
| 6 | write::{PropertyObject, write_revision}, | |
| 6 | write::{PropertyObject, write_revision_on}, | |
| 7 | 7 | }; |
| 8 | 8 | use serde::{Deserialize, Serialize}; |
| 9 | 9 | use std::collections::{BTreeMap, BTreeSet}; |
| ... | ... | @@ -105,7 +105,7 @@ impl OutlineEdit { |
| 105 | 105 | pending.extend(parents.get(&id).into_iter().flatten().copied()); |
| 106 | 106 | } |
| 107 | 107 | let modified = current_timestamps()?.0.to_le_bytes(); |
| 108 | write_revision(source, space, |raw| { | |
| 108 | write_revision_on(&index, space, |raw| { | |
| 109 | 109 | let mut target = PropertyObject::from_object(&raw.objects[&object])?; |
| 110 | 110 | target.set( |
| 111 | 111 | &values |
| ... | ... | @@ -139,6 +139,7 @@ impl OutlineEdit { |
| 139 | 139 | #[cfg(test)] |
| 140 | 140 | mod tests { |
| 141 | 141 | use super::*; |
| 142 | use crate::write::write_revision; | |
| 142 | 143 | |
| 143 | 144 | #[test] |
| 144 | 145 | fn protection_on_the_target_or_ancestor_prevents_layout_edits() { |
crates/onestore/src/page/write.rs+46-26| ... | ... | @@ -339,25 +339,34 @@ pub(crate) fn write_page( |
| 339 | 339 | if author.contains('\0') { |
| 340 | 340 | return Err(invalid("Choose an author name without NUL")); |
| 341 | 341 | } |
| 342 | let store = Store::parse(source)?; | |
| 343 | let index = RevisionIndex::parse(&store)?; | |
| 344 | index.validate_current()?; | |
| 345 | let document = Document::parse(&index)?; | |
| 346 | let pages = document.pages_in(space)?; | |
| 347 | let [page] = pages.as_slice() else { | |
| 348 | return Err(invalid("Choose an object space containing one active page")); | |
| 342 | // The parsed source is released before the writers parse their own images. | |
| 343 | let (page, before, existing) = { | |
| 344 | let store = Store::parse(source)?; | |
| 345 | let index = RevisionIndex::parse(&store)?; | |
| 346 | index.validate_current()?; | |
| 347 | let document = Document::parse(&index)?; | |
| 348 | let pages = document.pages_in(space)?; | |
| 349 | let [page] = pages.as_slice() else { | |
| 350 | return Err(invalid("Choose an object space containing one active page")); | |
| 351 | }; | |
| 352 | let existing = index | |
| 353 | .resolve_active(space)? | |
| 354 | .objects | |
| 355 | .keys() | |
| 356 | .copied() | |
| 357 | .collect(); | |
| 358 | (*page, Page::from_space(&document, space)?, existing) | |
| 349 | 359 | }; |
| 350 | let before = Page::from_space(&document, space)?; | |
| 351 | let raw = index.resolve_active(space)?; | |
| 352 | 360 | let mut lowering = Lowering { |
| 353 | 361 | image: source.to_vec(), |
| 362 | current: Some(before.clone()), | |
| 354 | 363 | space, |
| 355 | page: *page, | |
| 364 | page, | |
| 356 | 365 | author, |
| 357 | 366 | alias: BTreeMap::new(), |
| 358 | 367 | built: BTreeSet::new(), |
| 359 | 368 | }; |
| 360 | lowering.run(&before, after, &raw.objects.keys().copied().collect())?; | |
| 369 | lowering.run(&before, after, &existing)?; | |
| 361 | 370 | if lowering.image == source { |
| 362 | 371 | return Ok(lowering.image); |
| 363 | 372 | } |
| ... | ... | @@ -453,6 +462,8 @@ impl<'a> View<'a> { |
| 453 | 462 | |
| 454 | 463 | struct Lowering<'a> { |
| 455 | 464 | image: Vec<u8>, |
| 465 | /// The page as `image` stores it, until a writer changes the image. | |
| 466 | current: Option<Page>, | |
| 456 | 467 | space: ExGuid, |
| 457 | 468 | page: ExGuid, |
| 458 | 469 | author: &'a str, |
| ... | ... | @@ -478,15 +489,22 @@ impl Lowering<'_> { |
| 478 | 489 | } |
| 479 | 490 | |
| 480 | 491 | fn apply(&mut self, edit: impl FnOnce(&[u8]) -> Result<Vec<u8>, Error>) -> Result<(), Error> { |
| 481 | self.image = edit(&self.image)?; | |
| 492 | let image = edit(&self.image)?; | |
| 493 | if image != self.image { | |
| 494 | self.image = image; | |
| 495 | self.current = None; | |
| 496 | } | |
| 482 | 497 | Ok(()) |
| 483 | 498 | } |
| 484 | 499 | |
| 485 | fn current(&self) -> Result<Page, Error> { | |
| 486 | let store = Store::parse(&self.image)?; | |
| 487 | let index = RevisionIndex::parse(&store)?; | |
| 488 | let document = Document::parse(&index)?; | |
| 489 | Page::from_space(&document, self.space) | |
| 500 | fn current(&mut self) -> Result<Page, Error> { | |
| 501 | if self.current.is_none() { | |
| 502 | let store = Store::parse(&self.image)?; | |
| 503 | let index = RevisionIndex::parse(&store)?; | |
| 504 | let document = Document::parse(&index)?; | |
| 505 | self.current = Some(Page::from_space(&document, self.space)?); | |
| 506 | } | |
| 507 | Ok(self.current.clone().unwrap()) | |
| 490 | 508 | } |
| 491 | 509 | |
| 492 | 510 | fn run( |
| ... | ... | @@ -2921,10 +2939,6 @@ pub(crate) fn squash( |
| 2921 | 2939 | .iter() |
| 2922 | 2940 | .map(|(model, image)| (*image, *model)) |
| 2923 | 2941 | .collect(); |
| 2924 | let applied_store = Store::parse(applied)?; | |
| 2925 | let applied_index = RevisionIndex::parse(&applied_store)?; | |
| 2926 | // Payloads the typed edits embedded travel into the squashed transaction as well. | |
| 2927 | let source_store = Store::parse(source)?; | |
| 2928 | 2942 | let declared = |store: &Store<'_>| -> Vec<[u8; 16]> { |
| 2929 | 2943 | store |
| 2930 | 2944 | .lists |
| ... | ... | @@ -2934,7 +2948,10 @@ pub(crate) fn squash( |
| 2934 | 2948 | .filter_map(|node| node.payload.get(..16).and_then(|g| g.try_into().ok())) |
| 2935 | 2949 | .collect() |
| 2936 | 2950 | }; |
| 2937 | let existing = declared(&source_store); | |
| 2951 | let existing = declared(&Store::parse(source)?); | |
| 2952 | let applied_store = Store::parse(applied)?; | |
| 2953 | let applied_index = RevisionIndex::parse(&applied_store)?; | |
| 2954 | // Payloads the typed edits embedded travel into the squashed transaction as well. | |
| 2938 | 2955 | let mut payloads = Vec::new(); |
| 2939 | 2956 | for guid in declared(&applied_store) { |
| 2940 | 2957 | if !existing.contains(&guid) { |
| ... | ... | @@ -3018,10 +3035,13 @@ pub(crate) fn squash( |
| 3018 | 3035 | } |
| 3019 | 3036 | Ok(changes) |
| 3020 | 3037 | }; |
| 3021 | match protection { | |
| 3022 | Some(_) => crate::write::append_revisions(source, &payloads, protection, edit), | |
| 3023 | None => crate::write::write_revisions_with_payloads(source, &payloads, edit), | |
| 3024 | } | |
| 3038 | let validate = protection.is_none(); | |
| 3039 | let output = crate::write::build(source, &payloads, protection, validate, edit)?; | |
| 3040 | // The parsed images are released before the result is parsed. | |
| 3041 | drop(applied_index); | |
| 3042 | drop(applied_store); | |
| 3043 | crate::write::check(&output, validate)?; | |
| 3044 | Ok(output) | |
| 3025 | 3045 | } |
| 3026 | 3046 | |
| 3027 | 3047 | fn remap(object: &mut PropertyObject, rename: &BTreeMap<ExGuid, ExGuid>) -> Result<(), Error> { |
crates/onestore/src/paragraph.rs+3-3| ... | ... | @@ -3,7 +3,7 @@ use crate::{ |
| 3 | 3 | create::{current_timestamps, properties, string}, |
| 4 | 4 | document::{Document, Element, Kind}, |
| 5 | 5 | edit::{editable_parents, update_title}, |
| 6 | write::{PropertyObject, fresh_guid, write_revision}, | |
| 6 | write::{PropertyObject, fresh_guid, write_revision_on}, | |
| 7 | 7 | }; |
| 8 | 8 | use serde::{Deserialize, Serialize}; |
| 9 | 9 | use std::{ |
| ... | ... | @@ -252,7 +252,7 @@ impl ParagraphSplit { |
| 252 | 252 | object.set(&[(0x14001d7a, &modified)])?; |
| 253 | 253 | } |
| 254 | 254 | update_title(&store, &raw, view, &pages, &mut changed)?; |
| 255 | write_revision(source, space, |_| Ok(changed)) | |
| 255 | write_revision_on(&index, space, |_| Ok(changed)) | |
| 256 | 256 | } |
| 257 | 257 | } |
| 258 | 258 | |
| ... | ... | @@ -682,7 +682,7 @@ impl ParagraphJoin { |
| 682 | 682 | pending.extend(parents.get(&id).into_iter().flatten().copied()); |
| 683 | 683 | } |
| 684 | 684 | update_title(&store, &raw, view, &pages, &mut changed)?; |
| 685 | write_revision(source, space, |_| Ok(changed)) | |
| 685 | write_revision_on(&index, space, |_| Ok(changed)) | |
| 686 | 686 | } |
| 687 | 687 | } |
| 688 | 688 |
crates/onestore/src/store.rs+1| ... | ... | @@ -505,6 +505,7 @@ impl<'a> Store<'a> { |
| 505 | 505 | pending.push(reference); |
| 506 | 506 | } |
| 507 | 507 | } |
| 508 | nodes.shrink_to_fit(); | |
| 508 | 509 | lists.insert(list_id.unwrap(), NodeList { fragments, nodes }); |
| 509 | 510 | } |
| 510 | 511 | Ok(Self { |
crates/onestore/src/tree.rs+3-2| ... | ... | @@ -3,7 +3,7 @@ use crate::{ |
| 3 | 3 | create::{current_timestamps, properties, string}, |
| 4 | 4 | document::{Document, Kind, Revision}, |
| 5 | 5 | edit::{editable_parents, update_title}, |
| 6 | write::{PropertyObject, fresh_guid, write_revision}, | |
| 6 | write::{PropertyObject, fresh_guid, write_revision_on}, | |
| 7 | 7 | }; |
| 8 | 8 | use serde::{Deserialize, Serialize}; |
| 9 | 9 | use std::{ |
| ... | ... | @@ -371,7 +371,7 @@ impl TreeEdit { |
| 371 | 371 | changed.insert(self.object, object); |
| 372 | 372 | } |
| 373 | 373 | update_title(&store, &raw, view, &pages, &mut changed)?; |
| 374 | write_revision(source, space, |_| Ok(changed)) | |
| 374 | write_revision_on(&index, space, |_| Ok(changed)) | |
| 375 | 375 | } |
| 376 | 376 | } |
| 377 | 377 | |
| ... | ... | @@ -416,6 +416,7 @@ fn checked_path( |
| 416 | 416 | #[cfg(test)] |
| 417 | 417 | mod tests { |
| 418 | 418 | use super::*; |
| 419 | use crate::write::write_revision; | |
| 419 | 420 | use crate::{Insertion, PreparedEdit}; |
| 420 | 421 | |
| 421 | 422 | #[test] |
crates/onestore/src/write.rs+82-15| ... | ... | @@ -305,8 +305,11 @@ pub fn replace_property_bytes( |
| 305 | 305 | message: "Property does not contain scalar bytes", |
| 306 | 306 | }); |
| 307 | 307 | } |
| 308 | let store = Store::parse(source)?; | |
| 309 | let index = RevisionIndex::parse(&store)?; | |
| 310 | index.validate_current()?; | |
| 308 | 311 | replace_objects( |
| 309 | source, | |
| 312 | &index, | |
| 310 | 313 | space, |
| 311 | 314 | &[ObjectEdit { |
| 312 | 315 | object: object_id, |
| ... | ... | @@ -316,12 +319,13 @@ pub fn replace_property_bytes( |
| 316 | 319 | ) |
| 317 | 320 | } |
| 318 | 321 | |
| 322 | /// Patches objects of a source the caller has parsed and validated. | |
| 319 | 323 | pub(crate) fn replace_objects( |
| 320 | source: &[u8], | |
| 324 | index: &RevisionIndex<'_>, | |
| 321 | 325 | space: ExGuid, |
| 322 | 326 | edits: &[ObjectEdit<'_>], |
| 323 | 327 | ) -> Result<Vec<u8>> { |
| 324 | write_revision(source, space, |revision| { | |
| 328 | write_revision_on(index, space, |revision| { | |
| 325 | 329 | let mut changed = BTreeMap::new(); |
| 326 | 330 | for edit in edits { |
| 327 | 331 | if changed.contains_key(&edit.object) { |
| ... | ... | @@ -788,23 +792,86 @@ pub(crate) fn write_revisions_with_payloads( |
| 788 | 792 | payloads: &[([u8; 16], &[u8])], |
| 789 | 793 | edit: impl FnOnce(&RevisionIndex<'_>) -> Result<BTreeMap<ExGuid, RevisionEdit>>, |
| 790 | 794 | ) -> Result<Vec<u8>> { |
| 791 | let store = Store::parse(source)?; | |
| 792 | RevisionIndex::parse(&store)?.validate_current()?; | |
| 793 | let output = append_revisions(source, payloads, None, edit)?; | |
| 794 | let store = Store::parse(&output)?; | |
| 795 | RevisionIndex::parse(&store)?.validate_current()?; | |
| 796 | Ok(output) | |
| 795 | publish(source, payloads, None, true, edit) | |
| 797 | 796 | } |
| 798 | 797 | |
| 799 | 798 | /// `write_revisions_with_payloads` without validating that current revisions are |
| 800 | 799 | /// complete: a protected section is validated by unlocking it, through `protection`. |
| 800 | #[cfg(feature = "protected")] | |
| 801 | 801 | pub(crate) fn append_revisions( |
| 802 | 802 | source: &[u8], |
| 803 | 803 | payloads: &[([u8; 16], &[u8])], |
| 804 | 804 | protection: Option<&dyn Protection>, |
| 805 | 805 | edit: impl FnOnce(&RevisionIndex<'_>) -> Result<BTreeMap<ExGuid, RevisionEdit>>, |
| 806 | ) -> Result<Vec<u8>> { | |
| 807 | publish(source, payloads, protection, false, edit) | |
| 808 | } | |
| 809 | ||
| 810 | fn publish( | |
| 811 | source: &[u8], | |
| 812 | payloads: &[([u8; 16], &[u8])], | |
| 813 | protection: Option<&dyn Protection>, | |
| 814 | validate: bool, | |
| 815 | edit: impl FnOnce(&RevisionIndex<'_>) -> Result<BTreeMap<ExGuid, RevisionEdit>>, | |
| 816 | ) -> Result<Vec<u8>> { | |
| 817 | // The parsed source is released before the result is parsed. | |
| 818 | let output = build(source, payloads, protection, validate, edit)?; | |
| 819 | check(&output, validate)?; | |
| 820 | Ok(output) | |
| 821 | } | |
| 822 | ||
| 823 | /// Parses a written image, and with `validate` requires its current revisions complete. | |
| 824 | pub(crate) fn check(output: &[u8], validate: bool) -> Result<()> { | |
| 825 | let store = Store::parse(output)?; | |
| 826 | let index = RevisionIndex::parse(&store)?; | |
| 827 | if validate { | |
| 828 | index.validate_current()?; | |
| 829 | } | |
| 830 | Ok(()) | |
| 831 | } | |
| 832 | ||
| 833 | /// The written image, unchecked: `check` follows once the caller has released whatever | |
| 834 | /// `edit` borrowed. | |
| 835 | pub(crate) fn build( | |
| 836 | source: &[u8], | |
| 837 | payloads: &[([u8; 16], &[u8])], | |
| 838 | protection: Option<&dyn Protection>, | |
| 839 | validate: bool, | |
| 840 | edit: impl FnOnce(&RevisionIndex<'_>) -> Result<BTreeMap<ExGuid, RevisionEdit>>, | |
| 806 | 841 | ) -> Result<Vec<u8>> { |
| 807 | 842 | let store = Store::parse(source)?; |
| 843 | let index = RevisionIndex::parse(&store)?; | |
| 844 | if validate { | |
| 845 | index.validate_current()?; | |
| 846 | } | |
| 847 | build_on(&index, payloads, protection, edit) | |
| 848 | } | |
| 849 | ||
| 850 | /// `write_revision` on a source the caller has parsed and validated. | |
| 851 | pub(crate) fn write_revision_on( | |
| 852 | index: &RevisionIndex<'_>, | |
| 853 | space: ExGuid, | |
| 854 | edit: impl FnOnce(&crate::ResolvedRevision<'_>) -> Result<BTreeMap<ExGuid, PropertyObject>>, | |
| 855 | ) -> Result<Vec<u8>> { | |
| 856 | let output = build_on(index, &[], None, |index| { | |
| 857 | let revision = index.resolve(space, index.active(space)?)?; | |
| 858 | Ok(BTreeMap::from([( | |
| 859 | space, | |
| 860 | RevisionEdit::Update(edit(&revision)?), | |
| 861 | )])) | |
| 862 | })?; | |
| 863 | check(&output, true)?; | |
| 864 | Ok(output) | |
| 865 | } | |
| 866 | ||
| 867 | fn build_on( | |
| 868 | index: &RevisionIndex<'_>, | |
| 869 | payloads: &[([u8; 16], &[u8])], | |
| 870 | protection: Option<&dyn Protection>, | |
| 871 | edit: impl FnOnce(&RevisionIndex<'_>) -> Result<BTreeMap<ExGuid, RevisionEdit>>, | |
| 872 | ) -> Result<Vec<u8>> { | |
| 873 | let store = index.store; | |
| 874 | let source = store.data; | |
| 808 | 875 | let is_section = store.header.file_type == FileType::Section; |
| 809 | 876 | if !store.checksum_mismatches.is_empty() { |
| 810 | 877 | return Err(Error { |
| ... | ... | @@ -812,13 +879,14 @@ pub(crate) fn append_revisions( |
| 812 | 879 | message: "Cannot write a file with transaction checksum damage", |
| 813 | 880 | }); |
| 814 | 881 | } |
| 815 | let index = RevisionIndex::parse(&store)?; | |
| 816 | 882 | let resolve = |space, rid| match protection { |
| 817 | 883 | Some(protection) => protection.resolve(space, rid), |
| 818 | 884 | None => index.resolve(space, rid), |
| 819 | 885 | }; |
| 820 | let changes = edit(&index)?; | |
| 821 | let mut output = source.to_vec(); | |
| 886 | let changes = edit(index)?; | |
| 887 | // Room for the revision, so appending does not double the image. | |
| 888 | let mut output = Vec::with_capacity(source.len() + source.len() / 16 + (1 << 16)); | |
| 889 | output.extend_from_slice(source); | |
| 822 | 890 | // Native files reserve 1 KiB per transaction-log fragment; a fragment that ends the file |
| 823 | 891 | // keeps that room before the first new chunk. |
| 824 | 892 | let tail = store.transaction_fragments.last().unwrap().chunk; |
| ... | ... | @@ -1348,12 +1416,12 @@ pub(crate) fn append_revisions( |
| 1348 | 1416 | let space_node = root |
| 1349 | 1417 | .nodes |
| 1350 | 1418 | .iter() |
| 1351 | .find(|node| node.id == 8 && node.fields(&store).exguid() == Ok(space)) | |
| 1419 | .find(|node| node.id == 8 && node.fields(store).exguid() == Ok(space)) | |
| 1352 | 1420 | .ok_or(Error { |
| 1353 | 1421 | offset: 0, |
| 1354 | 1422 | message: "Object space is absent from the root list", |
| 1355 | 1423 | })?; |
| 1356 | let space_list = space_node.referenced_list(&store)?; | |
| 1424 | let space_list = space_node.referenced_list(store)?; | |
| 1357 | 1425 | let revision_node = space_list.iter().rfind(|node| node.id == 0x10).unwrap(); |
| 1358 | 1426 | let Some(Reference::NodeList(manifest_reference)) = revision_node.reference else { |
| 1359 | 1427 | unreachable!() |
| ... | ... | @@ -1522,6 +1590,5 @@ pub(crate) fn append_revisions( |
| 1522 | 1590 | message: "File generation counter is exhausted", |
| 1523 | 1591 | })?; |
| 1524 | 1592 | output[228..236].copy_from_slice(&generation.to_le_bytes()); |
| 1525 | RevisionIndex::parse(&Store::parse(&output)?)?; | |
| 1526 | 1593 | Ok(output) |
| 1527 | 1594 | } |
tools/cache_images.py created+17| ... | ... | @@ -0,0 +1,17 @@ |
| 1 | """The cache's working image, stored as the byte ranges where it differs from the base.""" | |
| 2 | import struct | |
| 3 | ||
| 4 | ||
| 5 | def working(connection): | |
| 6 | base, patch = connection.execute('SELECT base, working FROM replica WHERE id=1').fetchone() | |
| 7 | if not patch: return base | |
| 8 | length, = struct.unpack_from('<Q', patch) | |
| 9 | image = bytearray(base[:length].ljust(length, b'\0')) | |
| 10 | at = 8 | |
| 11 | while at < len(patch): | |
| 12 | offset, size = struct.unpack_from('<QQ', patch, at) | |
| 13 | at += 16 | |
| 14 | assert at + size <= len(patch) and offset + size <= length, 'Damaged working image' | |
| 15 | image[offset:offset + size] = patch[at:at + size] | |
| 16 | at += size | |
| 17 | return bytes(image) |
tools/offline_publication_crash.py+2-1| ... | ... | @@ -1,6 +1,7 @@ |
| 1 | 1 | #!/usr/bin/env python3 |
| 2 | 2 | """Kill owned processes across local-cache/remote-file publication boundaries.""" |
| 3 | 3 | import argparse |
| 4 | import cache_images | |
| 4 | 5 | import hashlib |
| 5 | 6 | import json |
| 6 | 7 | import os |
| ... | ... | @@ -186,7 +187,7 @@ def run(source, output): |
| 186 | 187 | try: |
| 187 | 188 | assert connection.execute('PRAGMA quick_check').fetchall() == [('ok',)] |
| 188 | 189 | assert connection.execute('PRAGMA foreign_key_check').fetchall() == [] |
| 189 | local_image = save_image(output, connection.execute('SELECT working FROM replica WHERE id=1').fetchone()[0]) | |
| 190 | local_image = save_image(output, cache_images.working(connection)) | |
| 190 | 191 | finally: |
| 191 | 192 | connection.close() |
| 192 | 193 | result = {'case': index, 'phase': phase, 'during_confirmation': confirmation, 'before_status': before['status'], 'after_status': after['status'], |
tools/test_offline_history.py+2-2| ... | ... | @@ -141,7 +141,7 @@ class OfflineLedgerTests(unittest.TestCase): |
| 141 | 141 | connection = sqlite3.connect(output / 'rust' / f'{actor}.sqlite') |
| 142 | 142 | connection.executescript('CREATE TABLE receipts(edit_id INTEGER, revision TEXT); CREATE TABLE edits(id INTEGER); CREATE TABLE attempt(id INTEGER); CREATE TABLE conflicts(id INTEGER); CREATE TABLE replica(id INTEGER, base BLOB, working BLOB);') |
| 143 | 143 | connection.executemany('INSERT INTO receipts VALUES (?,?)', [(op+1, f'{actor}-{op}') for op in range(2)]) |
| 144 | connection.execute('INSERT INTO replica VALUES (1, ?, ?)', (b'opaque image', b'opaque image')) | |
| 144 | connection.execute('INSERT INTO replica VALUES (1, ?, ?)', (b'opaque image', b'')) | |
| 145 | 145 | connection.commit() |
| 146 | 146 | connection.close() |
| 147 | 147 | for i in range(4): |
| ... | ... | @@ -166,7 +166,7 @@ class OfflineLedgerTests(unittest.TestCase): |
| 166 | 166 | connection = sqlite3.connect(output / 'rust/w0.sqlite') |
| 167 | 167 | for sql, undo in [("UPDATE receipts SET revision='wrong' WHERE edit_id=1", "UPDATE receipts SET revision='w0-0' WHERE edit_id=1"), |
| 168 | 168 | ('INSERT INTO attempt VALUES (1)', 'DELETE FROM attempt'), |
| 169 | ("UPDATE replica SET working=X'00'", "UPDATE replica SET working=base")]: | |
| 169 | ("UPDATE replica SET working=X'00'", "UPDATE replica SET working=X''")]: | |
| 170 | 170 | connection.execute(sql) |
| 171 | 171 | connection.commit() |
| 172 | 172 | with self.assertRaises(AssertionError): verify(output) |
tools/verify_offline.py+1-1| ... | ... | @@ -46,7 +46,7 @@ def verify(output, max_gap=120): |
| 46 | 46 | persisted = dict(connection.execute('SELECT edit_id, revision FROM receipts')) |
| 47 | 47 | assert persisted == {id: event['revision'] for id, event in receipts.items()}, 'SQLite receipts differ from observed acknowledgements' |
| 48 | 48 | base, working = connection.execute('SELECT base, working FROM replica WHERE id=1').fetchone() |
| 49 | assert base == working, 'A drained cache retained a divergent local branch' | |
| 49 | assert working == b'', 'A drained cache retained a divergent local branch' | |
| 50 | 50 | caches[actor] = {'receipts': len(persisted), 'image_bytes': len(base), 'image_sha256': hashlib.sha256(base).hexdigest(), 'database_sha256': hashlib.sha256(path.read_bytes()).hexdigest()} |
| 51 | 51 | finally: |
| 52 | 52 | connection.close() |