1//! The service's journal: every device's lines in `journal_line`, in the order they arrived. 2 3use ::log::{error, info, warn}; 4use rusqlite::params; 5use whiskers_ports::{Appended, DeviceCursor, DeviceName, EntryLine, Journal, LogCursor, PULL_LIMIT, Page, StoreError, StoredLine}; 6 7use crate::log::insert_line; 8use crate::{DbError, Store}; 9 10pub struct SqliteJournal { 11 store: Store, 12} 13 14impl SqliteJournal { 15 pub fn new(store: Store) -> Self { 16 Self { store } 17 } 18} 19 20fn held_by(conn: &rusqlite::Connection, device: &DeviceName) -> Result<u64, rusqlite::Error> { 21 conn.query_row("SELECT COALESCE(MAX(seq) + 1, 0) FROM journal_line WHERE device = ?1", [device.as_str()], |r| r.get::<_, i64>(0)).map(|n| n as u64) 22} 23 24impl Journal for SqliteJournal { 25 async fn held(&self, device: &DeviceName) -> Result<DeviceCursor, StoreError> { 26 held_by(&self.store.lock(), device).map(DeviceCursor::new).map_err(|e| { 27 error!("journal: cannot count {}'s lines: {e}", device.as_str()); 28 StoreError::Unavailable 29 }) 30 } 31 32 async fn append(&self, device: &DeviceName, at: DeviceCursor, lines: &[EntryLine]) -> Result<Appended, StoreError> { 33 // One transaction under the one lock: counting and appending are a single step, so two pushes at one 34 // cursor take exactly one. 35 let appended = self.store.transaction(|tx| { 36 let held = held_by(tx, device)?; 37 if lines.is_empty() { 38 return Ok(Appended::Accepted { held: DeviceCursor::new(held) }); 39 } 40 if at.get() != held { 41 warn!("journal append for {} at {} but the journal holds {held}", device.as_str(), at.get()); 42 return Ok(Appended::Misaligned { held: DeviceCursor::new(held) }); 43 } 44 for l in lines { 45 insert_line(tx, device.as_str(), l.as_str())?; 46 } 47 let now = DeviceCursor::new(held).after(lines.len()); 48 info!("journal: {} line(s) appended for {}; it now has {}", lines.len(), device.as_str(), now.get()); 49 Ok(Appended::Accepted { held: now }) 50 }); 51 appended.map_err(|e: DbError| { 52 error!("journal append for {} failed: {e}", device.as_str()); 53 StoreError::Unavailable 54 }) 55 } 56 57 async fn pull(&self, since: LogCursor) -> Result<Page, StoreError> { 58 let conn = self.store.lock(); 59 let read = || -> Result<Page, DbError> { 60 let total: i64 = conn.query_row("SELECT COUNT(*) FROM journal_line", [], |r| r.get(0))?; 61 let from = i64::try_from(since.get()).unwrap_or(i64::MAX).min(total); 62 let mut st = conn.prepare("SELECT device, entry FROM journal_line ORDER BY position LIMIT ?1 OFFSET ?2")?; 63 let rows = st.query_map(params![PULL_LIMIT as i64, from], |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?)))?; 64 let mut lines = Vec::new(); 65 for r in rows { 66 let (device, entry) = r?; 67 let (Ok(device), Ok(entry)) = (DeviceName::new(&device), EntryLine::parse(&entry)) else { 68 return Err(DbError::Damaged("a journal line is not a device's entry".into())); 69 }; 70 lines.push(StoredLine::of(&device, &entry)); 71 } 72 Ok(Page { from: LogCursor::new(from as u64), total: LogCursor::new(total as u64), lines }) 73 }; 74 read().map_err(|e| { 75 error!("journal pull failed: {e}"); 76 match e { 77 DbError::Damaged(_) => StoreError::Damaged, 78 _ => StoreError::Unavailable, 79 } 80 }) 81 } 82}