use bytes::Bytes;
use http_body_util::{BodyExt, Full};
use hyper::header::HeaderValue;
use hyper::{Response, StatusCode};
use zygo_core::pool::Outcome;
use zygo_core::supervisor::{ControlError, Response as Reply};
use super::request::REQUEST_ID_HEADER;
pub(super) type ApiBody = http_body_util::combinators::BoxBody<Bytes, std::io::Error>;
pub(super) fn whole(bytes: Bytes) -> ApiBody {
Full::new(bytes).map_err(|never| match never {}).boxed()
}
#[derive(Debug)]
pub(super) struct HttpError {
pub(super) status: StatusCode,
pub(super) body: serde_json::Value,
close: bool,
}
impl HttpError {
pub(super) fn new(status: StatusCode, message: impl std::fmt::Display) -> HttpError {
HttpError {
status,
body: serde_json::json!({ "error": message.to_string() }),
close: false,
}
}
pub(super) fn closing(status: StatusCode, message: impl std::fmt::Display) -> HttpError {
HttpError {
close: true,
..HttpError::new(status, message)
}
}
pub(super) fn into_response(self) -> Response<ApiBody> {
let mut response = json(self.status, &self.body);
if self.close {
response
.headers_mut()
.insert(hyper::header::CONNECTION, HeaderValue::from_static("close"));
}
response
}
}
impl From<anyhow::Error> for HttpError {
fn from(e: anyhow::Error) -> HttpError {
HttpError::new(StatusCode::INTERNAL_SERVER_ERROR, format!("{e:#}"))
}
}
pub(super) fn reply_to_json(reply: Reply) -> (StatusCode, serde_json::Value) {
match reply {
Reply::Executed { outcome } => outcome_to_json(*outcome),
Reply::Busy {
name,
in_flight,
queued,
limit,
} => (
StatusCode::TOO_MANY_REQUESTS,
serde_json::json!({
"error": format!("`{name}` is at its concurrency limit"),
"in_flight": in_flight, "queued": queued, "limit": limit,
}),
),
Reply::Warmed { name, state } => (
StatusCode::OK,
serde_json::json!({ "name": name, "state": state }),
),
Reply::Functions { functions } => (
StatusCode::OK,
serde_json::json!({ "functions": functions }),
),
Reply::Served {
name,
runtime,
rss_kb,
imports_ms,
warm_ms,
warnings,
change,
} => (
StatusCode::OK,
serde_json::json!({
"name": name, "runtime": runtime, "rss_kb": rss_kb,
"imports_ms": imports_ms, "warm_ms": warm_ms,
"warnings": warnings, "change": change,
}),
),
Reply::Stopped { names } => (StatusCode::OK, serde_json::json!({ "stopped": names })),
Reply::Drained {
in_flight,
grace_ms,
} => (
StatusCode::OK,
serde_json::json!({
"drained": in_flight == 0,
"in_flight": in_flight,
"grace_ms": grace_ms,
}),
),
Reply::Secrets { names } => (
StatusCode::OK,
serde_json::json!({ "secrets": names }),
),
Reply::Cancelled { id, started } => (
StatusCode::OK,
serde_json::json!({
"cancelled": true,
"request_id": id,
"started": started,
}),
),
Reply::Script {
digest,
size,
existed,
} => (
StatusCode::OK,
serde_json::json!({ "sha256": digest, "size": size, "existed": existed }),
),
Reply::RuntimeServed {
name,
runtime,
warm,
rss_kb,
imports_ms,
warm_ms,
warnings,
change,
} => (
StatusCode::OK,
serde_json::json!({
"name": name, "runtime": runtime, "warm": warm, "rss_kb": rss_kb,
"imports_ms": imports_ms, "warm_ms": warm_ms,
"warnings": warnings, "change": change,
}),
),
Reply::Runtimes { runtimes } => {
(StatusCode::OK, serde_json::json!({ "runtimes": runtimes }))
}
Reply::Logs {
name,
entries,
next,
} => (
StatusCode::OK,
serde_json::json!({ "name": name, "entries": entries, "next": next }),
),
Reply::Error { code, message } => {
let status = match code {
ControlError::NotFound => StatusCode::NOT_FOUND,
ControlError::BadSpec => StatusCode::BAD_REQUEST,
ControlError::WarmFailed => StatusCode::SERVICE_UNAVAILABLE,
ControlError::Unauthorised => StatusCode::FORBIDDEN,
ControlError::AboveCeiling => StatusCode::UNPROCESSABLE_ENTITY,
ControlError::DepsBuilding => StatusCode::SERVICE_UNAVAILABLE,
ControlError::VersionMismatch
| ControlError::CallFailed
| ControlError::BadMessage => StatusCode::INTERNAL_SERVER_ERROR,
};
(
status,
serde_json::json!({ "error": message, "code": code.as_str() }),
)
}
other => (
StatusCode::INTERNAL_SERVER_ERROR,
serde_json::json!({ "error": format!("unexpected reply {other:?}") }),
),
}
}
pub(super) fn reply_to_response(reply: Reply) -> Response<ApiBody> {
let (status, body) = reply_to_json(reply);
let mut response = json(status, &body);
if status == StatusCode::TOO_MANY_REQUESTS {
response
.headers_mut()
.insert("retry-after", hyper::header::HeaderValue::from_static("1"));
}
if body.get("code") == Some(&serde_json::json!("deps_building")) {
response
.headers_mut()
.insert("retry-after", hyper::header::HeaderValue::from_static("5"));
}
if let Some(id) = body.get("request_id").and_then(|v| v.as_str())
&& let Ok(value) = hyper::header::HeaderValue::from_str(id)
{
response.headers_mut().insert(REQUEST_ID_HEADER, value);
}
response
}
pub(super) fn outcome_to_json(outcome: Outcome) -> (StatusCode, serde_json::Value) {
let metrics = serde_json::json!({
"wall_ms": outcome.metrics.wall_ms,
"cpu_ms": outcome.metrics.cpu_ms,
"peak_rss_kb": outcome.metrics.peak_rss_kb,
});
let id = outcome.id;
if outcome.cancelled {
return (
StatusCode::from_u16(499).expect("a valid status"),
serde_json::json!({
"error": "the request was cancelled",
"cancelled": true,
"request_id": id,
"stdout": outcome.stdout,
"stderr": outcome.stderr,
"metrics": metrics,
}),
);
}
if outcome.timed_out {
return (
StatusCode::REQUEST_TIMEOUT,
serde_json::json!({
"error": "the request exceeded the function's timeout and was killed",
"request_id": id,
"stderr": outcome.stderr,
"metrics": metrics,
}),
);
}
if outcome.stuck {
return (
StatusCode::GATEWAY_TIMEOUT,
serde_json::json!({
"error": "the sandbox stopped reporting this request and it was killed",
"stuck": true,
"request_id": id,
"stdout": outcome.stdout,
"stderr": outcome.stderr,
"metrics": metrics,
}),
);
}
if let Some(error) = outcome.error {
return (
StatusCode::INTERNAL_SERVER_ERROR,
serde_json::json!({
"error": error,
"request_id": id,
"stdout": outcome.stdout,
"stderr": outcome.stderr,
"exit_code": outcome.exit_code,
"metrics": metrics,
}),
);
}
let mut body = serde_json::json!({
"result": outcome.result,
"request_id": id,
"stdout": outcome.stdout,
"stderr": outcome.stderr,
"metrics": metrics,
});
if let Some(tar) = outcome.workspace {
body["workspace"] = tar.into();
}
(StatusCode::OK, body)
}
pub(super) fn json(status: StatusCode, body: &serde_json::Value) -> Response<ApiBody> {
Response::builder()
.status(status)
.header("content-type", "application/json")
.body(whole(Bytes::from(
serde_json::to_vec(body).unwrap_or_else(|_| b"{}".to_vec()),
)))
.expect("a valid response")
}
#[cfg(test)]
mod tests {
use super::*;
use zygo_core::protocol::Metrics;
fn outcome(error: Option<&str>, timed_out: bool) -> Outcome {
Outcome {
tenant: "default".into(),
function: "resize".into(),
script: None,
id: "00000001".into(),
cancelled: false,
stuck: false,
workspace: None,
exit_code: if error.is_some() { 1 } else { 0 },
result: serde_json::json!({ "ok": true }),
stdout: "hi\n".into(),
stderr: String::new(),
error: error.map(str::to_string),
metrics: Metrics {
wall_ms: 1.5,
..Default::default()
},
timed_out,
}
}
#[test]
fn the_status_codes_are_the_design_documents() {
let (s, body) = outcome_to_json(outcome(None, false));
assert_eq!(s, StatusCode::OK);
assert_eq!(body["result"]["ok"], true);
assert_eq!(body["metrics"]["wall_ms"], 1.5);
let (s, body) = outcome_to_json(outcome(Some("ZeroDivisionError"), false));
assert_eq!(s, StatusCode::INTERNAL_SERVER_ERROR);
assert_eq!(body["error"], "ZeroDivisionError");
let (s, _) = outcome_to_json(outcome(None, true));
assert_eq!(s, StatusCode::REQUEST_TIMEOUT);
let (s, body) = reply_to_json(Reply::Busy {
name: "f".into(),
in_flight: 4,
queued: 16,
limit: 4,
});
assert_eq!(s, StatusCode::TOO_MANY_REQUESTS);
assert_eq!(body["limit"], 4);
}
#[test]
fn a_deadline_kill_is_a_408_not_a_500() {
let killed = Outcome {
error: Some("killed by SIGKILL (out of memory, or the deadline expired)".into()),
exit_code: 137,
..outcome(None, true)
};
assert_eq!(outcome_to_json(killed).0, StatusCode::REQUEST_TIMEOUT);
let oom = Outcome {
error: Some("killed by SIGKILL (out of memory, or the deadline expired)".into()),
exit_code: 137,
..outcome(None, false)
};
assert_eq!(outcome_to_json(oom).0, StatusCode::INTERNAL_SERVER_ERROR);
}
#[test]
fn control_errors_map_to_distinct_statuses() {
let status = |code| reply_to_json(Reply::error(code, "x")).0;
assert_eq!(status(ControlError::NotFound), StatusCode::NOT_FOUND);
assert_eq!(status(ControlError::BadSpec), StatusCode::BAD_REQUEST);
assert_eq!(
status(ControlError::WarmFailed),
StatusCode::SERVICE_UNAVAILABLE
);
assert_eq!(
status(ControlError::CallFailed),
StatusCode::INTERNAL_SERVER_ERROR
);
}
#[test]
fn the_new_replies_have_statuses_of_their_own() {
let (status, body) = reply_to_json(Reply::Served {
name: "resize".into(),
runtime: "python3.12".into(),
rss_kb: 2048,
imports_ms: 40.0,
warm_ms: 120.0,
warnings: vec!["no timeout set".into()],
change: zygo_core::supervisor::Change::Replaced,
});
assert_eq!(status, StatusCode::OK);
assert_eq!(body["change"], "replaced");
assert_eq!(body["warnings"][0], "no timeout set");
let (status, body) = reply_to_json(Reply::Stopped {
names: vec!["resize".into()],
});
assert_eq!(status, StatusCode::OK);
assert_eq!(body["stopped"][0], "resize");
let (status, body) = reply_to_json(Reply::Logs {
name: "resize".into(),
entries: Vec::new(),
next: 7,
});
assert_eq!(status, StatusCode::OK);
assert_eq!(body["next"], 7);
}
}