Skip to main content

laminar_db/subscription/
registry.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, OnceLock};
10
11use arrow_array::RecordBatch;
12use laminar_core::state::CheckpointAttempt;
13use parking_lot::{Mutex, RwLock};
14use tokio::sync::watch;
15
16pub(super) const MAX_LIVE_BATCH_BYTES: usize = 8 * 1024 * 1024;
17
18// Retention is a replay policy, not permission for unbounded live buffering.
19// Two maximum-sized batches allow a portal to make progress while keeping the
20// process-wide exposure finite when user retention is disabled.
21const INTERNAL_LIVE_LOG_BYTES: usize = MAX_LIVE_BATCH_BYTES * 2;
22const PROCESS_SUBSCRIPTION_BYTES: usize = 256 * 1024 * 1024;
23const ENTRY_ACCOUNTING_BYTES: usize = 64;
24const BARRIER_ENTRY_BYTES: usize = ENTRY_ACCOUNTING_BYTES + std::mem::size_of::<u64>() * 3;
25
26static PROCESS_SUBSCRIPTION_BUDGET: OnceLock<Arc<SubscriptionMemoryBudget>> = OnceLock::new();
27
28struct SubscriptionMemoryBudget {
29    limit: usize,
30    used: AtomicUsize,
31}
32
33struct ChargedUpdateInner {
34    update: Option<MvUpdate>,
35    bytes: usize,
36    budget: Arc<SubscriptionMemoryBudget>,
37}
38
39impl Drop for ChargedUpdateInner {
40    fn drop(&mut self) {
41        {
42            let _update = self.update.take();
43        }
44        self.budget.release(self.bytes);
45    }
46}
47
48#[derive(Clone)]
49pub(super) struct ChargedUpdate(Arc<ChargedUpdateInner>);
50
51impl ChargedUpdate {
52    /// Takes ownership of capacity that was reserved before this value was built.
53    fn from_reserved(
54        update: MvUpdate,
55        bytes: usize,
56        budget: Arc<SubscriptionMemoryBudget>,
57    ) -> Self {
58        Self(Arc::new(ChargedUpdateInner {
59            update: Some(update),
60            bytes,
61            budget,
62        }))
63    }
64
65    pub(super) fn as_ref(&self) -> &MvUpdate {
66        self.0
67            .update
68            .as_ref()
69            .expect("charged subscription update is present until final drop")
70    }
71}
72
73impl SubscriptionMemoryBudget {
74    fn new(limit: usize) -> Self {
75        Self {
76            limit,
77            used: AtomicUsize::new(0),
78        }
79    }
80
81    fn try_reserve(&self, bytes: usize) -> bool {
82        self.used
83            .fetch_update(Ordering::AcqRel, Ordering::Acquire, |used| {
84                (bytes <= self.limit.saturating_sub(used)).then(|| used.saturating_add(bytes))
85            })
86            .is_ok()
87    }
88
89    fn release(&self, bytes: usize) {
90        if bytes == 0 {
91            return;
92        }
93        let released = self
94            .used
95            .fetch_update(Ordering::AcqRel, Ordering::Acquire, |used| {
96                used.checked_sub(bytes)
97            });
98        debug_assert!(released.is_ok(), "subscription memory released twice");
99    }
100
101    #[cfg(test)]
102    fn used(&self) -> usize {
103        self.used.load(Ordering::Acquire)
104    }
105}
106
107#[derive(Debug)]
108pub(crate) enum MvUpdate {
109    Batch(RecordBatch),
110    Barrier {
111        epoch: u64,
112        checkpoint_id: u64,
113        /// First shared-log sequence not certified by this checkpoint.
114        through_sequence: u64,
115    },
116    Error(String),
117}
118
119/// Where a new subscriber should start reading.
120#[derive(Clone, Copy, Debug)]
121pub enum SubscribeStart {
122    /// See only entries sequenced after attachment.
123    Tail,
124    /// Replay entries strictly after the retained barrier with `epoch == n`.
125    AsOfEpoch(u64),
126}
127
128#[derive(Debug)]
129pub(crate) enum SubscriptionOpenError {
130    ReplayPruned {
131        /// Earliest barrier epoch eligible for replay; `0` if none is retained.
132        earliest_retained: u64,
133    },
134    EpochNotCommitted {
135        requested: u64,
136        latest_committed: Option<u64>,
137    },
138    Capacity {
139        attached: usize,
140        limit: usize,
141    },
142}
143
144struct SequencedUpdate {
145    sequence: u64,
146    bytes: usize,
147    update: ChargedUpdate,
148}
149
150struct StreamLog {
151    inner: Mutex<StreamLogInner>,
152    wake: watch::Sender<()>,
153    budget: Arc<SubscriptionMemoryBudget>,
154}
155
156struct StreamLogInner {
157    entries: VecDeque<SequencedUpdate>,
158    bytes: usize,
159    retention_cap: usize,
160    retention_floor: u64,
161    retention_bytes: usize,
162    next_sequence: u64,
163    next_reader_id: u64,
164    readers: HashMap<u64, u64>,
165    latest_committed_epoch: Option<u64>,
166    terminal_error: Option<String>,
167    reserved_marker: Option<CheckpointAttempt>,
168}
169
170#[must_use]
171enum AppendOutcome {
172    Stored,
173    Unobserved,
174    Rejected(String),
175}
176
177#[must_use]
178enum ReservedAppendOutcome {
179    Stored,
180    Unobserved,
181    Rejected(String),
182}
183
184impl StreamLog {
185    fn new(
186        retention_cap: usize,
187        budget: Arc<SubscriptionMemoryBudget>,
188        latest_committed_epoch: Option<u64>,
189    ) -> Self {
190        Self::new_at(retention_cap, budget, latest_committed_epoch, 0)
191    }
192
193    fn new_at(
194        retention_cap: usize,
195        budget: Arc<SubscriptionMemoryBudget>,
196        latest_committed_epoch: Option<u64>,
197        next_sequence: u64,
198    ) -> Self {
199        let (wake, _) = watch::channel(());
200        Self {
201            inner: Mutex::new(StreamLogInner {
202                entries: VecDeque::new(),
203                bytes: 0,
204                retention_cap,
205                retention_floor: next_sequence,
206                retention_bytes: 0,
207                next_sequence,
208                next_reader_id: 0,
209                readers: HashMap::new(),
210                latest_committed_epoch,
211                terminal_error: None,
212                reserved_marker: None,
213            }),
214            wake,
215            budget,
216        }
217    }
218
219    fn append(&self, update: MvUpdate) -> AppendOutcome {
220        let entry_bytes = approx_size(&update);
221        let mut inner = self.inner.lock();
222        debug_assert!(!matches!(&update, MvUpdate::Barrier { .. }));
223        if let Some(message) = &inner.terminal_error {
224            return AppendOutcome::Rejected(message.clone());
225        }
226        if inner.retention_cap == 0 && inner.readers.is_empty() {
227            return AppendOutcome::Unobserved;
228        }
229
230        let storage_cap = storage_cap(inner.retention_cap);
231        if entry_bytes > storage_cap {
232            let message = format!(
233                "subscription entry requires {entry_bytes} bytes, exceeding the internal shared-log limit of {storage_cap} bytes"
234            );
235            return self.reject_append(inner, message);
236        }
237
238        let required_sequence_space = if inner.reserved_marker.is_some() {
239            2
240        } else {
241            1
242        };
243        if inner
244            .next_sequence
245            .checked_add(required_sequence_space)
246            .is_none()
247        {
248            return self.reject_append(inner, "subscription sequence space exhausted".into());
249        }
250        let Some(next_sequence) = inner.next_sequence.checked_add(1) else {
251            unreachable!("sequence-space preflight admitted the next sequence");
252        };
253
254        reclaim_consumed_prefix(&mut inner);
255        let available_locally = storage_cap.saturating_sub(entry_bytes);
256        evict_to_byte_target(&mut inner, available_locally);
257
258        if !self.budget.try_reserve(entry_bytes) {
259            let message = format!(
260                "subscription process memory budget exhausted ({} byte limit); update was not delivered",
261                self.budget.limit
262            );
263            return self.reject_append(inner, message);
264        }
265
266        let sequence = inner.next_sequence;
267        inner.entries.push_back(SequencedUpdate {
268            sequence,
269            bytes: entry_bytes,
270            update: ChargedUpdate::from_reserved(update, entry_bytes, Arc::clone(&self.budget)),
271        });
272        inner.next_sequence = next_sequence;
273        inner.bytes = inner.bytes.saturating_add(entry_bytes);
274        retain_appended_entry(&mut inner, sequence, entry_bytes);
275        reclaim_consumed_prefix(&mut inner);
276        drop(inner);
277        self.wake.send_modify(|()| {});
278        AppendOutcome::Stored
279    }
280
281    fn reject_append(
282        &self,
283        mut inner: parking_lot::MutexGuard<'_, StreamLogInner>,
284        message: String,
285    ) -> AppendOutcome {
286        if inner.reserved_marker.is_some() {
287            return AppendOutcome::Rejected(message);
288        }
289        Self::terminate_locked(&mut inner, message.clone());
290        drop(inner);
291        self.wake.send_modify(|()| {});
292        AppendOutcome::Rejected(message)
293    }
294
295    fn reserve_marker(&self, attempt: CheckpointAttempt) -> Result<Option<u64>, String> {
296        let mut inner = self.inner.lock();
297        if inner.terminal_error.is_some() || (inner.retention_cap == 0 && inner.readers.is_empty())
298        {
299            return Ok(None);
300        }
301        if let Some(existing) = inner.reserved_marker {
302            return Err(format!(
303                "subscription marker already reserved for epoch={} checkpoint_id={}",
304                existing.epoch, existing.checkpoint_id
305            ));
306        }
307        if inner.next_sequence.checked_add(1).is_none() {
308            return Err("subscription sequence space cannot admit a checkpoint marker".into());
309        }
310        let through_sequence = inner.next_sequence;
311        inner.reserved_marker = Some(attempt);
312        Ok(Some(through_sequence))
313    }
314
315    fn cancel_marker(&self, attempt: CheckpointAttempt) -> bool {
316        let mut inner = self.inner.lock();
317        if inner.reserved_marker != Some(attempt) {
318            return false;
319        }
320        inner.reserved_marker = None;
321        true
322    }
323
324    fn append_reserved_marker(
325        &self,
326        attempt: CheckpointAttempt,
327        through_sequence: u64,
328    ) -> ReservedAppendOutcome {
329        let mut inner = self.inner.lock();
330        if inner.reserved_marker != Some(attempt) {
331            return ReservedAppendOutcome::Rejected(format!(
332                "subscription marker reservation mismatch for epoch={} checkpoint_id={}",
333                attempt.epoch, attempt.checkpoint_id
334            ));
335        }
336        inner.reserved_marker = None;
337
338        if let Some(message) = &inner.terminal_error {
339            return ReservedAppendOutcome::Rejected(format!(
340                "subscription log terminated before checkpoint marker publication: {message}"
341            ));
342        }
343        inner.latest_committed_epoch = Some(
344            inner
345                .latest_committed_epoch
346                .map_or(attempt.epoch, |latest| latest.max(attempt.epoch)),
347        );
348        if inner.retention_cap == 0 && inner.readers.is_empty() {
349            return ReservedAppendOutcome::Unobserved;
350        }
351        if through_sequence > inner.next_sequence {
352            return ReservedAppendOutcome::Rejected(format!(
353                "subscription checkpoint cut cursor {through_sequence} is ahead of log cursor {}",
354                inner.next_sequence
355            ));
356        }
357        let Some(next_sequence) = inner.next_sequence.checked_add(1) else {
358            return ReservedAppendOutcome::Rejected(
359                "subscription sequence space exhausted before checkpoint marker publication".into(),
360            );
361        };
362
363        let storage_cap = storage_cap(inner.retention_cap);
364        if BARRIER_ENTRY_BYTES > storage_cap {
365            return ReservedAppendOutcome::Rejected(format!(
366                "subscription checkpoint marker requires {BARRIER_ENTRY_BYTES} bytes, exceeding the internal shared-log limit of {storage_cap} bytes"
367            ));
368        }
369        reclaim_consumed_prefix(&mut inner);
370        evict_to_byte_target(&mut inner, storage_cap.saturating_sub(BARRIER_ENTRY_BYTES));
371
372        let sequence = inner.next_sequence;
373        inner.entries.push_back(SequencedUpdate {
374            sequence,
375            bytes: BARRIER_ENTRY_BYTES,
376            update: ChargedUpdate::from_reserved(
377                MvUpdate::Barrier {
378                    epoch: attempt.epoch,
379                    checkpoint_id: attempt.checkpoint_id,
380                    through_sequence,
381                },
382                BARRIER_ENTRY_BYTES,
383                Arc::clone(&self.budget),
384            ),
385        });
386        inner.next_sequence = next_sequence;
387        inner.bytes = inner.bytes.saturating_add(BARRIER_ENTRY_BYTES);
388        retain_appended_entry(&mut inner, sequence, BARRIER_ENTRY_BYTES);
389        reclaim_consumed_prefix(&mut inner);
390        drop(inner);
391        self.wake.send_modify(|()| {});
392        ReservedAppendOutcome::Stored
393    }
394
395    fn subscribe(
396        self: &Arc<Self>,
397        start: SubscribeStart,
398    ) -> Result<SubscriptionReader, SubscriptionOpenError> {
399        let mut inner = self.inner.lock();
400        let attached = inner.readers.len();
401        if attached >= super::MAX_SUBSCRIBERS_PER_MV {
402            return Err(SubscriptionOpenError::Capacity {
403                attached,
404                limit: super::MAX_SUBSCRIBERS_PER_MV,
405            });
406        }
407
408        let (cursor, skip_barrier) = match (inner.terminal_error.is_some(), start) {
409            (true, _) | (false, SubscribeStart::Tail) => (inner.next_sequence, None),
410            (false, SubscribeStart::AsOfEpoch(epoch)) => {
411                let (cursor, barrier_sequence) = cursor_after_retained_epoch(&inner, epoch)?;
412                (cursor, Some((epoch, barrier_sequence)))
413            }
414        };
415        let wake = self.wake.subscribe();
416        let mut reader_id = inner.next_reader_id;
417        while inner.readers.contains_key(&reader_id) {
418            reader_id = reader_id.wrapping_add(1);
419        }
420        inner.next_reader_id = reader_id.wrapping_add(1);
421        inner.readers.insert(reader_id, cursor);
422        drop(inner);
423
424        Ok(SubscriptionReader {
425            log: Arc::clone(self),
426            reader_id,
427            cursor,
428            skip_barrier,
429            wake,
430            registered: true,
431        })
432    }
433
434    fn set_retention_cap(&self, retention_cap: usize) {
435        let mut inner = self.inner.lock();
436        let previous_head = head_sequence(&inner);
437        inner.retention_cap = retention_cap;
438        recompute_retention_suffix(&mut inner);
439        reclaim_consumed_prefix(&mut inner);
440        evict_to_byte_target(&mut inner, storage_cap(retention_cap));
441        let head_changed = head_sequence(&inner) != previous_head;
442        drop(inner);
443        if head_changed {
444            self.wake.send_modify(|()| {});
445        }
446    }
447
448    fn terminate(&self, message: &str) {
449        let mut inner = self.inner.lock();
450        Self::terminate_locked(&mut inner, message.to_owned());
451        drop(inner);
452        self.wake.send_modify(|()| {});
453    }
454
455    fn terminate_and_replace(&self, message: &str, latest_committed_epoch: Option<u64>) -> Self {
456        let mut inner = self.inner.lock();
457        debug_assert!(
458            inner.reserved_marker.is_none(),
459            "subscription generation invalidated before its checkpoint marker was released"
460        );
461        let retention_cap = inner.retention_cap;
462        let next_sequence = inner.next_sequence;
463        Self::terminate_locked(&mut inner, message.to_owned());
464        drop(inner);
465        self.wake.send_modify(|()| {});
466        Self::new_at(
467            retention_cap,
468            Arc::clone(&self.budget),
469            latest_committed_epoch,
470            next_sequence,
471        )
472    }
473
474    fn subscriber_count(&self) -> usize {
475        self.inner.lock().readers.len()
476    }
477
478    fn observe_committed_epoch(&self, epoch: u64) {
479        let mut inner = self.inner.lock();
480        inner.latest_committed_epoch = Some(
481            inner
482                .latest_committed_epoch
483                .map_or(epoch, |latest| latest.max(epoch)),
484        );
485    }
486
487    fn terminate_locked(inner: &mut StreamLogInner, message: String) {
488        if inner.terminal_error.is_none() {
489            inner.terminal_error = Some(message);
490        }
491        clear_entries(inner);
492    }
493}
494
495impl Drop for StreamLog {
496    fn drop(&mut self) {
497        let inner = self.inner.get_mut();
498        debug_assert!(
499            inner.reserved_marker.is_none(),
500            "subscription log dropped with a live checkpoint marker reservation"
501        );
502    }
503}
504
505fn storage_cap(retention_cap: usize) -> usize {
506    retention_cap.max(INTERNAL_LIVE_LOG_BYTES)
507}
508
509fn head_sequence(inner: &StreamLogInner) -> u64 {
510    inner
511        .entries
512        .front()
513        .map_or(inner.next_sequence, |entry| entry.sequence)
514}
515
516fn clear_entries(inner: &mut StreamLogInner) {
517    inner.entries.clear();
518    inner.bytes = 0;
519    inner.retention_floor = inner.next_sequence;
520    inner.retention_bytes = 0;
521}
522
523fn pop_front(inner: &mut StreamLogInner) -> Option<SequencedUpdate> {
524    let entry = inner.entries.pop_front()?;
525    inner.bytes = inner.bytes.saturating_sub(entry.bytes);
526    if entry.sequence >= inner.retention_floor {
527        debug_assert_eq!(
528            entry.sequence, inner.retention_floor,
529            "retention floor skipped a retained entry"
530        );
531        inner.retention_bytes = inner.retention_bytes.saturating_sub(entry.bytes);
532        inner.retention_floor = entry.sequence.saturating_add(1);
533        if inner.retention_bytes == 0 {
534            inner.retention_floor = inner.next_sequence;
535        }
536    }
537    Some(entry)
538}
539
540fn evict_to_byte_target(inner: &mut StreamLogInner, target: usize) {
541    while inner.bytes > target {
542        let Some(_evicted) = pop_front(inner) else {
543            inner.bytes = 0;
544            inner.retention_floor = inner.next_sequence;
545            inner.retention_bytes = 0;
546            break;
547        };
548    }
549}
550
551fn reclaim_consumed_prefix(inner: &mut StreamLogInner) {
552    let reader_floor = inner
553        .readers
554        .values()
555        .copied()
556        .min()
557        .unwrap_or(inner.next_sequence);
558    let protected_floor = inner.retention_floor.min(reader_floor);
559
560    while inner
561        .entries
562        .front()
563        .is_some_and(|entry| entry.sequence < protected_floor)
564    {
565        let Some(_evicted) = pop_front(inner) else {
566            break;
567        };
568    }
569}
570
571fn calculate_retention_suffix(inner: &StreamLogInner) -> (u64, usize) {
572    if inner.retention_cap == 0 {
573        return (inner.next_sequence, 0);
574    }
575
576    let mut bytes = 0usize;
577    let mut floor = inner.next_sequence;
578    for entry in inner.entries.iter().rev() {
579        let with_entry = bytes.saturating_add(entry.bytes);
580        if with_entry > inner.retention_cap {
581            break;
582        }
583        bytes = with_entry;
584        floor = entry.sequence;
585    }
586    (floor, bytes)
587}
588
589fn recompute_retention_suffix(inner: &mut StreamLogInner) {
590    let (floor, bytes) = calculate_retention_suffix(inner);
591    inner.retention_floor = floor;
592    inner.retention_bytes = bytes;
593}
594
595fn retained_entry_bytes(inner: &StreamLogInner, sequence: u64) -> Option<usize> {
596    let head = head_sequence(inner);
597    if sequence < head || sequence >= inner.next_sequence {
598        return None;
599    }
600    let index = usize::try_from(sequence.saturating_sub(head)).ok()?;
601    inner.entries.get(index).map(|entry| entry.bytes)
602}
603
604fn retain_appended_entry(inner: &mut StreamLogInner, sequence: u64, bytes: usize) {
605    if inner.retention_cap == 0 {
606        inner.retention_floor = inner.next_sequence;
607        inner.retention_bytes = 0;
608        return;
609    }
610
611    if inner.retention_floor == sequence {
612        debug_assert_eq!(inner.retention_bytes, 0);
613        inner.retention_floor = sequence;
614    }
615    inner.retention_bytes = inner.retention_bytes.saturating_add(bytes);
616    while inner.retention_bytes > inner.retention_cap {
617        let Some(expired_bytes) = retained_entry_bytes(inner, inner.retention_floor) else {
618            debug_assert!(false, "retention floor is outside the shared log");
619            inner.retention_floor = inner.next_sequence;
620            inner.retention_bytes = 0;
621            return;
622        };
623        inner.retention_bytes = inner.retention_bytes.saturating_sub(expired_bytes);
624        inner.retention_floor = inner.retention_floor.saturating_add(1);
625    }
626    if inner.retention_bytes == 0 {
627        inner.retention_floor = inner.next_sequence;
628    }
629}
630
631fn cursor_after_retained_epoch(
632    inner: &StreamLogInner,
633    requested_epoch: u64,
634) -> Result<(u64, u64), SubscriptionOpenError> {
635    let Some(latest_committed) = inner.latest_committed_epoch else {
636        return Err(SubscriptionOpenError::EpochNotCommitted {
637            requested: requested_epoch,
638            latest_committed: None,
639        });
640    };
641    if requested_epoch > latest_committed {
642        return Err(SubscriptionOpenError::EpochNotCommitted {
643            requested: requested_epoch,
644            latest_committed: Some(latest_committed),
645        });
646    }
647    if inner.retention_cap == 0 {
648        return Err(SubscriptionOpenError::ReplayPruned {
649            earliest_retained: 0,
650        });
651    }
652
653    let mut cursor = None;
654    let mut earliest_retained = u64::MAX;
655    for entry in inner
656        .entries
657        .iter()
658        .filter(|entry| entry.sequence >= inner.retention_floor)
659    {
660        if let MvUpdate::Barrier {
661            epoch,
662            through_sequence,
663            ..
664        } = entry.update.as_ref()
665        {
666            if *through_sequence < inner.retention_floor {
667                continue;
668            }
669            earliest_retained = earliest_retained.min(*epoch);
670            if *epoch == requested_epoch {
671                cursor = Some((*through_sequence, entry.sequence));
672            }
673        }
674    }
675
676    if let Some(cursor) = cursor {
677        return Ok(cursor);
678    }
679    if earliest_retained == u64::MAX || requested_epoch < earliest_retained {
680        return Err(SubscriptionOpenError::ReplayPruned {
681            earliest_retained: if earliest_retained == u64::MAX {
682                0
683            } else {
684                earliest_retained
685            },
686        });
687    }
688    Err(SubscriptionOpenError::EpochNotCommitted {
689        requested: requested_epoch,
690        latest_committed: Some(latest_committed),
691    })
692}
693
694pub(super) fn approx_size(update: &MvUpdate) -> usize {
695    let payload = match update {
696        MvUpdate::Batch(batch) => batch.get_array_memory_size(),
697        MvUpdate::Barrier { .. } => return BARRIER_ENTRY_BYTES,
698        MvUpdate::Error(message) => message.len(),
699    };
700    payload.saturating_add(ENTRY_ACCOUNTING_BYTES)
701}
702
703pub(super) enum SubscriptionRead {
704    Update {
705        sequence: u64,
706        update: ChargedUpdate,
707    },
708    Lagged(u64),
709    Terminal(String),
710}
711
712enum TryRead {
713    Ready(SubscriptionRead),
714    Pending,
715}
716
717pub(crate) struct SubscriptionReader {
718    log: Arc<StreamLog>,
719    reader_id: u64,
720    cursor: u64,
721    /// `(epoch, physical sequence)` of the retained progress marker that an
722    /// AS-OF reader must skip while still replaying post-cut rows sequenced before it.
723    skip_barrier: Option<(u64, u64)>,
724    wake: watch::Receiver<()>,
725    registered: bool,
726}
727
728impl std::fmt::Debug for SubscriptionReader {
729    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
730        formatter
731            .debug_struct("SubscriptionReader")
732            .field("reader_id", &self.reader_id)
733            .field("cursor", &self.cursor)
734            .field("skip_barrier", &self.skip_barrier)
735            .field("registered", &self.registered)
736            .finish_non_exhaustive()
737    }
738}
739
740impl SubscriptionReader {
741    pub(super) async fn next(&mut self) -> SubscriptionRead {
742        loop {
743            match self.try_read() {
744                TryRead::Ready(read) => return read,
745                TryRead::Pending => {
746                    if self.wake.changed().await.is_err() {
747                        return SubscriptionRead::Terminal(
748                            "subscription shared log closed unexpectedly".into(),
749                        );
750                    }
751                }
752            }
753        }
754    }
755
756    fn try_read(&mut self) -> TryRead {
757        let mut inner = self.log.inner.lock();
758        if let Some(message) = &inner.terminal_error {
759            return TryRead::Ready(SubscriptionRead::Terminal(message.clone()));
760        }
761        if inner.readers.get(&self.reader_id).copied() != Some(self.cursor) {
762            return TryRead::Ready(SubscriptionRead::Terminal(
763                "subscription reader cursor registration invariant failed".into(),
764            ));
765        }
766        let head = head_sequence(&inner);
767        if self.cursor < head {
768            let mut skipped = head.saturating_sub(self.cursor);
769            if self
770                .skip_barrier
771                .is_some_and(|(_, sequence)| sequence >= self.cursor && sequence < head)
772            {
773                skipped = skipped.saturating_sub(1);
774                self.skip_barrier = None;
775                if skipped == 0 {
776                    self.cursor = head;
777                    inner.readers.insert(self.reader_id, self.cursor);
778                    reclaim_consumed_prefix(&mut inner);
779                    drop(inner);
780                    return self.try_read();
781                }
782            }
783            return TryRead::Ready(SubscriptionRead::Lagged(skipped));
784        }
785
786        if self.cursor < inner.next_sequence {
787            let Ok(index) = usize::try_from(self.cursor.saturating_sub(head)) else {
788                return TryRead::Ready(SubscriptionRead::Terminal(
789                    "subscription shared log index exceeds addressable memory".into(),
790                ));
791            };
792            let Some(entry) = inner.entries.get(index) else {
793                return TryRead::Ready(SubscriptionRead::Terminal(
794                    "subscription shared log sequence invariant failed".into(),
795                ));
796            };
797            if entry.sequence != self.cursor {
798                return TryRead::Ready(SubscriptionRead::Terminal(
799                    "subscription shared log is not contiguous".into(),
800                ));
801            }
802            let sequence = entry.sequence;
803            let update = entry.update.clone();
804            let skip = self.skip_barrier.is_some_and(|(epoch, sequence)| {
805                sequence == entry.sequence
806                    &&
807                matches!(update.as_ref(), MvUpdate::Barrier { epoch: seen, .. } if *seen == epoch)
808            });
809            self.cursor = self.cursor.saturating_add(1);
810            inner.readers.insert(self.reader_id, self.cursor);
811            reclaim_consumed_prefix(&mut inner);
812            drop(inner);
813            if skip {
814                self.skip_barrier = None;
815                return self.try_read();
816            }
817            return TryRead::Ready(SubscriptionRead::Update { sequence, update });
818        }
819
820        TryRead::Pending
821    }
822
823    fn release(&mut self) {
824        if !self.registered {
825            return;
826        }
827        let mut inner = self.log.inner.lock();
828        inner.readers.remove(&self.reader_id);
829        reclaim_consumed_prefix(&mut inner);
830        drop(inner);
831        self.registered = false;
832    }
833}
834
835impl Drop for SubscriptionReader {
836    fn drop(&mut self) {
837        self.release();
838    }
839}
840
841struct ReservedMarker {
842    log: Arc<StreamLog>,
843    through_sequence: u64,
844}
845
846struct ReservedCut {
847    attempt: CheckpointAttempt,
848    markers: Vec<ReservedMarker>,
849    reserved_bytes: usize,
850}
851
852fn require_canonical_checkpoint_attempt(attempt: CheckpointAttempt) -> Result<(), String> {
853    if attempt.is_canonical() {
854        return Ok(());
855    }
856    Err(format!(
857        "subscription checkpoint cut requires one nonzero canonical checkpoint ID; received epoch={} checkpoint_id={}",
858        attempt.epoch, attempt.checkpoint_id
859    ))
860}
861
862#[derive(Default)]
863struct RegistryLifecycle {
864    latest_committed_epoch: Option<u64>,
865    pending_cut: Option<ReservedCut>,
866}
867
868pub(crate) struct SubscriptionRegistry {
869    lifecycle: Mutex<RegistryLifecycle>,
870    streams: RwLock<HashMap<String, Arc<StreamLog>>>,
871    budget: Arc<SubscriptionMemoryBudget>,
872}
873
874impl SubscriptionRegistry {
875    pub(crate) fn new() -> Self {
876        let budget =
877            Arc::clone(PROCESS_SUBSCRIPTION_BUDGET.get_or_init(|| {
878                Arc::new(SubscriptionMemoryBudget::new(PROCESS_SUBSCRIPTION_BYTES))
879            }));
880        Self::with_budget(budget)
881    }
882
883    fn with_budget(budget: Arc<SubscriptionMemoryBudget>) -> Self {
884        Self {
885            lifecycle: Mutex::new(RegistryLifecycle::default()),
886            streams: RwLock::new(HashMap::new()),
887            budget,
888        }
889    }
890
891    #[cfg(test)]
892    pub(super) fn with_storage_budget(limit: usize) -> Self {
893        Self::with_budget(Arc::new(SubscriptionMemoryBudget::new(limit)))
894    }
895
896    /// Set the AS-OF retention budget. A zero budget still keeps a small
897    /// internal live suffix for attached portals but is never replay-eligible.
898    pub(crate) fn configure(&self, name: &str, cap: usize) {
899        let lifecycle = self.lifecycle.lock();
900        let log = self.get_or_create_locked(name, &lifecycle);
901        log.set_retention_cap(cap);
902    }
903
904    pub(crate) fn subscribe(
905        &self,
906        name: &str,
907        start: SubscribeStart,
908    ) -> Result<SubscriptionReader, SubscriptionOpenError> {
909        self.get_or_create(name).subscribe(start)
910    }
911
912    pub(crate) fn send_batch(&self, name: &str, batch: RecordBatch) -> Result<(), String> {
913        let Some(log) = self.streams.read().get(name).cloned() else {
914            return Ok(());
915        };
916        let bytes = batch.get_array_memory_size();
917        let outcome = if bytes > MAX_LIVE_BATCH_BYTES {
918            log.append(MvUpdate::Error(format!(
919                "subscription batch is {bytes} Arrow bytes, exceeding the internal live-batch limit of {MAX_LIVE_BATCH_BYTES} bytes; rows were not delivered"
920            )))
921        } else {
922            log.append(MvUpdate::Batch(batch))
923        };
924        match outcome {
925            AppendOutcome::Stored | AppendOutcome::Unobserved => Ok(()),
926            AppendOutcome::Rejected(message) => Err(format!(
927                "subscription update for '{name}' was rejected: {message}"
928            )),
929        }
930    }
931
932    /// Snapshot every live object's exact cursor at the aligned checkpoint cut.
933    pub(crate) fn reserve_cut(&self, attempt: CheckpointAttempt) -> Result<(), String> {
934        require_canonical_checkpoint_attempt(attempt)?;
935        let mut lifecycle = self.lifecycle.lock();
936        if let Some(existing) = &lifecycle.pending_cut {
937            return Err(format!(
938                "subscription cut already reserved for epoch={} checkpoint_id={}; cannot reserve epoch={} checkpoint_id={}",
939                existing.attempt.epoch,
940                existing.attempt.checkpoint_id,
941                attempt.epoch,
942                attempt.checkpoint_id
943            ));
944        }
945
946        let streams = self.streams.read();
947        let mut markers = Vec::with_capacity(streams.len());
948        for log in streams.values() {
949            match log.reserve_marker(attempt) {
950                Ok(Some(through_sequence)) => markers.push(ReservedMarker {
951                    log: Arc::clone(log),
952                    through_sequence,
953                }),
954                Ok(None) => {}
955                Err(error) => {
956                    for marker in &markers {
957                        let cancelled = marker.log.cancel_marker(attempt);
958                        debug_assert!(cancelled);
959                    }
960                    return Err(format!(
961                        "subscription cut reservation failed for epoch={} checkpoint_id={}: {error}",
962                        attempt.epoch, attempt.checkpoint_id
963                    ));
964                }
965            }
966        }
967        drop(streams);
968
969        let Some(reserved_bytes) = markers.len().checked_mul(BARRIER_ENTRY_BYTES) else {
970            for marker in &markers {
971                let cancelled = marker.log.cancel_marker(attempt);
972                debug_assert!(cancelled);
973            }
974            return Err("subscription checkpoint marker reservation size overflow".into());
975        };
976        if !self.budget.try_reserve(reserved_bytes) {
977            for marker in &markers {
978                let cancelled = marker.log.cancel_marker(attempt);
979                debug_assert!(cancelled);
980            }
981            return Err(format!(
982                "subscription checkpoint markers require {reserved_bytes} bytes, exceeding available process subscription memory ({} byte limit)",
983                self.budget.limit
984            ));
985        }
986        lifecycle.pending_cut = Some(ReservedCut {
987            attempt,
988            markers,
989            reserved_bytes,
990        });
991        Ok(())
992    }
993
994    /// Resolve a previously reserved cut after the checkpoint is durable.
995    pub(crate) fn commit_cut(&self, attempt: CheckpointAttempt) -> Result<(), String> {
996        require_canonical_checkpoint_attempt(attempt)?;
997        let mut lifecycle = self.lifecycle.lock();
998        let Some(pending) = &lifecycle.pending_cut else {
999            return Err(format!(
1000                "subscription cut missing for committed epoch={} checkpoint_id={}",
1001                attempt.epoch, attempt.checkpoint_id
1002            ));
1003        };
1004        if pending.attempt != attempt {
1005            return Err(format!(
1006                "subscription cut mismatch: reserved epoch={} checkpoint_id={}, committed epoch={} checkpoint_id={}",
1007                pending.attempt.epoch,
1008                pending.attempt.checkpoint_id,
1009                attempt.epoch,
1010                attempt.checkpoint_id
1011            ));
1012        }
1013        let cut = lifecycle
1014            .pending_cut
1015            .take()
1016            .expect("exact pending subscription cut was checked");
1017
1018        // The lifecycle guard freezes membership while a shared streams guard
1019        // leaves existing-log lookups available during marker publication.
1020        let logs = self.streams.read();
1021        lifecycle.latest_committed_epoch = Some(
1022            lifecycle
1023                .latest_committed_epoch
1024                .map_or(attempt.epoch, |current| current.max(attempt.epoch)),
1025        );
1026        for log in logs.values() {
1027            log.observe_committed_epoch(attempt.epoch);
1028        }
1029        let reserved_bytes = cut.reserved_bytes;
1030        let mut unused_reserved_bytes = 0usize;
1031        let mut failures = Vec::new();
1032        for marker in cut.markers {
1033            match marker
1034                .log
1035                .append_reserved_marker(attempt, marker.through_sequence)
1036            {
1037                ReservedAppendOutcome::Stored => {}
1038                ReservedAppendOutcome::Unobserved => {
1039                    unused_reserved_bytes = unused_reserved_bytes
1040                        .checked_add(BARRIER_ENTRY_BYTES)
1041                        .expect("reserved subscription marker accounting overflowed");
1042                }
1043                ReservedAppendOutcome::Rejected(error) => {
1044                    unused_reserved_bytes = unused_reserved_bytes
1045                        .checked_add(BARRIER_ENTRY_BYTES)
1046                        .expect("reserved subscription marker accounting overflowed");
1047                    failures.push(error);
1048                }
1049            }
1050        }
1051        drop(logs);
1052        debug_assert!(unused_reserved_bytes <= reserved_bytes);
1053        self.budget.release(unused_reserved_bytes);
1054        if failures.is_empty() {
1055            Ok(())
1056        } else {
1057            Err(format!(
1058                "subscription checkpoint marker publication failed for epoch={} checkpoint_id={}: {}",
1059                attempt.epoch,
1060                attempt.checkpoint_id,
1061                failures.join("; ")
1062            ))
1063        }
1064    }
1065
1066    pub(crate) fn abort_cut(&self, attempt: CheckpointAttempt) {
1067        if !attempt.is_canonical() {
1068            return;
1069        }
1070        let mut lifecycle = self.lifecycle.lock();
1071        if lifecycle
1072            .pending_cut
1073            .as_ref()
1074            .is_some_and(|cut| cut.attempt == attempt)
1075        {
1076            let cut = lifecycle
1077                .pending_cut
1078                .take()
1079                .expect("exact pending subscription cut was checked");
1080            release_reserved_cut(&cut, &self.budget);
1081        }
1082    }
1083
1084    /// End the current in-memory delivery generation before recovery can replay it.
1085    /// Existing readers receive a terminal error; replacement logs retain their
1086    /// configured byte caps and continue that object's current in-process cursor.
1087    pub(crate) fn invalidate_all(&self, reason: &str) {
1088        let mut lifecycle = self.lifecycle.lock();
1089        if let Some(cut) = lifecycle.pending_cut.take() {
1090            release_reserved_cut(&cut, &self.budget);
1091        }
1092        let mut streams = self.streams.write();
1093        let latest_committed_epoch = lifecycle.latest_committed_epoch;
1094        for log in streams.values_mut() {
1095            let replacement = log.terminate_and_replace(reason, latest_committed_epoch);
1096            *log = Arc::new(replacement);
1097        }
1098    }
1099
1100    #[cfg(test)]
1101    pub(crate) fn broadcast_barrier(&self, epoch: u64, checkpoint_id: u64) {
1102        let attempt = CheckpointAttempt::new(epoch, checkpoint_id);
1103        self.reserve_cut(attempt).unwrap();
1104        self.commit_cut(attempt).unwrap();
1105    }
1106
1107    pub(crate) fn drop_name(&self, name: &str) -> bool {
1108        let mut lifecycle = self.lifecycle.lock();
1109        let mut streams = self.streams.write();
1110        let Some(log) = streams.remove(name) else {
1111            return false;
1112        };
1113        if let Some(cut) = &mut lifecycle.pending_cut {
1114            let attempt = cut.attempt;
1115            let before = cut.markers.len();
1116            cut.markers.retain(|marker| {
1117                if Arc::ptr_eq(&marker.log, &log) {
1118                    let cancelled = marker.log.cancel_marker(attempt);
1119                    debug_assert!(cancelled);
1120                    false
1121                } else {
1122                    true
1123                }
1124            });
1125            let removed = before.saturating_sub(cut.markers.len());
1126            let released = removed
1127                .checked_mul(BARRIER_ENTRY_BYTES)
1128                .expect("reserved subscription marker accounting overflowed");
1129            cut.reserved_bytes = cut
1130                .reserved_bytes
1131                .checked_sub(released)
1132                .expect("released more subscription marker bytes than were reserved");
1133            self.budget.release(released);
1134        }
1135        log.terminate("object dropped");
1136        true
1137    }
1138
1139    pub(crate) fn subscriber_count(&self, name: &str) -> usize {
1140        self.streams
1141            .read()
1142            .get(name)
1143            .map_or(0, |log| log.subscriber_count())
1144    }
1145
1146    pub(crate) fn contains_name(&self, name: &str) -> bool {
1147        self.streams.read().contains_key(name)
1148    }
1149
1150    fn get_or_create(&self, name: &str) -> Arc<StreamLog> {
1151        if let Some(log) = self.streams.read().get(name) {
1152            return Arc::clone(log);
1153        }
1154        let lifecycle = self.lifecycle.lock();
1155        self.get_or_create_locked(name, &lifecycle)
1156    }
1157
1158    fn get_or_create_locked(&self, name: &str, lifecycle: &RegistryLifecycle) -> Arc<StreamLog> {
1159        Arc::clone(
1160            self.streams
1161                .write()
1162                .entry(name.to_owned())
1163                .or_insert_with(|| {
1164                    Arc::new(StreamLog::new(
1165                        0,
1166                        Arc::clone(&self.budget),
1167                        lifecycle.latest_committed_epoch,
1168                    ))
1169                }),
1170        )
1171    }
1172
1173    #[cfg(test)]
1174    fn head_sequence(&self, name: &str) -> Option<u64> {
1175        self.streams
1176            .read()
1177            .get(name)
1178            .map(|log| head_sequence(&log.inner.lock()))
1179    }
1180
1181    #[cfg(test)]
1182    fn next_sequence(&self, name: &str) -> Option<u64> {
1183        self.streams
1184            .read()
1185            .get(name)
1186            .map(|log| log.inner.lock().next_sequence)
1187    }
1188
1189    #[cfg(test)]
1190    pub(super) fn charged_bytes(&self) -> usize {
1191        self.budget.used()
1192    }
1193
1194    #[cfg(test)]
1195    fn assert_retention_cache(&self, name: &str) {
1196        let streams = self.streams.read();
1197        let log = streams.get(name).expect("test stream log must exist");
1198        let inner = log.inner.lock();
1199        assert_eq!(
1200            (inner.retention_floor, inner.retention_bytes),
1201            calculate_retention_suffix(&inner),
1202            "cached retention suffix diverged from a cold reverse scan"
1203        );
1204    }
1205}
1206
1207fn release_reserved_cut(cut: &ReservedCut, budget: &SubscriptionMemoryBudget) {
1208    debug_assert_eq!(
1209        cut.reserved_bytes,
1210        cut.markers
1211            .len()
1212            .checked_mul(BARRIER_ENTRY_BYTES)
1213            .expect("reserved subscription marker accounting overflowed"),
1214        "subscription marker reservation accounting diverged"
1215    );
1216    for marker in &cut.markers {
1217        let cancelled = marker.log.cancel_marker(cut.attempt);
1218        debug_assert!(cancelled);
1219    }
1220    budget.release(cut.reserved_bytes);
1221}
1222
1223impl Drop for SubscriptionRegistry {
1224    fn drop(&mut self) {
1225        if let Some(cut) = self.lifecycle.get_mut().pending_cut.take() {
1226            release_reserved_cut(&cut, &self.budget);
1227        }
1228    }
1229}
1230
1231impl Default for SubscriptionRegistry {
1232    fn default() -> Self {
1233        Self::new()
1234    }
1235}
1236
1237#[cfg(test)]
1238mod tests {
1239    use std::sync::Arc as StdArc;
1240
1241    use arrow_array::{ArrayRef, Int64Array};
1242    use arrow_schema::{DataType, Field, Schema};
1243
1244    use super::*;
1245
1246    fn batch(ids: Vec<i64>) -> RecordBatch {
1247        let schema = StdArc::new(Schema::new(vec![Field::new("id", DataType::Int64, false)]));
1248        RecordBatch::try_new(schema, vec![StdArc::new(Int64Array::from(ids))]).unwrap()
1249    }
1250
1251    fn earliest_retained(error: SubscriptionOpenError) -> u64 {
1252        match error {
1253            SubscriptionOpenError::ReplayPruned { earliest_retained } => earliest_retained,
1254            SubscriptionOpenError::EpochNotCommitted { .. } => {
1255                panic!("expected replay-pruned error")
1256            }
1257            SubscriptionOpenError::Capacity { .. } => panic!("expected replay-pruned error"),
1258        }
1259    }
1260
1261    async fn next_update(reader: &mut SubscriptionReader) -> ChargedUpdate {
1262        match reader.next().await {
1263            SubscriptionRead::Update { update, .. } => update,
1264            SubscriptionRead::Lagged(skipped) => panic!("unexpected gap of {skipped} entries"),
1265            SubscriptionRead::Terminal(message) => panic!("unexpected terminal error: {message}"),
1266        }
1267    }
1268
1269    #[tokio::test]
1270    async fn tail_starts_at_atomic_attach_cut() {
1271        let registry = SubscriptionRegistry::new();
1272        registry.configure("mv", 1 << 20);
1273        registry.send_batch("mv", batch(vec![1])).unwrap();
1274        let mut reader = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1275        registry.send_batch("mv", batch(vec![2])).unwrap();
1276
1277        let update = next_update(&mut reader).await;
1278        assert!(
1279            matches!(update.as_ref(), MvUpdate::Batch(batch) if batch.column(0).as_any().downcast_ref::<Int64Array>().unwrap().value(0) == 2)
1280        );
1281    }
1282
1283    #[tokio::test]
1284    async fn as_of_starts_strictly_after_exact_retained_barrier() {
1285        let registry = SubscriptionRegistry::new();
1286        registry.configure("mv", 1 << 20);
1287        registry.broadcast_barrier(1, 1);
1288        registry.send_batch("mv", batch(vec![10])).unwrap();
1289        registry.broadcast_barrier(2, 2);
1290
1291        let mut reader = registry
1292            .subscribe("mv", SubscribeStart::AsOfEpoch(1))
1293            .unwrap();
1294        assert!(matches!(
1295            next_update(&mut reader).await.as_ref(),
1296            MvUpdate::Batch(_)
1297        ));
1298        assert!(matches!(
1299            next_update(&mut reader).await.as_ref(),
1300            MvUpdate::Barrier {
1301                epoch: 2,
1302                checkpoint_id: 2,
1303                ..
1304            }
1305        ));
1306    }
1307
1308    #[tokio::test]
1309    async fn delayed_commit_preserves_the_aligned_cut_cursor() {
1310        let registry = SubscriptionRegistry::new();
1311        registry.configure("mv", 1 << 20);
1312        registry.send_batch("mv", batch(vec![10])).unwrap();
1313        let mut live = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1314
1315        let attempt = CheckpointAttempt::canonical(1);
1316        registry.reserve_cut(attempt).unwrap();
1317        registry.send_batch("mv", batch(vec![20])).unwrap();
1318        registry.commit_cut(attempt).unwrap();
1319
1320        assert!(matches!(
1321            next_update(&mut live).await.as_ref(),
1322            MvUpdate::Batch(batch)
1323                if batch.column(0).as_any().downcast_ref::<Int64Array>().unwrap().value(0) == 20
1324        ));
1325        assert!(matches!(
1326            next_update(&mut live).await.as_ref(),
1327            MvUpdate::Barrier {
1328                epoch: 1,
1329                checkpoint_id: 1,
1330                through_sequence: 1,
1331            }
1332        ));
1333
1334        let mut replay = registry
1335            .subscribe("mv", SubscribeStart::AsOfEpoch(1))
1336            .unwrap();
1337        assert!(matches!(
1338            next_update(&mut replay).await.as_ref(),
1339            MvUpdate::Batch(batch)
1340                if batch.column(0).as_any().downcast_ref::<Int64Array>().unwrap().value(0) == 20
1341        ));
1342    }
1343
1344    #[test]
1345    fn cut_reservation_fails_before_capture_when_marker_budget_is_full() {
1346        let sample = batch(vec![1, 2, 3]);
1347        let entry_bytes = approx_size(&MvUpdate::Batch(sample.clone()));
1348        let registry = SubscriptionRegistry::with_storage_budget(entry_bytes);
1349        registry.configure("mv", 1 << 20);
1350        registry.send_batch("mv", sample).unwrap();
1351        let attempt = CheckpointAttempt::canonical(1);
1352
1353        let error = registry.reserve_cut(attempt).unwrap_err();
1354
1355        assert!(error.contains("checkpoint markers require"));
1356        assert_eq!(registry.charged_bytes(), entry_bytes);
1357        let log = registry.streams.read().get("mv").cloned().unwrap();
1358        assert_eq!(log.inner.lock().reserved_marker, None);
1359        assert!(registry.lifecycle.lock().pending_cut.is_none());
1360    }
1361
1362    #[test]
1363    fn abort_releases_marker_headroom_for_the_next_attempt() {
1364        let registry = SubscriptionRegistry::with_storage_budget(BARRIER_ENTRY_BYTES);
1365        registry.configure("mv", BARRIER_ENTRY_BYTES);
1366        let first = CheckpointAttempt::canonical(1);
1367        registry.reserve_cut(first).unwrap();
1368        assert_eq!(registry.charged_bytes(), BARRIER_ENTRY_BYTES);
1369
1370        registry.abort_cut(first);
1371        assert_eq!(registry.charged_bytes(), 0);
1372
1373        let second = CheckpointAttempt::canonical(2);
1374        registry.reserve_cut(second).unwrap();
1375        assert_eq!(registry.charged_bytes(), BARRIER_ENTRY_BYTES);
1376        registry.abort_cut(second);
1377        assert_eq!(registry.charged_bytes(), 0);
1378    }
1379
1380    #[test]
1381    fn conflicting_attempt_cannot_steal_the_reserved_cut() {
1382        let registry = SubscriptionRegistry::with_storage_budget(BARRIER_ENTRY_BYTES);
1383        registry.configure("mv", BARRIER_ENTRY_BYTES);
1384        let reserved = CheckpointAttempt::canonical(1);
1385        let conflicting = CheckpointAttempt::canonical(2);
1386        registry.reserve_cut(reserved).unwrap();
1387
1388        assert!(registry.reserve_cut(conflicting).is_err());
1389        assert!(registry.commit_cut(conflicting).is_err());
1390        assert_eq!(registry.charged_bytes(), BARRIER_ENTRY_BYTES);
1391        assert_eq!(
1392            registry
1393                .lifecycle
1394                .lock()
1395                .pending_cut
1396                .as_ref()
1397                .map(|cut| cut.attempt),
1398            Some(reserved)
1399        );
1400
1401        registry.commit_cut(reserved).unwrap();
1402        assert_eq!(registry.next_sequence("mv"), Some(1));
1403        assert_eq!(registry.charged_bytes(), BARRIER_ENTRY_BYTES);
1404    }
1405
1406    #[test]
1407    fn noncanonical_cut_attempts_cannot_mutate_registry_state() {
1408        let registry = SubscriptionRegistry::with_storage_budget(BARRIER_ENTRY_BYTES);
1409        registry.configure("mv", BARRIER_ENTRY_BYTES);
1410        let invalid = CheckpointAttempt::new(1, 2);
1411        let log = registry.streams.read().get("mv").cloned().unwrap();
1412
1413        let error = registry.reserve_cut(invalid).unwrap_err();
1414
1415        assert!(error.contains("canonical checkpoint ID"));
1416        assert!(registry.lifecycle.lock().pending_cut.is_none());
1417        assert_eq!(log.inner.lock().reserved_marker, None);
1418        assert_eq!(registry.charged_bytes(), 0);
1419
1420        let canonical = CheckpointAttempt::canonical(1);
1421        registry.reserve_cut(canonical).unwrap();
1422        assert!(registry.commit_cut(invalid).is_err());
1423        registry.abort_cut(invalid);
1424
1425        assert_eq!(
1426            registry
1427                .lifecycle
1428                .lock()
1429                .pending_cut
1430                .as_ref()
1431                .map(|cut| cut.attempt),
1432            Some(canonical)
1433        );
1434        assert_eq!(log.inner.lock().reserved_marker, Some(canonical));
1435        assert_eq!(registry.charged_bytes(), BARRIER_ENTRY_BYTES);
1436
1437        registry.abort_cut(canonical);
1438        assert!(registry.lifecycle.lock().pending_cut.is_none());
1439        assert_eq!(log.inner.lock().reserved_marker, None);
1440        assert_eq!(registry.charged_bytes(), 0);
1441    }
1442
1443    #[test]
1444    fn sequence_exhaustion_rejects_cut_reservation_without_claiming_budget() {
1445        let registry = SubscriptionRegistry::with_storage_budget(BARRIER_ENTRY_BYTES);
1446        registry.configure("mv", BARRIER_ENTRY_BYTES);
1447        let log = registry.streams.read().get("mv").cloned().unwrap();
1448        {
1449            let mut inner = log.inner.lock();
1450            inner.next_sequence = u64::MAX;
1451            inner.retention_floor = u64::MAX;
1452        }
1453
1454        let error = registry
1455            .reserve_cut(CheckpointAttempt::canonical(1))
1456            .unwrap_err();
1457
1458        assert!(error.contains("sequence space"));
1459        assert_eq!(registry.charged_bytes(), 0);
1460        assert_eq!(log.inner.lock().reserved_marker, None);
1461    }
1462
1463    #[test]
1464    fn terminal_logs_do_not_claim_checkpoint_marker_capacity() {
1465        let registry = SubscriptionRegistry::with_storage_budget(BARRIER_ENTRY_BYTES);
1466        registry.configure("mv", BARRIER_ENTRY_BYTES);
1467        let log = registry.streams.read().get("mv").cloned().unwrap();
1468        log.terminate("injected terminal state");
1469        let attempt = CheckpointAttempt::canonical(1);
1470
1471        registry.reserve_cut(attempt).unwrap();
1472
1473        assert_eq!(registry.charged_bytes(), 0);
1474        assert_eq!(
1475            registry
1476                .lifecycle
1477                .lock()
1478                .pending_cut
1479                .as_ref()
1480                .map(|cut| cut.markers.len()),
1481            Some(0)
1482        );
1483        registry.commit_cut(attempt).unwrap();
1484    }
1485
1486    #[test]
1487    fn post_cut_data_cannot_consume_the_reserved_marker_sequence() {
1488        let sample = batch(vec![1]);
1489        let entry_bytes = approx_size(&MvUpdate::Batch(sample.clone()));
1490        let registry = SubscriptionRegistry::with_storage_budget(
1491            BARRIER_ENTRY_BYTES.saturating_add(entry_bytes),
1492        );
1493        registry.configure("mv", 1 << 20);
1494        let log = registry.streams.read().get("mv").cloned().unwrap();
1495        {
1496            let mut inner = log.inner.lock();
1497            inner.next_sequence = u64::MAX - 1;
1498            inner.retention_floor = u64::MAX - 1;
1499        }
1500        let attempt = CheckpointAttempt::canonical(1);
1501        registry.reserve_cut(attempt).unwrap();
1502
1503        let error = registry.send_batch("mv", sample).unwrap_err();
1504
1505        assert!(error.contains("sequence space exhausted"));
1506        {
1507            let inner = log.inner.lock();
1508            assert_eq!(inner.next_sequence, u64::MAX - 1);
1509            assert_eq!(inner.reserved_marker, Some(attempt));
1510            assert!(inner.terminal_error.is_none());
1511        }
1512        registry.commit_cut(attempt).unwrap();
1513        let inner = log.inner.lock();
1514        assert_eq!(inner.next_sequence, u64::MAX);
1515        assert_eq!(inner.reserved_marker, None);
1516    }
1517
1518    #[test]
1519    fn reserved_marker_survives_process_budget_contention() {
1520        let registry = SubscriptionRegistry::with_storage_budget(BARRIER_ENTRY_BYTES);
1521        registry.configure("mv", BARRIER_ENTRY_BYTES);
1522        let log = registry.streams.read().get("mv").cloned().unwrap();
1523        let attempt = CheckpointAttempt::canonical(1);
1524        registry.reserve_cut(attempt).unwrap();
1525
1526        let error = registry.send_batch("mv", batch(vec![1])).unwrap_err();
1527
1528        assert!(error.contains("process memory budget exhausted"));
1529        {
1530            let inner = log.inner.lock();
1531            assert_eq!(inner.next_sequence, 0);
1532            assert_eq!(inner.reserved_marker, Some(attempt));
1533            assert!(inner.terminal_error.is_none());
1534        }
1535        registry.commit_cut(attempt).unwrap();
1536        assert_eq!(registry.next_sequence("mv"), Some(1));
1537        assert_eq!(registry.charged_bytes(), BARRIER_ENTRY_BYTES);
1538    }
1539
1540    #[test]
1541    fn unobserved_commit_releases_reserved_marker_bytes() {
1542        let registry = SubscriptionRegistry::with_storage_budget(BARRIER_ENTRY_BYTES);
1543        let reader = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1544        let attempt = CheckpointAttempt::canonical(1);
1545        registry.reserve_cut(attempt).unwrap();
1546        assert_eq!(registry.charged_bytes(), BARRIER_ENTRY_BYTES);
1547
1548        drop(reader);
1549        registry.commit_cut(attempt).unwrap();
1550
1551        assert_eq!(registry.charged_bytes(), 0);
1552        assert_eq!(registry.next_sequence("mv"), Some(0));
1553    }
1554
1555    #[test]
1556    fn rejected_marker_commit_releases_reserved_bytes_and_reports_failure() {
1557        let registry = SubscriptionRegistry::with_storage_budget(BARRIER_ENTRY_BYTES);
1558        registry.configure("mv", BARRIER_ENTRY_BYTES);
1559        let log = registry.streams.read().get("mv").cloned().unwrap();
1560        let attempt = CheckpointAttempt::canonical(1);
1561        registry.reserve_cut(attempt).unwrap();
1562        log.inner.lock().terminal_error = Some("injected terminal state".into());
1563
1564        let error = registry.commit_cut(attempt).unwrap_err();
1565
1566        assert!(error.contains("terminated before checkpoint marker publication"));
1567        assert_eq!(registry.charged_bytes(), 0);
1568        assert_eq!(log.inner.lock().reserved_marker, None);
1569    }
1570
1571    #[test]
1572    fn invalidation_and_object_drop_release_exact_marker_reservations() {
1573        let budget = StdArc::new(SubscriptionMemoryBudget::new(BARRIER_ENTRY_BYTES));
1574        let registry = SubscriptionRegistry::with_budget(StdArc::clone(&budget));
1575        registry.configure("mv", BARRIER_ENTRY_BYTES);
1576        let first = CheckpointAttempt::canonical(1);
1577        registry.reserve_cut(first).unwrap();
1578
1579        registry.invalidate_all("injected recovery");
1580        assert_eq!(budget.used(), 0);
1581        assert!(registry.lifecycle.lock().pending_cut.is_none());
1582
1583        let second = CheckpointAttempt::canonical(2);
1584        registry.reserve_cut(second).unwrap();
1585        assert!(registry.drop_name("mv"));
1586        assert_eq!(budget.used(), 0);
1587        registry.commit_cut(second).unwrap();
1588    }
1589
1590    #[tokio::test]
1591    async fn recreated_object_is_outside_the_dropped_objects_reserved_cut() {
1592        let registry = SubscriptionRegistry::with_storage_budget(1 << 20);
1593        registry.configure("mv", 1 << 20);
1594        let mut dropped_reader = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1595        let attempt = CheckpointAttempt::canonical(1);
1596        registry.reserve_cut(attempt).unwrap();
1597
1598        assert!(registry.drop_name("mv"));
1599        assert!(matches!(
1600            dropped_reader.next().await,
1601            SubscriptionRead::Terminal(ref error) if error == "object dropped"
1602        ));
1603
1604        registry.configure("mv", 1 << 20);
1605        let mut recreated_reader = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1606        registry.commit_cut(attempt).unwrap();
1607        assert_eq!(registry.next_sequence("mv"), Some(0));
1608        assert!(matches!(recreated_reader.try_read(), TryRead::Pending));
1609
1610        registry.send_batch("mv", batch(vec![7])).unwrap();
1611        assert!(matches!(
1612            recreated_reader.next().await,
1613            SubscriptionRead::Update {
1614                sequence: 0,
1615                update,
1616            } if matches!(update.as_ref(), MvUpdate::Batch(_))
1617        ));
1618    }
1619
1620    #[test]
1621    fn registry_drop_releases_an_unresolved_marker_reservation() {
1622        let budget = StdArc::new(SubscriptionMemoryBudget::new(BARRIER_ENTRY_BYTES));
1623        {
1624            let registry = SubscriptionRegistry::with_budget(StdArc::clone(&budget));
1625            registry.configure("mv", BARRIER_ENTRY_BYTES);
1626            registry
1627                .reserve_cut(CheckpointAttempt::canonical(1))
1628                .unwrap();
1629            assert_eq!(budget.used(), BARRIER_ENTRY_BYTES);
1630        }
1631        assert_eq!(budget.used(), 0);
1632    }
1633
1634    #[tokio::test]
1635    async fn recovery_replacement_continues_the_current_object_sequence() {
1636        let registry = SubscriptionRegistry::new();
1637        registry.configure("mv", 1 << 20);
1638        let mut before_recovery = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1639
1640        registry.send_batch("mv", batch(vec![10])).unwrap();
1641        assert!(matches!(
1642            before_recovery.next().await,
1643            SubscriptionRead::Update {
1644                sequence: 0,
1645                update,
1646            } if matches!(update.as_ref(), MvUpdate::Batch(_))
1647        ));
1648        let abandoned = CheckpointAttempt::canonical(1);
1649        registry.reserve_cut(abandoned).unwrap();
1650
1651        registry.invalidate_all("injected recovery");
1652        assert!(matches!(
1653            before_recovery.next().await,
1654            SubscriptionRead::Terminal(message) if message == "injected recovery"
1655        ));
1656        assert!(registry.commit_cut(abandoned).is_err());
1657
1658        let mut after_recovery = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1659        registry.send_batch("mv", batch(vec![10])).unwrap();
1660        assert!(matches!(
1661            after_recovery.next().await,
1662            SubscriptionRead::Update {
1663                sequence: 1,
1664                update,
1665            } if matches!(update.as_ref(), MvUpdate::Batch(_))
1666        ));
1667
1668        let committed = CheckpointAttempt::canonical(2);
1669        registry.reserve_cut(committed).unwrap();
1670        registry.commit_cut(committed).unwrap();
1671        assert!(matches!(
1672            after_recovery.next().await,
1673            SubscriptionRead::Update {
1674                sequence: 2,
1675                update,
1676            } if matches!(
1677                update.as_ref(),
1678                MvUpdate::Barrier {
1679                    epoch: 2,
1680                    checkpoint_id: 2,
1681                    through_sequence: 2,
1682                }
1683            )
1684        ));
1685    }
1686
1687    #[test]
1688    fn as_of_classifies_future_missing_and_pruned_epochs() {
1689        let registry = SubscriptionRegistry::with_storage_budget(1 << 20);
1690        registry.configure("mv", 1 << 20);
1691        registry.broadcast_barrier(5, 5);
1692        registry.broadcast_barrier(7, 7);
1693
1694        assert!(matches!(
1695            registry
1696                .subscribe("mv", SubscribeStart::AsOfEpoch(8))
1697                .unwrap_err(),
1698            SubscriptionOpenError::EpochNotCommitted {
1699                requested: 8,
1700                latest_committed: Some(7)
1701            }
1702        ));
1703        assert!(matches!(
1704            registry
1705                .subscribe("mv", SubscribeStart::AsOfEpoch(6))
1706                .unwrap_err(),
1707            SubscriptionOpenError::EpochNotCommitted {
1708                requested: 6,
1709                latest_committed: Some(7)
1710            }
1711        ));
1712
1713        registry.configure(
1714            "mv",
1715            approx_size(&MvUpdate::Barrier {
1716                epoch: 0,
1717                checkpoint_id: 0,
1718                through_sequence: 0,
1719            }),
1720        );
1721        assert_eq!(
1722            earliest_retained(
1723                registry
1724                    .subscribe("mv", SubscribeStart::AsOfEpoch(5))
1725                    .unwrap_err()
1726            ),
1727            7
1728        );
1729    }
1730
1731    #[test]
1732    fn as_of_knows_latest_epoch_without_a_stored_log_entry() {
1733        let registry = SubscriptionRegistry::with_storage_budget(1 << 20);
1734        registry.broadcast_barrier(11, 11);
1735
1736        assert!(matches!(
1737            registry
1738                .subscribe("late", SubscribeStart::AsOfEpoch(12))
1739                .unwrap_err(),
1740            SubscriptionOpenError::EpochNotCommitted {
1741                requested: 12,
1742                latest_committed: Some(11)
1743            }
1744        ));
1745        assert_eq!(
1746            earliest_retained(
1747                registry
1748                    .subscribe("late", SubscribeStart::AsOfEpoch(11))
1749                    .unwrap_err()
1750            ),
1751            0
1752        );
1753    }
1754
1755    #[test]
1756    fn zero_retention_never_enables_as_of() {
1757        let registry = SubscriptionRegistry::new();
1758        let _reader = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1759        registry.broadcast_barrier(1, 1);
1760        registry.send_batch("mv", batch(vec![1])).unwrap();
1761
1762        let error = registry
1763            .subscribe("mv", SubscribeStart::AsOfEpoch(1))
1764            .unwrap_err();
1765        assert_eq!(earliest_retained(error), 0);
1766    }
1767
1768    #[tokio::test]
1769    async fn cached_retention_floor_matches_cold_scan_after_hot_path_updates() {
1770        let registry = SubscriptionRegistry::with_storage_budget(1 << 20);
1771        registry.configure("mv", 512);
1772        let mut reader = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1773
1774        for value in 0..32_i64 {
1775            registry.send_batch("mv", batch(vec![value])).unwrap();
1776            registry.assert_retention_cache("mv");
1777        }
1778        for _ in 0..16 {
1779            let _ = next_update(&mut reader).await;
1780            registry.assert_retention_cache("mv");
1781        }
1782
1783        registry.configure("mv", 4096);
1784        registry.assert_retention_cache("mv");
1785        registry.configure("mv", 128);
1786        registry.assert_retention_cache("mv");
1787    }
1788
1789    #[test]
1790    fn disabling_retention_without_readers_releases_the_log() {
1791        let registry = SubscriptionRegistry::new();
1792        registry.configure("mv", 1 << 20);
1793        let values: ArrayRef = StdArc::new(Int64Array::from(vec![1, 2, 3]));
1794        let schema = StdArc::new(Schema::new(vec![Field::new("id", DataType::Int64, false)]));
1795        registry
1796            .send_batch(
1797                "mv",
1798                RecordBatch::try_new(schema, vec![StdArc::clone(&values)]).unwrap(),
1799            )
1800            .unwrap();
1801        assert_eq!(StdArc::strong_count(&values), 2);
1802
1803        registry.configure("mv", 0);
1804
1805        assert_eq!(StdArc::strong_count(&values), 1);
1806    }
1807
1808    #[test]
1809    fn as_of_readers_do_not_clone_retained_arrow_batches() {
1810        let registry = SubscriptionRegistry::new();
1811        registry.configure("mv", 1 << 20);
1812        registry.broadcast_barrier(1, 1);
1813        let values: ArrayRef = StdArc::new(Int64Array::from(vec![1, 2, 3]));
1814        let schema = StdArc::new(Schema::new(vec![Field::new("id", DataType::Int64, false)]));
1815        registry
1816            .send_batch(
1817                "mv",
1818                RecordBatch::try_new(schema, vec![StdArc::clone(&values)]).unwrap(),
1819            )
1820            .unwrap();
1821        let owners_before_attach = StdArc::strong_count(&values);
1822
1823        let readers = (0..super::super::MAX_SUBSCRIBERS_PER_MV)
1824            .map(|_| {
1825                registry
1826                    .subscribe("mv", SubscribeStart::AsOfEpoch(1))
1827                    .unwrap()
1828            })
1829            .collect::<Vec<_>>();
1830
1831        assert_eq!(readers.len(), super::super::MAX_SUBSCRIBERS_PER_MV);
1832        assert_eq!(StdArc::strong_count(&values), owners_before_attach);
1833    }
1834
1835    #[tokio::test]
1836    async fn terminal_drop_releases_retained_storage_immediately() {
1837        let registry = SubscriptionRegistry::with_storage_budget(1 << 20);
1838        registry.configure("mv", 1 << 20);
1839        registry.broadcast_barrier(1, 1);
1840        let values: ArrayRef = StdArc::new(Int64Array::from(vec![1, 2, 3]));
1841        let schema = StdArc::new(Schema::new(vec![Field::new("id", DataType::Int64, false)]));
1842        registry
1843            .send_batch(
1844                "mv",
1845                RecordBatch::try_new(schema, vec![StdArc::clone(&values)]).unwrap(),
1846            )
1847            .unwrap();
1848        let mut reader = registry
1849            .subscribe("mv", SubscribeStart::AsOfEpoch(1))
1850            .unwrap();
1851        assert!(registry.charged_bytes() > 0);
1852        assert_eq!(StdArc::strong_count(&values), 2);
1853
1854        assert!(registry.drop_name("mv"));
1855
1856        assert_eq!(registry.charged_bytes(), 0);
1857        assert_eq!(StdArc::strong_count(&values), 1);
1858        assert!(matches!(
1859            reader.next().await,
1860            SubscriptionRead::Terminal(message) if message == "object dropped"
1861        ));
1862    }
1863
1864    #[tokio::test]
1865    async fn local_log_reclaims_while_the_charge_follows_reader_updates() {
1866        let registry = SubscriptionRegistry::with_storage_budget(1 << 20);
1867        let mut first = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1868        let mut second = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1869        registry.send_batch("mv", batch(vec![1, 2, 3])).unwrap();
1870        let log = registry.streams.read().get("mv").cloned().unwrap();
1871        let retained = registry.charged_bytes();
1872        assert!(retained > 0);
1873
1874        let first_frame = next_update(&mut first).await;
1875        assert!(matches!(first_frame.as_ref(), MvUpdate::Batch(_)));
1876        assert_eq!(registry.charged_bytes(), retained);
1877        assert_eq!(log.inner.lock().bytes, retained);
1878
1879        let second_frame = next_update(&mut second).await;
1880        assert!(matches!(second_frame.as_ref(), MvUpdate::Batch(_)));
1881        assert_eq!(registry.charged_bytes(), retained);
1882        assert_eq!(log.inner.lock().bytes, 0);
1883
1884        drop(first_frame);
1885        assert_eq!(registry.charged_bytes(), retained);
1886        drop(second_frame);
1887        assert_eq!(registry.charged_bytes(), 0);
1888    }
1889
1890    #[tokio::test]
1891    async fn process_budget_contention_fails_without_claim_and_release_is_reusable() {
1892        let sample = batch(vec![1, 2, 3]);
1893        let entry_bytes = approx_size(&MvUpdate::Batch(sample.clone()));
1894        let budget = StdArc::new(SubscriptionMemoryBudget::new(entry_bytes));
1895        let first_registry = SubscriptionRegistry::with_budget(StdArc::clone(&budget));
1896        let contender_registry = SubscriptionRegistry::with_budget(budget);
1897        first_registry.configure("first", entry_bytes.saturating_mul(4));
1898        contender_registry.configure("contender", entry_bytes.saturating_mul(4));
1899
1900        first_registry.send_batch("first", sample.clone()).unwrap();
1901        assert_eq!(first_registry.charged_bytes(), entry_bytes);
1902        let contender_sequence = contender_registry.next_sequence("contender").unwrap();
1903        assert!(contender_registry
1904            .send_batch("contender", sample.clone())
1905            .is_err());
1906        assert_eq!(
1907            contender_registry.next_sequence("contender"),
1908            Some(contender_sequence),
1909            "failed admission must not claim a sequence"
1910        );
1911        let mut contender = contender_registry
1912            .subscribe("contender", SubscribeStart::Tail)
1913            .unwrap();
1914        assert!(matches!(
1915            contender.next().await,
1916            SubscriptionRead::Terminal(message)
1917                if message.contains("process memory budget exhausted")
1918        ));
1919
1920        assert!(first_registry.drop_name("first"));
1921        assert_eq!(contender_registry.charged_bytes(), 0);
1922        contender_registry.configure("replacement", entry_bytes.saturating_mul(4));
1923        contender_registry
1924            .send_batch("replacement", sample)
1925            .unwrap();
1926        assert_eq!(contender_registry.charged_bytes(), entry_bytes);
1927        assert_eq!(contender_registry.next_sequence("replacement"), Some(1));
1928    }
1929
1930    #[tokio::test]
1931    async fn as_of_cursor_reports_exact_gap_after_live_byte_eviction() {
1932        let registry = SubscriptionRegistry::new();
1933        registry.configure("mv", 1024);
1934        registry.broadcast_barrier(1, 1);
1935        let mut reader = registry
1936            .subscribe("mv", SubscribeStart::AsOfEpoch(1))
1937            .unwrap();
1938
1939        let values_per_batch = (MAX_LIVE_BATCH_BYTES / 2) / std::mem::size_of::<i64>();
1940        for value in 0..6_i64 {
1941            registry
1942                .send_batch("mv", batch(vec![value; values_per_batch]))
1943                .unwrap();
1944        }
1945
1946        let head = registry.head_sequence("mv").unwrap();
1947        let expected = head.saturating_sub(1);
1948        assert!(
1949            expected > 0,
1950            "test must evict entries beyond the AS-OF cursor"
1951        );
1952        assert!(matches!(
1953            reader.next().await,
1954            SubscriptionRead::Lagged(skipped) if skipped == expected
1955        ));
1956    }
1957
1958    #[tokio::test]
1959    async fn dropping_name_is_a_visible_terminal_error() {
1960        let registry = SubscriptionRegistry::new();
1961        let mut reader = registry.subscribe("mv", SubscribeStart::Tail).unwrap();
1962        assert!(registry.drop_name("mv"));
1963        assert!(matches!(
1964            reader.next().await,
1965            SubscriptionRead::Terminal(message) if message == "object dropped"
1966        ));
1967    }
1968
1969    #[tokio::test]
1970    async fn oversized_batch_is_one_explicit_sequence_not_claim_then_evict() {
1971        let registry = SubscriptionRegistry::new();
1972        registry.configure("mv", INTERNAL_LIVE_LOG_BYTES);
1973        registry.broadcast_barrier(1, 1);
1974        let mut reader = registry
1975            .subscribe("mv", SubscribeStart::AsOfEpoch(1))
1976            .unwrap();
1977        let before = registry.next_sequence("mv").unwrap();
1978        let values = vec![0_i64; MAX_LIVE_BATCH_BYTES / std::mem::size_of::<i64>() + 1];
1979
1980        registry.send_batch("mv", batch(values)).unwrap();
1981
1982        assert_eq!(registry.next_sequence("mv"), Some(before + 1));
1983        assert!(matches!(
1984            next_update(&mut reader).await.as_ref(),
1985            MvUpdate::Error(message) if message.contains("rows were not delivered")
1986        ));
1987    }
1988
1989    #[test]
1990    fn subscriber_cap_is_atomic_across_65_simultaneous_attempts() {
1991        let registry = StdArc::new(SubscriptionRegistry::new());
1992        let attempts = super::super::MAX_SUBSCRIBERS_PER_MV + 1;
1993        let start = StdArc::new(std::sync::Barrier::new(attempts));
1994        let handles = (0..attempts)
1995            .map(|_| {
1996                let registry = StdArc::clone(&registry);
1997                let start = StdArc::clone(&start);
1998                std::thread::spawn(move || {
1999                    start.wait();
2000                    registry.subscribe("mv", SubscribeStart::Tail)
2001                })
2002            })
2003            .collect::<Vec<_>>();
2004
2005        let results = handles
2006            .into_iter()
2007            .map(|handle| handle.join().unwrap())
2008            .collect::<Vec<_>>();
2009        let successes = results.iter().filter(|result| result.is_ok()).count();
2010        let capacity_failures = results
2011            .iter()
2012            .filter(|result| matches!(result, Err(SubscriptionOpenError::Capacity { .. })))
2013            .count();
2014        assert_eq!(successes, super::super::MAX_SUBSCRIBERS_PER_MV);
2015        assert_eq!(capacity_failures, 1);
2016    }
2017}