whiskers.git / backend / worker / src / documents.rs

What a household's object does with the documents it holds, with no Cloudflare in it.

The Durable Object (object.rs) is a thin shell round [serve]: it reads a request, hands over a [DocStore] backed by its SQLite, and writes back the answer. Everything that decides anything is here, in plain Rust that runs (and is tested) natively:

  • the merge is the core's. [Held::absorb] calls MemoryDoc::merge, Household::merge and ChatState::merge and nothing else; no merge rule is written in this crate.
  • the merge is atomic because nothing in it waits. [merge_in] is synchronous: it reads the stored document, merges and writes it back in one turn of the object's single thread, with no await between the read and the write. That is what makes it atomic on Cloudflare and under celld alike (celld's docs: "an event that awaits can overlap another event unless blockConcurrencyWhile() closes the input gate", so an atomic merge must not await).
  • a failed write leaves the stored document as it was, and the caller is told.
16use log::{debug, error, info, warn};
17use serde::de::DeserializeOwned;
18use serde::{Deserialize, Serialize};
19use whiskers_core::{ChatState, Household, MemoryDoc, MemorySnapshot};
20use whiskers_ports::{Merged, StoreError};

Where an object keeps its documents: one text per document name. Synchronous on purpose (see the module docs).

24pub trait DocStore {
25    fn load(&self, name: &str) -> Result<Option<String>, StoreError>;
26    fn save(&self, name: &str, json: &str) -> Result<(), StoreError>;
27}

A document a household's object holds, and how it takes another copy in.

30pub trait Held: Default + Serialize + DeserializeOwned {

The document's name in storage and in the object's paths.

32    const NAME: &'static str;

What a device sends and is sent back: the document on the wire.

34    type Wire: Serialize + DeserializeOwned;

Takes theirs in with the core's rule; returns how many things changed.

36    fn absorb(&mut self, theirs: Self::Wire) -> usize;
37    fn wire(&self) -> Self::Wire;
38}
40impl Held for MemoryDoc {
41    const NAME: &'static str = "memory";
42    type Wire = MemorySnapshot;
43    fn absorb(&mut self, theirs: MemorySnapshot) -> usize {
44        self.merge(theirs)
45    }
46    fn wire(&self) -> MemorySnapshot {
47        self.snapshot()
48    }
49}
50
51impl Held for Household {
52    const NAME: &'static str = "household";
53    type Wire = Household;
54    fn absorb(&mut self, theirs: Household) -> usize {
55        usize::from(self.merge(&theirs))
56    }
57    fn wire(&self) -> Household {
58        self.clone()
59    }
60}
61
62impl Held for ChatState {
63    const NAME: &'static str = "chat";
64    type Wire = ChatState;
65    fn absorb(&mut self, theirs: ChatState) -> usize {
66        usize::from(self.merge(theirs))
67    }
68    fn wire(&self) -> ChatState {
69        self.clone()
70    }
71}
72
73fn load<D: Held>(store: &impl DocStore) -> Result<D, StoreError> {
74    match store.load(D::NAME)? {
75        None => {
76            debug!("{}: nothing stored yet; starting empty", D::NAME);
77            Ok(D::default())
78        }
79        Some(json) => serde_json::from_str(&json).map_err(|e| {
80            // The error text of a parse can quote the document; log the kind and the length only.
81            error!("{}: the stored document does not read back ({:?}, {} bytes)", D::NAME, e.classify(), json.len());
82            StoreError::Damaged
83        }),
84    }
85}

Takes theirs into the stored document and returns the document as it now stands. Writes only when the stored document actually differs afterwards (a merge can change what is kept without changing what it counts, such as an identity forgotten that was never held).

90pub fn merge_in<D: Held>(store: &impl DocStore, theirs: D::Wire) -> Result<Merged<D::Wire>, StoreError> {
91    let mut doc: D = load(store)?;
92    let before = serde_json::to_string(&doc).map_err(|_| StoreError::Damaged)?;
93    let changed = doc.absorb(theirs);
94    let after = serde_json::to_string(&doc).map_err(|_| StoreError::Damaged)?;
95    if after != before {
96        store.save(D::NAME, &after).inspect_err(|_| error!("{}: the merge could not be kept", D::NAME))?;
97        debug!("{}: kept {} bytes", D::NAME, after.len());
98    }
99    Ok(Merged { document: doc.wire(), changed })
100}

The stored document as it stands, changing nothing.

103pub fn current_of<D: Held>(store: &impl DocStore) -> Result<D::Wire, StoreError> {
104    Ok(load::<D>(store)?.wire())
105}

The object's answer to a merge, as it crosses from the object to the Worker.

108#[derive(Debug, Serialize, Deserialize)]
109pub struct Reply<W> {
110    pub document: W,
111    pub changed: usize,
112}

What the object says back to the Worker: a status and a JSON body (empty on a refusal).

115#[derive(Debug, PartialEq, Eq)]
116pub struct Answer {
117    pub status: u16,
118    pub body: String,
119}
121impl Answer {
122    fn empty(status: u16) -> Self {
123        Self { status, body: String::new() }
124    }
125
126    fn of<T: Serialize>(value: &T) -> Self {
127        match serde_json::to_string(value) {
128            Ok(body) => Self { status: 200, body },
129            Err(e) => {
130                error!("a reply does not serialise: {e}");
131                Self::empty(500)
132            }
133        }
134    }
135
136    fn of_store(e: StoreError) -> Self {
137        Self::empty(match e {
138            StoreError::Unavailable => 503,
139            StoreError::Damaged => 500,
140        })
141    }
142}

What an object's request is: /<document>/merge (body: the device's copy) or /<document>/current.

145fn serve_one<D: Held>(store: &impl DocStore, op: &str, body: &str) -> Answer {
146    match op {
147        "merge" => match serde_json::from_str::<D::Wire>(body) {
148            Ok(theirs) => match merge_in::<D>(store, theirs) {
149                Ok(m) => Answer::of(&Reply { document: m.document, changed: m.changed }),
150                Err(e) => Answer::of_store(e),
151            },
152            Err(e) => {
153                warn!("{}: a copy sent to the object does not read ({:?})", D::NAME, e.classify());
154                Answer::empty(400)
155            }
156        },
157        "current" => match current_of::<D>(store) {
158            Ok(doc) => Answer::of(&doc),
159            Err(e) => Answer::of_store(e),
160        },
161        _ => {
162            warn!("{}: no such operation", D::NAME);
163            Answer::empty(404)
164        }
165    }
166}

Answers one request to a household's object. path is /<document>/<operation>.

169pub fn serve(store: &impl DocStore, path: &str, body: &str) -> Answer {
170    let mut parts = path.trim_start_matches('/').split('/');
171    let (Some(name), Some(op), None) = (parts.next(), parts.next(), parts.next()) else {
172        warn!("the object was asked for a path of the wrong shape");
173        return Answer::empty(404);
174    };
175    info!("object: {name} {op}");
176    match name {
177        MemoryDoc::NAME => serve_one::<MemoryDoc>(store, op, body),
178        Household::NAME => serve_one::<Household>(store, op, body),
179        ChatState::NAME => serve_one::<ChatState>(store, op, body),
180        _ => {
181            warn!("the object holds no such document");
182            Answer::empty(404)
183        }
184    }
185}
187#[cfg(test)]
188pub(crate) mod testing {
189    use std::cell::RefCell;
190    use std::collections::BTreeMap;
191
192    use super::*;

A household's object without Cloudflare: a map, which can be told to fail its next write.

195    #[derive(Default)]
196    pub struct MapStore {
197        pub docs: RefCell<BTreeMap<String, String>>,
198        pub fail_writes: std::cell::Cell<bool>,
199    }
201    impl DocStore for MapStore {
202        fn load(&self, name: &str) -> Result<Option<String>, StoreError> {
203            Ok(self.docs.borrow().get(name).cloned())
204        }
205        fn save(&self, name: &str, json: &str) -> Result<(), StoreError> {
206            if self.fail_writes.get() {
207                return Err(StoreError::Unavailable);
208            }
209            self.docs.borrow_mut().insert(name.to_owned(), json.to_owned());
210            Ok(())
211        }
212    }
213}
214
215#[cfg(test)]
216mod tests {
217    use whiskers_conformance::suite::build::{fact, snapshot as snap};
218    use whiskers_conformance::suite::fixture::{ChatFaultFixture, ChatFixture, HouseholdFixture, MemoryFixture, Suite};
219    use whiskers_conformance::suite::hubs;
220    use whiskers_ports::{Hub, run_ready};
221
222    use super::testing::MapStore;
223    use super::*;

The shared logic as a Hub, over a map instead of an object's SQLite: so the real conformance laws judge merge_in natively, before any host is involved.

227    #[derive(Clone, Default)]
228    struct StoreHub<D>(std::rc::Rc<MapStore>, std::marker::PhantomData<D>);
230    impl<D: Held> Hub for StoreHub<D> {
231        type Document = D::Wire;
232        async fn merge(&self, theirs: D::Wire) -> Result<Merged<D::Wire>, StoreError> {
233            merge_in::<D>(&*self.0, theirs)
234        }
235        async fn current(&self) -> Result<D::Wire, StoreError> {
236            current_of::<D>(&*self.0)
237        }
238    }

The Worker's conformance harness: the three hubs, and nothing else. It has no fixture for the journal, the pictures or the allowance, so the suite has no way to select their laws for it.

242    struct Hubs;
244    macro_rules! hub_fixture {
245        ($fixture:ident, $document:ty) => {
246            impl $fixture for Hubs {
247                type Port = StoreHub<$document>;
248                fn fresh(&self) -> StoreHub<$document> {
249                    StoreHub::default()
250                }
251                fn restart(&self, hub: StoreHub<$document>) -> StoreHub<$document> {
252                    hub
253                }
254            }
255        };
256    }
257
258    hub_fixture!(MemoryFixture, MemoryDoc);
259    hub_fixture!(HouseholdFixture, Household);
260    hub_fixture!(ChatFixture, ChatState);
261
262    impl ChatFaultFixture for Hubs {
263        type Port = StoreHub<ChatState>;
264        fn fresh_refusing(&self) -> (StoreHub<ChatState>, hubs::WriteSwitch) {
265            let hub = StoreHub::<ChatState>::default();
266            let store = hub.0.clone();
267            (hub, Box::new(move |refuse| store.fail_writes.set(refuse)))
268        }
269        fn restart_refusing(&self, hub: StoreHub<ChatState>) -> StoreHub<ChatState> {
270            hub
271        }
272    }
273
274    #[test]
275    fn the_shared_logic_obeys_every_law_of_the_hubs_and_runs_no_other() {
276        let suite = Suite::new(Hubs).memory().household().chat().chat_write_faults();
277        assert_eq!(suite.ran(), ["memory", "household", "chat", "chat write faults"]);
278    }
279
280    #[test]
281    fn a_merge_never_waits() {
282        // The atomicity argument in the module docs: the whole merge is one synchronous turn.
283        let hub = StoreHub::<MemoryDoc>::default();
284        run_ready(hub.merge(snap(vec![fact("g1", "a bunny")], &[]))).unwrap();
285    }
286
287    #[test]
288    fn a_merge_is_kept_and_read_back() {
289        let store = MapStore::default();
290        let m = merge_in::<MemoryDoc>(&store, snap(vec![fact("g1", "a bunny")], &[])).unwrap();
291        assert_eq!((m.changed, m.document.facts.len()), (1, 1));
292        assert_eq!(current_of::<MemoryDoc>(&store).unwrap(), m.document);
293        assert!(store.docs.borrow().contains_key("memory"));
294    }
295
296    #[test]
297    fn a_merge_that_changes_nothing_writes_nothing() {
298        let store = MapStore::default();
299        merge_in::<MemoryDoc>(&store, snap(vec![fact("g1", "a bunny")], &[])).unwrap();
300        store.fail_writes.set(true);
301        let again = merge_in::<MemoryDoc>(&store, snap(vec![fact("g1", "a bunny")], &[])).unwrap();
302        assert_eq!(again.changed, 0, "a repeat needs no write, so a failing store is not asked");
303    }
304
305    #[test]
306    fn a_forgotten_identity_never_held_is_still_kept() {
307        // The count says nothing changed, yet the document did: it must be written.
308        let store = MapStore::default();
309        let m = merge_in::<MemoryDoc>(&store, snap(vec![], &["g9"])).unwrap();
310        assert_eq!(m.changed, 0);
311        assert!(m.document.forgotten.contains(&"g9".to_owned()));
312        assert_eq!(current_of::<MemoryDoc>(&store).unwrap().forgotten, ["g9"]);
313    }
314
315    #[test]
316    fn a_failed_write_leaves_the_document_as_it_was_and_says_so() {
317        let store = MapStore::default();
318        let first = merge_in::<MemoryDoc>(&store, snap(vec![fact("g1", "a bunny")], &[])).unwrap();
319        store.fail_writes.set(true);
320        let err = merge_in::<MemoryDoc>(&store, snap(vec![fact("g2", "purple")], &[])).unwrap_err();
321        assert_eq!(err, StoreError::Unavailable);
322        store.fail_writes.set(false);
323        assert_eq!(current_of::<MemoryDoc>(&store).unwrap(), first.document, "g2 was not adopted");
324    }
325
326    #[test]
327    fn a_stored_document_that_does_not_read_is_damaged_and_is_not_overwritten() {
328        let store = MapStore::default();
329        store.docs.borrow_mut().insert("memory".into(), "{not json".into());
330        assert_eq!(current_of::<MemoryDoc>(&store).unwrap_err(), StoreError::Damaged);
331        assert_eq!(merge_in::<MemoryDoc>(&store, snap(vec![fact("g1", "x")], &[])).unwrap_err(), StoreError::Damaged);
332        assert_eq!(store.docs.borrow()["memory"], "{not json", "a person has to look; nothing is destroyed");
333    }
334
335    #[test]
336    fn the_object_answers_by_status() {
337        let store = MapStore::default();
338        let body = serde_json::to_string(&snap(vec![fact("g1", "a bunny")], &[])).unwrap();
339        let ok = serve(&store, "/memory/merge", &body);
340        assert_eq!(ok.status, 200);
341        let reply: Reply<MemorySnapshot> = serde_json::from_str(&ok.body).unwrap();
342        assert_eq!((reply.changed, reply.document.facts.len()), (1, 1));
343        assert_eq!(serve(&store, "/memory/merge", "nonsense").status, 400);
344        assert_eq!(serve(&store, "/memory/nope", "").status, 404);
345        assert_eq!(serve(&store, "/journal/merge", "").status, 404);
346        assert_eq!(serve(&store, "/memory", "").status, 404);
347        assert_eq!(serve(&store, "/memory/merge/extra", "").status, 404);
348        store.fail_writes.set(true);
349        let other = serde_json::to_string(&snap(vec![fact("g2", "purple")], &[])).unwrap();
350        assert_eq!(serve(&store, "/memory/merge", &other).status, 503);
351        store.docs.borrow_mut().insert("household".into(), "x".into());
352        assert_eq!(serve(&store, "/household/current", "").status, 500);
353    }
354
355    #[test]
356    fn the_three_documents_do_not_touch_each_other() {
357        let store = MapStore::default();
358        let chat = ChatState { summary: "s".into(), version: 3, ..ChatState::default() };
359        merge_in::<ChatState>(&store, chat.clone()).unwrap();
360        merge_in::<MemoryDoc>(&store, snap(vec![fact("g1", "a bunny")], &[])).unwrap();
361        assert_eq!(current_of::<ChatState>(&store).unwrap(), chat);
362        assert_eq!(current_of::<Household>(&store).unwrap(), Household::default());
363    }
364}