whiskers.git / backend / worker / src / documents.rs
1//! What a household's object does with the documents it holds, with no Cloudflare in it.
2//!
3//! The Durable Object (`object.rs`) is a thin shell round [`serve`]: it reads a request, hands over a
4//! [`DocStore`] backed by its SQLite, and writes back the answer. Everything that decides anything is
5//! here, in plain Rust that runs (and is tested) natively:
6//!
7//! - **the merge is the core's.** [`Held::absorb`] calls `MemoryDoc::merge`, `Household::merge` and
8//!   `ChatState::merge` and nothing else; no merge rule is written in this crate.
9//! - **the merge is atomic because nothing in it waits.** [`merge_in`] is synchronous: it reads the
10//!   stored document, merges and writes it back in one turn of the object's single thread, with no
11//!   `await` between the read and the write. That is what makes it atomic on Cloudflare and under
12//!   celld alike (celld's docs: "an event that awaits can overlap another event unless
13//!   `blockConcurrencyWhile()` closes the input gate", so an atomic merge must not await).
14//! - **a failed write leaves the stored document as it was**, and the caller is told.
15
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};
21
22/// Where an object keeps its documents: one text per document name. Synchronous on purpose (see the
23/// 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}
28
29/// A document a household's object holds, and how it takes another copy in.
30pub trait Held: Default + Serialize + DeserializeOwned {
31    /// The document's name in storage and in the object's paths.
32    const NAME: &'static str;
33    /// What a device sends and is sent back: the document on the wire.
34    type Wire: Serialize + DeserializeOwned;
35    /// 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}
39
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}
86
87/// Takes `theirs` into the stored document and returns the document as it now stands. Writes only
88/// when the stored document actually differs afterwards (a merge can change what is kept without
89/// 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}
101
102/// 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}
106
107/// 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}
113
114/// 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}
120
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}
143
144/// 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}
167
168/// 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}
186
187#[cfg(test)]
188pub(crate) mod testing {
189    use std::cell::RefCell;
190    use std::collections::BTreeMap;
191
192    use super::*;
193
194    /// 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    }
200
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::*;
224
225    /// The shared logic as a `Hub`, over a map instead of an object's SQLite: so the real conformance
226    /// 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>);
229
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    }
239
240    /// The Worker's conformance harness: the three hubs, and nothing else. It has no fixture for the journal,
241    /// the pictures or the allowance, so the suite has no way to select their laws for it.
242    struct Hubs;
243
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}