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}