1//! The `Hub` port on a household's Durable Object: a hub that does nothing itself but ask the 2//! household's object to merge, and so is a handful of lines whichever document it is. 3 4use std::marker::PhantomData; 5 6use log::{debug, error}; 7use whiskers_ports::{Hub, Merged, StoreError}; 8use worker::{Method, ObjectNamespace, Request, RequestInit}; 9 10use crate::documents::{Held, Reply}; 11use crate::household::HouseholdId; 12 13/// The Durable Object namespace binding, as `wrangler.jsonc` names it. 14pub const OBJECTS: &str = "HOUSEHOLD_OBJECT"; 15 16/// A household's hub for the document `D`, kept in that household's object. 17pub struct ObjectHub<D> { 18 namespace: ObjectNamespace, 19 household: HouseholdId, 20 document: PhantomData<D>, 21} 22 23impl<D: Held> ObjectHub<D> { 24 pub fn new(namespace: ObjectNamespace, household: HouseholdId) -> Self { 25 Self { namespace, household, document: PhantomData } 26 } 27 28 /// Asks the object `/<document>/<operation>` and returns its body, or why it did not answer. 29 async fn ask(&self, operation: &str, body: String) -> Result<String, StoreError> { 30 let unavailable = |what: &str, e: &dyn std::fmt::Display| { 31 error!("{}: {what}: {e}", D::NAME); 32 StoreError::Unavailable 33 }; 34 let stub = self 35 .namespace 36 .id_from_name(self.household.as_str()) 37 .and_then(|id| id.get_stub()) 38 .map_err(|e| unavailable("the household's object cannot be addressed", &e))?; 39 let mut init = RequestInit::new(); 40 init.with_method(Method::Post).with_body(Some(body.into())); 41 // The host in the address is never resolved: a stub sends the request to its object. 42 let request = Request::new_with_init(&format!("https://household.invalid/{}/{operation}", D::NAME), &init) 43 .map_err(|e| unavailable("the request to the object cannot be made", &e))?; 44 let mut response = stub.fetch_with_request(request).await.map_err(|e| unavailable("the object did not answer", &e))?; 45 let status = response.status_code(); 46 debug!("{}: the object answered {status}", D::NAME); 47 match status { 48 200 => response.text().await.map_err(|e| unavailable("the object's answer cannot be read", &e)), 49 500 => Err(StoreError::Damaged), 50 _ => Err(StoreError::Unavailable), 51 } 52 } 53} 54 55impl<D: Held> Hub for ObjectHub<D> { 56 type Document = D::Wire; 57 58 async fn merge(&self, theirs: D::Wire) -> Result<Merged<D::Wire>, StoreError> { 59 let body = serde_json::to_string(&theirs).map_err(|_| StoreError::Damaged)?; 60 let text = self.ask("merge", body).await?; 61 serde_json::from_str::<Reply<D::Wire>>(&text).map(|r| Merged { document: r.document, changed: r.changed }).map_err(|e| { 62 error!("{}: the object's reply does not read ({:?})", D::NAME, e.classify()); 63 StoreError::Damaged 64 }) 65 } 66 67 async fn current(&self) -> Result<D::Wire, StoreError> { 68 let text = self.ask("current", String::new()).await?; 69 serde_json::from_str(&text).map_err(|e| { 70 error!("{}: the object's document does not read ({:?})", D::NAME, e.classify()); 71 StoreError::Damaged 72 }) 73 } 74} 75 76/// The test mode's window onto a household's object: `/_hub/<document>/<operation>` is passed to the 77/// object as `/<document>/<operation>` and its answer returned as it is, `changed` included, which the 78/// public routes do not carry. Only `lib.rs`'s `Households::is_testing` lets a request here. 79pub async fn through( 80 namespace: &ObjectNamespace, 81 household: &HouseholdId, 82 path: &str, 83 body: String, 84) -> worker::Result<worker::Response> { 85 let stub = namespace.id_from_name(household.as_str())?.get_stub()?; 86 let mut init = RequestInit::new(); 87 init.with_method(Method::Post).with_body(Some(body.into())); 88 stub.fetch_with_request(Request::new_with_init(&format!("https://household.invalid{path}"), &init)?).await 89}