Skip to main content

uqa_client/notifications/http/
subscription.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7use super::{
8    attempt::Connection, inbox::Inbox, reconnect, HttpNotificationError as Error,
9    HttpNotificationOptions, NotificationCancellation,
10};
11use crate::notifications::{NotificationReady, SubscriptionRequest};
12use reqwest::Url;
13use secrecy::SecretString;
14use std::{fmt, sync::Arc};
15use tokio::{sync::oneshot, task::JoinHandle};
16use uqa_core::notifications::{NotificationEvent, NotificationFailureKind, NotificationIdentity};
17
18/// A bounded independently read HTTP subscription. Closing joins its one worker; dropping signals cancellation and aborts that same task. Received events never execute SQL.
19pub struct HttpNotificationSubscription {
20    inbox: Arc<Inbox>,
21    worker: WorkerOwner,
22    initial_ready: NotificationReady,
23    visible_identity: NotificationIdentity,
24    closed: bool,
25}
26
27impl HttpNotificationSubscription {
28    pub fn initial_ready(&self) -> &NotificationReady {
29        &self.initial_ready
30    }
31    /// Identity visible to this consumer; advances when Reconnected is consumed, never ahead of its gap event.
32    pub fn identity(&self) -> &NotificationIdentity {
33        &self.visible_identity
34    }
35    pub fn cancellation(&self) -> NotificationCancellation {
36        self.worker.cancellation.clone()
37    }
38
39    /// Cancelling this receive future leaves the subscription alive. Close/drop or its explicit cancellation signal stops the owned transport.
40    pub async fn next_event(&mut self) -> Result<Option<NotificationEvent>, Error> {
41        if self.closed {
42            return Ok(None);
43        }
44        loop {
45            let notified = self.inbox.changed.notified();
46            tokio::pin!(notified);
47            notified.as_mut().enable();
48            if self.worker.cancellation.is_cancelled() {
49                return Err(Error::cancelled());
50            }
51            if let Some(result) = self.inbox.take() {
52                if let Ok(Some(NotificationEvent::Reconnected { identity })) = &result {
53                    self.visible_identity = identity.clone();
54                }
55                return result;
56            }
57            tokio::select! {
58                biased;
59                () = self.worker.cancellation.cancelled() => return Err(Error::cancelled()),
60                () = &mut notified => {}
61            }
62        }
63    }
64
65    /// Completes local worker/response cleanup. It does not acknowledge a remote listener-cleanup transaction.
66    pub async fn close(&mut self) -> Result<(), Error> {
67        self.worker.cancellation.cancel();
68        self.inbox.finish(Ok(()));
69        self.worker.join().await?;
70        self.closed = true;
71        Ok(())
72    }
73}
74
75impl fmt::Debug for HttpNotificationSubscription {
76    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
77        formatter
78            .debug_struct("HttpNotificationSubscription")
79            .field("identity", &self.visible_identity)
80            .field("closed", &self.closed)
81            .finish_non_exhaustive()
82    }
83}
84
85struct WorkerOwner {
86    cancellation: NotificationCancellation,
87    task: Option<JoinHandle<()>>,
88}
89
90impl WorkerOwner {
91    async fn join(&mut self) -> Result<(), Error> {
92        if let Some(task) = self.task.as_mut() {
93            let result = task.await;
94            self.task = None;
95            result.map_err(|_| Error::local(NotificationFailureKind::SourceUnavailable))?;
96        }
97        Ok(())
98    }
99}
100
101impl Drop for WorkerOwner {
102    fn drop(&mut self) {
103        self.cancellation.cancel();
104        if let Some(task) = &self.task {
105            task.abort();
106        }
107    }
108}
109
110pub(super) struct Completion {
111    initial: Option<oneshot::Sender<Result<NotificationReady, Error>>>,
112    pub inbox: Arc<Inbox>,
113    cancellation: NotificationCancellation,
114    finished: bool,
115}
116
117impl Completion {
118    pub fn ready(&mut self, ready: &NotificationReady) -> Result<(), Error> {
119        if self.cancellation.is_cancelled() {
120            return Err(Error::cancelled());
121        }
122        if let Some(initial) = self.initial.take() {
123            initial
124                .send(Ok(ready.clone()))
125                .map_err(|_| Error::cancelled())
126        } else {
127            self.inbox.push(NotificationEvent::Reconnected {
128                identity: ready.identity.clone(),
129            })
130        }
131    }
132
133    fn finish(&mut self, error: Error) {
134        self.inbox.finish(Err(error.clone()));
135        if let Some(initial) = self.initial.take() {
136            let _ = initial.send(Err(error));
137        }
138        self.finished = true;
139    }
140}
141
142impl Drop for Completion {
143    fn drop(&mut self) {
144        if !self.finished {
145            self.finish(if self.cancellation.is_cancelled() {
146                Error::cancelled()
147            } else {
148                Error::local(NotificationFailureKind::SourceUnavailable)
149            });
150        }
151    }
152}
153
154pub(crate) async fn subscribe(
155    endpoint: Url,
156    credential: SecretString,
157    channels: &[&str],
158    options: HttpNotificationOptions,
159    cancellation: NotificationCancellation,
160) -> Result<HttpNotificationSubscription, Error> {
161    if cancellation.is_cancelled() {
162        return Err(Error::cancelled());
163    }
164    options.validate()?;
165    tokio::runtime::Handle::try_current().map_err(|_| Error::invalid_options())?;
166    let request = Arc::new(
167        SubscriptionRequest::new(channels, options.channel_limit()).map_err(Error::request)?,
168    );
169    let inbox = Arc::new(Inbox::new(&options)?);
170    let connection = Connection::new(endpoint, credential, request, options)?;
171    let (sender, receiver) = oneshot::channel();
172    let mut completion = Completion {
173        initial: Some(sender),
174        inbox: Arc::clone(&inbox),
175        cancellation: cancellation.clone(),
176        finished: false,
177    };
178    let worker_cancel = cancellation.clone();
179    let task = tokio::spawn(async move {
180        let error = tokio::select! {
181            biased;
182            () = worker_cancel.cancelled() => Error::cancelled(),
183            error = reconnect::run(&connection, &mut completion) => error,
184        };
185        completion.finish(error);
186    });
187    let mut worker = WorkerOwner {
188        cancellation,
189        task: Some(task),
190    };
191    let ready = receiver
192        .await
193        .unwrap_or_else(|_| Err(Error::local(NotificationFailureKind::SourceUnavailable)));
194    match ready {
195        Ok(initial_ready) => Ok(HttpNotificationSubscription {
196            inbox,
197            worker,
198            visible_identity: initial_ready.identity.clone(),
199            initial_ready,
200            closed: false,
201        }),
202        Err(error) => {
203            worker.cancellation.cancel();
204            worker.join().await?;
205            Err(error)
206        }
207    }
208}