whiskers.git / crates / whiskers-store / src / reflections.rs
1//! What waits to be thought about, in `pending_reflection`.
2//!
3//! After a turn the memory work reads the exchange and decides what to remember. That can fail (the model cannot be
4//! reached) or be cut short (the process stops), so the exchange is kept here until it has been dealt with, and what
5//! is here is still here when the app starts again. The rows hold her words, so the schema deletes them: all of them
6//! when the parents forget a memory, and a turn's when that turn is emptied (see `SCHEMA.md`).
7
8use ::log::{debug, error, info, warn};
9use rusqlite::params;
10use whiskers_core::{PictureId, ReflectInput};
11
12use crate::{DbError, Store};
13
14/// A waiting exchange with the number it is kept under, so that exactly this one can be taken off.
15#[derive(Clone, Debug, PartialEq, Eq)]
16pub struct Waiting {
17    pub id: u64,
18    pub input: ReflectInput,
19}
20
21/// The queue of one database. A handle is only a [`Store`].
22#[derive(Clone)]
23pub struct SqliteReflections {
24    store: Store,
25}
26
27impl SqliteReflections {
28    pub fn new(store: Store) -> Self {
29        Self { store }
30    }
31
32    /// Keeps `input` at the back of the queue, and drops the oldest while more than `limit` wait, so that a long
33    /// outage does not queue the whole evening. Returns how many were dropped to make room.
34    pub fn push(&self, input: &ReflectInput, now_ms: u64, limit: usize) -> Result<usize, DbError> {
35        let pictures = serde_json::to_string(&input.pictures.iter().map(|p| p.0.as_str()).collect::<Vec<_>>()).map_err(|e| DbError::Damaged(e.to_string()))?;
36        let dropped = self.store.transaction(|tx| {
37            tx.execute(
38                "INSERT INTO pending_reflection (turn, heard, said, description, pictures, queued_at_ms) VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
39                params![input.turn, input.heard, input.said, input.description, pictures, now_ms as i64],
40            )?;
41            let over: i64 = tx.query_row("SELECT MAX(COUNT(*) - ?1, 0) FROM pending_reflection", [limit as i64], |r| r.get(0))?;
42            if over > 0 {
43                tx.execute("DELETE FROM pending_reflection WHERE id IN (SELECT id FROM pending_reflection ORDER BY id LIMIT ?1)", [over])?;
44            }
45            Ok(over as usize)
46        })?;
47        if dropped > 0 {
48            warn!("{dropped} exchange(s) waiting for the memory were dropped to make room");
49        }
50        debug!("an exchange waits for the memory");
51        Ok(dropped)
52    }
53
54    /// Everything that waits, oldest first. Nothing is taken off: an exchange stays until [`done`](Self::done) says
55    /// it was dealt with, so a stop in between loses nothing.
56    pub fn waiting(&self) -> Result<Vec<Waiting>, DbError> {
57        let conn = self.store.lock();
58        let mut st = conn.prepare("SELECT id, turn, heard, said, description, pictures FROM pending_reflection ORDER BY id")?;
59        let rows = st.query_map([], |r| {
60            Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?, r.get::<_, String>(2)?, r.get::<_, String>(3)?, r.get::<_, Option<String>>(4)?, r.get::<_, String>(5)?))
61        })?;
62        let mut out = Vec::new();
63        for row in rows {
64            let (id, turn, heard, said, description, pictures) = row?;
65            let pictures: Vec<String> = serde_json::from_str(&pictures).map_err(|e| DbError::Damaged(format!("a waiting exchange's pictures: {e}")))?;
66            out.push(Waiting {
67                id: id as u64,
68                input: ReflectInput { turn, heard, said, description, pictures: pictures.into_iter().map(PictureId).collect() },
69            });
70        }
71        Ok(out)
72    }
73
74    /// Whether the exchange of `turn` still waits. False once it was dealt with, and once a forgetting dropped it.
75    pub fn has_turn(&self, turn: &str) -> Result<bool, DbError> {
76        Ok(self.store.lock().query_row("SELECT EXISTS (SELECT 1 FROM pending_reflection WHERE turn = ?1)", [turn], |r| r.get::<_, i64>(0))? == 1)
77    }
78
79    /// How many wait.
80    pub fn len(&self) -> Result<usize, DbError> {
81        Ok(self.store.lock().query_row("SELECT COUNT(*) FROM pending_reflection", [], |r| r.get::<_, i64>(0))? as usize)
82    }
83
84    pub fn is_empty(&self) -> Result<bool, DbError> {
85        Ok(self.len()? == 0)
86    }
87
88    /// The exchange kept under `id` was dealt with (or is not to be tried again): it is deleted, and what was kept of her
89    /// words is overwritten.
90    pub fn done(&self, id: u64) -> Result<(), DbError> {
91        self.store.transaction(|tx| {
92            tx.execute("DELETE FROM pending_reflection WHERE id = ?1", [id as i64])?;
93            Ok(())
94        })
95    }
96
97    /// The write-ahead log may still hold the words of exchanges that were dealt with. Emptied here, rather than after
98    /// each one, because every exchange passes through it.
99    pub fn tidy(&self) {
100        if let Err(e) = self.store.scrub() {
101            error!("the write-ahead log could not be emptied after the waiting exchanges were dealt with: {e}");
102        } else {
103            info!("the waiting exchanges' old bytes are out of the write-ahead log");
104        }
105    }
106}