From d956049b4eeac2458cf646e61eb29a90390ff59f Mon Sep 17 00:00:00 2001 From: clover caruso Date: Fri, 11 Sep 2026 04:01:57 -0700 Subject: [PATCH] perf: decode only the selected synchronization head Avoid materializing the complete pending queue for every publication step. Retain full cache-open validation and share row validation with queue inspection. Add malformed-tail coverage and make the cancellation test observe actual SQLite lock release. Verified the notebook suite and an isolated 1000-edit workload with every receipt and recovery image retained. Assisted-by: gpt-6-astra --- crates/notebook/src/lib.rs | 32 ++++++++++++++----------- crates/notebook/src/sync.rs | 8 ++++++- crates/notebook/tests/sync.rs | 35 ++++++++++++++++++++++++++++ crates/notebook/tests/sync_worker.rs | 14 +++++++---- 4 files changed, 70 insertions(+), 19 deletions(-) diff --git a/crates/notebook/src/lib.rs b/crates/notebook/src/lib.rs index 9793056e4ec3e3f25e196febc341e672e31f3477..0e4cd47db025b66e4136562154dd7dfe011acbae 100644 --- a/crates/notebook/src/lib.rs +++ b/crates/notebook/src/lib.rs @@ -453,24 +453,28 @@ fn pending(connection: &Connection) -> Result> { let mut rows = query.query([])?; let mut edits = Vec::new(); while let Some(row) = rows.next()? { - let edit = PendingEdit { - id: u64::try_from(row.get::<_, i64>(0)?).map_err(io::Error::other)?, - space: row.get::<_, String>(1)?.parse()?, - operation: serde_json::from_str(&row.get::<_, String>(2)?) - .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?, - }; - if matches!(&edit.operation, Operation::CreatePage(page) if page.space() != edit.space) { - return Err(io::Error::new( - io::ErrorKind::InvalidData, - "Cached page identity differs from its creation intent", - ) - .into()); - } - edits.push(edit); + edits.push(pending_edit(row)?); } Ok(edits) } +fn pending_edit(row: &rusqlite::Row<'_>) -> Result { + let edit = PendingEdit { + id: u64::try_from(row.get::<_, i64>(0)?).map_err(io::Error::other)?, + space: row.get::<_, String>(1)?.parse()?, + operation: serde_json::from_str(&row.get::<_, String>(2)?) + .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?, + }; + if matches!(&edit.operation, Operation::CreatePage(page) if page.space() != edit.space) { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Cached page identity differs from its creation intent", + ) + .into()); + } + Ok(edit) +} + fn cache_connection(path: &Path) -> Result { let connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_WRITE)?; connection.busy_timeout(Duration::ZERO)?; diff --git a/crates/notebook/src/sync.rs b/crates/notebook/src/sync.rs index 9c9426dbc8536bcde85cb7ff97a593a16beef255..56b71f181a3cf4db24010ae7a15e3a8cbe09367b 100644 --- a/crates/notebook/src/sync.rs +++ b/crates/notebook/src/sync.rs @@ -93,7 +93,13 @@ impl Replica { ) .into()); } - let Some(intent) = pending(&transaction)?.into_iter().next() else { + let intent = { + let mut query = transaction + .prepare("SELECT id, space, operation FROM edits ORDER BY id LIMIT 1")?; + let mut rows = query.query([])?; + rows.next()?.map(pending_edit).transpose()? + }; + let Some(intent) = intent else { transaction.execute( "UPDATE replica SET base=?1, working=?1 WHERE id=1", [&snapshot], diff --git a/crates/notebook/tests/sync.rs b/crates/notebook/tests/sync.rs index b29297f1b3d5b755dcca93c3e2d48286725d36a6..047b59de8873a97452f134398402d525a4724027 100644 --- a/crates/notebook/tests/sync.rs +++ b/crates/notebook/tests/sync.rs @@ -995,3 +995,38 @@ fn restore_conflicts_keep_dependent_work_and_reject_stale_review() { assert_eq!(archive.pending().unwrap(), pending); assert_eq!(archive.status(head).unwrap(), Some(conflict)); } + +#[test] +fn cache_open_validates_the_tail_before_head_only_synchronization() { + for damage in [ + "operation='{}'", + "space=(SELECT space FROM edits ORDER BY id LIMIT 1)", + ] { + let directory = tempfile::tempdir().unwrap(); + let path = directory.path().join("cache.sqlite"); + let source = onestore::create_section("tail.one", "Original", "Fixture").unwrap(); + let (_, object, _) = text(&source); + let cache = Replica::create(&path, &source).unwrap(); + save(&cache, object, 0..0, "Head ").unwrap(); + let page = onestore::PageCreation::new(None, Some("Tail"), "Fixture").unwrap(); + let tail = cache + .create_page(&cache.snapshot().unwrap(), &page) + .unwrap() + .unwrap(); + assert_eq!(cache.pending().unwrap().len(), 2); + drop(cache); + let database = rusqlite::Connection::open(&path).unwrap(); + database + .execute( + &format!("UPDATE edits SET {damage} WHERE id=?1"), + [i64::try_from(tail).unwrap()], + ) + .unwrap(); + drop(database); + let damaged = std::fs::read(&path).unwrap(); + assert!( + matches!(Replica::open(&path), Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidData) + ); + assert_eq!(std::fs::read(&path).unwrap(), damaged); + } +} diff --git a/crates/notebook/tests/sync_worker.rs b/crates/notebook/tests/sync_worker.rs index 5ebcfed085021e70843727723a03d088b035f881..aa85f40166a8e544cb10f40910826ddb66a62b60 100644 --- a/crates/notebook/tests/sync_worker.rs +++ b/crates/notebook/tests/sync_worker.rs @@ -279,14 +279,20 @@ fn dropping_during_publication_is_nonblocking_and_retains_ownership_until_recove |_| {}, ); assert!(matches!(start, Err(error) if error.kind() == io::ErrorKind::WouldBlock)); - let weak = Arc::downgrade(&cache); drop(cache); + assert!(matches!(Replica::open(&path), Err(Error::Database(error)) + if error.sqlite_error_code() == Some(rusqlite::ErrorCode::DatabaseBusy))); resume_tx.send(()).unwrap(); - while weak.upgrade().is_some() { + let cache = Arc::new(loop { + match Replica::open(&path) { + Ok(cache) => break cache, + Err(Error::Database(error)) + if error.sqlite_error_code() == Some(rusqlite::ErrorCode::DatabaseBusy) => {} + Err(error) => panic!("{error}"), + } assert!(stopped.elapsed() < Duration::from_secs(5)); std::thread::sleep(Duration::from_millis(1)); - } - let cache = Arc::new(Replica::open(&path).unwrap()); + }); assert!(matches!( cache.status(id).unwrap(), Some(EditStatus::AwaitingConfirmation { .. }) -- 2.54.0