hubs.rsannotatedhubs.rssource179 lines · 6.3 KB · raw
1//! The three hubs as files: the master copy of what Whiskers remembers, of the grown-ups' choices and
2//! of the one session. Each is one document behind a mutex (so a merge is atomic: it runs wholly before
3//! or wholly after any other) and replaced atomically on disk (so a power cut leaves the old file or the
4//! new one). The merge rules are the core's; this only keeps the result.
5
6use std::path::{Path, PathBuf};
7use std::sync::{Mutex, MutexGuard, PoisonError};
8
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}
19
20/// What Whiskers remembers, in one JSON file.
21pub struct FileMemoryHub {
22    path: PathBuf,
23    memory: Mutex<JsonMemory>,
24}
25
26impl FileMemoryHub {
27    /// 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    }
43
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}
67
68/// The grown-ups' choices and the day's time, in one JSON file.
69pub struct FileHouseholdHub {
70    path: PathBuf,
71    household: Mutex<TimeKeeper>,
72}
73
74impl FileHouseholdHub {
75    /// 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    }
80
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}
109
110/// 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}
115
116impl FileChatHub {
117    /// Opens the file; one that cannot be read starts the session empty, as it always has (the device
118    /// 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}
144
145impl<S: ChatStore> FileChatHub<S> {
146    /// A hub holding `state`, which is already what `store` holds, and keeping every adopted copy in
147    /// `store`. `path` is only what [`path`](Self::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    }
151
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}