projectx_client/realtime/
event_flow.rs1use super::{
7 Arc, AtomicBool, AtomicUsize, Notify, Ordering, ParkingMutex, RealtimeError, RealtimeEvent,
8 fmt, mpsc,
9};
10
11#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
13pub struct RealtimeGeneration(pub(super) u64);
14
15#[derive(Clone, Debug, PartialEq)]
17#[non_exhaustive]
18pub struct RealtimeMessage {
19 pub generation: RealtimeGeneration,
21 pub event: RealtimeEvent,
23}
24
25pub(super) struct EventFlow {
26 pub(super) state: ParkingMutex<EventFlowState>,
27 pub(super) queued_weight: AtomicUsize,
28 overflowed: AtomicBool,
29 changed: Notify,
30}
31
32#[derive(Default)]
33pub(super) struct EventFlowState {
34 pub(super) active_generation: Option<u64>,
35 overflow: Option<OverflowFence>,
36}
37
38struct OverflowFence {
41 generation: u64,
42 start_event: Option<RealtimeEvent>,
43 terminal_event: Option<RealtimeEvent>,
44 generation_ended: bool,
45 acknowledged: bool,
46}
47
48struct EventReservation {
49 flow: Arc<EventFlow>,
50 weight: usize,
51}
52
53impl Drop for EventReservation {
54 fn drop(&mut self) {
55 let previous = self
56 .flow
57 .queued_weight
58 .fetch_sub(self.weight, Ordering::AcqRel);
59 debug_assert!(previous >= self.weight);
60 }
61}
62
63pub(super) struct EventEnvelope {
64 message: RealtimeMessage,
65 reservation: EventReservation,
66}
67
68#[derive(Clone, Copy, Debug, Eq, PartialEq)]
69pub(super) enum PublishOutcome {
70 Published,
71 StaleGeneration,
72}
73
74impl EventEnvelope {
75 fn into_message(self) -> RealtimeMessage {
76 let Self {
77 message,
78 reservation,
79 } = self;
80 drop(reservation);
81 message
82 }
83}
84
85impl EventFlow {
86 pub(super) fn new() -> Arc<Self> {
87 Arc::new(Self {
88 state: ParkingMutex::new(EventFlowState::default()),
89 queued_weight: AtomicUsize::new(0),
90 overflowed: AtomicBool::new(false),
91 changed: Notify::new(),
92 })
93 }
94
95 pub(super) fn has_unacknowledged_gap(&self) -> bool {
96 self.overflowed.load(Ordering::Acquire)
97 }
98
99 pub(super) fn start_generation(&self, generation: u64) -> Result<(), RealtimeError> {
100 let mut state = self.state.lock();
101 if state.overflow.is_some() {
102 return Err(RealtimeError::TransportGapPending);
103 }
104 if state.active_generation.is_some() {
105 return Err(RealtimeError::ConnectionCancelled);
106 }
107 state.active_generation = Some(generation);
108 Ok(())
109 }
110
111 pub(super) fn publish(
112 self: &Arc<Self>,
113 events: &mpsc::Sender<EventEnvelope>,
114 generation: u64,
115 event: RealtimeEvent,
116 weight: usize,
117 ) -> Result<PublishOutcome, RealtimeError> {
118 let mut state = self.state.lock();
119 if state.active_generation != Some(generation) {
120 return Ok(PublishOutcome::StaleGeneration);
121 }
122 if state.overflow.is_some() {
123 return self.retain_overflow(&mut state, generation, event);
124 }
125 self.queued_weight.fetch_add(weight, Ordering::AcqRel);
126 let envelope = EventEnvelope {
127 message: RealtimeMessage {
128 generation: RealtimeGeneration(generation),
129 event,
130 },
131 reservation: EventReservation {
132 flow: Arc::clone(self),
133 weight,
134 },
135 };
136 match events.try_send(envelope) {
137 Ok(()) => Ok(PublishOutcome::Published),
138 Err(mpsc::error::TrySendError::Full(envelope)) => {
139 self.retain_overflow(&mut state, generation, envelope.into_message().event)
140 }
141 Err(mpsc::error::TrySendError::Closed(_)) => Err(RealtimeError::EventReceiverClosed),
142 }
143 }
144
145 fn retain_overflow(
146 &self,
147 state: &mut EventFlowState,
148 generation: u64,
149 event: RealtimeEvent,
150 ) -> Result<PublishOutcome, RealtimeError> {
151 let fence = state.overflow.get_or_insert(OverflowFence {
152 generation,
153 start_event: None,
154 terminal_event: None,
155 generation_ended: false,
156 acknowledged: false,
157 });
158 let retained = match event {
159 RealtimeEvent::Connected | RealtimeEvent::Reconnected => {
160 fence.start_event = Some(event);
161 true
162 }
163 RealtimeEvent::Disconnected => {
164 fence.terminal_event = Some(event);
165 true
166 }
167 _ => false,
168 };
169 self.overflowed.store(true, Ordering::Release);
170 self.changed.notify_waiters();
171 if retained {
172 Ok(PublishOutcome::Published)
173 } else {
174 Err(RealtimeError::EventQueueFull)
175 }
176 }
177
178 pub(super) fn mark_gap(&self, generation: u64) {
179 let mut state = self.state.lock();
180 if state.active_generation == Some(generation) {
181 let _ = self.retain_overflow(&mut state, generation, RealtimeEvent::TransportGap);
182 }
183 }
184
185 pub(super) fn finish_generation(&self, generation: u64) {
186 let mut state = self.state.lock();
187 if state.active_generation != Some(generation) {
188 return;
189 }
190 state.active_generation = None;
191 if let Some(fence) = state.overflow.as_mut() {
192 fence.generation_ended = true;
193 }
194 drop(state);
195 self.changed.notify_waiters();
196 }
197
198 fn release_acknowledged(&self, state: &mut EventFlowState) {
199 if state.overflow.as_ref().is_some_and(|fence| {
200 fence.acknowledged && fence.start_event.is_none() && fence.terminal_event.is_none()
201 }) {
202 state.overflow = None;
203 self.overflowed.store(false, Ordering::Release);
204 }
205 }
206
207 fn acknowledge_gap(&self) {
208 let mut state = self.state.lock();
209 if let Some(fence) = state.overflow.as_mut() {
210 fence.acknowledged = true;
211 }
212 self.release_acknowledged(&mut state);
213 drop(state);
214 self.changed.notify_waiters();
215 }
216}
217
218pub struct RealtimeEventReceiver {
220 pub(super) events: mpsc::Receiver<EventEnvelope>,
221 pub(super) flow: Arc<EventFlow>,
222 pub(super) gap_reported: bool,
223}
224
225impl RealtimeEventReceiver {
226 #[must_use]
228 pub fn is_closed(&self) -> bool {
229 self.events.is_closed()
230 }
231
232 pub async fn recv(&mut self) -> Option<RealtimeEvent> {
235 self.recv_message().await.map(|message| message.event)
236 }
237
238 pub async fn recv_message(&mut self) -> Option<RealtimeMessage> {
241 loop {
242 let changed = self.flow.changed.notified();
243 tokio::pin!(changed);
244 let _ = changed.as_mut().enable();
245 {
246 let mut state = self.flow.state.lock();
247 if let Ok(envelope) = self.events.try_recv() {
248 return Some(envelope.into_message());
249 }
250 if let Some(fence) = state.overflow.as_mut() {
251 let event = if let Some(event) = fence.start_event.take() {
252 Some(event)
253 } else if !fence.acknowledged && !self.gap_reported {
254 self.gap_reported = true;
255 Some(RealtimeEvent::TransportGap)
256 } else if fence.generation_ended {
257 fence.terminal_event.take()
258 } else {
259 None
260 };
261 if let Some(event) = event {
262 let message = RealtimeMessage {
263 generation: RealtimeGeneration(fence.generation),
264 event,
265 };
266 self.flow.release_acknowledged(&mut state);
267 return Some(message);
268 }
269 }
270 if self.events.is_closed() {
271 return None;
272 }
273 }
274 tokio::select! {
275 biased;
276 envelope = self.events.recv() => {
277 if let Some(envelope) = envelope { return Some(envelope.into_message()); }
278 }
279 () = &mut changed => {}
280 }
281 }
282 }
283
284 pub fn acknowledge_transport_gap(&mut self) {
287 if self.gap_reported {
288 self.flow.acknowledge_gap();
289 self.gap_reported = false;
290 }
291 }
292}
293
294impl fmt::Debug for RealtimeEventReceiver {
295 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
296 formatter
297 .debug_struct("RealtimeEventReceiver")
298 .field("overflowed", &self.flow.has_unacknowledged_gap())
299 .field(
300 "queued_weight",
301 &self.flow.queued_weight.load(Ordering::Acquire),
302 )
303 .field("gap_reported", &self.gap_reported)
304 .finish_non_exhaustive()
305 }
306}