lora_database/database/
changes.rs1use 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 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 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 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 pub fn set_change_retention(&self, batches: usize) {
92 self.changes.set_retention(batches);
93 }
94
95 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 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 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 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}