Skip to main content

uqa_core/
notifications.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Notification values shared by embedded listeners and remote clients.
8//!
9//! These types do not register listeners, encode a wire protocol, or establish delivery. A producer must preserve its effective registration boundary and emit contiguous, positive sequences within each epoch. A replacement epoch does not imply replay.
10
11mod identity;
12
13pub use identity::{InvalidNotificationIdentity, NotificationEpoch, NotificationRequestId};
14
15/// One committed SQL notification waiting for this session.
16#[derive(Debug, Clone, PartialEq, Eq)]
17pub struct SQLNotification {
18    /// Stable backend process identifier of the sending SQL session.
19    pub process_id: i32,
20    /// Subscribed SQL channel that received the message.
21    pub channel: String,
22    /// Sender-provided payload, or the empty string when `NOTIFY` omitted it.
23    pub payload: String,
24}
25
26/// The identity of one ready subscription, independent of its delivery adapter.
27#[derive(Debug, Clone, PartialEq, Eq)]
28pub struct NotificationIdentity {
29    /// Fresh identity for this registration; it is not a durable cursor.
30    pub epoch: NotificationEpoch,
31    /// Present only when delivery has an HTTP request identity.
32    pub request_id: Option<NotificationRequestId>,
33}
34
35/// A closed failure category shared by language bindings and delivery adapters.
36///
37/// Adapters retain the original cause separately. A category does not grant retry authority: authentication, protocol and backpressure failures are terminal by default, and transport retries require the caller's remaining reconnect budget.
38#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
39pub enum NotificationFailureKind {
40    InvalidRequest,
41    Authentication,
42    AuthorityRevoked,
43    Unsupported,
44    Capacity,
45    Backpressure,
46    Protocol,
47    SourceUnavailable,
48    Transport,
49    Timeout,
50    ServerDraining,
51    Cancelled,
52    SequenceExhausted,
53}
54
55impl NotificationFailureKind {
56    /// Stable, content-free category for cross-language errors and diagnostics.
57    pub const fn code(self) -> &'static str {
58        match self {
59            Self::InvalidRequest => "NOTIFICATION_INVALID_REQUEST",
60            Self::Authentication => "NOTIFICATION_AUTHENTICATION",
61            Self::AuthorityRevoked => "NOTIFICATION_AUTHORITY_REVOKED",
62            Self::Unsupported => "NOTIFICATION_UNSUPPORTED",
63            Self::Capacity => "NOTIFICATION_CAPACITY",
64            Self::Backpressure => "NOTIFICATION_BACKPRESSURE",
65            Self::Protocol => "NOTIFICATION_PROTOCOL",
66            Self::SourceUnavailable => "NOTIFICATION_SOURCE_UNAVAILABLE",
67            Self::Transport => "NOTIFICATION_TRANSPORT",
68            Self::Timeout => "NOTIFICATION_TIMEOUT",
69            Self::ServerDraining => "NOTIFICATION_SERVER_DRAINING",
70            Self::Cancelled => "NOTIFICATION_CANCELLED",
71            Self::SequenceExhausted => "NOTIFICATION_SEQUENCE_EXHAUSTED",
72        }
73    }
74}
75
76/// An ordered subscription observation, including any loss of continuity.
77///
78/// Readiness precedes the first event. `ResyncRequired` names the old identity and precedes replacement data; `Reconnected` names the new ready identity. Embedded listeners do not reconnect implicitly. Channel and payload content is omitted from `Debug`; applications access it explicitly through the notification value.
79#[derive(Clone, PartialEq, Eq)]
80pub enum NotificationEvent {
81    Notification {
82        identity: NotificationIdentity,
83        /// Positive, contiguous counter within `identity.epoch`, starting at one. JavaScript adapters must preserve this exact integer as `bigint`.
84        sequence: u64,
85        notification: SQLNotification,
86    },
87    ResyncRequired {
88        identity: NotificationIdentity,
89        cause: NotificationFailureKind,
90    },
91    Reconnected {
92        identity: NotificationIdentity,
93    },
94}
95
96impl std::fmt::Debug for NotificationEvent {
97    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
98        match self {
99            Self::Notification {
100                identity, sequence, ..
101            } => f
102                .debug_struct("Notification")
103                .field("identity", identity)
104                .field("sequence", sequence)
105                .finish_non_exhaustive(),
106            Self::ResyncRequired { identity, cause } => f
107                .debug_struct("ResyncRequired")
108                .field("identity", identity)
109                .field("cause", cause)
110                .finish(),
111            Self::Reconnected { identity } => f
112                .debug_struct("Reconnected")
113                .field("identity", identity)
114                .finish(),
115        }
116    }
117}
118
119#[cfg(test)]
120mod tests;