Skip to main content

lora_database/changes/
mod.rs

1//! Committed-change feed.
2//!
3//! [`Database::changes`](crate::Database::changes) opens a [`ChangeFeed`]:
4//! an ordered stream of [`ChangeBatch`]es, one per committed write
5//! (auto-commit query, explicit transaction, streamed write, admin mutator,
6//! `clear()`, snapshot restore). Rolled-back work never appears.
7//!
8//! # LSNs
9//!
10//! Every batch carries an `lsn`, a strictly increasing resume token.
11//! WAL-backed databases use the LSN of the transaction's `TxCommit` record,
12//! so tokens stay valid across restarts. In-memory databases use a
13//! per-process commit counter.
14//!
15//! # Capture cost
16//!
17//! Capture is off until the first feed opens. From then on every write
18//! builds its batch (net changes per entity plus property maps) under the
19//! writer lock and keeps the last [`DEFAULT_RETENTION`] batches in memory
20//! so feeds can resume from a recent LSN without touching disk.
21//!
22//! # Back-pressure
23//!
24//! Each feed has a bounded queue. Writers never wait for readers: when a
25//! queue is full the feed ends with `LORA_CHANGES_LAGGED` after the
26//! buffered batches are drained, and the consumer resumes from the last LSN
27//! it processed.
28
29mod build;
30mod history;
31
32use std::collections::VecDeque;
33use std::sync::atomic::{AtomicBool, Ordering};
34use std::sync::{Arc, Condvar, Mutex, MutexGuard};
35use std::time::Duration;
36
37use lora_store::{InMemoryGraph, MutationEvent, NodeId, Properties, RelationshipId};
38
39pub(crate) use build::{build_changes, PreImageSink, PreImages};
40pub(crate) use history::{HistoryReplay, HistorySources};
41
42use crate::error::{LoraError, LoraErrorCode};
43
44/// Batches kept in memory for same-process resume.
45pub const DEFAULT_RETENTION: usize = 1024;
46
47/// Default per-feed queue capacity, in batches.
48pub const DEFAULT_FEED_BUFFER: usize = 1024;
49
50/// Net effect of one committed transaction on one entity.
51///
52/// Created and updated entities carry their state after the commit.
53/// Deleted entities carry their last committed state. An entity created and
54/// deleted inside the same transaction is not reported.
55#[derive(Debug, Clone, PartialEq)]
56pub enum Change {
57    NodeCreated {
58        id: NodeId,
59        labels: Vec<String>,
60        properties: Properties,
61    },
62    NodeUpdated {
63        id: NodeId,
64        labels: Vec<String>,
65        properties: Properties,
66        /// Keys the transaction set (last operation per key wins).
67        set_keys: Vec<String>,
68        /// Keys the transaction removed (last operation per key wins).
69        removed_keys: Vec<String>,
70        added_labels: Vec<String>,
71        removed_labels: Vec<String>,
72    },
73    NodeDeleted {
74        id: NodeId,
75        labels: Vec<String>,
76        properties: Properties,
77    },
78    RelationshipCreated {
79        id: RelationshipId,
80        rel_type: String,
81        start: NodeId,
82        end: NodeId,
83        properties: Properties,
84    },
85    RelationshipUpdated {
86        id: RelationshipId,
87        rel_type: String,
88        start: NodeId,
89        end: NodeId,
90        properties: Properties,
91        set_keys: Vec<String>,
92        removed_keys: Vec<String>,
93    },
94    RelationshipDeleted {
95        id: RelationshipId,
96        rel_type: String,
97        start: NodeId,
98        end: NodeId,
99        properties: Properties,
100    },
101    /// The whole graph was replaced (`clear()` or a snapshot restore).
102    /// Consumers should drop anything they derived from earlier batches.
103    Reset,
104}
105
106/// Every change one committed write made, in the order it first touched
107/// each entity.
108#[derive(Debug, Clone, PartialEq)]
109pub struct ChangeBatch {
110    /// Strictly increasing resume token.
111    pub lsn: u64,
112    pub changes: Vec<Change>,
113}
114
115/// Options for [`Database::changes`](crate::Database::changes).
116#[derive(Debug, Clone, Copy)]
117pub struct ChangeFeedOptions {
118    /// Resume after this LSN: the feed first yields every retained batch
119    /// with a greater LSN, then live batches. `None` starts with the next
120    /// commit.
121    pub from_lsn: Option<u64>,
122    /// Maximum number of undelivered live batches before the feed ends with
123    /// `LORA_CHANGES_LAGGED`.
124    pub buffer_size: usize,
125}
126
127impl Default for ChangeFeedOptions {
128    fn default() -> Self {
129        Self {
130            from_lsn: None,
131            buffer_size: DEFAULT_FEED_BUFFER,
132        }
133    }
134}
135
136/// Result of polling a [`ChangeFeed`].
137#[derive(Debug, Clone)]
138pub enum ChangePoll {
139    /// The next batch in commit order.
140    Batch(Arc<ChangeBatch>),
141    /// Nothing buffered right now.
142    Pending,
143    /// The feed was closed (by the consumer or because the database shut
144    /// down). No further batches follow.
145    Closed,
146}
147
148#[derive(Debug, Clone, Copy, PartialEq, Eq)]
149enum SubStatus {
150    Open,
151    Lagged,
152    Closed,
153}
154
155type Waker = Arc<dyn Fn() + Send + Sync>;
156
157struct SubState {
158    queue: VecDeque<Arc<ChangeBatch>>,
159    capacity: usize,
160    status: SubStatus,
161    waker: Option<Waker>,
162}
163
164pub(crate) struct Subscriber {
165    state: Mutex<SubState>,
166    ready: Condvar,
167}
168
169impl Subscriber {
170    fn lock(&self) -> MutexGuard<'_, SubState> {
171        self.state.lock().unwrap_or_else(|p| p.into_inner())
172    }
173
174    fn is_open(&self) -> bool {
175        self.lock().status == SubStatus::Open
176    }
177
178    fn push(&self, batch: &Arc<ChangeBatch>) {
179        let waker = {
180            let mut state = self.lock();
181            if state.status != SubStatus::Open {
182                return;
183            }
184            if state.queue.len() >= state.capacity {
185                state.status = SubStatus::Lagged;
186                state.waker.clone()
187            } else {
188                state.queue.push_back(batch.clone());
189                if state.queue.len() == 1 {
190                    state.waker.clone()
191                } else {
192                    None
193                }
194            }
195        };
196        self.ready.notify_all();
197        if let Some(waker) = waker {
198            waker();
199        }
200    }
201
202    fn close(&self) {
203        let waker = {
204            let mut state = self.lock();
205            if state.status == SubStatus::Closed {
206                return;
207            }
208            state.status = SubStatus::Closed;
209            state.queue.clear();
210            state.waker.take()
211        };
212        self.ready.notify_all();
213        if let Some(waker) = waker {
214            waker();
215        }
216    }
217}
218
219struct HubState {
220    /// LSN of the newest published batch (or the capture start point).
221    last_lsn: u64,
222    /// Recent batches, oldest first.
223    ring: VecDeque<Arc<ChangeBatch>>,
224    /// Every published batch with an LSN above this one is in `ring`.
225    ring_floor: u64,
226    retention: usize,
227    subscribers: Vec<Arc<Subscriber>>,
228    closed: bool,
229}
230
231/// Per-database fan-out point. Write paths publish into it while holding
232/// the writer lock, so batches arrive in commit order.
233pub(crate) struct ChangeHub {
234    active: AtomicBool,
235    state: Mutex<HubState>,
236}
237
238impl Default for ChangeHub {
239    fn default() -> Self {
240        Self {
241            active: AtomicBool::new(false),
242            state: Mutex::new(HubState {
243                last_lsn: 0,
244                ring: VecDeque::new(),
245                ring_floor: 0,
246                retention: DEFAULT_RETENTION,
247                subscribers: Vec::new(),
248                closed: false,
249            }),
250        }
251    }
252}
253
254impl ChangeHub {
255    fn lock(&self) -> MutexGuard<'_, HubState> {
256        self.state.lock().unwrap_or_else(|p| p.into_inner())
257    }
258
259    /// Whether writes should capture changes. Read under the writer lock.
260    #[inline]
261    pub(crate) fn is_active(&self) -> bool {
262        self.active.load(Ordering::Acquire)
263    }
264
265    /// Turn capture on. The caller holds the writer lock and passes the
266    /// LSN of the newest commit (WAL) or `None` for in-memory databases,
267    /// so no commit can slip between "not captured" and "captured".
268    pub(crate) fn activate(&self, head: Option<u64>) {
269        let mut state = self.lock();
270        if self.is_active() {
271            return;
272        }
273        let head = head.unwrap_or(state.last_lsn).max(state.last_lsn);
274        state.last_lsn = head;
275        state.ring_floor = head;
276        self.active.store(true, Ordering::Release);
277    }
278
279    pub(crate) fn set_retention(&self, batches: usize) {
280        let mut state = self.lock();
281        state.retention = batches;
282        trim_ring(&mut state);
283    }
284
285    pub(crate) fn head(&self) -> Option<u64> {
286        self.is_active().then(|| self.lock().last_lsn)
287    }
288
289    /// Publish one committed write. `lsn` is the WAL commit LSN; `None`
290    /// allocates the next in-memory counter value. Empty batches (reads,
291    /// catalog-only writes) are dropped.
292    pub(crate) fn publish(&self, lsn: Option<u64>, changes: Vec<Change>) {
293        if changes.is_empty() {
294            return;
295        }
296        let mut state = self.lock();
297        if state.closed {
298            return;
299        }
300        let lsn = lsn.unwrap_or(state.last_lsn + 1).max(state.last_lsn + 1);
301        state.last_lsn = lsn;
302        let batch = Arc::new(ChangeBatch { lsn, changes });
303        state.ring.push_back(batch.clone());
304        trim_ring(&mut state);
305        state.subscribers.retain(|sub| {
306            sub.push(&batch);
307            sub.is_open()
308        });
309    }
310
311    /// Publish a `Reset` for a graph replacement that has no WAL record
312    /// (snapshot restore). On a WAL-backed database the batch takes an LSN
313    /// between the last commit and the next one; when two resets arrive
314    /// without a commit between them the second is folded into the first.
315    pub(crate) fn publish_reset(&self, wal_next_lsn: Option<u64>) {
316        match wal_next_lsn {
317            None => self.publish(None, vec![Change::Reset]),
318            Some(next) => {
319                let last = self.lock().last_lsn;
320                let lsn = next.max(last + 1);
321                // The next WAL commit record takes `next + 2`.
322                if lsn <= next.saturating_add(1) {
323                    self.publish(Some(lsn), vec![Change::Reset]);
324                }
325            }
326        }
327    }
328
329    /// End every feed. Called when the database shuts down.
330    pub(crate) fn close(&self) {
331        let subscribers = {
332            let mut state = self.lock();
333            state.closed = true;
334            std::mem::take(&mut state.subscribers)
335        };
336        for sub in subscribers {
337            sub.close();
338        }
339    }
340}
341
342fn trim_ring(state: &mut HubState) {
343    while state.ring.len() > state.retention {
344        if let Some(evicted) = state.ring.pop_front() {
345            state.ring_floor = evicted.lsn;
346        }
347    }
348}
349
350/// How a new feed catches up before it switches to live batches.
351pub(crate) enum CatchUp {
352    None,
353    /// Rebuild `(from, floor]` from the WAL, then continue with `ring`.
354    Wal(Box<HistoryReplay>),
355}
356
357/// Outcome of [`ChangeHub::subscribe`] before any history is planned.
358pub(crate) enum Resume {
359    /// Everything needed is in memory.
360    Ready(ChangeFeed),
361    /// `from` is older than the in-memory window; the caller must plan a
362    /// WAL replay of `(from, floor]` and then call [`ChangeFeed::with_history`].
363    NeedsHistory {
364        feed: ChangeFeed,
365        from: u64,
366        floor: u64,
367    },
368}
369
370impl ChangeHub {
371    /// Register a feed. Runs under the hub lock only, so it never waits for
372    /// an open transaction once capture is active.
373    pub(crate) fn subscribe(
374        self: &Arc<Self>,
375        options: ChangeFeedOptions,
376    ) -> Result<Resume, LoraError> {
377        let mut state = self.lock();
378        if state.closed {
379            return Err(LoraError::new(
380                LoraErrorCode::TransactionFailure,
381                "database is closed",
382            ));
383        }
384        let head = state.last_lsn;
385        let sub = Arc::new(Subscriber {
386            state: Mutex::new(SubState {
387                queue: VecDeque::new(),
388                capacity: options.buffer_size.max(1),
389                status: SubStatus::Open,
390                waker: None,
391            }),
392            ready: Condvar::new(),
393        });
394        let mut feed = ChangeFeed {
395            sub: sub.clone(),
396            backlog: VecDeque::new(),
397            history: CatchUp::None,
398            last_lsn: options.from_lsn,
399            errored: false,
400        };
401
402        let resume = match options.from_lsn {
403            None => Resume::Ready(feed),
404            Some(from) if from > head => {
405                return Err(LoraError::new(
406                    LoraErrorCode::ChangesTruncated,
407                    format!(
408                        "change feed cannot resume from LSN {from} because the newest committed LSN is {head}; start a new feed without `fromLsn` and re-read current state"
409                    ),
410                ));
411            }
412            Some(from) => {
413                feed.backlog = state
414                    .ring
415                    .iter()
416                    .filter(|batch| batch.lsn > from)
417                    .cloned()
418                    .collect();
419                if from >= state.ring_floor {
420                    Resume::Ready(feed)
421                } else {
422                    Resume::NeedsHistory {
423                        feed,
424                        from,
425                        floor: state.ring_floor,
426                    }
427                }
428            }
429        };
430        state.subscribers.push(sub);
431        Ok(resume)
432    }
433}
434
435/// Ordered stream of committed [`ChangeBatch`]es from one database.
436///
437/// Poll it with [`Self::poll`] (non-blocking) or [`Self::next_timeout`].
438/// Dropping the feed unsubscribes it.
439pub struct ChangeFeed {
440    sub: Arc<Subscriber>,
441    /// Retained batches past `from_lsn`, delivered before live ones.
442    backlog: VecDeque<Arc<ChangeBatch>>,
443    history: CatchUp,
444    last_lsn: Option<u64>,
445    errored: bool,
446}
447
448/// Cloneable handle that closes a [`ChangeFeed`] from another thread.
449#[derive(Clone)]
450pub struct ChangeFeedCloser {
451    sub: Arc<Subscriber>,
452}
453
454impl ChangeFeedCloser {
455    pub fn close(&self) {
456        self.sub.close();
457    }
458}
459
460impl ChangeFeed {
461    pub(crate) fn with_history(mut self, history: HistoryReplay) -> Self {
462        self.history = CatchUp::Wal(Box::new(history));
463        self
464    }
465
466    /// LSN of the last batch this feed delivered (or the `from_lsn` it was
467    /// opened with). Pass it as `from_lsn` to resume.
468    pub fn last_lsn(&self) -> Option<u64> {
469        self.last_lsn
470    }
471
472    /// Handle for closing this feed from another thread.
473    pub fn closer(&self) -> ChangeFeedCloser {
474        ChangeFeedCloser {
475            sub: self.sub.clone(),
476        }
477    }
478
479    /// Close the feed. Later polls return [`ChangePoll::Closed`].
480    pub fn close(&self) {
481        self.sub.close();
482    }
483
484    /// Install a callback that fires when the feed goes from empty to
485    /// non-empty, lags, or closes. It runs on the writing thread while the
486    /// writer lock is held, so it must be quick and must not call back into
487    /// the database.
488    pub fn set_waker(&self, waker: impl Fn() + Send + Sync + 'static) {
489        let notify_now = {
490            let mut state = self.sub.lock();
491            state.waker = Some(Arc::new(waker));
492            (!state.queue.is_empty() || state.status != SubStatus::Open)
493                .then(|| state.waker.clone())
494                .flatten()
495        };
496        if let Some(waker) = notify_now {
497            waker();
498        }
499    }
500
501    fn deliver(&mut self, batch: Arc<ChangeBatch>) -> ChangePoll {
502        self.last_lsn = Some(batch.lsn);
503        ChangePoll::Batch(batch)
504    }
505
506    /// Next batch without blocking. History replay (resuming from an LSN
507    /// older than the in-memory window) does its work inside this call.
508    ///
509    /// Fails with `LORA_CHANGES_LAGGED` once the feed has fallen behind and
510    /// its buffered batches are drained; the feed is closed afterwards.
511    pub fn poll(&mut self) -> Result<ChangePoll, LoraError> {
512        if self.errored {
513            return Ok(ChangePoll::Closed);
514        }
515        if let CatchUp::Wal(history) = &mut self.history {
516            let sub = self.sub.clone();
517            let stop = move || sub.lock().status == SubStatus::Closed;
518            match history.next_batch(&stop) {
519                Ok(Some(batch)) => {
520                    // The ring backlog starts where the history ends.
521                    let lsn = batch.lsn;
522                    self.backlog.retain(|b| b.lsn > lsn);
523                    return Ok(self.deliver(Arc::new(batch)));
524                }
525                Ok(None) => self.history = CatchUp::None,
526                Err(err) => {
527                    self.errored = true;
528                    self.sub.close();
529                    return Err(err);
530                }
531            }
532        }
533        if let Some(batch) = self.backlog.pop_front() {
534            if self.last_lsn.is_none_or(|last| batch.lsn > last) {
535                return Ok(self.deliver(batch));
536            }
537            return self.poll();
538        }
539
540        let mut state = self.sub.lock();
541        if let Some(batch) = state.queue.pop_front() {
542            drop(state);
543            return Ok(self.deliver(batch));
544        }
545        match state.status {
546            SubStatus::Open => Ok(ChangePoll::Pending),
547            SubStatus::Closed => Ok(ChangePoll::Closed),
548            SubStatus::Lagged => {
549                state.status = SubStatus::Closed;
550                drop(state);
551                self.errored = true;
552                Err(self.lagged_error())
553            }
554        }
555    }
556
557    fn lagged_error(&self) -> LoraError {
558        let capacity = self.sub.lock().capacity;
559        let hint = match self.last_lsn {
560            Some(lsn) => format!("resume with `fromLsn` {lsn}"),
561            None => "start a new feed and re-read current state".to_string(),
562        };
563        LoraError::new(
564            LoraErrorCode::ChangesLagged,
565            format!("change feed fell more than {capacity} batches behind the writers; {hint}"),
566        )
567    }
568
569    /// Block until a batch arrives, the feed closes, or `timeout` passes
570    /// (then [`ChangePoll::Pending`]).
571    pub fn next_timeout(&mut self, timeout: Duration) -> Result<ChangePoll, LoraError> {
572        let deadline = std::time::Instant::now() + timeout;
573        loop {
574            match self.poll()? {
575                ChangePoll::Pending => {}
576                other => return Ok(other),
577            }
578            let now = std::time::Instant::now();
579            if now >= deadline {
580                return Ok(ChangePoll::Pending);
581            }
582            let state = self.sub.lock();
583            if state.queue.is_empty() && state.status == SubStatus::Open {
584                let _ = self
585                    .sub
586                    .ready
587                    .wait_timeout(state, deadline - now)
588                    .unwrap_or_else(|p| p.into_inner());
589            }
590        }
591    }
592}
593
594impl Drop for ChangeFeed {
595    fn drop(&mut self) {
596        self.sub.close();
597    }
598}
599
600/// Build and publish the batch for one committed write against an
601/// `InMemoryGraph`. `pre` holds the records the write deleted.
602pub(crate) fn publish_committed(
603    hub: &ChangeHub,
604    lsn: Option<u64>,
605    events: &[MutationEvent],
606    pre: &PreImages,
607    post: &InMemoryGraph,
608) {
609    if events.is_empty() {
610        return;
611    }
612    hub.publish(lsn, build_changes(events, pre, post));
613}
614
615/// Recorder that buffers events for change capture on databases without a
616/// WAL (the WAL recorder already buffers them otherwise).
617#[derive(Default)]
618pub(crate) struct CaptureRecorder {
619    events: Mutex<Vec<MutationEvent>>,
620}
621
622impl CaptureRecorder {
623    pub(crate) fn take(&self) -> Vec<MutationEvent> {
624        std::mem::take(&mut *self.events.lock().unwrap_or_else(|p| p.into_inner()))
625    }
626}
627
628impl lora_store::MutationRecorder for CaptureRecorder {
629    fn record(&self, event: MutationEvent) {
630        self.events
631            .lock()
632            .unwrap_or_else(|p| p.into_inner())
633            .push(event);
634    }
635}