sync.rsannotatedsync.rssource120 lines · 4.4 KB · raw

/memory/sync, /household/sync, /chat/sync: a device sends what it has and gets the hub's converged copy back.

4use ::log::{debug, error, info, warn};
5use whiskers_core::{ChatState, Household, MemorySnapshot};
6use whiskers_ports::{HasChat, HasHousehold, HasMemory, Hub as _};
8use super::{Route, kind_of, sealed};
9use crate::Response;

POST /memory/sync. Needs the memory hub.

12pub struct MemorySync;

POST /household/sync. Needs the household hub.

14pub struct HouseholdSync;

POST /chat/sync. Needs the chat hub.

16pub struct ChatSync;
18impl sealed::Sealed for MemorySync {}
19impl sealed::Sealed for HouseholdSync {}
20impl sealed::Sealed for ChatSync {}
21
22impl<C: HasMemory> Route<C> for MemorySync {
23    const PATHS: &'static [&'static str] = &["/memory/sync"];
24    async fn answer(&self, adapter: &C, _: &str, body: &str) -> Response {
25        memory_sync(adapter, body).await
26    }
27}
28
29impl<C: HasHousehold> Route<C> for HouseholdSync {
30    const PATHS: &'static [&'static str] = &["/household/sync"];
31    async fn answer(&self, adapter: &C, _: &str, body: &str) -> Response {
32        household_sync(adapter, body).await
33    }
34}
35
36impl<C: HasChat> Route<C> for ChatSync {
37    const PATHS: &'static [&'static str] = &["/chat/sync"];
38    async fn answer(&self, adapter: &C, _: &str, body: &str) -> Response {
39        chat_sync(adapter, body).await
40    }
41}
42
43fn bad_request(route: &str, e: &serde_json::Error) -> Response {
44    warn!("{route}: the request body does not parse: {}", kind_of(e));
45    Response::empty(400)
46}
47
48#[cfg(test)]
49mod tests {
50    use crate::routes::kind_of;
51
52    #[test]
53    fn a_parse_failure_is_described_without_quoting_the_body() {
54        let e = serde_json::from_str::<whiskers_core::MemorySnapshot>(r#"{"facts":"a secret thing she said"}"#).unwrap_err();
55        assert!(e.to_string().contains("secret"), "serde quotes it: {e}");
56        let said = kind_of(&e);
57        assert!(!said.contains("secret"), "{said}");
58        assert!(said.contains("Data"), "{said}");
59    }
60}
61
62async fn memory_sync<C: HasMemory>(adapter: &C, body: &str) -> Response {
63    let theirs = match serde_json::from_str::<MemorySnapshot>(body) {
64        Ok(theirs) => theirs,
65        Err(e) => return bad_request("/memory/sync", &e),
66    };
67    info!("/memory/sync: a device sent {} fact(s), {} forgotten", theirs.facts.len(), theirs.forgotten.len());
68    match adapter.memory().merge(theirs).await {
69        Ok(merged) => {
70            info!("/memory/sync: merged, {} change(s)", merged.changed);
71            Response::of(&merged.document)
72        }
73        Err(e) => {
74            error!("/memory/sync: the merge failed: {e}");
75            Response::empty(500)
76        }
77    }
78}
79
80async fn household_sync<C: HasHousehold>(adapter: &C, body: &str) -> Response {
81    let theirs = match serde_json::from_str::<Household>(body) {
82        Ok(theirs) => theirs,
83        Err(e) => return bad_request("/household/sync", &e),
84    };
85    match adapter.household().merge(theirs).await {
86        Ok(merged) => {
87            info!("/household/sync: merged, changed = {}", merged.changed > 0);
88            Response::of(&merged.document)
89        }
90        Err(e) => {
91            error!("/household/sync: the merge failed: {e}");
92            Response::empty(500)
93        }
94    }
95}
96
97async fn chat_sync<C: HasChat>(adapter: &C, body: &str) -> Response {
98    let theirs = match serde_json::from_str::<ChatState>(body) {
99        Ok(theirs) => theirs,
100        Err(e) => return bad_request("/chat/sync", &e),
101    };
102    let their_version = theirs.version;
103    match adapter.chat().merge(theirs).await {
104        Ok(merged) if merged.changed > 0 => {
105            info!("/chat/sync: adopted the device's copy v{their_version}; the service is now at v{}", merged.document.version);
106            Response::of(&merged.document)
107        }
108        Ok(merged) => {
109            debug!("/chat/sync: the service's copy v{} is at least as new as the device's v{their_version}", merged.document.version);
110            Response::of(&merged.document)
111        }
112        Err(e) => {
113            // The device's copy could not be kept. Like every other route, that is an error the
114            // device can tell from "yours was older": the hub's document is unchanged (a hub keeps
115            // before it adopts), and the device will send its copy again at the next sync.
116            error!("/chat/sync: the device's copy v{their_version} could not be kept: {e}");
117            Response::empty(500)
118        }
119    }
120}