1mod 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
44pub const DEFAULT_RETENTION: usize = 1024;
46
47pub const DEFAULT_FEED_BUFFER: usize = 1024;
49
50#[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 set_keys: Vec<String>,
68 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 Reset,
104}
105
106#[derive(Debug, Clone, PartialEq)]
109pub struct ChangeBatch {
110 pub lsn: u64,
112 pub changes: Vec<Change>,
113}
114
115#[derive(Debug, Clone, Copy)]
117pub struct ChangeFeedOptions {
118 pub from_lsn: Option<u64>,
122 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#[derive(Debug, Clone)]
138pub enum ChangePoll {
139 Batch(Arc<ChangeBatch>),
141 Pending,
143 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 last_lsn: u64,
222 ring: VecDeque<Arc<ChangeBatch>>,
224 ring_floor: u64,
226 retention: usize,
227 subscribers: Vec<Arc<Subscriber>>,
228 closed: bool,
229}
230
231pub(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 #[inline]
261 pub(crate) fn is_active(&self) -> bool {
262 self.active.load(Ordering::Acquire)
263 }
264
265 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 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 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 if lsn <= next.saturating_add(1) {
323 self.publish(Some(lsn), vec![Change::Reset]);
324 }
325 }
326 }
327 }
328
329 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
350pub(crate) enum CatchUp {
352 None,
353 Wal(Box<HistoryReplay>),
355}
356
357pub(crate) enum Resume {
359 Ready(ChangeFeed),
361 NeedsHistory {
364 feed: ChangeFeed,
365 from: u64,
366 floor: u64,
367 },
368}
369
370impl ChangeHub {
371 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
435pub struct ChangeFeed {
440 sub: Arc<Subscriber>,
441 backlog: VecDeque<Arc<ChangeBatch>>,
443 history: CatchUp,
444 last_lsn: Option<u64>,
445 errored: bool,
446}
447
448#[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 pub fn last_lsn(&self) -> Option<u64> {
469 self.last_lsn
470 }
471
472 pub fn closer(&self) -> ChangeFeedCloser {
474 ChangeFeedCloser {
475 sub: self.sub.clone(),
476 }
477 }
478
479 pub fn close(&self) {
481 self.sub.close();
482 }
483
484 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 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 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 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
600pub(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#[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}