use crate::{
builder, grpc,
pubsub::{api, PublisherClient, SubscriberClient},
};
const MAX_MESSAGE_SIZE: usize = 10 * 1024 * 1024;
#[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)]
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 builder::ClientBuilder {
async fn pubsub_authed_service(
&self,
config: PubSubConfig,
) -> Result<grpc::DefaultGrpcImpl, BuildError> {
let connection = tonic::transport::Endpoint::new(config.endpoint)?
.connect()
.await?;
Ok(grpc::DefaultGrpcImpl::new(
connection,
self.auth.clone(),
config.auth_scopes,
))
}
pub async fn build_pubsub_publisher(
&self,
config: PubSubConfig,
) -> Result<PublisherClient, BuildError> {
Ok(PublisherClient::from_raw_api(
api::publisher_client::PublisherClient::new(self.pubsub_authed_service(config).await?)
.max_decoding_message_size(MAX_MESSAGE_SIZE),
))
}
pub async fn build_pubsub_subscriber(
&self,
config: PubSubConfig,
) -> Result<SubscriberClient, BuildError> {
Ok(SubscriberClient::from_raw_api(
api::subscriber_client::SubscriberClient::new(
self.pubsub_authed_service(config).await?,
)
.max_decoding_message_size(MAX_MESSAGE_SIZE),
))
}
}
#[cfg(test)]
mod test {
use super::*;
#[test]
fn config_default() {
let config = PubSubConfig::default();
assert_eq!(config.endpoint, "https://pubsub.googleapis.com/v1");
}
}