lora_wal/recorder/recorder.rs
1//! [`WalRecorder`] — adapter from `MutationRecorder` to the durable
2//! [`Wal`].
3//!
4//! Lifecycle, viewed from `lora-database::Database::execute_with_params`:
5//!
6//! 1. Acquire the store write lock.
7//! 2. `recorder.arm()` — marks the recorder as inside-a-query but
8//! appends nothing to the WAL yet. A pure read query that fires
9//! no `MutationEvent` therefore touches the WAL zero times.
10//! 3. Run analyze + compile + execute. The executor mutates the
11//! in-memory store, which fires `MutationRecorder::record` for each
12//! primitive mutation. The adapter buffers those events in memory.
13//! 4. On Ok: `recorder.commit()` drains the buffered events and hands
14//! them to [`Wal::commit_tx`], which writes `TxBegin` +
15//! `MutationBatch` + `TxCommit` in one critical section and applies
16//! the configured single-thread flush policy. A read-only query returns
17//! `WroteCommit::No` and the WAL never wakes up.
18//! 5. On Err / panic: `recorder.abort()` discards the buffered events.
19//! Because `commit_tx` writes the begin/batch/commit triple
20//! atomically, an aborted query has *no* on-disk presence — there
21//! is no `TxBegin` to pair with a later `TxAbort`, so the WAL stays
22//! consistent without an explicit abort marker.
23//! 6. Before returning, the host inspects `recorder.poisoned()` once.
24//! If `Some`, the query fails loudly with a durability error so
25//! the caller can act on it; the WAL is now refusing further
26//! appends until the operator restarts the database, which
27//! recovers from the last consistent snapshot + WAL.
28//!
29//! ### Hot-path cost
30//!
31//! `record` is called once per primitive mutation. It takes only the
32//! recorder mutex and pushes the event into a query-local buffer; the
33//! WAL mutex, framing, checksum, and segment append happen once at
34//! commit time inside `Wal::commit_tx`.
35//!
36//! ### When `record` fires after a failed in-memory mutation
37//!
38//! `InMemoryGraph::emit` only calls the recorder *after* the mutation
39//! has been committed to the in-memory state. If the subsequent WAL
40//! append fails, the live in-memory store is briefly ahead of disk:
41//! the next query sees the partial state, but the next query also
42//! observes `poisoned() = Some(_)` and is rejected. Recovery from a
43//! snapshot + WAL after operator restart will not include the failed
44//! mutation, so durable state stays consistent. The cost is "the live
45//! process is wrong until the next restart"; the gain is that the
46//! storage trait does not need to learn about durability.
47
48use std::sync::{Arc, Mutex, MutexGuard};
49
50use lora_store::{MutationEvent, MutationRecorder};
51
52use super::errors::{WalBufferedCommitError, WalPoisonError, WroteCommit};
53use super::mirror::WalMirror;
54use crate::errors::WalError;
55use crate::lsn::Lsn;
56use crate::wal::Wal;
57
58#[derive(Default)]
59struct RecorderState {
60 /// True between `arm()` and the matching `commit()` / `abort()`.
61 /// Marks the host's critical section so [`MutationRecorder::record`]
62 /// knows whether to buffer fresh events or poison itself for an
63 /// out-of-scope call.
64 armed: bool,
65 /// Query-local mutation buffer. Drained by `commit()` and passed
66 /// to [`Wal::commit_tx`] as one batched WAL transaction; cleared
67 /// by `abort()` without ever touching the durable log.
68 buffer: Vec<MutationEvent>,
69 /// Sticky failure flag. Once set, [`MutationRecorder::record`]
70 /// becomes a no-op (we cannot append safely) and `poisoned`
71 /// surfaces the message.
72 poisoned: Option<String>,
73}
74
75/// Adapter that lets a [`Wal`] act as a [`MutationRecorder`] on
76/// [`lora_store::InMemoryGraph::set_mutation_recorder`].
77pub struct WalRecorder {
78 wal: Arc<Wal>,
79 mirror: Option<Arc<dyn WalMirror>>,
80 state: Mutex<RecorderState>,
81}
82
83impl WalRecorder {
84 pub fn new(wal: Arc<Wal>) -> Self {
85 Self::new_with_mirror(wal, None)
86 }
87
88 pub fn new_with_mirror(wal: Arc<Wal>, mirror: Option<Arc<dyn WalMirror>>) -> Self {
89 Self {
90 wal,
91 mirror,
92 state: Mutex::new(RecorderState::default()),
93 }
94 }
95
96 /// Underlying log handle. Exposed so admin paths
97 /// (`Database::checkpoint_to`, `truncate_up_to`) can hit the WAL
98 /// directly without going through the recorder's transaction
99 /// state machine.
100 pub fn wal(&self) -> &Arc<Wal> {
101 &self.wal
102 }
103
104 fn state_lock(&self) -> MutexGuard<'_, RecorderState> {
105 match self.state.lock() {
106 Ok(state) => state,
107 Err(poisoned) => {
108 let mut state = poisoned.into_inner();
109 state.poisoned.get_or_insert_with(|| {
110 "WalRecorder state lock was poisoned; buffered durability state is suspect"
111 .into()
112 });
113 state
114 }
115 }
116 }
117
118 /// Mark the recorder as inside a query critical section. No WAL
119 /// I/O happens here — `Wal::begin` is deferred until the first
120 /// mutation event fires. A pure read query that never produces a
121 /// `MutationEvent` therefore costs the WAL nothing: no record
122 /// allocation, no buffer drain, no `fsync`.
123 ///
124 /// Errors with [`WalError::Poisoned`] if a prior failure has
125 /// poisoned the recorder, or if the host is double-arming
126 /// (`arm` already in effect).
127 pub fn arm(&self) -> Result<(), WalError> {
128 let mut state = self.state_lock();
129 if state.poisoned.is_some() {
130 return Err(WalError::Poisoned);
131 }
132 if state.armed {
133 state.poisoned = Some("WalRecorder::arm called while already armed".into());
134 return Err(WalError::Poisoned);
135 }
136 state.armed = true;
137 state.buffer.clear();
138 Ok(())
139 }
140
141 /// Drain the buffered events and commit them as one durable WAL
142 /// transaction.
143 ///
144 /// Routes through [`Wal::commit_tx`], which encodes
145 /// `TxBegin` + `MutationBatch` + `TxCommit` in a single critical
146 /// section and applies the configured flush policy. Under `GroupSync`,
147 /// bytes are written before this method returns; storage durability is
148 /// completed by the background flusher or an explicit sync boundary.
149 ///
150 /// Returns:
151 /// - [`WroteCommit::Yes`] when mutation events fired and the WAL
152 /// wrote the begin/batch/commit triple.
153 /// - [`WroteCommit::No`] when no mutations fired during the query
154 /// and no records were written.
155 pub fn commit(&self) -> Result<WroteCommit, WalError> {
156 let events = self.take_armed_buffer()?;
157 if events.is_empty() {
158 return Ok(WroteCommit::No);
159 }
160
161 self.wal.commit_tx(events).inspect_err(|e| {
162 self.state_lock()
163 .poisoned
164 .get_or_insert_with(|| e.to_string());
165 })
166 }
167
168 /// Like [`Self::commit`], but also hands back the committed events and
169 /// the LSN of the `TxCommit` record. Returns `None` when the query
170 /// fired no mutations. Change feeds use this; it clones the event
171 /// buffer, so the plain [`Self::commit`] stays the default.
172 pub fn commit_capture(&self) -> Result<Option<(Lsn, Vec<MutationEvent>)>, WalError> {
173 let events = self.take_armed_buffer()?;
174 if events.is_empty() {
175 return Ok(None);
176 }
177 let captured = events.clone();
178 let lsn = self.wal.commit_tx_lsn(events).inspect_err(|e| {
179 self.state_lock()
180 .poisoned
181 .get_or_insert_with(|| e.to_string());
182 })?;
183 Ok(lsn.map(|lsn| (lsn, captured)))
184 }
185
186 fn take_armed_buffer(&self) -> Result<Vec<MutationEvent>, WalError> {
187 let events = {
188 let mut state = self.state_lock();
189 if state.poisoned.is_some() {
190 return Err(WalError::Poisoned);
191 }
192 if !state.armed {
193 state.poisoned = Some("WalRecorder::commit called without an armed query".into());
194 return Err(WalError::Poisoned);
195 }
196 state.armed = false;
197 std::mem::take(&mut state.buffer)
198 };
199 Ok(events)
200 }
201
202 /// Commit an explicit transaction's externally-buffered mutation
203 /// events as one durable WAL transaction.
204 ///
205 /// Used by `lora-database`'s [`Transaction`] flow, which keeps its
206 /// own `Vec<MutationEvent>` per transaction (statements may
207 /// rollback to a savepoint, which the recorder's flat buffer
208 /// cannot model). At commit time the host hands the accumulated
209 /// events here and we route them through [`Wal::commit_tx`] in one
210 /// call.
211 ///
212 /// [`Transaction`]: lora_database::Transaction
213 pub fn commit_events(
214 &self,
215 events: impl IntoIterator<Item = MutationEvent>,
216 ) -> Result<WroteCommit, WalBufferedCommitError> {
217 Ok(match self.commit_events_lsn(events)? {
218 Some(_) => WroteCommit::Yes,
219 None => WroteCommit::No,
220 })
221 }
222
223 /// Like [`Self::commit_events`], but returns the LSN of the `TxCommit`
224 /// record (`None` when there were no events and nothing was written).
225 pub fn commit_events_lsn(
226 &self,
227 events: impl IntoIterator<Item = MutationEvent>,
228 ) -> Result<Option<Lsn>, WalBufferedCommitError> {
229 self.ensure_not_poisoned()
230 .map_err(|e| WalBufferedCommitError::Poisoned(e.reason().to_string()))?;
231
232 let events: Vec<MutationEvent> = events.into_iter().collect();
233 if events.is_empty() {
234 return Ok(None);
235 }
236
237 self.wal.commit_tx_lsn(events).map_err(|e| {
238 self.state_lock()
239 .poisoned
240 .get_or_insert_with(|| e.to_string());
241 WalBufferedCommitError::Commit(super::errors::WalCommitError::Commit(e))
242 })
243 }
244
245 /// Discard any buffered events and disarm the recorder.
246 ///
247 /// Because [`Wal::commit_tx`] writes the entire begin/batch/commit
248 /// triple atomically, an aborted query never has any on-disk
249 /// presence — there is no half-written transaction to follow up
250 /// with a `TxAbort` marker. The returned bool indicates whether
251 /// the query observed any mutations (so the host can decide
252 /// whether to quarantine the live in-memory graph).
253 pub fn abort(&self) -> Result<bool, WalError> {
254 let mut state = self.state_lock();
255 if state.poisoned.is_some() {
256 return Err(WalError::Poisoned);
257 }
258 // Tolerate "abort without arm" — the host calls abort in
259 // unwind paths and we'd rather no-op than poison.
260 state.armed = false;
261 let had_buffered_events = !state.buffer.is_empty();
262 state.buffer.clear();
263 Ok(had_buffered_events)
264 }
265
266 /// Flush the WAL — write the pending buffer to the OS.
267 pub fn flush(&self) -> Result<(), WalError> {
268 let mut state = self.state_lock();
269 if state.poisoned.is_some() {
270 return Err(WalError::Poisoned);
271 }
272 self.wal.flush().inspect_err(|e| {
273 state.poisoned = Some(e.to_string());
274 })?;
275 if let Some(mirror) = &self.mirror {
276 mirror.persist(self.wal.dir()).inspect_err(|e| {
277 state.poisoned = Some(e.to_string());
278 })?;
279 }
280 Ok(())
281 }
282
283 /// Force the underlying WAL to write, `fsync`, and advance its
284 /// durable fence regardless of the configured sync mode. Admin
285 /// paths use this when they need a durability point immediately.
286 pub fn force_fsync(&self) -> Result<(), WalError> {
287 let mut state = self.state_lock();
288 if state.poisoned.is_some() {
289 return Err(WalError::Poisoned);
290 }
291 self.wal.force_fsync().inspect_err(|e| {
292 state.poisoned = Some(e.to_string());
293 })?;
294 if let Some(mirror) = &self.mirror {
295 mirror.persist_force(self.wal.dir()).inspect_err(|e| {
296 state.poisoned = Some(e.to_string());
297 })?;
298 }
299 Ok(())
300 }
301
302 /// Force only the underlying WAL to storage durability, without invoking
303 /// the optional mirror. Container-backed callers use this when they need to
304 /// build a single richer mirror refresh (for example snapshot + WAL delta)
305 /// after the WAL bytes are durable.
306 pub fn force_fsync_wal_only(&self) -> Result<(), WalError> {
307 let mut state = self.state_lock();
308 if state.poisoned.is_some() {
309 return Err(WalError::Poisoned);
310 }
311 self.wal.force_fsync().inspect_err(|e| {
312 state.poisoned = Some(e.to_string());
313 })?;
314 Ok(())
315 }
316
317 /// Append a `Checkpoint` marker. Used by the checkpoint admin
318 /// path after a successful snapshot rename — the marker doubles
319 /// as the log-side fence the next replay will trust.
320 pub fn checkpoint_marker(&self, snapshot_lsn: Lsn) -> Result<Lsn, WalError> {
321 let mut state = self.state_lock();
322 if state.poisoned.is_some() {
323 return Err(WalError::Poisoned);
324 }
325 self.wal.checkpoint_marker(snapshot_lsn).inspect_err(|e| {
326 state.poisoned = Some(e.to_string());
327 })
328 }
329
330 /// Drop sealed segments at or below `fence_lsn`. Forwards to
331 /// [`Wal::truncate_up_to`].
332 pub fn truncate_up_to(&self, fence_lsn: Lsn) -> Result<(), WalError> {
333 // Archive-backed databases must stay self-contained. Until snapshot
334 // checkpoint payloads are stored inside the archive too, preserving the
335 // full WAL history is the only safe way to let the archive recover by
336 // itself after a checkpoint marker.
337 if let Some(mirror) = &self.mirror {
338 mirror.persist_force(self.wal.dir())?;
339 return Ok(());
340 }
341 self.wal.truncate_up_to(fence_lsn)?;
342 Ok(())
343 }
344
345 /// True iff the recorder has already failed an append, **or** the WAL has
346 /// latched a durability failure. Cheap to poll under the store lock.
347 pub fn is_poisoned(&self) -> bool {
348 self.poisoned_reason().is_some()
349 }
350
351 pub fn poisoned_reason(&self) -> Option<String> {
352 let state = self.state_lock();
353 if let Some(msg) = state.poisoned.clone() {
354 return Some(msg);
355 }
356 self.wal.bg_failure()
357 }
358
359 pub fn ensure_not_poisoned(&self) -> Result<(), WalPoisonError> {
360 if let Some(reason) = self.poisoned_reason() {
361 return Err(WalPoisonError { reason });
362 }
363 Ok(())
364 }
365
366 /// Quarantine the recorder after the host detects that the live
367 /// in-memory graph may no longer match durable state. Once poisoned,
368 /// future query arms fail until the database is restarted from a
369 /// snapshot + WAL.
370 pub fn poison(&self, reason: impl Into<String>) {
371 let mut state = self.state_lock();
372 state.poisoned.get_or_insert_with(|| reason.into());
373 state.armed = false;
374 state.buffer.clear();
375 }
376
377 /// Test helper: clear the poisoned flag and disarm. Production
378 /// code should not call this — once the WAL is poisoned the right
379 /// move is to fail loudly and let the operator restart from the
380 /// last snapshot + WAL.
381 #[doc(hidden)]
382 pub fn clear_poisoned_for_tests(&self) {
383 let mut state = self.state_lock();
384 state.poisoned = None;
385 state.armed = false;
386 state.buffer.clear();
387 }
388}
389
390impl MutationRecorder for WalRecorder {
391 fn record(&self, event: MutationEvent) {
392 let mut state = self.state_lock();
393 if state.poisoned.is_some() {
394 return;
395 }
396 if !state.armed {
397 state.poisoned.get_or_insert_with(|| {
398 "MutationRecorder::record fired outside an armed query".into()
399 });
400 return;
401 }
402 state.buffer.push(event);
403 }
404
405 fn poisoned(&self) -> Option<String> {
406 // Surface a latched WAL failure too — the recorder is the host's
407 // single point of contact for "is the WAL still safe to commit
408 // through?".
409 self.poisoned_reason()
410 }
411}