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}