| 1 | use super::*; |
| 2 | use crate::{ |
| 3 | resolve::{Merged, Version}, |
| 4 | working::{Request, Sealed}, |
| 5 | }; |
| 6 | use onestore::{CommitError, CommitState, RevisionIndex, Stamp, Store, Transaction}; |
| 7 | use rusqlite::OptionalExtension; |
| 8 | use std::{ |
| 9 | collections::BTreeMap, |
| 10 | sync::{MutexGuard, TryLockError}, |
| 11 | }; |
| 12 | |
| 13 | /// A single remote file with fresh reads and native-compatible guarded publication. |
| 14 | /// Errors retain publication state; confirmation checks the stamp, flushes, and notifies |
| 15 | /// cached readers. |
| 16 | pub trait Remote { |
| 17 | /// Completes a non-blocking operation that returned `WouldBlock`; retrying that |
| 18 | /// operation then takes its result. Blocking providers leave this absent. |
| 19 | fn pending(&mut self) -> Option<std::pin::Pin<Box<dyn Future<Output = ()> + '_>>> { |
| 20 | None |
| 21 | } |
| 22 | fn read(&mut self) -> io::Result<Vec<u8>>; |
| 23 | /// The file's stamp without reading its body or coordinating with writers. While it is |
| 24 | /// the last observed image's, synchronization neither reads nor revalidates the file. |
| 25 | fn stamp(&mut self) -> io::Result<Stamp>; |
| 26 | fn publish(&mut self, transaction: &Transaction) -> std::result::Result<(), CommitError>; |
| 27 | fn accepts_edits(&self) -> bool { |
| 28 | false |
| 29 | } |
| 30 | fn publish_edits( |
| 31 | &mut self, |
| 32 | transaction: &Transaction, |
| 33 | edits: &[PendingEdit], |
| 34 | revisions: &BTreeMap<ExGuid, ExGuid>, |
| 35 | ) -> std::result::Result<(), CommitError> { |
| 36 | let _ = (edits, revisions); |
| 37 | self.publish(transaction) |
| 38 | } |
| 39 | /// Confirms that the file still has `base`'s stamp and is durable (`onestore::confirm`). |
| 40 | fn confirm(&mut self, base: &Stamp) -> std::result::Result<(), CommitError>; |
| 41 | /// The versions a file provider keeps beside the file, as iCloud Drive keeps the commits |
| 42 | /// that lost to another device's (unresolved conflict versions); none by default. |
| 43 | /// Synchronization merges each into the file, then retires it. |
| 44 | fn versions(&mut self) -> io::Result<Vec<Version>> { |
| 45 | Ok(Vec::new()) |
| 46 | } |
| 47 | /// A version's image. |
| 48 | fn version(&mut self, id: &str) -> io::Result<Vec<u8>> { |
| 49 | Err(io::Error::new( |
| 50 | io::ErrorKind::NotFound, |
| 51 | format!("No version {id}"), |
| 52 | )) |
| 53 | } |
| 54 | /// Retires a version the file now holds everything of; with `keep`, one that is another |
| 55 | /// section or cannot be read, kept first as a file of its own beside this one. |
| 56 | fn retire(&mut self, id: &str, keep: bool) -> io::Result<()> { |
| 57 | let _ = (id, keep); |
| 58 | Ok(()) |
| 59 | } |
| 60 | } |
| 61 | |
| 62 | pub(crate) trait Waiting { |
| 63 | fn waiting(&self) -> bool; |
| 64 | } |
| 65 | |
| 66 | impl Waiting for io::Error { |
| 67 | fn waiting(&self) -> bool { |
| 68 | self.kind() == io::ErrorKind::WouldBlock |
| 69 | } |
| 70 | } |
| 71 | |
| 72 | impl Waiting for CommitError { |
| 73 | fn waiting(&self) -> bool { |
| 74 | self.error.kind() == io::ErrorKind::WouldBlock |
| 75 | } |
| 76 | } |
| 77 | |
| 78 | pub(crate) async fn awaited<R: Remote, T, E: Waiting>( |
| 79 | remote: &mut R, |
| 80 | mut operation: impl FnMut(&mut R) -> std::result::Result<T, E>, |
| 81 | ) -> std::result::Result<T, E> { |
| 82 | loop { |
| 83 | let result = operation(remote); |
| 84 | if result.as_ref().is_err_and(Waiting::waiting) |
| 85 | && let Some(pending) = remote.pending() |
| 86 | { |
| 87 | pending.await; |
| 88 | continue; |
| 89 | } |
| 90 | return result; |
| 91 | } |
| 92 | } |
| 93 | |
| 94 | #[derive(Debug, Clone, PartialEq, Eq)] |
| 95 | pub enum EditStatus { |
| 96 | Pending, |
| 97 | /// Retained publication attempt; this revision alone may be insufficient to confirm it. |
| 98 | AwaitingConfirmation { |
| 99 | revision: ExGuid, |
| 100 | }, |
| 101 | /// Historical confirmation; later remote edits or restores may remove the effect. |
| 102 | Published { |
| 103 | revision: ExGuid, |
| 104 | }, |
| 105 | /// Retired unpublished by a reviewed release; the archive at this path holds the |
| 106 | /// edit, its attempt evidence and both images. |
| 107 | Archived { |
| 108 | archive: String, |
| 109 | }, |
| 110 | } |
| 111 | |
| 112 | /// What one synchronization step did: the state its batch reached, named by the batch's |
| 113 | /// newest edit, and the pages a remote change replaced. |
| 114 | #[derive(Debug, Clone, Default, PartialEq, Eq)] |
| 115 | pub struct Synced { |
| 116 | pub edit: Option<(u64, EditStatus)>, |
| 117 | pub changed: Vec<ExGuid>, |
| 118 | } |
| 119 | |
| 120 | /// The queue as a synchronization step starts it. |
| 121 | struct State { |
| 122 | base: Stamp, |
| 123 | /// The observed remote image's stamp while it is not the base. |
| 124 | remote: Option<Stamp>, |
| 125 | /// The first batch waiting on a remote answer to an attempt. |
| 126 | blocked: Option<Blocked>, |
| 127 | queued: bool, |
| 128 | } |
| 129 | |
| 130 | struct Blocked { |
| 131 | batch: i64, |
| 132 | revisions: Option<BTreeMap<ExGuid, ExGuid>>, |
| 133 | } |
| 134 | |
| 135 | fn state(connection: &Connection) -> Result<State> { |
| 136 | let base = base::base_stamp(connection)?; |
| 137 | let blocked = connection |
| 138 | .query_row( |
| 139 | "SELECT id, revisions FROM batches WHERE attempted=1 ORDER BY id LIMIT 1", |
| 140 | [], |
| 141 | |row| Ok((row.get::<_, i64>(0)?, row.get::<_, Option<String>>(1)?)), |
| 142 | ) |
| 143 | .optional()? |
| 144 | .map(|(batch, revisions)| { |
| 145 | Ok::<_, Error>(Blocked { |
| 146 | batch, |
| 147 | revisions: revisions |
| 148 | .map(|revisions| decode_revisions(&revisions)) |
| 149 | .transpose()?, |
| 150 | }) |
| 151 | }) |
| 152 | .transpose()?; |
| 153 | Ok(State { |
| 154 | base, |
| 155 | remote: base::stamp(connection, base::Image::Remote)?, |
| 156 | blocked, |
| 157 | queued: connection |
| 158 | .query_row("SELECT EXISTS(SELECT 1 FROM batches)", [], |row| row.get(0))?, |
| 159 | }) |
| 160 | } |
| 161 | |
| 162 | fn decode_revisions(encoded: &str) -> Result<BTreeMap<ExGuid, ExGuid>> { |
| 163 | let revisions: BTreeMap<ExGuid, ExGuid> = serde_json::from_str(encoded) |
| 164 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?; |
| 165 | if revisions |
| 166 | .iter() |
| 167 | .any(|(space, revision)| space.guid == [0; 16] || revision.guid == [0; 16]) |
| 168 | { |
| 169 | return Err(io::Error::new( |
| 170 | io::ErrorKind::InvalidData, |
| 171 | "Cached publication evidence is incomplete", |
| 172 | ) |
| 173 | .into()); |
| 174 | } |
| 175 | Ok(revisions) |
| 176 | } |
| 177 | |
| 178 | /// The revision a batch's seal stored for the first page an edit changes. |
| 179 | fn revision_of(edit: &onestore::op::Edit, revisions: &BTreeMap<ExGuid, ExGuid>) -> Result<ExGuid> { |
| 180 | edit.ops |
| 181 | .iter() |
| 182 | .find_map(|op| match op { |
| 183 | onestore::op::Op::Page { space, .. } => revisions.get(space), |
| 184 | onestore::op::Op::Section(_) => None, |
| 185 | }) |
| 186 | .or_else(|| revisions.values().next()) |
| 187 | .copied() |
| 188 | .ok_or_else(|| { |
| 189 | io::Error::new( |
| 190 | io::ErrorKind::InvalidData, |
| 191 | "Cached publication evidence is incomplete", |
| 192 | ) |
| 193 | .into() |
| 194 | }) |
| 195 | } |
| 196 | |
| 197 | /// The newest edit of a batch. |
| 198 | fn newest(connection: &Connection, batch: i64) -> Result<u64> { |
| 199 | unsigned( |
| 200 | connection.query_row("SELECT max(id) FROM edits WHERE batch=?1", [batch], |row| { |
| 201 | row.get(0) |
| 202 | })?, |
| 203 | ) |
| 204 | } |
| 205 | |
| 206 | impl Replica { |
| 207 | /// Returns a durable receipt or the persisted state of a locally acknowledged edit. |
| 208 | pub fn status(&self, id: u64) -> Result<Option<EditStatus>> { |
| 209 | status(&*self.lock()?, id, self.section.key.as_ref()) |
| 210 | } |
| 211 | |
| 212 | /// The last observed remote image. Observation alone does not acknowledge any pending |
| 213 | /// edit's remote durability. |
| 214 | pub(crate) fn remote_snapshot(&self) -> Result<Vec<u8>> { |
| 215 | let connection = self.lock()?; |
| 216 | match base::read(&connection, base::Image::Remote)? { |
| 217 | Some(image) => Ok(image), |
| 218 | None => base::base(&connection), |
| 219 | } |
| 220 | } |
| 221 | |
| 222 | /// Publishes the oldest unpublished batch, rebasing the queue first when the remote |
| 223 | /// changed. Reads the remote image only when its stamp moved; network I/O holds |
| 224 | /// synchronization ownership without holding the cache mutex. Uncertain attempts are |
| 225 | /// never replayed. |
| 226 | /// Versions the remote keeps beside the file merge into it first, each published as one |
| 227 | /// more revision and then retired (`resolve.rs`). |
| 228 | pub fn sync_once(&self, remote: &mut impl Remote) -> Result<Synced> { |
| 229 | crate::task::ready(self.sync_once_async(remote))? |
| 230 | } |
| 231 | |
| 232 | /// The same guarded synchronization step, awaiting non-blocking remote operations. |
| 233 | pub async fn sync_once_async(&self, remote: &mut impl Remote) -> Result<Synced> { |
| 234 | let result = self.sync_once_inner(remote).await; |
| 235 | let durable = crate::fs::durable().await; |
| 236 | let synced = result?; |
| 237 | durable?; |
| 238 | Ok(synced) |
| 239 | } |
| 240 | |
| 241 | /// Ownership uses `try_lock`; the cache mutex is released before every await. |
| 242 | #[allow(clippy::await_holding_lock)] |
| 243 | async fn sync_once_inner(&self, remote: &mut impl Remote) -> Result<Synced> { |
| 244 | let _owner = self.sync_owner()?; |
| 245 | for version in awaited(remote, Remote::versions) |
| 246 | .await |
| 247 | .map_err(Error::RemoteIo)? |
| 248 | { |
| 249 | let image = awaited(remote, |remote| remote.version(&version.id)) |
| 250 | .await |
| 251 | .map_err(Error::RemoteIo)?; |
| 252 | let current = awaited(remote, Remote::read) |
| 253 | .await |
| 254 | .map_err(Error::RemoteIo)?; |
| 255 | let device = version.device.as_deref().unwrap_or("Another device"); |
| 256 | let keep = match crate::resolve::merge(&current, &image, device) { |
| 257 | Ok(Merged::Held) => false, |
| 258 | Ok(Merged::Publish(transaction)) => { |
| 259 | awaited(remote, |remote| remote.publish(&transaction)).await?; |
| 260 | false |
| 261 | } |
| 262 | // A version this cannot merge is kept whole rather than lost. |
| 263 | Ok(Merged::Foreign) | Err(Error::Document(_) | Error::Rejected(_)) => true, |
| 264 | Err(error) => return Err(error), |
| 265 | }; |
| 266 | awaited(remote, |remote| remote.retire(&version.id, keep)) |
| 267 | .await |
| 268 | .map_err(Error::RemoteIo)?; |
| 269 | } |
| 270 | let state = state(&*self.lock()?)?; |
| 271 | let batched = self.section.key.is_none() |
| 272 | && state.queued |
| 273 | && state.blocked.is_none() |
| 274 | && remote.accepts_edits(); |
| 275 | let observed = awaited(remote, Remote::stamp) |
| 276 | .await |
| 277 | .map_err(Error::RemoteIo)?; |
| 278 | if let Some(blocked) = &state.blocked |
| 279 | && observed == *state.remote.as_ref().unwrap_or(&state.base) |
| 280 | { |
| 281 | // Nothing new to decide: the remote is as it was when the batch blocked. |
| 282 | let id = newest(&*self.lock()?, blocked.batch)?; |
| 283 | return Ok(Synced { |
| 284 | edit: Some(( |
| 285 | id, |
| 286 | status(&*self.lock()?, id, self.section.key.as_ref())? |
| 287 | .unwrap_or(EditStatus::Pending), |
| 288 | )), |
| 289 | changed: Vec::new(), |
| 290 | }); |
| 291 | } |
| 292 | let image = match observed { |
| 293 | observed if observed == state.base || batched => None, |
| 294 | _ => { |
| 295 | let image = awaited(remote, Remote::read) |
| 296 | .await |
| 297 | .map_err(Error::RemoteIo)?; |
| 298 | // Protected elsewhere: a section written anew, which only its key reads. |
| 299 | if self.section.key.is_none() && crate::discover::locked(&Store::parse(&image)?) { |
| 300 | return Err(Error::RemoteIo(io::Error::new( |
| 301 | io::ErrorKind::Unsupported, |
| 302 | "Password protected", |
| 303 | ))); |
| 304 | } |
| 305 | // A read image is compared whole: a stamp stands for it only when read alone. |
| 306 | (base::base(&*self.lock()?)? != image).then_some(image) |
| 307 | } |
| 308 | }; |
| 309 | let mut changed = Vec::new(); |
| 310 | if let Some(image) = image { |
| 311 | if let Some(Blocked { batch, revisions }) = &state.blocked { |
| 312 | let id = newest(&*self.lock()?, *batch)?; |
| 313 | let revisions = revisions.clone().unwrap_or_default(); |
| 314 | let store = Store::parse(&image)?; |
| 315 | if !store.checksum_mismatches.is_empty() { |
| 316 | return Err(io::Error::new( |
| 317 | io::ErrorKind::InvalidData, |
| 318 | "Notebook transaction checksum damage", |
| 319 | ) |
| 320 | .into()); |
| 321 | } |
| 322 | let index = RevisionIndex::parse(&store)?; |
| 323 | index.validate_current()?; |
| 324 | if index.root != self.root { |
| 325 | return Err(io::Error::new( |
| 326 | io::ErrorKind::InvalidInput, |
| 327 | "Remote snapshot belongs to another document", |
| 328 | ) |
| 329 | .into()); |
| 330 | } |
| 331 | let observed = revisions.iter().all(|(space, revision)| { |
| 332 | index |
| 333 | .spaces |
| 334 | .get(space) |
| 335 | .is_some_and(|space| space.revisions.contains_key(revision)) |
| 336 | }); |
| 337 | let sealed = working::sealed(&*self.lock()?)? |
| 338 | .filter(|sealed| sealed.batch == *batch) |
| 339 | .and_then(|sealed| sealed.transaction); |
| 340 | // Without its revisions the attempt still counts once the remote holds |
| 341 | // every page it changed as it changed them; its receipts then name the |
| 342 | // remote's revisions. |
| 343 | let receipts = if observed { |
| 344 | None |
| 345 | } else { |
| 346 | let equal = self.holds(&image, sealed.as_ref(), &revisions)?; |
| 347 | if !equal { |
| 348 | base::write(&*self.lock()?, base::Image::Remote, &image)?; |
| 349 | return Ok(Synced { |
| 350 | edit: Some(( |
| 351 | id, |
| 352 | status(&*self.lock()?, id, self.section.key.as_ref())? |
| 353 | .unwrap_or(EditStatus::Pending), |
| 354 | )), |
| 355 | changed, |
| 356 | }); |
| 357 | } |
| 358 | Some( |
| 359 | revisions |
| 360 | .keys() |
| 361 | .filter_map(|space| Some((*space, index.active(*space).ok()?))) |
| 362 | .collect(), |
| 363 | ) |
| 364 | }; |
| 365 | let stamp = Stamp::of(&image)?; |
| 366 | if let Err(error) = awaited(remote, |remote| remote.confirm(&stamp)).await { |
| 367 | if error.state == CommitState::Committed { |
| 368 | self.acknowledge(*batch, sealed.as_ref(), receipts.as_ref())?; |
| 369 | } |
| 370 | return Err(error.into()); |
| 371 | } |
| 372 | self.acknowledge(*batch, sealed.as_ref(), receipts.as_ref())?; |
| 373 | let revision = self.receipt(id)?; |
| 374 | changed = self.rebase(Some(image))?; |
| 375 | return Ok(Synced { |
| 376 | edit: Some((id, EditStatus::Published { revision })), |
| 377 | changed, |
| 378 | }); |
| 379 | } |
| 380 | changed = self.rebase(Some(image))?; |
| 381 | } else if let Some(blocked) = &state.blocked { |
| 382 | if state.remote.is_some() { |
| 383 | // The remote is back at the base: what blocked the queue is gone. |
| 384 | base::clear(&*self.lock()?, base::Image::Remote)?; |
| 385 | } |
| 386 | let id = newest(&*self.lock()?, blocked.batch)?; |
| 387 | return Ok(Synced { |
| 388 | edit: Some(( |
| 389 | id, |
| 390 | status(&*self.lock()?, id, self.section.key.as_ref())? |
| 391 | .unwrap_or(EditStatus::Pending), |
| 392 | )), |
| 393 | changed, |
| 394 | }); |
| 395 | } else if !state.queued { |
| 396 | return Ok(Synced::default()); |
| 397 | } |
| 398 | // The cache lock is released before the section thread, which takes it, seals. |
| 399 | let waiting = working::sealed(&*self.lock()?)?; |
| 400 | let sealed = match waiting { |
| 401 | Some(sealed) => { |
| 402 | if sealed.transaction.is_some() { |
| 403 | self.lock()? |
| 404 | .execute("UPDATE batches SET attempted=1 WHERE id=?1", [sealed.batch])?; |
| 405 | } |
| 406 | Some(sealed) |
| 407 | } |
| 408 | None => self.ask(|reply| Request::Seal { reply })?, |
| 409 | }; |
| 410 | let Some(Sealed { batch, transaction }) = sealed else { |
| 411 | return Ok(Synced { |
| 412 | edit: None, |
| 413 | changed, |
| 414 | }); |
| 415 | }; |
| 416 | let id = newest(&*self.lock()?, batch)?; |
| 417 | let Some(transaction) = transaction else { |
| 418 | // Edits that changed nothing are published once the remote's image is durable. |
| 419 | let base = base::base_stamp(&*self.lock()?)?; |
| 420 | if let Err(error) = awaited(remote, |remote| remote.confirm(&base)).await { |
| 421 | if error.state == CommitState::Committed { |
| 422 | self.acknowledge(batch, None, None)?; |
| 423 | } |
| 424 | return Err(error.into()); |
| 425 | } |
| 426 | self.acknowledge(batch, None, None)?; |
| 427 | let revision = self.receipt(id)?; |
| 428 | return Ok(Synced { |
| 429 | edit: Some((id, EditStatus::Published { revision })), |
| 430 | changed, |
| 431 | }); |
| 432 | }; |
| 433 | if let Err(error) = crate::fs::durable().await { |
| 434 | self.lock()? |
| 435 | .execute("UPDATE batches SET attempted=0 WHERE id=?1", [batch])?; |
| 436 | return Err(error.into()); |
| 437 | } |
| 438 | let published = if self.section.key.is_none() && remote.accepts_edits() { |
| 439 | let (edits, revisions) = { |
| 440 | let connection = self.lock()?; |
| 441 | let edits = queue::load(&connection, None, Some(batch))?; |
| 442 | let revisions = decode_revisions(&connection.query_row( |
| 443 | "SELECT revisions FROM batches WHERE id=?1", |
| 444 | [batch], |
| 445 | |row| row.get::<_, String>(0), |
| 446 | )?)?; |
| 447 | (edits, revisions) |
| 448 | }; |
| 449 | awaited(remote, |remote| { |
| 450 | remote.publish_edits(&transaction, &edits, &revisions) |
| 451 | }) |
| 452 | .await |
| 453 | } else { |
| 454 | awaited(remote, |remote| remote.publish(&transaction)).await |
| 455 | }; |
| 456 | match published { |
| 457 | Ok(()) => {} |
| 458 | Err(error) if error.state == CommitState::NotCommitted => { |
| 459 | self.lock()? |
| 460 | .execute("UPDATE batches SET attempted=0 WHERE id=?1", [batch])?; |
| 461 | return Err(error.into()); |
| 462 | } |
| 463 | Err(error) if error.state == CommitState::Committed => { |
| 464 | self.acknowledge(batch, Some(&transaction), None)?; |
| 465 | return Err(error.into()); |
| 466 | } |
| 467 | Err(error) => return Err(error.into()), |
| 468 | } |
| 469 | self.acknowledge(batch, Some(&transaction), None)?; |
| 470 | let revision = self.receipt(id)?; |
| 471 | if batched { |
| 472 | changed.extend( |
| 473 | self.rebase(Some( |
| 474 | awaited(remote, Remote::read) |
| 475 | .await |
| 476 | .map_err(Error::RemoteIo)?, |
| 477 | ))?, |
| 478 | ); |
| 479 | changed.sort(); |
| 480 | changed.dedup(); |
| 481 | } |
| 482 | Ok(Synced { |
| 483 | edit: Some((id, EditStatus::Published { revision })), |
| 484 | changed, |
| 485 | }) |
| 486 | } |
| 487 | |
| 488 | /// Whether `image` holds every page in `revisions` as the base with `sealed` has it. |
| 489 | fn holds( |
| 490 | &self, |
| 491 | image: &[u8], |
| 492 | sealed: Option<&Transaction>, |
| 493 | revisions: &BTreeMap<ExGuid, ExGuid>, |
| 494 | ) -> Result<bool> { |
| 495 | let base = base::base(&*self.lock()?)?; |
| 496 | let (local, remote) = (onestore::Arena::default(), onestore::Arena::default()); |
| 497 | let mut local = working::open(&local, base, self.section.key.as_ref())?; |
| 498 | if let Some(sealed) = sealed { |
| 499 | local.replay(sealed)?; |
| 500 | } |
| 501 | let mut remote = working::open(&remote, image.to_vec(), self.section.key.as_ref())?; |
| 502 | // A page's versions change without its page changing. |
| 503 | let versions = |section: &mut onestore::Section<'_>| { |
| 504 | section.versions().map(|versions| { |
| 505 | versions |
| 506 | .into_iter() |
| 507 | .filter(|(page, _)| revisions.contains_key(page)) |
| 508 | .collect::<Vec<_>>() |
| 509 | }) |
| 510 | }; |
| 511 | Ok(revisions.keys().all( |
| 512 | |space| matches!((local.page(*space), remote.page(*space)), (Ok(a), Ok(b)) if a == b), |
| 513 | ) && versions(&mut local)? == versions(&mut remote)?) |
| 514 | } |
| 515 | |
| 516 | /// Asks the section thread to replay the queue on `image`, returning the pages the |
| 517 | /// remote changed. |
| 518 | fn rebase(&self, image: Option<Vec<u8>>) -> Result<Vec<ExGuid>> { |
| 519 | self.ask(|reply| Request::Rebase { image, reply }) |
| 520 | } |
| 521 | |
| 522 | fn receipt(&self, id: u64) -> Result<ExGuid> { |
| 523 | match status(&*self.lock()?, id, self.section.key.as_ref())? { |
| 524 | Some(EditStatus::Published { revision }) => Ok(revision), |
| 525 | _ => Err(io::Error::other("A published edit has no receipt").into()), |
| 526 | } |
| 527 | } |
| 528 | |
| 529 | /// Records a published batch: the base becomes the image it left, each of its edits |
| 530 | /// gets a receipt naming the revisions its seal stored, or `revisions`, and it leaves |
| 531 | /// the queue. |
| 532 | fn acknowledge( |
| 533 | &self, |
| 534 | batch: i64, |
| 535 | transaction: Option<&Transaction>, |
| 536 | revisions: Option<&BTreeMap<ExGuid, ExGuid>>, |
| 537 | ) -> Result<()> { |
| 538 | let mut connection = self.lock()?; |
| 539 | let database = connection.transaction_with_behavior(TransactionBehavior::Immediate)?; |
| 540 | if let Some(transaction) = transaction { |
| 541 | base::publish(&database, transaction)?; |
| 542 | } |
| 543 | let revisions = match revisions { |
| 544 | Some(revisions) => revisions.clone(), |
| 545 | None => decode_revisions(&database.query_row( |
| 546 | "SELECT revisions FROM batches WHERE id=?1", |
| 547 | [batch], |
| 548 | |row| row.get::<_, String>(0), |
| 549 | )?)?, |
| 550 | }; |
| 551 | for queued in queue::load(&database, self.section.key.as_ref(), Some(batch))? { |
| 552 | database.execute( |
| 553 | "INSERT INTO receipts(edit_id, revision) VALUES (?1, ?2)", |
| 554 | params![ |
| 555 | signed(queued.id)?, |
| 556 | revision_of(&queued.edit, &revisions)?.to_string() |
| 557 | ], |
| 558 | )?; |
| 559 | } |
| 560 | database.execute("DELETE FROM batches WHERE id=?1", [batch])?; |
| 561 | queue::collect(&database)?; |
| 562 | database.commit()?; |
| 563 | Ok(()) |
| 564 | } |
| 565 | |
| 566 | /// Retires the uncertain attempt of the batch holding edit `id` after review. The |
| 567 | /// queue is first exported to `archive` (a new file), which is the record: no receipt |
| 568 | /// is written. `Mine` publishes the batch again against the current remote; `Theirs` |
| 569 | /// abandons every unpublished edit, which become `Archived`, and the local pages return |
| 570 | /// to the remote's. |
| 571 | pub fn release(&self, id: u64, archive: &Path, resolution: Resolution) -> Result<()> { |
| 572 | let owner = self.sync_owner()?; |
| 573 | let batch: Option<i64> = self |
| 574 | .lock()? |
| 575 | .query_row( |
| 576 | "SELECT id FROM batches WHERE attempted=1 AND id=(SELECT batch FROM edits WHERE id=?1)", |
| 577 | [signed(id)?], |
| 578 | |row| row.get(0), |
| 579 | ) |
| 580 | .optional()?; |
| 581 | let Some(batch) = batch else { |
| 582 | return Err(io::Error::new( |
| 583 | io::ErrorKind::InvalidInput, |
| 584 | "The edit has no uncertain attempt", |
| 585 | ) |
| 586 | .into()); |
| 587 | }; |
| 588 | self.export_recovery(archive)?; |
| 589 | { |
| 590 | let mut connection = self.lock()?; |
| 591 | let transaction = |
| 592 | connection.transaction_with_behavior(TransactionBehavior::Immediate)?; |
| 593 | match resolution { |
| 594 | Resolution::Mine => { |
| 595 | transaction.execute("UPDATE batches SET attempted=0 WHERE id=?1", [batch])?; |
| 596 | } |
| 597 | Resolution::Theirs => { |
| 598 | transaction.execute( |
| 599 | "INSERT INTO archived(edit_id, archive) SELECT id, ?1 FROM edits", |
| 600 | [archive.to_string_lossy().into_owned()], |
| 601 | )?; |
| 602 | transaction.execute("DELETE FROM batches", [])?; |
| 603 | if let Some(remote) = base::read(&transaction, base::Image::Remote)? { |
| 604 | base::write(&transaction, base::Image::Base, &remote)?; |
| 605 | base::clear(&transaction, base::Image::Remote)?; |
| 606 | } |
| 607 | queue::collect(&transaction)?; |
| 608 | } |
| 609 | } |
| 610 | transaction.commit()?; |
| 611 | } |
| 612 | if resolution == Resolution::Theirs { |
| 613 | self.ask(|reply| Request::Reopen { reply })?; |
| 614 | } |
| 615 | drop(owner); |
| 616 | self.wake_sync(); |
| 617 | Ok(()) |
| 618 | } |
| 619 | |
| 620 | fn sync_owner(&self) -> Result<MutexGuard<'_, ()>> { |
| 621 | Ok(self |
| 622 | .synchronization |
| 623 | .try_lock() |
| 624 | .map_err(|error| match error { |
| 625 | TryLockError::WouldBlock => io::Error::from(io::ErrorKind::WouldBlock), |
| 626 | TryLockError::Poisoned(_) => { |
| 627 | io::Error::other("Synchronization owner panicked; reopen the cache") |
| 628 | } |
| 629 | })?) |
| 630 | } |
| 631 | |
| 632 | /// Whether nothing is queued and the remote still has the base's stamp, or the queue |
| 633 | /// is blocked on a remote that has not changed since; reads neither image and does not |
| 634 | /// lock the cache during remote I/O. |
| 635 | pub(crate) async fn settled(&self, remote: &mut impl Remote) -> Result<bool> { |
| 636 | let state = state(&*self.lock()?)?; |
| 637 | let expected = match &state.blocked { |
| 638 | Some(_) => state.remote.unwrap_or(state.base), |
| 639 | None if !state.queued => state.base, |
| 640 | None => return Ok(false), |
| 641 | }; |
| 642 | Ok(awaited(remote, Remote::stamp) |
| 643 | .await |
| 644 | .map_err(Error::RemoteIo)? |
| 645 | == expected |
| 646 | && awaited(remote, Remote::versions) |
| 647 | .await |
| 648 | .map_err(Error::RemoteIo)? |
| 649 | .is_empty()) |
| 650 | } |
| 651 | } |
| 652 | |
| 653 | pub(crate) fn status( |
| 654 | connection: &Connection, |
| 655 | id: u64, |
| 656 | key: Option<&onestore::protected::Key>, |
| 657 | ) -> Result<Option<EditStatus>> { |
| 658 | let id = signed(id)?; |
| 659 | if let Some(revision) = connection |
| 660 | .query_row( |
| 661 | "SELECT revision FROM receipts WHERE edit_id=?1", |
| 662 | [id], |
| 663 | |row| row.get::<_, String>(0), |
| 664 | ) |
| 665 | .optional()? |
| 666 | { |
| 667 | return Ok(Some(EditStatus::Published { |
| 668 | revision: revision.parse()?, |
| 669 | })); |
| 670 | } |
| 671 | if let Some(archive) = connection |
| 672 | .query_row( |
| 673 | "SELECT archive FROM archived WHERE edit_id=?1", |
| 674 | [id], |
| 675 | |row| row.get::<_, String>(0), |
| 676 | ) |
| 677 | .optional()? |
| 678 | { |
| 679 | return Ok(Some(EditStatus::Archived { archive })); |
| 680 | } |
| 681 | let record: Option<(bool, Option<String>, String)> = connection |
| 682 | .query_row( |
| 683 | "SELECT batches.attempted, batches.revisions, edits.edit |
| 684 | FROM edits JOIN batches ON batches.id=edits.batch WHERE edits.id=?1", |
| 685 | [id], |
| 686 | |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), |
| 687 | ) |
| 688 | .optional()?; |
| 689 | Ok(match record { |
| 690 | None => None, |
| 691 | Some((true, revisions, edit)) => { |
| 692 | let revisions = decode_revisions(revisions.as_deref().unwrap_or("{}"))?; |
| 693 | let edit = queue::parse(key, &edit)?; |
| 694 | let revision = revision_of(&edit, &revisions)?; |
| 695 | Some(EditStatus::AwaitingConfirmation { revision }) |
| 696 | } |
| 697 | Some((false, ..)) => Some(EditStatus::Pending), |
| 698 | }) |
| 699 | } |