use futures::Stream;
use google_cloud_pubsub::subscriber::{MessageStream, ShutdownToken};
use ruststream::Subscriber;
use tokio::sync::mpsc;
use crate::broker::Core;
use crate::error::{PubSubError, box_err};
use crate::message::PubSubMessage;
use crate::subscription::PubSubSubscription;
const CHANNEL_CAPACITY: usize = 16;
pub struct PubSubSubscriber {
subscription: String,
rx: mpsc::Receiver<Result<PubSubMessage, PubSubError>>,
shutdown: ShutdownToken,
}
impl std::fmt::Debug for PubSubSubscriber {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PubSubSubscriber")
.field("subscription", &self.subscription)
.finish_non_exhaustive()
}
}
impl PubSubSubscriber {
#[must_use]
pub fn subscription(&self) -> &str {
&self.subscription
}
pub(crate) fn open(core: &Core, descriptor: &PubSubSubscription) -> Self {
let name = core.subscription_name(descriptor.subscription());
let mut builder = core.subscriber.subscribe(name.clone());
if let Some(messages) = descriptor.max_outstanding_value() {
builder = builder.set_max_outstanding_messages(messages);
}
if let Some(extension) = descriptor.ack_extension_value() {
builder = builder.set_max_lease_extension(extension);
}
let stream = builder.build();
let shutdown = stream.shutdown_token();
let (tx, rx) = mpsc::channel(CHANNEL_CAPACITY);
tokio::spawn(pump(stream, tx, name.clone()));
Self {
subscription: name,
rx,
shutdown,
}
}
}
impl Drop for PubSubSubscriber {
fn drop(&mut self) {
let token = self.shutdown.clone();
if let Ok(handle) = tokio::runtime::Handle::try_current() {
handle.spawn(async move {
token.shutdown().await;
});
}
}
}
impl Subscriber for PubSubSubscriber {
type Message = PubSubMessage;
type Error = PubSubError;
fn stream(&mut self) -> impl Stream<Item = Result<PubSubMessage, PubSubError>> + Send + '_ {
futures::stream::poll_fn(move |cx| self.rx.poll_recv(cx))
}
}
async fn pump(
mut stream: MessageStream,
out: mpsc::Sender<Result<PubSubMessage, PubSubError>>,
subscription: String,
) {
while let Some(item) = stream.next().await {
match item {
Ok((message, handler)) => {
if out
.send(Ok(PubSubMessage::new(message, handler)))
.await
.is_err()
{
break;
}
}
Err(err) => {
let _ = out
.send(Err(PubSubError::Receive {
subscription: subscription.clone(),
source: box_err(err),
}))
.await;
break;
}
}
}
}