whiskers.git / crates / whiskers-store / src / reflections.rs

What waits to be thought about, in pending_reflection.

After a turn the memory work reads the exchange and decides what to remember. That can fail (the model cannot be reached) or be cut short (the process stops), so the exchange is kept here until it has been dealt with, and what is here is still here when the app starts again. The rows hold her words, so the schema deletes them: all of them when the parents forget a memory, and a turn's when that turn is emptied (see SCHEMA.md).

8use ::log::{debug, error, info, warn};
9use rusqlite::params;
10use whiskers_core::{PictureId, ReflectInput};
12use crate::{DbError, Store};

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}

The queue of one database. A handle is only a [Store].

22#[derive(Clone)]
23pub struct SqliteReflections {
24    store: Store,
25}
27impl SqliteReflections {
28    pub fn new(store: Store) -> Self {
29        Self { store }
30    }

Keeps input at the back of the queue, and drops the oldest while more than limit wait, so that a long 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    }

Everything that waits, oldest first. Nothing is taken off: an exchange stays until done says 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    }

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    }

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    }
84    pub fn is_empty(&self) -> Result<bool, DbError> {
85        Ok(self.len()? == 0)
86    }

The exchange kept under id was dealt with (or is not to be tried again): it is deleted, and what was kept of her 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    }

The write-ahead log may still hold the words of exchanges that were dealt with. Emptied here, rather than after 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}