use crate::api::{ExecRequest, Gates, MetricsQuery, NewToken};
use crate::error::{Error, Result};
use crate::shell::{ShellEvent, ShellExit, ShellOptions};
use crate::transport::Auth;
use crate::types::*;
use serde_json::Value;
use std::sync::Arc;
use std::time::Duration;
struct Rt(Option<tokio::runtime::Runtime>);
impl Drop for Rt {
fn drop(&mut self) {
if let Some(rt) = self.0.take() {
rt.shutdown_background();
}
}
}
impl std::ops::Deref for Rt {
type Target = tokio::runtime::Runtime;
fn deref(&self) -> &Self::Target {
self.0.as_ref().expect("the runtime outlives every use of it")
}
}
fn block_on<F: std::future::Future>(rt: &Rt, f: F) -> Result<F::Output> {
if tokio::runtime::Handle::try_current().is_ok() {
return Err(Error::Invalid(
"hws::blocking cannot be used inside an async runtime — use \
hws::Client instead"
.into(),
));
}
Ok(rt.block_on(f))
}
pub struct Client {
inner: crate::Client,
rt: Arc<Rt>,
}
impl std::fmt::Debug for Client {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
self.inner.fmt(f)
}
}
macro_rules! run {
($self:expr, $call:expr) => {
block_on(&$self.rt, $call)?
};
}
impl Client {
pub fn new(server: impl Into<String>, auth: Auth) -> Result<Self> {
Self::builder(server).auth(auth).build()
}
pub fn builder(server: impl Into<String>) -> ClientBuilder {
ClientBuilder {
inner: crate::Client::builder(server),
}
}
pub fn into_async(self) -> crate::Client {
self.inner
}
pub fn connect(
server: &str,
user: Option<&str>,
password: Option<&str>,
token: Option<&str>,
insecure: bool,
timeout: Duration,
) -> Result<Self> {
let auth = match (token, user, password) {
(Some(t), _, _) => Auth::Token(t.to_string()),
(None, Some(u), Some(p)) => Auth::Basic {
user: u.to_string(),
password: p.to_string(),
},
_ => Auth::None,
};
Self::builder(server)
.auth(auth)
.insecure(insecure)
.timeout(timeout)
.build()
}
pub fn connect_with_token(
server: &str,
token: &str,
insecure: bool,
timeout: Duration,
) -> Result<Self> {
Self::builder(server)
.token(token)
.insecure(insecure)
.timeout(timeout)
.build()
}
pub fn server(&self) -> &str {
self.inner.server().unwrap_or_default()
}
pub fn has_credentials(&self) -> bool {
self.inner.has_credentials()
}
pub fn status_of(&self, path: &str) -> Result<u16> {
run!(self, self.inner.probe(path))
}
pub fn status_detail_of(&self, path: &str) -> Result<(u16, Option<String>)> {
run!(self, self.inner.probe_detail(path))
}
pub fn healthz(&self) -> Result<()> {
run!(self, self.inner.healthz())
}
pub fn gates(server: &str, insecure: bool, timeout: Duration) -> Result<Gates> {
let probe = Self::builder(server).insecure(insecure).timeout(timeout).build()?;
block_on(&probe.rt, crate::Client::gates(server, insecure))?
}
pub fn gates_of(client: &Self) -> Result<Gates> {
block_on(&client.rt, crate::Client::gates(client.server(), false))?
}
pub fn deployments(&self) -> Result<Vec<DeploymentStatus>> {
run!(self, self.inner.deployments())
}
pub fn deployment(&self, id: &str) -> Result<DeploymentStatus> {
run!(self, self.inner.deployment(id))
}
pub fn deployment_with_timeout(
&self,
id: &str,
timeout: std::time::Duration,
) -> Result<DeploymentStatus> {
run!(self, self.inner.deployment_with_timeout(id, timeout))
}
pub fn deployment_exists(&self, id: &str) -> Result<bool> {
run!(self, self.inner.deployment_exists(id))
}
pub fn create_deployment(&self, spec: &Value) -> Result<DeploymentStatus> {
run!(self, self.inner.create_deployment(spec))
}
pub fn replace_deployment(&self, id: &str, spec: &Value) -> Result<DeploymentStatus> {
run!(self, self.inner.replace_deployment(id, spec))
}
pub fn patch_scaling(&self, id: &str, patch: &Value) -> Result<DeploymentStatus> {
run!(self, self.inner.patch_scaling(id, patch))
}
pub fn delete_deployment(&self, id: &str) -> Result<()> {
run!(self, self.inner.delete_deployment(id))
}
pub fn evict_vm(&self, id: &str, sandbox: &str, force: bool) -> Result<EvictOutcome> {
run!(self, self.inner.evict_vm(id, sandbox, force))
}
pub fn cordon_upstream(
&self,
id: &str,
upstream: &str,
force: bool,
reason: Option<&str>,
) -> Result<UpstreamTrafficStatus> {
run!(
self,
self.inner.cordon_upstream(id, upstream, force, reason)
)
}
pub fn uncordon_upstream(&self, id: &str, upstream: &str) -> Result<UpstreamTrafficStatus> {
run!(self, self.inner.uncordon_upstream(id, upstream))
}
pub fn deployment_ids(&self, page: usize) -> Result<Vec<String>> {
run!(self, self.inner.deployment_ids(page))
}
pub fn exec(&self, id: &str, req: &ExecRequest) -> Result<ExecOutput> {
run!(self, self.inner.exec(id, req))
}
pub fn shell(&self, id: &str, opts: &ShellOptions) -> Result<Shell> {
let inner = run!(self, self.inner.shell(id, opts))?;
Ok(Shell {
inner,
rt: self.rt.clone(),
})
}
pub fn workflows(&self) -> Result<Vec<WorkflowView>> {
run!(self, self.inner.workflows())
}
pub fn workflow(&self, id: &str) -> Result<WorkflowView> {
run!(self, self.inner.workflow(id))
}
pub fn create_workflow(&self, spec: &Value) -> Result<WorkflowView> {
run!(self, self.inner.create_workflow(spec))
}
pub fn replace_workflow(&self, id: &str, spec: &Value) -> Result<WorkflowView> {
run!(self, self.inner.replace_workflow(id, spec))
}
pub fn delete_workflow(&self, id: &str) -> Result<()> {
run!(self, self.inner.delete_workflow(id))
}
pub fn secrets(&self) -> Result<Vec<SecretSummary>> {
run!(self, self.inner.secrets())
}
pub fn secrets_in(&self, namespace: Option<&str>) -> Result<Vec<SecretSummary>> {
run!(self, self.inner.secrets_in(namespace))
}
pub fn secret(&self, id: &str) -> Result<SecretSummary> {
run!(self, self.inner.secret(id))
}
pub fn secret_in(&self, namespace: Option<&str>, id: &str) -> Result<SecretSummary> {
run!(self, self.inner.secret_in(namespace, id))
}
pub fn secret_exists(&self, id: &str) -> Result<bool> {
run!(self, self.inner.secret_exists(id))
}
pub fn secret_exists_in(&self, namespace: Option<&str>, id: &str) -> Result<bool> {
run!(self, self.inner.secret_exists_in(namespace, id))
}
pub fn put_secret(&self, spec: &Value) -> Result<SecretSummary> {
run!(self, self.inner.put_secret(spec))
}
pub fn patch_secret(&self, id: &str, patch: &Value) -> Result<SecretSummary> {
run!(self, self.inner.patch_secret(id, patch))
}
pub fn patch_secret_in(
&self,
namespace: Option<&str>,
id: &str,
patch: &Value,
) -> Result<SecretSummary> {
run!(self, self.inner.patch_secret_in(namespace, id, patch))
}
pub fn delete_secret(&self, id: &str, force: bool) -> Result<()> {
run!(self, self.inner.delete_secret(id, force))
}
pub fn delete_secret_in(&self, namespace: Option<&str>, id: &str, force: bool) -> Result<()> {
run!(self, self.inner.delete_secret_in(namespace, id, force))
}
pub fn auth_providers(&self, namespace: Option<&str>) -> Result<Vec<AuthProviderView>> {
run!(self, self.inner.auth_providers(namespace))
}
pub fn auth_provider(&self, namespace: &str, name: &str) -> Result<AuthProviderView> {
run!(self, self.inner.auth_provider(namespace, name))
}
pub fn auth_provider_exists(&self, namespace: &str, name: &str) -> Result<bool> {
run!(self, self.inner.auth_provider_exists(namespace, name))
}
pub fn create_auth_provider(&self, spec: &Value) -> Result<AuthProviderView> {
run!(self, self.inner.create_auth_provider(spec))
}
pub fn delete_auth_provider(&self, namespace: &str, name: &str) -> Result<()> {
run!(self, self.inner.delete_auth_provider(namespace, name))
}
pub fn mint_token(&self, req: &NewToken) -> Result<MintedToken> {
run!(self, self.inner.mint_token(req))
}
pub fn tokens(&self) -> Result<Vec<TokenSummary>> {
run!(self, self.inner.tokens())
}
pub fn token(&self, id: &str) -> Result<TokenSummary> {
run!(self, self.inner.token(id))
}
pub fn patch_token(&self, id: &str, patch: &Value) -> Result<TokenSummary> {
run!(self, self.inner.patch_token(id, patch))
}
pub fn revoke_token(&self, id: &str) -> Result<()> {
run!(self, self.inner.revoke_token(id))
}
pub fn namespaces(&self) -> Result<Vec<NamespaceEntry>> {
run!(self, self.inner.namespaces())
}
pub fn create_namespace(&self, spec: &Value) -> Result<Value> {
run!(self, self.inner.create_namespace(spec))
}
pub fn delete_namespace(&self, name: &str) -> Result<()> {
run!(self, self.inner.delete_namespace(name))
}
pub fn feeds(&self) -> Result<Vec<FeedIndexEntry>> {
run!(self, self.inner.feeds())
}
pub fn feed_events(&self, namespace: &str) -> Result<Vec<FeedEvent>> {
run!(self, self.inner.feed_events(namespace))
}
pub fn feed_rss(&self, namespace: &str) -> Result<String> {
run!(self, self.inner.feed_rss(namespace))
}
pub fn start_build(&self, id: &str, git_ref: Option<&str>) -> Result<JobRecord> {
run!(self, self.inner.start_build(id, git_ref))
}
pub fn start_pull(&self, id: &str, artifact_ref: Option<&str>, force: bool) -> Result<JobRecord> {
run!(self, self.inner.start_pull(id, artifact_ref, force))
}
pub fn start_mount_pull(&self, id: &str, force: bool) -> Result<JobRecord> {
run!(self, self.inner.start_mount_pull(id, force))
}
pub fn start_update(&self, id: &str) -> Result<JobRecord> {
run!(self, self.inner.start_update(id))
}
pub fn jobs(&self) -> Result<Vec<JobRecord>> {
run!(self, self.inner.jobs())
}
pub fn deployment_jobs(&self, id: &str) -> Result<Vec<JobRecord>> {
run!(self, self.inner.deployment_jobs(id))
}
pub fn job(&self, job_id: &str) -> Result<JobRecord> {
run!(self, self.inner.job(job_id))
}
pub fn wait_for_job(
&self,
job_id: &str,
timeout: Duration,
on_progress: impl FnMut(crate::wait::JobProgress<'_>) + Send,
) -> Result<JobRecord> {
run!(
self,
self.inner
.wait_for_job(job_id)
.timeout(timeout)
.on_progress(on_progress)
.await_done()
)
}
pub fn wait_for_ready(
&self,
id: &str,
timeout: Duration,
on_progress: impl FnMut(crate::wait::PoolProgress) + Send,
) -> Result<DeploymentStatus> {
run!(
self,
self.inner
.wait_for_ready(id)
.timeout(timeout)
.on_progress(on_progress)
.await_ready()
)
}
pub fn metrics(&self, query: &MetricsQuery) -> Result<MetricsResponse> {
run!(self, self.inner.metrics(query))
}
pub fn certs(&self) -> Result<Vec<CertStatus>> {
run!(self, self.inner.certs())
}
pub fn disks(&self) -> Result<DiskInventory> {
run!(self, self.inner.disks())
}
pub fn plugins(&self) -> Result<Vec<PluginView>> {
run!(self, self.inner.plugins())
}
pub fn plugin(&self, id: &str) -> Result<PluginView> {
run!(self, self.inner.plugin(id))
}
pub fn set_plugin(&self, id: &str, enabled: bool, config: Option<&Value>) -> Result<PluginView> {
run!(self, self.inner.set_plugin(id, enabled, config))
}
}
pub struct ClientBuilder {
inner: crate::ClientBuilder,
}
impl ClientBuilder {
pub fn auth(mut self, auth: Auth) -> Self {
self.inner = self.inner.auth(auth);
self
}
pub fn token(mut self, token: impl Into<String>) -> Self {
self.inner = self.inner.token(token);
self
}
pub fn basic(mut self, user: impl Into<String>, password: impl Into<String>) -> Self {
self.inner = self.inner.basic(user, password);
self
}
pub fn timeout(mut self, d: Duration) -> Self {
self.inner = self.inner.timeout(d);
self
}
pub fn insecure(mut self, yes: bool) -> Self {
self.inner = self.inner.insecure(yes);
self
}
pub fn build(self) -> Result<Client> {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| Error::Invalid(format!("could not start a runtime: {e}")))?;
Ok(Client {
inner: self.inner.build()?,
rt: Arc::new(Rt(Some(rt))),
})
}
}
pub struct Shell {
inner: crate::Shell,
rt: Arc<Rt>,
}
impl Shell {
pub fn sandbox_id(&self) -> &str {
self.inner.sandbox_id()
}
pub fn write(&mut self, bytes: &[u8]) -> Result<()> {
block_on(&self.rt, self.inner.write(bytes))?
}
pub fn resize(&mut self, cols: u16, rows: u16) -> Result<()> {
block_on(&self.rt, self.inner.resize(cols, rows))?
}
pub fn next_timeout(&mut self, d: Duration) -> Option<ShellEvent> {
block_on(&self.rt, async {
tokio::time::timeout(d, self.inner.next()).await.ok().flatten()
})
.ok()
.flatten()
}
pub fn exit(&self) -> Option<&ShellExit> {
self.inner.exit()
}
pub fn close(&mut self) -> Result<()> {
block_on(&self.rt, self.inner.close())?
}
}
impl Iterator for Shell {
type Item = ShellEvent;
fn next(&mut self) -> Option<ShellEvent> {
block_on(&self.rt, self.inner.next()).ok().flatten()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_client_can_be_built_and_answers_without_a_runtime() {
let c = Client::builder("127.0.0.1:1").timeout(Duration::from_millis(50)).build().unwrap();
let e = c.healthz().unwrap_err();
assert!(matches!(e, Error::Transport(_)), "{e:?}");
}
#[tokio::test]
async fn calling_it_from_inside_a_runtime_is_refused_not_a_panic() {
let c = tokio::task::spawn_blocking(|| {
Client::builder("127.0.0.1:1")
.timeout(Duration::from_millis(50))
.build()
.unwrap()
})
.await
.unwrap();
let e = c.healthz().unwrap_err();
assert!(matches!(&e, Error::Invalid(m) if m.contains("async runtime")), "{e:?}");
assert!(e.to_string().contains("hws::Client"), "it should name the fix: {e}");
}
}
pub struct Raw<'a> {
client: &'a Client,
}
impl Client {
pub fn raw(&self) -> Raw<'_> {
Raw { client: self }
}
}
macro_rules! raw_blocking {
($($name:ident),* $(,)?) => {
$(
pub fn $name(&self) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().$name())?
}
)*
};
}
macro_rules! raw_blocking_id {
($($name:ident),* $(,)?) => {
$(
pub fn $name(&self, id: &str) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().$name(id))?
}
)*
};
}
impl Raw<'_> {
raw_blocking!(deployments, secrets, tokens, jobs, certs, workflows, feeds, disks, namespaces, plugins);
pub fn deployments_in(&self, namespace: &str) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().deployments_in(namespace))?
}
raw_blocking_id!(deployment, secret, token, job, deployment_jobs, spec, workflow);
pub fn metrics(&self, query: &MetricsQuery) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().metrics(query))?
}
pub fn create_deployment(&self, spec: &Value) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().create_deployment(spec))?
}
pub fn replace_deployment(&self, id: &str, spec: &Value) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().replace_deployment(id, spec))?
}
pub fn patch_scaling(&self, id: &str, patch: &Value) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().patch_scaling(id, patch))?
}
pub fn put_secret(&self, spec: &Value) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().put_secret(spec))?
}
pub fn patch_secret(&self, id: &str, patch: &Value) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().patch_secret(id, patch))?
}
pub fn patch_secret_in(
&self,
namespace: Option<&str>,
id: &str,
patch: &Value,
) -> Result<Value> {
block_on(
&self.client.rt,
self.client.inner.raw().patch_secret_in(namespace, id, patch),
)?
}
pub fn feed_events(&self, namespace: &str) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().feed_events(namespace))?
}
pub fn secrets_in(&self, namespace: Option<&str>) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().secrets_in(namespace))?
}
pub fn secret_in(&self, namespace: Option<&str>, id: &str) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().secret_in(namespace, id))?
}
pub fn auth_providers(&self, namespace: Option<&str>) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().auth_providers(namespace))?
}
pub fn auth_provider(&self, namespace: &str, name: &str) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().auth_provider(namespace, name))?
}
pub fn mint_token(&self, req: &NewToken) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().mint_token(req))?
}
pub fn patch_token(&self, id: &str, patch: &Value) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().patch_token(id, patch))?
}
pub fn evict_vm(&self, id: &str, sandbox: &str, force: bool) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().evict_vm(id, sandbox, force))?
}
pub fn cordon_upstream(
&self,
id: &str,
upstream: &str,
force: bool,
reason: Option<&str>,
) -> Result<Value> {
block_on(
&self.client.rt,
self.client
.inner
.raw()
.cordon_upstream(id, upstream, force, reason),
)?
}
pub fn uncordon_upstream(&self, id: &str, upstream: &str) -> Result<Value> {
block_on(
&self.client.rt,
self.client.inner.raw().uncordon_upstream(id, upstream),
)?
}
pub fn start_build(&self, id: &str, git_ref: Option<&str>) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().start_build(id, git_ref))?
}
pub fn start_pull(&self, id: &str, r: Option<&str>, force: bool) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().start_pull(id, r, force))?
}
pub fn start_mount_pull(&self, id: &str, force: bool) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().start_mount_pull(id, force))?
}
pub fn start_update(&self, id: &str) -> Result<Value> {
block_on(&self.client.rt, self.client.inner.raw().start_update(id))?
}
}