1//! A device's parents' log in `journal_line`: append-only, and beside it the copy of every other device's lines 2//! that sync pulls from the service. 3 4use std::collections::BTreeMap; 5 6use ::log::{debug, error, info, trace, warn}; 7use rusqlite::{Transaction, params}; 8use whiskers_core::{Entry, Log, LogError}; 9 10use crate::{DbError, Store}; 11 12/// The picture names a log entry shows (the ones she showed Whiskers), so a picture a logged conversation 13/// still needs is not deleted with a memory made from it. 14pub(crate) fn heard_pictures(entry_json: &str) -> Vec<String> { 15 serde_json::from_str::<serde_json::Value>(entry_json) 16 .ok() 17 .and_then(|v| v.pointer("/event/Heard/pictures").and_then(|p| p.as_array()).map(|a| a.iter().filter_map(|x| x.as_str().map(str::to_owned)).collect())) 18 .unwrap_or_default() 19} 20 21/// Adds one line to the log as `device`'s next. Returns the line's number among that device's. 22pub(crate) fn insert_line(tx: &Transaction<'_>, device: &str, entry_json: &str) -> Result<u64, DbError> { 23 let seq: i64 = tx.query_row("SELECT COALESCE(MAX(seq) + 1, 0) FROM journal_line WHERE device = ?1", [device], |r| r.get(0))?; 24 tx.execute("INSERT INTO journal_line (device, seq, entry) VALUES (?1, ?2, ?3)", params![device, seq, entry_json])?; 25 // The pictures are read from the line as stored, not as written: the schema empties a line of a forgotten turn 26 // as it arrives, and a picture it no longer names must not be kept for it. 27 let stored: String = tx.query_row("SELECT entry FROM journal_line WHERE device = ?1 AND seq = ?2", params![device, seq], |r| r.get(0))?; 28 for p in heard_pictures(&stored) { 29 tx.execute("INSERT OR IGNORE INTO journal_picture (device, seq, picture_id) VALUES (?1, ?2, ?3)", params![device, seq, p])?; 30 } 31 Ok(seq as u64) 32} 33 34/// This device's log, and the replica of the others'. A handle is only a name and a [`Store`], so it is cheap to 35/// make another over the same database. 36#[derive(Clone)] 37pub struct SqliteLog { 38 store: Store, 39 device: String, 40} 41 42/// The most days one page of [`SqliteLog::days`] holds. 43pub const MAX_DAYS_PAGE: usize = 366; 44 45/// How far past the end of a range an exchange that began inside it can run: a turn is a few calls over the network. 46pub const EXCHANGE_SLACK_MS: u64 = 10 * 60 * 1000; 47 48/// One day of the log, as the parents' page of days shows it. 49#[derive(Clone, Copy, Debug, PartialEq, Eq)] 50pub struct LogDay { 51 /// The moment the household's day began, in milliseconds since the epoch (UTC). 52 pub start_ms: u64, 53 /// How many lines the day still has. 54 pub lines: usize, 55 /// How many lines of the day the household's retention setting has emptied. 56 pub cleared: usize, 57} 58 59/// What sync keeps of the service's log (every device's lines, in the service's order). 60pub type Replica = SqliteLog; 61 62impl SqliteLog { 63 pub fn new(store: Store, device: &str) -> Self { 64 Self { store, device: device.to_owned() } 65 } 66 67 /// How many lines this device has written. 68 pub fn own_count(&self) -> Result<usize, LogError> { 69 self.store 70 .lock() 71 .query_row("SELECT COUNT(*) FROM journal_line WHERE device = ?1", [&self.device], |r| r.get::<_, i64>(0)) 72 .map(|n| n as usize) 73 .map_err(|e| LogError(e.to_string())) 74 } 75 76 /// This device's own lines from the `from`th on (counting from zero), at most `limit`, as the JSON each was 77 /// written in. 78 pub fn own_lines(&self, from: usize, limit: usize) -> Result<Vec<String>, LogError> { 79 let conn = self.store.lock(); 80 let mut st = conn 81 .prepare("SELECT entry FROM journal_line WHERE device = ?1 AND seq >= ?2 ORDER BY seq LIMIT ?3") 82 .map_err(|e| LogError(e.to_string()))?; 83 let rows = st.query_map(params![self.device, from as i64, limit as i64], |r| r.get(0)).map_err(|e| LogError(e.to_string()))?; 84 rows.collect::<Result<Vec<String>, _>>().map_err(|e| LogError(e.to_string())) 85 } 86 87 /// How many lines of the service's log this device has taken in (its place in it). 88 pub fn pulled(&self) -> Result<usize, LogError> { 89 let conn = self.store.lock(); 90 let n: Option<String> = rusqlite::OptionalExtension::optional(conn.query_row("SELECT value FROM meta WHERE key = 'journal_pulled'", [], |r| r.get(0))) 91 .map_err(|e| LogError(e.to_string()))?; 92 Ok(n.and_then(|n| n.parse().ok()).unwrap_or(0)) 93 } 94 95 /// Takes in lines of the service's log that continue exactly where the copy ends. Each is 96 /// `{"device":..,"entry":{..}}`; the ones from this device are already here and only move the place. 97 /// Returns how many lines were taken in. 98 pub fn take(&self, from: usize, lines: &[String]) -> Result<usize, LogError> { 99 let held = self.pulled()?; 100 if from != held { 101 warn!("take: lines begin at {from} but the copy ends at {held}; taking none"); 102 return Ok(0); 103 } 104 let me = self.device.clone(); 105 self.store 106 .transaction(|tx| { 107 let mut kept = 0; 108 for l in lines { 109 match serde_json::from_str::<serde_json::Value>(l) { 110 Ok(v) => match (v.get("device").and_then(|d| d.as_str()), v.get("entry")) { 111 (Some(d), Some(e)) if d != me => { 112 insert_line(tx, d, &e.to_string())?; 113 kept += 1; 114 } 115 (Some(_), Some(_)) => {} 116 _ => warn!("take: a line of the log has no device or entry ({} bytes); skipped", l.len()), 117 }, 118 Err(e) => warn!("take: a line of the log is not JSON ({} bytes): {e}", l.len()), 119 } 120 } 121 tx.execute( 122 "INSERT INTO meta (key, value) VALUES ('journal_pulled', ?1) ON CONFLICT (key) DO UPDATE SET value = excluded.value", 123 [(held + lines.len()).to_string()], 124 )?; 125 Ok(kept) 126 }) 127 .map_err(|e| { 128 error!("take: the pulled lines were not kept: {e}"); 129 LogError(e.to_string()) 130 }) 131 } 132 133 /// Throws the copy of the other devices' lines away, so the log is taken again from the start (the copy 134 /// was longer than the log it copies, so it was not a copy of it). 135 pub fn rebuild_replica(&self) -> Result<(), LogError> { 136 warn!("the copy of the service's log is being thrown away and taken again"); 137 let me = self.device.clone(); 138 self.store 139 .transaction(|tx| { 140 tx.execute("DELETE FROM journal_picture WHERE device <> ?1", [&me])?; 141 tx.execute("DELETE FROM journal_line WHERE device <> ?1", [&me])?; 142 tx.execute("DELETE FROM meta WHERE key = 'journal_pulled'", [])?; 143 Ok(()) 144 }) 145 .map_err(|e| LogError(e.to_string())) 146 } 147 148 /// The entries written from `from_ms` up to (not including) `to_ms`, grouped by the device that wrote them (this 149 /// device's first), and how many lines could not be read. It reads what the index by time finds and nothing 150 /// else: this is how the parents' view reads a day, never the whole log. An exchange is a run of one device's 151 /// entries, so a caller that wants the exchanges that begin in a range asks for a little past its end (see 152 /// [`EXCHANGE_SLACK_MS`]). 153 pub fn read_range(&self, from_ms: u64, to_ms: u64) -> Result<(Vec<Vec<Entry>>, usize), LogError> { 154 debug!("reading the log for device {} from {from_ms} to {to_ms}", self.device); 155 let conn = self.store.lock(); 156 let mut st = conn 157 .prepare("SELECT device, entry FROM journal_line WHERE at_ms >= ?1 AND at_ms < ?2 ORDER BY device, seq") 158 .map_err(|e| LogError(e.to_string()))?; 159 let (from, to) = (i64::try_from(from_ms).unwrap_or(i64::MAX), i64::try_from(to_ms).unwrap_or(i64::MAX)); 160 let rows = st.query_map(params![from, to], |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))).map_err(|e| LogError(e.to_string()))?; 161 let mut unreadable = 0; 162 let mut mine = Vec::new(); 163 let mut others: BTreeMap<String, Vec<Entry>> = BTreeMap::new(); 164 for row in rows { 165 let (device, text) = row.map_err(|e| LogError(e.to_string()))?; 166 match serde_json::from_str::<Entry>(&text) { 167 Ok(e) if device == self.device => mine.push(e), 168 Ok(e) => others.entry(device).or_default().push(e), 169 Err(e) => { 170 unreadable += 1; 171 warn!("unreadable line in the log ({} bytes): {e}", text.len()); 172 } 173 } 174 } 175 let mut groups = vec![mine]; 176 groups.extend(others.into_values()); 177 info!("log range: {} group(s), {unreadable} unreadable line(s)", groups.len()); 178 Ok((groups, unreadable)) 179 } 180 181 /// Up to `limit` days (at most [`MAX_DAYS_PAGE`]) that have lines, newest first, from the one before `before_ms` 182 /// back: the page the parents turn through. A day is the household's own, `utc_offset_minutes` from UTC. Each 183 /// says how many lines it has and whether the household's retention setting has cleared some. One index lookup and 184 /// one count per day, however long the log. 185 pub fn days(&self, utc_offset_minutes: i32, before_ms: u64, limit: usize) -> Result<Vec<LogDay>, LogError> { 186 const DAY_MS: i64 = 86_400_000; 187 let offset = i64::from(utc_offset_minutes) * 60_000; 188 let conn = self.store.lock(); 189 let mut upto = i64::try_from(before_ms).unwrap_or(i64::MAX); 190 let mut days = Vec::new(); 191 for _ in 0..limit.min(MAX_DAYS_PAGE) { 192 let last: Option<i64> = conn 193 .query_row("SELECT MAX(at_ms) FROM journal_line WHERE at_ms < ?1", [upto], |r| r.get(0)) 194 .map_err(|e| LogError(e.to_string()))?; 195 let Some(last) = last else { break }; 196 let start = (last + offset).div_euclid(DAY_MS) * DAY_MS - offset; 197 let (lines, cleared): (i64, i64) = conn 198 .query_row( 199 "SELECT COUNT(*), COALESCE(SUM(json_extract(entry, '$.event') IS 'Expired'), 0) FROM journal_line WHERE at_ms >= ?1 AND at_ms < ?2", 200 [start, start + DAY_MS], 201 |r| Ok((r.get(0)?, r.get(1)?)), 202 ) 203 .map_err(|e| LogError(e.to_string()))?; 204 days.push(LogDay { start_ms: start.max(0) as u64, lines: (lines - cleared) as usize, cleared: cleared as usize }); 205 upto = start; 206 } 207 Ok(days) 208 } 209 210 /// Every picture name some line of the log shows, once each. Read from the index of pictures by line, not 211 /// from the lines. 212 pub fn pictures_shown(&self) -> Result<Vec<String>, LogError> { 213 let conn = self.store.lock(); 214 let mut st = conn.prepare("SELECT DISTINCT picture_id FROM journal_picture ORDER BY picture_id").map_err(|e| LogError(e.to_string()))?; 215 let rows = st.query_map([], |r| r.get(0)).map_err(|e| LogError(e.to_string()))?; 216 rows.collect::<Result<Vec<String>, _>>().map_err(|e| LogError(e.to_string())) 217 } 218 219 /// Empties every line written before `cutoff_ms` that is not already empty: its content goes (the line stays as 220 /// [`Event::Expired`], with its place and its time, so the counts that sync is made of do not move) and the 221 /// pictures only those lines showed are deleted. Returns how many lines were emptied. The caller says that it did 222 /// (see [`Event::Cleared`]); nothing here is silent, but it is not this method's to write. 223 pub fn expire_before(&self, cutoff_ms: u64) -> Result<usize, LogError> { 224 let cutoff = i64::try_from(cutoff_ms).unwrap_or(i64::MAX); 225 let emptied = self 226 .store 227 .transaction(|tx| { 228 // The pictures go first, while the lines that show them can still be found by their time. 229 tx.execute("DELETE FROM journal_picture WHERE (device, seq) IN (SELECT device, seq FROM journal_line WHERE at_ms < ?1 AND live)", [cutoff])?; 230 let n = tx.execute( 231 "UPDATE journal_line SET entry = json_object('at_ms', at_ms, 'event', 'Expired') WHERE at_ms < ?1 AND live", 232 [cutoff], 233 )?; 234 Ok(n) 235 }) 236 .map_err(|e| { 237 error!("the log could not be cleared of old lines: {e}"); 238 LogError(e.to_string()) 239 })?; 240 if emptied > 0 { 241 info!("log: {emptied} line(s) older than {cutoff_ms} ms were emptied"); 242 // What was emptied must not wait in the write-ahead log. 243 self.store.scrub().map_err(|e| LogError(e.to_string()))?; 244 } 245 Ok(emptied) 246 } 247 248 /// Every entry on this device, grouped by the device that wrote it (this device's first), and how many 249 /// lines could not be read. An exchange is a run of one device's entries. **This reads the whole log.** The 250 /// parents' view uses [`read_range`](Self::read_range); this is for tests and for a tool that really means all of it. 251 pub fn read_groups(&self) -> Result<(Vec<Vec<Entry>>, usize), LogError> { 252 debug!("reading the log for device {}", self.device); 253 let conn = self.store.lock(); 254 let mut st = conn.prepare("SELECT device, entry FROM journal_line ORDER BY device, seq").map_err(|e| LogError(e.to_string()))?; 255 let mut rows = st.query([]).map_err(|e| LogError(e.to_string()))?; 256 let mut unreadable = 0; 257 let mut mine = Vec::new(); 258 let mut others: BTreeMap<String, Vec<Entry>> = BTreeMap::new(); 259 while let Some(r) = rows.next().map_err(|e| LogError(e.to_string()))? { 260 let (device, text): (String, String) = (r.get(0).map_err(|e| LogError(e.to_string()))?, r.get(1).map_err(|e| LogError(e.to_string()))?); 261 match serde_json::from_str::<Entry>(&text) { 262 Ok(e) if device == self.device => mine.push(e), 263 Ok(e) => others.entry(device).or_default().push(e), 264 Err(e) => { 265 unreadable += 1; 266 warn!("unreadable line in the log ({} bytes): {e}", text.len()); 267 } 268 } 269 } 270 let mut groups = vec![mine]; 271 groups.extend(others.into_values()); 272 info!("log: {} group(s), {unreadable} unreadable line(s)", groups.len()); 273 Ok((groups, unreadable)) 274 } 275 276 /// This device's own entries, oldest first. 277 pub fn own_entries(&self) -> Result<Vec<Entry>, LogError> { 278 let (mut groups, _) = self.read_groups()?; 279 Ok(std::mem::take(&mut groups[0])) 280 } 281} 282 283impl Log for SqliteLog { 284 fn append(&mut self, entry: &Entry) -> Result<(), LogError> { 285 let json = serde_json::to_string(entry).map_err(|e| { 286 error!("journal entry does not serialize: {e}"); 287 LogError(e.to_string()) 288 })?; 289 trace!("journal append of {} bytes", json.len()); 290 self.store.transaction(|tx| insert_line(tx, &self.device, &json).map(|_| ())).map_err(|e| { 291 error!("journal append of {} bytes failed: {e}", json.len()); 292 LogError(e.to_string()) 293 }) 294 } 295}