hyperhive/hive-c0re/src/dashboard/lifecycle_ops.rs

459 lines
17 KiB
Rust

//! Container lifecycle endpoints for the dashboard.
//!
//! Rebuild / restart / start / stop (hard + graceful) / update-all all
//! submit DAGs to the job queue (`job_queue::submit`), so each shows a
//! visible queued→running transient on the dashboard — a direct
//! sub-second start/stop only flashed the badge. Start/stop also
//! persist the agent's `wanted` power intent before submitting; the
//! DAG's `Reconcile` converges to it. Destroy delegates to
//! `actions::destroy` (optionally purging).
use axum::{
extract::{Form, Path as AxumPath, Query, State},
http::StatusCode,
response::{IntoResponse, Response},
};
use serde::Deserialize;
use utoipa::{IntoParams, ToSchema};
/// Query params for `post_kill` / `post_restart`. `?graceful=1` routes to
/// the graceful-stop/-restart orchestration (quiesce the harness, flush
/// `/state`, then container stop/restart) instead of an immediate hard
/// action. Defaults false → today's hard kill/restart.
#[derive(Deserialize, IntoParams)]
pub(super) struct GracefulParams {
#[serde(default)]
graceful: bool,
}
use super::{AppState, Ident, error_response, guard_agent_name, strip_container_prefix};
use crate::job_queue::{Source, submit};
use crate::{actions, lifecycle};
/// `POST /api/rebuild/{name}` — queue a rebuild DAG for `name`.
#[utoipa::path(
post,
path = "/api/rebuild/{name}",
params(("name" = String, Path, description = "agent name")),
responses(
(status = 200, description = "rebuild queued", body = String),
(status = 400, description = "bad agent name"),
(status = 404, description = "no such agent"),
),
tag = "lifecycle_ops"
)]
pub(super) async fn post_rebuild(
State(state): State<AppState>,
AxumPath(name): AxumPath<String>,
) -> Response {
let logical = strip_container_prefix(&name);
if let Some(reject) = guard_agent_name(&state, &logical).await {
return reject;
}
submit::rebuild(
&state.coord,
&logical,
Source::Manual,
"manual via dashboard ↻ R3BU1LD button".to_owned(),
);
(StatusCode::OK, "ok").into_response()
}
/// `POST /api/kill/{name}?graceful=1` — stop `name`, hard by default or
/// gracefully (quiesce → drain → stop) when `graceful=1`.
#[utoipa::path(
post,
path = "/api/kill/{name}",
params(
("name" = String, Path, description = "agent name"),
GracefulParams,
),
responses(
(status = 200, description = "stop queued/performed", body = String),
(status = 400, description = "bad agent name"),
(status = 404, description = "no such agent"),
),
tag = "lifecycle_ops"
)]
pub(super) async fn post_kill(
State(state): State<AppState>,
AxumPath(name): AxumPath<String>,
Query(params): Query<GracefulParams>,
) -> Response {
let logical = strip_container_prefix(&name);
if let Some(reject) = guard_agent_name(&state, &logical).await {
return reject;
}
if params.graceful {
// Graceful stop: submit the quiesce DAG (signal the harness →
// one stop-checkpoint turn → drain → container stop, with a
// timeout fallback to a hard stop). The agent's lifecycle
// lease keeps it from racing an in-flight rebuild for the same
// agent, and per-node progress surfaces on the queue snapshot.
submit::graceful_stop(
&state.coord,
&logical,
Source::Manual,
"manual via dashboard graceful stop".to_owned(),
)
.await;
return (StatusCode::OK, "ok").into_response();
}
// Manager is stoppable from the dashboard like any other
// agent. The host's dashboard server keeps running (it's
// hive-c0re, not the manager container), per-agent approvals
// submitted by other sub-agents still process through the
// host-side approval queue without the manager up, and
// operator-driven meta-input updates work from the dashboard
// either way. The MCP-surface self-kill guard in
// `socket_server.rs::Request::Kill` stays in place: a
// manager calling Kill on its own container is self-suicide
// mid-call, not a legitimate operator action.
submit::stop(
&state.coord,
&logical,
Source::Manual,
"manual via dashboard stop".to_owned(),
)
.await;
(StatusCode::OK, "ok").into_response()
}
/// `POST /api/restart/{name}?graceful=1` — restart `name`, hard by default
/// or gracefully (quiesce → drain → restart) when `graceful=1`.
#[utoipa::path(
post,
path = "/api/restart/{name}",
params(
("name" = String, Path, description = "agent name"),
GracefulParams,
),
responses(
(status = 200, description = "restart queued/performed", body = String),
(status = 400, description = "bad agent name"),
(status = 404, description = "no such agent"),
),
tag = "lifecycle_ops"
)]
pub(super) async fn post_restart(
State(state): State<AppState>,
AxumPath(name): AxumPath<String>,
Query(params): Query<GracefulParams>,
) -> Response {
let logical = strip_container_prefix(&name);
if let Some(reject) = guard_agent_name(&state, &logical).await {
return reject;
}
if params.graceful {
submit::graceful_restart(
&state.coord,
&logical,
Source::Manual,
"manual via dashboard graceful restart".to_owned(),
)
.await;
return (StatusCode::OK, "ok").into_response();
}
submit::restart(
&state.coord,
&logical,
Source::Manual,
"manual via dashboard ↺ R3START button".to_owned(),
)
.await;
(StatusCode::OK, "ok").into_response()
}
/// Query params for `post_start`. `?paused=1` writes the pause marker
/// before (or instead of) starting — see `post_start`'s doc.
#[derive(Deserialize, IntoParams)]
pub(super) struct StartParams {
#[serde(default)]
paused: bool,
}
/// `POST /api/start/{name}?paused=1` — start `name`, optionally paused.
///
/// Plain `?paused=1` mirrors `hivectl agent <name> start --paused`: if
/// `name` is already running, this just writes the pause marker in place
/// and returns without submitting a start DAG (nothing to start). If it's
/// down, the marker is written *before* the start DAG is submitted, so
/// the container comes up paused rather than racing the harness's own
/// pause-gate poll against an already-in-flight start.
#[utoipa::path(
post,
path = "/api/start/{name}",
params(
("name" = String, Path, description = "agent name"),
StartParams,
),
responses(
(status = 200, description = "start queued (or paused in place)", body = String),
(status = 400, description = "bad agent name"),
(status = 404, description = "no such agent"),
(status = 500, description = "pause marker write failed"),
),
tag = "lifecycle_ops"
)]
pub(super) async fn post_start(
State(state): State<AppState>,
AxumPath(name): AxumPath<String>,
Query(params): Query<StartParams>,
) -> Response {
let logical = strip_container_prefix(&name);
if let Some(reject) = guard_agent_name(&state, &logical).await {
return reject;
}
if params.paused {
let ident = match Ident::parse(&logical) {
Ok(i) => i,
Err(e) => {
return (StatusCode::BAD_REQUEST, format!("bad agent name: {e}")).into_response();
}
};
let already_running = lifecycle::is_running(&logical).await;
if let Err(e) = crate::coordinator::Coordinator::set_paused(&ident, true).await {
return error_response(&format!("pause {logical}: {e}"));
}
state.coord.rescan_containers_and_emit().await;
if already_running {
// Already up — pausing in place is the whole request, no DAG
// to submit.
return (StatusCode::OK, "ok").into_response();
}
}
submit::start(
&state.coord,
&logical,
Source::Manual,
"manual via dashboard start".to_owned(),
)
.await;
(StatusCode::OK, "ok").into_response()
}
/// `POST /api/pause/{name}` — write the pause marker for `name`.
///
/// Unlike the lifecycle ops above this is not a DAG: it writes a single
/// marker file, which the harness stats at the top of its serve loop.
/// Works on stopped containers too (the marker is sticky and takes effect
/// when the container next boots). Triggers an immediate rescan so the
/// `paused` badge flips on the dashboard without waiting for the next
/// periodic sweep.
#[utoipa::path(
post,
path = "/api/pause/{name}",
params(("name" = String, Path, description = "agent name")),
responses(
(status = 200, description = "pause marker written", body = String),
(status = 400, description = "bad agent name"),
(status = 404, description = "no such agent"),
(status = 500, description = "marker write failed"),
),
tag = "lifecycle_ops"
)]
pub(super) async fn post_pause(
State(state): State<AppState>,
AxumPath(name): AxumPath<String>,
) -> Response {
let logical = strip_container_prefix(&name);
if let Some(reject) = guard_agent_name(&state, &logical).await {
return reject;
}
let ident = match Ident::parse(&logical) {
Ok(i) => i,
Err(e) => return (StatusCode::BAD_REQUEST, format!("bad agent name: {e}")).into_response(),
};
if let Err(e) = crate::coordinator::Coordinator::set_paused(&ident, true).await {
return error_response(&format!("pause {logical}: {e}"));
}
state.coord.rescan_containers_and_emit().await;
(StatusCode::OK, "ok").into_response()
}
/// `POST /api/resume/{name}` — remove the pause marker for `name`.
///
/// The inverse of `post_pause`. Removing a non-existent marker is a no-op
/// (idempotent). Triggers an immediate rescan so the paused badge clears.
#[utoipa::path(
post,
path = "/api/resume/{name}",
params(("name" = String, Path, description = "agent name")),
responses(
(status = 200, description = "pause marker removed", body = String),
(status = 400, description = "bad agent name"),
(status = 404, description = "no such agent"),
(status = 500, description = "marker removal failed"),
),
tag = "lifecycle_ops"
)]
pub(super) async fn post_resume(
State(state): State<AppState>,
AxumPath(name): AxumPath<String>,
) -> Response {
let logical = strip_container_prefix(&name);
if let Some(reject) = guard_agent_name(&state, &logical).await {
return reject;
}
let ident = match Ident::parse(&logical) {
Ok(i) => i,
Err(e) => return (StatusCode::BAD_REQUEST, format!("bad agent name: {e}")).into_response(),
};
if let Err(e) = crate::coordinator::Coordinator::set_paused(&ident, false).await {
return error_response(&format!("resume {logical}: {e}"));
}
state.coord.rescan_containers_and_emit().await;
(StatusCode::OK, "ok").into_response()
}
/// Form fields for `post_resource_limits`. Both fields are optional strings;
/// an empty value clears the per-agent override for that field, falling back
/// to the hive-wide default.
#[derive(Deserialize, Default, ToSchema)]
pub(super) struct ResourceLimitsForm {
#[serde(default)]
cpu_quota: String,
#[serde(default)]
memory_max: String,
}
/// `POST /api/resource-limits/{name}` — write per-agent CPU/memory limit
/// overrides for `name`.
///
/// An empty `cpu_quota` or `memory_max` field clears that field's override,
/// falling back to the hive-wide default. Both empty together removes the
/// agent's entry entirely. The new drop-in is written immediately — the
/// limits take effect on the next container start or restart. Triggers an
/// immediate rescan so `ContainerView.cpu_quota`/`memory_max` update on
/// the dashboard via SSE without waiting for the next periodic sweep.
#[utoipa::path(
post,
path = "/api/resource-limits/{name}",
params(("name" = String, Path, description = "agent name")),
request_body(content = ResourceLimitsForm, content_type = "application/x-www-form-urlencoded"),
responses(
(status = 200, description = "limits written", body = String),
(status = 400, description = "bad agent name"),
(status = 404, description = "no such agent"),
(status = 422, description = "invalid cpu_quota/memory_max value"),
(status = 500, description = "commit or drop-in write failed"),
),
tag = "lifecycle_ops"
)]
pub(super) async fn post_resource_limits(
State(state): State<AppState>,
AxumPath(name): AxumPath<String>,
Form(form): Form<ResourceLimitsForm>,
) -> Response {
let logical = strip_container_prefix(&name);
if let Some(reject) = guard_agent_name(&state, &logical).await {
return reject;
}
let ident = match Ident::parse(&logical) {
Ok(i) => i,
Err(e) => return (StatusCode::BAD_REQUEST, format!("bad agent name: {e}")).into_response(),
};
let cpu_quota = if form.cpu_quota.is_empty() {
None
} else {
Some(form.cpu_quota.as_str())
};
let memory_max = if form.memory_max.is_empty() {
None
} else {
Some(form.memory_max.as_str())
};
if let Some(v) = cpu_quota
&& let Err(e) = crate::resource_limits::validate_cpu_quota(v)
{
return (StatusCode::UNPROCESSABLE_ENTITY, e).into_response();
}
if let Some(v) = memory_max
&& let Err(e) = crate::resource_limits::validate_memory_max(v)
{
return (StatusCode::UNPROCESSABLE_ENTITY, e).into_response();
}
let limits = crate::resource_limits::AgentLimits {
cpu_quota: cpu_quota.map(str::to_owned),
memory_max: memory_max.map(str::to_owned),
};
if let Err(e) = crate::meta::commit_resource_limits(ident.as_str(), &limits).await {
return error_response(&format!("set limits {logical}: {e:#}"));
}
let agent_dir = crate::paths::agent_runtime_dir(ident.as_str());
let hive = state.coord.hive_env();
let paths = crate::coordinator::Coordinator::agent_paths(ident.as_str(), agent_dir);
if let Err(e) = crate::lifecycle::write_dropins(ident.as_str(), &hive, &paths).await {
return error_response(&format!("write_dropins {logical}: {e:#}"));
}
state.coord.rescan_containers_and_emit().await;
(StatusCode::OK, "ok").into_response()
}
/// `POST /api/update-all` — queue a rebuild DAG for every live agent
/// container.
#[utoipa::path(
post,
path = "/api/update-all",
responses((status = 200, description = "rebuilds queued", body = String)),
tag = "lifecycle_ops"
)]
pub(super) async fn post_update_all(State(state): State<AppState>) -> Response {
let containers = lifecycle::list().await.unwrap_or_default();
for container in containers {
let Some(logical) = container
.strip_prefix(lifecycle::AGENT_PREFIX)
.map(str::to_owned)
else {
continue;
};
submit::rebuild(
&state.coord,
&logical,
Source::Manual,
"manual via dashboard 🌀 UPDATE ALL".to_owned(),
);
}
(StatusCode::OK, "ok").into_response()
}
#[derive(Deserialize, Default, ToSchema)]
pub(super) struct DestroyForm {
#[serde(default)]
purge: Option<String>,
}
/// `POST /api/destroy/{name}` — destroy `name`'s container. Form field
/// `purge` (any non-empty value, e.g. `"on"`) also wipes the retained
/// state dir instead of leaving a tombstone.
#[utoipa::path(
post,
path = "/api/destroy/{name}",
params(("name" = String, Path, description = "agent name")),
request_body(content = DestroyForm, content_type = "application/x-www-form-urlencoded"),
responses(
(status = 200, description = "destroyed", body = String),
(status = 400, description = "bad agent name"),
(status = 404, description = "no such agent"),
(status = 500, description = "destroy failed"),
),
tag = "lifecycle_ops"
)]
pub(super) async fn post_destroy(
State(state): State<AppState>,
AxumPath(name): AxumPath<String>,
Form(form): Form<DestroyForm>,
) -> Response {
if let Some(reject) = guard_agent_name(&state, &name).await {
return reject;
}
// Checkbox semantics: any non-empty value (axum sends "on") = purge.
let purge = form.purge.as_deref().is_some_and(|v| !v.is_empty());
// `actions::destroy` rescans the container list on success, so the
// `ContainerRemoved` event lands before we return 200. The matching
// form carries `data-no-refresh`.
match actions::destroy(&state.coord, &name, purge).await {
Ok(()) => (StatusCode::OK, "ok").into_response(),
Err(e) => error_response(&format!("destroy {name} failed: {e:#}")),
}
}