Skip to main content

uqa_client/notifications/
wire.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Validated wire observations and closed diagnostic metadata.
8
9use super::{NotificationTiming, ProtocolError};
10use std::fmt;
11use uqa_core::notifications::{NotificationEvent, NotificationFailureKind, NotificationIdentity};
12
13#[derive(Clone, Debug, PartialEq, Eq)]
14pub struct NotificationReady {
15    pub identity: NotificationIdentity,
16    pub accepted_channel_count: usize,
17    pub timing: NotificationTiming,
18}
19
20/// A syntactically bounded remote failure. Unknown codes remain available explicitly but never appear in diagnostic formatting or acquire automatic retry authority.
21#[derive(Clone, PartialEq, Eq)]
22pub struct ServerFailure {
23    pub identity: NotificationIdentity,
24    code: Box<str>,
25    pub retryable: bool,
26}
27
28impl ServerFailure {
29    pub(super) fn new(
30        identity: NotificationIdentity,
31        code: String,
32        retryable: bool,
33    ) -> Result<Self, ProtocolError> {
34        if code.is_empty()
35            || code.len() > 64
36            || !code
37                .bytes()
38                .all(|byte| byte.is_ascii_uppercase() || byte.is_ascii_digit() || byte == b'_')
39        {
40            return Err(ProtocolError::InvalidFields);
41        }
42        Ok(Self {
43            identity,
44            code: code.into_boxed_str(),
45            retryable,
46        })
47    }
48
49    pub fn code(&self) -> &str {
50        &self.code
51    }
52
53    pub fn known_kind(&self) -> Option<NotificationFailureKind> {
54        use NotificationFailureKind as Kind;
55        [
56            Kind::InvalidRequest,
57            Kind::Authentication,
58            Kind::AuthorityRevoked,
59            Kind::Unsupported,
60            Kind::Capacity,
61            Kind::Backpressure,
62            Kind::Protocol,
63            Kind::SourceUnavailable,
64            Kind::Transport,
65            Kind::Timeout,
66            Kind::ServerDraining,
67            Kind::Cancelled,
68            Kind::SequenceExhausted,
69        ]
70        .into_iter()
71        .find(|kind| kind.code() == self.code.as_ref())
72    }
73}
74
75impl fmt::Debug for ServerFailure {
76    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
77        formatter
78            .debug_struct("ServerFailure")
79            .field("identity", &self.identity)
80            .field("kind", &self.known_kind())
81            .field("retryable", &self.retryable)
82            .finish_non_exhaustive()
83    }
84}
85
86#[derive(Clone, Debug, PartialEq, Eq)]
87pub enum NotificationWireEvent {
88    Ready(NotificationReady),
89    /// The decoder only produces the Core `Notification` variant; gap and reconnection events belong to transport lifecycle.
90    Notification(NotificationEvent),
91    Error(ServerFailure),
92    ServerDraining {
93        identity: NotificationIdentity,
94    },
95    /// An empty/comment-only SSE block; it consumes no notification sequence.
96    Heartbeat,
97}