uqa_client/notifications/http/
subscription.rs1use 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
18pub 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 pub fn identity(&self) -> &NotificationIdentity {
33 &self.visible_identity
34 }
35 pub fn cancellation(&self) -> NotificationCancellation {
36 self.worker.cancellation.clone()
37 }
38
39 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 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}