sui_rpc/client/ledger_streams/
subscription.rs1use 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 buffered_subscription_frames: VecDeque<LiveFrame<A::Item, Progress<A::Cursor>>>,
33 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 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}