hubs.rsannotatedhubs.rssource218 lines · 8.7 KB · raw

The three hubs of the service in one database: the master copy of what Whiskers remembers, of the grown-ups' choices and of the one session. Each is one document behind a mutex (so a merge is atomic: it runs wholly before or wholly after any other) and written in one transaction (so a power cut leaves the old state or the new one). The merge rules are the core's; this only keeps the result.

6use std::sync::{Mutex, MutexGuard, PoisonError};
8use ::log::{debug, error, info, trace};
9use whiskers_core::{ChatState, ChatStore, Household, HouseholdStore, Memory as _, MemorySnapshot, TimeKeeper};
10use whiskers_ports::{Hub, Merged, StoreError};
11
12use crate::{SqliteChat, SqliteHousehold, SqliteMemory, Store};
13
14fn lock<T>(m: &Mutex<T>) -> MutexGuard<'_, T> {
15    // A panic elsewhere must not stop the hub.
16    m.lock().unwrap_or_else(PoisonError::into_inner)
17}

What Whiskers remembers.

20pub struct SqliteMemoryHub {
21    memory: Mutex<SqliteMemory>,
22}
24impl SqliteMemoryHub {

Opens the memory kept in store. What cannot be read is an error: a hub that started empty would then overwrite it.

27    pub fn open(store: Store) -> Result<Self, StoreError> {
28        match SqliteMemory::open(store, "service") {
29            Ok(m) => {
30                info!("shared memory opened");
31                Ok(Self { memory: Mutex::new(m) })
32            }
33            Err(e) => {
34                error!("cannot open the shared memory: {}", e.0);
35                Err(StoreError::Unavailable)
36            }
37        }
38    }
39}
41impl Hub for SqliteMemoryHub {
42    type Document = MemorySnapshot;
43
44    async fn merge(&self, theirs: MemorySnapshot) -> Result<Merged<MemorySnapshot>, StoreError> {
45        let mut hub = lock(&self.memory);
46        match hub.merge(theirs) {
47            Ok(changed) => Ok(Merged { document: hub.snapshot(), changed }),
48            Err(e) => {
49                error!("memory hub: the merge failed: {}", e.0);
50                Err(StoreError::Unavailable)
51            }
52        }
53    }
54
55    async fn current(&self) -> Result<MemorySnapshot, StoreError> {
56        Ok(lock(&self.memory).snapshot())
57    }
58}

The grown-ups' choices and the day's time.

61pub struct SqliteHouseholdHub {
62    household: Mutex<TimeKeeper>,
63}
65impl SqliteHouseholdHub {

A database with no choices starts with the default ones; one that cannot be read does too (see TimeKeeper::new), and says so in the log.

68    pub fn open(store: Store) -> Self {
69        debug!("household hub opening");
70        Self::with_store(Box::new(SqliteHousehold::new(store)))
71    }

A hub keeping its document in store, which a test may make fail.

74    pub fn with_store(store: Box<dyn HouseholdStore>) -> Self {
75        Self { household: Mutex::new(TimeKeeper::new(store, "service")) }
76    }
77}
79impl Hub for SqliteHouseholdHub {
80    type Document = Household;
81
82    async fn merge(&self, theirs: Household) -> Result<Merged<Household>, StoreError> {
83        let mut hub = lock(&self.household);
84        match hub.merge_durably(&theirs) {
85            Ok(changed) => Ok(Merged { document: hub.household(), changed: usize::from(changed) }),
86            Err(e) => {
87                error!("household hub: the merge could not be kept: {}", e.0);
88                Err(StoreError::Unavailable)
89            }
90        }
91    }
92
93    async fn current(&self) -> Result<Household, StoreError> {
94        Ok(lock(&self.household).household())
95    }
96}
97
98struct ChatInner<S> {
99    state: ChatState,
100    store: S,
101}

The one session. (S is where a copy is kept: the database, unless a test breaks it.)

104pub struct SqliteChatHub<S: ChatStore = SqliteChat> {
105    chat: Mutex<ChatInner<S>>,
106}
108impl SqliteChatHub {

Opens the session kept in store; one that cannot be read starts empty, as it always has (the device that has the session will bring it back, since it is the newer).

111    pub fn open(store: Store) -> Self {
112        let chat = SqliteChat::new(store);
113        let state = chat.load().unwrap_or_else(|e| {
114            error!("cannot load the shared chat: {}; starting empty", e.0);
115            ChatState::default()
116        });
117        info!("shared chat opened: version {}, {} turns", state.version, state.turns.len());
118        Self::with_store(state, chat)
119    }
120}
122impl<S: ChatStore> SqliteChatHub<S> {

A hub holding state, which is already what store holds, and keeping every adopted copy in store.

124    pub fn with_store(state: ChatState, store: S) -> Self {
125        Self { chat: Mutex::new(ChatInner { state, store }) }
126    }
127}
129impl<S: ChatStore> Hub for SqliteChatHub<S> {
130    type Document = ChatState;
131
132    async fn merge(&self, theirs: ChatState) -> Result<Merged<ChatState>, StoreError> {
133        let mut inner = lock(&self.chat);
134        // The core's rule: the newer copy wins whole, unless it was not emptied for a forgetting the other was, in which
135        // case what the hub keeps is an empty session (see `ChatState::merge`).
136        let mut merged = inner.state.clone();
137        if !merged.merge(theirs) {
138            trace!("chat hub: the service's copy v{} is at least as new as the device's", inner.state.version);
139            return Ok(Merged { document: inner.state.clone(), changed: 0 });
140        }
141        debug!("chat hub: taking v{} over v{}", merged.version, inner.state.version);
142        // Kept first, adopted after: what is in memory and what is stored must not disagree.
143        if let Err(e) = inner.store.save(&merged) {
144            error!("chat hub: the merged copy could not be saved: {}", e.0);
145            return Err(StoreError::Unavailable);
146        }
147        inner.state = merged;
148        Ok(Merged { document: inner.state.clone(), changed: 1 })
149    }
150
151    async fn current(&self) -> Result<ChatState, StoreError> {
152        Ok(lock(&self.chat).state.clone())
153    }
154}
155
156#[cfg(test)]
157mod tests {
158    use super::*;
159    use whiskers_core::{HouseholdConfig, TimeKeeper};
160    use whiskers_ports::run_ready;
161
162    fn writes(store: &Store) -> u64 {
163        store.lock().total_changes()
164    }

The service reports a change, and writes, only when what it was sent differs from what it holds. A device's counter of the time it was used rises while it is used, and that is a real change; the same document sent again (which is what every sync is when nobody used the app) is not.

169    #[test]
170    fn the_same_household_sent_again_reports_no_change_and_writes_nothing() {
171        let store = Store::open_in_memory().unwrap();
172        let hub = SqliteHouseholdHub::open(store.clone());
173        let mut phone = TimeKeeper::new(Box::new(SqliteHousehold::new(Store::open_in_memory().unwrap())), "phone");
174        phone.set_config(HouseholdConfig::default(), 100);
175        phone.tick(20_000, 3_000);
176
177        let first = run_ready(hub.merge(phone.household())).unwrap();
178        assert_eq!(first.changed, 1, "the hub had none of it");
179        let after_first = writes(&store);
180
181        let second = run_ready(hub.merge(phone.household())).unwrap();
182        assert_eq!(second.changed, 0, "the same input a second time");
183        assert_eq!(writes(&store), after_first, "and nothing was written");
184        // The phone takes the hub's copy back (which is what every sync does) and sends that: still nothing.
185        phone.merge(&second.document);
186        let third = run_ready(hub.merge(phone.household())).unwrap();
187        assert_eq!((third.changed, writes(&store)), (0, after_first));
188
189        // Time spent is a real change, and is written once.
190        phone.tick(20_000, 2_000);
191        let fourth = run_ready(hub.merge(phone.household())).unwrap();
192        assert_eq!(fourth.changed, 1);
193        let after_fourth = writes(&store);
194        assert!(after_fourth > after_first);
195        assert_eq!(run_ready(hub.merge(phone.household())).unwrap().changed, 0);
196        assert_eq!(writes(&store), after_fourth);
197    }
199    #[test]
200    fn what_the_hub_reads_back_from_its_database_is_what_it_was_sent() {
201        let path = std::env::temp_dir().join(format!("whiskers-hub-readback-{}.db", std::process::id()));
202        let _ = std::fs::remove_file(&path);
203        let mut phone = TimeKeeper::new(Box::new(SqliteHousehold::new(Store::open_in_memory().unwrap())), "phone");
204        phone.set_config(HouseholdConfig::default(), 100);
205        phone.set_utc_offset(-240, 150);
206        phone.tick(20_000, 1_000);
207        {
208            let hub = SqliteHouseholdHub::open(Store::open(&path).unwrap());
209            assert_eq!(run_ready(hub.merge(phone.household())).unwrap().changed, 1);
210        }
211        // A service that restarted holds what it was sent, to the letter: the next sync is not a change.
212        let store = Store::open(&path).unwrap();
213        let hub = SqliteHouseholdHub::open(store.clone());
214        let before = writes(&store);
215        assert_eq!(run_ready(hub.merge(phone.household())).unwrap().changed, 0);
216        assert_eq!(writes(&store), before);
217    }
218}