1//! The section thread: the only owner of the parsed section, which is the cached base
2//! image with each sealed batch replayed and the open batch's edits applied. Edits apply
3//! here and are written in one SQLite transaction per burst; the sync thread asks it to seal
4//! and, when the remote changed, to rebase the queue. In the browser, which has one thread,
5//! the section is served as each request is sent, and a rebuild runs as it is asked for.
6
7use crate::{Result, base, lock, merge, queue, worker::Signal};
8use onestore::{
9 Arena, ExGuid, Section, Transaction,
10 op::{Edit, OpError},
11 page::Page,
12 protected::Key,
13};
14use rusqlite::{Connection, OptionalExtension, TransactionBehavior, params};
15use std::{
16 collections::{BTreeMap, BTreeSet, VecDeque},
17 io,
18 sync::{Arc, Mutex, Weak, mpsc},
19};
20
21pub(crate) type Reply<T> = Box<dyn FnOnce(Result<T>) + Send>;
22
23pub(crate) enum Request {
24 Apply {
25 author: String,
26 edit: Edit,
27 reply: Reply<u64>,
28 },
29 Page {
30 space: ExGuid,
31 reply: Reply<Page>,
32 },
33 Pages {
34 reply: Reply<Vec<(ExGuid, String, u32)>>,
35 },
36 Conflicts {
37 reply: Reply<Vec<(ExGuid, Vec<onestore::ConflictPage>)>>,
38 },
39 Versions {
40 reply: Reply<Vec<(ExGuid, Vec<onestore::PageVersion>)>>,
41 },
42 Version {
43 space: ExGuid,
44 version: ExGuid,
45 reply: Reply<Page>,
46 },
47 /// Answers once the edits before it are written.
48 Flush {
49 reply: Reply<()>,
50 },
51 /// Seals the open batch, when no sealed batch waits for publication, recording the
52 /// attempt to publish it.
53 Seal {
54 reply: Reply<Option<Sealed>>,
55 },
56 /// Replays the queue on `image` (the stored remote image, or the base, when `None`);
57 /// answers with the pages the remote changed.
58 Rebase {
59 image: Option<Vec<u8>>,
60 reply: Reply<Vec<ExGuid>>,
61 },
62 /// Rereads the queue after the sync thread replaced it.
63 Reopen {
64 reply: Reply<()>,
65 },
66 /// From the thread that rebuilt the section: hand it the requests.
67 #[cfg_attr(target_arch = "wasm32", allow(dead_code))]
68 Handover(mpsc::SyncSender<Takeover>),
69 /// From the thread that rebuilt nothing: carry on with the section as it is.
70 #[cfg_attr(target_arch = "wasm32", allow(dead_code))]
71 Resume,
72}
73
74/// The request channel and the requests held while the section was rebuilt.
75type Takeover = (mpsc::Receiver<Request>, VecDeque<Request>);
76
77/// A sealed batch: the transaction publishing it, or none when its edits changed nothing.
78pub(crate) struct Sealed {
79 pub batch: i64,
80 pub transaction: Option<Transaction>,
81}
82
83impl Request {
84 fn fail(self, message: &str) {
85 let error = || io::Error::other(message.to_owned()).into();
86 match self {
87 Self::Apply { reply, .. } => reply(Err(error())),
88 Self::Page { reply, .. } => reply(Err(error())),
89 Self::Pages { reply } => reply(Err(error())),
90 Self::Conflicts { reply } => reply(Err(error())),
91 Self::Versions { reply } => reply(Err(error())),
92 Self::Version { reply, .. } => reply(Err(error())),
93 Self::Flush { reply } => reply(Err(error())),
94 Self::Seal { reply } => reply(Err(error())),
95 Self::Rebase { reply, .. } => reply(Err(error())),
96 Self::Reopen { reply } => reply(Err(error())),
97 Self::Handover(_) | Self::Resume => {}
98 }
99 }
100}
101
102/// The section thread as the replica sees it. Rereading the section after the queue was
103/// replaced (a rebase, a released attempt) happens on a new thread while the current one
104/// keeps answering page reads from the section as it was; the new thread then takes the
105/// requests over, so reads never wait for a rebuild.
106pub(crate) struct Thread {
107 pub(crate) connection: Mutex<Connection>,
108 /// The sync worker, which a durable burst of edits wakes.
109 pub(crate) worker: Mutex<Weak<Signal>>,
110 /// Taken when the replica drops, which ends every section thread.
111 sender: Mutex<Option<mpsc::Sender<Request>>>,
112 #[cfg(not(target_arch = "wasm32"))]
113 threads: Mutex<Vec<std::thread::JoinHandle<()>>>,
114 /// A password-protected section's key, under which its queue is sealed too.
115 pub(crate) key: Option<Key>,
116}
117
118impl Thread {
119 pub(crate) fn send(&self, request: Request) -> Result<()> {
120 self.sender
121 .lock()
122 .ok()
123 .and_then(|sender| sender.as_ref()?.send(request).ok())
124 .ok_or_else(|| io::Error::other("The section thread stopped"))?;
125 #[cfg(target_arch = "wasm32")]
126 serving::serve(self);
127 Ok(())
128 }
129
130 /// Ends the section threads and waits for them, so the cache is released.
131 pub(crate) fn stop(&self) {
132 if let Ok(mut sender) = self.sender.lock() {
133 sender.take();
134 }
135 #[cfg(target_arch = "wasm32")]
136 serving::serve(self);
137 #[cfg(not(target_arch = "wasm32"))]
138 while let Some(thread) = self
139 .threads
140 .lock()
141 .ok()
142 .and_then(|mut threads| threads.pop())
143 {
144 let _ = thread.join();
145 }
146 }
147
148 #[cfg(not(target_arch = "wasm32"))]
149 fn start<F: Future<Output = ()>>(
150 self: &Arc<Self>,
151 run: impl FnOnce(Arc<Self>) -> F + Send + 'static,
152 ) -> Result<()> {
153 let shared = Arc::clone(self);
154 let thread = std::thread::Builder::new()
155 .name("onestore-section".into())
156 .spawn(move || crate::task::complete(run(shared)))?;
157 self.threads
158 .lock()
159 .map_err(|_| io::Error::other("The section thread panicked"))?
160 .push(thread);
161 Ok(())
162 }
163
164 #[cfg(target_arch = "wasm32")]
165 fn start<F: Future<Output = ()> + 'static>(
166 self: &Arc<Self>,
167 run: impl FnOnce(Arc<Self>) -> F,
168 ) -> Result<()> {
169 serving::start(self, Box::pin(run(Arc::clone(self))));
170 Ok(())
171 }
172}
173
174/// The browser's section threads: each replica's `run`, polled whenever a request is sent
175/// until it waits for the next.
176#[cfg(target_arch = "wasm32")]
177mod serving {
178 use super::Thread;
179 use std::{cell::RefCell, collections::BTreeMap, future::Future, pin::Pin, task};
180
181 type Serving = Pin<Box<dyn Future<Output = ()>>>;
182
183 thread_local! {
184 /// By the address of each replica's `Thread`, which its `run` keeps alive.
185 static SERVING: RefCell<BTreeMap<usize, Serving>> = const { RefCell::new(BTreeMap::new()) };
186 }
187
188 pub(super) fn start(thread: &Thread, run: Serving) {
189 SERVING.with_borrow_mut(|serving| serving.insert(std::ptr::from_ref(thread).addr(), run));
190 serve(thread);
191 }
192
193 /// Serves what `thread` was sent. A request sent while it is being served, as from a
194 /// reply, waits for the serving already under way, which takes it next.
195 pub(super) fn serve(thread: &Thread) {
196 let id = std::ptr::from_ref(thread).addr();
197 let Some(mut run) = SERVING.with_borrow_mut(|serving| serving.remove(&id)) else {
198 return;
199 };
200 let mut context = task::Context::from_waker(task::Waker::noop());
201 if run.as_mut().poll(&mut context).is_pending() {
202 SERVING.with_borrow_mut(|serving| serving.insert(id, run));
203 }
204 }
205}
206
207/// The next request, or none once the replica dropped.
208#[cfg(not(target_arch = "wasm32"))]
209async fn receive(requests: &mpsc::Receiver<Request>) -> Option<Request> {
210 requests.recv().ok()
211}
212
213/// The next request, waiting for `serving::serve` where none is queued, or none once the
214/// replica dropped.
215#[cfg(target_arch = "wasm32")]
216async fn receive(requests: &mpsc::Receiver<Request>) -> Option<Request> {
217 std::future::poll_fn(|_| match requests.try_recv() {
218 Ok(request) => std::task::Poll::Ready(Some(request)),
219 Err(mpsc::TryRecvError::Disconnected) => std::task::Poll::Ready(None),
220 Err(mpsc::TryRecvError::Empty) => std::task::Poll::Pending,
221 })
222 .await
223}
224
225/// Starts the section thread once the queue opens.
226pub(crate) fn spawn(connection: Connection, key: Option<Key>) -> Result<(Arc<Thread>, ExGuid)> {
227 let (sender, requests) = mpsc::channel();
228 let shared = Arc::new(Thread {
229 connection: Mutex::new(connection),
230 worker: Mutex::new(Weak::new()),
231 sender: Mutex::new(Some(sender)),
232 #[cfg(not(target_arch = "wasm32"))]
233 threads: Mutex::new(Vec::new()),
234 key,
235 });
236 let (ready, opened) = mpsc::sync_channel(1);
237 shared.start(move |shared| run(shared, requests, VecDeque::new(), Some(ready)))?;
238 match opened.recv() {
239 Ok(Ok(root)) => Ok((shared, root)),
240 Ok(Err(error)) => {
241 shared.stop();
242 Err(error)
243 }
244 Err(_) => {
245 shared.stop();
246 Err(io::Error::other("The section thread panicked").into())
247 }
248 }
249}
250
251enum Next {
252 Stop,
253 Reopen,
254 Handover(mpsc::SyncSender<Takeover>, VecDeque<Request>),
255}
256
257async fn run(
258 shared: Arc<Thread>,
259 requests: mpsc::Receiver<Request>,
260 mut backlog: VecDeque<Request>,
261 mut ready: Option<mpsc::SyncSender<Result<ExGuid>>>,
262) {
263 loop {
264 let arena = Arena::default();
265 let working = match Working::open(&arena, &shared.connection, shared.key.as_ref()) {
266 Ok(working) => working,
267 Err(error) => {
268 if let Some(ready) = ready.take() {
269 let _ = ready.send(Err(error));
270 } else {
271 fail(error, backlog, requests);
272 }
273 return;
274 }
275 };
276 if let Some(ready) = ready.take() {
277 let _ = ready.send(Ok(working.section.root()));
278 }
279 match working.serve(&shared, &requests, &mut backlog).await {
280 Next::Stop => return,
281 Next::Reopen => {}
282 Next::Handover(to, held) => {
283 let _ = to.send((requests, held));
284 return;
285 }
286 }
287 }
288}
289
290/// Answers every request with `error` until the replica drops.
291fn fail(error: crate::Error, backlog: VecDeque<Request>, requests: mpsc::Receiver<Request>) {
292 let message = error.to_string();
293 for request in backlog.into_iter().chain(requests.iter()) {
294 request.fail(&message);
295 }
296}
297
298/// What a rebuilding thread does before it rereads the section.
299#[cfg_attr(target_arch = "wasm32", allow(dead_code))]
300enum Job {
301 Rebase {
302 image: Option<Vec<u8>>,
303 reply: Reply<Vec<ExGuid>>,
304 },
305 Reopen {
306 reply: Reply<()>,
307 },
308}
309
310/// Runs `job`, rereads the section and takes the requests over from the thread that
311/// started it; a failed rebase changed nothing that thread serves.
312#[cfg(not(target_arch = "wasm32"))]
313fn build(shared: Arc<Thread>, job: Job, signal: mpsc::Sender<Request>) {
314 let answer: Box<dyn FnOnce() + Send> = match job {
315 Job::Reopen { reply } => Box::new(move || reply(Ok(()))),
316 Job::Rebase { image, reply } => {
317 match rebase(&shared.connection, image, shared.key.as_ref()) {
318 Ok(changed) => Box::new(move || reply(Ok(changed))),
319 Err(error) => {
320 reply(Err(error));
321 let _ = signal.send(Request::Resume);
322 return;
323 }
324 }
325 }
326 };
327 let arena = Arena::default();
328 let opened = Working::open(&arena, &shared.connection, shared.key.as_ref());
329 let (to, from) = mpsc::sync_channel(1);
330 if signal.send(Request::Handover(to)).is_err() {
331 return;
332 }
333 drop(signal);
334 answer();
335 let Ok((requests, mut backlog)) = from.recv() else {
336 return;
337 };
338 let working = match opened {
339 Ok(working) => working,
340 Err(error) => return fail(error, backlog, requests),
341 };
342 match crate::task::complete(working.serve(&shared, &requests, &mut backlog)) {
343 Next::Stop => {}
344 Next::Reopen => crate::task::complete(run(shared, requests, backlog, None)),
345 Next::Handover(to, held) => {
346 let _ = to.send((requests, held));
347 }
348 }
349}
350
351/// `build` where there is one thread: runs `job` before the section is reread, which the
352/// serving thread does once true is returned; a failed rebase changed nothing it serves.
353#[cfg(target_arch = "wasm32")]
354fn inline(shared: &Thread, job: Job) -> bool {
355 match job {
356 Job::Reopen { reply } => reply(Ok(())),
357 Job::Rebase { image, reply } => {
358 match rebase(&shared.connection, image, shared.key.as_ref()) {
359 Ok(changed) => reply(Ok(changed)),
360 Err(error) => {
361 reply(Err(error));
362 return false;
363 }
364 }
365 }
366 }
367 true
368}
369
370/// An edit applied to the section and not yet written, with who waits for it.
371struct Accepted {
372 author: String,
373 edit: Edit,
374 reply: Reply<u64>,
375}
376
377struct Working<'a> {
378 section: Section<'a>,
379 /// The batch collecting edits; sealed batches before it are replayed.
380 open: Option<i64>,
381 /// Spaces the open batch's edits change.
382 touched: BTreeSet<ExGuid>,
383}
384
385impl<'a> Working<'a> {
386 fn open(arena: &'a Arena, connection: &Mutex<Connection>, key: Option<&Key>) -> Result<Self> {
387 let (section, open, touched) = replay(arena, &*lock(connection)?, key)?;
388 Ok(Self {
389 section,
390 open,
391 touched,
392 })
393 }
394
395 async fn serve(
396 mut self,
397 shared: &Arc<Thread>,
398 requests: &mpsc::Receiver<Request>,
399 backlog: &mut VecDeque<Request>,
400 ) -> Next {
401 let connection = &shared.connection;
402 // While another thread rebuilds the section, reads answer from this one and every
403 // other request waits for the rebuilt section.
404 let mut held: Option<VecDeque<Request>> = None;
405 loop {
406 let first = match backlog.pop_front() {
407 Some(request) => request,
408 None => match receive(requests).await {
409 Some(request) => request,
410 None => return Next::Stop,
411 },
412 };
413 let mut burst: VecDeque<Request> = VecDeque::from([first]);
414 burst.extend(backlog.drain(..));
415 burst.extend(requests.try_iter());
416 let mut accepted = Vec::new();
417 // Requests that write the queue beyond the burst's edits run after its reads.
418 let mut later = VecDeque::new();
419 while let Some(request) = burst.pop_front().or_else(|| later.pop_front()) {
420 if let Some(waiting) = &mut held {
421 match request {
422 Request::Page { space, reply } => {
423 reply(self.section.page(space).map_err(Into::into))
424 }
425 Request::Pages { reply } => reply(self.section.pages().map_err(Into::into)),
426 Request::Conflicts { reply } => {
427 reply(self.section.conflicts().map_err(Into::into))
428 }
429 Request::Versions { reply } => {
430 reply(self.section.versions().map_err(Into::into))
431 }
432 Request::Version {
433 space,
434 version,
435 reply,
436 } => reply(self.section.version(space, version).map_err(Into::into)),
437 Request::Handover(to) => {
438 let mut waiting = held.take().unwrap_or_default();
439 waiting.extend(burst.drain(..).chain(later.drain(..)));
440 return Next::Handover(to, waiting);
441 }
442 Request::Resume => {
443 let waiting = held.take().unwrap_or_default();
444 burst = waiting.into_iter().chain(burst.drain(..)).collect();
445 }
446 other => waiting.push_back(other),
447 }
448 continue;
449 }
450 let reopen = match request {
451 Request::Apply {
452 author,
453 edit,
454 reply,
455 } => self
456 .accept(author, edit, reply, &mut accepted, shared)
457 .err()
458 .unwrap_or(false),
459 Request::Page { space, reply } => {
460 reply(self.section.page(space).map_err(Into::into));
461 false
462 }
463 Request::Pages { reply } => {
464 reply(self.section.pages().map_err(Into::into));
465 false
466 }
467 Request::Conflicts { reply } => {
468 reply(self.section.conflicts().map_err(Into::into));
469 false
470 }
471 Request::Versions { reply } => {
472 reply(self.section.versions().map_err(Into::into));
473 false
474 }
475 Request::Version {
476 space,
477 version,
478 reply,
479 } => {
480 reply(self.section.version(space, version).map_err(Into::into));
481 false
482 }
483 request @ (Request::Flush { .. } | Request::Seal { .. })
484 if !burst.is_empty() =>
485 {
486 later.push_back(request);
487 false
488 }
489 Request::Flush { reply } => {
490 let written = self.flush(shared, &mut accepted);
491 reply(if written {
492 Ok(())
493 } else {
494 Err(io::Error::other("The queue could not be written").into())
495 });
496 !written
497 }
498 Request::Seal { reply } => {
499 if !self.flush(shared, &mut accepted) {
500 reply(Err(
501 io::Error::other("The queue could not be written").into()
502 ));
503 true
504 } else {
505 let sealed = self.seal(connection);
506 let failed = sealed.is_err();
507 reply(sealed);
508 failed
509 }
510 }
511 Request::Rebase { image, reply } => {
512 if self.flush(shared, &mut accepted) {
513 #[cfg(target_arch = "wasm32")]
514 let reopen = inline(shared, Job::Rebase { image, reply });
515 #[cfg(not(target_arch = "wasm32"))]
516 let reopen = {
517 held = self
518 .rebuild(shared, Job::Rebase { image, reply })
519 .then(VecDeque::new);
520 false
521 };
522 reopen
523 } else {
524 reply(Err(
525 io::Error::other("The queue could not be written").into()
526 ));
527 true
528 }
529 }
530 Request::Reopen { reply } => {
531 self.flush(shared, &mut accepted);
532 #[cfg(target_arch = "wasm32")]
533 let reopen = inline(shared, Job::Reopen { reply });
534 #[cfg(not(target_arch = "wasm32"))]
535 let reopen = {
536 held = self
537 .rebuild(shared, Job::Reopen { reply })
538 .then(VecDeque::new);
539 false
540 };
541 reopen
542 }
543 Request::Handover(_) | Request::Resume => false,
544 };
545 if reopen {
546 backlog.extend(burst.drain(..).chain(later.drain(..)));
547 return Next::Reopen;
548 }
549 }
550 if !self.flush(shared, &mut accepted) {
551 return Next::Reopen;
552 }
553 }
554 }
555
556 /// Starts a thread that runs `job` and rereads the section; false when the replica is
557 /// stopping, which answers the job's request.
558 #[cfg(not(target_arch = "wasm32"))]
559 fn rebuild(&self, shared: &Arc<Thread>, job: Job) -> bool {
560 let signal = shared
561 .sender
562 .lock()
563 .ok()
564 .and_then(|sender| sender.as_ref().cloned());
565 let Some(signal) = signal else {
566 match job {
567 Job::Rebase { reply, .. } => {
568 reply(Err(io::Error::other("The section thread stopped").into()))
569 }
570 Job::Reopen { reply } => {
571 reply(Err(io::Error::other("The section thread stopped").into()))
572 }
573 }
574 return false;
575 };
576 shared
577 .start(move |shared| async move { build(shared, job, signal) })
578 .is_ok()
579 }
580
581 /// Applies an edit, keeping it to be written with the burst. A refused edit is answered
582 /// here; `Err(true)` when it failed part way and the section must be reread.
583 fn accept(
584 &mut self,
585 author: String,
586 edit: Edit,
587 reply: Reply<u64>,
588 accepted: &mut Vec<Accepted>,
589 shared: &Thread,
590 ) -> std::result::Result<(), bool> {
591 match self.section.apply(&author, &edit) {
592 Ok(()) => {
593 self.touched
594 .extend(queue::spaces(&edit, self.section.root()));
595 accepted.push(Accepted {
596 author,
597 edit,
598 reply,
599 });
600 Ok(())
601 }
602 Err(error) => {
603 let broken = matches!(error, OpError::Failed(_));
604 if broken {
605 self.flush(shared, accepted);
606 }
607 reply(Err(error.into()));
608 Err(broken)
609 }
610 }
611 }
612
613 /// Writes the accepted edits in one transaction, then answers their senders; false
614 /// when the write failed and the section holds edits the queue lacks.
615 fn flush(&mut self, shared: &Thread, accepted: &mut Vec<Accepted>) -> bool {
616 if accepted.is_empty() {
617 return true;
618 }
619 let written = (|| -> Result<Vec<u64>> {
620 let mut connection = lock(&shared.connection)?;
621 let transaction =
622 connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
623 let batch = match self.open {
624 Some(batch) => batch,
625 None => {
626 transaction.execute("INSERT INTO batches DEFAULT VALUES", [])?;
627 transaction.last_insert_rowid()
628 }
629 };
630 let mut ids = Vec::new();
631 for edit in accepted.iter() {
632 ids.push(queue::insert(
633 &transaction,
634 shared.key.as_ref(),
635 None,
636 batch,
637 &edit.author,
638 &edit.edit,
639 )?);
640 }
641 transaction.commit()?;
642 self.open = Some(batch);
643 Ok(ids)
644 })();
645 let ok = written.is_ok();
646 match written {
647 Ok(ids) => {
648 for (edit, id) in accepted.drain(..).zip(ids) {
649 (edit.reply)(Ok(id));
650 }
651 crate::edited(&shared.worker);
652 }
653 Err(error) => {
654 let message = error.to_string();
655 for edit in accepted.drain(..) {
656 (edit.reply)(Err(io::Error::other(message.clone()).into()));
657 }
658 }
659 }
660 ok
661 }
662
663 /// Seals the open batch unless a sealed batch still waits for publication.
664 fn seal(&mut self, connection: &Mutex<Connection>) -> Result<Option<Sealed>> {
665 let Some(batch) = self.open else {
666 return Ok(None);
667 };
668 {
669 let connection = lock(connection)?;
670 let waiting: bool = connection.query_row(
671 "SELECT EXISTS(SELECT 1 FROM batches WHERE sealed IS NOT NULL)",
672 [],
673 |row| row.get(0),
674 )?;
675 if waiting {
676 return Ok(None);
677 }
678 }
679 let transaction = self.section.seal()?;
680 // A batch whose edits change nothing is recorded against the section's root.
681 let root = self.section.root();
682 let revisions: BTreeMap<ExGuid, ExGuid> = self
683 .section
684 .newest()
685 .filter(|(space, _)| {
686 self.touched.contains(space) || (self.touched.is_empty() && *space == root)
687 })
688 .collect();
689 // The sync thread publishes what it seals at once: the attempt is recorded with it.
690 lock(connection)?.execute(
691 "UPDATE batches SET sealed=?1, revisions=?2, attempted=?3 WHERE id=?4",
692 params![
693 serde_json::to_string(&transaction).map_err(io::Error::other)?,
694 serde_json::to_string(&revisions).map_err(io::Error::other)?,
695 transaction.is_some(),
696 batch
697 ],
698 )?;
699 self.open = None;
700 self.touched.clear();
701 Ok(Some(Sealed { batch, transaction }))
702 }
703}
704
705/// Replays every queued edit on `image`, which becomes the base; each page whose local
706/// version the remote could not take gains a conflict page holding it, queued as a new edit.
707/// Returns the pages the remote changed.
708fn rebase(
709 connection: &Mutex<Connection>,
710 image: Option<Vec<u8>>,
711 key: Option<&Key>,
712) -> Result<Vec<ExGuid>> {
713 let (base_image, image, edits) = {
714 let connection = lock(connection)?;
715 if connection.query_row(
716 "SELECT EXISTS(SELECT 1 FROM batches WHERE attempted=1)",
717 [],
718 |row| row.get::<_, bool>(0),
719 )? {
720 return Err(io::Error::new(
721 io::ErrorKind::InvalidInput,
722 "An uncertain attempt is never rebased",
723 )
724 .into());
725 }
726 let base_image = base::base(&connection)?;
727 let image = match image {
728 Some(image) => image,
729 None => {
730 base::read(&connection, base::Image::Remote)?.unwrap_or_else(|| base_image.clone())
731 }
732 };
733 (base_image, image, queue::load(&connection, key, None)?)
734 };
735 let old_arena = Arena::default();
736 let mut old = open(&old_arena, base_image, key)?;
737 let before: BTreeMap<ExGuid, ExGuid> = old.revisions().collect();
738 let remote_arena = Arena::default();
739 let remote = open(&remote_arena, image.clone(), key)?;
740 if old.root() != remote.root() {
741 return Err(io::Error::new(
742 io::ErrorKind::InvalidInput,
743 "Remote snapshot belongs to another document",
744 )
745 .into());
746 }
747 let after: BTreeMap<ExGuid, ExGuid> = remote.revisions().collect();
748 // The pages as the queue leaves them, read once a page conflicts: O(section).
749 let local_arena = Arena::default();
750 let mut local = None;
751 // Pages the remote already holds as the local edits leave them: their ops are done.
752 let mut converged = BTreeSet::new();
753 let (rewritten, added) = loop {
754 let arena = Arena::default();
755 let mut new = open(&arena, image.clone(), key)?;
756 let merged = merge::rebase(&mut old, &mut new, &edits, &converged)?;
757 if !merged.conflicts.is_empty() && local.is_none() {
758 local = Some(replay(&local_arena, &*lock(connection)?, key)?.0);
759 }
760 let pages: BTreeMap<ExGuid, Page> = merged
761 .conflicts
762 .keys()
763 .filter_map(|space| Some((*space, local.as_ref()?.page(*space).ok()?)))
764 .collect();
765 let settled: Vec<ExGuid> = pages
766 .iter()
767 .filter(|(space, page)| remote.page(**space).is_ok_and(|remote| remote == **page))
768 .map(|(space, _)| *space)
769 .collect();
770 if !settled.is_empty() {
771 converged.extend(settled);
772 continue;
773 }
774 let mut added = Vec::new();
775 for (space, (author, objects)) in &merged.conflicts {
776 // A page the queue went on to delete keeps no version of its own.
777 let (Some(page), Some(local)) = (pages.get(space), local.as_mut()) else {
778 continue;
779 };
780 let edit = merge::conflict_page(
781 &mut new,
782 local,
783 &merged.moved,
784 *space,
785 page,
786 author,
787 objects,
788 crate::now(),
789 None,
790 )?;
791 added.push((author.clone(), edit));
792 }
793 break (merged.rewritten, added);
794 };
795 let mut connection = lock(connection)?;
796 let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
797 base::write(&transaction, base::Image::Base, &image)?;
798 base::clear(&transaction, base::Image::Remote)?;
799 transaction.execute("INSERT INTO batches DEFAULT VALUES", [])?;
800 let batch = transaction.last_insert_rowid();
801 transaction.execute("UPDATE edits SET batch=?1", [batch])?;
802 transaction.execute("DELETE FROM batches WHERE id<>?1", [batch])?;
803 for (id, edit) in rewritten {
804 queue::rewrite(&transaction, key, id, &edit)?;
805 }
806 for (author, edit) in &added {
807 queue::insert(&transaction, key, None, batch, author, edit)?;
808 }
809 if transaction.query_row("SELECT count(*) FROM edits", [], |row| row.get::<_, i64>(0))? == 0 {
810 transaction.execute("DELETE FROM batches", [])?;
811 }
812 queue::collect(&transaction)?;
813 transaction.commit()?;
814 let mut changed: BTreeSet<ExGuid> = before
815 .iter()
816 .filter(|(space, rid)| after.get(space) != Some(rid))
817 .map(|(space, _)| *space)
818 .chain(
819 after
820 .keys()
821 .filter(|space| !before.contains_key(space))
822 .copied(),
823 )
824 .collect();
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 }
834 Ok(changed.into_iter().collect())
835}
836
837/// Batches in order with their sealed transactions: `None` while open, `Some(None)` when
838/// sealing stored nothing.
839fn batches(connection: &Connection) -> Result<Vec<(i64, Option<Option<Transaction>>)>> {
840 let mut query = connection.prepare("SELECT id, sealed FROM batches ORDER BY id")?;
841 let mut rows = query.query([])?;
842 let mut batches = Vec::new();
843 while let Some(row) = rows.next()? {
844 let sealed: Option<String> = row.get(1)?;
845 batches.push((
846 row.get(0)?,
847 sealed
848 .map(|sealed| serde_json::from_str(&sealed))
849 .transpose()
850 .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?,
851 ));
852 }
853 Ok(batches)
854}
855
856/// The cached base with each sealed batch replayed and the open batch's edits applied,
857/// with the open batch and the spaces its edits change.
858fn replay<'a>(
859 arena: &'a Arena,
860 connection: &Connection,
861 key: Option<&Key>,
862) -> Result<(Section<'a>, Option<i64>, BTreeSet<ExGuid>)> {
863 let image = base::base(connection)?;
864 let mut section = open(arena, image, key)?;
865 let mut open = None;
866 let mut touched = BTreeSet::new();
867 for (batch, sealed) in batches(connection)? {
868 if open.is_some() {
869 return Err(io::Error::new(
870 io::ErrorKind::InvalidData,
871 "A cached batch follows the open batch",
872 )
873 .into());
874 }
875 match sealed {
876 Some(Some(transaction)) => section.replay(&transaction)?,
877 Some(None) => {}
878 None => {
879 for queued in queue::load(connection, key, Some(batch))? {
880 section
881 .apply(&queued.author, &queued.edit)
882 .map_err(|error| {
883 io::Error::new(
884 io::ErrorKind::InvalidData,
885 format!("Queued edit {} no longer applies: {error}", queued.id),
886 )
887 })?;
888 touched.extend(queue::spaces(&queued.edit, section.root()));
889 }
890 open = Some(batch);
891 }
892 }
893 }
894 Ok((section, open, touched))
895}
896
897/// The image the queue leaves, its unsealed edits sealed as one more revision; that
898/// revision's identities are fresh on every call. O(section), for tests and recovery.
899pub(crate) fn image(connection: &Connection, key: Option<&Key>) -> Result<Vec<u8>> {
900 let arena = Arena::default();
901 let (mut section, ..) = replay(&arena, connection, key)?;
902 section.seal()?;
903 Ok(section.image())
904}
905
906/// A section image, a password-protected one under `key`.
907pub(crate) fn open<'a>(arena: &'a Arena, image: Vec<u8>, key: Option<&Key>) -> Result<Section<'a>> {
908 Ok(match key {
909 Some(key) => Section::unlock(arena, image, key)?,
910 None => Section::open(arena, image)?,
911 })
912}
913
914/// The sealed batch waiting for publication, if any.
915pub(crate) fn sealed(connection: &Connection) -> Result<Option<Sealed>> {
916 let row: Option<(i64, String)> = connection
917 .query_row(
918 "SELECT id, sealed FROM batches WHERE sealed IS NOT NULL ORDER BY id LIMIT 1",
919 [],
920 |row| Ok((row.get(0)?, row.get(1)?)),
921 )
922 .optional()?;
923 row.map(|(batch, sealed)| {
924 Ok(Sealed {
925 batch,
926 transaction: serde_json::from_str(&sealed)
927 .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?,
928 })
929 })
930 .transpose()
931}