1use super::*;
2use crate::{PendingEdit, merge};
3use onestore::{Arena, ExGuid, Section, op::Edit};
4use std::collections::BTreeSet;
5
6#[derive(serde::Serialize, serde::Deserialize)]
7#[serde(deny_unknown_fields)]
8pub(super) struct Edits {
9 pub edits: Vec<(String, Edit)>,
10 pub revisions: BTreeMap<ExGuid, ExGuid>,
11}
12
13pub(super) fn encode(edits: &Edits, limit: usize) -> io::Result<Vec<u8>> {
14 struct Buffer {
15 bytes: Vec<u8>,
16 limit: usize,
17 }
18 impl io::Write for Buffer {
19 fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
20 if bytes.len() > self.limit.saturating_sub(self.bytes.len()) {
21 return Err(io::ErrorKind::FileTooLarge.into());
22 }
23 self.bytes.extend_from_slice(bytes);
24 Ok(bytes.len())
25 }
26
27 fn flush(&mut self) -> io::Result<()> {
28 Ok(())
29 }
30 }
31 let mut buffer = Buffer {
32 bytes: Vec::new(),
33 limit,
34 };
35 serde_json::to_writer(&mut buffer, edits).map_err(|error| {
36 io::Error::new(
37 error.io_error_kind().unwrap_or(io::ErrorKind::InvalidData),
38 error,
39 )
40 })?;
41 Ok(buffer.bytes)
42}
43
44pub(super) struct Waiting {
45 peer: [u8; 16],
46 request: Request,
47 reply: mpsc::Sender<Reply>,
48}
49
50fn retry(message: &str) -> Error {
51 Error::Remote(CommitError {
52 state: CommitState::NotCommitted,
53 error: io::Error::new(io::ErrorKind::ResourceBusy, message.to_owned()),
54 })
55}
56
57impl Served {
58 pub(super) fn batch(self: &Arc<Self>, peer: &[u8; 16], request: Request) -> Result<Reply> {
59 let path = request.path.clone();
60 if !self.writers.lock().unwrap().contains_key(&path) {
61 self.image(&path)?;
62 }
63 let mut writers = self.writers.lock().unwrap();
64 let writer = writers.entry(path.clone()).or_insert_with(|| {
65 let (send, receive) = mpsc::sync_channel::<Waiting>(64);
66 let served = Arc::downgrade(self);
67 thread::spawn(move || {
68 while let Ok(first) = receive.recv() {
69 let mut waiting = vec![first];
70 if let Ok(next) = receive.recv_timeout(Duration::from_millis(5)) {
71 waiting.push(next);
72 }
73 waiting.extend(receive.try_iter().take(62));
74 let Some(served) = served.upgrade() else {
75 return;
76 };
77 served.write_batches(&path, waiting);
78 }
79 });
80 send
81 });
82 let (reply, receive) = mpsc::channel();
83 writer
84 .try_send(Waiting {
85 peer: *peer,
86 request,
87 reply,
88 })
89 .map_err(|_| retry("The section's writer is busy"))?;
90 drop(writers);
91 receive.recv_timeout(TIMEOUT).map_err(|_| {
92 Error::Remote(CommitError {
93 state: CommitState::Unknown,
94 error: io::Error::new(
95 io::ErrorKind::TimedOut,
96 "The section's writer did not answer",
97 ),
98 })
99 })
100 }
101
102 fn write_batches(&self, path: &str, waiting: Vec<Waiting>) {
103 let before = match self.image(path) {
104 Ok(image) => image,
105 Err(error) => {
106 for batch in waiting {
107 let _ = batch.reply.send(failed(0, retry(&error.to_string())));
108 }
109 return;
110 }
111 };
112 let mut image = Arc::clone(&before);
113 let mut combined: Option<Transaction> = None;
114 let mut accepted = Vec::new();
115 for batch in waiting {
116 match self.prepare_batch(path, &image, &batch) {
117 Ok(Some(transaction)) => {
118 let mut next = (*image).clone();
119 let extended =
120 transaction
121 .apply(&mut next)
122 .and_then(|()| match &mut combined {
123 Some(combined) => combined.extend(transaction),
124 None => {
125 combined = Some(transaction);
126 Ok(())
127 }
128 });
129 match extended {
130 Ok(()) => {
131 image = Arc::new(next);
132 accepted.push(batch);
133 }
134 Err(error) => {
135 let _ = batch.reply.send(failed(0, error.into()));
136 }
137 }
138 }
139 Ok(None) => accepted.push(batch),
140 Err(error) => {
141 let _ = batch.reply.send(failed(0, error));
142 }
143 }
144 }
145 if accepted.is_empty() {
146 return;
147 }
148 let result = match &combined {
149 Some(transaction) => self.storage.commit(path, transaction),
150 None => self
151 .storage
152 .confirm(path, &Stamp::of(&image).expect("section stamp"))
153 .map_err(Error::from)
154 .and_then(|()| {
155 image = self.image(path).map_err(|error| CommitError {
156 state: CommitState::Committed,
157 error: match error {
158 Error::Io(error) | Error::RemoteIo(error) => error,
159 Error::Remote(error) => error.error,
160 error => io::Error::other(error),
161 },
162 })?;
163 Ok(())
164 }),
165 };
166 if result.is_ok() {
167 self.tell_delta(path, &before, &image, None);
168 self.keep(path, Arc::clone(&image));
169 self.changed(&[path.to_owned()]);
170 }
171 for batch in accepted {
172 let reply = match &result {
173 Ok(()) => Reply {
174 stamp: Some((&Stamp::of(&image).expect("section stamp")).into()),
175 ..Reply::default()
176 },
177 Err(Error::Remote(error)) => failed(
178 0,
179 Error::Remote(CommitError {
180 state: error.state,
181 error: io::Error::new(error.error.kind(), error.error.to_string()),
182 }),
183 ),
184 Err(error) => failed(0, retry(&error.to_string())),
185 };
186 let _ = batch.reply.send(reply);
187 }
188 }
189
190 fn prepare_batch(
191 &self,
192 path: &str,
193 image: &[u8],
194 batch: &Waiting,
195 ) -> Result<Option<Transaction>> {
196 if !self.guests.lock().unwrap().contains_key(&batch.peer) {
197 return Err(refused(
198 io::ErrorKind::PermissionDenied,
199 "The device is no longer connected",
200 ));
201 }
202 let stamp: Stamp = batch
203 .request
204 .stamp
205 .as_ref()
206 .ok_or_else(|| retry("No batch base"))?
207 .try_into()?;
208 let edits: Edits = serde_json::from_slice(&self.carried(&batch.peer, &batch.request)?)
209 .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
210 if edits.revisions.is_empty()
211 || edits
212 .revisions
213 .values()
214 .any(|revision| revision.guid == [0; 16])
215 {
216 return Err(retry("No batch revisions"));
217 }
218 let store = Store::parse(image)?;
219 let index = RevisionIndex::parse(&store)?;
220 if edits.revisions.iter().all(|(space, revision)| {
221 index
222 .spaces
223 .get(space)
224 .is_some_and(|space| space.revisions.contains_key(revision))
225 }) {
226 return Ok(None);
227 }
228 let names = edits
229 .revisions
230 .iter()
231 .filter(|(space, revision)| {
232 !index
233 .spaces
234 .get(space)
235 .is_some_and(|space| space.revisions.contains_key(revision))
236 })
237 .map(|(space, revision)| (*space, *revision))
238 .collect();
239 let arena = Arena::default();
240 let mut next = Section::open(&arena, image.to_vec())?;
241 let queued: Vec<PendingEdit> = edits
242 .edits
243 .into_iter()
244 .enumerate()
245 .map(|(id, (author, edit))| PendingEdit {
246 id: id as u64,
247 author,
248 edit,
249 })
250 .collect();
251 if Stamp::of(image)? == stamp {
252 for edit in &queued {
253 next.apply(&edit.author, &edit.edit)?;
254 }
255 } else {
256 let base = self
257 .images
258 .lock()
259 .unwrap()
260 .iter()
261 .find_map(|(held, at, image)| {
262 (held == path && *at == stamp).then(|| Arc::clone(image))
263 })
264 .ok_or_else(|| retry("The batch's base is no longer held"))?;
265 let old_arena = Arena::default();
266 let mut old = Section::open(&old_arena, (*base).clone())?;
267 if old.root() != next.root() {
268 return Err(retry("The batch belongs to another section"));
269 }
270 let merged = merge::rebase(&mut old, &mut next, &queued, &BTreeSet::new())?;
271 if !merged.conflicts.is_empty() {
272 let local_arena = Arena::default();
273 let mut local = Section::open(&local_arena, (*base).clone())?;
274 for edit in &queued {
275 local.apply(&edit.author, &edit.edit)?;
276 }
277 for (space, (author, objects)) in &merged.conflicts {
278 let page = local.page(*space)?;
279 if next.page(*space).is_ok_and(|remote| remote == page) {
280 return Err(retry("The batch's edits already landed"));
281 }
282 merge::conflict_page(
283 &mut next,
284 &mut local,
285 &merged.moved,
286 *space,
287 &page,
288 author,
289 objects,
290 crate::now(),
291 None,
292 )?;
293 }
294 }
295 }
296 let transaction = next.seal_as(&names)?;
297 if edits.revisions.iter().any(|(space, revision)| {
298 !next
299 .newest()
300 .any(|(held, at)| held == *space && at == *revision)
301 }) {
302 return Err(refused(
303 io::ErrorKind::Unsupported,
304 "The batch needs a guest's commit",
305 ));
306 }
307 Ok(transaction)
308 }
309}
310
311#[cfg(test)]
312mod tests {
313 use super::*;
314 use onestore::{
315 op::{Op, PageOp},
316 page::{Attachment, PageObject},
317 };
318
319 #[cfg(not(target_arch = "wasm32"))]
320 #[test]
321 fn missing_sections_do_not_retain_writers_and_can_be_created_later() {
322 let directory = tempfile::tempdir().unwrap();
323 let folder = directory.path().join("notebook");
324 std::fs::create_dir(&folder).unwrap();
325 let served = Arc::new(Served {
326 storage: crate::session::Notebook::open(&folder, directory.path().join("cache"))
327 .unwrap()
328 .into_storage(),
329 images: Mutex::default(),
330 snapshots: Mutex::default(),
331 puts: Mutex::default(),
332 guests: Mutex::default(),
333 writers: Mutex::default(),
334 host: Mutex::default(),
335 room: Mutex::default(),
336 });
337 let peer = [1; 16];
338 let (line, _outbound) = mpsc::channel();
339 let (queue, _requests) = mpsc::sync_channel(1);
340 served.guests.lock().unwrap().insert(
341 peer,
342 Admitted {
343 line: crate::live::Line(line),
344 queue,
345 starts: BURST,
346 counted: Instant::now(),
347 held: Vec::new(),
348 },
349 );
350 for at in 0..128 {
351 let reply = served.handle(
352 &peer,
353 kind::EDITS,
354 &minicbor::to_vec(Request {
355 path: format!("Missing{at}.one"),
356 ..Request::default()
357 })
358 .unwrap(),
359 );
360 assert!(reply.failure.is_some());
361 }
362 assert!(served.writers.lock().unwrap().is_empty());
363
364 let path = "Missing0.one";
365 let image = onestore::create_section(path, "Original text", "Fixture").unwrap();
366 std::fs::write(folder.join(path), &image).unwrap();
367 let store = Store::parse(&image).unwrap();
368 let index = RevisionIndex::parse(&store).unwrap();
369 let document = onestore::document::Document::parse(&index).unwrap();
370 let (space, text) = document
371 .spaces
372 .iter()
373 .find_map(|(space, object)| {
374 object.revisions[&object.contexts[&ExGuid::default()]]
375 .nodes
376 .iter()
377 .find_map(|(id, node)| {
378 matches!(node.kind, onestore::document::Kind::RichText { .. })
379 .then_some((*space, *id))
380 })
381 })
382 .unwrap();
383 let edit = Edit {
384 at: 134_000_000_000_000_000,
385 ops: vec![Op::Page {
386 space,
387 op: PageOp::Text {
388 text,
389 range: 13..13,
390 with: "X".into(),
391 },
392 }],
393 };
394 let arena = Arena::default();
395 let mut section = Section::open(&arena, image).unwrap();
396 section.apply("Alice", &edit).unwrap();
397 let transaction = section.seal().unwrap().unwrap();
398 let request = Request {
399 path: path.into(),
400 stamp: Some(transaction.base().into()),
401 bytes: Some(
402 encode(
403 &Edits {
404 edits: vec![("Alice".into(), edit)],
405 revisions: section.newest().filter(|(at, _)| *at == space).collect(),
406 },
407 LIMIT,
408 )
409 .unwrap(),
410 ),
411 ..Request::default()
412 };
413 let first = served.batch(&peer, request.clone()).unwrap();
414 assert!(first.failure.is_none(), "{:?}", first.failure);
415 let committed = std::fs::read(folder.join(path)).unwrap();
416 let stamp = Stamp::of(&committed).unwrap();
417 assert_eq!(
418 Stamp::try_from(first.stamp.as_ref().unwrap()).unwrap(),
419 stamp
420 );
421 let read_arena = Arena::default();
422 let reread = Section::open(&read_arena, committed.clone()).unwrap();
423 assert_eq!(reread.page(space).unwrap(), section.page(space).unwrap());
424 let repeated = served.batch(&peer, request).unwrap();
425 assert!(repeated.failure.is_none(), "{:?}", repeated.failure);
426 let replayed = std::fs::read(folder.join(path)).unwrap();
427 let confirmed = Stamp::of(&replayed).unwrap();
428 assert_eq!(
429 Stamp::try_from(repeated.stamp.as_ref().unwrap()).unwrap(),
430 confirmed
431 );
432 {
433 let images = served.images.lock().unwrap();
434 let (_, held, image) = images.iter().find(|(held, _, _)| held == path).unwrap();
435 assert_eq!(*held, confirmed);
436 assert_eq!(**image, replayed);
437 }
438 assert_eq!(&replayed[1024..], &committed[1024..]);
439 let replay_arena = Arena::default();
440 let replay = Section::open(&replay_arena, replayed).unwrap();
441 assert_eq!(replay.page(space).unwrap(), section.page(space).unwrap());
442 assert_eq!(served.writers.lock().unwrap().len(), 1);
443 }
444
445 #[test]
446 fn attachment_encoding_stops_at_the_upload_budget() {
447 let id = ExGuid {
448 guid: [1; 16],
449 n: 1,
450 };
451 let bytes: Arc<[u8]> = (0..4096)
452 .map(|at| (at % 251) as u8)
453 .collect::<Vec<_>>()
454 .into();
455 let edits = Edits {
456 edits: vec![(
457 "Ada".into(),
458 Edit {
459 at: 0,
460 ops: vec![Op::Page {
461 space: id,
462 op: PageOp::Add {
463 object: PageObject::Attachment(Attachment {
464 id,
465 filename: "payload.bin".into(),
466 source_path: None,
467 size: None,
468 layout: Default::default(),
469 bytes: Some(Arc::clone(&bytes)),
470 preview: None,
471 recording: None,
472 tags: Vec::new(),
473 }),
474 before: None,
475 },
476 }],
477 },
478 )],
479 revisions: BTreeMap::new(),
480 };
481 let encoded = serde_json::to_vec(&edits).unwrap();
482 assert_eq!(encode(&edits, encoded.len()).unwrap(), encoded);
483 let decoded: Edits = serde_json::from_slice(&encoded).unwrap();
484 assert_eq!(decoded.edits, edits.edits);
485 assert_eq!(
486 encode(&edits, encoded.len() - 1).unwrap_err().kind(),
487 io::ErrorKind::FileTooLarge
488 );
489 assert_eq!(
490 encode(&edits, 128).unwrap_err().kind(),
491 io::ErrorKind::FileTooLarge
492 );
493 }
494}