Skip to main content

uqa_client/notifications/
decoder.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Strict identities, readiness and contiguous sequence admission over bounded SSE framing.
8
9mod envelope;
10
11use super::{
12    framing::{self, Framer},
13    NotificationReady, NotificationWireEvent, ProtocolError, SubscriptionRequest, TimerLimits,
14};
15use std::sync::Arc;
16use uqa_core::notifications::NotificationRequestId;
17
18/// At most one observation is returned. The caller retains and resubmits the unconsumed suffix of its HTTP chunk; no chunk-sized copy or event queue is hidden in the decoder.
19#[derive(Debug)]
20pub struct DecodeStep {
21    pub consumed: usize,
22    pub event: Option<NotificationWireEvent>,
23}
24
25/// A single response epoch. This validates the wire contract; HTTP security, actual server readiness, clocks and reconnection belong to the transport adapter.
26pub struct NotificationDecoder {
27    request: Arc<SubscriptionRequest>,
28    expected_request_id: NotificationRequestId,
29    timer_limits: TimerLimits,
30    framer: Framer,
31    ready: Option<NotificationReady>,
32    sequence: u64,
33    terminal: bool,
34    ended: bool,
35    failure: Option<ProtocolError>,
36}
37
38impl NotificationDecoder {
39    pub fn new(
40        request: Arc<SubscriptionRequest>,
41        expected_request_id: NotificationRequestId,
42        timer_limits: TimerLimits,
43    ) -> Result<Self, ProtocolError> {
44        Ok(Self {
45            request,
46            expected_request_id,
47            timer_limits,
48            framer: Framer::new()?,
49            ready: None,
50            sequence: 0,
51            terminal: false,
52            ended: false,
53            failure: None,
54        })
55    }
56
57    /// Returns zero consumed bytes only for empty input or a completed CR-delimited frame that was waiting to disambiguate its exact byte limit. In the latter case an event is returned, so callers can make progress without discarding input.
58    pub fn decode(&mut self, input: &[u8]) -> Result<DecodeStep, ProtocolError> {
59        if let Some(error) = self.failure {
60            return Err(error);
61        }
62        match self.decode_inner(input) {
63            Ok(step) => Ok(step),
64            Err(error) => self.fail(error),
65        }
66    }
67
68    fn decode_inner(&mut self, input: &[u8]) -> Result<DecodeStep, ProtocolError> {
69        if self.ended {
70            return Err(ProtocolError::UnexpectedEnd);
71        }
72        if self.terminal {
73            let consumed = self.framer.consume_terminal_lf(input);
74            return if consumed == input.len() {
75                Ok(DecodeStep {
76                    consumed,
77                    event: None,
78                })
79            } else {
80                Err(ProtocolError::EventOrder)
81            };
82        }
83        let step = self.framer.next(input)?;
84        let event = if step.complete {
85            Some(self.admit_frame()?)
86        } else {
87            None
88        };
89        Ok(DecodeStep {
90            consumed: step.consumed,
91            event,
92        })
93    }
94
95    /// Call at EOF, then again if this returns an event. EOF can finalize a bare CR at the exact frame limit, but never dispatches an unterminated event. A live stream without a terminal frame ends with `UnexpectedEnd`.
96    pub fn finish(&mut self) -> Result<Option<NotificationWireEvent>, ProtocolError> {
97        if let Some(error) = self.failure {
98            return Err(error);
99        }
100        self.ended = true;
101        let result = if self.framer.complete_at_end() {
102            self.admit_frame().map(Some)
103        } else if self.terminal && self.framer.frame().is_empty() {
104            Ok(None)
105        } else {
106            Err(self.framer.end_error())
107        };
108        result.or_else(|error| self.fail(error))
109    }
110
111    pub fn ready(&self) -> Option<&NotificationReady> {
112        self.ready.as_ref()
113    }
114
115    pub fn is_terminal(&self) -> bool {
116        self.terminal || self.failure.is_some()
117    }
118
119    fn fail<T>(&mut self, error: ProtocolError) -> Result<T, ProtocolError> {
120        self.failure = Some(error);
121        self.framer.clear();
122        Err(error)
123    }
124
125    fn admit_frame(&mut self) -> Result<NotificationWireEvent, ProtocolError> {
126        let event = match framing::fields(self.framer.frame())? {
127            Some(fields) => envelope::admit(
128                &fields,
129                &self.request,
130                &self.expected_request_id,
131                self.timer_limits,
132                self.ready.as_ref(),
133                self.sequence,
134            )?,
135            None => NotificationWireEvent::Heartbeat,
136        };
137        match &event {
138            NotificationWireEvent::Ready(ready) => self.ready = Some(ready.clone()),
139            NotificationWireEvent::Notification(
140                uqa_core::notifications::NotificationEvent::Notification { sequence, .. },
141            ) => self.sequence = *sequence,
142            NotificationWireEvent::Error(_) | NotificationWireEvent::ServerDraining { .. } => {
143                self.terminal = true;
144            }
145            _ => {}
146        }
147        self.framer.clear();
148        Ok(event)
149    }
150}
151
152#[cfg(test)]
153mod tests;