1use 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
18const 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 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 through_sequence: u64,
115 },
116 Error(String),
117}
118
119#[derive(Clone, Copy, Debug)]
121pub enum SubscribeStart {
122 Tail,
124 AsOfEpoch(u64),
126}
127
128#[derive(Debug)]
129pub(crate) enum SubscriptionOpenError {
130 ReplayPruned {
131 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 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 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 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 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 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 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(®istry);
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}