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}