hub.rsannotatedhub.rssource89 lines · 3.9 KB · raw

The Hub port on a household's Durable Object: a hub that does nothing itself but ask the household's object to merge, and so is a handful of lines whichever document it is.

4use std::marker::PhantomData;
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;

The Durable Object namespace binding, as wrangler.jsonc names it.

14pub const OBJECTS: &str = "HOUSEHOLD_OBJECT";

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}
23impl<D: Held> ObjectHub<D> {
24    pub fn new(namespace: ObjectNamespace, household: HouseholdId) -> Self {
25        Self { namespace, household, document: PhantomData }
26    }

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}
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}

The test mode's window onto a household's object: /_hub/<document>/<operation> is passed to the object as /<document>/<operation> and its answer returned as it is, changed included, which the 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}