The journal: the one append-only log of everything said to Whiskers, from every device, in the order the hub received it. Parents read it as a single conversation wherever it happened.

Nothing in it is ever rewritten, so it cannot conflict. A device tells the hub how many of its own lines the hub already holds and sends the rest; any device reads the whole log from the position it last reached. Two different counts are in play and they are different types on purpose: a [DeviceCursor] counts one device's lines, a [LogCursor] is a position in the whole log.

9use ::log::trace;
10use serde::{Deserialize, Serialize};
12use crate::error::StoreError;

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;

At most this many lines go back in one pull; the caller asks again for the rest.

18pub const PULL_LIMIT: usize = 2000;

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);
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}

One entry a device wrote: a single JSON object on a single line. The text is kept compact, in the key order the device wrote it, so what is stored is what was sent.

43#[derive(Clone, Debug, PartialEq, Eq)]
44pub struct EntryLine(String);
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}

A line of the log as stored and as pulled: {"device":"..","entry":{..}}. Built in one place so every backend stores the same bytes.

67#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
68pub struct StoredLine(String);
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    }

A line read back from storage that of produced. Adapters use this to hand lines out again; nothing else should have a reason to make one.

78    pub fn from_storage(raw: String) -> Self {
79        Self(raw)
80    }
82    pub fn as_str(&self) -> &str {
83        &self.0
84    }
85
86    pub fn into_string(self) -> String {
87        self.0
88    }
89}

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);
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}

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);
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}

What an append did.

127#[derive(Clone, Copy, Debug, PartialEq, Eq)]
128pub enum Appended {

The lines were taken (or there were none), and the journal now holds this many of the device's.

130    Accepted { held: DeviceCursor },

The device thought the journal held a different number than it does, so nothing was taken. The device is told the real count and sends again from there. This is how a device that was behind (a lost reply, a restored backup) or ahead recovers without duplicating or skipping a line.

134    Misaligned { held: DeviceCursor },
135}
137impl Appended {
138    pub fn held(self) -> DeviceCursor {
139        match self {
140            Appended::Accepted { held } | Appended::Misaligned { held } => held,
141        }
142    }
143}

A slice of the log.

146#[derive(Clone, Debug, PartialEq, Eq)]
147pub struct Page {

Where lines begin. Never past total: a caller that thought it had read further than the log goes (its copy is longer than the log, as a damaged copy is) is told where the log really ends, and sees from is less than what it held, and rebuilds its copy from here.

151    pub from: LogCursor,

The whole log's length.

153    pub total: LogCursor,
154    pub lines: Vec<StoredLine>,
155}
157impl Page {

The page since asks for from a log of all. One definition, so no backend can read the rule 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}
168impl Page {

Whether the caller's copy of the log is longer than the log: it held held lines and this page begins before that. The only reason a page can start earlier than asked is that the log is shorter than the caller believed (a duplicated copy, say), so its copy is wrong 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}

The journal.

Contract, for every adapter:

  • Append-only. A stored line is never changed, moved or removed.
  • Atomic per call. An append either stores all its lines or none, and two appends overlapping (even from one device, even at the same cursor) behave as if one ran first: exactly one is Accepted; the other sees the first's lines and is Misaligned. No line is ever duplicated.
  • Aligned only. Lines are taken only when at equals what the journal holds for that device. Appending no lines is a question, and is always Accepted { held }.
  • Order is arrival order, across devices.
  • 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 {

How many of device's lines the journal holds.

193    async fn held(&self, device: &DeviceName) -> Result<DeviceCursor, StoreError>;
195    async fn append(&self, device: &DeviceName, at: DeviceCursor, lines: &[EntryLine]) -> Result<Appended, StoreError>;

The log from since on, as [Page::of] defines it.

198    async fn pull(&self, since: LogCursor) -> Result<Page, StoreError>;
199}
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}