use std::sync::Arc;
use hyper::{Response, StatusCode};
use zygo_core::supervisor::{Request as Control, Response as Reply};
use super::reply::{ApiBody, HttpError, json, reply_to_response};
use super::{Api, control};
use zygo_core::spec::Layer;
use zygo_core::supervisor::WorkspaceRequest;
use super::request::CallParams;
use super::stream::exec_streaming;
use super::usage::count_usage;
#[derive(Debug, serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct ServeRuntimeRequest {
name: String,
#[serde(default)]
layer: Layer,
#[serde(default)]
base_dir: Option<std::path::PathBuf>,
#[serde(default)]
deps: Option<String>,
}
pub(super) async fn serve_runtime(
api: &Arc<Api>,
body: &[u8],
tenant: Option<String>,
) -> Result<Response<ApiBody>, HttpError> {
let request: ServeRuntimeRequest = serde_json::from_slice(body).map_err(|e| {
HttpError::new(
StatusCode::BAD_REQUEST,
format!("body is not a runtime definition: {e}"),
)
})?;
if let Some(base_dir) = &request.base_dir
&& !base_dir.is_absolute()
{
return Err(HttpError::new(
StatusCode::BAD_REQUEST,
format!(
"`base_dir` must be absolute, and `{}` is not\n \
→ it names a directory on the host this API runs on",
base_dir.display()
),
));
}
let base_dir = request
.base_dir
.unwrap_or_else(|| std::path::PathBuf::from("/"));
let deps = request.deps;
let private_net = api.private_net;
let reply = control(api, move |c| {
Ok(c.send(&Control::ServeRuntime {
tenant,
name: request.name,
spec: None,
layer: Box::new(request.layer),
base_dir,
deps,
allow_host_net: false,
allow_private_net: private_net,
allow_unlimited: false,
})?)
})
.await?;
Ok(reply_to_response(reply))
}
pub(super) async fn runtimes(
api: &Arc<Api>,
tenant: Option<String>,
) -> Result<Response<ApiBody>, HttpError> {
let reply = control(api, |c| Ok(c.send(&Control::Runtimes)?)).await?;
match (reply, tenant) {
(Reply::Runtimes { runtimes }, Some(id)) => {
let runtimes: Vec<_> = runtimes.into_iter().filter(|r| r.tenant == id).collect();
Ok(json(
StatusCode::OK,
&serde_json::json!({ "runtimes": runtimes }),
))
}
(other, _) => Ok(reply_to_response(other)),
}
}
pub(super) async fn stop_runtime(
api: &Arc<Api>,
name: String,
) -> Result<Response<ApiBody>, HttpError> {
let wanted = name.clone();
let reply = control(api, move |c| Ok(c.send(&Control::StopRuntime { name })?)).await?;
match reply {
Reply::Stopped { names } if names.is_empty() => Err(HttpError::new(
StatusCode::NOT_FOUND,
format!("no runtime named `{wanted}`"),
)),
other => Ok(reply_to_response(other)),
}
}
#[derive(Debug, serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct CallRuntimeRequest {
script: ScriptRef,
#[serde(default)]
event: serde_json::Value,
#[serde(default)]
entry_point: Option<String>,
#[serde(default)]
workspace: Option<WorkspaceBody>,
}
#[derive(Debug, Default, serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct WorkspaceBody {
#[serde(default)]
inline: Option<String>,
#[serde(default)]
blob: Option<String>,
}
impl From<WorkspaceBody> for WorkspaceRequest {
fn from(body: WorkspaceBody) -> WorkspaceRequest {
WorkspaceRequest {
inline: body.inline,
blob: body.blob,
collect: false,
}
}
}
pub(super) fn with_out(workspace: Option<WorkspaceRequest>, out: bool) -> Option<WorkspaceRequest> {
match (workspace, out) {
(Some(w), out) => Some(WorkspaceRequest { collect: out, ..w }),
(None, true) => Some(WorkspaceRequest {
collect: true,
..WorkspaceRequest::default()
}),
(None, false) => None,
}
}
#[derive(Debug, serde::Deserialize)]
#[serde(untagged)]
enum ScriptRef {
Digest(String),
Source { source: String },
}
pub(super) async fn call_runtime(
api: &Arc<Api>,
name: String,
body: &[u8],
params: CallParams,
) -> Result<Response<ApiBody>, HttpError> {
let CallParams {
timeout_ms,
tenant,
key,
streaming,
out,
} = params;
let request: CallRuntimeRequest = serde_json::from_slice(body).map_err(|e| {
HttpError::new(
StatusCode::BAD_REQUEST,
format!(
"body must be {{\"script\": \"sha256:…\" | {{\"source\": \"…\"}}, \
\"event\": …}}: {e}"
),
)
})?;
let mut script = match request.script {
ScriptRef::Digest(digest) => zygo_core::protocol::Script {
path: None,
source: None,
digest: Some(digest),
entry_point: None,
},
ScriptRef::Source { source } => zygo_core::protocol::Script::inline(source),
};
script.entry_point = request.entry_point;
let call = Control::ExecScript {
key,
runtime: name,
script,
event: request.event,
timeout_ms,
tenant,
stream: streaming,
workspace: with_out(request.workspace.map(Into::into), out),
};
if streaming {
return exec_streaming(api, call).await;
}
let reply = control(api, move |c| Ok(c.send(&call)?)).await?;
count_usage(api, &reply);
Ok(reply_to_response(reply))
}