use std::ops::Range;
use crate::{
builder::{AuthFlow, ClientBuilder, ClientBuilderConfig},
pubsub,
};
use rand::{self, Rng};
use std::path::Path;
type BoxError = Box<dyn std::error::Error + Send + Sync + 'static>;
const PUBSUB_PROJECT_ID: &str = "test-project";
const PORT_RANGE: Range<usize> = 8000..12000;
const HOST: &str = "localhost";
const PUBSUB_CLI_RETRY: usize = 100;
const CLIENT_CONNECT_RETRY: usize = 50;
pub struct EmulatorClient {
_child: tokio::process::Child,
port: String,
_temp: tempdir::TempDir,
builder: ClientBuilder,
project_name: String,
}
impl EmulatorClient {
pub async fn new() -> Result<Self, BoxError> {
Self::with_project(PUBSUB_PROJECT_ID).await
}
pub async fn with_project(project_name: impl Into<String>) -> Result<Self, BoxError> {
let temp = tempdir::TempDir::new("pubsub_emulator")?;
let project_name = project_name.into();
let (child, port) = start_emulator(temp.path(), &project_name).await?;
let mut err: Option<tonic::transport::Error> = None;
for _ in 0..CLIENT_CONNECT_RETRY {
match create_schema_client(&port).await {
Err(e) => {
err = Some(e);
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
Ok(_) => {
err = None;
break;
}
}
}
if let Some(err) = err {
return Err(err.into());
}
let config = ClientBuilderConfig {
auth_flow: AuthFlow::NoAuth,
};
let builder = ClientBuilder::new(config).await?;
Ok(Self {
_child: child,
port,
builder,
_temp: temp,
project_name,
})
}
pub fn endpoint(&self) -> String {
format!("http://{}:{}/v1", HOST, self.port)
}
pub fn project(&self) -> &str {
&self.project_name
}
pub fn builder(&self) -> &ClientBuilder {
&self.builder
}
pub async fn create_topic(&self, topic_name: impl AsRef<str>) -> Result<(), BoxError> {
let config = pubsub::PubSubConfig {
endpoint: self.endpoint(),
..pubsub::PubSubConfig::default()
};
let mut publisher = self.builder().build_pubsub_publisher(config).await?;
publisher
.create_topic(pubsub::api::Topic {
name: pubsub::ProjectTopicName::new(self.project(), topic_name.as_ref()).into(),
..pubsub::api::Topic::default()
})
.await?;
Ok(())
}
}
impl Drop for EmulatorClient {
fn drop(&mut self) {
tokio::spawn(
tokio::process::Command::new("pkill")
.arg("-f")
.arg(format!(".*pubsub.*--port={}.*", self.port))
.output(),
);
}
}
async fn start_emulator(
tmp_dir: &Path,
project_name: &str,
) -> Result<(tokio::process::Child, String), std::io::Error> {
let mut err: Option<_> = None;
let mut rng = rand::thread_rng();
for _ in 0..PUBSUB_CLI_RETRY {
let port = format!("{}", rng.gen_range(PORT_RANGE));
match start_emulator_once(&port, tmp_dir, project_name) {
Ok(child) => return Ok((child, port)),
Err(e) => {
err = Some(e);
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
}
}
Err(err.unwrap())
}
fn start_emulator_once(
port: &str,
tmp_dir: &Path,
project_name: &str,
) -> Result<tokio::process::Child, std::io::Error> {
tokio::process::Command::new("gcloud")
.arg("beta")
.arg("emulators")
.arg("pubsub")
.arg("start")
.arg("--project")
.arg(project_name)
.arg("--host-port")
.arg(format!("{}:{}", HOST, port))
.arg("--data-dir")
.arg(tmp_dir)
.arg("--verbosity")
.arg("debug")
.kill_on_drop(true)
.spawn()
}
async fn create_schema_client(
port: &str,
) -> Result<
pubsub::api::schema_service_client::SchemaServiceClient<
impl tonic::client::GrpcService<
tonic::body::BoxBody,
Error = tonic::transport::Error,
ResponseBody = tonic::transport::Body,
>,
>,
tonic::transport::Error,
> {
pubsub::api::schema_service_client::SchemaServiceClient::connect(format!(
"http://{}:{}",
HOST, port
))
.await
}