use crate::{
auth::grpc::{AuthGrpcService, OAuthTokenSource},
retry_policy::RetryPredicate,
};
use std::fmt::Display;
pub use ::tonic::Status as Error;
pub use client_builder::{BuildError, MakeConnection, PubSubConfig, Uri};
pub use publish_sink::{PublishConfig, PublishError, PublishTopicSink, SinkError};
pub use streaming_subscription::{
AcknowledgeError, AcknowledgeToken, ModifyAcknowledgeError, StreamSubscription,
StreamSubscriptionConfig,
};
pub(crate) mod client_builder;
mod publish_sink;
mod streaming_subscription;
#[cfg(feature = "emulators")]
#[cfg_attr(docsrs, doc(cfg(feature = "emulators")))]
pub mod emulator;
#[allow(rustdoc::broken_intra_doc_links, rustdoc::bare_urls, missing_docs)]
pub mod api {
include!("../generated/google.pubsub.v1.rs");
pub use prost_types::{Duration, FieldMask, Timestamp};
}
#[derive(Debug, Clone)]
pub struct PublisherClient<C = crate::DefaultConnector> {
inner: api::publisher_client::PublisherClient<
AuthGrpcService<tonic::transport::Channel, OAuthTokenSource<C>>,
>,
}
impl<C> PublisherClient<C>
where
C: crate::Connect + Clone + Send + Sync + 'static,
{
pub fn publish_topic_sink(&mut self, topic: ProjectTopicName) -> PublishTopicSink<C> {
self.publish_topic_sink_config(topic, PublishConfig::default())
}
pub fn publish_topic_sink_config(
&mut self,
topic: ProjectTopicName,
config: PublishConfig,
) -> PublishTopicSink<C> {
PublishTopicSink::new(self.inner.clone(), topic, config)
}
}
impl<C> std::ops::Deref for PublisherClient<C> {
type Target = api::publisher_client::PublisherClient<
AuthGrpcService<tonic::transport::Channel, OAuthTokenSource<C>>,
>;
fn deref(&self) -> &Self::Target {
&self.inner
}
}
impl<C> std::ops::DerefMut for PublisherClient<C> {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.inner
}
}
#[derive(Debug, Clone)]
pub struct SubscriberClient<C = crate::DefaultConnector> {
inner: api::subscriber_client::SubscriberClient<
AuthGrpcService<tonic::transport::Channel, OAuthTokenSource<C>>,
>,
}
impl<C> SubscriberClient<C>
where
C: crate::Connect + Clone + Send + Sync + 'static,
{
pub fn stream_subscription(
&mut self,
subscription: ProjectSubscriptionName,
config: StreamSubscriptionConfig,
) -> StreamSubscription<C> {
StreamSubscription::new(self.inner.clone(), subscription.into(), config)
}
}
impl<C> std::ops::Deref for SubscriberClient<C> {
type Target = api::subscriber_client::SubscriberClient<
AuthGrpcService<tonic::transport::Channel, OAuthTokenSource<C>>,
>;
fn deref(&self) -> &Self::Target {
&self.inner
}
}
impl<C> std::ops::DerefMut for SubscriberClient<C> {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.inner
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct ProjectSubscriptionName(String);
impl ProjectSubscriptionName {
pub fn new(project_name: impl Display, subscription_name: impl Display) -> Self {
Self(format!(
"projects/{project}/subscriptions/{subscription}",
project = project_name,
subscription = subscription_name
))
}
}
impl From<ProjectSubscriptionName> for String {
fn from(from: ProjectSubscriptionName) -> String {
from.0
}
}
impl std::fmt::Display for ProjectSubscriptionName {
fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
std::fmt::Display::fmt(&self.0, f)
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct ProjectTopicName(String);
impl ProjectTopicName {
pub fn new(project_name: impl Display, topic_name: impl Display) -> Self {
Self(format!(
"projects/{project}/topics/{topic}",
project = project_name,
topic = topic_name,
))
}
}
impl From<ProjectTopicName> for String {
fn from(from: ProjectTopicName) -> String {
from.0
}
}
impl std::fmt::Display for ProjectTopicName {
fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
std::fmt::Display::fmt(&self.0, f)
}
}
#[derive(Debug, Default, Clone)]
pub struct PubSubRetryCheck {
_priv: (),
}
impl PubSubRetryCheck {
pub fn new() -> Self {
Self { _priv: () }
}
}
impl RetryPredicate<Error> for PubSubRetryCheck {
fn is_retriable(&self, error: &Error) -> bool {
use tonic::Code;
match error.code() {
Code::DeadlineExceeded
| Code::Internal
| Code::Cancelled
| Code::ResourceExhausted
| Code::Aborted
| Code::Unknown => true,
Code::Unavailable => {
let is_shutdown = error.message().contains("Server shutdownNow invoked");
!is_shutdown
}
_ => false,
}
}
}