1//! The one log, as `journal.jsonl`: every device's entries in the order they arrived, each line 2//! `{"device":..,"entry":{..}}`. Nothing is ever rewritten. A mutex makes an append atomic with 3//! respect to every other call, and each append is one write that is flushed to disk before it returns. 4 5use std::fs::{File, OpenOptions}; 6use std::io::Write as _; 7use std::path::{Path, PathBuf}; 8use std::sync::{Mutex, MutexGuard, PoisonError}; 9 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, 15 /// Held across every read-then-append: that is the whole of the journal's atomicity. 16 lock: Mutex<()>, 17} 18 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 } 38 39 /// 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 } 53 54 /// 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} 60 61/// Writes `text` at the end of `file` and flushes it; if that fails part-way the file is cut back to 62/// `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} 72 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}