whiskers.git / crates / whiskersd / src / journal.rs

The one log, as journal.jsonl: every device's entries in the order they arrived, each line {"device":..,"entry":{..}}. Nothing is ever rewritten. A mutex makes an append atomic with respect to every other call, and each append is one write that is flushed to disk before it returns.

5use std::fs::{File, OpenOptions};
6use std::io::Write as _;
7use std::path::{Path, PathBuf};
8use std::sync::{Mutex, MutexGuard, PoisonError};
10use log::{debug, error, info, warn};
11use whiskers_ports::{Appended, DeviceCursor, DeviceName, EntryLine, Journal, LogCursor, Page, StoreError, StoredLine};
12
13pub struct FileJournal {
14    dir: PathBuf,

Held across every read-then-append: that is the whole of the journal's atomicity.

16    lock: Mutex<()>,
17}
19impl FileJournal {
20    pub fn open(dir: &Path) -> Self {
21        if let Err(e) = std::fs::create_dir_all(dir) {
22            error!("cannot create the journal directory {}: {e}", dir.display());
23        }
24        Self { dir: dir.to_owned(), lock: Mutex::new(()) }
25    }
26
27    pub fn dir(&self) -> &Path {
28        &self.dir
29    }
30
31    fn path(&self) -> PathBuf {
32        self.dir.join("journal.jsonl")
33    }
34
35    fn locked(&self) -> MutexGuard<'_, ()> {
36        self.lock.lock().unwrap_or_else(PoisonError::into_inner)
37    }

The complete lines of the log (a last line with no newline is a write in progress).

40    fn lines(&self) -> Vec<String> {
41        let path = self.path();
42        let bytes = std::fs::read(&path).unwrap_or_else(|e| {
43            debug!("journal {} not readable ({e}); treated as empty", path.display());
44            Vec::new()
45        });
46        let text = String::from_utf8_lossy(&bytes);
47        let complete = match text.rfind('\n') {
48            Some(i) => &text[..=i],
49            None => "",
50        };
51        complete.lines().map(str::to_owned).collect()
52    }

How many of the log's lines came from this device: where its next line must continue.

55    fn held_in(lines: &[String], device: &DeviceName) -> DeviceCursor {
56        let n = lines.iter().filter(|l| serde_json::from_str::<serde_json::Value>(l).is_ok_and(|e| e["device"] == device.as_str())).count();
57        DeviceCursor::new(n as u64)
58    }
59}

Writes text at the end of file and flushes it; if that fails part-way the file is cut back to before, so an append is all or nothing.

63fn append_whole(file: &mut File, before: u64, text: &str) -> std::io::Result<()> {
64    let wrote = file.write_all(text.as_bytes()).and_then(|()| file.sync_data());
65    if wrote.is_err()
66        && let Err(e) = file.set_len(before)
67    {
68        error!("journal: could not cut back a failed append: {e}");
69    }
70    wrote
71}
73impl Journal for FileJournal {
74    async fn held(&self, device: &DeviceName) -> Result<DeviceCursor, StoreError> {
75        let _guard = self.locked();
76        Ok(Self::held_in(&self.lines(), device))
77    }
78
79    async fn append(&self, device: &DeviceName, at: DeviceCursor, lines: &[EntryLine]) -> Result<Appended, StoreError> {
80        let _guard = self.locked();
81        let all = self.lines();
82        let held = Self::held_in(&all, device);
83        if lines.is_empty() {
84            return Ok(Appended::Accepted { held });
85        }
86        if at != held {
87            warn!("journal append for {} at {} but the journal holds {}", device.as_str(), at.get(), held.get());
88            return Ok(Appended::Misaligned { held });
89        }
90        let mut text = String::new();
91        for l in lines {
92            text.push_str(StoredLine::of(device, l).as_str());
93            text.push('\n');
94        }
95        let path = self.path();
96        let mut file = OpenOptions::new().create(true).append(true).open(&path).map_err(|e| {
97            error!("journal append for {} failed to open {}: {e}", device.as_str(), path.display());
98            StoreError::Unavailable
99        })?;
100        let before = file.metadata().map(|m| m.len()).map_err(|e| {
101            error!("journal append for {} failed to measure the log: {e}", device.as_str());
102            StoreError::Unavailable
103        })?;
104        append_whole(&mut file, before, &text).map_err(|e| {
105            error!("journal append for {} failed: {e}", device.as_str());
106            StoreError::Unavailable
107        })?;
108        info!("journal: {} line(s) appended for {}; it now has {}", lines.len(), device.as_str(), held.after(lines.len()).get());
109        Ok(Appended::Accepted { held: held.after(lines.len()) })
110    }
111
112    async fn pull(&self, since: LogCursor) -> Result<Page, StoreError> {
113        let _guard = self.locked();
114        let all: Vec<StoredLine> = self.lines().into_iter().map(StoredLine::from_storage).collect();
115        Ok(Page::of(&all, since))
116    }
117}
118
119#[cfg(test)]
120mod tests {
121    use super::*;
122    use whiskers_ports::run_ready;
123
124    fn scratch(name: &str) -> PathBuf {
125        let d = std::env::temp_dir().join(format!("whiskers-journal-{name}-{}", std::process::id()));
126        let _ = std::fs::remove_dir_all(&d);
127        d
128    }
129
130    #[test]
131    fn a_half_written_last_line_is_not_a_line_yet() {
132        let dir = scratch("half");
133        let j = FileJournal::open(&dir);
134        std::fs::write(j.path(), "{\"device\":\"phone1\",\"entry\":{\"n\":0}}\n{\"device\":\"phone1\",\"en").unwrap();
135        let p = run_ready(j.pull(LogCursor::START)).unwrap();
136        assert_eq!(p.total.get(), 1, "only the complete line counts");
137        assert_eq!(run_ready(j.held(&DeviceName::new("phone1").unwrap())).unwrap().get(), 1);
138        let _ = std::fs::remove_dir_all(dir);
139    }
140
141    #[test]
142    fn many_threads_pushing_at_one_cursor_take_exactly_one_push() {
143        let dir = scratch("race");
144        let j = std::sync::Arc::new(FileJournal::open(&dir));
145        let device = DeviceName::new("phone1").unwrap();
146        let accepted: usize = (0..8)
147            .map(|_| {
148                let (j, device) = (j.clone(), device.clone());
149                std::thread::spawn(move || {
150                    let lines = [EntryLine::parse(r#"{"n":1}"#).unwrap(), EntryLine::parse(r#"{"n":2}"#).unwrap()];
151                    matches!(run_ready(j.append(&device, DeviceCursor::new(0), &lines)).unwrap(), Appended::Accepted { .. })
152                })
153            })
154            .collect::<Vec<_>>()
155            .into_iter()
156            .map(|t| usize::from(t.join().unwrap()))
157            .sum();
158        assert_eq!(accepted, 1, "the old service, with no lock, took the same lines more than once");
159        assert_eq!(run_ready(j.pull(LogCursor::START)).unwrap().total.get(), 2);
160        let _ = std::fs::remove_dir_all(dir);
161    }
162}