hub.rsannotatedhub.rssource89 lines · 3.9 KB · raw
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}