A device's parents' log in journal_line: append-only, and beside it the copy of every other device's lines
that sync pulls from the service.
4use std::collections::BTreeMap;
The picture names a log entry shows (the ones she showed Whiskers), so a picture a logged conversation 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}
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}
This device's log, and the replica of the others'. A handle is only a name and a [Store], so it is cheap to
make another over the same database.
The most days one page of [SqliteLog::days] holds.
43pub const MAX_DAYS_PAGE: usize = 366;
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;
One day of the log, as the parents' page of days shows it.
The moment the household's day began, in milliseconds since the epoch (UTC).
52 pub start_ms: u64,
How many lines the day still has.
54 pub lines: usize,
What sync keeps of the service's log (every device's lines, in the service's order).
60pub type Replica = SqliteLog;
How many lines this device has written.
This device's own lines from the fromth on (counting from zero), at most limit, as the JSON each was
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 }
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 }
Takes in lines of the service's log that continue exactly where the copy ends. Each is
{"device":..,"entry":{..}}; the ones from this device are already here and only move the place.
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 }
Throws the copy of the other devices' lines away, so the log is taken again from the start (the copy 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 }
The entries written from from_ms up to (not including) to_ms, grouped by the device that wrote them (this
device's first), and how many lines could not be read. It reads what the index by time finds and nothing
else: this is how the parents' view reads a day, never the whole log. An exchange is a run of one device's
entries, so a caller that wants the exchanges that begin in a range asks for a little past its end (see
[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 }
Up to limit days (at most [MAX_DAYS_PAGE]) that have lines, newest first, from the one before before_ms
back: the page the parents turn through. A day is the household's own, utc_offset_minutes from UTC. Each
says how many lines it has and whether the household's retention setting has cleared some. One index lookup and
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 }
Every picture name some line of the log shows, once each. Read from the index of pictures by line, not 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 }
Empties every line written before cutoff_ms that is not already empty: its content goes (the line stays as
[Event::Expired], with its place and its time, so the counts that sync is made of do not move) and the
pictures only those lines showed are deleted. Returns how many lines were emptied. The caller says that it did
(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 }
Every entry on this device, grouped by the device that wrote it (this device's first), and how many
lines could not be read. An exchange is a run of one device's entries. This reads the whole log. The
parents' view uses 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 }
This device's own entries, oldest first.
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}