Skip to main content

laminar_db/subscription/registry/
mod.rs

1//! Per-object shared log for `SUBSCRIBE`.
2//!
3//! Publication and attachment use the same mutex. Portals retain only a
4//! sequence cursor and a wake-up receiver. Arrow allocations are charged once
5//! and shared by the byte-bounded log and any frame currently in transit.
6
7use std::collections::{HashMap, VecDeque};
8use std::sync::atomic::{AtomicUsize, Ordering};
9use std::sync::Arc;
10
11use arrow_array::RecordBatch;
12use laminar_core::checkpoint::CheckpointAttempt;
13use parking_lot::Mutex;
14use tokio::sync::watch;
15
16mod lifecycle;
17mod reader;
18
19pub(crate) use lifecycle::SubscriptionRegistry;
20pub(super) use reader::SubscriptionRead;
21pub(crate) use reader::SubscriptionReader;
22#[cfg(test)]
23use reader::TryRead;
24
25pub(super) const MAX_LIVE_BATCH_BYTES: usize = 8 * 1024 * 1024;
26
27// Retention is a replay policy, not permission for unbounded live buffering.
28// Two maximum-sized batches allow a portal to make progress while keeping the
29// process-wide exposure finite when user retention is disabled.
30const INTERNAL_LIVE_LOG_BYTES: usize = MAX_LIVE_BATCH_BYTES * 2;
31const ENTRY_ACCOUNTING_BYTES: usize = 64;
32const BARRIER_ENTRY_BYTES: usize = ENTRY_ACCOUNTING_BYTES + std::mem::size_of::<u64>() * 3;
33
34struct SubscriptionMemoryBudget {
35    limit: usize,
36    used: AtomicUsize,
37}
38
39struct ChargedUpdateInner {
40    update: Option<MvUpdate>,
41    bytes: usize,
42    budget: Arc<SubscriptionMemoryBudget>,
43}
44
45impl Drop for ChargedUpdateInner {
46    fn drop(&mut self) {
47        {
48            let _update = self.update.take();
49        }
50        self.budget.release(self.bytes);
51    }
52}
53
54#[derive(Clone)]
55pub(super) struct ChargedUpdate(Arc<ChargedUpdateInner>);
56
57impl ChargedUpdate {
58    /// Takes ownership of capacity that was reserved before this value was built.
59    fn from_reserved(
60        update: MvUpdate,
61        bytes: usize,
62        budget: Arc<SubscriptionMemoryBudget>,
63    ) -> Self {
64        Self(Arc::new(ChargedUpdateInner {
65            update: Some(update),
66            bytes,
67            budget,
68        }))
69    }
70
71    pub(super) fn as_ref(&self) -> &MvUpdate {
72        self.0
73            .update
74            .as_ref()
75            .expect("charged subscription update is present until final drop")
76    }
77}
78
79impl SubscriptionMemoryBudget {
80    fn new(limit: usize) -> Self {
81        Self {
82            limit,
83            used: AtomicUsize::new(0),
84        }
85    }
86
87    fn try_reserve(&self, bytes: usize) -> bool {
88        self.used
89            .fetch_update(Ordering::AcqRel, Ordering::Acquire, |used| {
90                (bytes <= self.limit.saturating_sub(used)).then(|| used.saturating_add(bytes))
91            })
92            .is_ok()
93    }
94
95    fn release(&self, bytes: usize) {
96        if bytes == 0 {
97            return;
98        }
99        let released = self
100            .used
101            .fetch_update(Ordering::AcqRel, Ordering::Acquire, |used| {
102                used.checked_sub(bytes)
103            });
104        debug_assert!(released.is_ok(), "subscription memory released twice");
105    }
106
107    #[cfg(test)]
108    fn used(&self) -> usize {
109        self.used.load(Ordering::Acquire)
110    }
111}
112
113#[derive(Debug)]
114pub(crate) enum MvUpdate {
115    Batch(RecordBatch),
116    Barrier {
117        epoch: u64,
118        checkpoint_id: u64,
119        /// First shared-log sequence not certified by this checkpoint.
120        through_sequence: u64,
121    },
122    Error(String),
123}
124
125/// Where a new subscriber should start reading.
126#[derive(Clone, Copy, Debug)]
127pub enum SubscribeStart {
128    /// See only entries sequenced after attachment.
129    Tail,
130    /// Replay entries strictly after the retained barrier with `epoch == n`.
131    AsOfEpoch(u64),
132}
133
134#[derive(Debug)]
135pub(crate) enum SubscriptionOpenError {
136    ReplayPruned {
137        /// Earliest barrier epoch eligible for replay; `0` if none is retained.
138        earliest_retained: u64,
139    },
140    EpochNotCommitted {
141        requested: u64,
142        latest_committed: Option<u64>,
143    },
144    Capacity {
145        attached: usize,
146        limit: usize,
147    },
148}
149
150struct SequencedUpdate {
151    sequence: u64,
152    bytes: usize,
153    update: ChargedUpdate,
154}
155
156struct StreamLog {
157    inner: Mutex<StreamLogInner>,
158    wake: watch::Sender<()>,
159    budget: Arc<SubscriptionMemoryBudget>,
160}
161
162struct StreamLogInner {
163    entries: VecDeque<SequencedUpdate>,
164    bytes: usize,
165    retention_cap: usize,
166    retention_floor: u64,
167    retention_bytes: usize,
168    next_sequence: u64,
169    next_reader_id: u64,
170    readers: HashMap<u64, u64>,
171    latest_committed_epoch: Option<u64>,
172    terminal_error: Option<String>,
173    reserved_marker: Option<CheckpointAttempt>,
174}
175
176#[must_use]
177enum AppendOutcome {
178    Stored,
179    Unobserved,
180    Rejected(String),
181}
182
183#[must_use]
184enum ReservedAppendOutcome {
185    Stored,
186    Unobserved,
187    Rejected(String),
188}
189
190impl StreamLog {
191    fn new(
192        retention_cap: usize,
193        budget: Arc<SubscriptionMemoryBudget>,
194        latest_committed_epoch: Option<u64>,
195    ) -> Self {
196        Self::new_at(retention_cap, budget, latest_committed_epoch, 0)
197    }
198
199    fn new_at(
200        retention_cap: usize,
201        budget: Arc<SubscriptionMemoryBudget>,
202        latest_committed_epoch: Option<u64>,
203        next_sequence: u64,
204    ) -> Self {
205        let (wake, _) = watch::channel(());
206        Self {
207            inner: Mutex::new(StreamLogInner {
208                entries: VecDeque::new(),
209                bytes: 0,
210                retention_cap,
211                retention_floor: next_sequence,
212                retention_bytes: 0,
213                next_sequence,
214                next_reader_id: 0,
215                readers: HashMap::new(),
216                latest_committed_epoch,
217                terminal_error: None,
218                reserved_marker: None,
219            }),
220            wake,
221            budget,
222        }
223    }
224
225    fn append(&self, update: MvUpdate) -> AppendOutcome {
226        let entry_bytes = approx_size(&update);
227        let mut inner = self.inner.lock();
228        debug_assert!(!matches!(&update, MvUpdate::Barrier { .. }));
229        if let Some(message) = &inner.terminal_error {
230            return AppendOutcome::Rejected(message.clone());
231        }
232        if inner.retention_cap == 0 && inner.readers.is_empty() {
233            return AppendOutcome::Unobserved;
234        }
235
236        let storage_cap = storage_cap(inner.retention_cap);
237        if entry_bytes > storage_cap {
238            let message = format!(
239                "subscription entry requires {entry_bytes} bytes, exceeding the internal shared-log limit of {storage_cap} bytes"
240            );
241            return self.reject_append(inner, message);
242        }
243
244        let required_sequence_space = if inner.reserved_marker.is_some() {
245            2
246        } else {
247            1
248        };
249        if inner
250            .next_sequence
251            .checked_add(required_sequence_space)
252            .is_none()
253        {
254            return self.reject_append(inner, "subscription sequence space exhausted".into());
255        }
256        let Some(next_sequence) = inner.next_sequence.checked_add(1) else {
257            unreachable!("sequence-space preflight admitted the next sequence");
258        };
259
260        reclaim_consumed_prefix(&mut inner);
261        let available_locally = storage_cap.saturating_sub(entry_bytes);
262        evict_to_byte_target(&mut inner, available_locally);
263
264        if !self.budget.try_reserve(entry_bytes) {
265            let message = format!(
266                "subscription process memory budget exhausted ({} byte limit); update was not delivered",
267                self.budget.limit
268            );
269            return self.reject_append(inner, message);
270        }
271
272        let sequence = inner.next_sequence;
273        inner.entries.push_back(SequencedUpdate {
274            sequence,
275            bytes: entry_bytes,
276            update: ChargedUpdate::from_reserved(update, entry_bytes, Arc::clone(&self.budget)),
277        });
278        inner.next_sequence = next_sequence;
279        inner.bytes = inner.bytes.saturating_add(entry_bytes);
280        retain_appended_entry(&mut inner, sequence, entry_bytes);
281        reclaim_consumed_prefix(&mut inner);
282        drop(inner);
283        self.wake.send_modify(|()| {});
284        AppendOutcome::Stored
285    }
286
287    fn reject_append(
288        &self,
289        mut inner: parking_lot::MutexGuard<'_, StreamLogInner>,
290        message: String,
291    ) -> AppendOutcome {
292        if inner.reserved_marker.is_some() {
293            return AppendOutcome::Rejected(message);
294        }
295        Self::terminate_locked(&mut inner, message.clone());
296        drop(inner);
297        self.wake.send_modify(|()| {});
298        AppendOutcome::Rejected(message)
299    }
300
301    fn reserve_marker(&self, attempt: CheckpointAttempt) -> Result<Option<u64>, String> {
302        let mut inner = self.inner.lock();
303        if inner.terminal_error.is_some() || (inner.retention_cap == 0 && inner.readers.is_empty())
304        {
305            return Ok(None);
306        }
307        if let Some(existing) = inner.reserved_marker {
308            return Err(format!(
309                "subscription marker already reserved for epoch={} checkpoint_id={}",
310                existing.epoch, existing.checkpoint_id
311            ));
312        }
313        if inner.next_sequence.checked_add(1).is_none() {
314            return Err("subscription sequence space cannot admit a checkpoint marker".into());
315        }
316        let through_sequence = inner.next_sequence;
317        inner.reserved_marker = Some(attempt);
318        Ok(Some(through_sequence))
319    }
320
321    fn cancel_marker(&self, attempt: CheckpointAttempt) -> bool {
322        let mut inner = self.inner.lock();
323        if inner.reserved_marker != Some(attempt) {
324            return false;
325        }
326        inner.reserved_marker = None;
327        true
328    }
329
330    fn append_reserved_marker(
331        &self,
332        attempt: CheckpointAttempt,
333        through_sequence: u64,
334    ) -> ReservedAppendOutcome {
335        let mut inner = self.inner.lock();
336        if inner.reserved_marker != Some(attempt) {
337            return ReservedAppendOutcome::Rejected(format!(
338                "subscription marker reservation mismatch for epoch={} checkpoint_id={}",
339                attempt.epoch, attempt.checkpoint_id
340            ));
341        }
342        inner.reserved_marker = None;
343
344        if let Some(message) = &inner.terminal_error {
345            return ReservedAppendOutcome::Rejected(format!(
346                "subscription log terminated before checkpoint marker publication: {message}"
347            ));
348        }
349        inner.latest_committed_epoch = Some(
350            inner
351                .latest_committed_epoch
352                .map_or(attempt.epoch, |latest| latest.max(attempt.epoch)),
353        );
354        if inner.retention_cap == 0 && inner.readers.is_empty() {
355            return ReservedAppendOutcome::Unobserved;
356        }
357        if through_sequence > inner.next_sequence {
358            return ReservedAppendOutcome::Rejected(format!(
359                "subscription checkpoint cut cursor {through_sequence} is ahead of log cursor {}",
360                inner.next_sequence
361            ));
362        }
363        let Some(next_sequence) = inner.next_sequence.checked_add(1) else {
364            return ReservedAppendOutcome::Rejected(
365                "subscription sequence space exhausted before checkpoint marker publication".into(),
366            );
367        };
368
369        let storage_cap = storage_cap(inner.retention_cap);
370        if BARRIER_ENTRY_BYTES > storage_cap {
371            return ReservedAppendOutcome::Rejected(format!(
372                "subscription checkpoint marker requires {BARRIER_ENTRY_BYTES} bytes, exceeding the internal shared-log limit of {storage_cap} bytes"
373            ));
374        }
375        reclaim_consumed_prefix(&mut inner);
376        evict_to_byte_target(&mut inner, storage_cap.saturating_sub(BARRIER_ENTRY_BYTES));
377
378        let sequence = inner.next_sequence;
379        inner.entries.push_back(SequencedUpdate {
380            sequence,
381            bytes: BARRIER_ENTRY_BYTES,
382            update: ChargedUpdate::from_reserved(
383                MvUpdate::Barrier {
384                    epoch: attempt.epoch,
385                    checkpoint_id: attempt.checkpoint_id,
386                    through_sequence,
387                },
388                BARRIER_ENTRY_BYTES,
389                Arc::clone(&self.budget),
390            ),
391        });
392        inner.next_sequence = next_sequence;
393        inner.bytes = inner.bytes.saturating_add(BARRIER_ENTRY_BYTES);
394        retain_appended_entry(&mut inner, sequence, BARRIER_ENTRY_BYTES);
395        reclaim_consumed_prefix(&mut inner);
396        drop(inner);
397        self.wake.send_modify(|()| {});
398        ReservedAppendOutcome::Stored
399    }
400
401    fn subscribe(
402        self: &Arc<Self>,
403        start: SubscribeStart,
404    ) -> Result<SubscriptionReader, SubscriptionOpenError> {
405        let mut inner = self.inner.lock();
406        let attached = inner.readers.len();
407        if attached >= super::MAX_SUBSCRIBERS_PER_MV {
408            return Err(SubscriptionOpenError::Capacity {
409                attached,
410                limit: super::MAX_SUBSCRIBERS_PER_MV,
411            });
412        }
413
414        let (cursor, skip_barrier) = match (inner.terminal_error.is_some(), start) {
415            (true, _) | (false, SubscribeStart::Tail) => (inner.next_sequence, None),
416            (false, SubscribeStart::AsOfEpoch(epoch)) => {
417                let (cursor, barrier_sequence) = cursor_after_retained_epoch(&inner, epoch)?;
418                (cursor, Some((epoch, barrier_sequence)))
419            }
420        };
421        let wake = self.wake.subscribe();
422        let mut reader_id = inner.next_reader_id;
423        while inner.readers.contains_key(&reader_id) {
424            reader_id = reader_id.wrapping_add(1);
425        }
426        inner.next_reader_id = reader_id.wrapping_add(1);
427        inner.readers.insert(reader_id, cursor);
428        drop(inner);
429
430        Ok(SubscriptionReader::attached(
431            Arc::clone(self),
432            reader_id,
433            cursor,
434            skip_barrier,
435            wake,
436        ))
437    }
438
439    fn set_retention_cap(&self, retention_cap: usize) {
440        let mut inner = self.inner.lock();
441        let previous_head = head_sequence(&inner);
442        inner.retention_cap = retention_cap;
443        recompute_retention_suffix(&mut inner);
444        reclaim_consumed_prefix(&mut inner);
445        evict_to_byte_target(&mut inner, storage_cap(retention_cap));
446        let head_changed = head_sequence(&inner) != previous_head;
447        drop(inner);
448        if head_changed {
449            self.wake.send_modify(|()| {});
450        }
451    }
452
453    fn terminate(&self, message: &str) {
454        let mut inner = self.inner.lock();
455        Self::terminate_locked(&mut inner, message.to_owned());
456        drop(inner);
457        self.wake.send_modify(|()| {});
458    }
459
460    fn terminate_and_replace(&self, message: &str, latest_committed_epoch: Option<u64>) -> Self {
461        let mut inner = self.inner.lock();
462        debug_assert!(
463            inner.reserved_marker.is_none(),
464            "subscription generation invalidated before its checkpoint marker was released"
465        );
466        let retention_cap = inner.retention_cap;
467        let next_sequence = inner.next_sequence;
468        Self::terminate_locked(&mut inner, message.to_owned());
469        drop(inner);
470        self.wake.send_modify(|()| {});
471        Self::new_at(
472            retention_cap,
473            Arc::clone(&self.budget),
474            latest_committed_epoch,
475            next_sequence,
476        )
477    }
478
479    fn subscriber_count(&self) -> usize {
480        self.inner.lock().readers.len()
481    }
482
483    fn observe_committed_epoch(&self, epoch: u64) {
484        let mut inner = self.inner.lock();
485        inner.latest_committed_epoch = Some(
486            inner
487                .latest_committed_epoch
488                .map_or(epoch, |latest| latest.max(epoch)),
489        );
490    }
491
492    fn terminate_locked(inner: &mut StreamLogInner, message: String) {
493        if inner.terminal_error.is_none() {
494            inner.terminal_error = Some(message);
495        }
496        clear_entries(inner);
497    }
498}
499
500impl Drop for StreamLog {
501    fn drop(&mut self) {
502        let inner = self.inner.get_mut();
503        debug_assert!(
504            inner.reserved_marker.is_none(),
505            "subscription log dropped with a live checkpoint marker reservation"
506        );
507    }
508}
509
510fn storage_cap(retention_cap: usize) -> usize {
511    retention_cap.max(INTERNAL_LIVE_LOG_BYTES)
512}
513
514fn head_sequence(inner: &StreamLogInner) -> u64 {
515    inner
516        .entries
517        .front()
518        .map_or(inner.next_sequence, |entry| entry.sequence)
519}
520
521fn clear_entries(inner: &mut StreamLogInner) {
522    inner.entries.clear();
523    inner.bytes = 0;
524    inner.retention_floor = inner.next_sequence;
525    inner.retention_bytes = 0;
526}
527
528fn pop_front(inner: &mut StreamLogInner) -> Option<SequencedUpdate> {
529    let entry = inner.entries.pop_front()?;
530    inner.bytes = inner.bytes.saturating_sub(entry.bytes);
531    if entry.sequence >= inner.retention_floor {
532        debug_assert_eq!(
533            entry.sequence, inner.retention_floor,
534            "retention floor skipped a retained entry"
535        );
536        inner.retention_bytes = inner.retention_bytes.saturating_sub(entry.bytes);
537        inner.retention_floor = entry.sequence.saturating_add(1);
538        if inner.retention_bytes == 0 {
539            inner.retention_floor = inner.next_sequence;
540        }
541    }
542    Some(entry)
543}
544
545fn evict_to_byte_target(inner: &mut StreamLogInner, target: usize) {
546    while inner.bytes > target {
547        let Some(_evicted) = pop_front(inner) else {
548            inner.bytes = 0;
549            inner.retention_floor = inner.next_sequence;
550            inner.retention_bytes = 0;
551            break;
552        };
553    }
554}
555
556fn reclaim_consumed_prefix(inner: &mut StreamLogInner) {
557    let reader_floor = inner
558        .readers
559        .values()
560        .copied()
561        .min()
562        .unwrap_or(inner.next_sequence);
563    let protected_floor = inner.retention_floor.min(reader_floor);
564
565    while inner
566        .entries
567        .front()
568        .is_some_and(|entry| entry.sequence < protected_floor)
569    {
570        let Some(_evicted) = pop_front(inner) else {
571            break;
572        };
573    }
574}
575
576fn calculate_retention_suffix(inner: &StreamLogInner) -> (u64, usize) {
577    if inner.retention_cap == 0 {
578        return (inner.next_sequence, 0);
579    }
580
581    let mut bytes = 0usize;
582    let mut floor = inner.next_sequence;
583    for entry in inner.entries.iter().rev() {
584        let with_entry = bytes.saturating_add(entry.bytes);
585        if with_entry > inner.retention_cap {
586            break;
587        }
588        bytes = with_entry;
589        floor = entry.sequence;
590    }
591    (floor, bytes)
592}
593
594fn recompute_retention_suffix(inner: &mut StreamLogInner) {
595    let (floor, bytes) = calculate_retention_suffix(inner);
596    inner.retention_floor = floor;
597    inner.retention_bytes = bytes;
598}
599
600fn retained_entry_bytes(inner: &StreamLogInner, sequence: u64) -> Option<usize> {
601    let head = head_sequence(inner);
602    if sequence < head || sequence >= inner.next_sequence {
603        return None;
604    }
605    let index = usize::try_from(sequence.saturating_sub(head)).ok()?;
606    inner.entries.get(index).map(|entry| entry.bytes)
607}
608
609fn retain_appended_entry(inner: &mut StreamLogInner, sequence: u64, bytes: usize) {
610    if inner.retention_cap == 0 {
611        inner.retention_floor = inner.next_sequence;
612        inner.retention_bytes = 0;
613        return;
614    }
615
616    if inner.retention_floor == sequence {
617        debug_assert_eq!(inner.retention_bytes, 0);
618        inner.retention_floor = sequence;
619    }
620    inner.retention_bytes = inner.retention_bytes.saturating_add(bytes);
621    while inner.retention_bytes > inner.retention_cap {
622        let Some(expired_bytes) = retained_entry_bytes(inner, inner.retention_floor) else {
623            debug_assert!(false, "retention floor is outside the shared log");
624            inner.retention_floor = inner.next_sequence;
625            inner.retention_bytes = 0;
626            return;
627        };
628        inner.retention_bytes = inner.retention_bytes.saturating_sub(expired_bytes);
629        inner.retention_floor = inner.retention_floor.saturating_add(1);
630    }
631    if inner.retention_bytes == 0 {
632        inner.retention_floor = inner.next_sequence;
633    }
634}
635
636fn cursor_after_retained_epoch(
637    inner: &StreamLogInner,
638    requested_epoch: u64,
639) -> Result<(u64, u64), SubscriptionOpenError> {
640    let Some(latest_committed) = inner.latest_committed_epoch else {
641        return Err(SubscriptionOpenError::EpochNotCommitted {
642            requested: requested_epoch,
643            latest_committed: None,
644        });
645    };
646    if requested_epoch > latest_committed {
647        return Err(SubscriptionOpenError::EpochNotCommitted {
648            requested: requested_epoch,
649            latest_committed: Some(latest_committed),
650        });
651    }
652    if inner.retention_cap == 0 {
653        return Err(SubscriptionOpenError::ReplayPruned {
654            earliest_retained: 0,
655        });
656    }
657
658    let mut cursor = None;
659    let mut earliest_retained = u64::MAX;
660    for entry in inner
661        .entries
662        .iter()
663        .filter(|entry| entry.sequence >= inner.retention_floor)
664    {
665        if let MvUpdate::Barrier {
666            epoch,
667            through_sequence,
668            ..
669        } = entry.update.as_ref()
670        {
671            if *through_sequence < inner.retention_floor {
672                continue;
673            }
674            earliest_retained = earliest_retained.min(*epoch);
675            if *epoch == requested_epoch {
676                cursor = Some((*through_sequence, entry.sequence));
677            }
678        }
679    }
680
681    if let Some(cursor) = cursor {
682        return Ok(cursor);
683    }
684    if earliest_retained == u64::MAX || requested_epoch < earliest_retained {
685        return Err(SubscriptionOpenError::ReplayPruned {
686            earliest_retained: if earliest_retained == u64::MAX {
687                0
688            } else {
689                earliest_retained
690            },
691        });
692    }
693    Err(SubscriptionOpenError::EpochNotCommitted {
694        requested: requested_epoch,
695        latest_committed: Some(latest_committed),
696    })
697}
698
699pub(super) fn approx_size(update: &MvUpdate) -> usize {
700    let payload = match update {
701        MvUpdate::Batch(batch) => batch.get_array_memory_size(),
702        MvUpdate::Barrier { .. } => return BARRIER_ENTRY_BYTES,
703        MvUpdate::Error(message) => message.len(),
704    };
705    payload.saturating_add(ENTRY_ACCOUNTING_BYTES)
706}
707
708#[cfg(test)]
709mod tests;