From 9a82329b95a69c3d3e1844133a9d3045db885758 Mon Sep 17 00:00:00 2001 From: clover caruso Date: Sat, 3 Oct 2026 18:42:36 -0700 Subject: [PATCH] feat: batch Live Share guests' edits at the host Retain each guest batch's revision identities while rebasing its ops at a per-section writer and publishing consecutive revisions in one guarded commit. Confirm lost receipts against the host's revision history. Assisted-by: gpt-6.1-sol --- crates/notebook/examples/live_crowd.rs | 5 + crates/notebook/src/live/share.rs | 86 +++++++- crates/notebook/src/live/share/batch.rs | 258 ++++++++++++++++++++++++ crates/notebook/src/live/wire.rs | 5 + crates/notebook/src/merge.rs | 24 ++- crates/notebook/src/sync.rs | 38 +++- crates/notebook/src/working.rs | 8 + crates/notebook/tests/live_batch.rs | 179 ++++++++++++++++ crates/onestore/src/commit.rs | 59 ++++++ 9 files changed, 647 insertions(+), 15 deletions(-) create mode 100644 crates/notebook/src/live/share/batch.rs create mode 100644 crates/notebook/tests/live_batch.rs diff --git a/crates/notebook/examples/live_crowd.rs b/crates/notebook/examples/live_crowd.rs index 57ed5d7282faa91063a3ed7e68e08cdf46c3c10e..362ea4ba5852cadaf3c8fa66e05cd14bb53f062d 100644 --- a/crates/notebook/examples/live_crowd.rs +++ b/crates/notebook/examples/live_crowd.rs @@ -141,6 +141,11 @@ fn host(url: &str, paragraphs: usize) { println!("code {}", host.code().unwrap()); let _ = std::io::stdin().read_to_end(&mut Vec::new()); let image = std::fs::read(&file).unwrap(); + if let Some(output) = std::env::var_os("SNOWBOUND_LIVE_EVIDENCE") { + let output = Path::new(&output); + std::fs::create_dir_all(output).unwrap(); + std::fs::copy(&file, output.join("Garden.one")).unwrap(); + } let size = image.len(); let arena = onestore::Arena::default(); let mut section = onestore::Section::open(&arena, image).unwrap(); diff --git a/crates/notebook/src/live/share.rs b/crates/notebook/src/live/share.rs index bea88c2a0668175924ab2988e0a1479fe7b8fefc..1bc7b04bbd6a9281f062119a896f8325137329f9 100644 --- a/crates/notebook/src/live/share.rs +++ b/crates/notebook/src/live/share.rs @@ -1,8 +1,7 @@ //! Live Share: a notebook one Snowbound holds, opened on others through a short code. The -//! host serves its notebook's storage verbs (`session::Storage`) to each peer in the share's -//! room; a guest runs the replica, queue and merge it runs on an SMB share against those -//! verbs, so offline queueing, rebases and conflict pages work as there, and the host's files -//! stay what its own storage writes. A guest first meets the host in the code's room, where +//! host serves its notebook's storage verbs (`session::Storage`) and batches guests' ops into +//! guarded publications. A guest runs the same replica and durable queue as on an SMB share; +//! protected sections and older peers use transactions. A guest first meets the host in the code's room, where //! the host welcomes it with the share's room and secret; a new share has a new secret, so //! stopping retires every guest. Large bodies travel a chunk at a time, each answered before //! the next, so a relay never holds much for a slow peer. @@ -29,6 +28,8 @@ use std::{ time::{Duration, Instant}, }; +mod batch; + /// The most bytes one message of a read or an upload carries. const CHUNK: usize = 128 << 10; /// The most bytes of chunks a guest has asked for and not yet been given. @@ -234,6 +235,7 @@ impl Host { snapshots: Mutex::default(), puts: Mutex::default(), guests: Mutex::default(), + writers: Mutex::default(), host: Mutex::default(), room: OnceLock::new(), }); @@ -393,6 +395,7 @@ struct Served { /// Bytes a later request carries, by guest and upload. puts: Mutex>>, guests: Mutex>, + writers: Mutex>>, /// Hears the paths guests changed, as the host's own notebook should. host: Mutex>, /// The share's room, to tell guests what changed. @@ -524,7 +527,7 @@ impl Served { guest.held.truncate(IMAGES); } - fn handle(&self, peer: &[u8; 16], kind: u16, body: &[u8]) -> Reply { + fn handle(self: &Arc, peer: &[u8; 16], kind: u16, body: &[u8]) -> Reply { let request = match minicbor::decode::(body) { Ok(request) => request, Err(_) => { @@ -541,7 +544,7 @@ impl Served { } } - fn answer(&self, peer: &[u8; 16], kind: u16, request: Request) -> Result { + fn answer(self: &Arc, peer: &[u8; 16], kind: u16, request: Request) -> Result { let path = request.path.as_str(); if !(path.is_empty() && matches!(kind, kind::LIST | kind::PUT) || allowed(path)) || request.to.as_deref().is_some_and(|to| !allowed(to)) @@ -606,6 +609,7 @@ impl Served { self.changed(&[path.to_owned()]); done } + kind::EDITS => self.batch(peer, request)?, kind::CONFIRM => { self.storage.confirm(path, &stamp()?)?; done @@ -1474,6 +1478,7 @@ pub struct HostedRemote { path: String, /// The stamp last asked for, which an image the guest holds may already have. seen: Option, + rejected: bool, } impl HostedRemote { @@ -1482,12 +1487,81 @@ impl HostedRemote { guest: Arc::clone(guest), path: path.to_owned(), seen: None, + rejected: false, } } } impl crate::Remote for HostedRemote { + fn accepts_edits(&self) -> bool { + !self.rejected + && self + .guest + .host() + .is_some_and(|host| host.ops == Some(1) && host.kinds.contains(&kind::EDITS)) + } + + fn publish_edits( + &mut self, + transaction: &Transaction, + edits: &[crate::PendingEdit], + revisions: &BTreeMap, + ) -> std::result::Result<(), CommitError> { + let bytes = serde_json::to_vec(&batch::Edits { + edits: edits + .iter() + .map(|edit| (edit.author.clone(), edit.edit.clone())) + .collect(), + revisions: revisions.clone(), + }) + .map_err(|error| CommitError { + state: CommitState::NotCommitted, + error: io::Error::other(error), + })?; + let mut request = Request { + path: self.path.clone(), + stamp: Some(transaction.base().into()), + ..Request::default() + }; + self.guest + .carry(&mut request, bytes) + .map_err(Failed::commit)?; + let result = self.guest.request(kind::EDITS, request); + match result { + Ok(reply) => { + if let Some(stamp) = reply.stamp { + let stamp = Stamp::try_from(&stamp).map_err(|error| CommitError { + state: CommitState::Unknown, + error, + })?; + self.guest + .inner + .current + .lock() + .unwrap() + .insert(self.path.clone(), stamp); + } + self.seen = None; + Ok(()) + } + Err(Failed::Refused(failure)) + if wire::error_kind(failure.kind) == io::ErrorKind::Unsupported => + { + self.publish(transaction) + } + Err(error) => { + let error = error.commit(); + if error.state == CommitState::NotCommitted { + self.rejected = true; + self.guest.inner.current.lock().unwrap().remove(&self.path); + } + Err(error) + } + } + } + fn read(&mut self) -> io::Result> { + self.rejected = false; if let Some(image) = (self.seen.as_ref()).and_then(|seen| self.guest.held(&self.path, seen)) { return Ok(image); diff --git a/crates/notebook/src/live/share/batch.rs b/crates/notebook/src/live/share/batch.rs new file mode 100644 index 0000000000000000000000000000000000000000..92dd7e92c4eda78e7f85eae4bffc08263a99a0fa --- /dev/null +++ b/crates/notebook/src/live/share/batch.rs @@ -0,0 +1,258 @@ +use super::*; +use crate::{PendingEdit, merge}; +use onestore::{Arena, ExGuid, Section, op::Edit}; +use std::collections::BTreeSet; + +#[derive(serde::Serialize, serde::Deserialize)] +#[serde(deny_unknown_fields)] +pub(super) struct Edits { + pub edits: Vec<(String, Edit)>, + pub revisions: BTreeMap, +} + +pub(super) struct Waiting { + peer: [u8; 16], + request: Request, + reply: mpsc::Sender, +} + +fn retry(message: &str) -> Error { + Error::Remote(CommitError { + state: CommitState::NotCommitted, + error: io::Error::new(io::ErrorKind::ResourceBusy, message.to_owned()), + }) +} + +impl Served { + pub(super) fn batch(self: &Arc, peer: &[u8; 16], request: Request) -> Result { + let path = request.path.clone(); + let mut writers = self.writers.lock().unwrap(); + let writer = writers.entry(path.clone()).or_insert_with(|| { + let (send, receive) = mpsc::sync_channel::(64); + let served = Arc::downgrade(self); + thread::spawn(move || { + while let Ok(first) = receive.recv() { + let mut waiting = vec![first]; + if let Ok(next) = receive.recv_timeout(Duration::from_millis(5)) { + waiting.push(next); + } + waiting.extend(receive.try_iter().take(62)); + let Some(served) = served.upgrade() else { + return; + }; + served.write_batches(&path, waiting); + } + }); + send + }); + let (reply, receive) = mpsc::channel(); + writer + .try_send(Waiting { + peer: *peer, + request, + reply, + }) + .map_err(|_| retry("The section's writer is busy"))?; + drop(writers); + receive.recv_timeout(TIMEOUT).map_err(|_| { + Error::Remote(CommitError { + state: CommitState::Unknown, + error: io::Error::new( + io::ErrorKind::TimedOut, + "The section's writer did not answer", + ), + }) + }) + } + + fn write_batches(&self, path: &str, waiting: Vec) { + let before = match self.image(path) { + Ok(image) => image, + Err(error) => { + for batch in waiting { + let _ = batch.reply.send(failed(0, retry(&error.to_string()))); + } + return; + } + }; + let mut image = Arc::clone(&before); + let mut combined: Option = None; + let mut accepted = Vec::new(); + for batch in waiting { + match self.prepare_batch(path, &image, &batch) { + Ok(Some(transaction)) => { + let mut next = (*image).clone(); + let extended = + transaction + .apply(&mut next) + .and_then(|()| match &mut combined { + Some(combined) => combined.extend(transaction), + None => { + combined = Some(transaction); + Ok(()) + } + }); + match extended { + Ok(()) => { + image = Arc::new(next); + accepted.push(batch); + } + Err(error) => { + let _ = batch.reply.send(failed(0, error.into())); + } + } + } + Ok(None) => accepted.push(batch), + Err(error) => { + let _ = batch.reply.send(failed(0, error)); + } + } + } + if accepted.is_empty() { + return; + } + let result = match &combined { + Some(transaction) => self.storage.commit(path, transaction), + None => self + .storage + .confirm(path, &Stamp::of(&image).expect("section stamp")) + .map_err(Error::from), + }; + if result.is_ok() { + self.tell_delta(path, &before, &image, None); + self.keep(path, Arc::clone(&image)); + self.changed(&[path.to_owned()]); + } + for batch in accepted { + let reply = match &result { + Ok(()) => Reply { + stamp: Some((&Stamp::of(&image).expect("section stamp")).into()), + ..Reply::default() + }, + Err(Error::Remote(error)) => failed( + 0, + Error::Remote(CommitError { + state: error.state, + error: io::Error::new(error.error.kind(), error.error.to_string()), + }), + ), + Err(error) => failed(0, retry(&error.to_string())), + }; + let _ = batch.reply.send(reply); + } + } + + fn prepare_batch( + &self, + path: &str, + image: &[u8], + batch: &Waiting, + ) -> Result> { + let stamp: Stamp = batch + .request + .stamp + .as_ref() + .ok_or_else(|| retry("No batch base"))? + .try_into()?; + let edits: Edits = serde_json::from_slice(&self.carried(&batch.peer, &batch.request)?) + .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?; + if edits.revisions.is_empty() + || edits + .revisions + .values() + .any(|revision| revision.guid == [0; 16]) + { + return Err(retry("No batch revisions")); + } + let store = Store::parse(image)?; + let index = RevisionIndex::parse(&store)?; + if edits.revisions.iter().all(|(space, revision)| { + index + .spaces + .get(space) + .is_some_and(|space| space.revisions.contains_key(revision)) + }) { + return Ok(None); + } + let names = edits + .revisions + .iter() + .filter(|(space, revision)| { + !index + .spaces + .get(space) + .is_some_and(|space| space.revisions.contains_key(revision)) + }) + .map(|(space, revision)| (*space, *revision)) + .collect(); + let arena = Arena::default(); + let mut next = Section::open(&arena, image.to_vec())?; + let queued: Vec = edits + .edits + .into_iter() + .enumerate() + .map(|(id, (author, edit))| PendingEdit { + id: id as u64, + author, + edit, + }) + .collect(); + if Stamp::of(image)? == stamp { + for edit in &queued { + next.apply(&edit.author, &edit.edit)?; + } + } else { + let base = self + .images + .lock() + .unwrap() + .iter() + .find_map(|(held, at, image)| { + (held == path && *at == stamp).then(|| Arc::clone(image)) + }) + .ok_or_else(|| retry("The batch's base is no longer held"))?; + let old_arena = Arena::default(); + let mut old = Section::open(&old_arena, (*base).clone())?; + if old.root() != next.root() { + return Err(retry("The batch belongs to another section")); + } + let merged = merge::rebase(&mut old, &mut next, &queued, &BTreeSet::new())?; + if !merged.conflicts.is_empty() { + let local_arena = Arena::default(); + let mut local = Section::open(&local_arena, (*base).clone())?; + for edit in &queued { + local.apply(&edit.author, &edit.edit)?; + } + for (space, (author, objects)) in &merged.conflicts { + let page = local.page(*space)?; + if next.page(*space).is_ok_and(|remote| remote == page) { + return Err(retry("The batch's edits already landed")); + } + merge::conflict_page( + &mut next, + &mut local, + &merged.moved, + *space, + &page, + author, + objects, + crate::now(), + None, + )?; + } + } + } + let transaction = next.seal_as(&names)?; + if edits.revisions.iter().any(|(space, revision)| { + !next + .newest() + .any(|(held, at)| held == *space && at == *revision) + }) { + return Err(refused( + io::ErrorKind::Unsupported, + "The batch needs a guest's commit", + )); + } + Ok(transaction) + } +} diff --git a/crates/notebook/src/live/wire.rs b/crates/notebook/src/live/wire.rs index 9cf8d2c396b7a7be1354bf954020f5e24f40f4c2..15b1e98286b6665b12058e2a9c076b000d3fd34b 100644 --- a/crates/notebook/src/live/wire.rs +++ b/crates/notebook/src/live/wire.rs @@ -65,6 +65,7 @@ pub mod kind { pub const EXISTS: u16 = 271; /// A chunk of the bytes a later request carries. pub const PUT: u16 = 272; + pub const EDITS: u16 = 273; pub const REPLY: u16 = 511; } @@ -94,6 +95,7 @@ pub const KNOWN: &[u16] = &[ kind::READ_FILE, kind::EXISTS, kind::PUT, + kind::EDITS, kind::REPLY, ]; @@ -118,6 +120,8 @@ pub struct Hello { /// The share this peer hosts, whose storage requests it answers. #[cbor(n(5), with = "minicbor::bytes")] pub serves: Option<[u8; 16]>, + #[n(6)] + pub ops: Option, } impl Hello { @@ -132,6 +136,7 @@ impl Hello { app: format!("Snowbound {}", env!("CARGO_PKG_VERSION")), kinds: KNOWN.to_vec(), serves: None, + ops: Some(1), }) } } diff --git a/crates/notebook/src/merge.rs b/crates/notebook/src/merge.rs index 745fc682809edc89c7f40585d206e9f48633b063..cd5d59645950181c300e3d1ca1b238cca132a95b 100644 --- a/crates/notebook/src/merge.rs +++ b/crates/notebook/src/merge.rs @@ -120,7 +120,7 @@ pub(crate) fn rebase( } Op::Section(section_op) => { advance(&mut state.local, section_op); - let mut kept = state.form(new, section_op)?.map(Op::Section); + let mut kept = state.form(old, new, section_op)?.map(Op::Section); if let Some(form) = &kept && !apply(&mut state, new, form)? { @@ -425,11 +425,15 @@ impl Replay { Some(diff) => diff, None => { let diff = match (self.before.get(&space), self.after.get(&space)) { - (Some(before), Some(after)) if before != after => Some(Diff { - old: Index::of(&old.page(space)?), - new: Index::of(&new.page(space)?), - regions: BTreeMap::new(), - }), + (Some(before), Some(after)) => { + let (old, new) = (old.page(space)?, new.page(space)?); + // A host's merged batch retains its guest's revision identity. + (before != after || old != new).then(|| Diff { + old: Index::of(&old), + new: Index::of(&new), + regions: BTreeMap::new(), + }) + } _ => None, }; self.diffs.entry(space).or_insert(diff) @@ -469,7 +473,12 @@ impl Replay { /// `anchor`, the page edits of pages the remote still has and did not move, removals of /// pages the remote still has (a page it changed is not removed), a conflict page's /// content as a page of its own where the remote removed its page. - fn form(&self, new: &mut Section<'_>, op: &SectionOp) -> Result> { + fn form( + &self, + old: &mut Section<'_>, + new: &mut Section<'_>, + op: &SectionOp, + ) -> Result> { let conflicts: BTreeSet = new .conflicts()? .into_iter() @@ -539,6 +548,7 @@ impl Replay { .iter() .filter(|space| { self.listed(**space) && self.before.get(space) == self.after.get(space) + && matches!((old.page(**space), new.page(**space)), (Ok(a), Ok(b)) if a == b) || conflicts.contains(space) }) .copied() diff --git a/crates/notebook/src/sync.rs b/crates/notebook/src/sync.rs index 2f0f0a0628b751d31a9bc0597bb254b4b0ca16d3..2d73ab2faf045bcd6d3dd1b672f3c100c2d65f95 100644 --- a/crates/notebook/src/sync.rs +++ b/crates/notebook/src/sync.rs @@ -19,6 +19,18 @@ pub trait Remote { /// the last observed image's, synchronization neither reads nor revalidates the file. fn stamp(&mut self) -> io::Result; fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError>; + fn accepts_edits(&self) -> bool { + false + } + fn publish_edits( + &mut self, + transaction: &Transaction, + edits: &[PendingEdit], + revisions: &BTreeMap, + ) -> std::result::Result<(), CommitError> { + let _ = (edits, revisions); + self.publish(transaction) + } /// Confirms that the file still has `base`'s stamp and is durable (`onestore::confirm`). fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError>; /// The versions a file provider keeps beside the file, as iCloud Drive keeps the commits @@ -195,6 +207,10 @@ impl Replica { remote.retire(&version.id, keep).map_err(Error::RemoteIo)?; } let state = state(&*self.lock()?)?; + let batched = self.section.key.is_none() + && state.queued + && state.blocked.is_none() + && remote.accepts_edits(); let observed = remote.stamp().map_err(Error::RemoteIo)?; if let Some(blocked) = &state.blocked && observed == *state.remote.as_ref().unwrap_or(&state.base) @@ -211,7 +227,7 @@ impl Replica { }); } let image = match observed { - observed if observed == state.base => None, + observed if observed == state.base || batched => None, _ => { let image = remote.read().map_err(Error::RemoteIo)?; // Protected elsewhere: a section written anew, which only its key reads. @@ -348,7 +364,20 @@ impl Replica { changed, }); }; - match remote.publish(&transaction) { + let published = if self.section.key.is_none() && remote.accepts_edits() { + let connection = self.lock()?; + let edits = queue::load(&connection, None, Some(batch))?; + let revisions = decode_revisions(&connection.query_row( + "SELECT revisions FROM batches WHERE id=?1", + [batch], + |row| row.get::<_, String>(0), + )?)?; + drop(connection); + remote.publish_edits(&transaction, &edits, &revisions) + } else { + remote.publish(&transaction) + }; + match published { Ok(()) => {} Err(error) if error.state == CommitState::NotCommitted => { self.lock()? @@ -363,6 +392,11 @@ impl Replica { } self.acknowledge(batch, Some(&transaction), None)?; let revision = self.receipt(id)?; + if batched { + changed.extend(self.rebase(Some(remote.read().map_err(Error::RemoteIo)?))?); + changed.sort(); + changed.dedup(); + } Ok(Synced { edit: Some((id, EditStatus::Published { revision })), changed, diff --git a/crates/notebook/src/working.rs b/crates/notebook/src/working.rs index 521a9c44c388a32036dc715472de2836d26ba29d..df13e01b9cfbc3bff6ff2733b4be1e161ee8715d 100644 --- a/crates/notebook/src/working.rs +++ b/crates/notebook/src/working.rs @@ -823,6 +823,14 @@ fn rebase( ) .collect(); changed.extend(converged); + for (space, rid) in &before { + if after.get(space) == Some(rid) + && let (Ok(before), Ok(after)) = (old.page(*space), remote.page(*space)) + && before != after + { + changed.insert(*space); + } + } Ok(changed.into_iter().collect()) } diff --git a/crates/notebook/tests/live_batch.rs b/crates/notebook/tests/live_batch.rs new file mode 100644 index 0000000000000000000000000000000000000000..2945cfe0a6cab863c6beb05943bbd297c8315f23 --- /dev/null +++ b/crates/notebook/tests/live_batch.rs @@ -0,0 +1,179 @@ +#![cfg(feature = "live")] + +use notebook::{ + Replica, + live::share::{HostedRemote, Sharing}, +}; +use onestore::op::{Edit, Op, PageOp}; +use std::sync::{Arc, Barrier}; + +#[path = "support/live.rs"] +mod live; +use live::*; + +#[test] +fn a_lost_batch_receipt_is_confirmed_without_repeating_its_edit() { + use notebook::{PendingEdit, Remote}; + use onestore::{CommitError, CommitState, ExGuid, Stamp, Transaction}; + use std::{collections::BTreeMap, io}; + + struct Lost(HostedRemote, bool); + impl Remote for Lost { + fn read(&mut self) -> io::Result> { + self.0.read() + } + fn stamp(&mut self) -> io::Result { + self.0.stamp() + } + fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> { + self.0.publish(transaction) + } + fn confirm(&mut self, stamp: &Stamp) -> Result<(), CommitError> { + self.0.confirm(stamp) + } + fn accepts_edits(&self) -> bool { + true + } + fn publish_edits( + &mut self, + transaction: &Transaction, + edits: &[PendingEdit], + revisions: &BTreeMap, + ) -> Result<(), CommitError> { + self.0.publish_edits(transaction, edits, revisions)?; + if std::mem::take(&mut self.1) { + return Err(CommitError { + state: CommitState::Unknown, + error: io::ErrorKind::BrokenPipe.into(), + }); + } + Ok(()) + } + } + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let url = relay(Default::default()); + let sharing = Sharing::new("").unwrap(); + let host = host(&folder, &directory.path().join("host"), &sharing, &url); + let (guest, _) = guest("Alice", &code(&host), &url, &directory.path().join("alice")); + let file = folder.join("Garden.one"); + let image = std::fs::read(&file).unwrap(); + let (space, text, _) = server::text(&image); + let replica = Replica::open_or_create(directory.path().join("replica.sqlite"), None, || { + Ok(image.clone()) + }) + .unwrap(); + let id = replica + .apply( + "Alice", + Edit { + at: 134_000_000_000_000_000, + ops: vec![Op::Page { + space, + op: PageOp::Text { + text, + range: 13..13, + with: "X".into(), + }, + }], + }, + ) + .unwrap(); + let mut remote = Lost(HostedRemote::new(&guest, "Garden.one"), true); + assert!(matches!( + replica.sync_once(&mut remote), + Err(notebook::Error::Remote(CommitError { + state: CommitState::Unknown, + .. + })) + )); + replica.sync_once(&mut remote).unwrap(); + assert!(matches!( + replica.status(id).unwrap(), + Some(notebook::EditStatus::Published { .. }) + )); + assert_eq!( + server::text(&std::fs::read(file).unwrap()).2, + "Original textX" + ); + assert!(replica.pending().unwrap().is_empty()); +} + +#[test] +fn concurrent_batches_keep_receipts_and_rebase_later_keystrokes() { + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let url = relay(Default::default()); + let sharing = Sharing::new("").unwrap(); + let host = host(&folder, &directory.path().join("host"), &sharing, &url); + let code = code(&host); + let (alice, _) = guest("Alice", &code, &url, &directory.path().join("alice")); + let (bob, _) = guest("Bob", &code, &url, &directory.path().join("bob")); + let file = folder.join("Garden.one"); + let original = std::fs::read(&file).unwrap(); + let (space, text, _) = server::text(&original); + let edit = |with: &str, at| Edit { + at: 134_000_000_000_000_000, + ops: vec![Op::Page { + space, + op: PageOp::Text { + text, + range: at..at, + with: with.into(), + }, + }], + }; + let a = Arc::new( + Replica::open_or_create(directory.path().join("a.sqlite"), None, || { + Ok(original.clone()) + }) + .unwrap(), + ); + let b = Arc::new( + Replica::open_or_create(directory.path().join("b.sqlite"), None, || { + Ok(original.clone()) + }) + .unwrap(), + ); + let first = a.apply("Alice", edit("A", 13)).unwrap(); + let second = b.apply("Bob", edit("B", 0)).unwrap(); + let barrier = Arc::new(Barrier::new(2)); + std::thread::scope(|scope| { + for (replica, guest) in [(&a, &alice), (&b, &bob)] { + let barrier = Arc::clone(&barrier); + scope.spawn(move || { + barrier.wait(); + replica + .sync_once(&mut HostedRemote::new(guest, "Garden.one")) + .unwrap(); + }); + } + }); + assert!(matches!( + a.status(first).unwrap(), + Some(notebook::EditStatus::Published { .. }) + )); + assert!(matches!( + b.status(second).unwrap(), + Some(notebook::EditStatus::Published { .. }) + )); + let next = server::text(&a.snapshot().unwrap()).2.find('A').unwrap() as u32 + 1; + let third = a.apply("Alice", edit("a", next)).unwrap(); + let mut remote = HostedRemote::new(&alice, "Garden.one"); + for _ in 0..5 { + a.sync_once(&mut remote).unwrap(); + if matches!( + a.status(third).unwrap(), + Some(notebook::EditStatus::Published { .. }) + ) { + break; + } + } + let result = std::fs::read(&file).unwrap(); + let result_text = server::text(&result).2; + assert!(result_text.contains("Aa"), "{result_text}"); + assert!(result_text.contains('B'), "{result_text}"); + let arena = onestore::Arena::default(); + let mut section = onestore::Section::open(&arena, result).unwrap(); + assert!(section.conflicts().unwrap().is_empty()); +} diff --git a/crates/onestore/src/commit.rs b/crates/onestore/src/commit.rs index 213e6f9d42e073ac8a6262577ca2d4bcf4a7c0af..003324d821df2ab3894c062d1068361868a2c167 100644 --- a/crates/onestore/src/commit.rs +++ b/crates/onestore/src/commit.rs @@ -492,6 +492,38 @@ impl TryFrom for Transaction { } impl Transaction { + /// Combines a following transaction into one guarded publication, retaining every revision. + pub fn extend(&mut self, next: Transaction) -> Result<(), crate::Error> { + let length = self.base.length + self.append.len() as u64; + if next.base.header != self.header || next.base.length != length { + return Err(crate::Error { + offset: 0, + message: "Transactions are not consecutive", + }); + } + if next.patches.iter().any(|(offset, bytes)| { + *offset < 1024 || offset.saturating_add(bytes.len() as u64) > length + }) { + return Err(crate::Error { + offset: 0, + message: "A patch outside the base's data", + }); + } + for (offset, bytes) in next.patches { + let earlier = (self.base.length.saturating_sub(offset) as usize).min(bytes.len()); + if earlier != 0 { + self.patches.push((offset, bytes[..earlier].to_vec())); + } + if earlier < bytes.len() { + let at = (offset + earlier as u64 - self.base.length) as usize; + self.append[at..at + bytes.len() - earlier].copy_from_slice(&bytes[earlier..]); + } + } + self.append.extend_from_slice(&next.append); + self.header = next.header; + Ok(()) + } + /// The image this transaction applies to. pub fn base(&self) -> &Stamp { &self.base @@ -658,6 +690,33 @@ mod tests { } } + #[test] + fn consecutive_transactions_publish_the_same_bytes_together() { + let original = vec![1; 2048]; + let first = Transaction { + base: Stamp::of(&original).unwrap(), + append: vec![2; 32], + patches: vec![(1024, vec![3; 8])], + header: [4; 1024], + }; + let mut separate = original.clone(); + first.apply(&mut separate).unwrap(); + let second = Transaction { + base: Stamp::of(&separate).unwrap(), + append: vec![5; 16], + patches: vec![(1028, vec![6; 8]), (2044, vec![7; 12])], + header: [8; 1024], + }; + second.apply(&mut separate).unwrap(); + let mut combined = first.clone(); + assert!(combined.extend(first).is_err()); + combined.extend(second).unwrap(); + let combined = Transaction::from_bytes(&combined.to_bytes()).unwrap(); + let mut together = original; + combined.apply(&mut together).unwrap(); + assert_eq!(together, separate); + } + #[test] fn transactions_read_back_from_bytes() { let transaction = Transaction { -- 2.54.0