mj_controller/database/
writer.rs1use super::*;
2
3pub(super) const DATABASE_WRITE_QUEUE_CAPACITY: usize = 256;
4
5pub(super) type DatabaseWriteJob = Box<dyn FnOnce(Result<&mut Connection>) + Send + 'static>;
9
10pub(super) enum DatabaseWriterMessage {
11 Run {
12 label: &'static str,
13 job: DatabaseWriteJob,
14 },
15 Shutdown,
16}
17
18#[derive(Clone)]
24pub struct DatabaseWriter {
25 pub(super) id: u64,
26 pub(super) sender: SyncSender<DatabaseWriterMessage>,
27}
28
29impl std::fmt::Debug for DatabaseWriter {
30 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
31 formatter
32 .debug_struct("DatabaseWriter")
33 .field("id", &self.id)
34 .finish_non_exhaustive()
35 }
36}
37
38impl DatabaseWriter {
39 pub(super) fn execute<T, F>(&self, label: &'static str, operation: F) -> Result<T>
40 where
41 T: Send + 'static,
42 F: FnOnce(&mut Connection) -> Result<T> + Send + 'static,
43 {
44 let (reply_tx, reply_rx) = sync_channel(1);
45 self.sender
46 .send(DatabaseWriterMessage::Run {
47 label,
48 job: Box::new(move |connection| {
49 let reply = match connection {
50 Ok(connection) => operation(connection),
51 Err(error) => Err(error),
55 };
56 let _ = reply_tx.send(reply);
57 }),
58 })
59 .map_err(|_| {
60 anyhow::anyhow!("submit database writer operation {label}: writer stopped")
61 })?;
62 reply_rx
63 .recv()
64 .with_context(|| format!("database writer stopped during {label}"))?
65 }
66}
67
68pub struct DatabaseWriterOwner {
74 pub(super) writer: DatabaseWriter,
75 pub(super) thread: Option<JoinHandle<()>>,
76 pub(super) stopped: Receiver<Result<()>>,
77}
78
79impl std::fmt::Debug for DatabaseWriterOwner {
80 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
81 formatter
82 .debug_struct("DatabaseWriterOwner")
83 .field("writer", &self.writer)
84 .finish_non_exhaustive()
85 }
86}
87
88impl DatabaseWriterOwner {
89 pub fn shutdown(mut self) -> Result<()> {
90 self.shutdown_inner()
91 }
92
93 pub(super) fn shutdown_inner(&mut self) -> Result<()> {
94 if self.thread.is_none() {
95 return Ok(());
96 }
97 clear_database_writer(self.writer.id);
98 let send_result = self.writer.sender.send(DatabaseWriterMessage::Shutdown);
99 let worker_result = self
100 .stopped
101 .recv()
102 .context("database writer stopped without reporting its result")?;
103 let join_result = self
104 .thread
105 .take()
106 .expect("database writer thread checked above")
107 .join();
108 if let Err(panic) = join_result {
109 std::panic::resume_unwind(panic);
110 }
111 match (send_result, worker_result) {
112 (_, Err(error)) => Err(error),
113 (Err(_), Ok(())) => bail!("request database writer shutdown: writer stopped"),
114 (Ok(()), Ok(())) => Ok(()),
115 }
116 }
117}
118
119impl Drop for DatabaseWriterOwner {
120 fn drop(&mut self) {
121 if let Err(error) = self.shutdown_inner() {
122 tracing::error!(%error, "database writer did not shut down cleanly");
123 }
124 }
125}
126
127pub(super) fn database_writer_slot() -> &'static Mutex<Option<DatabaseWriter>> {
128 static WRITER: OnceLock<Mutex<Option<DatabaseWriter>>> = OnceLock::new();
129 WRITER.get_or_init(|| Mutex::new(None))
130}
131
132pub(crate) fn database_writer_installed() -> bool {
136 database_writer_slot()
137 .lock()
138 .unwrap_or_else(PoisonError::into_inner)
139 .is_some()
140}
141
142pub(super) fn clear_database_writer(id: u64) {
143 let mut installed = database_writer_slot()
144 .lock()
145 .unwrap_or_else(PoisonError::into_inner);
146 if installed.as_ref().is_some_and(|writer| writer.id == id) {
147 *installed = None;
148 }
149}
150
151#[doc(hidden)]
167#[must_use = "the writer stops when this owner is dropped"]
168pub fn install_isolated_test_writer() -> DatabaseWriterOwner {
169 start_database_writer().expect("install the writer for an isolated test child")
170}
171
172pub fn start_database_writer() -> Result<DatabaseWriterOwner> {
173 start_database_writer_at(&database_path(), true)
174}
175
176pub(super) fn start_database_writer_at(
177 path: &Path,
178 install_globally: bool,
179) -> Result<DatabaseWriterOwner> {
180 static NEXT_WRITER_ID: AtomicU64 = AtomicU64::new(1);
181
182 let connection = schema::open_writer(path)?;
183 let mut observed_revision = schema::read_schema_state(&connection)?.revision;
184 let path = path.to_owned();
185 let (sender, receiver) = sync_channel(DATABASE_WRITE_QUEUE_CAPACITY);
186 let (stopped_tx, stopped) = sync_channel(1);
187 let id = NEXT_WRITER_ID.fetch_add(1, Ordering::Relaxed);
188 let writer = DatabaseWriter { id, sender };
189 if install_globally {
190 let mut installed = database_writer_slot()
191 .lock()
192 .unwrap_or_else(PoisonError::into_inner);
193 ensure!(installed.is_none(), "database writer is already running");
194 *installed = Some(writer.clone());
195 }
196 let thread = match thread::Builder::new()
197 .name("hel-database-writer".to_owned())
198 .spawn(move || {
199 let mut connection = connection;
200 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
201 loop {
202 match receiver.recv() {
203 Ok(DatabaseWriterMessage::Run { label, job }) => {
204 tracing::trace!(operation = label, "running database writer operation");
205 match writer_schema_state(
208 &path,
209 &connection,
210 label,
211 &mut observed_revision,
212 ) {
213 Ok(()) => job(Ok(&mut connection)),
214 Err(error) => job(Err(error)),
215 }
216 }
217 Ok(DatabaseWriterMessage::Shutdown) => break Ok(()),
218 Err(error) => {
219 break Err(error).context("database writer queue disconnected");
220 }
221 }
222 }
223 }))
224 .unwrap_or_else(|panic| {
225 let detail = panic
226 .downcast_ref::<&str>()
227 .copied()
228 .or_else(|| panic.downcast_ref::<String>().map(String::as_str))
229 .unwrap_or("unknown panic payload");
230 Err(anyhow::anyhow!("database writer thread panicked: {detail}"))
231 });
232 clear_database_writer(id);
233 let _ = stopped_tx.send(result);
234 }) {
235 Ok(thread) => thread,
236 Err(error) => {
237 if install_globally {
238 clear_database_writer(id);
239 }
240 return Err(error).context("spawn database writer thread");
241 }
242 };
243 Ok(DatabaseWriterOwner {
244 writer,
245 thread: Some(thread),
246 stopped,
247 })
248}
249
250pub(super) fn writer_schema_state(
252 path: &Path,
253 connection: &Connection,
254 label: &'static str,
255 observed_revision: &mut i64,
256) -> Result<()> {
257 let result: Result<()> = (|| {
258 let state = schema::read_schema_state(connection)?;
259 if state.revision < *observed_revision {
260 return Err(StoreSchemaMismatch {
261 found: state.revision,
262 supported: SCHEMA_VERSION,
263 reason: StoreSchemaMismatchReason::Rollback {
264 previous: *observed_revision,
265 },
266 }
267 .into());
268 }
269 *observed_revision = state.revision;
270 state.ensure_supported()
271 })();
272 if let Err(error) = &result {
273 tracing::error!(
274 operation = label,
275 path = %path.display(),
276 error = %error,
277 "could not establish store compatibility; refusing the operation"
278 );
279 }
280 result.with_context(|| {
281 format!(
282 "check database compatibility before {label} at {}",
283 path.display()
284 )
285 })
286}
287
288pub(super) fn submit_database_write<T, F>(label: &'static str, operation: F) -> Result<T>
289where
290 T: Send + 'static,
291 F: FnOnce(&mut Connection) -> Result<T> + Send + 'static,
292{
293 let writer = database_writer_slot()
294 .lock()
295 .unwrap_or_else(PoisonError::into_inner)
296 .clone();
297 if let Some(writer) = writer {
298 writer.execute(label, operation)
299 } else {
300 bail!("database writer is not available for operation {label}")
308 }
309}