| 1 | use super::*; |
| 2 | use rusqlite::backup::{Backup, StepResult}; |
| 3 | use std::collections::BTreeMap; |
| 4 | |
| 5 | const RECOVERY_ID: u32 = 0x4f4e4552; |
| 6 | |
| 7 | /// Counts and image sizes without notebook text, paths, authors or credentials. |
| 8 | #[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)] |
| 9 | pub struct RecoverySummary { |
| 10 | pub queued_edits: u64, |
| 11 | pub uncertain_edits: u64, |
| 12 | pub published_receipts: u64, |
| 13 | /// Bytes of the image the queued edits apply to. |
| 14 | pub base_bytes: u64, |
| 15 | /// Bytes of the observed remote image, while it is not the base. |
| 16 | pub remote_bytes: u64, |
| 17 | pub cached_assets: u64, |
| 18 | pub cached_asset_bytes: u64, |
| 19 | } |
| 20 | |
| 21 | /// Read-only recovery evidence; it cannot publish or acknowledge an edit. |
| 22 | pub struct Recovery { |
| 23 | connection: Connection, |
| 24 | /// A password-protected section's key, which opens its sealed queue. |
| 25 | key: Option<Key>, |
| 26 | } |
| 27 | |
| 28 | impl Recovery { |
| 29 | /// Opens an exported archive without migration or conversion into a writable replica. |
| 30 | pub fn open(path: impl AsRef<Path>) -> Result<Self> { |
| 31 | Self::open_with(path.as_ref(), None) |
| 32 | } |
| 33 | |
| 34 | /// `open` for an archive of a password-protected section, under its `key`. |
| 35 | pub fn open_unlocked(path: impl AsRef<Path>, key: &Key) -> Result<Self> { |
| 36 | Self::open_with(path.as_ref(), Some(key.clone())) |
| 37 | } |
| 38 | |
| 39 | fn open_with(path: &Path, key: Option<Key>) -> Result<Self> { |
| 40 | let connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_ONLY)?; |
| 41 | connection.busy_timeout(Duration::ZERO)?; |
| 42 | connection.execute_batch("BEGIN")?; |
| 43 | let application: u32 = |
| 44 | connection.pragma_query_value(None, "application_id", |row| row.get(0))?; |
| 45 | let version: u32 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?; |
| 46 | if application != RECOVERY_ID { |
| 47 | return Err( |
| 48 | io::Error::new(io::ErrorKind::InvalidData, "Not a recovery archive").into(), |
| 49 | ); |
| 50 | } |
| 51 | // A schema-15 archive differs only in columns nothing reads. |
| 52 | if ![schema::PREVIOUS, schema::VERSION].contains(&version) { |
| 53 | return Err(io::Error::new( |
| 54 | io::ErrorKind::InvalidData, |
| 55 | format!( |
| 56 | "Archive schema version {version} is not the supported version {}", |
| 57 | schema::VERSION |
| 58 | ), |
| 59 | ) |
| 60 | .into()); |
| 61 | } |
| 62 | let recovery = Self { connection, key }; |
| 63 | recovery.snapshot()?; |
| 64 | for edit in recovery.pending()? { |
| 65 | recovery.status(edit.id)?; |
| 66 | } |
| 67 | receipts(&recovery.connection)?; |
| 68 | Ok(recovery) |
| 69 | } |
| 70 | |
| 71 | pub fn summary(&self) -> Result<RecoverySummary> { |
| 72 | summary(&self.connection) |
| 73 | } |
| 74 | |
| 75 | /// The image the archived queue leaves, as `Replica::snapshot` gives it. |
| 76 | pub fn snapshot(&self) -> Result<Vec<u8>> { |
| 77 | working::image(&self.connection, self.key.as_ref()) |
| 78 | } |
| 79 | |
| 80 | /// The last observed remote image. |
| 81 | pub fn remote_snapshot(&self) -> Result<Vec<u8>> { |
| 82 | match base::read(&self.connection, base::Image::Remote)? { |
| 83 | Some(image) => Ok(image), |
| 84 | None => base::base(&self.connection), |
| 85 | } |
| 86 | } |
| 87 | |
| 88 | pub fn pending(&self) -> Result<Vec<PendingEdit>> { |
| 89 | pending(&self.connection, self.key.as_ref()) |
| 90 | } |
| 91 | |
| 92 | pub fn status(&self, id: u64) -> Result<Option<EditStatus>> { |
| 93 | sync::status(&self.connection, id, self.key.as_ref()) |
| 94 | } |
| 95 | |
| 96 | pub fn receipts(&self) -> Result<BTreeMap<u64, ExGuid>> { |
| 97 | receipts(&self.connection) |
| 98 | } |
| 99 | |
| 100 | /// Reads a previously downloaded external payload without accessing its former server. |
| 101 | pub fn cached_asset(&self, filename: &str, limit: usize) -> Result<Option<Vec<u8>>> { |
| 102 | assets::cached(&self.connection, &assets::key(filename)?, limit) |
| 103 | } |
| 104 | } |
| 105 | |
| 106 | impl Replica { |
| 107 | /// Summarizes durable state without exposing notebook content. |
| 108 | pub fn recovery_summary(&self) -> Result<RecoverySummary> { |
| 109 | summary(&*self.lock()?) |
| 110 | } |
| 111 | |
| 112 | /// Exports a consistent archive to a new local path without changing the live queue. |
| 113 | /// Archives contain notebook content and cannot be opened as writable replicas. |
| 114 | /// An error after publication can leave a complete archive at the destination. |
| 115 | pub fn export_recovery(&self, path: impl AsRef<Path>) -> Result<()> { |
| 116 | export(&*self.lock()?, path.as_ref(), false) |
| 117 | } |
| 118 | } |
| 119 | |
| 120 | /// Copies the cache to `path` as a recovery archive; `replace` overwrites an existing file. |
| 121 | pub(crate) fn export(source: &Connection, path: &Path, replace: bool) -> Result<()> { |
| 122 | if path.file_name().is_none() { |
| 123 | return Err(io::Error::new( |
| 124 | io::ErrorKind::InvalidInput, |
| 125 | "Provide a new recovery archive filename", |
| 126 | ) |
| 127 | .into()); |
| 128 | } |
| 129 | let parent = path |
| 130 | .parent() |
| 131 | .filter(|parent| !parent.as_os_str().is_empty()) |
| 132 | .unwrap_or_else(|| Path::new(".")); |
| 133 | let temporary = tempfile::Builder::new() |
| 134 | .prefix(".onestore-recovery-") |
| 135 | .tempfile_in(parent)?; |
| 136 | let mut destination = cache_connection(temporary.path())?; |
| 137 | { |
| 138 | let backup = Backup::new(source, &mut destination)?; |
| 139 | if backup.step(-1)? != StepResult::Done { |
| 140 | return Err(io::Error::new( |
| 141 | io::ErrorKind::WouldBlock, |
| 142 | "Recovery export could not acquire the database snapshot", |
| 143 | ) |
| 144 | .into()); |
| 145 | } |
| 146 | } |
| 147 | // One self-contained file: the copied header says WAL, as the cache does. |
| 148 | destination.pragma_update(None, "journal_mode", "DELETE")?; |
| 149 | destination.pragma_update(None, "application_id", RECOVERY_ID)?; |
| 150 | destination.close().map_err(|(_, error)| error)?; |
| 151 | temporary.as_file().sync_all()?; |
| 152 | if replace { |
| 153 | temporary.persist(path).map_err(|error| error.error)?; |
| 154 | } else { |
| 155 | temporary |
| 156 | .persist_noclobber(path) |
| 157 | .map_err(|error| error.error)?; |
| 158 | } |
| 159 | // Windows opens no folder as a file; NTFS journals the rename. |
| 160 | #[cfg(unix)] |
| 161 | std::fs::File::open(parent)?.sync_all()?; |
| 162 | Ok(()) |
| 163 | } |
| 164 | |
| 165 | fn summary(connection: &Connection) -> Result<RecoverySummary> { |
| 166 | let (cached_assets, cached_asset_bytes): (i64, i64) = connection.query_row( |
| 167 | "SELECT count(*),coalesce(sum(length(data)),0) FROM assets", |
| 168 | [], |
| 169 | |row| Ok((row.get(0)?, row.get(1)?)), |
| 170 | )?; |
| 171 | let length = |image| -> Result<u64> { |
| 172 | Ok(base::stamp(connection, image)?.map_or(0, |stamp| stamp.length)) |
| 173 | }; |
| 174 | let (queued_edits, uncertain_edits, published_receipts): (i64, i64, i64) = |
| 175 | connection.query_row( |
| 176 | "SELECT (SELECT count(*) FROM edits), |
| 177 | (SELECT count(*) FROM edits JOIN batches ON batches.id=edits.batch WHERE attempted=1), |
| 178 | (SELECT count(*) FROM receipts)", |
| 179 | [], |
| 180 | |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), |
| 181 | )?; |
| 182 | Ok(RecoverySummary { |
| 183 | queued_edits: unsigned(queued_edits)?, |
| 184 | uncertain_edits: unsigned(uncertain_edits)?, |
| 185 | published_receipts: unsigned(published_receipts)?, |
| 186 | base_bytes: length(base::Image::Base)?, |
| 187 | remote_bytes: length(base::Image::Remote)?, |
| 188 | cached_assets: unsigned(cached_assets)?, |
| 189 | cached_asset_bytes: unsigned(cached_asset_bytes)?, |
| 190 | }) |
| 191 | } |
| 192 | |
| 193 | fn receipts(connection: &Connection) -> Result<BTreeMap<u64, ExGuid>> { |
| 194 | let mut statement = |
| 195 | connection.prepare("SELECT edit_id, revision FROM receipts ORDER BY edit_id")?; |
| 196 | let mut rows = statement.query([])?; |
| 197 | let mut receipts = BTreeMap::new(); |
| 198 | while let Some(row) = rows.next()? { |
| 199 | receipts.insert(unsigned(row.get(0)?)?, row.get::<_, String>(1)?.parse()?); |
| 200 | } |
| 201 | Ok(receipts) |
| 202 | } |