1//! Section images kept as 16 KiB chunk rows, so a publication rewrites the chunks its
2//! transaction touches instead of the image: `base`, the image the queued edits apply to,
3//! and `remote`, the last observed remote image while it is not the base.
4
5use crate::Result;
6use onestore::{Stamp, Transaction};
7use rusqlite::{Connection, OptionalExtension, params};
8use std::io;
9
10const CHUNK: usize = 16 * 1024;
11
12#[derive(Clone, Copy)]
13pub(crate) enum Image {
14 Base,
15 Remote,
16}
17
18impl Image {
19 fn table(self) -> &'static str {
20 match self {
21 Self::Base => "base",
22 Self::Remote => "remote",
23 }
24 }
25}
26
27fn damaged() -> crate::Error {
28 io::Error::new(io::ErrorKind::InvalidData, "Damaged cached image").into()
29}
30
31/// The base image, which every cache holds.
32pub(crate) fn base(connection: &Connection) -> Result<Vec<u8>> {
33 read(connection, Image::Base)?.ok_or_else(missing)
34}
35
36/// The base image's stamp.
37pub(crate) fn base_stamp(connection: &Connection) -> Result<Stamp> {
38 stamp(connection, Image::Base)?.ok_or_else(missing)
39}
40
41fn missing() -> crate::Error {
42 io::Error::new(io::ErrorKind::InvalidData, "The cache has no base image").into()
43}
44
45/// The whole image, or `None` when none is stored.
46pub(crate) fn read(connection: &Connection, image: Image) -> Result<Option<Vec<u8>>> {
47 let mut query = connection.prepare_cached(&format!(
48 "SELECT chunk, bytes FROM {} ORDER BY chunk",
49 image.table()
50 ))?;
51 let mut rows = query.query([])?;
52 let mut bytes = Vec::new();
53 let mut count = 0;
54 while let Some(row) = rows.next()? {
55 let chunk: i64 = row.get(0)?;
56 let part = row.get_ref(1)?.as_blob().map_err(|_| damaged())?;
57 if chunk != count || bytes.len() % CHUNK != 0 || part.is_empty() || part.len() > CHUNK {
58 return Err(damaged());
59 }
60 bytes.extend_from_slice(part);
61 count += 1;
62 }
63 Ok((count > 0).then_some(bytes))
64}
65
66/// The stored image's header and length, reading its first and last chunks.
67pub(crate) fn stamp(connection: &Connection, image: Image) -> Result<Option<Stamp>> {
68 let table = image.table();
69 let Some((header, last)): Option<(Vec<u8>, i64)> = connection
70 .query_row(
71 &format!("SELECT substr(bytes, 1, 1024), (SELECT max(chunk) FROM {table}) FROM {table} WHERE chunk=0"),
72 [],
73 |row| Ok((row.get(0)?, row.get(1)?)),
74 )
75 .optional()?
76 else {
77 return Ok(None);
78 };
79 let tail: i64 = connection.query_row(
80 &format!("SELECT length(bytes) FROM {table} WHERE chunk=?1"),
81 [last],
82 |row| row.get(0),
83 )?;
84 if !(1..=CHUNK as i64).contains(&tail) {
85 return Err(damaged());
86 }
87 Ok(Some(Stamp {
88 header: header.try_into().map_err(|_| damaged())?,
89 length: u64::try_from(last)
90 .map_err(|_| damaged())?
91 .checked_mul(CHUNK as u64)
92 .and_then(|length| length.checked_add(tail as u64))
93 .ok_or_else(damaged)?,
94 }))
95}
96
97/// Stores `bytes` as the image, rewriting only the chunks that differ.
98pub(crate) fn write(connection: &Connection, image: Image, bytes: &[u8]) -> Result<()> {
99 let table = image.table();
100 let count = bytes.len().div_ceil(CHUNK);
101 connection.execute(
102 &format!("DELETE FROM {table} WHERE chunk>=?1"),
103 [count as i64],
104 )?;
105 let mut stored =
106 connection.prepare_cached(&format!("SELECT bytes FROM {table} WHERE chunk=?1"))?;
107 let mut replace = connection.prepare_cached(&format!(
108 "INSERT INTO {table}(chunk, bytes) VALUES (?1, ?2) ON CONFLICT(chunk) DO UPDATE SET bytes=excluded.bytes"
109 ))?;
110 for (index, part) in bytes.chunks(CHUNK).enumerate() {
111 let current: Option<Vec<u8>> = stored
112 .query_row([index as i64], |row| row.get(0))
113 .optional()?;
114 if current.as_deref() != Some(part) {
115 replace.execute(params![index as i64, part])?;
116 }
117 }
118 Ok(())
119}
120
121pub(crate) fn clear(connection: &Connection, image: Image) -> Result<()> {
122 connection.execute(&format!("DELETE FROM {}", image.table()), [])?;
123 Ok(())
124}
125
126/// Commits `transaction` to the base image the way it commits to the file, rewriting only
127/// the chunks it writes.
128pub(crate) fn publish(connection: &Connection, transaction: &Transaction) -> Result<()> {
129 if stamp(connection, Image::Base)?.as_ref() != Some(transaction.base()) {
130 return Err(io::Error::new(
131 io::ErrorKind::InvalidData,
132 "The published transaction is not on the cached base image",
133 )
134 .into());
135 }
136 let mut stored = connection.prepare_cached("SELECT bytes FROM base WHERE chunk=?1")?;
137 let mut replace = connection.prepare_cached(
138 "INSERT INTO base(chunk, bytes) VALUES (?1, ?2) ON CONFLICT(chunk) DO UPDATE SET bytes=excluded.bytes",
139 )?;
140 let mut chunks: std::collections::BTreeMap<usize, Vec<u8>> = Default::default();
141 for (offset, bytes) in transaction.writes() {
142 let mut offset = usize::try_from(offset).map_err(|_| damaged())?;
143 let mut bytes = bytes;
144 while !bytes.is_empty() {
145 let index = offset / CHUNK;
146 let chunk = match chunks.entry(index) {
147 std::collections::btree_map::Entry::Occupied(entry) => entry.into_mut(),
148 std::collections::btree_map::Entry::Vacant(entry) => entry.insert(
149 stored
150 .query_row([index as i64], |row| row.get(0))
151 .optional()?
152 .unwrap_or_default(),
153 ),
154 };
155 let at = offset % CHUNK;
156 let take = bytes.len().min(CHUNK - at);
157 if chunk.len() < at {
158 return Err(damaged());
159 }
160 if chunk.len() < at + take {
161 chunk.resize(at + take, 0);
162 }
163 chunk[at..at + take].copy_from_slice(&bytes[..take]);
164 offset += take;
165 bytes = &bytes[take..];
166 }
167 }
168 for (index, bytes) in chunks {
169 replace.execute(params![index as i64, bytes])?;
170 }
171 Ok(())
172}
173
174#[cfg(test)]
175mod tests {
176 use super::*;
177
178 #[test]
179 fn damaged_chunk_lengths_do_not_wrap_the_stamp() {
180 for image in [Image::Base, Image::Remote] {
181 for (last, size) in [(i64::MAX, 1), (1, 0), (1, CHUNK + 1)] {
182 let connection = connection();
183 connection
184 .execute(
185 &format!("INSERT INTO {} VALUES (0, ?1), (?2, ?3)", image.table()),
186 params![vec![0u8; CHUNK], last, vec![0u8; size]],
187 )
188 .unwrap();
189 assert!(stamp(&connection, image).is_err());
190 }
191 }
192 }
193
194 fn connection() -> Connection {
195 let connection = Connection::open_in_memory().unwrap();
196 connection.execute_batch(crate::schema::QUEUE).unwrap();
197 connection
198 }
199
200 /// Types `text` at the start of the page's first body paragraph.
201 pub(crate) fn typed(section: &mut onestore::Section<'_>, space: onestore::ExGuid, text: &str) {
202 use onestore::op::{Edit, Op, PageOp};
203 let page = section.page(space).unwrap();
204 let target = page
205 .objects
206 .iter()
207 .find_map(|object| match object {
208 onestore::page::PageObject::Outline(outline) => outline
209 .paragraphs
210 .iter()
211 .find_map(|p| p.text().map(|t| t.id)),
212 _ => None,
213 })
214 .unwrap();
215 section
216 .apply(
217 "Author",
218 &Edit {
219 at: 133_000_000_000_000_000,
220 ops: vec![Op::Page {
221 space,
222 op: PageOp::Text {
223 text: target,
224 range: 0..0,
225 with: text.into(),
226 },
227 }],
228 },
229 )
230 .unwrap();
231 }
232
233 #[test]
234 fn a_published_transaction_rewrites_only_its_chunks_and_matches_the_file() {
235 let source =
236 include_bytes!("../../../corpus/outline-edit/before/notebook/synthetic.one").to_vec();
237 assert!(source.len() > 4 * CHUNK);
238 let arena = onestore::Arena::default();
239 let mut section = onestore::Section::open(&arena, source.clone()).unwrap();
240 let space = section.pages().unwrap()[0].0;
241 let connection = connection();
242 write(&connection, Image::Base, &source).unwrap();
243 assert_eq!(read(&connection, Image::Base).unwrap().unwrap(), source);
244 assert_eq!(
245 stamp(&connection, Image::Base).unwrap().unwrap(),
246 Stamp::of(&source).unwrap()
247 );
248 let mut image = source.clone();
249 for round in 0..3 {
250 typed(&mut section, space, &format!("{round}"));
251 let transaction = section.seal().unwrap().unwrap();
252 publish(&connection, &transaction).unwrap();
253 transaction.apply(&mut image).unwrap();
254 assert_eq!(read(&connection, Image::Base).unwrap().unwrap(), image);
255 assert_eq!(
256 stamp(&connection, Image::Base).unwrap().unwrap(),
257 Stamp::of(&image).unwrap()
258 );
259 assert!(publish(&connection, &transaction).is_err());
260 }
261 let shorter = &image[..CHUNK + 5];
262 write(&connection, Image::Base, shorter).unwrap();
263 assert_eq!(read(&connection, Image::Base).unwrap().unwrap(), shorter);
264 clear(&connection, Image::Remote).unwrap();
265 assert!(read(&connection, Image::Remote).unwrap().is_none());
266 }
267}