Skip to main content

mj_controller/database/
writer.rs

1use super::*;
2
3pub(super) const DATABASE_WRITE_QUEUE_CAPACITY: usize = 256;
4
5/// A queued write, handed either the writer's connection or the reason it
6/// must not be used. The job -- not the lane -- decides what a refusal means
7/// to its caller.
8pub(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/// Cloneable submission handle for the daemon's ordered SQLite write lane.
19///
20/// Calling [`DatabaseWriter::execute`] is synchronous and may apply bounded
21/// backpressure, so async and UI callers must invoke database mutations from
22/// their existing supervised blocking tasks.
23#[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                        // The mismatch travels as the operation's own failure,
52                        // so a refused write reports why rather than the
53                        // writer-stopped message a dropped reply would give.
54                        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
68/// Owns the daemon's writer thread and persistent SQLite connection.
69///
70/// The owner is deliberately not cloneable. Dropping it removes the global
71/// submission handle, drains accepted work in FIFO order, and joins the
72/// thread before releasing the connection.
73pub 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
132/// Whether this process holds the database writer. Only the daemon does, so
133/// work that merely keeps a daemon-owned record current, and that other
134/// processes also run, asks this instead of failing the write there.
135pub(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/// Install the process-wide writer for a test that owns its data directory.
152///
153/// Production installs this once, in the daemon, after `ControllerStoreGuard`
154/// establishes exclusivity, and the daemon is then the only process that
155/// writes. A test may do the same only because it re-execs itself with its own
156/// `MJ_DATA_DIR` and is therefore alone in its process — which is exactly why
157/// the tests that need this are shaped that way.
158///
159/// The returned owner has to be held for the rest of the test: dropping it
160/// stops the writer, and the next write fails with the message above.
161///
162/// This fixture is compiled unconditionally and hidden from the documentation
163/// because the controller crate's tests need it and a `#[cfg(test)]` item is
164/// invisible to another crate. It is a thin wrapper over
165/// [`start_database_writer`], so nothing test-only leaks into the library.
166#[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                            // Recheck compatibility even after startup, and
206                            // remember forward progress to detect rollback.
207                            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
250/// Refuse incompatible stores, rollback, and unreadable metadata before a job.
251pub(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        // There is one way to write, and this is not it. In production the
301        // daemon installs the writer after `ControllerStoreGuard` establishes
302        // exclusivity, and it is the only process that writes; a caller
303        // reaching here has no exclusivity and would be competing with
304        // whatever does. This used to open `database_path()` directly, which
305        // meant any process without a writer silently wrote to — and migrated
306        // — the real user database as a side effect of doing something else.
307        bail!("database writer is not available for operation {label}")
308    }
309}