| 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 | |
| 7 | use crate::{Result, base, lock, merge, queue, worker::Signal}; |
| 8 | use onestore::{ |
| 9 | Arena, ExGuid, Section, Transaction, |
| 10 | op::{Edit, OpError}, |
| 11 | page::Page, |
| 12 | protected::Key, |
| 13 | }; |
| 14 | use rusqlite::{Connection, OptionalExtension, TransactionBehavior, params}; |
| 15 | use std::{ |
| 16 | collections::{BTreeMap, BTreeSet, VecDeque}, |
| 17 | io, |
| 18 | sync::{Arc, Mutex, Weak, mpsc}, |
| 19 | }; |
| 20 | |
| 21 | pub(crate) type Reply<T> = Box<dyn FnOnce(Result<T>) + Send>; |
| 22 | |
| 23 | pub(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. |
| 75 | type Takeover = (mpsc::Receiver<Request>, VecDeque<Request>); |
| 76 | |
| 77 | /// A sealed batch: the transaction publishing it, or none when its edits changed nothing. |
| 78 | pub(crate) struct Sealed { |
| 79 | pub batch: i64, |
| 80 | pub transaction: Option<Transaction>, |
| 81 | } |
| 82 | |
| 83 | impl 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. |
| 106 | pub(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 | |
| 118 | impl 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")] |
| 177 | mod 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"))] |
| 209 | async 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")] |
| 216 | async 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. |
| 226 | pub(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 | |
| 251 | enum Next { |
| 252 | Stop, |
| 253 | Reopen, |
| 254 | Handover(mpsc::SyncSender<Takeover>, VecDeque<Request>), |
| 255 | } |
| 256 | |
| 257 | async 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. |
| 291 | fn 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))] |
| 300 | enum 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"))] |
| 313 | fn 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")] |
| 354 | fn 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. |
| 371 | struct Accepted { |
| 372 | author: String, |
| 373 | edit: Edit, |
| 374 | reply: Reply<u64>, |
| 375 | } |
| 376 | |
| 377 | struct 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 | |
| 385 | impl<'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. |
| 708 | fn 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. |
| 839 | fn 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. |
| 858 | fn 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. |
| 899 | pub(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`. |
| 907 | pub(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. |
| 915 | pub(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 | } |