use crate::grpc::{Body, BoxBody, Bytes, DefaultGrpcImpl, GrpcService, StdError};
use crate::retry_policy::RetryPredicate;
use std::fmt::Display;
use tracing::debug_span;
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,
unreachable_pub
)]
pub mod api {
include!("../generated/google.pubsub.v1.rs");
pub use prost_types::{Duration, FieldMask, Timestamp};
}
#[derive(Debug, Clone)]
pub struct PublisherClient<S = DefaultGrpcImpl> {
inner: api::publisher_client::PublisherClient<S>,
}
impl<S> PublisherClient<S> {
pub fn from_raw_api(client: api::publisher_client::PublisherClient<S>) -> Self {
PublisherClient { inner: client }
}
pub fn raw_api(&self) -> &api::publisher_client::PublisherClient<S> {
&self.inner
}
pub fn raw_api_mut(&mut self) -> &mut api::publisher_client::PublisherClient<S> {
&mut self.inner
}
}
impl<S> PublisherClient<S>
where
S: GrpcService<BoxBody> + Clone,
S::Error: Into<StdError>,
S::ResponseBody: Body<Data = Bytes> + Send + 'static,
<S::ResponseBody as Body>::Error: Into<StdError> + Send,
{
pub fn publish_topic_sink(
&mut self,
topic: ProjectTopicName,
config: PublishConfig,
) -> PublishTopicSink<S> {
PublishTopicSink::new(self.inner.clone(), topic, config)
}
}
#[derive(Debug, Clone)]
pub struct SubscriberClient<S = DefaultGrpcImpl> {
inner: api::subscriber_client::SubscriberClient<S>,
}
impl<S> SubscriberClient<S> {
pub fn from_raw_api(client: api::subscriber_client::SubscriberClient<S>) -> Self {
Self { inner: client }
}
pub fn raw_api(&self) -> &api::subscriber_client::SubscriberClient<S> {
&self.inner
}
pub fn raw_api_mut(&mut self) -> &mut api::subscriber_client::SubscriberClient<S> {
&mut self.inner
}
}
impl<S> SubscriberClient<S>
where
S: GrpcService<BoxBody> + Clone,
S::Error: Into<StdError>,
S::ResponseBody: Body<Data = Bytes> + Send + 'static,
<S::ResponseBody as Body>::Error: Into<StdError> + Send,
{
pub fn stream_subscription(
&mut self,
subscription: ProjectSubscriptionName,
config: StreamSubscriptionConfig,
) -> StreamSubscription<S> {
let sub_name: String = subscription.clone().into();
let span = debug_span!("create_subscription", topic = sub_name);
let _guard = span.enter();
StreamSubscription::new(
[
self.inner.clone(),
self.inner.clone(),
self.inner.clone(),
self.inner.clone(),
],
subscription.into(),
config,
)
}
}
#[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,
}
}
}