use crate::{
auth::grpc,
builder,
pubsub::{api, PublisherClient, SubscriberClient},
};
#[allow(unreachable_pub)] pub use http::Uri;
#[allow(unreachable_pub)]
pub use tower::make::MakeConnection;
config_default! {
#[derive(Debug, Clone, Eq, PartialEq, Hash, serde::Deserialize)]
#[non_exhaustive]
pub struct PubSubConfig {
@default("https: pub endpoint: String,
@default(vec!["https: pub auth_scopes: Vec<String>,
}
}
#[derive(Debug, thiserror::Error)]
#[error(transparent)]
pub struct BuildError(#[from] tonic::transport::Error);
impl<C> builder::ClientBuilder<C>
where
C: MakeConnection<Uri> + crate::Connect + Clone + Send + Sync + 'static,
C::Connection: Unpin + Send + 'static,
C::Future: Send + 'static,
Box<dyn std::error::Error + Send + Sync + 'static>: From<C::Error>,
{
async fn pubsub_authed_service(
&self,
config: PubSubConfig,
) -> Result<
grpc::AuthGrpcService<tonic::transport::Channel, grpc::OAuthTokenSource<C>>,
BuildError,
> {
let connection = tonic::transport::Endpoint::new(config.endpoint)?
.connect_with_connector(self.connector.clone())
.await?;
Ok(grpc::oauth_grpc(
connection,
self.auth.clone(),
config.auth_scopes,
))
}
pub async fn build_pubsub_publisher(
&self,
config: PubSubConfig,
) -> Result<PublisherClient<C>, BuildError> {
Ok(PublisherClient {
inner: api::publisher_client::PublisherClient::new(
self.pubsub_authed_service(config).await?,
),
})
}
pub async fn build_pubsub_subscriber(
&self,
config: PubSubConfig,
) -> Result<SubscriberClient<C>, BuildError> {
Ok(SubscriberClient {
inner: api::subscriber_client::SubscriberClient::new(
self.pubsub_authed_service(config).await?,
),
})
}
}
#[cfg(test)]
mod test {
use super::*;
#[test]
fn config_default() {
let config = PubSubConfig::default();
assert_eq!(config.endpoint, "https://pubsub.googleapis.com/v1");
}
}