hub.rsannotatedhub.rssource76 lines · 2.6 KB · raw
1//! A running host's hubs as the `Hub` port, so the real conformance laws judge a running host.
2//!
3//! The laws assert `changed`, which the wire (`POST /memory/sync` and the others) does not carry, so they
4//! cannot be run through the routes. A Worker started in its test mode (`TRUST_HOUSEHOLD_HEADER=1`, from
5//! `.dev.vars`) therefore serves `/_hub/<document>/merge` and `/_hub/<document>/current`: the household
6//! object's own protocol, `{"document":..,"changed":n}`, with nothing between the request and the object
7//! but the household header. Everything else about the hub (the Durable Object, its SQLite, the merge in
8//! the core) is the production path.
9
10use std::marker::PhantomData;
11
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;
18
19/// 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}
32
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}