Skip to main content

sui_rpc/client/ledger_streams/
subscription.rs

1use std::collections::VecDeque;
2
3use tonic::Status;
4use tonic::codegen::BoxStream;
5
6use super::super::Result;
7use super::adapter::LiveFrame;
8use super::adapter::Progress;
9use super::adapter::ProgressAdvance;
10use super::adapter::SubscriptionAdapter;
11use super::list::ListMachine;
12use super::observability::LedgerStreamEvent;
13use super::observability::LedgerStreamObservability;
14use super::observability::LedgerStreamStage;
15use super::types::LedgerStreamConfig;
16
17pub(super) struct LiveSubscription<A: SubscriptionAdapter> {
18    pub(super) stream: BoxStream<A::SubscribeResponse>,
19    pub(super) last_seen: Progress<A::Cursor>,
20}
21
22pub(super) enum BufferedSubscriptionState<A: SubscriptionAdapter> {
23    Active(LiveSubscription<A>),
24    Failed(Status),
25    DroppedAtBufferLimit,
26}
27
28pub(super) struct GapReplay<A: SubscriptionAdapter> {
29    pub(super) list_machine: ListMachine<A>,
30    pub(super) subscription_state: BufferedSubscriptionState<A>,
31    /// Item frames retain payloads; adjacent progress-only frames coalesce once per contiguous run.
32    buffered_subscription_frames: VecDeque<LiveFrame<A::Item, Progress<A::Cursor>>>,
33    /// Number of payload-bearing frames currently buffered, bounded by `max_buffered_live_items`
34    /// (empty progress-only frames do not count toward this limit).
35    buffered_subscription_items: usize,
36    replay_item_frontier: Option<A::ItemPosition>,
37    observability: LedgerStreamObservability,
38}
39
40impl<A: SubscriptionAdapter> GapReplay<A> {
41    pub(super) fn new(
42        list_machine: ListMachine<A>,
43        live: LiveSubscription<A>,
44        first: LiveFrame<A::Item, Progress<A::Cursor>>,
45        committed_item_position: Option<A::ItemPosition>,
46        config: &LedgerStreamConfig,
47        replays_boundary_item: bool,
48    ) -> Self {
49        let observability = list_machine.observability.clone();
50        let mut gap = Self {
51            list_machine,
52            subscription_state: BufferedSubscriptionState::Active(live),
53            buffered_subscription_frames: VecDeque::new(),
54            buffered_subscription_items: 0,
55            observability,
56            replay_item_frontier: committed_item_position,
57        };
58        if !replays_boundary_item {
59            gap.retain_frame(first, Some(config));
60        }
61        gap
62    }
63
64    fn retain_frame(
65        &mut self,
66        frame: LiveFrame<A::Item, Progress<A::Cursor>>,
67        config: Option<&LedgerStreamConfig>,
68    ) {
69        if frame.item.is_none()
70            && let Some(deferred_tail) = self.buffered_subscription_frames.back_mut()
71            && deferred_tail.item.is_none()
72        {
73            deferred_tail.progress = frame.progress;
74            return;
75        }
76
77        if frame.item.is_some() {
78            self.buffered_subscription_items = self.buffered_subscription_items.saturating_add(1);
79        }
80        self.buffered_subscription_frames.push_back(frame);
81        if let Some(config) = config
82            && self.buffered_subscription_items >= config.max_buffered_live_items.get()
83            && matches!(
84                &self.subscription_state,
85                BufferedSubscriptionState::Active(_)
86            )
87        {
88            self.subscription_state = BufferedSubscriptionState::DroppedAtBufferLimit;
89            let buffered_items = self.buffered_subscription_items;
90            let limit = config.max_buffered_live_items.get();
91            self.observability
92                .emit(|| LedgerStreamEvent::SubscriptionBufferLimitReached {
93                    family: A::FAMILY,
94                    buffered_items,
95                    limit,
96                });
97        }
98    }
99
100    pub(super) fn has_deferred_progress(&self, progress: &Progress<A::Cursor>) -> bool {
101        self.buffered_subscription_frames
102            .iter()
103            .any(|frame| frame.progress.same_position(progress))
104    }
105
106    pub(super) fn replayed_item_was_already_emitted(&self, item: &A::Item) -> bool {
107        self.replay_item_frontier
108            .as_ref()
109            .is_some_and(|frontier| A::item_position(item) <= frontier)
110    }
111
112    pub(super) fn record_replayed_item(&mut self, item: &A::Item) {
113        let position = A::item_position(item);
114        if self
115            .replay_item_frontier
116            .as_ref()
117            .is_none_or(|frontier| position > frontier)
118        {
119            self.replay_item_frontier = Some(position.clone());
120        }
121    }
122
123    /// Returns whether a valid buffered frame advanced subscription progress, controlling retry
124    /// backoff reset.
125    pub(super) fn buffer_subscription_result(
126        &mut self,
127        result: Option<Result<A::SubscribeResponse>>,
128        item_required: bool,
129        config: &LedgerStreamConfig,
130    ) -> bool {
131        let BufferedSubscriptionState::Active(live) = &mut self.subscription_state else {
132            return false;
133        };
134        let frame = match result {
135            Some(Ok(response)) => match A::parse_live(response, item_required) {
136                Ok(frame) => frame,
137                Err(status) => {
138                    self.subscription_state = BufferedSubscriptionState::Failed(status);
139                    return false;
140                }
141            },
142            Some(Err(status)) => {
143                self.observability
144                    .emit(|| LedgerStreamEvent::SubscriptionStreamInterrupted {
145                        family: A::FAMILY,
146                        stage: LedgerStreamStage::GapRecovery,
147                        status: status.clone(),
148                    });
149                self.subscription_state = BufferedSubscriptionState::Failed(status);
150                return false;
151            }
152            None => {
153                let status = Status::unavailable("subscription stream ended unexpectedly");
154                self.observability
155                    .emit(|| LedgerStreamEvent::SubscriptionStreamInterrupted {
156                        family: A::FAMILY,
157                        stage: LedgerStreamStage::GapRecovery,
158                        status: status.clone(),
159                    });
160                self.subscription_state = BufferedSubscriptionState::Failed(status);
161                return false;
162            }
163        };
164
165        match live
166            .last_seen
167            .classify_consecutive_subscription_progress(&frame.progress, frame.item.is_some())
168        {
169            Err(status) => {
170                self.subscription_state = BufferedSubscriptionState::Failed(status);
171                false
172            }
173            Ok(ProgressAdvance::Unchanged) => false,
174            Ok(ProgressAdvance::CheckpointCoverageAdvanced | ProgressAdvance::CursorAdvanced) => {
175                live.last_seen = frame.progress.clone();
176                self.retain_frame(frame, Some(config));
177                true
178            }
179        }
180    }
181
182    pub(super) fn into_buffered_subscription_drain(self) -> BufferedSubscriptionDrain<A> {
183        BufferedSubscriptionDrain {
184            subscription_state: self.subscription_state,
185            buffered_subscription_frames: self.buffered_subscription_frames,
186            replay_item_frontier: self.replay_item_frontier,
187        }
188    }
189}
190
191pub(super) struct BufferedSubscriptionDrain<A: SubscriptionAdapter> {
192    pub(super) subscription_state: BufferedSubscriptionState<A>,
193    pub(super) buffered_subscription_frames: VecDeque<LiveFrame<A::Item, Progress<A::Cursor>>>,
194    pub(super) replay_item_frontier: Option<A::ItemPosition>,
195}