Make agent creation swarm-only and refuse a name placed on another hive
swarm-controller's POST /api/agents now refuses (409) a name the swarm
has already placed on a different hive: a non-Destroyed declaration in
that hive's wanted state, or a SetAgentWanted node still queued for it.
The same name on the same hive is that agent being re-created and goes
through. A wanted state that cannot be read refuses (503/500) instead of
reading as "placed nowhere". Creations are serialised from that read to
the graph insert so two concurrent creations of one name cannot both
pass.
Hive-level creation is removed: hivectl `agent create` / `request-create`,
HostRequest::Spawn / RequestSpawn, the dashboard POST /api/request-spawn
route, and ApprovalKind::Spawn with its approve/resolve arms and the
approval-carrying `templates::spawn`. The swarm path (deploy request or
wanted-state sweep -> queue_first_deploy -> templates::first_deploy) used
none of them. Old `spawn` approval rows are skipped by collect_lenient,
as `init_config` rows were in a3b672d1.
policy.rs's comment on agent_object_name stated swarm-wide name
uniqueness as a fact; it now says where it is enforced and what that
check cannot see.
Refs #4396
This commit is contained in:
parent
1d8ec00ddc
commit
5785c0024c
35 changed files with 376 additions and 434 deletions
|
|
@ -22,7 +22,6 @@ use crate::lifecycle;
|
|||
/// FinalizeDeploy`, plus an `AfterAny` `DeployTail`, under a
|
||||
/// resource-holding root; ~30-90s)
|
||||
/// - `UpdateMetaInputs` → a `MetaUpdate` DAG (fan-out on completion)
|
||||
/// - `Spawn` → a `Spawn` DAG (`Create → WriteDropin → Reconcile`)
|
||||
///
|
||||
/// Every queued kind — deploys included — resolves its approval row via
|
||||
/// [`resolve_approval_dag`] when the DAG settles terminal.
|
||||
|
|
@ -55,25 +54,6 @@ pub async fn approve(coord: Arc<Coordinator>, id: i64) -> Result<()> {
|
|||
coord.emit_rebuild_queue_snapshot();
|
||||
Ok(())
|
||||
}
|
||||
ApprovalKind::Spawn => {
|
||||
// The spawn's tail `Reconcile` starts the container, so the
|
||||
// new agent's power intent is `Up` from the outset.
|
||||
if let Err(e) = coord
|
||||
.power
|
||||
.set(approval.agent.as_str(), crate::power::Wanted::Up)
|
||||
{
|
||||
tracing::warn!(agent = %approval.agent, error = ?e, "agent_power: seed on spawn failed");
|
||||
}
|
||||
let inserted = coord.job_queue.insert_job(|b| {
|
||||
crate::job_queue::templates::spawn(b, approval.agent.as_str(), id);
|
||||
Vec::new()
|
||||
});
|
||||
if let Err(e) = inserted {
|
||||
return Err(e.context("insert spawn dag"));
|
||||
}
|
||||
coord.emit_rebuild_queue_snapshot();
|
||||
Ok(())
|
||||
}
|
||||
ApprovalKind::SchedulePrompt => {
|
||||
// No queue card for SchedulePrompt — the work is a single
|
||||
// sqlite insert, the actual "running" lifetime lives on
|
||||
|
|
@ -525,30 +505,16 @@ pub(crate) async fn resolve_approval_dag(
|
|||
TerminalState::Failed => Err(anyhow::anyhow!("{}", error.unwrap_or("job dag failed"))),
|
||||
};
|
||||
let mut terminal_tag = None;
|
||||
match approval.kind {
|
||||
ApprovalKind::Spawn => {
|
||||
// Post-spawn forge bookkeeping (config repo mirror, meta
|
||||
// access) — warn-only, then the resolution events + a rescan so
|
||||
// the dashboard reflects the post-spawn state either way.
|
||||
if result.is_ok() {
|
||||
forge_after_first_spawn(coord, approval.agent.as_str()).await;
|
||||
} else {
|
||||
coord.rescan_containers_and_emit().await;
|
||||
crate::dashboard::emit_tombstones_snapshot(coord).await;
|
||||
}
|
||||
if approval.kind == ApprovalKind::MergeConfigPr {
|
||||
terminal_tag = deploy_terminal_tag(approval.agent.as_str(), approval_id, outcome).await;
|
||||
// On a failed deploy, surface the failing build log back onto the
|
||||
// PR so the manager sees why it was rejected without leaving the
|
||||
// forge. Posted here rather than inside a node because this is the
|
||||
// one place that holds the DAG's definitive error — a `MergeVerify`
|
||||
// rejection and a `DeployApply` build failure both land here.
|
||||
if let Err(e) = &result {
|
||||
post_merge_failure_to_pr(coord, &approval, e).await;
|
||||
}
|
||||
ApprovalKind::MergeConfigPr => {
|
||||
terminal_tag = deploy_terminal_tag(approval.agent.as_str(), approval_id, outcome).await;
|
||||
// On a failed deploy, surface the failing build log back onto the
|
||||
// PR so the manager sees why it was rejected without leaving the
|
||||
// forge. Posted here rather than inside a node because this is the
|
||||
// one place that holds the DAG's definitive error — a `MergeVerify`
|
||||
// rejection and a `DeployApply` build failure both land here.
|
||||
if let Err(e) = &result {
|
||||
post_merge_failure_to_pr(coord, &approval, e).await;
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
if let Err(e) = finish_approval(coord, &approval, result, terminal_tag).await {
|
||||
tracing::warn!(approval_id, error = ?e, "approval dag resolved with failure");
|
||||
|
|
@ -604,26 +570,6 @@ fn fetch_approval_for_worker(
|
|||
Ok(approval)
|
||||
}
|
||||
|
||||
/// Forge bookkeeping run once after the very first container spawn:
|
||||
/// mirror the applied repo and grant read access to core/meta. The
|
||||
/// agent's forge user and token are swarm-controller's, not this
|
||||
/// hive's. Also rescans containers so the dashboard reflects the post-spawn state.
|
||||
async fn forge_after_first_spawn(coord: &Arc<Coordinator>, agent: &str) {
|
||||
if let Err(e) = crate::forge::ensure_config_repo(agent).await {
|
||||
tracing::warn!(%agent, error = ?e, "forge: ensure_config_repo after first spawn failed");
|
||||
}
|
||||
if let Some(core_token) = crate::forge::core_token()
|
||||
&& let Err(e) = crate::forge::meta_read_access(agent, &core_token).await
|
||||
{
|
||||
tracing::warn!(%agent, error = ?e, "forge: meta_read_access after first spawn failed");
|
||||
}
|
||||
if let Err(e) = crate::forge::ensure_meta_remote(agent).await {
|
||||
tracing::warn!(%agent, error = ?e, "forge: ensure_meta_remote after first spawn failed");
|
||||
}
|
||||
coord.rescan_containers_and_emit().await;
|
||||
crate::dashboard::emit_tombstones_snapshot(coord).await;
|
||||
}
|
||||
|
||||
async fn finish_approval(
|
||||
coord: &Coordinator,
|
||||
approval: &hive_sh4re::approvals::Approval,
|
||||
|
|
@ -670,35 +616,14 @@ async fn finish_approval(
|
|||
note: note.clone(),
|
||||
description: approval.description.clone(),
|
||||
});
|
||||
// For spawn/rebuild approvals, also surface the underlying action so the
|
||||
// For rebuild approvals, also surface the underlying action so the
|
||||
// manager knows whether the lifecycle step succeeded. The
|
||||
// ApprovalResolved event already carries the same `ok` signal but
|
||||
// separating it lets the manager react to the lifecycle change
|
||||
// without having to special-case approvals.
|
||||
match approval.kind {
|
||||
ApprovalKind::Spawn => {
|
||||
let summary = if ok {
|
||||
format!("agent '{}' spawned", approval.agent)
|
||||
} else {
|
||||
format!(
|
||||
"agent '{}' spawn FAILED: {}",
|
||||
approval.agent,
|
||||
note.as_deref().unwrap_or("unknown error")
|
||||
)
|
||||
};
|
||||
let _ = coord
|
||||
.push_todo_submitter(
|
||||
approval.id,
|
||||
"core",
|
||||
Some(format!("spawned:{}", approval.agent)),
|
||||
summary,
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
// MergeConfigPr ends in a container rebuild — surface a Rebuilt
|
||||
// lifecycle event. (It is never a first spawn — the agent already
|
||||
// exists — so it never needs the Spawned arm above.)
|
||||
// lifecycle event.
|
||||
ApprovalKind::MergeConfigPr => {
|
||||
let summary = crate::coordinator::rebuilt_todo_summary(
|
||||
approval.agent.as_str(),
|
||||
|
|
@ -738,8 +663,8 @@ async fn finish_approval(
|
|||
///
|
||||
/// Caller-specific bits stay OUT of here: fetching the PR head, the
|
||||
/// `verify_commit` gate, and the ff-merge. The agent always already exists here
|
||||
/// (a merge is never a first spawn), so there's no `sync_agents` step — the
|
||||
/// operator `Spawn` flow owns first-time meta registration.
|
||||
/// (a merge is never a first deploy), so there's no `sync_agents` step — the
|
||||
/// first-deploy DAG's `Provision` node owns first-time meta registration.
|
||||
async fn prepare_applied_target(
|
||||
agent: &str,
|
||||
applied_dir: &std::path::Path,
|
||||
|
|
|
|||
|
|
@ -841,7 +841,7 @@ impl Coordinator {
|
|||
/// Call after any mutation that could affect what
|
||||
/// `nixos-container list` returns or what a row's
|
||||
/// `running` / `needs_update` / `needs_login` / `deployed_sha`
|
||||
/// resolves to — lifecycle ops, destroy, approve (post-spawn),
|
||||
/// resolves to — lifecycle ops, destroy,
|
||||
/// rebuild, meta-update, and the crash-watcher's periodic poll.
|
||||
/// Cheap when nothing changed (one `nixos-container list` + a
|
||||
/// `HashMap` diff + zero emits).
|
||||
|
|
|
|||
|
|
@ -83,11 +83,6 @@ pub(super) fn gc_orphans(coord: &Coordinator, approvals: Vec<Approval>) -> Vec<A
|
|||
approvals
|
||||
.into_iter()
|
||||
.filter(|a| {
|
||||
// A Spawn approval is for a not-yet-existent agent; the proposed
|
||||
// dir is supposed to be missing, so absence is not orphanhood.
|
||||
if a.kind == hive_sh4re::approvals::ApprovalKind::Spawn {
|
||||
return true;
|
||||
}
|
||||
if Coordinator::agent_proposed_dir(&a.agent).exists() {
|
||||
true
|
||||
} else {
|
||||
|
|
|
|||
|
|
@ -232,56 +232,3 @@ pub(super) async fn post_op_send(
|
|||
// no full-state refresh in between.
|
||||
(axum::http::StatusCode::OK, "ok").into_response()
|
||||
}
|
||||
|
||||
#[derive(Deserialize, ToSchema)]
|
||||
pub(super) struct RequestSpawnForm {
|
||||
name: String,
|
||||
}
|
||||
|
||||
/// Queue a spawn approval for `name`.
|
||||
#[utoipa::path(
|
||||
post,
|
||||
path = "/api/request-spawn",
|
||||
request_body(content = RequestSpawnForm, content_type = "application/x-www-form-urlencoded"),
|
||||
responses(
|
||||
(status = 200, description = "spawn approval queued", body = String),
|
||||
(status = 500, description = "missing name, or the approval submit failed"),
|
||||
),
|
||||
tag = "misc_api"
|
||||
)]
|
||||
pub(super) async fn post_request_spawn(
|
||||
State(state): State<AppState>,
|
||||
Form(form): Form<RequestSpawnForm>,
|
||||
) -> Response {
|
||||
let name = form.name.trim().to_owned();
|
||||
if name.is_empty() {
|
||||
return error_response("spawn: `name` required");
|
||||
}
|
||||
match state.coord.approvals.submit_kind(
|
||||
&name,
|
||||
hive_sh4re::approvals::ApprovalKind::Spawn,
|
||||
"",
|
||||
None,
|
||||
"operator",
|
||||
None,
|
||||
) {
|
||||
Ok(id) => {
|
||||
tracing::info!(%id, %name, "operator: spawn approval queued via dashboard");
|
||||
// Phase 5b: notify the dashboard event channel so live
|
||||
// subscribers can append the row without a snapshot
|
||||
// refetch. Spawn approvals carry no sha.
|
||||
state
|
||||
.coord
|
||||
.emit_approval_added(crate::coordinator::ApprovalAdded {
|
||||
id,
|
||||
agent: &name,
|
||||
approval_kind: "spawn",
|
||||
sha_short: None,
|
||||
description: None,
|
||||
pr_number: None,
|
||||
});
|
||||
(StatusCode::OK, "ok").into_response()
|
||||
}
|
||||
Err(e) => error_response(&format!("request-spawn {name} failed: {e:#}")),
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -149,7 +149,6 @@ pub async fn serve(
|
|||
.routes(routes!(misc_api::api_stats_hive))
|
||||
.routes(routes!(misc_api::api_container_resources))
|
||||
.routes(routes!(misc_api::post_mark_all_read))
|
||||
.routes(routes!(misc_api::post_request_spawn))
|
||||
.routes(routes!(misc_api::post_op_send))
|
||||
.routes(routes!(build_logs::get_build_logs_all))
|
||||
.routes(routes!(build_logs::get_build_log_for_node))
|
||||
|
|
@ -386,7 +385,6 @@ mod router_build_probe {
|
|||
.routes(routes!(misc_api::api_stats_hive))
|
||||
.routes(routes!(misc_api::api_container_resources))
|
||||
.routes(routes!(misc_api::post_mark_all_read))
|
||||
.routes(routes!(misc_api::post_request_spawn))
|
||||
.routes(routes!(misc_api::post_op_send))
|
||||
.routes(routes!(build_logs::get_build_logs_all))
|
||||
.routes(routes!(build_logs::get_build_log_for_node))
|
||||
|
|
|
|||
|
|
@ -140,7 +140,7 @@ struct ApprovalHistoryView {
|
|||
agent: String,
|
||||
kind: &'static str,
|
||||
/// First 12 chars of the canonical sha (preferred) or
|
||||
/// manager-supplied ref. None for resolved spawn approvals.
|
||||
/// manager-supplied ref. None when the approval carries neither.
|
||||
sha_short: Option<String>,
|
||||
/// `approved` / `denied` / `failed`.
|
||||
status: &'static str,
|
||||
|
|
@ -404,7 +404,7 @@ fn build_transient_views(
|
|||
}
|
||||
|
||||
/// Render each pending approval into its dashboard view (short sha for
|
||||
/// `MergeConfigPr`, just the name for `Spawn`).
|
||||
/// `MergeConfigPr`).
|
||||
/// Project a resolved sqlite row into the lean shape the dashboard
|
||||
/// history tab consumes — no `diff_html` (rendering 30 of them
|
||||
/// per /api/state poll would mean 30 git diffs per refresh).
|
||||
|
|
@ -439,16 +439,6 @@ fn build_approval_views(approvals: Vec<Approval>) -> Vec<ApprovalView> {
|
|||
let mut out = Vec::with_capacity(approvals.len());
|
||||
for a in approvals {
|
||||
out.push(match a.kind {
|
||||
hive_sh4re::approvals::ApprovalKind::Spawn => ApprovalView {
|
||||
id: a.id,
|
||||
agent: a.agent.to_string(),
|
||||
kind: "spawn",
|
||||
sha_short: None,
|
||||
description: a.description,
|
||||
pr_number: None,
|
||||
commit_ref: None,
|
||||
requested_at: a.requested_at,
|
||||
},
|
||||
hive_sh4re::approvals::ApprovalKind::UpdateMetaInputs => ApprovalView {
|
||||
id: a.id,
|
||||
agent: a.agent.to_string(),
|
||||
|
|
|
|||
|
|
@ -31,7 +31,7 @@ pub struct TombstoneView {
|
|||
}
|
||||
|
||||
/// State-dir names that don't appear in the live container list. Each
|
||||
/// one surfaces in the dashboard as a row with R3V1V3 + PURG3 actions.
|
||||
/// one surfaces in the dashboard as a row with a PURG3 action.
|
||||
///
|
||||
/// ⚠️ **This lists every agent whose container is absent, not only destroyed
|
||||
/// ones** — a mid-spawn agent (state dir seeded by `Provision`, container not
|
||||
|
|
|
|||
|
|
@ -52,7 +52,7 @@ pub enum DashboardEvent {
|
|||
/// enough to render the dashboard row without a `/api/state`
|
||||
/// refetch.
|
||||
///
|
||||
/// The approval's own kind (`"merge_config_pr"` / `"spawn"`) lives
|
||||
/// The approval's own kind (`"merge_config_pr"` / `"schedule_prompt"`) lives
|
||||
/// on `approval_kind` rather than `kind` because the latter is taken
|
||||
/// by the serde tag identifying which `DashboardEvent` variant
|
||||
/// this is.
|
||||
|
|
@ -119,8 +119,8 @@ pub enum DashboardEvent {
|
|||
name: String,
|
||||
transient_kind: String,
|
||||
},
|
||||
/// One container row changed — new container appeared (post-spawn
|
||||
/// finalise), an existing one flipped `running` / `needs_update` /
|
||||
/// One container row changed — a new container appeared, an existing
|
||||
/// one flipped `running` / `needs_update` /
|
||||
/// `sha`, etc. Clients upsert by `container.name`. Payload carries
|
||||
/// the full row so cold-loaded clients and event-driven clients
|
||||
/// converge on the same render.
|
||||
|
|
@ -136,10 +136,8 @@ pub enum DashboardEvent {
|
|||
/// `nixos-container destroy` (operator-driven or otherwise) on the
|
||||
/// next rescan.
|
||||
ContainerRemoved { seq: u64, name: String },
|
||||
/// Full snapshot of the tombstones list. Emitted on every
|
||||
/// mutation that could add / remove a tombstone: destroy
|
||||
/// (with or without purge), purge-tombstone, spawn approval
|
||||
/// (which can consume a tombstone of the same name). Snapshot
|
||||
/// Full snapshot of the tombstones list. Emitted by destroy (with
|
||||
/// or without purge) and purge-tombstone. Snapshot
|
||||
/// shape (not diff) because the list is tiny (single-digit
|
||||
/// typical) and recomputing avoids the add/remove races a
|
||||
/// per-row event would have.
|
||||
|
|
|
|||
|
|
@ -122,7 +122,7 @@ pub enum NodeKind {
|
|||
/// node-inventory row for its three responsibilities and why it isn't
|
||||
/// named `AbortDeploy`.
|
||||
DeployTail { agent: String, approval_id: i64 },
|
||||
/// Tail node of an approval-carrying DAG (spawn / opaque deploy /
|
||||
/// Tail node of an approval-carrying DAG (opaque deploy /
|
||||
/// config-PR merge): resolve the approval row from how the work ended.
|
||||
ResolveApproval {
|
||||
approval_id: i64,
|
||||
|
|
|
|||
|
|
@ -440,27 +440,21 @@ pub fn approval_deploy(builder: &JobBuilder, agent: &str, approval_id: i64) {
|
|||
resolve_approval_tails(builder, approval_id, window);
|
||||
}
|
||||
|
||||
/// First-deploy spawn (approval-driven): `Provision` (proposed/applied
|
||||
/// repos, state subvolume, meta registration) then `Create`
|
||||
/// (`nixos-container create`), drop-in write, then `Reconcile` starts
|
||||
/// the container (`wanted = Up` written at approve time). All-or-nothing:
|
||||
/// `Provision` (lease-exempt, precedes the container) is the group root;
|
||||
/// `Create` (child) owns the agent lease; `WriteDropin` + `Reconcile`
|
||||
/// (children of `Create`) borrow it. A failure cancel-cascades the rest —
|
||||
/// unlike rebuild there's no recovery-reconcile (nothing to converge if the
|
||||
/// container was never created). Closed by a `ResolveApproval` tail root edged
|
||||
/// `AfterAny` onto `Provision` — the DAG's only other group-root, so its roll-up
|
||||
/// already carries the whole cascade.
|
||||
pub fn spawn(builder: &JobBuilder, agent: &str, approval_id: i64) {
|
||||
let provision = spawn_nodes(builder, agent);
|
||||
resolve_approval_tails(builder, approval_id, provision);
|
||||
}
|
||||
|
||||
/// The spawn subgraph with no tail, returning its group root.
|
||||
/// First deploy of an agent this hive has never seen, asked for by the swarm:
|
||||
/// `Provision` (proposed/applied repos, state subvolume, meta registration)
|
||||
/// then `Create` (`nixos-container create`), drop-in write, then `Reconcile`
|
||||
/// starts the container (`wanted = Up`, seeded by
|
||||
/// `swarm_status::queue_first_deploy`). All-or-nothing: `Provision`
|
||||
/// (lease-exempt, precedes the container) is the group root; `Create` (child)
|
||||
/// owns the agent lease; `WriteDropin` + `Reconcile` (children of `Create`)
|
||||
/// borrow it. A failure cancel-cascades the rest — unlike rebuild there's no
|
||||
/// recovery-reconcile (nothing to converge if the container was never created).
|
||||
///
|
||||
/// Split out for the same reason [`rebuild_nodes`] is: two callers want the
|
||||
/// same four nodes and disagree only about what closes them.
|
||||
pub(crate) fn spawn_nodes<'a>(builder: &'a JobBuilder, agent: &str) -> Handle<'a> {
|
||||
/// No approval tail: the operator authorised the creation at swarm level, and
|
||||
/// the deploy request carries that authorisation.
|
||||
///
|
||||
/// Returns the group root so a caller can wait on the whole subtree.
|
||||
pub fn first_deploy(builder: &JobBuilder, agent: &str) -> Vec<hive_jobq::NodeGuid> {
|
||||
let a = || agent.to_owned();
|
||||
let provision = builder
|
||||
.node(NodeKind::Provision { agent: a() })
|
||||
|
|
@ -479,20 +473,7 @@ pub(crate) fn spawn_nodes<'a>(builder: &'a JobBuilder, agent: &str) -> Handle<'a
|
|||
.needs(Resource::Agent(a()))
|
||||
.part_of(create)
|
||||
.after_ok(dropin);
|
||||
provision
|
||||
}
|
||||
|
||||
/// First deploy of an agent this hive has never seen, asked for by the swarm.
|
||||
///
|
||||
/// [`spawn`] without the approval tail, and the absence is the point rather
|
||||
/// than an omission: that flow exists because an operator used to approve the
|
||||
/// spawn *at the hive*. When the swarm asks, the operator has already clicked
|
||||
/// create at swarm level — the deploy request carries that authorisation, and a
|
||||
/// second gate here would be asking the same person the same question twice.
|
||||
///
|
||||
/// Returns the group root so a caller can wait on the whole subtree.
|
||||
pub fn first_deploy(builder: &JobBuilder, agent: &str) -> Vec<hive_jobq::NodeGuid> {
|
||||
vec![spawn_nodes(builder, agent).guid()]
|
||||
vec![provision.guid()]
|
||||
}
|
||||
|
||||
/// Teardown: `Stop` → `DestroyContainer` → (`PurgeState`) → `DestroyBookkeeping`.
|
||||
|
|
|
|||
|
|
@ -1624,10 +1624,10 @@ fn pause_shape_signal_drain() {
|
|||
}
|
||||
|
||||
#[test]
|
||||
fn spawn_shape_provision_create_dropin_reconcile() {
|
||||
fn first_deploy_shape_provision_create_dropin_reconcile() {
|
||||
let q = JobQueue::new(1);
|
||||
insert(&q, |builder| {
|
||||
templates::spawn(builder, "newbie", 7);
|
||||
templates::first_deploy(builder, "newbie");
|
||||
});
|
||||
assert_eq!(
|
||||
declared_shape(&q),
|
||||
|
|
@ -1636,14 +1636,6 @@ fn spawn_shape_provision_create_dropin_reconcile() {
|
|||
row("create", Some("provision"), &[]),
|
||||
row("write_dropin", Some("create"), &[]),
|
||||
row("reconcile", Some("create"), &[("write_dropin", "done")]),
|
||||
// One tail per outcome, each edged to accept only that one — so
|
||||
// *which* tail the graph lets run already is the answer, and
|
||||
// nothing branches at runtime. The three differ **only** in their
|
||||
// accepted outcome, which is why `declared_shape` spells the
|
||||
// outcome set out instead of bucketing it.
|
||||
row("resolve_approval", None, &[("provision", "done")]),
|
||||
row("resolve_approval", None, &[("provision", "failed")]),
|
||||
row("resolve_approval", None, &[("provision", "cancelled")]),
|
||||
]
|
||||
);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -109,20 +109,6 @@ async fn write_response(
|
|||
async fn dispatch(req: &HostRequest, coord: Arc<Coordinator>) -> HostResponse {
|
||||
let result: anyhow::Result<HostResponse> = async {
|
||||
Ok(match req {
|
||||
HostRequest::Spawn { name } => handle_spawn(&coord, name.as_str()).await?,
|
||||
HostRequest::RequestSpawn { name } => {
|
||||
tracing::info!(%name, "request_spawn");
|
||||
let id = coord.approvals.submit_kind(
|
||||
name.as_str(),
|
||||
hive_sh4re::approvals::ApprovalKind::Spawn,
|
||||
"",
|
||||
None,
|
||||
"operator",
|
||||
None,
|
||||
)?;
|
||||
tracing::info!(%id, %name, "spawn approval queued");
|
||||
HostResponse::success()
|
||||
}
|
||||
HostRequest::Kill { name } => submit_single(&coord, name.as_str(), Verb::Kill).await,
|
||||
HostRequest::Restart { name } => {
|
||||
submit_single(&coord, name.as_str(), Verb::Restart).await
|
||||
|
|
@ -278,49 +264,6 @@ async fn dispatch(req: &HostRequest, coord: Arc<Coordinator>) -> HostResponse {
|
|||
}
|
||||
}
|
||||
|
||||
/// Create + start the container for `name` and bind its MCP listener. On a
|
||||
/// failed spawn nothing was registered: post a swarm notice and return the error.
|
||||
async fn handle_spawn(coord: &Arc<Coordinator>, name: &str) -> Result<HostResponse> {
|
||||
tracing::info!(%name, "spawn");
|
||||
let agent_dir = crate::paths::agent_runtime_dir(name);
|
||||
let hive = coord.hive_env();
|
||||
let paths = Coordinator::agent_paths(name, agent_dir)?;
|
||||
// lifecycle::spawn creates the runtime dir internally before start.
|
||||
// MCP listener registration is event-driven: bind immediately on
|
||||
// success so the harness can connect on its first turn without
|
||||
// waiting for any poll interval.
|
||||
match lifecycle::spawn(name, &hive, &paths).await {
|
||||
Ok(()) => {
|
||||
if let Err(e) = coord.power.set(name, crate::power::Wanted::Up) {
|
||||
tracing::warn!(%name, error = ?e, "agent_power: set wanted=up failed");
|
||||
}
|
||||
// Bind the MCP listener now that the container is starting up.
|
||||
// The harness connects to this socket on its first turn.
|
||||
coord.register_agent(name)?;
|
||||
crate::swarm_notices::notify(
|
||||
"core",
|
||||
Some(format!("spawned:{name}")),
|
||||
format!("agent '{name}' spawned"),
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
Err(e) => {
|
||||
// Spawn failed: register_agent was never called, so there is
|
||||
// nothing to unregister. Notify the swarm and propagate.
|
||||
crate::swarm_notices::notify(
|
||||
"core",
|
||||
Some(format!("spawned:{name}")),
|
||||
format!("agent '{name}' spawn FAILED: {e:#}"),
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
return Err(e);
|
||||
}
|
||||
}
|
||||
Ok(HostResponse::success())
|
||||
}
|
||||
|
||||
/// `hivectl pause|resume` / the dashboard toggle: write or remove the
|
||||
/// agent's pause marker.
|
||||
///
|
||||
|
|
|
|||
|
|
@ -31,7 +31,7 @@ pub(crate) async fn submit_merge_config_pr(
|
|||
anyhow::bail!(
|
||||
"applied repo missing for agent '{agent}' (expected at {}) — \
|
||||
merge_config_pr requires the agent to be fully provisioned; \
|
||||
spawn the agent first (operator spawn) before opening config PRs",
|
||||
create it first (`swarmctl agent create`) before opening config PRs",
|
||||
applied_dir.display()
|
||||
);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,7 +1,7 @@
|
|||
//! Approval queue. Requests are submitted by an agent
|
||||
//! (`RequestSchedulePrompt`), the config-PR webhook (`MergeConfigPr`), or
|
||||
//! the operator (`Spawn`); the user approves/denies via the host admin CLI;
|
||||
//! on approval the host runs the corresponding action.
|
||||
//! (`RequestSchedulePrompt`) or the config-PR webhook (`MergeConfigPr`); the
|
||||
//! user approves/denies via the host admin CLI; on approval the host runs the
|
||||
//! corresponding action.
|
||||
//!
|
||||
//! `UpdateMetaInputs` rows are legacy: the MCP tool that queued them was
|
||||
//! removed and nothing produces the kind any more. The variant and
|
||||
|
|
@ -78,7 +78,7 @@ impl Approvals {
|
|||
/// Insert a new pending approval row. `fetched_sha` may be supplied
|
||||
/// when the sha is already known at submission time (e.g. `MergeConfigPr`
|
||||
/// fetches the PR head before inserting), making the insert + sha-set
|
||||
/// atomic. Pass `None` when the kind carries no sha (e.g. `Spawn`).
|
||||
/// atomic. Pass `None` when the kind carries no sha (e.g. `SchedulePrompt`).
|
||||
pub fn submit_kind(
|
||||
&self,
|
||||
agent: &str,
|
||||
|
|
@ -382,7 +382,6 @@ fn row_to_approval(row: &rusqlite::Row<'_>) -> rusqlite::Result<Approval> {
|
|||
// Column order: id, agent, kind, commit_ref, requested_at, status, resolved_at, note, fetched_sha, description.
|
||||
let kind: String = row.get(2)?;
|
||||
let kind = match kind.as_str() {
|
||||
"spawn" => ApprovalKind::Spawn,
|
||||
"update_meta_inputs" => ApprovalKind::UpdateMetaInputs,
|
||||
"schedule_prompt" => ApprovalKind::SchedulePrompt,
|
||||
"merge_config_pr" => ApprovalKind::MergeConfigPr,
|
||||
|
|
@ -435,7 +434,6 @@ fn row_to_approval(row: &rusqlite::Row<'_>) -> rusqlite::Result<Approval> {
|
|||
|
||||
fn kind_from_str(s: &str) -> Result<ApprovalKind> {
|
||||
Ok(match s {
|
||||
"spawn" => ApprovalKind::Spawn,
|
||||
"update_meta_inputs" => ApprovalKind::UpdateMetaInputs,
|
||||
"schedule_prompt" => ApprovalKind::SchedulePrompt,
|
||||
"merge_config_pr" => ApprovalKind::MergeConfigPr,
|
||||
|
|
@ -467,7 +465,7 @@ mod tests {
|
|||
None,
|
||||
)
|
||||
.unwrap();
|
||||
db.submit_kind("b", ApprovalKind::Spawn, "", None, "b", None)
|
||||
db.submit_kind("b", ApprovalKind::SchedulePrompt, "", None, "b", None)
|
||||
.unwrap();
|
||||
db.submit_kind("c", ApprovalKind::UpdateMetaInputs, "[]", None, "c", None)
|
||||
.unwrap();
|
||||
|
|
@ -508,7 +506,14 @@ mod tests {
|
|||
// final — re-cancelling errors instead of silently overwriting.
|
||||
let (_dir, _path, db) = open_temp();
|
||||
let id = db
|
||||
.submit_kind("a", ApprovalKind::Spawn, "deadbeef", None, "a", None)
|
||||
.submit_kind(
|
||||
"a",
|
||||
ApprovalKind::MergeConfigPr,
|
||||
"deadbeef",
|
||||
None,
|
||||
"a",
|
||||
None,
|
||||
)
|
||||
.unwrap();
|
||||
db.mark_cancelled(id, "manager").expect("first cancel");
|
||||
let err = db
|
||||
|
|
@ -567,7 +572,7 @@ mod tests {
|
|||
let raw = Connection::open(&path).unwrap();
|
||||
raw.execute(
|
||||
"INSERT INTO approvals (agent, kind, commit_ref, requested_at, status)
|
||||
VALUES ('old', 'spawn', '', 0, 'pending')",
|
||||
VALUES ('old', 'merge_config_pr', '', 0, 'pending')",
|
||||
[],
|
||||
)
|
||||
.unwrap();
|
||||
|
|
|
|||
|
|
@ -370,10 +370,9 @@ async fn handle_deploy_request(
|
|||
/// (`power::Store::get_or_seed`), and at that point the container is freshly
|
||||
/// created but not started — so an unseeded row locks the agent to `Offline`
|
||||
/// on its very first reconcile and `Reconcile` never emits the `Start` node.
|
||||
/// Setting the row up front closes that window, the same way
|
||||
/// `actions::approve`'s `ApprovalKind::Spawn` arm does for the
|
||||
/// operator-approved path. A second caller open-coding the insert would lose
|
||||
/// exactly that, and the agent would come up stopped for no visible reason.
|
||||
/// Setting the row up front closes that window. A second caller open-coding
|
||||
/// the insert would lose exactly that, and the agent would come up stopped for
|
||||
/// no visible reason.
|
||||
pub(crate) fn queue_first_deploy(
|
||||
coord: &std::sync::Arc<crate::coordinator::Coordinator>,
|
||||
agent: &str,
|
||||
|
|
|
|||
|
|
@ -6,8 +6,8 @@
|
|||
//! re-register all running agents.
|
||||
//!
|
||||
//! After startup, listeners are managed event-driven:
|
||||
//! - `run_start` (the job queue's `Start` node) and `server::handle_spawn`
|
||||
//! call `register_agent` once the container is started.
|
||||
//! - `run_start` (the job queue's `Start` node) calls `register_agent`
|
||||
//! once the container is started.
|
||||
//! - `kill`/`destroy` paths call `unregister_agent`.
|
||||
//!
|
||||
//! No recurring poll is needed because c0re owns the listener lifecycle. An
|
||||
|
|
|
|||
Loading…
Reference in a new issue