| author | |
| committer | |
| log | d956049b4eeac2458cf646e61eb29a90390ff59f |
| tree | 2372ed00ca4d2c674364fa189f10150b52129f97 |
| parent | 613bb6e8c7022b768402c9801dc069a9c5bb78bd |
| signature | Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU |
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-astra4 files changed, 70 insertions(+), 19 deletions(-)
crates/notebook/src/lib.rs+18-14| ... | @@ -453,24 +453,28 @@ fn pending(connection: &Connection) -> Result<Vec<PendingEdit>> { | ... | @@ -453,24 +453,28 @@ fn pending(connection: &Connection) -> Result<Vec<PendingEdit>> { |
| 453 | let mut rows = query.query([])?; | 453 | let mut rows = query.query([])?; |
| 454 | let mut edits = Vec::new(); | 454 | let mut edits = Vec::new(); |
| 455 | while let Some(row) = rows.next()? { | 455 | while let Some(row) = rows.next()? { |
| 456 | let edit = PendingEdit { | 456 | edits.push(pending_edit(row)?); |
| 457 | id: u64::try_from(row.get::<_, i64>(0)?).map_err(io::Error::other)?, | ||
| 458 | space: row.get::<_, String>(1)?.parse()?, | ||
| 459 | operation: serde_json::from_str(&row.get::<_, String>(2)?) | ||
| 460 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?, | ||
| 461 | }; | ||
| 462 | if matches!(&edit.operation, Operation::CreatePage(page) if page.space() != edit.space) { | ||
| 463 | return Err(io::Error::new( | ||
| 464 | io::ErrorKind::InvalidData, | ||
| 465 | "Cached page identity differs from its creation intent", | ||
| 466 | ) | ||
| 467 | .into()); | ||
| 468 | } | ||
| 469 | edits.push(edit); | ||
| 470 | } | 457 | } |
| 471 | Ok(edits) | 458 | Ok(edits) |
| 472 | } | 459 | } |
| 473 | 460 | ||
| 461 | fn pending_edit(row: &rusqlite::Row<'_>) -> Result<PendingEdit> { | ||
| 462 | let edit = PendingEdit { | ||
| 463 | id: u64::try_from(row.get::<_, i64>(0)?).map_err(io::Error::other)?, | ||
| 464 | space: row.get::<_, String>(1)?.parse()?, | ||
| 465 | operation: serde_json::from_str(&row.get::<_, String>(2)?) | ||
| 466 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?, | ||
| 467 | }; | ||
| 468 | if matches!(&edit.operation, Operation::CreatePage(page) if page.space() != edit.space) { | ||
| 469 | return Err(io::Error::new( | ||
| 470 | io::ErrorKind::InvalidData, | ||
| 471 | "Cached page identity differs from its creation intent", | ||
| 472 | ) | ||
| 473 | .into()); | ||
| 474 | } | ||
| 475 | Ok(edit) | ||
| 476 | } | ||
| 477 | |||
| 474 | fn cache_connection(path: &Path) -> Result<Connection> { | 478 | fn cache_connection(path: &Path) -> Result<Connection> { |
| 475 | let connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_WRITE)?; | 479 | let connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_WRITE)?; |
| 476 | connection.busy_timeout(Duration::ZERO)?; | 480 | connection.busy_timeout(Duration::ZERO)?; |
crates/notebook/src/sync.rs+7-1| ... | @@ -93,7 +93,13 @@ impl Replica { | ... | @@ -93,7 +93,13 @@ impl Replica { |
| 93 | ) | 93 | ) |
| 94 | .into()); | 94 | .into()); |
| 95 | } | 95 | } |
| 96 | let Some(intent) = pending(&transaction)?.into_iter().next() else { | 96 | let intent = { |
| 97 | let mut query = transaction | ||
| 98 | .prepare("SELECT id, space, operation FROM edits ORDER BY id LIMIT 1")?; | ||
| 99 | let mut rows = query.query([])?; | ||
| 100 | rows.next()?.map(pending_edit).transpose()? | ||
| 101 | }; | ||
| 102 | let Some(intent) = intent else { | ||
| 97 | transaction.execute( | 103 | transaction.execute( |
| 98 | "UPDATE replica SET base=?1, working=?1 WHERE id=1", | 104 | "UPDATE replica SET base=?1, working=?1 WHERE id=1", |
| 99 | [&snapshot], | 105 | [&snapshot], |
crates/notebook/tests/sync.rs+35| ... | @@ -995,3 +995,38 @@ fn restore_conflicts_keep_dependent_work_and_reject_stale_review() { | ... | @@ -995,3 +995,38 @@ fn restore_conflicts_keep_dependent_work_and_reject_stale_review() { |
| 995 | assert_eq!(archive.pending().unwrap(), pending); | 995 | assert_eq!(archive.pending().unwrap(), pending); |
| 996 | assert_eq!(archive.status(head).unwrap(), Some(conflict)); | 996 | assert_eq!(archive.status(head).unwrap(), Some(conflict)); |
| 997 | } | 997 | } |
| 998 | |||
| 999 | #[test] | ||
| 1000 | fn cache_open_validates_the_tail_before_head_only_synchronization() { | ||
| 1001 | for damage in [ | ||
| 1002 | "operation='{}'", | ||
| 1003 | "space=(SELECT space FROM edits ORDER BY id LIMIT 1)", | ||
| 1004 | ] { | ||
| 1005 | let directory = tempfile::tempdir().unwrap(); | ||
| 1006 | let path = directory.path().join("cache.sqlite"); | ||
| 1007 | let source = onestore::create_section("tail.one", "Original", "Fixture").unwrap(); | ||
| 1008 | let (_, object, _) = text(&source); | ||
| 1009 | let cache = Replica::create(&path, &source).unwrap(); | ||
| 1010 | save(&cache, object, 0..0, "Head ").unwrap(); | ||
| 1011 | let page = onestore::PageCreation::new(None, Some("Tail"), "Fixture").unwrap(); | ||
| 1012 | let tail = cache | ||
| 1013 | .create_page(&cache.snapshot().unwrap(), &page) | ||
| 1014 | .unwrap() | ||
| 1015 | .unwrap(); | ||
| 1016 | assert_eq!(cache.pending().unwrap().len(), 2); | ||
| 1017 | drop(cache); | ||
| 1018 | let database = rusqlite::Connection::open(&path).unwrap(); | ||
| 1019 | database | ||
| 1020 | .execute( | ||
| 1021 | &format!("UPDATE edits SET {damage} WHERE id=?1"), | ||
| 1022 | [i64::try_from(tail).unwrap()], | ||
| 1023 | ) | ||
| 1024 | .unwrap(); | ||
| 1025 | drop(database); | ||
| 1026 | let damaged = std::fs::read(&path).unwrap(); | ||
| 1027 | assert!( | ||
| 1028 | matches!(Replica::open(&path), Err(Error::Io(error)) if error.kind() == io::ErrorKind::InvalidData) | ||
| 1029 | ); | ||
| 1030 | assert_eq!(std::fs::read(&path).unwrap(), damaged); | ||
| 1031 | } | ||
| 1032 | } |
crates/notebook/tests/sync_worker.rs+10-4| ... | @@ -279,14 +279,20 @@ fn dropping_during_publication_is_nonblocking_and_retains_ownership_until_recove | ... | @@ -279,14 +279,20 @@ fn dropping_during_publication_is_nonblocking_and_retains_ownership_until_recove |
| 279 | |_| {}, | 279 | |_| {}, |
| 280 | ); | 280 | ); |
| 281 | assert!(matches!(start, Err(error) if error.kind() == io::ErrorKind::WouldBlock)); | 281 | assert!(matches!(start, Err(error) if error.kind() == io::ErrorKind::WouldBlock)); |
| 282 | let weak = Arc::downgrade(&cache); | ||
| 283 | drop(cache); | 282 | drop(cache); |
| 283 | assert!(matches!(Replica::open(&path), Err(Error::Database(error)) | ||
| 284 | if error.sqlite_error_code() == Some(rusqlite::ErrorCode::DatabaseBusy))); | ||
| 284 | resume_tx.send(()).unwrap(); | 285 | resume_tx.send(()).unwrap(); |
| 285 | while weak.upgrade().is_some() { | 286 | let cache = Arc::new(loop { |
| 287 | match Replica::open(&path) { | ||
| 288 | Ok(cache) => break cache, | ||
| 289 | Err(Error::Database(error)) | ||
| 290 | if error.sqlite_error_code() == Some(rusqlite::ErrorCode::DatabaseBusy) => {} | ||
| 291 | Err(error) => panic!("{error}"), | ||
| 292 | } | ||
| 286 | assert!(stopped.elapsed() < Duration::from_secs(5)); | 293 | assert!(stopped.elapsed() < Duration::from_secs(5)); |
| 287 | std::thread::sleep(Duration::from_millis(1)); | 294 | std::thread::sleep(Duration::from_millis(1)); |
| 288 | } | 295 | }); |
| 289 | let cache = Arc::new(Replica::open(&path).unwrap()); | ||
| 290 | assert!(matches!( | 296 | assert!(matches!( |
| 291 | cache.status(id).unwrap(), | 297 | cache.status(id).unwrap(), |
| 292 | Some(EditStatus::AwaitingConfirmation { .. }) | 298 | Some(EditStatus::AwaitingConfirmation { .. }) |