mod auth;
mod host;
mod listen;
mod reply;
mod request;
mod routes_fn;
mod routes_runtimes;
mod routes_scripts_blobs_deps;
mod routes_tenants;
mod stream;
mod usage;
use std::convert::Infallible;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Instant;
use anyhow::Context;
use hyper::body::Incoming;
use hyper::{Method, Request, Response, StatusCode};
use zygo_core::spec::{ApiAuth, Spec};
use zygo_core::supervisor::Request as Control;
use zygo_core::supervisor::client::Client;
use crate::cli::{ApiArgs, Cli};
use crate::output::Style;
use auth::authorise;
use host::{drain, healthz, metrics, snapshot};
use listen::{Listen, serve};
use reply::{ApiBody, HttpError, json};
use request::{CallParams, Query, parse_event, read_body, request_key, timeout_header};
use routes_fn::{batch, cancel, exec, list, logs, one_shot, serve_fn, stats, stop, warm};
use routes_runtimes::{call_runtime, runtimes, serve_runtime, stop_runtime, with_out};
use routes_scripts_blobs_deps::{
delete_blob, delete_deps, delete_script, deps, get_blob, get_script, put_blob, put_deps,
put_script,
};
use routes_tenants::{
create_tenant, delete_secret, delete_tenant, list_tokens, mint_token, put_secret, revoke_token,
secret_names, set_limits, tenants,
};
use stream::exec_streaming;
use usage::{Usage, deliver_usage};
pub const API_VERSION: u32 = 1;
pub(super) const MAX_IDLE_CLIENTS: usize = 32;
struct Api {
paths: zygo_core::Paths,
exe: std::path::PathBuf,
token: Option<String>,
deploy: bool,
clients: std::sync::Mutex<Vec<Client>>,
usage: std::sync::Mutex<Usage>,
health: std::sync::Mutex<Option<(Instant, serde_json::Value)>>,
started: Instant,
requests: AtomicU64,
errors: AtomicU64,
}
pub fn run(cli: &Cli, args: &ApiArgs) -> anyhow::Result<u8> {
if args.openapi {
crate::output::json(&super::openapi::document())?;
return Ok(0);
}
let spec = Spec::discover(args.spec_file.path())?.unwrap_or_default();
let api_spec = spec.api.clone().unwrap_or_default();
let listen = Listen::parse(args.listen.as_deref().unwrap_or(&api_spec.listen))?;
let auth = if args.no_auth {
ApiAuth::None
} else {
api_spec.auth
};
let token = match auth {
ApiAuth::Bearer => Some(
std::env::var("ZYGO_API_TOKEN")
.ok()
.filter(|t| !t.is_empty())
.context(
"bearer auth is on and ZYGO_API_TOKEN is not set\n \
→ export ZYGO_API_TOKEN=<a long random string>, \
or pass --no-auth for a unix socket or loopback listener",
)?,
),
ApiAuth::None => {
anyhow::ensure!(
listen.allows_no_auth(),
"refusing to serve without authentication on {listen}\n \
→ an unauthenticated API on a reachable address lets anyone on the \
network run code as you; listen on 127.0.0.1 or a unix socket, \
or set ZYGO_API_TOKEN and use bearer auth"
);
None
}
};
let paths = super::paths(cli);
let exe = std::env::current_exe().context("cannot find this binary to start a supervisor")?;
let first = Client::connect_or_start(&paths, &exe)?;
let api = Arc::new(Api {
paths,
exe,
token,
deploy: args.allow_deploy,
clients: std::sync::Mutex::new(vec![first]),
usage: std::sync::Mutex::new(Usage::default()),
health: std::sync::Mutex::new(None),
started: Instant::now(),
requests: AtomicU64::new(0),
errors: AtomicU64::new(0),
});
let exporter = match &args.otlp_endpoint {
Some(endpoint) => Some(super::otlp::Exporter::new(endpoint, args.otlp_interval.0)?),
None => None,
};
let style = Style::stderr();
eprintln!(
"{} {listen} {} {}",
style.dim("api"),
style.dim(if api.token.is_some() {
"bearer auth"
} else {
"no auth"
}),
style.dim(if api.deploy {
"deploy on: callers may serve, stop and run"
} else {
"call-only: serve, stop and run are refused"
})
);
if let Some(exporter) = &exporter {
eprintln!(
"{} {} every {:?}",
style.dim("otlp"),
exporter.url,
exporter.interval
);
}
let runtime = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()?;
if let Some(exporter) = exporter {
let started = std::time::SystemTime::now();
let for_export = Arc::clone(&api);
runtime.spawn(super::otlp::run(exporter, started, move || {
let api = Arc::clone(&for_export);
async move { snapshot(&api).await }
}));
}
if let Some(url) = &args.usage_webhook {
let url = reqwest::Url::parse(url)
.with_context(|| format!("`{url}` is not a URL for --usage-webhook"))?;
eprintln!(
"{} {url} every {:?}",
style.dim("usage"),
args.usage_interval.get()
);
let for_usage = Arc::clone(&api);
runtime.spawn(deliver_usage(for_usage, url, args.usage_interval.get()));
}
runtime.block_on(serve(listen, api))?;
Ok(0)
}
async fn control<T, F>(api: &Arc<Api>, f: F) -> anyhow::Result<T>
where
T: Send + 'static,
F: FnOnce(&mut Client) -> anyhow::Result<T> + Send + 'static,
{
let api = Arc::clone(api);
tokio::task::spawn_blocking(move || {
let mut client = match api.clients.lock().expect("clients").pop() {
Some(client) => client,
None => Client::connect_or_start(&api.paths, &api.exe)?,
};
let result = f(&mut client);
if result.is_ok() {
let mut idle = api.clients.lock().expect("clients");
if idle.len() < MAX_IDLE_CLIENTS {
idle.push(client);
}
}
result
})
.await
.context("the control request panicked")?
}
async fn handle(req: Request<Incoming>, api: Arc<Api>) -> Result<Response<ApiBody>, Infallible> {
api.requests.fetch_add(1, Ordering::Relaxed);
let response = match route(req, &api).await {
Ok(response) => response,
Err(e) => e.into_response(),
};
if response.status().is_server_error() {
api.errors.fetch_add(1, Ordering::Relaxed);
}
Ok(response)
}
async fn route(req: Request<Incoming>, api: &Arc<Api>) -> Result<Response<ApiBody>, HttpError> {
if req.method() == Method::GET && req.uri().path() == "/healthz" {
return healthz(api).await;
}
let actor = authorise(&req, api)?;
let path = req.uri().path().to_string();
let segments: Vec<&str> = path
.trim_matches('/')
.split('/')
.filter(|s| !s.is_empty())
.collect();
let tenant = actor.tenant().map(str::to_string);
match (req.method(), segments.as_slice()) {
(&Method::GET, ["fn"]) => list(api, tenant).await,
(&Method::POST, ["tenants"]) => {
let body = read_body(req).await?;
actor.operator_only("creating a tenant")?;
create_tenant(api, &body).await
}
(&Method::GET, ["tenants"]) => {
actor.operator_only("listing the tenants")?;
tenants(api, None).await
}
(&Method::GET, ["tenants", id]) => {
let id = id.to_string();
if actor.tenant() != Some(id.as_str()) {
actor.operator_only("reading another tenant")?;
}
tenants(api, Some(id)).await
}
(&Method::DELETE, ["tenants", id]) => {
let id = id.to_string();
actor.may_deploy()?;
delete_tenant(api, id).await
}
(&Method::PATCH, ["tenants", id, "limits"]) => {
let id = id.to_string();
let body = read_body(req).await?;
actor.may_deploy()?;
set_limits(api, id, &body).await
}
(&Method::GET, ["tenants", id, "secrets"]) => {
let id = id.to_string();
if actor.tenant() != Some(id.as_str()) {
actor.operator_only("reading another tenant's secrets")?;
}
secret_names(api, id).await
}
(&Method::PUT, ["tenants", id, "secrets", name]) => {
let (id, name) = (id.to_string(), name.to_string());
let body = read_body(req).await?;
actor.may_deploy()?;
put_secret(api, id, name, &body).await
}
(&Method::DELETE, ["tenants", id, "secrets", name]) => {
let (id, name) = (id.to_string(), name.to_string());
actor.may_deploy()?;
delete_secret(api, id, name).await
}
(&Method::POST, ["tenants", id, "tokens"]) => {
let id = id.to_string();
actor.may_deploy()?;
mint_token(api, Some(id)).await
}
(&Method::POST, ["tokens"]) => {
actor.may_deploy()?;
mint_token(api, None).await
}
(&Method::GET, ["tokens"]) => {
actor.may_deploy()?;
list_tokens(api).await
}
(&Method::DELETE, ["tokens", id]) => {
let id = id.to_string();
actor.may_deploy()?;
revoke_token(api, id).await
}
(&Method::DELETE, ["requests", id]) => {
let id = id.to_string();
cancel(api, id, tenant).await
}
(&Method::POST, ["drain"]) => {
let query = Query::parse(req.uri().query().unwrap_or(""));
let grace_ms = query.number("grace_ms")?.unwrap_or(30_000);
actor.may_deploy()?;
drain(api, grace_ms).await
}
(&Method::GET, ["metrics"]) => metrics(api, tenant).await,
(&Method::GET, ["version"]) => Ok(json(
StatusCode::OK,
&serde_json::json!({
"version": env!("CARGO_PKG_VERSION"),
"api": API_VERSION,
"control": zygo_core::supervisor::CONTROL_VERSION,
"deploy": actor.deploy && actor.is_operator(),
}),
)),
(&Method::PUT, ["fn", name]) => {
let name = name.to_string();
let body = read_body(req).await?;
actor.may_deploy()?;
serve_fn(api, name, &body, tenant).await
}
(&Method::DELETE, ["fn", name]) => {
let name = name.to_string();
actor.may_deploy()?;
stop(api, name).await
}
(&Method::POST, ["run"]) => {
let body = read_body(req).await?;
actor.may_deploy()?;
one_shot(api, &body).await
}
(&Method::GET, ["fn", name, "logs"]) => {
let name = name.to_string();
let query = req.uri().query().unwrap_or("").to_string();
logs(api, name, &query, tenant).await
}
(&Method::POST, ["fn", name]) => {
let name = name.to_string();
let timeout_ms = timeout_header(&req)?;
let key = request_key(&req)?;
let query = Query::parse(req.uri().query().unwrap_or(""));
let streaming = query.flag("stream")?;
let out = query.flag("out")?;
let workspace = with_out(query.blob("workspace")?, out);
let body = read_body(req).await?;
let event = parse_event(&body)?;
if streaming {
return exec_streaming(
api,
Control::Exec {
name,
event,
timeout_ms,
tenant,
key,
stream: true,
workspace,
},
)
.await;
}
exec(api, name, event, timeout_ms, tenant, key, workspace).await
}
(&Method::POST, ["fn", name, "batch"]) => {
let name = name.to_string();
let timeout_ms = timeout_header(&req)?;
let body = read_body(req).await?;
let events: Vec<serde_json::Value> = serde_json::from_slice(&body).map_err(|e| {
HttpError::new(
StatusCode::BAD_REQUEST,
format!("body must be a JSON array of events: {e}"),
)
})?;
batch(api, name, events, timeout_ms, tenant).await
}
(&Method::GET, ["fn", name, "stats"]) => stats(api, name.to_string(), tenant).await,
(&Method::POST, ["fn", name, "warm"]) => warm(api, name.to_string(), tenant).await,
(&Method::GET, ["runtimes"]) => runtimes(api, tenant).await,
(&Method::POST, ["runtimes"]) => {
let body = read_body(req).await?;
actor.may_deploy()?;
serve_runtime(api, &body, tenant).await
}
(&Method::DELETE, ["runtimes", name]) => {
let name = name.to_string();
actor.may_deploy()?;
stop_runtime(api, name).await
}
(&Method::POST, ["runtimes", name, "call"]) => {
let name = name.to_string();
let params = CallParams::parse(&req, tenant)?;
let body = read_body(req).await?;
call_runtime(api, name, &body, params).await
}
(&Method::PUT, ["scripts"]) => {
let body = read_body(req).await?;
put_script(api, &body, tenant).await
}
(&Method::PUT, ["blobs"]) => {
let body = read_body(req).await?;
put_blob(api, &body).await
}
(&Method::GET, ["blobs", digest]) => get_blob(api, digest.to_string()).await,
(&Method::DELETE, ["blobs", digest]) => {
let digest = digest.to_string();
actor.may_deploy()?;
delete_blob(api, digest).await
}
(&Method::POST, ["deps"]) => {
let body = read_body(req).await?;
put_deps(api, &body, tenant).await
}
(&Method::GET, ["deps"]) => deps(api, None, tenant).await,
(&Method::GET, ["deps", id]) => deps(api, Some(id.to_string()), tenant).await,
(&Method::DELETE, ["deps", id]) => {
let id = id.to_string();
actor.may_deploy()?;
delete_deps(api, id).await
}
(&Method::GET, ["scripts", digest]) => get_script(api, digest.to_string()).await,
(&Method::DELETE, ["scripts", digest]) => {
let digest = digest.to_string();
actor.may_deploy()?;
delete_script(api, digest).await
}
(_, ["fn", ..])
| (_, ["requests", ..])
| (_, ["drain"])
| (_, ["metrics"])
| (_, ["version"])
| (_, ["run"])
| (_, ["runtimes", ..])
| (_, ["tenants", ..])
| (_, ["tokens", ..])
| (_, ["blobs", ..])
| (_, ["deps", ..])
| (_, ["scripts", ..]) => Err(HttpError::new(
StatusCode::METHOD_NOT_ALLOWED,
format!("{} {}", req.method(), path),
)),
_ => Err(HttpError::new(
StatusCode::NOT_FOUND,
format!("no route {path}"),
)),
}
}