hubs.rsannotatedhubs.rssource179 lines · 6.3 KB · raw

The three hubs as files: 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 replaced atomically on disk (so a power cut leaves the old file or the new one). The merge rules are the core's; this only keeps the result.

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

What Whiskers remembers, in one JSON file.

21pub struct FileMemoryHub {
22    path: PathBuf,
23    memory: Mutex<JsonMemory>,
24}
26impl FileMemoryHub {

Opens the file, moving a damaged one aside and starting empty (the service must be able to start).

28    pub fn open(path: &Path) -> Result<Self, StoreError> {
29        match JsonMemory::open_or_set_aside(path) {
30            Ok((m, aside)) => {
31                if let Some(a) = aside {
32                    warn!("the shared memory was damaged; kept at {}", a.display());
33                }
34                info!("shared memory opened at {}", path.display());
35                Ok(Self { path: path.to_owned(), memory: Mutex::new(m) })
36            }
37            Err(e) => {
38                error!("cannot open the shared memory at {}: {}", path.display(), e.0);
39                Err(StoreError::Unavailable)
40            }
41        }
42    }
44    pub fn path(&self) -> &Path {
45        &self.path
46    }
47}
48
49impl Hub for FileMemoryHub {
50    type Document = MemorySnapshot;
51
52    async fn merge(&self, theirs: MemorySnapshot) -> Result<Merged<MemorySnapshot>, StoreError> {
53        let mut hub = lock(&self.memory);
54        match hub.merge(theirs) {
55            Ok(changed) => Ok(Merged { document: hub.snapshot(), changed }),
56            Err(e) => {
57                error!("memory hub: the merge failed: {}", e.0);
58                Err(StoreError::Unavailable)
59            }
60        }
61    }
62
63    async fn current(&self) -> Result<MemorySnapshot, StoreError> {
64        Ok(lock(&self.memory).snapshot())
65    }
66}

The grown-ups' choices and the day's time, in one JSON file.

69pub struct FileHouseholdHub {
70    path: PathBuf,
71    household: Mutex<TimeKeeper>,
72}
74impl FileHouseholdHub {

A missing or unreadable file starts with the default choices (see TimeKeeper::open).

76    pub fn open(path: &Path) -> Self {
77        debug!("household hub at {}", path.display());
78        Self { path: path.to_owned(), household: Mutex::new(TimeKeeper::open(path, "service")) }
79    }
81    pub fn path(&self) -> &Path {
82        &self.path
83    }
84}
85
86impl Hub for FileHouseholdHub {
87    type Document = Household;
88
89    async fn merge(&self, theirs: Household) -> Result<Merged<Household>, StoreError> {
90        let mut hub = lock(&self.household);
91        match hub.merge_durably(&theirs) {
92            Ok(changed) => Ok(Merged { document: hub.household(), changed: usize::from(changed) }),
93            Err(e) => {
94                error!("household hub: the merge could not be kept: {e}");
95                Err(StoreError::Unavailable)
96            }
97        }
98    }
99
100    async fn current(&self) -> Result<Household, StoreError> {
101        Ok(lock(&self.household).household())
102    }
103}
104
105struct ChatInner<S> {
106    state: ChatState,
107    store: S,
108}

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

111pub struct FileChatHub<S: ChatStore = JsonChat> {
112    path: PathBuf,
113    chat: Mutex<ChatInner<S>>,
114}
116impl FileChatHub {

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

119    pub fn open(path: &Path) -> Self {
120        let state = match JsonChat::open_or_set_aside(path) {
121            Ok((c, aside)) => {
122                if let Some(a) = aside {
123                    warn!("the shared chat was damaged; kept at {}", a.display());
124                }
125                match c.load() {
126                    Ok(state) => {
127                        info!("shared chat opened: version {}, {} turns", state.version, state.turns.len());
128                        state
129                    }
130                    Err(e) => {
131                        error!("cannot load the shared chat: {}; starting empty", e.0);
132                        ChatState::default()
133                    }
134                }
135            }
136            Err(e) => {
137                error!("cannot open the shared chat at {}: {}; starting empty", path.display(), e.0);
138                ChatState::default()
139            }
140        };
141        Self::with_store(path, state, JsonChat::open(path))
142    }
143}
145impl<S: ChatStore> FileChatHub<S> {

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

148    pub fn with_store(path: &Path, state: ChatState, store: S) -> Self {
149        Self { path: path.to_owned(), chat: Mutex::new(ChatInner { state, store }) }
150    }
152    pub fn path(&self) -> &Path {
153        &self.path
154    }
155}
156
157impl<S: ChatStore> Hub for FileChatHub<S> {
158    type Document = ChatState;
159
160    async fn merge(&self, theirs: ChatState) -> Result<Merged<ChatState>, StoreError> {
161        let mut inner = lock(&self.chat);
162        if !theirs.is_newer_than(&inner.state) {
163            trace!("chat hub: the service's copy v{} is at least as new as v{}", inner.state.version, theirs.version);
164            return Ok(Merged { document: inner.state.clone(), changed: 0 });
165        }
166        debug!("chat hub: adopting v{} over v{}", theirs.version, inner.state.version);
167        // Kept first, adopted after: what is in memory and what is on disk must not disagree.
168        if let Err(e) = inner.store.save(&theirs) {
169            error!("chat hub: the device's copy could not be saved: {}", e.0);
170            return Err(StoreError::Unavailable);
171        }
172        inner.state = theirs;
173        Ok(Merged { document: inner.state.clone(), changed: 1 })
174    }
175
176    async fn current(&self) -> Result<ChatState, StoreError> {
177        Ok(lock(&self.chat).state.clone())
178    }
179}