hub.rsannotatedhub.rssource76 lines · 2.6 KB · raw

A running host's hubs as the Hub port, so the real conformance laws judge a running host.

The laws assert changed, which the wire (POST /memory/sync and the others) does not carry, so they cannot be run through the routes. A Worker started in its test mode (TRUST_HOUSEHOLD_HEADER=1, from .dev.vars) therefore serves /_hub/<document>/merge and /_hub/<document>/current: the household object's own protocol, {"document":..,"changed":n}, with nothing between the request and the object but the household header. Everything else about the hub (the Durable Object, its SQLite, the merge in the core) is the production path.

10use std::marker::PhantomData;
12use serde::de::DeserializeOwned;
13use serde::{Deserialize, Serialize};
14use whiskers_core::{ChatState, Household, MemorySnapshot};
15use whiskers_ports::{Hub, Merged, StoreError};
16
17use crate::host::Host;

A document a household's object holds, by the name it is addressed under.

20pub trait Named: Serialize + DeserializeOwned {
21    const NAME: &'static str;
22}
23impl Named for MemorySnapshot {
24    const NAME: &'static str = "memory";
25}
26impl Named for Household {
27    const NAME: &'static str = "household";
28}
29impl Named for ChatState {
30    const NAME: &'static str = "chat";
31}
33#[derive(Deserialize)]
34struct Reply<D> {
35    document: D,
36    changed: usize,
37}
38
39pub struct HttpHub<D> {
40    host: Host,
41    household: String,
42    document: PhantomData<D>,
43}
44
45impl<D: Named> HttpHub<D> {
46    pub fn new(host: &Host, household: &str) -> Self {
47        Self { host: host.clone(), household: household.to_owned(), document: PhantomData }
48    }
49
50    fn ask(&self, op: &str, body: &str) -> Result<String, StoreError> {
51        let reply = self.host.post(&format!("/_hub/{}/{op}", D::NAME), &self.household, body).map_err(|e| {
52            eprintln!("hostcheck: {e}");
53            StoreError::Unavailable
54        })?;
55        match reply.status {
56            200 => Ok(reply.body),
57            500 => Err(StoreError::Damaged),
58            _ => Err(StoreError::Unavailable),
59        }
60    }
61}
62
63impl<D: Named> Hub for HttpHub<D> {
64    type Document = D;
65
66    async fn merge(&self, theirs: D) -> Result<Merged<D>, StoreError> {
67        let body = serde_json::to_string(&theirs).map_err(|_| StoreError::Damaged)?;
68        let text = self.ask("merge", &body)?;
69        let reply: Reply<D> = serde_json::from_str(&text).map_err(|_| StoreError::Damaged)?;
70        Ok(Merged { document: reply.document, changed: reply.changed })
71    }
72
73    async fn current(&self) -> Result<D, StoreError> {
74        serde_json::from_str(&self.ask("current", "")?).map_err(|_| StoreError::Damaged)
75    }
76}