1use super::*;
2use rusqlite::backup::{Backup, StepResult};
3use std::collections::BTreeMap;
4
5const 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)]
9pub 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.
22pub struct Recovery {
23 connection: Connection,
24 /// A password-protected section's key, which opens its sealed queue.
25 key: Option<Key>,
26}
27
28impl 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
106impl 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.
121pub(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
165fn 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
193fn 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}