uqa_client/notifications/
decoder.rs1mod 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#[derive(Debug)]
20pub struct DecodeStep {
21 pub consumed: usize,
22 pub event: Option<NotificationWireEvent>,
23}
24
25pub 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 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 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;