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).
12use crate::{DbError, Store};
A waiting exchange with the number it is kept under, so that exactly this one can be taken off.
The queue of one database. A handle is only a [Store].
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.
How many wait.
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.
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.