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