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] callsMemoryDoc::merge,Household::mergeandChatState::mergeand 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 noawaitbetween 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 unlessblockConcurrencyWhile()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.
Where an object keeps its documents: one text per document name. Synchronous on purpose (see the module docs).
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.
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.
The object's answer to a merge, as it crosses from the object to the Worker.
What the object says back to the Worker: a status and a JSON body (empty on a refusal).
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}
A household's object without Cloudflare: a map, which can be told to fail its next write.
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.
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}