Skip to main content

projectx_client/realtime/
event_flow.rs

1// SPDX-FileCopyrightText: 2026 Kevin Monaghan
2// SPDX-License-Identifier: MIT-0
3
4//! Bounded data delivery with retained continuity and lifecycle boundaries.
5
6use super::{
7    Arc, AtomicBool, AtomicUsize, Notify, Ordering, ParkingMutex, RealtimeError, RealtimeEvent,
8    fmt, mpsc,
9};
10
11/// Identity of one ready socket, scoped to its owning real-time client.
12#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
13pub struct RealtimeGeneration(pub(super) u64);
14
15/// An observation attributed to the exact socket that produced it.
16#[derive(Clone, Debug, PartialEq)]
17#[non_exhaustive]
18pub struct RealtimeMessage {
19    /// The originating socket's identity.
20    pub generation: RealtimeGeneration,
21    /// The observation or lifecycle boundary.
22    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
38// Two retained lifecycle slots suffice: an overflowing generation may start
39// and end, but another generation cannot start until this fence is consumed.
40struct 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
218/// Single-consumer bounded receiver for real-time events.
219pub 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    /// Returns whether every producer for this event stream is gone.
227    #[must_use]
228    pub fn is_closed(&self) -> bool {
229        self.events.is_closed()
230    }
231
232    /// Receives an event without its socket identity.
233    /// Use [`Self::recv_message`] when work may overlap transport replacement.
234    pub async fn recv(&mut self) -> Option<RealtimeEvent> {
235        self.recv_message().await.map(|message| message.event)
236    }
237
238    /// Receives the accepted prefix, then retained gap and lifecycle boundaries.
239    /// A disconnected boundary proves both socket tasks have stopped.
240    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    /// Resumes data admission after the caller installs its recovery boundary.
285    /// This never closes a socket. Completions and keepalives continue while fenced.
286    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}