use std::{ffi::OsString, ops::Range};
use crate::builder::{AuthFlow, ClientBuilder, ClientBuilderConfig};
use futures::future::BoxFuture;
use rand::{self, Rng};
use tracing::debug;
type BoxError = Box<dyn std::error::Error + Send + Sync + 'static>;
const PROJECT_ID: &str = "test-project";
const PORT_RANGE: Range<usize> = 8000..12000;
pub(crate) const HOST: &str = "localhost";
const CLI_RETRY: usize = 100;
const CLIENT_CONNECT_RETRY: usize = 50;
#[derive(Clone)]
pub(crate) struct EmulatorData {
pub(crate) gcloud_param: &'static str,
pub(crate) extra_args: Vec<OsString>,
pub(crate) kill_pattern: &'static str,
pub(crate) availability_check: fn(&str) -> BoxFuture<Result<(), tonic::transport::Error>>,
}
pub(crate) struct EmulatorClient {
_child: tokio::process::Child,
port: String,
builder: ClientBuilder,
project_name: String,
data: EmulatorData,
}
impl EmulatorClient {
pub(crate) async fn new(data: EmulatorData) -> Result<Self, BoxError> {
Self::with_project(data, PROJECT_ID).await
}
pub(crate) async fn with_project(
data: EmulatorData,
project_name: impl Into<String>,
) -> Result<Self, BoxError> {
let project_name = project_name.into();
let (child, port) =
start_emulator(data.gcloud_param, &project_name, &data.extra_args).await?;
debug!("Started emulator");
let mut err: Option<tonic::transport::Error> = None;
let check = data.availability_check;
for _ in 0..CLIENT_CONNECT_RETRY {
match check(&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,
project_name,
data,
})
}
pub(crate) fn endpoint(&self) -> String {
format!("http://{}:{}/v1", HOST, self.port)
}
pub(crate) fn project(&self) -> &str {
&self.project_name
}
pub(crate) fn builder(&self) -> &ClientBuilder {
&self.builder
}
}
impl Drop for EmulatorClient {
fn drop(&mut self) {
tokio::spawn(
tokio::process::Command::new("pkill")
.arg("-f")
.arg(format!(
".*{}.*--port={}.*",
self.data.kill_pattern, self.port
))
.output(),
);
}
}
async fn start_emulator(
emulator_name: &str,
project_name: &str,
extra_args: &[OsString],
) -> Result<(tokio::process::Child, String), std::io::Error> {
let mut err: Option<_> = None;
let mut rng = rand::thread_rng();
for _ in 0..CLI_RETRY {
let port = format!("{}", rng.gen_range(PORT_RANGE));
match start_emulator_once(emulator_name, &port, project_name, extra_args) {
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(
emulator_name: &str,
port: &str,
project_name: &str,
extra_args: &[OsString],
) -> Result<tokio::process::Child, std::io::Error> {
tokio::process::Command::new("gcloud")
.arg("beta")
.arg("emulators")
.arg(emulator_name)
.arg("start")
.arg("--project")
.arg(project_name)
.arg("--host-port")
.arg(format!("{}:{}", HOST, port))
.args(extra_args)
.arg("--verbosity")
.arg("debug")
.kill_on_drop(true)
.spawn()
}