| author | |
| committer | |
| log | 9a82329b95a69c3d3e1844133a9d3045db885758 |
| tree | f53cd6bc7d1bd6a23a06177a3bdccdca6b1cec0f |
| parent | b82493ec83dcf2de59b1b32a287055d86bba7fec |
| signature | Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU |
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-sol9 files changed, 647 insertions(+), 15 deletions(-)
crates/notebook/examples/live_crowd.rs+5| ... | ... | @@ -141,6 +141,11 @@ fn host(url: &str, paragraphs: usize) { |
| 141 | 141 | println!("code {}", host.code().unwrap()); |
| 142 | 142 | let _ = std::io::stdin().read_to_end(&mut Vec::new()); |
| 143 | 143 | let image = std::fs::read(&file).unwrap(); |
| 144 | if let Some(output) = std::env::var_os("SNOWBOUND_LIVE_EVIDENCE") { | |
| 145 | let output = Path::new(&output); | |
| 146 | std::fs::create_dir_all(output).unwrap(); | |
| 147 | std::fs::copy(&file, output.join("Garden.one")).unwrap(); | |
| 148 | } | |
| 144 | 149 | let size = image.len(); |
| 145 | 150 | let arena = onestore::Arena::default(); |
| 146 | 151 | let mut section = onestore::Section::open(&arena, image).unwrap(); |
crates/notebook/src/live/share.rs+80-6| ... | ... | @@ -1,8 +1,7 @@ |
| 1 | 1 | //! Live Share: a notebook one Snowbound holds, opened on others through a short code. The |
| 2 | //! host serves its notebook's storage verbs (`session::Storage`) to each peer in the share's | |
| 3 | //! room; a guest runs the replica, queue and merge it runs on an SMB share against those | |
| 4 | //! verbs, so offline queueing, rebases and conflict pages work as there, and the host's files | |
| 5 | //! stay what its own storage writes. A guest first meets the host in the code's room, where | |
| 2 | //! host serves its notebook's storage verbs (`session::Storage`) and batches guests' ops into | |
| 3 | //! guarded publications. A guest runs the same replica and durable queue as on an SMB share; | |
| 4 | //! protected sections and older peers use transactions. A guest first meets the host in the code's room, where | |
| 6 | 5 | //! the host welcomes it with the share's room and secret; a new share has a new secret, so |
| 7 | 6 | //! stopping retires every guest. Large bodies travel a chunk at a time, each answered before |
| 8 | 7 | //! the next, so a relay never holds much for a slow peer. |
| ... | ... | @@ -29,6 +28,8 @@ use std::{ |
| 29 | 28 | time::{Duration, Instant}, |
| 30 | 29 | }; |
| 31 | 30 | |
| 31 | mod batch; | |
| 32 | ||
| 32 | 33 | /// The most bytes one message of a read or an upload carries. |
| 33 | 34 | const CHUNK: usize = 128 << 10; |
| 34 | 35 | /// The most bytes of chunks a guest has asked for and not yet been given. |
| ... | ... | @@ -234,6 +235,7 @@ impl Host { |
| 234 | 235 | snapshots: Mutex::default(), |
| 235 | 236 | puts: Mutex::default(), |
| 236 | 237 | guests: Mutex::default(), |
| 238 | writers: Mutex::default(), | |
| 237 | 239 | host: Mutex::default(), |
| 238 | 240 | room: OnceLock::new(), |
| 239 | 241 | }); |
| ... | ... | @@ -393,6 +395,7 @@ struct Served { |
| 393 | 395 | /// Bytes a later request carries, by guest and upload. |
| 394 | 396 | puts: Mutex<ByGuest<Vec<u8>>>, |
| 395 | 397 | guests: Mutex<BTreeMap<[u8; 16], Admitted>>, |
| 398 | writers: Mutex<HashMap<String, mpsc::SyncSender<batch::Waiting>>>, | |
| 396 | 399 | /// Hears the paths guests changed, as the host's own notebook should. |
| 397 | 400 | host: Mutex<Option<crate::session::Listener>>, |
| 398 | 401 | /// The share's room, to tell guests what changed. |
| ... | ... | @@ -524,7 +527,7 @@ impl Served { |
| 524 | 527 | guest.held.truncate(IMAGES); |
| 525 | 528 | } |
| 526 | 529 | |
| 527 | fn handle(&self, peer: &[u8; 16], kind: u16, body: &[u8]) -> Reply { | |
| 530 | fn handle(self: &Arc<Self>, peer: &[u8; 16], kind: u16, body: &[u8]) -> Reply { | |
| 528 | 531 | let request = match minicbor::decode::<Request>(body) { |
| 529 | 532 | Ok(request) => request, |
| 530 | 533 | Err(_) => { |
| ... | ... | @@ -541,7 +544,7 @@ impl Served { |
| 541 | 544 | } |
| 542 | 545 | } |
| 543 | 546 | |
| 544 | fn answer(&self, peer: &[u8; 16], kind: u16, request: Request) -> Result<Reply> { | |
| 547 | fn answer(self: &Arc<Self>, peer: &[u8; 16], kind: u16, request: Request) -> Result<Reply> { | |
| 545 | 548 | let path = request.path.as_str(); |
| 546 | 549 | if !(path.is_empty() && matches!(kind, kind::LIST | kind::PUT) || allowed(path)) |
| 547 | 550 | || request.to.as_deref().is_some_and(|to| !allowed(to)) |
| ... | ... | @@ -606,6 +609,7 @@ impl Served { |
| 606 | 609 | self.changed(&[path.to_owned()]); |
| 607 | 610 | done |
| 608 | 611 | } |
| 612 | kind::EDITS => self.batch(peer, request)?, | |
| 609 | 613 | kind::CONFIRM => { |
| 610 | 614 | self.storage.confirm(path, &stamp()?)?; |
| 611 | 615 | done |
| ... | ... | @@ -1474,6 +1478,7 @@ pub struct HostedRemote { |
| 1474 | 1478 | path: String, |
| 1475 | 1479 | /// The stamp last asked for, which an image the guest holds may already have. |
| 1476 | 1480 | seen: Option<Stamp>, |
| 1481 | rejected: bool, | |
| 1477 | 1482 | } |
| 1478 | 1483 | |
| 1479 | 1484 | impl HostedRemote { |
| ... | ... | @@ -1482,12 +1487,81 @@ impl HostedRemote { |
| 1482 | 1487 | guest: Arc::clone(guest), |
| 1483 | 1488 | path: path.to_owned(), |
| 1484 | 1489 | seen: None, |
| 1490 | rejected: false, | |
| 1485 | 1491 | } |
| 1486 | 1492 | } |
| 1487 | 1493 | } |
| 1488 | 1494 | |
| 1489 | 1495 | impl crate::Remote for HostedRemote { |
| 1496 | fn accepts_edits(&self) -> bool { | |
| 1497 | !self.rejected | |
| 1498 | && self | |
| 1499 | .guest | |
| 1500 | .host() | |
| 1501 | .is_some_and(|host| host.ops == Some(1) && host.kinds.contains(&kind::EDITS)) | |
| 1502 | } | |
| 1503 | ||
| 1504 | fn publish_edits( | |
| 1505 | &mut self, | |
| 1506 | transaction: &Transaction, | |
| 1507 | edits: &[crate::PendingEdit], | |
| 1508 | revisions: &BTreeMap<onestore::ExGuid, onestore::ExGuid>, | |
| 1509 | ) -> std::result::Result<(), CommitError> { | |
| 1510 | let bytes = serde_json::to_vec(&batch::Edits { | |
| 1511 | edits: edits | |
| 1512 | .iter() | |
| 1513 | .map(|edit| (edit.author.clone(), edit.edit.clone())) | |
| 1514 | .collect(), | |
| 1515 | revisions: revisions.clone(), | |
| 1516 | }) | |
| 1517 | .map_err(|error| CommitError { | |
| 1518 | state: CommitState::NotCommitted, | |
| 1519 | error: io::Error::other(error), | |
| 1520 | })?; | |
| 1521 | let mut request = Request { | |
| 1522 | path: self.path.clone(), | |
| 1523 | stamp: Some(transaction.base().into()), | |
| 1524 | ..Request::default() | |
| 1525 | }; | |
| 1526 | self.guest | |
| 1527 | .carry(&mut request, bytes) | |
| 1528 | .map_err(Failed::commit)?; | |
| 1529 | let result = self.guest.request(kind::EDITS, request); | |
| 1530 | match result { | |
| 1531 | Ok(reply) => { | |
| 1532 | if let Some(stamp) = reply.stamp { | |
| 1533 | let stamp = Stamp::try_from(&stamp).map_err(|error| CommitError { | |
| 1534 | state: CommitState::Unknown, | |
| 1535 | error, | |
| 1536 | })?; | |
| 1537 | self.guest | |
| 1538 | .inner | |
| 1539 | .current | |
| 1540 | .lock() | |
| 1541 | .unwrap() | |
| 1542 | .insert(self.path.clone(), stamp); | |
| 1543 | } | |
| 1544 | self.seen = None; | |
| 1545 | Ok(()) | |
| 1546 | } | |
| 1547 | Err(Failed::Refused(failure)) | |
| 1548 | if wire::error_kind(failure.kind) == io::ErrorKind::Unsupported => | |
| 1549 | { | |
| 1550 | self.publish(transaction) | |
| 1551 | } | |
| 1552 | Err(error) => { | |
| 1553 | let error = error.commit(); | |
| 1554 | if error.state == CommitState::NotCommitted { | |
| 1555 | self.rejected = true; | |
| 1556 | self.guest.inner.current.lock().unwrap().remove(&self.path); | |
| 1557 | } | |
| 1558 | Err(error) | |
| 1559 | } | |
| 1560 | } | |
| 1561 | } | |
| 1562 | ||
| 1490 | 1563 | fn read(&mut self) -> io::Result<Vec<u8>> { |
| 1564 | self.rejected = false; | |
| 1491 | 1565 | if let Some(image) = (self.seen.as_ref()).and_then(|seen| self.guest.held(&self.path, seen)) |
| 1492 | 1566 | { |
| 1493 | 1567 | return Ok(image); |
crates/notebook/src/live/share/batch.rs created+258| ... | ... | @@ -0,0 +1,258 @@ |
| 1 | use super::*; | |
| 2 | use crate::{PendingEdit, merge}; | |
| 3 | use onestore::{Arena, ExGuid, Section, op::Edit}; | |
| 4 | use std::collections::BTreeSet; | |
| 5 | ||
| 6 | #[derive(serde::Serialize, serde::Deserialize)] | |
| 7 | #[serde(deny_unknown_fields)] | |
| 8 | pub(super) struct Edits { | |
| 9 | pub edits: Vec<(String, Edit)>, | |
| 10 | pub revisions: BTreeMap<ExGuid, ExGuid>, | |
| 11 | } | |
| 12 | ||
| 13 | pub(super) struct Waiting { | |
| 14 | peer: [u8; 16], | |
| 15 | request: Request, | |
| 16 | reply: mpsc::Sender<Reply>, | |
| 17 | } | |
| 18 | ||
| 19 | fn retry(message: &str) -> Error { | |
| 20 | Error::Remote(CommitError { | |
| 21 | state: CommitState::NotCommitted, | |
| 22 | error: io::Error::new(io::ErrorKind::ResourceBusy, message.to_owned()), | |
| 23 | }) | |
| 24 | } | |
| 25 | ||
| 26 | impl Served { | |
| 27 | pub(super) fn batch(self: &Arc<Self>, peer: &[u8; 16], request: Request) -> Result<Reply> { | |
| 28 | let path = request.path.clone(); | |
| 29 | let mut writers = self.writers.lock().unwrap(); | |
| 30 | let writer = writers.entry(path.clone()).or_insert_with(|| { | |
| 31 | let (send, receive) = mpsc::sync_channel::<Waiting>(64); | |
| 32 | let served = Arc::downgrade(self); | |
| 33 | thread::spawn(move || { | |
| 34 | while let Ok(first) = receive.recv() { | |
| 35 | let mut waiting = vec![first]; | |
| 36 | if let Ok(next) = receive.recv_timeout(Duration::from_millis(5)) { | |
| 37 | waiting.push(next); | |
| 38 | } | |
| 39 | waiting.extend(receive.try_iter().take(62)); | |
| 40 | let Some(served) = served.upgrade() else { | |
| 41 | return; | |
| 42 | }; | |
| 43 | served.write_batches(&path, waiting); | |
| 44 | } | |
| 45 | }); | |
| 46 | send | |
| 47 | }); | |
| 48 | let (reply, receive) = mpsc::channel(); | |
| 49 | writer | |
| 50 | .try_send(Waiting { | |
| 51 | peer: *peer, | |
| 52 | request, | |
| 53 | reply, | |
| 54 | }) | |
| 55 | .map_err(|_| retry("The section's writer is busy"))?; | |
| 56 | drop(writers); | |
| 57 | receive.recv_timeout(TIMEOUT).map_err(|_| { | |
| 58 | Error::Remote(CommitError { | |
| 59 | state: CommitState::Unknown, | |
| 60 | error: io::Error::new( | |
| 61 | io::ErrorKind::TimedOut, | |
| 62 | "The section's writer did not answer", | |
| 63 | ), | |
| 64 | }) | |
| 65 | }) | |
| 66 | } | |
| 67 | ||
| 68 | fn write_batches(&self, path: &str, waiting: Vec<Waiting>) { | |
| 69 | let before = match self.image(path) { | |
| 70 | Ok(image) => image, | |
| 71 | Err(error) => { | |
| 72 | for batch in waiting { | |
| 73 | let _ = batch.reply.send(failed(0, retry(&error.to_string()))); | |
| 74 | } | |
| 75 | return; | |
| 76 | } | |
| 77 | }; | |
| 78 | let mut image = Arc::clone(&before); | |
| 79 | let mut combined: Option<Transaction> = None; | |
| 80 | let mut accepted = Vec::new(); | |
| 81 | for batch in waiting { | |
| 82 | match self.prepare_batch(path, &image, &batch) { | |
| 83 | Ok(Some(transaction)) => { | |
| 84 | let mut next = (*image).clone(); | |
| 85 | let extended = | |
| 86 | transaction | |
| 87 | .apply(&mut next) | |
| 88 | .and_then(|()| match &mut combined { | |
| 89 | Some(combined) => combined.extend(transaction), | |
| 90 | None => { | |
| 91 | combined = Some(transaction); | |
| 92 | Ok(()) | |
| 93 | } | |
| 94 | }); | |
| 95 | match extended { | |
| 96 | Ok(()) => { | |
| 97 | image = Arc::new(next); | |
| 98 | accepted.push(batch); | |
| 99 | } | |
| 100 | Err(error) => { | |
| 101 | let _ = batch.reply.send(failed(0, error.into())); | |
| 102 | } | |
| 103 | } | |
| 104 | } | |
| 105 | Ok(None) => accepted.push(batch), | |
| 106 | Err(error) => { | |
| 107 | let _ = batch.reply.send(failed(0, error)); | |
| 108 | } | |
| 109 | } | |
| 110 | } | |
| 111 | if accepted.is_empty() { | |
| 112 | return; | |
| 113 | } | |
| 114 | let result = match &combined { | |
| 115 | Some(transaction) => self.storage.commit(path, transaction), | |
| 116 | None => self | |
| 117 | .storage | |
| 118 | .confirm(path, &Stamp::of(&image).expect("section stamp")) | |
| 119 | .map_err(Error::from), | |
| 120 | }; | |
| 121 | if result.is_ok() { | |
| 122 | self.tell_delta(path, &before, &image, None); | |
| 123 | self.keep(path, Arc::clone(&image)); | |
| 124 | self.changed(&[path.to_owned()]); | |
| 125 | } | |
| 126 | for batch in accepted { | |
| 127 | let reply = match &result { | |
| 128 | Ok(()) => Reply { | |
| 129 | stamp: Some((&Stamp::of(&image).expect("section stamp")).into()), | |
| 130 | ..Reply::default() | |
| 131 | }, | |
| 132 | Err(Error::Remote(error)) => failed( | |
| 133 | 0, | |
| 134 | Error::Remote(CommitError { | |
| 135 | state: error.state, | |
| 136 | error: io::Error::new(error.error.kind(), error.error.to_string()), | |
| 137 | }), | |
| 138 | ), | |
| 139 | Err(error) => failed(0, retry(&error.to_string())), | |
| 140 | }; | |
| 141 | let _ = batch.reply.send(reply); | |
| 142 | } | |
| 143 | } | |
| 144 | ||
| 145 | fn prepare_batch( | |
| 146 | &self, | |
| 147 | path: &str, | |
| 148 | image: &[u8], | |
| 149 | batch: &Waiting, | |
| 150 | ) -> Result<Option<Transaction>> { | |
| 151 | let stamp: Stamp = batch | |
| 152 | .request | |
| 153 | .stamp | |
| 154 | .as_ref() | |
| 155 | .ok_or_else(|| retry("No batch base"))? | |
| 156 | .try_into()?; | |
| 157 | let edits: Edits = serde_json::from_slice(&self.carried(&batch.peer, &batch.request)?) | |
| 158 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?; | |
| 159 | if edits.revisions.is_empty() | |
| 160 | || edits | |
| 161 | .revisions | |
| 162 | .values() | |
| 163 | .any(|revision| revision.guid == [0; 16]) | |
| 164 | { | |
| 165 | return Err(retry("No batch revisions")); | |
| 166 | } | |
| 167 | let store = Store::parse(image)?; | |
| 168 | let index = RevisionIndex::parse(&store)?; | |
| 169 | if edits.revisions.iter().all(|(space, revision)| { | |
| 170 | index | |
| 171 | .spaces | |
| 172 | .get(space) | |
| 173 | .is_some_and(|space| space.revisions.contains_key(revision)) | |
| 174 | }) { | |
| 175 | return Ok(None); | |
| 176 | } | |
| 177 | let names = edits | |
| 178 | .revisions | |
| 179 | .iter() | |
| 180 | .filter(|(space, revision)| { | |
| 181 | !index | |
| 182 | .spaces | |
| 183 | .get(space) | |
| 184 | .is_some_and(|space| space.revisions.contains_key(revision)) | |
| 185 | }) | |
| 186 | .map(|(space, revision)| (*space, *revision)) | |
| 187 | .collect(); | |
| 188 | let arena = Arena::default(); | |
| 189 | let mut next = Section::open(&arena, image.to_vec())?; | |
| 190 | let queued: Vec<PendingEdit> = edits | |
| 191 | .edits | |
| 192 | .into_iter() | |
| 193 | .enumerate() | |
| 194 | .map(|(id, (author, edit))| PendingEdit { | |
| 195 | id: id as u64, | |
| 196 | author, | |
| 197 | edit, | |
| 198 | }) | |
| 199 | .collect(); | |
| 200 | if Stamp::of(image)? == stamp { | |
| 201 | for edit in &queued { | |
| 202 | next.apply(&edit.author, &edit.edit)?; | |
| 203 | } | |
| 204 | } else { | |
| 205 | let base = self | |
| 206 | .images | |
| 207 | .lock() | |
| 208 | .unwrap() | |
| 209 | .iter() | |
| 210 | .find_map(|(held, at, image)| { | |
| 211 | (held == path && *at == stamp).then(|| Arc::clone(image)) | |
| 212 | }) | |
| 213 | .ok_or_else(|| retry("The batch's base is no longer held"))?; | |
| 214 | let old_arena = Arena::default(); | |
| 215 | let mut old = Section::open(&old_arena, (*base).clone())?; | |
| 216 | if old.root() != next.root() { | |
| 217 | return Err(retry("The batch belongs to another section")); | |
| 218 | } | |
| 219 | let merged = merge::rebase(&mut old, &mut next, &queued, &BTreeSet::new())?; | |
| 220 | if !merged.conflicts.is_empty() { | |
| 221 | let local_arena = Arena::default(); | |
| 222 | let mut local = Section::open(&local_arena, (*base).clone())?; | |
| 223 | for edit in &queued { | |
| 224 | local.apply(&edit.author, &edit.edit)?; | |
| 225 | } | |
| 226 | for (space, (author, objects)) in &merged.conflicts { | |
| 227 | let page = local.page(*space)?; | |
| 228 | if next.page(*space).is_ok_and(|remote| remote == page) { | |
| 229 | return Err(retry("The batch's edits already landed")); | |
| 230 | } | |
| 231 | merge::conflict_page( | |
| 232 | &mut next, | |
| 233 | &mut local, | |
| 234 | &merged.moved, | |
| 235 | *space, | |
| 236 | &page, | |
| 237 | author, | |
| 238 | objects, | |
| 239 | crate::now(), | |
| 240 | None, | |
| 241 | )?; | |
| 242 | } | |
| 243 | } | |
| 244 | } | |
| 245 | let transaction = next.seal_as(&names)?; | |
| 246 | if edits.revisions.iter().any(|(space, revision)| { | |
| 247 | !next | |
| 248 | .newest() | |
| 249 | .any(|(held, at)| held == *space && at == *revision) | |
| 250 | }) { | |
| 251 | return Err(refused( | |
| 252 | io::ErrorKind::Unsupported, | |
| 253 | "The batch needs a guest's commit", | |
| 254 | )); | |
| 255 | } | |
| 256 | Ok(transaction) | |
| 257 | } | |
| 258 | } |
crates/notebook/src/live/wire.rs+5| ... | ... | @@ -65,6 +65,7 @@ pub mod kind { |
| 65 | 65 | pub const EXISTS: u16 = 271; |
| 66 | 66 | /// A chunk of the bytes a later request carries. |
| 67 | 67 | pub const PUT: u16 = 272; |
| 68 | pub const EDITS: u16 = 273; | |
| 68 | 69 | pub const REPLY: u16 = 511; |
| 69 | 70 | } |
| 70 | 71 | |
| ... | ... | @@ -94,6 +95,7 @@ pub const KNOWN: &[u16] = &[ |
| 94 | 95 | kind::READ_FILE, |
| 95 | 96 | kind::EXISTS, |
| 96 | 97 | kind::PUT, |
| 98 | kind::EDITS, | |
| 97 | 99 | kind::REPLY, |
| 98 | 100 | ]; |
| 99 | 101 | |
| ... | ... | @@ -118,6 +120,8 @@ pub struct Hello { |
| 118 | 120 | /// The share this peer hosts, whose storage requests it answers. |
| 119 | 121 | #[cbor(n(5), with = "minicbor::bytes")] |
| 120 | 122 | pub serves: Option<[u8; 16]>, |
| 123 | #[n(6)] | |
| 124 | pub ops: Option<u16>, | |
| 121 | 125 | } |
| 122 | 126 | |
| 123 | 127 | impl Hello { |
| ... | ... | @@ -132,6 +136,7 @@ impl Hello { |
| 132 | 136 | app: format!("Snowbound {}", env!("CARGO_PKG_VERSION")), |
| 133 | 137 | kinds: KNOWN.to_vec(), |
| 134 | 138 | serves: None, |
| 139 | ops: Some(1), | |
| 135 | 140 | }) |
| 136 | 141 | } |
| 137 | 142 | } |
crates/notebook/src/merge.rs+17-7| ... | ... | @@ -120,7 +120,7 @@ pub(crate) fn rebase( |
| 120 | 120 | } |
| 121 | 121 | Op::Section(section_op) => { |
| 122 | 122 | advance(&mut state.local, section_op); |
| 123 | let mut kept = state.form(new, section_op)?.map(Op::Section); | |
| 123 | let mut kept = state.form(old, new, section_op)?.map(Op::Section); | |
| 124 | 124 | if let Some(form) = &kept |
| 125 | 125 | && !apply(&mut state, new, form)? |
| 126 | 126 | { |
| ... | ... | @@ -425,11 +425,15 @@ impl Replay { |
| 425 | 425 | Some(diff) => diff, |
| 426 | 426 | None => { |
| 427 | 427 | let diff = match (self.before.get(&space), self.after.get(&space)) { |
| 428 | (Some(before), Some(after)) if before != after => Some(Diff { | |
| 429 | old: Index::of(&old.page(space)?), | |
| 430 | new: Index::of(&new.page(space)?), | |
| 431 | regions: BTreeMap::new(), | |
| 432 | }), | |
| 428 | (Some(before), Some(after)) => { | |
| 429 | let (old, new) = (old.page(space)?, new.page(space)?); | |
| 430 | // A host's merged batch retains its guest's revision identity. | |
| 431 | (before != after || old != new).then(|| Diff { | |
| 432 | old: Index::of(&old), | |
| 433 | new: Index::of(&new), | |
| 434 | regions: BTreeMap::new(), | |
| 435 | }) | |
| 436 | } | |
| 433 | 437 | _ => None, |
| 434 | 438 | }; |
| 435 | 439 | self.diffs.entry(space).or_insert(diff) |
| ... | ... | @@ -469,7 +473,12 @@ impl Replay { |
| 469 | 473 | /// `anchor`, the page edits of pages the remote still has and did not move, removals of |
| 470 | 474 | /// pages the remote still has (a page it changed is not removed), a conflict page's |
| 471 | 475 | /// content as a page of its own where the remote removed its page. |
| 472 | fn form(&self, new: &mut Section<'_>, op: &SectionOp) -> Result<Option<SectionOp>> { | |
| 476 | fn form( | |
| 477 | &self, | |
| 478 | old: &mut Section<'_>, | |
| 479 | new: &mut Section<'_>, | |
| 480 | op: &SectionOp, | |
| 481 | ) -> Result<Option<SectionOp>> { | |
| 473 | 482 | let conflicts: BTreeSet<ExGuid> = new |
| 474 | 483 | .conflicts()? |
| 475 | 484 | .into_iter() |
| ... | ... | @@ -539,6 +548,7 @@ impl Replay { |
| 539 | 548 | .iter() |
| 540 | 549 | .filter(|space| { |
| 541 | 550 | self.listed(**space) && self.before.get(space) == self.after.get(space) |
| 551 | && matches!((old.page(**space), new.page(**space)), (Ok(a), Ok(b)) if a == b) | |
| 542 | 552 | || conflicts.contains(space) |
| 543 | 553 | }) |
| 544 | 554 | .copied() |
crates/notebook/src/sync.rs+36-2| ... | ... | @@ -19,6 +19,18 @@ pub trait Remote { |
| 19 | 19 | /// the last observed image's, synchronization neither reads nor revalidates the file. |
| 20 | 20 | fn stamp(&mut self) -> io::Result<Stamp>; |
| 21 | 21 | fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError>; |
| 22 | fn accepts_edits(&self) -> bool { | |
| 23 | false | |
| 24 | } | |
| 25 | fn publish_edits( | |
| 26 | &mut self, | |
| 27 | transaction: &Transaction, | |
| 28 | edits: &[PendingEdit], | |
| 29 | revisions: &BTreeMap<ExGuid, ExGuid>, | |
| 30 | ) -> std::result::Result<(), CommitError> { | |
| 31 | let _ = (edits, revisions); | |
| 32 | self.publish(transaction) | |
| 33 | } | |
| 22 | 34 | /// Confirms that the file still has `base`'s stamp and is durable (`onestore::confirm`). |
| 23 | 35 | fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError>; |
| 24 | 36 | /// The versions a file provider keeps beside the file, as iCloud Drive keeps the commits |
| ... | ... | @@ -195,6 +207,10 @@ impl Replica { |
| 195 | 207 | remote.retire(&version.id, keep).map_err(Error::RemoteIo)?; |
| 196 | 208 | } |
| 197 | 209 | let state = state(&*self.lock()?)?; |
| 210 | let batched = self.section.key.is_none() | |
| 211 | && state.queued | |
| 212 | && state.blocked.is_none() | |
| 213 | && remote.accepts_edits(); | |
| 198 | 214 | let observed = remote.stamp().map_err(Error::RemoteIo)?; |
| 199 | 215 | if let Some(blocked) = &state.blocked |
| 200 | 216 | && observed == *state.remote.as_ref().unwrap_or(&state.base) |
| ... | ... | @@ -211,7 +227,7 @@ impl Replica { |
| 211 | 227 | }); |
| 212 | 228 | } |
| 213 | 229 | let image = match observed { |
| 214 | observed if observed == state.base => None, | |
| 230 | observed if observed == state.base || batched => None, | |
| 215 | 231 | _ => { |
| 216 | 232 | let image = remote.read().map_err(Error::RemoteIo)?; |
| 217 | 233 | // Protected elsewhere: a section written anew, which only its key reads. |
| ... | ... | @@ -348,7 +364,20 @@ impl Replica { |
| 348 | 364 | changed, |
| 349 | 365 | }); |
| 350 | 366 | }; |
| 351 | match remote.publish(&transaction) { | |
| 367 | let published = if self.section.key.is_none() && remote.accepts_edits() { | |
| 368 | let connection = self.lock()?; | |
| 369 | let edits = queue::load(&connection, None, Some(batch))?; | |
| 370 | let revisions = decode_revisions(&connection.query_row( | |
| 371 | "SELECT revisions FROM batches WHERE id=?1", | |
| 372 | [batch], | |
| 373 | |row| row.get::<_, String>(0), | |
| 374 | )?)?; | |
| 375 | drop(connection); | |
| 376 | remote.publish_edits(&transaction, &edits, &revisions) | |
| 377 | } else { | |
| 378 | remote.publish(&transaction) | |
| 379 | }; | |
| 380 | match published { | |
| 352 | 381 | Ok(()) => {} |
| 353 | 382 | Err(error) if error.state == CommitState::NotCommitted => { |
| 354 | 383 | self.lock()? |
| ... | ... | @@ -363,6 +392,11 @@ impl Replica { |
| 363 | 392 | } |
| 364 | 393 | self.acknowledge(batch, Some(&transaction), None)?; |
| 365 | 394 | let revision = self.receipt(id)?; |
| 395 | if batched { | |
| 396 | changed.extend(self.rebase(Some(remote.read().map_err(Error::RemoteIo)?))?); | |
| 397 | changed.sort(); | |
| 398 | changed.dedup(); | |
| 399 | } | |
| 366 | 400 | Ok(Synced { |
| 367 | 401 | edit: Some((id, EditStatus::Published { revision })), |
| 368 | 402 | changed, |
crates/notebook/src/working.rs+8| ... | ... | @@ -823,6 +823,14 @@ fn rebase( |
| 823 | 823 | ) |
| 824 | 824 | .collect(); |
| 825 | 825 | changed.extend(converged); |
| 826 | for (space, rid) in &before { | |
| 827 | if after.get(space) == Some(rid) | |
| 828 | && let (Ok(before), Ok(after)) = (old.page(*space), remote.page(*space)) | |
| 829 | && before != after | |
| 830 | { | |
| 831 | changed.insert(*space); | |
| 832 | } | |
| 833 | } | |
| 826 | 834 | Ok(changed.into_iter().collect()) |
| 827 | 835 | } |
| 828 | 836 |
crates/notebook/tests/live_batch.rs created+179| ... | ... | @@ -0,0 +1,179 @@ |
| 1 | #![cfg(feature = "live")] | |
| 2 | ||
| 3 | use notebook::{ | |
| 4 | Replica, | |
| 5 | live::share::{HostedRemote, Sharing}, | |
| 6 | }; | |
| 7 | use onestore::op::{Edit, Op, PageOp}; | |
| 8 | use std::sync::{Arc, Barrier}; | |
| 9 | ||
| 10 | #[path = "support/live.rs"] | |
| 11 | mod live; | |
| 12 | use live::*; | |
| 13 | ||
| 14 | #[test] | |
| 15 | fn a_lost_batch_receipt_is_confirmed_without_repeating_its_edit() { | |
| 16 | use notebook::{PendingEdit, Remote}; | |
| 17 | use onestore::{CommitError, CommitState, ExGuid, Stamp, Transaction}; | |
| 18 | use std::{collections::BTreeMap, io}; | |
| 19 | ||
| 20 | struct Lost(HostedRemote, bool); | |
| 21 | impl Remote for Lost { | |
| 22 | fn read(&mut self) -> io::Result<Vec<u8>> { | |
| 23 | self.0.read() | |
| 24 | } | |
| 25 | fn stamp(&mut self) -> io::Result<Stamp> { | |
| 26 | self.0.stamp() | |
| 27 | } | |
| 28 | fn publish(&mut self, transaction: &Transaction) -> Result<(), CommitError> { | |
| 29 | self.0.publish(transaction) | |
| 30 | } | |
| 31 | fn confirm(&mut self, stamp: &Stamp) -> Result<(), CommitError> { | |
| 32 | self.0.confirm(stamp) | |
| 33 | } | |
| 34 | fn accepts_edits(&self) -> bool { | |
| 35 | true | |
| 36 | } | |
| 37 | fn publish_edits( | |
| 38 | &mut self, | |
| 39 | transaction: &Transaction, | |
| 40 | edits: &[PendingEdit], | |
| 41 | revisions: &BTreeMap<ExGuid, ExGuid>, | |
| 42 | ) -> Result<(), CommitError> { | |
| 43 | self.0.publish_edits(transaction, edits, revisions)?; | |
| 44 | if std::mem::take(&mut self.1) { | |
| 45 | return Err(CommitError { | |
| 46 | state: CommitState::Unknown, | |
| 47 | error: io::ErrorKind::BrokenPipe.into(), | |
| 48 | }); | |
| 49 | } | |
| 50 | Ok(()) | |
| 51 | } | |
| 52 | } | |
| 53 | let directory = tempfile::tempdir().unwrap(); | |
| 54 | let folder = notebook(directory.path()); | |
| 55 | let url = relay(Default::default()); | |
| 56 | let sharing = Sharing::new("").unwrap(); | |
| 57 | let host = host(&folder, &directory.path().join("host"), &sharing, &url); | |
| 58 | let (guest, _) = guest("Alice", &code(&host), &url, &directory.path().join("alice")); | |
| 59 | let file = folder.join("Garden.one"); | |
| 60 | let image = std::fs::read(&file).unwrap(); | |
| 61 | let (space, text, _) = server::text(&image); | |
| 62 | let replica = Replica::open_or_create(directory.path().join("replica.sqlite"), None, || { | |
| 63 | Ok(image.clone()) | |
| 64 | }) | |
| 65 | .unwrap(); | |
| 66 | let id = replica | |
| 67 | .apply( | |
| 68 | "Alice", | |
| 69 | Edit { | |
| 70 | at: 134_000_000_000_000_000, | |
| 71 | ops: vec![Op::Page { | |
| 72 | space, | |
| 73 | op: PageOp::Text { | |
| 74 | text, | |
| 75 | range: 13..13, | |
| 76 | with: "X".into(), | |
| 77 | }, | |
| 78 | }], | |
| 79 | }, | |
| 80 | ) | |
| 81 | .unwrap(); | |
| 82 | let mut remote = Lost(HostedRemote::new(&guest, "Garden.one"), true); | |
| 83 | assert!(matches!( | |
| 84 | replica.sync_once(&mut remote), | |
| 85 | Err(notebook::Error::Remote(CommitError { | |
| 86 | state: CommitState::Unknown, | |
| 87 | .. | |
| 88 | })) | |
| 89 | )); | |
| 90 | replica.sync_once(&mut remote).unwrap(); | |
| 91 | assert!(matches!( | |
| 92 | replica.status(id).unwrap(), | |
| 93 | Some(notebook::EditStatus::Published { .. }) | |
| 94 | )); | |
| 95 | assert_eq!( | |
| 96 | server::text(&std::fs::read(file).unwrap()).2, | |
| 97 | "Original textX" | |
| 98 | ); | |
| 99 | assert!(replica.pending().unwrap().is_empty()); | |
| 100 | } | |
| 101 | ||
| 102 | #[test] | |
| 103 | fn concurrent_batches_keep_receipts_and_rebase_later_keystrokes() { | |
| 104 | let directory = tempfile::tempdir().unwrap(); | |
| 105 | let folder = notebook(directory.path()); | |
| 106 | let url = relay(Default::default()); | |
| 107 | let sharing = Sharing::new("").unwrap(); | |
| 108 | let host = host(&folder, &directory.path().join("host"), &sharing, &url); | |
| 109 | let code = code(&host); | |
| 110 | let (alice, _) = guest("Alice", &code, &url, &directory.path().join("alice")); | |
| 111 | let (bob, _) = guest("Bob", &code, &url, &directory.path().join("bob")); | |
| 112 | let file = folder.join("Garden.one"); | |
| 113 | let original = std::fs::read(&file).unwrap(); | |
| 114 | let (space, text, _) = server::text(&original); | |
| 115 | let edit = |with: &str, at| Edit { | |
| 116 | at: 134_000_000_000_000_000, | |
| 117 | ops: vec![Op::Page { | |
| 118 | space, | |
| 119 | op: PageOp::Text { | |
| 120 | text, | |
| 121 | range: at..at, | |
| 122 | with: with.into(), | |
| 123 | }, | |
| 124 | }], | |
| 125 | }; | |
| 126 | let a = Arc::new( | |
| 127 | Replica::open_or_create(directory.path().join("a.sqlite"), None, || { | |
| 128 | Ok(original.clone()) | |
| 129 | }) | |
| 130 | .unwrap(), | |
| 131 | ); | |
| 132 | let b = Arc::new( | |
| 133 | Replica::open_or_create(directory.path().join("b.sqlite"), None, || { | |
| 134 | Ok(original.clone()) | |
| 135 | }) | |
| 136 | .unwrap(), | |
| 137 | ); | |
| 138 | let first = a.apply("Alice", edit("A", 13)).unwrap(); | |
| 139 | let second = b.apply("Bob", edit("B", 0)).unwrap(); | |
| 140 | let barrier = Arc::new(Barrier::new(2)); | |
| 141 | std::thread::scope(|scope| { | |
| 142 | for (replica, guest) in [(&a, &alice), (&b, &bob)] { | |
| 143 | let barrier = Arc::clone(&barrier); | |
| 144 | scope.spawn(move || { | |
| 145 | barrier.wait(); | |
| 146 | replica | |
| 147 | .sync_once(&mut HostedRemote::new(guest, "Garden.one")) | |
| 148 | .unwrap(); | |
| 149 | }); | |
| 150 | } | |
| 151 | }); | |
| 152 | assert!(matches!( | |
| 153 | a.status(first).unwrap(), | |
| 154 | Some(notebook::EditStatus::Published { .. }) | |
| 155 | )); | |
| 156 | assert!(matches!( | |
| 157 | b.status(second).unwrap(), | |
| 158 | Some(notebook::EditStatus::Published { .. }) | |
| 159 | )); | |
| 160 | let next = server::text(&a.snapshot().unwrap()).2.find('A').unwrap() as u32 + 1; | |
| 161 | let third = a.apply("Alice", edit("a", next)).unwrap(); | |
| 162 | let mut remote = HostedRemote::new(&alice, "Garden.one"); | |
| 163 | for _ in 0..5 { | |
| 164 | a.sync_once(&mut remote).unwrap(); | |
| 165 | if matches!( | |
| 166 | a.status(third).unwrap(), | |
| 167 | Some(notebook::EditStatus::Published { .. }) | |
| 168 | ) { | |
| 169 | break; | |
| 170 | } | |
| 171 | } | |
| 172 | let result = std::fs::read(&file).unwrap(); | |
| 173 | let result_text = server::text(&result).2; | |
| 174 | assert!(result_text.contains("Aa"), "{result_text}"); | |
| 175 | assert!(result_text.contains('B'), "{result_text}"); | |
| 176 | let arena = onestore::Arena::default(); | |
| 177 | let mut section = onestore::Section::open(&arena, result).unwrap(); | |
| 178 | assert!(section.conflicts().unwrap().is_empty()); | |
| 179 | } |
crates/onestore/src/commit.rs+59| ... | ... | @@ -492,6 +492,38 @@ impl TryFrom<Wire> for Transaction { |
| 492 | 492 | } |
| 493 | 493 | |
| 494 | 494 | impl Transaction { |
| 495 | /// Combines a following transaction into one guarded publication, retaining every revision. | |
| 496 | pub fn extend(&mut self, next: Transaction) -> Result<(), crate::Error> { | |
| 497 | let length = self.base.length + self.append.len() as u64; | |
| 498 | if next.base.header != self.header || next.base.length != length { | |
| 499 | return Err(crate::Error { | |
| 500 | offset: 0, | |
| 501 | message: "Transactions are not consecutive", | |
| 502 | }); | |
| 503 | } | |
| 504 | if next.patches.iter().any(|(offset, bytes)| { | |
| 505 | *offset < 1024 || offset.saturating_add(bytes.len() as u64) > length | |
| 506 | }) { | |
| 507 | return Err(crate::Error { | |
| 508 | offset: 0, | |
| 509 | message: "A patch outside the base's data", | |
| 510 | }); | |
| 511 | } | |
| 512 | for (offset, bytes) in next.patches { | |
| 513 | let earlier = (self.base.length.saturating_sub(offset) as usize).min(bytes.len()); | |
| 514 | if earlier != 0 { | |
| 515 | self.patches.push((offset, bytes[..earlier].to_vec())); | |
| 516 | } | |
| 517 | if earlier < bytes.len() { | |
| 518 | let at = (offset + earlier as u64 - self.base.length) as usize; | |
| 519 | self.append[at..at + bytes.len() - earlier].copy_from_slice(&bytes[earlier..]); | |
| 520 | } | |
| 521 | } | |
| 522 | self.append.extend_from_slice(&next.append); | |
| 523 | self.header = next.header; | |
| 524 | Ok(()) | |
| 525 | } | |
| 526 | ||
| 495 | 527 | /// The image this transaction applies to. |
| 496 | 528 | pub fn base(&self) -> &Stamp { |
| 497 | 529 | &self.base |
| ... | ... | @@ -658,6 +690,33 @@ mod tests { |
| 658 | 690 | } |
| 659 | 691 | } |
| 660 | 692 | |
| 693 | #[test] | |
| 694 | fn consecutive_transactions_publish_the_same_bytes_together() { | |
| 695 | let original = vec![1; 2048]; | |
| 696 | let first = Transaction { | |
| 697 | base: Stamp::of(&original).unwrap(), | |
| 698 | append: vec![2; 32], | |
| 699 | patches: vec![(1024, vec![3; 8])], | |
| 700 | header: [4; 1024], | |
| 701 | }; | |
| 702 | let mut separate = original.clone(); | |
| 703 | first.apply(&mut separate).unwrap(); | |
| 704 | let second = Transaction { | |
| 705 | base: Stamp::of(&separate).unwrap(), | |
| 706 | append: vec![5; 16], | |
| 707 | patches: vec![(1028, vec![6; 8]), (2044, vec![7; 12])], | |
| 708 | header: [8; 1024], | |
| 709 | }; | |
| 710 | second.apply(&mut separate).unwrap(); | |
| 711 | let mut combined = first.clone(); | |
| 712 | assert!(combined.extend(first).is_err()); | |
| 713 | combined.extend(second).unwrap(); | |
| 714 | let combined = Transaction::from_bytes(&combined.to_bytes()).unwrap(); | |
| 715 | let mut together = original; | |
| 716 | combined.apply(&mut together).unwrap(); | |
| 717 | assert_eq!(together, separate); | |
| 718 | } | |
| 719 | ||
| 661 | 720 | #[test] |
| 662 | 721 | fn transactions_read_back_from_bytes() { |
| 663 | 722 | let transaction = Transaction { |