use crate::error::{Error, Result};
use crate::transport::{Auth, HttpTransport, Method, Request, Response, Transport};
use crate::types::*;
use serde::de::DeserializeOwned;
use serde_json::{Value, json};
use std::sync::Arc;
use std::time::Duration;
pub const DEFAULT_TIMEOUT: Duration = Duration::from_secs(30);
pub const ASSUMED_COLD_START_SECS: u64 = 120;
const EXEC_MARGIN_SECS: u64 = 15;
fn secret_path(namespace: Option<&str>, id: &str, force: Option<bool>) -> String {
let mut path = format!("/secrets/{}", seg(id));
let mut query: Vec<String> = Vec::new();
if let Some(ns) = namespace {
query.push(format!("namespace={}", seg(ns)));
}
if let Some(force) = force {
query.push(format!("force={}", if force { "true" } else { "false" }));
}
if !query.is_empty() {
path.push('?');
path.push_str(&query.join("&"));
}
path
}
fn seg(s: &str) -> String {
let mut out = String::with_capacity(s.len());
for b in s.bytes() {
match b {
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => {
out.push(b as char)
}
_ => out.push_str(&format!("%{b:02X}")),
}
}
out
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Gates {
pub view: bool,
pub crud: bool,
}
impl Gates {
pub fn any(&self) -> bool {
self.view || self.crud
}
}
#[derive(Debug, Clone)]
pub struct Client {
transport: Arc<dyn Transport>,
ws: Option<Arc<WsConfig>>,
}
#[derive(Debug)]
pub(crate) struct WsConfig {
pub base: String,
pub auth: Auth,
pub insecure: bool,
}
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 {
server: server.into(),
auth: Auth::None,
timeout: DEFAULT_TIMEOUT,
insecure: false,
}
}
pub fn with_transport(transport: Arc<dyn Transport>) -> Self {
Self {
transport,
ws: None,
}
}
async fn send(&self, req: Request, kind: &'static str, name: &str) -> Result<Response> {
let r = self.transport.send(req).await?;
if r.is_success() {
Ok(r)
} else {
Err(Error::from_response(
r.status,
&r.body,
kind,
name,
self.transport.credential(),
))
}
}
async fn read<T: DeserializeOwned>(
&self,
req: Request,
kind: &'static str,
name: &str,
) -> Result<T> {
let r = self.send(req, kind, name).await?;
serde_json::from_str(&r.body).map_err(Error::Decode)
}
async fn unit(&self, req: Request, kind: &'static str, name: &str) -> Result<()> {
self.send(req, kind, name).await.map(|_| ())
}
pub async fn healthz(&self) -> Result<()> {
self.unit(Request::new(Method::Get, "/healthz"), "server", "")
.await
}
pub async fn gates(server: &str, insecure: bool) -> Result<Gates> {
let anon = Client::builder(server).insecure(insecure).build()?;
let refused = |e: &Error| matches!(e, Error::Unauthorized { .. });
let view = anon
.read::<Value>(Request::new(Method::Get, "/metrics"), "server", "")
.await
.err()
.as_ref()
.is_some_and(refused);
let crud = anon
.read::<Value>(Request::new(Method::Get, "/deployments"), "server", "")
.await
.err()
.as_ref()
.is_some_and(refused);
Ok(Gates { view, crud })
}
pub async fn deployments(&self) -> Result<Vec<DeploymentStatus>> {
self.read(Request::new(Method::Get, "/deployments"), "deployment", "")
.await
}
pub async fn deployment(&self, id: &str) -> Result<DeploymentStatus> {
self.read(
Request::new(Method::Get, format!("/deployments/{}", seg(id))),
"deployment",
id,
)
.await
}
pub async fn deployment_with_timeout(
&self,
id: &str,
timeout: Duration,
) -> Result<DeploymentStatus> {
self.read(
Request::new(Method::Get, format!("/deployments/{}", seg(id))).timeout(timeout),
"deployment",
id,
)
.await
}
pub async fn deployment_exists(&self, id: &str) -> Result<bool> {
match self.deployment(id).await {
Ok(_) => Ok(true),
Err(Error::NotFound { .. }) => Ok(false),
Err(e) => Err(e),
}
}
pub async fn create_deployment(&self, spec: &Value) -> Result<DeploymentStatus> {
let id = spec.get("id").and_then(Value::as_str).unwrap_or_default().to_string();
self.read(
Request::new(Method::Post, "/deployments").json(spec.clone()),
"deployment",
&id,
)
.await
}
pub async fn replace_deployment(&self, id: &str, spec: &Value) -> Result<DeploymentStatus> {
self.read(
Request::new(Method::Put, format!("/deployments/{}", seg(id))).json(spec.clone()),
"deployment",
id,
)
.await
}
pub async fn patch_scaling(&self, id: &str, patch: &Value) -> Result<DeploymentStatus> {
self.read(
Request::new(
Method::Patch,
format!("/deployments/{}/scaling", seg(id)),
)
.json(patch.clone()),
"deployment",
id,
)
.await
}
pub async fn delete_deployment(&self, id: &str) -> Result<()> {
self.unit(
Request::new(Method::Delete, format!("/deployments/{}", seg(id))),
"deployment",
id,
)
.await
}
pub async fn evict_vm(&self, id: &str, sandbox: &str, force: bool) -> Result<EvictOutcome> {
let path = format!(
"/deployments/{}/vms/{}?force={}",
seg(id),
seg(sandbox),
if force { "true" } else { "false" }
);
self.read(Request::new(Method::Delete, path), "vm", sandbox)
.await
}
pub async fn cordon_upstream(
&self,
id: &str,
upstream: &str,
force: bool,
reason: Option<&str>,
) -> Result<UpstreamTrafficStatus> {
let mut body = json!({ "force": force });
if let Some(reason) = reason {
body["reason"] = json!(reason);
}
self.read(
Request::new(
Method::Put,
format!(
"/deployments/{}/upstreams/{}/drain",
seg(id),
seg(upstream),
),
)
.json(body),
"upstream",
upstream,
)
.await
}
pub async fn uncordon_upstream(
&self,
id: &str,
upstream: &str,
) -> Result<UpstreamTrafficStatus> {
self.read(
Request::new(
Method::Delete,
format!(
"/deployments/{}/upstreams/{}/drain",
seg(id),
seg(upstream),
),
),
"upstream",
upstream,
)
.await
}
pub async fn exec(&self, id: &str, req: &ExecRequest) -> Result<ExecOutput> {
if req.command.trim().is_empty() {
return Err(Error::Invalid("a command to exec must not be blank".into()));
}
let mut body = json!({ "command": req.command, "wake": req.wake });
if let Some(cwd) = &req.cwd {
body["cwd"] = json!(cwd);
}
if let Some(env) = &req.env {
body["env"] = json!(env);
}
if let Some(t) = req.timeout_secs {
body["timeout_secs"] = json!(t);
}
self.read(
Request::new(Method::Post, format!("/deployments/{}/exec", seg(id)))
.json(body)
.timeout(req.patience()),
"deployment",
id,
)
.await
}
pub async fn workflows(&self) -> Result<Vec<WorkflowView>> {
let list: WorkflowList = self
.read(Request::new(Method::Get, "/workflows"), "workflow", "")
.await?;
Ok(list.workflows)
}
pub async fn workflow(&self, id: &str) -> Result<WorkflowView> {
self.read(
Request::new(Method::Get, format!("/workflows/{}", seg(id))),
"workflow",
id,
)
.await
}
pub async fn create_workflow(&self, spec: &Value) -> Result<WorkflowView> {
let id = spec
.get("id")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string();
self.read(
Request::new(Method::Post, "/workflows").json(spec.clone()),
"workflow",
&id,
)
.await
}
pub async fn replace_workflow(&self, id: &str, spec: &Value) -> Result<WorkflowView> {
self.read(
Request::new(Method::Put, format!("/workflows/{}", seg(id))).json(spec.clone()),
"workflow",
id,
)
.await
}
pub async fn delete_workflow(&self, id: &str) -> Result<()> {
self.unit(
Request::new(Method::Delete, format!("/workflows/{}", seg(id))),
"workflow",
id,
)
.await
}
pub async fn secrets(&self) -> Result<Vec<SecretSummary>> {
self.secrets_in(None).await
}
pub async fn secrets_in(&self, namespace: Option<&str>) -> Result<Vec<SecretSummary>> {
let path = match namespace {
Some(ns) => format!("/secrets?namespace={}", seg(ns)),
None => "/secrets".to_string(),
};
self.read(Request::new(Method::Get, path), "secret", "").await
}
pub async fn secret(&self, id: &str) -> Result<SecretSummary> {
self.secret_in(None, id).await
}
pub async fn secret_in(&self, namespace: Option<&str>, id: &str) -> Result<SecretSummary> {
self.read(
Request::new(Method::Get, secret_path(namespace, id, None)),
"secret",
id,
)
.await
}
pub async fn secret_exists(&self, id: &str) -> Result<bool> {
self.secret_exists_in(None, id).await
}
pub async fn secret_exists_in(&self, namespace: Option<&str>, id: &str) -> Result<bool> {
match self.secret_in(namespace, id).await {
Ok(_) => Ok(true),
Err(Error::NotFound { .. }) => Ok(false),
Err(e) => Err(e),
}
}
pub async fn put_secret(&self, spec: &Value) -> Result<SecretSummary> {
let id = spec.get("id").and_then(Value::as_str).unwrap_or_default().to_string();
self.read(
Request::new(Method::Post, "/secrets").json(spec.clone()),
"secret",
&id,
)
.await
}
pub async fn patch_secret(&self, id: &str, patch: &Value) -> Result<SecretSummary> {
self.patch_secret_in(None, id, patch).await
}
pub async fn patch_secret_in(
&self,
namespace: Option<&str>,
id: &str,
patch: &Value,
) -> Result<SecretSummary> {
self.read(
Request::new(Method::Patch, secret_path(namespace, id, None)).json(patch.clone()),
"secret",
id,
)
.await
}
pub async fn delete_secret(&self, id: &str, force: bool) -> Result<()> {
self.delete_secret_in(None, id, force).await
}
pub async fn delete_secret_in(
&self,
namespace: Option<&str>,
id: &str,
force: bool,
) -> Result<()> {
self.unit(
Request::new(Method::Delete, secret_path(namespace, id, Some(force))),
"secret",
id,
)
.await
}
pub async fn auth_providers(&self, namespace: Option<&str>) -> Result<Vec<AuthProviderView>> {
let path = match namespace {
Some(ns) => format!("/auth-providers?namespace={}", seg(ns)),
None => "/auth-providers".to_string(),
};
self.read(Request::new(Method::Get, path), "auth provider", "")
.await
}
pub async fn auth_provider(&self, namespace: &str, name: &str) -> Result<AuthProviderView> {
self.read(
Request::new(
Method::Get,
format!("/auth-providers/{}/{}", seg(namespace), seg(name)),
),
"auth provider",
name,
)
.await
}
pub async fn auth_provider_exists(&self, namespace: &str, name: &str) -> Result<bool> {
match self.auth_provider(namespace, name).await {
Ok(_) => Ok(true),
Err(Error::NotFound { .. }) => Ok(false),
Err(e) => Err(e),
}
}
pub async fn create_auth_provider(&self, spec: &Value) -> Result<AuthProviderView> {
let name = spec.get("name").and_then(Value::as_str).unwrap_or_default().to_string();
self.read(
Request::new(Method::Post, "/auth-providers").json(spec.clone()),
"auth provider",
&name,
)
.await
}
pub async fn delete_auth_provider(&self, namespace: &str, name: &str) -> Result<()> {
self.unit(
Request::new(
Method::Delete,
format!("/auth-providers/{}/{}", seg(namespace), seg(name)),
),
"auth provider",
name,
)
.await
}
pub async fn mint_token(&self, req: &NewToken) -> Result<MintedToken> {
if req.name.trim().is_empty() {
return Err(Error::Invalid(
"a token needs a name — it is how you know what to revoke".into(),
));
}
let mut body = json!({
"name": req.name,
"admin": req.admin,
"deployments": req.deployments,
});
if let Some(ns) = &req.namespace {
body["namespace"] = json!(ns);
}
if let Some(s) = req.expires_in_secs {
body["expires_in_secs"] = json!(s);
}
self.read(
Request::new(Method::Post, "/tokens").json(body),
"token",
&req.name,
)
.await
}
pub async fn tokens(&self) -> Result<Vec<TokenSummary>> {
self.read(Request::new(Method::Get, "/tokens"), "token", "")
.await
}
pub async fn namespaces(&self) -> Result<Vec<NamespaceEntry>> {
self.read(Request::new(Method::Get, "/namespaces"), "namespace", "")
.await
}
pub async fn create_namespace(&self, spec: &Value) -> Result<Value> {
let name = spec.get("name").and_then(Value::as_str).unwrap_or_default().to_string();
self.read(
Request::new(Method::Post, "/namespaces").json(spec.clone()),
"namespace",
&name,
)
.await
}
pub async fn delete_namespace(&self, name: &str) -> Result<()> {
self.unit(
Request::new(Method::Delete, format!("/namespaces/{}", seg(name))),
"namespace",
name,
)
.await
}
pub async fn feeds(&self) -> Result<Vec<FeedIndexEntry>> {
self.read(Request::new(Method::Get, "/feeds"), "feed", "")
.await
}
pub async fn feed_events(&self, namespace: &str) -> Result<Vec<FeedEvent>> {
self.read(
Request::new(Method::Get, format!("/feeds/{}?format=json", seg(namespace))),
"feed",
namespace,
)
.await
}
pub async fn feed_rss(&self, namespace: &str) -> Result<String> {
self.send(
Request::new(Method::Get, format!("/feeds/{}", seg(namespace))),
"feed",
namespace,
)
.await
.map(|r| r.body)
}
pub async fn token(&self, id: &str) -> Result<TokenSummary> {
self.read(
Request::new(Method::Get, format!("/tokens/{}", seg(id))),
"token",
id,
)
.await
}
pub async fn patch_token(&self, id: &str, patch: &Value) -> Result<TokenSummary> {
self.read(
Request::new(Method::Patch, format!("/tokens/{}", seg(id))).json(patch.clone()),
"token",
id,
)
.await
}
pub async fn revoke_token(&self, id: &str) -> Result<()> {
self.unit(
Request::new(Method::Delete, format!("/tokens/{}", seg(id))),
"token",
id,
)
.await
}
pub async fn start_build(&self, id: &str, git_ref: Option<&str>) -> Result<JobRecord> {
let body = git_ref.map_or_else(|| json!({}), |r| json!({ "ref": r }));
self.read(
Request::new(Method::Post, format!("/deployments/{}/build", seg(id))).json(body),
"deployment",
id,
)
.await
}
pub async fn start_pull(&self, id: &str, artifact_ref: Option<&str>, force: bool) -> Result<JobRecord> {
let mut body = json!({ "force": force });
if let Some(r) = artifact_ref {
body["ref"] = json!(r);
}
self.read(
Request::new(Method::Post, format!("/deployments/{}/pull", seg(id))).json(body),
"deployment",
id,
)
.await
}
pub async fn start_mount_pull(&self, id: &str, force: bool) -> Result<JobRecord> {
self.read(
Request::new(Method::Post, format!("/deployments/{}/mounts/pull", seg(id)))
.json(json!({ "force": force })),
"deployment",
id,
)
.await
}
pub async fn start_update(&self, id: &str) -> Result<JobRecord> {
self.read(
Request::new(Method::Post, format!("/deployments/{}/update", seg(id))).json(json!({})),
"deployment",
id,
)
.await
}
pub async fn jobs(&self) -> Result<Vec<JobRecord>> {
self.read(Request::new(Method::Get, "/jobs"), "job", "")
.await
}
pub async fn deployment_jobs(&self, id: &str) -> Result<Vec<JobRecord>> {
self.read(
Request::new(Method::Get, format!("/deployments/{}/jobs", seg(id))),
"deployment",
id,
)
.await
}
pub async fn job(&self, job_id: &str) -> Result<JobRecord> {
self.read(
Request::new(Method::Get, format!("/jobs/{}", seg(job_id))),
"job",
job_id,
)
.await
}
pub async fn metrics(&self, query: &MetricsQuery) -> Result<MetricsResponse> {
self.read(
Request::new(Method::Get, format!("/metrics{}", query.to_query_string())),
"server",
"",
)
.await
}
pub async fn certs(&self) -> Result<Vec<CertStatus>> {
self.read(Request::new(Method::Get, "/certs"), "certificate", "")
.await
}
pub async fn disks(&self) -> Result<DiskInventory> {
self.read(Request::new(Method::Get, "/disks"), "disk", "")
.await
}
pub async fn plugins(&self) -> Result<Vec<PluginView>> {
self.read(Request::new(Method::Get, "/api/plugins"), "plugin", "").await
}
pub async fn plugin(&self, id: &str) -> Result<PluginView> {
self.read(Request::new(Method::Get, format!("/api/plugins/{}", seg(id))), "plugin", id)
.await
}
pub async fn set_plugin(&self, id: &str, enabled: bool, config: Option<&Value>) -> Result<PluginView> {
let mut body = json!({ "enabled": enabled });
if let Some(c) = config {
body["config"] = c.clone();
}
self.read(
Request::new(Method::Put, format!("/api/plugins/{}", seg(id))).json(body),
"plugin",
id,
)
.await
}
pub(crate) fn ws(&self) -> Option<&Arc<WsConfig>> {
self.ws.as_ref()
}
pub async fn probe(&self, path: &str) -> Result<u16> {
Ok(self.probe_detail(path).await?.0)
}
pub async fn probe_detail(&self, path: &str) -> Result<(u16, Option<String>)> {
let r = self
.transport
.send(Request::new(Method::Get, path.to_string()))
.await?;
let detail = serde_json::from_str::<serde_json::Value>(&r.body)
.ok()
.and_then(|v| {
let field = |k: &str| {
v.get(k)
.and_then(|x| x.as_str())
.map(str::trim)
.filter(|x| !x.is_empty())
.map(str::to_owned)
};
match (field("error"), field("detail")) {
(Some(e), Some(d)) => Some(format!("{e} — {d}")),
(Some(e), None) => Some(e),
(None, Some(d)) => Some(d),
(None, None) => None,
}
});
Ok((r.status, detail))
}
pub fn server(&self) -> Option<&str> {
self.ws.as_ref().map(|w| w.base.as_str())
}
pub fn has_credentials(&self) -> bool {
self.transport.credential() != crate::error::Credential::None
}
pub fn raw(&self) -> Raw<'_> {
Raw(self)
}
}
#[derive(Debug, Clone, Copy)]
pub struct Raw<'a>(&'a Client);
macro_rules! raw_list {
($($name:ident => $kind:literal, $path:literal;)*) => {
$(
pub async fn $name(&self) -> Result<Value> {
self.0.read(Request::new(Method::Get, $path), $kind, "").await
}
)*
};
}
macro_rules! raw_item {
($($name:ident => $kind:literal, $fmt:literal;)*) => {
$(
pub async fn $name(&self, id: &str) -> Result<Value> {
self.0
.read(Request::new(Method::Get, format!($fmt, seg(id))), $kind, id)
.await
}
)*
};
}
impl Raw<'_> {
raw_list! {
deployments => "deployment", "/deployments";
secrets => "secret", "/secrets";
tokens => "token", "/tokens";
jobs => "job", "/jobs";
certs => "certificate", "/certs";
workflows => "workflow", "/workflows";
feeds => "feed", "/feeds";
namespaces => "namespace", "/namespaces";
disks => "disk", "/disks";
plugins => "plugin", "/api/plugins";
}
pub async fn deployments_in(&self, namespace: &str) -> Result<Value> {
self.0
.read(
Request::new(Method::Get, format!("/deployments?namespace={}", seg(namespace))),
"deployment",
"",
)
.await
}
pub async fn secrets_in(&self, namespace: Option<&str>) -> Result<Value> {
let path = match namespace {
Some(ns) => format!("/secrets?namespace={}", seg(ns)),
None => "/secrets".to_string(),
};
self.0.read(Request::new(Method::Get, path), "secret", "").await
}
pub async fn secret_in(&self, namespace: Option<&str>, id: &str) -> Result<Value> {
self.0
.read(
Request::new(Method::Get, secret_path(namespace, id, None)),
"secret",
id,
)
.await
}
pub async fn auth_providers(&self, namespace: Option<&str>) -> Result<Value> {
let path = match namespace {
Some(ns) => format!("/auth-providers?namespace={}", seg(ns)),
None => "/auth-providers".to_string(),
};
self.0
.read(Request::new(Method::Get, path), "auth provider", "")
.await
}
pub async fn auth_provider(&self, namespace: &str, name: &str) -> Result<Value> {
self.0
.read(
Request::new(
Method::Get,
format!("/auth-providers/{}/{}", seg(namespace), seg(name)),
),
"auth provider",
name,
)
.await
}
pub async fn feed_events(&self, namespace: &str) -> Result<Value> {
self.0
.read(
Request::new(Method::Get, format!("/feeds/{}?format=json", seg(namespace))),
"feed",
namespace,
)
.await
}
raw_item! {
deployment => "deployment", "/deployments/{}";
secret => "secret", "/secrets/{}";
token => "token", "/tokens/{}";
job => "job", "/jobs/{}";
deployment_jobs => "deployment", "/deployments/{}/jobs";
workflow => "workflow", "/workflows/{}";
}
pub async fn metrics(&self, query: &MetricsQuery) -> Result<Value> {
self.0
.read(
Request::new(Method::Get, format!("/metrics{}", query.to_query_string())),
"server",
"",
)
.await
}
pub async fn create_deployment(&self, spec: &Value) -> Result<Value> {
let id = spec.get("id").and_then(Value::as_str).unwrap_or_default().to_string();
self.0
.read(
Request::new(Method::Post, "/deployments").json(spec.clone()),
"deployment",
&id,
)
.await
}
pub async fn replace_deployment(&self, id: &str, spec: &Value) -> Result<Value> {
self.0
.read(
Request::new(Method::Put, format!("/deployments/{}", seg(id))).json(spec.clone()),
"deployment",
id,
)
.await
}
pub async fn patch_scaling(&self, id: &str, patch: &Value) -> Result<Value> {
self.0
.read(
Request::new(Method::Patch, format!("/deployments/{}/scaling", seg(id)))
.json(patch.clone()),
"deployment",
id,
)
.await
}
pub async fn put_secret(&self, spec: &Value) -> Result<Value> {
let id = spec.get("id").and_then(Value::as_str).unwrap_or_default().to_string();
self.0
.read(
Request::new(Method::Post, "/secrets").json(spec.clone()),
"secret",
&id,
)
.await
}
pub async fn patch_secret_in(
&self,
namespace: Option<&str>,
id: &str,
patch: &Value,
) -> Result<Value> {
self.0
.read(
Request::new(Method::Patch, secret_path(namespace, id, None)).json(patch.clone()),
"secret",
id,
)
.await
}
pub async fn patch_secret(&self, id: &str, patch: &Value) -> Result<Value> {
self.0
.read(
Request::new(Method::Patch, format!("/secrets/{}", seg(id))).json(patch.clone()),
"secret",
id,
)
.await
}
pub async fn mint_token(&self, req: &NewToken) -> Result<Value> {
let body = json!({
"name": req.name,
"admin": req.admin,
"namespace": req.namespace,
"deployments": req.deployments,
"expires_in_secs": req.expires_in_secs,
});
self.0
.read(
Request::new(Method::Post, "/tokens").json(body),
"token",
&req.name,
)
.await
}
pub async fn patch_token(&self, id: &str, patch: &Value) -> Result<Value> {
self.0
.read(
Request::new(Method::Patch, format!("/tokens/{}", seg(id))).json(patch.clone()),
"token",
id,
)
.await
}
pub async fn evict_vm(&self, id: &str, sandbox: &str, force: bool) -> Result<Value> {
self.0
.read(
Request::new(
Method::Delete,
format!(
"/deployments/{}/vms/{}?force={}",
seg(id),
seg(sandbox),
if force { "true" } else { "false" }
),
),
"vm",
sandbox,
)
.await
}
pub async fn cordon_upstream(
&self,
id: &str,
upstream: &str,
force: bool,
reason: Option<&str>,
) -> Result<Value> {
let mut body = json!({ "force": force });
if let Some(reason) = reason {
body["reason"] = json!(reason);
}
self.0
.read(
Request::new(
Method::Put,
format!(
"/deployments/{}/upstreams/{}/drain",
seg(id),
seg(upstream),
),
)
.json(body),
"upstream",
upstream,
)
.await
}
pub async fn uncordon_upstream(&self, id: &str, upstream: &str) -> Result<Value> {
self.0
.read(
Request::new(
Method::Delete,
format!(
"/deployments/{}/upstreams/{}/drain",
seg(id),
seg(upstream),
),
),
"upstream",
upstream,
)
.await
}
pub async fn start_build(&self, id: &str, git_ref: Option<&str>) -> Result<Value> {
let body = git_ref.map_or_else(|| json!({}), |r| json!({ "ref": r }));
self.0
.read(
Request::new(Method::Post, format!("/deployments/{}/build", seg(id))).json(body),
"deployment",
id,
)
.await
}
pub async fn start_pull(&self, id: &str, artifact_ref: Option<&str>, force: bool) -> Result<Value> {
let mut body = json!({ "force": force });
if let Some(r) = artifact_ref {
body["ref"] = json!(r);
}
self.0
.read(
Request::new(Method::Post, format!("/deployments/{}/pull", seg(id))).json(body),
"deployment",
id,
)
.await
}
pub async fn start_mount_pull(&self, id: &str, force: bool) -> Result<Value> {
self.0
.read(
Request::new(Method::Post, format!("/deployments/{}/mounts/pull", seg(id)))
.json(json!({ "force": force })),
"deployment",
id,
)
.await
}
pub async fn start_update(&self, id: &str) -> Result<Value> {
self.0
.read(
Request::new(Method::Post, format!("/deployments/{}/update", seg(id)))
.json(json!({})),
"deployment",
id,
)
.await
}
pub async fn spec(&self, id: &str) -> Result<Value> {
let status = self.deployment(id).await?;
status
.get("spec")
.cloned()
.ok_or_else(|| Error::Decode(serde::de::Error::custom("the response had no `spec`")))
}
}
pub struct ClientBuilder {
server: String,
auth: Auth,
timeout: Duration,
insecure: bool,
}
impl ClientBuilder {
pub fn auth(mut self, auth: Auth) -> Self {
self.auth = auth;
self
}
pub fn token(self, token: impl Into<String>) -> Self {
self.auth(Auth::Token(token.into()))
}
pub fn basic(self, user: impl Into<String>, password: impl Into<String>) -> Self {
self.auth(Auth::Basic {
user: user.into(),
password: password.into(),
})
}
pub fn timeout(mut self, d: Duration) -> Self {
self.timeout = d;
self
}
pub fn insecure(mut self, yes: bool) -> Self {
self.insecure = yes;
self
}
pub fn build(self) -> Result<Client> {
let base = crate::transport::normalize_base(&self.server);
let transport = HttpTransport::new(
base.clone(),
self.auth.clone(),
self.timeout,
self.insecure,
)?;
Ok(Client {
transport: Arc::new(transport),
ws: Some(Arc::new(WsConfig {
base,
auth: self.auth,
insecure: self.insecure,
})),
})
}
}
#[derive(Debug, Clone)]
pub struct ExecRequest {
pub command: String,
pub cwd: Option<String>,
pub env: Option<std::collections::BTreeMap<String, String>>,
pub timeout_secs: Option<u64>,
pub wake: bool,
pub patience: Option<Duration>,
}
impl ExecRequest {
pub fn new(command: impl Into<String>) -> Self {
Self {
command: command.into(),
cwd: None,
env: None,
timeout_secs: None,
wake: true,
patience: None,
}
}
pub fn cwd(mut self, cwd: impl Into<String>) -> Self {
self.cwd = Some(cwd.into());
self
}
pub fn env(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
self.env
.get_or_insert_with(Default::default)
.insert(key.into(), value.into());
self
}
pub fn timeout_secs(mut self, s: u64) -> Self {
self.timeout_secs = Some(s);
self
}
pub fn no_wake(mut self) -> Self {
self.wake = false;
self
}
pub fn patient_for(mut self, d: Duration) -> Self {
self.patience = Some(d);
self
}
pub fn patience(&self) -> Duration {
if let Some(d) = self.patience {
return d;
}
let command = self.timeout_secs.unwrap_or(60).clamp(1, 3600);
let cold = if self.wake { ASSUMED_COLD_START_SECS } else { 0 };
Duration::from_secs(command + cold + EXEC_MARGIN_SECS)
}
}
#[derive(Debug, Clone, Default)]
pub struct MetricsQuery {
pub deployment: Option<String>,
pub prefix: Option<String>,
pub summary: bool,
pub limit: Option<usize>,
pub offset: usize,
}
impl MetricsQuery {
pub fn new() -> Self {
Self::default()
}
pub fn deployment(mut self, id: impl Into<String>) -> Self {
self.deployment = Some(id.into());
self
}
pub fn prefix(mut self, p: impl Into<String>) -> Self {
self.prefix = Some(p.into());
self
}
pub fn summary(mut self, yes: bool) -> Self {
self.summary = yes;
self
}
pub fn page(mut self, offset: usize, limit: usize) -> Self {
self.offset = offset;
self.limit = Some(limit);
self
}
fn to_query_string(&self) -> String {
let mut parts: Vec<String> = Vec::new();
if let Some(d) = &self.deployment {
parts.push(format!("deployment={}", seg(d)));
}
if let Some(p) = &self.prefix {
parts.push(format!("prefix={}", seg(p)));
}
if self.summary {
parts.push("summary=true".into());
}
if let Some(l) = self.limit {
parts.push(format!("limit={l}"));
}
if self.offset > 0 {
parts.push(format!("offset={}", self.offset));
}
if parts.is_empty() {
String::new()
} else {
format!("?{}", parts.join("&"))
}
}
}
#[derive(Debug, Clone, Default)]
pub struct NewToken {
pub name: String,
pub admin: AdminScope,
pub namespace: Option<String>,
pub deployments: Vec<String>,
pub expires_in_secs: Option<u64>,
}
impl NewToken {
pub fn new(name: impl Into<String>) -> Self {
Self {
name: name.into(),
..Default::default()
}
}
pub fn admin(mut self, scope: AdminScope) -> Self {
self.admin = scope;
self
}
pub fn for_deployments<I, S>(mut self, ids: I) -> Self
where
I: IntoIterator<Item = S>,
S: Into<String>,
{
self.deployments = ids.into_iter().map(Into::into).collect();
self
}
pub fn fleet_wide(mut self) -> Self {
self.deployments = vec!["*".into()];
self
}
pub fn in_namespace(mut self, ns: impl Into<String>) -> Self {
self.namespace = Some(ns.into());
self
}
pub fn expires_in(mut self, d: Duration) -> Self {
self.expires_in_secs = Some(d.as_secs());
self
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::transport::stub::Stub;
fn client(stub: Stub) -> (Client, Arc<Stub>) {
let s = Arc::new(stub);
(Client::with_transport(s.clone()), s)
}
#[tokio::test]
async fn a_read_is_typed_and_a_write_is_not() {
let (c, stub) = client(Stub::new().json(
201,
json!({"spec": {"id": "demo", "routes": [], "unknown_future_field": 7},
"kind": "vm", "desired_replicas": 1, "ready": 0, "pending": 1,
"total_in_flight": 0, "vms": []}),
));
let spec = json!({"id": "demo", "routes": [], "unknown_future_field": 7});
let got = c.create_deployment(&spec).await.unwrap();
assert_eq!(got.kind, "vm");
assert_eq!(stub.calls()[0].body.as_ref().unwrap(), &spec);
}
#[tokio::test]
async fn a_failed_request_becomes_a_typed_error() {
let (c, _) = client(Stub::new().json(404, json!({"error": "no deployment \"demo\""})));
let e = c.deployment("demo").await.unwrap_err();
assert!(matches!(&e, Error::NotFound { kind: "deployment", name } if name == "demo"));
}
#[tokio::test]
async fn absence_is_a_bool_not_an_error_where_that_is_the_question() {
let (c, _) = client(Stub::new().json(404, json!({"error": "no deployment \"demo\""})));
assert!(!c.deployment_exists("demo").await.unwrap());
let (c, _) = client(Stub::new().reply(401, "authentication required\n"));
assert!(c.deployment_exists("demo").await.is_err());
}
#[tokio::test]
async fn a_nonzero_exit_is_not_an_error() {
let (c, _) = client(Stub::new().json(
200,
json!({"sandbox_id": "sb-1", "exit_code": 42, "stdout": "", "stderr": "nope\n",
"output": "nope\n"}),
));
let out = c.exec("demo", &ExecRequest::new("false")).await.unwrap();
assert_eq!(out.exit_code, 42);
assert!(!out.ok());
}
#[tokio::test]
async fn a_blank_command_never_reaches_the_wire() {
let (c, stub) = client(Stub::new());
let e = c.exec("demo", &ExecRequest::new(" ")).await.unwrap_err();
assert!(matches!(e, Error::Invalid(_)));
assert_eq!(stub.call_count(), 0, "nothing should have been sent");
}
#[test]
fn the_exec_deadline_outlasts_the_servers_worst_case() {
let waking = ExecRequest::new("x").timeout_secs(60);
assert!(
waking.patience() > Duration::from_secs(60 + ASSUMED_COLD_START_SECS),
"a waking exec must outlast command timeout + cold start"
);
let not_waking = ExecRequest::new("x").timeout_secs(60).no_wake();
assert!(not_waking.patience() < waking.patience());
assert!(not_waking.patience() > Duration::from_secs(60));
let absurd = ExecRequest::new("x").timeout_secs(99_999).no_wake();
assert!(absurd.patience() > Duration::from_secs(3600));
assert_eq!(
ExecRequest::new("x").patient_for(Duration::from_secs(5)).patience(),
Duration::from_secs(5),
"an explicit override wins"
);
}
#[tokio::test]
async fn query_booleans_are_spelled_the_only_way_app_lb_accepts() {
let (c, stub) = client(Stub::new().json(200, json!({"sandbox_id": "s", "outcome": "killed"})));
c.evict_vm("demo", "sb-1", true).await.unwrap();
assert!(stub.calls()[0].path.ends_with("?force=true"), "{:?}", stub.calls()[0].path);
let (c, stub) = client(Stub::new().json(200, json!({"sandbox_id": "s", "outcome": "killed"})));
c.evict_vm("demo", "sb-1", false).await.unwrap();
assert!(stub.calls()[0].path.ends_with("?force=false"));
}
#[tokio::test]
async fn ids_are_escaped_into_the_path() {
let (c, stub) = client(Stub::new().json(404, json!({"error": "no deployment"})));
let _ = c.deployment("a/b?c=d").await;
assert_eq!(stub.calls()[0].path, "/deployments/a%2Fb%3Fc%3Dd");
}
#[tokio::test]
async fn build_and_pull_always_send_a_json_body() {
let (c, stub) = client(Stub::new().json(202, json!({"id": "j1", "deployment": "d",
"kind": "image-build", "status": "running", "started_at": 0})));
c.start_build("demo", None).await.unwrap();
assert_eq!(stub.calls()[0].body, Some(json!({})));
let (c, stub) = client(Stub::new().json(202, json!({"id": "j1", "deployment": "d",
"kind": "image-build", "status": "running", "started_at": 0})));
c.start_build("demo", Some("v2")).await.unwrap();
assert_eq!(stub.calls()[0].body, Some(json!({"ref": "v2"})));
}
#[test]
fn a_metrics_query_serializes_only_what_was_asked_for() {
assert_eq!(MetricsQuery::new().to_query_string(), "");
assert_eq!(
MetricsQuery::new().deployment("sb-1").to_query_string(),
"?deployment=sb-1"
);
assert_eq!(
MetricsQuery::new().summary(true).page(20, 10).to_query_string(),
"?summary=true&limit=10&offset=20"
);
assert_eq!(
MetricsQuery::new().page(0, 10).to_query_string(),
"?limit=10"
);
}
#[tokio::test]
async fn minting_requires_a_name_before_anything_is_sent() {
let (c, stub) = client(Stub::new());
let e = c.mint_token(&NewToken::new(" ")).await.unwrap_err();
assert!(matches!(e, Error::Invalid(_)));
assert_eq!(stub.call_count(), 0);
}
#[tokio::test]
async fn a_minted_token_carries_its_secret_exactly_once() {
let (c, stub) = client(Stub::new().json(
201,
json!({"id": "abc", "name": "ci", "admin": "admin", "deployments": ["*"],
"created_at": 1, "token": "applb_abc_secret"}),
));
let t = c
.mint_token(&NewToken::new("ci").admin(AdminScope::Admin).fleet_wide())
.await
.unwrap();
assert_eq!(t.token, "applb_abc_secret");
assert_eq!(t.summary.id, "abc");
assert_eq!(
stub.calls()[0].body,
Some(json!({"name": "ci", "admin": "admin", "deployments": ["*"]}))
);
}
#[tokio::test]
async fn a_scoped_mint_does_not_quietly_become_fleet_wide() {
let (c, stub) = client(Stub::new().json(
201,
json!({"id": "abc", "name": "a", "admin": "none", "deployments": ["sb-1"],
"created_at": 1, "token": "applb_abc_s"}),
));
c.mint_token(&NewToken::new("a").for_deployments(["sb-1"]))
.await
.unwrap();
let body = stub.calls()[0].body.clone().unwrap();
assert_eq!(body["deployments"], json!(["sb-1"]));
assert_eq!(body["admin"], json!("none"), "the safe default, not admin");
}
#[tokio::test]
async fn gates_are_discovered_from_what_an_anonymous_caller_is_refused() {
let refused = Error::from_response(
401,
"authentication required\n",
"server",
"",
crate::error::Credential::None,
);
assert!(matches!(refused, Error::Unauthorized { .. }));
}
}