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}
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.
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.
A hub keeping its document in store, which a test may make fail.
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.)
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.
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}