log.rsannotatedlog.rssource295 lines · 15.1 KB · raw

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;
6use ::log::{debug, error, info, trace, warn};
7use rusqlite::{Transaction, params};
8use whiskers_core::{Entry, Log, LogError};
9
10use crate::{DbError, Store};

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.

36#[derive(Clone)]
37pub struct SqliteLog {
38    store: Store,
39    device: String,
40}

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.

49#[derive(Clone, Copy, Debug, PartialEq, Eq)]
50pub struct LogDay {

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,

How many lines of the day the household's retention setting has emptied.

56    pub cleared: usize,
57}

What sync keeps of the service's log (every device's lines, in the service's order).

60pub type Replica = SqliteLog;
62impl SqliteLog {
63    pub fn new(store: Store, device: &str) -> Self {
64        Self { store, device: device.to_owned() }
65    }

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    }

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.

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}
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}