Skip to main content

uqa_client/notifications/http/
error.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7use crate::notifications::{ProtocolError, ServerFailure};
8use std::error::Error as _;
9use std::{fmt, sync::Arc, time::Duration};
10use uqa_core::notifications::{NotificationFailureKind as Kind, NotificationRequestId};
11
12#[derive(Clone, Copy, Debug, PartialEq, Eq)]
13pub enum NotificationTimeoutStage {
14    Connection,
15    Readiness,
16    Idle,
17    Reconnect,
18}
19
20/// A retained failure with content-free Display/Debug. Private transport and server diagnostics are available only through explicit accessors.
21#[derive(Clone)]
22pub struct HttpNotificationError(Arc<Failure>);
23
24struct Failure {
25    kind: Kind,
26    retryable: bool,
27    cause: Cause,
28}
29
30enum Cause {
31    Local,
32    Protocol(ProtocolError),
33    Transport(reqwest::Error),
34    Timeout(NotificationTimeoutStage),
35    Remote(ServerFailure),
36    Response {
37        status: u16,
38        code: Box<str>,
39        message: Box<str>,
40        request_id: NotificationRequestId,
41        retry_after: Option<Duration>,
42    },
43    Exhausted {
44        original: HttpNotificationError,
45        last: HttpNotificationError,
46        attempts: u32,
47    },
48}
49
50impl HttpNotificationError {
51    pub fn kind(&self) -> Kind {
52        self.0.kind
53    }
54    pub fn code(&self) -> &'static str {
55        self.kind().code()
56    }
57    pub fn timeout_stage(&self) -> Option<NotificationTimeoutStage> {
58        match &self.0.cause {
59            Cause::Timeout(stage) => Some(*stage),
60            Cause::Transport(error) if error.is_timeout() && error.is_connect() => {
61                Some(NotificationTimeoutStage::Connection)
62            }
63            _ => None,
64        }
65    }
66    pub fn protocol_error(&self) -> Option<ProtocolError> {
67        match &self.0.cause {
68            Cause::Protocol(error) => Some(*error),
69            _ => None,
70        }
71    }
72    pub fn transport_error(&self) -> Option<&reqwest::Error> {
73        match &self.0.cause {
74            Cause::Transport(error) => Some(error),
75            _ => None,
76        }
77    }
78    pub fn server_failure(&self) -> Option<&ServerFailure> {
79        match &self.0.cause {
80            Cause::Remote(error) => Some(error),
81            _ => None,
82        }
83    }
84    pub fn http_status(&self) -> Option<u16> {
85        match &self.0.cause {
86            Cause::Response { status, .. } => Some(*status),
87            _ => None,
88        }
89    }
90    pub fn server_code(&self) -> Option<&str> {
91        match &self.0.cause {
92            Cause::Response { code, .. } => Some(code),
93            Cause::Remote(error) => Some(error.code()),
94            _ => None,
95        }
96    }
97    pub fn server_message(&self) -> Option<&str> {
98        match &self.0.cause {
99            Cause::Response { message, .. } => Some(message),
100            _ => None,
101        }
102    }
103    pub fn request_id(&self) -> Option<&NotificationRequestId> {
104        match &self.0.cause {
105            Cause::Response { request_id, .. } => Some(request_id),
106            Cause::Remote(error) => error.identity.request_id.as_ref(),
107            _ => None,
108        }
109    }
110    pub fn original_failure(&self) -> Option<&Self> {
111        match &self.0.cause {
112            Cause::Exhausted { original, .. } => Some(original),
113            _ => None,
114        }
115    }
116    pub fn last_attempt_failure(&self) -> Option<&Self> {
117        match &self.0.cause {
118            Cause::Exhausted { last, .. } => Some(last),
119            _ => None,
120        }
121    }
122    pub fn reconnect_attempts(&self) -> Option<u32> {
123        match self.0.cause {
124            Cause::Exhausted { attempts, .. } => Some(attempts),
125            _ => None,
126        }
127    }
128    pub(super) fn retry_after(&self) -> Option<Duration> {
129        match self.0.cause {
130            Cause::Response { retry_after, .. } => retry_after,
131            _ => None,
132        }
133    }
134    pub(super) fn retryable(&self) -> bool {
135        self.0.retryable
136    }
137    fn new(kind: Kind, retryable: bool, cause: Cause) -> Self {
138        Self(Arc::new(Failure {
139            kind,
140            retryable,
141            cause,
142        }))
143    }
144    pub(super) fn local(kind: Kind) -> Self {
145        Self::new(kind, false, Cause::Local)
146    }
147    pub(crate) fn invalid_options() -> Self {
148        Self::local(Kind::InvalidRequest)
149    }
150    pub(super) fn cancelled() -> Self {
151        Self::local(Kind::Cancelled)
152    }
153    pub(super) fn timeout(stage: NotificationTimeoutStage) -> Self {
154        Self::new(Kind::Timeout, true, Cause::Timeout(stage))
155    }
156    pub(super) fn protocol(error: ProtocolError) -> Self {
157        let kind = match error {
158            ProtocolError::UnexpectedEnd => Kind::Transport,
159            ProtocolError::Allocation => Kind::Capacity,
160            _ => Kind::Protocol,
161        };
162        Self::new(
163            kind,
164            error == ProtocolError::UnexpectedEnd,
165            Cause::Protocol(error),
166        )
167    }
168    pub(super) fn request(error: ProtocolError) -> Self {
169        Self::new(
170            if error == ProtocolError::Allocation {
171                Kind::Capacity
172            } else {
173                Kind::InvalidRequest
174            },
175            false,
176            Cause::Protocol(error),
177        )
178    }
179    pub(super) fn transport(error: reqwest::Error) -> Self {
180        let corrupt = corrupted_transport(&error);
181        let retryable = !error.is_builder() && !corrupt;
182        Self::new(
183            if error.is_builder() {
184                Kind::InvalidRequest
185            } else if corrupt {
186                Kind::Protocol
187            } else if error.is_timeout() {
188                Kind::Timeout
189            } else {
190                Kind::Transport
191            },
192            retryable,
193            Cause::Transport(error.without_url()),
194        )
195    }
196    pub(super) fn remote(error: ServerFailure) -> Self {
197        let kind = error.known_kind().unwrap_or(Kind::Protocol);
198        let retryable = error.retryable
199            && matches!(
200                kind,
201                Kind::Capacity
202                    | Kind::SourceUnavailable
203                    | Kind::Transport
204                    | Kind::Timeout
205                    | Kind::ServerDraining
206            );
207        Self::new(kind, retryable, Cause::Remote(error))
208    }
209    pub(super) fn draining() -> Self {
210        Self::new(Kind::ServerDraining, true, Cause::Local)
211    }
212    pub(super) fn http(
213        status: u16,
214        code: String,
215        message: String,
216        request_id: NotificationRequestId,
217        retry_after: Option<Duration>,
218    ) -> Self {
219        let kind = match status {
220            400 | 409 | 413 => Kind::InvalidRequest,
221            401 => Kind::Authentication,
222            403 => Kind::AuthorityRevoked,
223            404 | 405 | 501 => Kind::Unsupported,
224            429 => Kind::Capacity,
225            503 => Kind::SourceUnavailable,
226            _ => Kind::Protocol,
227        };
228        let retryable = matches!(
229            (status, code.as_str()),
230            (429, "NOTIFICATION_CAPACITY") | (503, "NOTIFICATION_SOURCE_UNAVAILABLE")
231        );
232        Self::new(
233            kind,
234            retryable,
235            Cause::Response {
236                status,
237                code: code.into_boxed_str(),
238                message: message.into_boxed_str(),
239                request_id,
240                retry_after,
241            },
242        )
243    }
244    pub(super) fn exhausted(original: Self, last: Self, attempts: u32) -> Self {
245        Self::new(
246            last.kind(),
247            false,
248            Cause::Exhausted {
249                original,
250                last,
251                attempts,
252            },
253        )
254    }
255}
256
257impl fmt::Display for HttpNotificationError {
258    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
259        formatter.write_str(self.code())
260    }
261}
262
263impl fmt::Debug for HttpNotificationError {
264    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
265        formatter
266            .debug_struct("HttpNotificationError")
267            .field("kind", &self.kind())
268            .field("stage", &self.timeout_stage())
269            .field("attempts", &self.reconnect_attempts())
270            .finish_non_exhaustive()
271    }
272}
273
274impl std::error::Error for HttpNotificationError {}
275
276fn corrupted_transport(error: &reqwest::Error) -> bool {
277    let mut cause = error.source();
278    while let Some(error) = cause {
279        if error
280            .downcast_ref::<hyper::Error>()
281            .is_some_and(hyper::Error::is_parse)
282            || error.downcast_ref::<std::io::Error>().is_some_and(|error| {
283                matches!(
284                    error.kind(),
285                    std::io::ErrorKind::InvalidData | std::io::ErrorKind::InvalidInput
286                )
287            })
288        {
289            return true;
290        }
291        cause = error.source();
292    }
293    false
294}