log.rsannotatedlog.rssource295 lines · 15.1 KB · raw
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}