Skip to main content

lora_database/database/
changes.rs

1//! Change-feed entry points on [`Database<InMemoryGraph>`].
2
3use std::any::Any;
4
5use lora_store::{GraphStorage, GraphStorageMut, InMemoryGraph, MutationEvent};
6use lora_wal::Lsn;
7
8use crate::changes::{
9    publish_committed, ChangeFeed, ChangeFeedOptions, HistoryReplay, HistorySources, PreImages,
10    Resume,
11};
12use crate::database::Database;
13use crate::error::LoraError;
14
15impl Database<InMemoryGraph> {
16    /// Open a feed of committed changes.
17    ///
18    /// Without `from_lsn` the feed starts with the next commit. With
19    /// `from_lsn` it first replays every batch after that LSN: from the
20    /// in-memory window when it is recent enough, otherwise (WAL-backed
21    /// databases only) by rebuilding history from the WAL. An LSN the
22    /// database no longer retains fails with `LORA_CHANGES_TRUNCATED`.
23    ///
24    /// The first call turns change capture on for the lifetime of the
25    /// database; it waits for the writer lock once so no commit is missed.
26    pub fn changes(&self, options: ChangeFeedOptions) -> Result<ChangeFeed, LoraError> {
27        if !self.changes.is_active() {
28            let _writer = self
29                .writer
30                .lock()
31                .unwrap_or_else(|poisoned| poisoned.into_inner());
32            self.activate_changes();
33        }
34        self.subscribe_changes(options)
35    }
36
37    /// Non-blocking [`Self::changes`]: `None` when this is the first feed on
38    /// the database and turning capture on would have to wait for a write
39    /// in progress. Bindings call it on their event-loop thread and fall
40    /// back to [`Self::changes`] on a worker.
41    pub fn try_changes(&self, options: ChangeFeedOptions) -> Option<Result<ChangeFeed, LoraError>> {
42        if !self.changes.is_active() {
43            let _writer = match self.writer.try_lock() {
44                Ok(guard) => guard,
45                Err(std::sync::TryLockError::Poisoned(poisoned)) => poisoned.into_inner(),
46                Err(std::sync::TryLockError::WouldBlock) => return None,
47            };
48            self.activate_changes();
49        }
50        Some(self.subscribe_changes(options))
51    }
52
53    /// Caller holds the writer lock.
54    fn activate_changes(&self) {
55        let head = self
56            .wal
57            .as_ref()
58            .map(|rec| rec.wal().next_lsn().raw().saturating_sub(1));
59        self.changes.activate(head);
60    }
61
62    fn subscribe_changes(&self, options: ChangeFeedOptions) -> Result<ChangeFeed, LoraError> {
63        match self.changes.subscribe(options)? {
64            Resume::Ready(feed) => Ok(feed),
65            Resume::NeedsHistory { feed, from, floor } => {
66                let Some(rec) = &self.wal else {
67                    return Err(LoraError::new(
68                        crate::error::LoraErrorCode::ChangesTruncated,
69                        format!(
70                            "change feed cannot resume from LSN {from} because this in-memory database only retains batches after LSN {floor}; start a new feed without `fromLsn` and re-read current state"
71                        ),
72                    ));
73                };
74                let container = match &self.named_archive {
75                    Some(archive) => archive.snapshot_bytes()?,
76                    None => None,
77                };
78                let sources = HistorySources {
79                    wal_dir: rec.wal().dir().to_path_buf(),
80                    snapshots: self.snapshots.clone(),
81                    container,
82                };
83                let history = HistoryReplay::plan(sources, from, floor)?;
84                Ok(feed.with_history(history))
85            }
86        }
87    }
88
89    /// Number of recent batches kept in memory for same-process resume
90    /// (default [`crate::DEFAULT_RETENTION`]).
91    pub fn set_change_retention(&self, batches: usize) {
92        self.changes.set_retention(batches);
93    }
94
95    /// LSN of the newest captured batch, or `None` while no feed has been
96    /// opened on this database.
97    pub fn changes_head(&self) -> Option<u64> {
98        self.changes.head()
99    }
100}
101
102impl<S> Database<S>
103where
104    S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
105{
106    /// Publish one committed write to the change feed. `pre` is the graph
107    /// before the write (read only when it deleted something), `post` the
108    /// graph after it. No-op for non-`InMemoryGraph` backends.
109    pub(crate) fn publish_changes(
110        &self,
111        lsn: Option<Lsn>,
112        events: &[MutationEvent],
113        pre: Option<&S>,
114        post: &S,
115    ) {
116        let pre = pre.and_then(|pre| (pre as &dyn Any).downcast_ref::<InMemoryGraph>());
117        self.publish_changes_with(lsn, events, &PreImages::for_events(events, pre), post);
118    }
119
120    /// [`Self::publish_changes`] with deleted records already collected.
121    pub(crate) fn publish_changes_with(
122        &self,
123        lsn: Option<Lsn>,
124        events: &[MutationEvent],
125        pre: &PreImages,
126        post: &S,
127    ) {
128        let Some(post) = (post as &dyn Any).downcast_ref::<InMemoryGraph>() else {
129            return;
130        };
131        publish_committed(&self.changes, lsn.map(Lsn::raw), events, pre, post);
132    }
133
134    /// Announce that the whole graph was replaced outside the WAL
135    /// (snapshot restore, admin swap).
136    pub(crate) fn publish_reset(&self) {
137        if !self.changes.is_active() {
138            return;
139        }
140        let next = self.wal.as_ref().map(|rec| rec.wal().next_lsn().raw());
141        self.changes.publish_reset(next);
142    }
143}