1//! The journal: the one append-only log of everything said to Whiskers, from every device, in the 2//! order the hub received it. Parents read it as a single conversation wherever it happened. 3//! 4//! Nothing in it is ever rewritten, so it cannot conflict. A device tells the hub how many of its 5//! own lines the hub already holds and sends the rest; any device reads the whole log from the 6//! position it last reached. Two different counts are in play and they are different types on purpose: 7//! a [`DeviceCursor`] counts one device's lines, a [`LogCursor`] is a position in the whole log. 8 9use ::log::trace; 10use serde::{Deserialize, Serialize}; 11 12use crate::error::StoreError; 13 14/// The longest device name accepted. A device's name is written into the log, so only plain ones get in. 15pub const DEVICE_NAME_MAX: usize = 64; 16 17/// At most this many lines go back in one pull; the caller asks again for the rest. 18pub const PULL_LIMIT: usize = 2000; 19 20/// A device's name as the log records it: 1 to 64 ASCII letters and digits. 21#[derive(Clone, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] 22pub struct DeviceName(String); 23 24#[derive(Clone, Copy, Debug, PartialEq, Eq)] 25pub struct NotADeviceName; 26 27impl DeviceName { 28 pub fn new(name: &str) -> Result<Self, NotADeviceName> { 29 if !name.is_empty() && name.len() <= DEVICE_NAME_MAX && name.chars().all(|c| c.is_ascii_alphanumeric()) { 30 Ok(Self(name.to_owned())) 31 } else { 32 Err(NotADeviceName) 33 } 34 } 35 36 pub fn as_str(&self) -> &str { 37 &self.0 38 } 39} 40 41/// One entry a device wrote: a single JSON object on a single line. The text is kept compact, in the 42/// key order the device wrote it, so what is stored is what was sent. 43#[derive(Clone, Debug, PartialEq, Eq)] 44pub struct EntryLine(String); 45 46#[derive(Clone, Copy, Debug, PartialEq, Eq)] 47pub struct NotAnEntry; 48 49impl EntryLine { 50 pub fn parse(raw: &str) -> Result<Self, NotAnEntry> { 51 if raw.contains('\n') { 52 return Err(NotAnEntry); 53 } 54 match serde_json::from_str::<serde_json::Value>(raw) { 55 Ok(v) if v.is_object() => Ok(Self(v.to_string())), 56 _ => Err(NotAnEntry), 57 } 58 } 59 60 pub fn as_str(&self) -> &str { 61 &self.0 62 } 63} 64 65/// A line of the log as stored and as pulled: `{"device":"..","entry":{..}}`. Built in one place so 66/// every backend stores the same bytes. 67#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] 68pub struct StoredLine(String); 69 70impl StoredLine { 71 pub fn of(device: &DeviceName, entry: &EntryLine) -> Self { 72 // The device name is plain ASCII, so quoting it needs no escaping. 73 Self(format!("{{\"device\":\"{}\",\"entry\":{}}}", device.as_str(), entry.as_str())) 74 } 75 76 /// A line read back from storage that [`of`](Self::of) produced. Adapters use this to hand lines 77 /// out again; nothing else should have a reason to make one. 78 pub fn from_storage(raw: String) -> Self { 79 Self(raw) 80 } 81 82 pub fn as_str(&self) -> &str { 83 &self.0 84 } 85 86 pub fn into_string(self) -> String { 87 self.0 88 } 89} 90 91/// How many of one device's lines the journal holds: where that device's next line must continue. 92#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, PartialOrd, Ord)] 93pub struct DeviceCursor(u64); 94 95impl DeviceCursor { 96 pub const fn new(n: u64) -> Self { 97 Self(n) 98 } 99 100 pub const fn get(self) -> u64 { 101 self.0 102 } 103 104 #[must_use] 105 pub const fn after(self, lines: usize) -> Self { 106 Self(self.0 + lines as u64) 107 } 108} 109 110/// A position in the whole log: how many lines (of every device) have been seen. 111#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, PartialOrd, Ord)] 112pub struct LogCursor(u64); 113 114impl LogCursor { 115 pub const START: LogCursor = LogCursor(0); 116 117 pub const fn new(n: u64) -> Self { 118 Self(n) 119 } 120 121 pub const fn get(self) -> u64 { 122 self.0 123 } 124} 125 126/// What an append did. 127#[derive(Clone, Copy, Debug, PartialEq, Eq)] 128pub enum Appended { 129 /// The lines were taken (or there were none), and the journal now holds this many of the device's. 130 Accepted { held: DeviceCursor }, 131 /// The device thought the journal held a different number than it does, so nothing was taken. The 132 /// device is told the real count and sends again from there. This is how a device that was behind 133 /// (a lost reply, a restored backup) or ahead recovers without duplicating or skipping a line. 134 Misaligned { held: DeviceCursor }, 135} 136 137impl Appended { 138 pub fn held(self) -> DeviceCursor { 139 match self { 140 Appended::Accepted { held } | Appended::Misaligned { held } => held, 141 } 142 } 143} 144 145/// A slice of the log. 146#[derive(Clone, Debug, PartialEq, Eq)] 147pub struct Page { 148 /// Where `lines` begin. Never past `total`: a caller that thought it had read further than the 149 /// log goes (its copy is longer than the log, as a damaged copy is) is told where the log 150 /// really ends, and sees `from` is less than what it held, and rebuilds its copy from here. 151 pub from: LogCursor, 152 /// The whole log's length. 153 pub total: LogCursor, 154 pub lines: Vec<StoredLine>, 155} 156 157impl Page { 158 /// The page `since` asks for from a log of `all`. One definition, so no backend can read the rule 159 /// differently: clamp to the end, then at most [`PULL_LIMIT`] lines. 160 pub fn of(all: &[StoredLine], since: LogCursor) -> Self { 161 let from = usize::try_from(since.get()).unwrap_or(usize::MAX).min(all.len()); 162 let lines: Vec<StoredLine> = all[from..].iter().take(PULL_LIMIT).cloned().collect(); 163 trace!("journal page from {from} of {}: {} line(s)", all.len(), lines.len()); 164 Self { from: LogCursor::new(from as u64), total: LogCursor::new(all.len() as u64), lines } 165 } 166} 167 168impl Page { 169 /// Whether the caller's copy of the log is longer than the log: it held `held` lines and this page 170 /// begins before that. The only reason a page can start earlier than asked is that the log is 171 /// shorter than the caller believed (a duplicated copy, say), so its copy is wrong 172 /// from `self.from` on and is rebuilt from this page rather than extended. 173 pub fn caller_is_ahead(&self, held: LogCursor) -> bool { 174 self.from < held 175 } 176} 177 178/// The journal. 179/// 180/// Contract, for every adapter: 181/// 182/// - **Append-only.** A stored line is never changed, moved or removed. 183/// - **Atomic per call.** An append either stores all its lines or none, and two appends overlapping 184/// (even from one device, even at the same cursor) behave as if one ran first: exactly one is 185/// `Accepted`; the other sees the first's lines and is `Misaligned`. No line is ever duplicated. 186/// - **Aligned only.** Lines are taken only when `at` equals what the journal holds for that device. 187/// Appending no lines is a question, and is always `Accepted { held }`. 188/// - **Order is arrival order,** across devices. 189/// - **Durable once it returns `Ok`.** 190#[expect(async_fn_in_trait, reason = "a Worker's futures hold JavaScript values and cannot be Send, so no Send bound may be required here")] 191pub trait Journal { 192 /// How many of `device`'s lines the journal holds. 193 async fn held(&self, device: &DeviceName) -> Result<DeviceCursor, StoreError>; 194 195 async fn append(&self, device: &DeviceName, at: DeviceCursor, lines: &[EntryLine]) -> Result<Appended, StoreError>; 196 197 /// The log from `since` on, as [`Page::of`] defines it. 198 async fn pull(&self, since: LogCursor) -> Result<Page, StoreError>; 199} 200 201#[cfg(test)] 202mod tests { 203 use super::*; 204 205 #[test] 206 fn only_plain_device_names_get_in() { 207 assert!(DeviceName::new("phone1").is_ok()); 208 for bad in ["", "bad name", "a-b", "../x", "caf\u{e9}"] { 209 assert!(DeviceName::new(bad).is_err(), "{bad:?}"); 210 } 211 assert!(DeviceName::new(&"a".repeat(64)).is_ok()); 212 assert!(DeviceName::new(&"a".repeat(65)).is_err()); 213 } 214 215 #[test] 216 fn an_entry_is_one_json_object_on_one_line() { 217 assert!(EntryLine::parse(r#"{"at_ms":3}"#).is_ok()); 218 for bad in ["3", "[1]", "nope", "{\"a\":\n1}", "\"s\"", "null"] { 219 assert!(EntryLine::parse(bad).is_err(), "{bad:?}"); 220 } 221 } 222 223 #[test] 224 fn an_entry_is_stored_compact_and_in_the_order_it_was_written() { 225 let e = EntryLine::parse(r#"{ "b": 1, "a": { "z": 1, "y": 2 } }"#).unwrap(); 226 assert_eq!(e.as_str(), r#"{"b":1,"a":{"z":1,"y":2}}"#); 227 let line = StoredLine::of(&DeviceName::new("phone1").unwrap(), &e); 228 assert_eq!(line.as_str(), r#"{"device":"phone1","entry":{"b":1,"a":{"z":1,"y":2}}}"#); 229 } 230 231 #[test] 232 fn a_page_never_starts_past_the_end_and_is_bounded() { 233 let all: Vec<StoredLine> = (0..PULL_LIMIT + 5).map(|i| StoredLine::from_storage(i.to_string())).collect(); 234 let p = Page::of(&all, LogCursor::START); 235 assert_eq!((p.from.get(), p.total.get(), p.lines.len()), (0, all.len() as u64, PULL_LIMIT)); 236 let p = Page::of(&all, LogCursor::new(PULL_LIMIT as u64)); 237 assert_eq!(p.lines.len(), 5); 238 let p = Page::of(&all, LogCursor::new(u64::MAX)); 239 assert_eq!((p.from.get(), p.lines.len()), (all.len() as u64, 0), "clamped to the end"); 240 } 241}