From 02dcf4d0288250086543725562601393834cd7ae Mon Sep 17 00:00:00 2001 From: damocles Date: Tue, 9 Jun 2026 11:43:31 +0200 Subject: [PATCH] refactor(#1474): extract dispatch arm logic in manager_server + hive-priv --- hive-c0re/src/manager_server.rs | 441 ++++++++++++++++++-------------- hive-priv/src/main.rs | 156 ++++++----- 2 files changed, 336 insertions(+), 261 deletions(-) diff --git a/hive-c0re/src/manager_server.rs b/hive-c0re/src/manager_server.rs index abd74a32..a9ea7e8c 100644 --- a/hive-c0re/src/manager_server.rs +++ b/hive-c0re/src/manager_server.rs @@ -74,11 +74,6 @@ async fn serve(stream: UnixStream, coord: Arc) -> Result<()> { } } -#[allow( - clippy::too_many_lines, - reason = "flat dispatch table: one match arm per manager-socket request \ - variant; splitting it would scatter the routing logic" -)] async fn dispatch(req: &ManagerRequest, coord: &Arc) -> ManagerResponse { // Delegate all variants shared with the agent socket to the common handler. if let Some(resp) = crate::agent_server::dispatch_shared(req, MANAGER_AGENT, coord).await { @@ -86,112 +81,16 @@ async fn dispatch(req: &ManagerRequest, coord: &Arc) -> ManagerResp } match req { ManagerRequest::RequestInitConfig { name, description } => { - tracing::info!(%name, "manager: request_init_config"); - match submit_init_config(coord, name, description.clone()) { - Ok(_id) => ManagerResponse::Ok, - Err(e) => ManagerResponse::Err { - message: format!("{e:#}"), - }, - } - } - ManagerRequest::Kill { name } => { - tracing::info!(%name, "manager: kill"); - let result: Result<()> = async { - lifecycle::kill(name).await?; - coord.unregister_agent(name); - Ok(()) - } - .await; - match result { - Ok(()) => { - coord.notify_manager(&hive_sh4re::HelperEvent::Killed { - agent: name.clone(), - }); - ManagerResponse::Ok - } - Err(e) => ManagerResponse::Err { - message: format!("{e:#}"), - }, - } - } - ManagerRequest::Start { name } => { - tracing::info!(%name, "manager: start"); - match lifecycle::start(name).await { - Ok(()) => { - coord.kick_agent(name, "container started"); - ManagerResponse::Ok - } - Err(e) => ManagerResponse::Err { - message: format!("{e:#}"), - }, - } - } - ManagerRequest::Restart { name } => { - tracing::info!(%name, "manager: enqueue restart"); - coord.rebuild_queue.enqueue( - crate::rebuild_queue::QueueKind::Restart, - name.to_owned(), - crate::rebuild_queue::QueueSource::Manual, - "manager `restart` tool".to_owned(), - None, - ); - coord.emit_rebuild_queue_snapshot(); - ManagerResponse::Ok - } - ManagerRequest::Update { name } => { - tracing::info!(%name, "manager: enqueue update"); - coord.rebuild_queue.enqueue( - crate::rebuild_queue::QueueKind::Rebuild, - name.to_owned(), - crate::rebuild_queue::QueueSource::Manual, - "manager `update` tool".to_owned(), - None, - ); - coord.emit_rebuild_queue_snapshot(); - ManagerResponse::Ok + handle_manager_init_config(coord, name, description.clone()) } + ManagerRequest::Kill { name } => handle_manager_kill(coord, name).await, + ManagerRequest::Start { name } => handle_manager_start(coord, name).await, + ManagerRequest::Restart { name } => handle_manager_restart(coord, name), + ManagerRequest::Update { name } => handle_manager_update(coord, name), ManagerRequest::RequestUpdateMetaInputs { inputs, description, - } => { - let label = if inputs.is_empty() { - "all inputs".to_string() - } else { - inputs.join(", ") - }; - tracing::info!(%label, "manager: request_update_meta_inputs"); - // Encode the inputs list as JSON and store it in commit_ref - // (there's no git commit involved; the field carries the - // payload for the approval handler to decode at run time). - let commit_ref = serde_json::to_string(inputs).unwrap_or_default(); - let id = match coord - .approvals - .submit_kind( - MANAGER_AGENT, - hive_sh4re::ApprovalKind::UpdateMetaInputs, - &commit_ref, - description.as_deref(), - ) - .map_err(|e| anyhow::anyhow!("{e:#}")) - { - Ok(id) => id, - Err(e) => { - return ManagerResponse::Err { - message: format!("queue update_meta_inputs approval: {e:#}"), - }; - } - }; - tracing::info!(%id, %label, "update_meta_inputs approval queued"); - coord.emit_approval_added( - id, - MANAGER_AGENT, - "update_meta_inputs", - None, - None, - description.clone(), - ); - ManagerResponse::Ok - } + } => handle_request_update_meta_inputs(coord, inputs, description.as_deref()), ManagerRequest::RequestSchedulePrompt(payload) => { handle_request_schedule_prompt(coord, MANAGER_AGENT, payload) } @@ -219,101 +118,24 @@ async fn dispatch(req: &ManagerRequest, coord: &Arc) -> ManagerResp targets_remove: targets_remove.clone(), }, ), - ManagerRequest::ListSchedules => match coord.scheduled_prompts.list() { - Ok(schedules) => ManagerResponse::Schedules { - schedules: schedules.into_iter().map(schedule_to_wire).collect(), - }, - Err(e) => ManagerResponse::Err { - message: format!("list scheduled prompts: {e:#}"), - }, - }, + ManagerRequest::ListSchedules => handle_list_schedules(coord), ManagerRequest::FireScheduleNow { id } => { handle_fire_schedule_now(coord, MANAGER_AGENT, *id).await } - ManagerRequest::GetLogs { agent, lines } => { - let n = lines.unwrap_or(50); - // `journalctl -M` wants the container name (`h-`), - // not the logical agent name. `container_name` adds the prefix. - // The `-M` read needs root, so it goes through hive-priv. - let machine = crate::lifecycle::container_name(agent); - tracing::info!(%agent, %machine, %n, "manager: get_logs"); - match crate::priv_client::read_container_journal( - &machine, - hive_sh4re::priv_proto::JournalQuery { - lines: n, - ..Default::default() - }, - ) - .await - { - Ok((stdout, stderr)) => { - let content = if stdout.is_empty() { stderr } else { stdout }; - ManagerResponse::Logs { content } - } - Err(e) => ManagerResponse::Err { - message: format!("get_logs: {e:#}"), - }, - } - } + ManagerRequest::GetLogs { agent, lines } => handle_get_logs(agent, *lines).await, ManagerRequest::RequestApplyCommit { agent, commit_ref, description, - } => { - tracing::info!(%agent, %commit_ref, "manager: request_apply_commit"); - match submit_apply_commit(coord, agent, commit_ref, description.as_deref()).await { - Ok((id, sha)) => { - tracing::info!(%id, %agent, manager_ref = %commit_ref, %sha, "approval queued + proposal tag planted"); - ManagerResponse::Ok - } - Err(e) => ManagerResponse::Err { - message: format!("{e:#}"), - }, - } - } + } => handle_manager_apply_commit(coord, agent, commit_ref, description.as_deref()).await, ManagerRequest::GetLooseEnds { agent } => { - let result = match agent.as_deref() { - Some("*") => { - // Hive-wide query requires query_agent_state capability. - if !crate::capabilities::has_cap( - MANAGER_AGENT, - hive_sh4re::Capability::QueryAgentState, - ) { - return ManagerResponse::Err { - message: - "query_agent_state capability required for hive-wide loose ends" - .into(), - }; - } - crate::loose_ends::hive_wide(coord) - } - Some(name) => crate::loose_ends::for_agent(coord, name), - None => crate::loose_ends::for_agent(coord, MANAGER_AGENT), - }; - match result { - Ok(loose_ends) => ManagerResponse::LooseEnds { loose_ends }, - Err(e) => ManagerResponse::Err { - message: format!("{e:#}"), - }, - } + handle_manager_loose_ends(coord, agent.as_deref()) } ManagerRequest::CountPendingReminders { agent } => { - let target = agent.as_deref().unwrap_or(MANAGER_AGENT); - match coord.broker.count_pending_reminders_for(target) { - Ok(count) => ManagerResponse::PendingRemindersCount { count }, - Err(e) => ManagerResponse::Err { - message: format!("{e:#}"), - }, - } + handle_manager_count_pending_reminders(coord, agent.as_deref()) } ManagerRequest::ReminderRollup { since_secs, agent } => { - let target = agent.as_deref().unwrap_or(MANAGER_AGENT); - match coord.broker.reminder_rollup_for(target, *since_secs) { - Ok(stats) => ManagerResponse::ReminderRollup(stats), - Err(e) => ManagerResponse::Err { - message: format!("{e:#}"), - }, - } + handle_manager_reminder_rollup(coord, agent.as_deref(), *since_secs) } _ => ManagerResponse::Err { message: "request not handled on manager socket".to_owned(), @@ -321,6 +143,245 @@ async fn dispatch(req: &ManagerRequest, coord: &Arc) -> ManagerResp } } +/// `RequestInitConfig` (manager socket) — queue an `InitConfig` +/// approval. No topology check: the manager can act on any agent. +fn handle_manager_init_config( + coord: &Arc, + name: &str, + description: Option, +) -> ManagerResponse { + tracing::info!(%name, "manager: request_init_config"); + match submit_init_config(coord, name, description) { + Ok(_id) => ManagerResponse::Ok, + Err(e) => ManagerResponse::Err { + message: format!("{e:#}"), + }, + } +} + +/// `Kill` (manager socket) — kill the container, unregister it, notify. +async fn handle_manager_kill(coord: &Arc, name: &str) -> ManagerResponse { + tracing::info!(%name, "manager: kill"); + let result: Result<()> = async { + lifecycle::kill(name).await?; + coord.unregister_agent(name); + Ok(()) + } + .await; + match result { + Ok(()) => { + coord.notify_manager(&hive_sh4re::HelperEvent::Killed { + agent: name.to_owned(), + }); + ManagerResponse::Ok + } + Err(e) => ManagerResponse::Err { + message: format!("{e:#}"), + }, + } +} + +/// `Start` (manager socket) — start the container, kick its next turn. +async fn handle_manager_start(coord: &Arc, name: &str) -> ManagerResponse { + tracing::info!(%name, "manager: start"); + match lifecycle::start(name).await { + Ok(()) => { + coord.kick_agent(name, "container started"); + ManagerResponse::Ok + } + Err(e) => ManagerResponse::Err { + message: format!("{e:#}"), + }, + } +} + +/// `Restart` (manager socket) — enqueue a restart. +fn handle_manager_restart(coord: &Arc, name: &str) -> ManagerResponse { + tracing::info!(%name, "manager: enqueue restart"); + coord.rebuild_queue.enqueue( + crate::rebuild_queue::QueueKind::Restart, + name.to_owned(), + crate::rebuild_queue::QueueSource::Manual, + "manager `restart` tool".to_owned(), + None, + ); + coord.emit_rebuild_queue_snapshot(); + ManagerResponse::Ok +} + +/// `Update` (manager socket) — enqueue a rebuild. +fn handle_manager_update(coord: &Arc, name: &str) -> ManagerResponse { + tracing::info!(%name, "manager: enqueue update"); + coord.rebuild_queue.enqueue( + crate::rebuild_queue::QueueKind::Rebuild, + name.to_owned(), + crate::rebuild_queue::QueueSource::Manual, + "manager `update` tool".to_owned(), + None, + ); + coord.emit_rebuild_queue_snapshot(); + ManagerResponse::Ok +} + +/// `RequestUpdateMetaInputs` — queue an `UpdateMetaInputs` approval +/// carrying the JSON-encoded input list in `commit_ref` (no git commit +/// is involved; the field is the payload the approval handler decodes). +fn handle_request_update_meta_inputs( + coord: &Arc, + inputs: &[String], + description: Option<&str>, +) -> ManagerResponse { + let label = if inputs.is_empty() { + "all inputs".to_string() + } else { + inputs.join(", ") + }; + tracing::info!(%label, "manager: request_update_meta_inputs"); + let commit_ref = serde_json::to_string(inputs).unwrap_or_default(); + let id = match coord + .approvals + .submit_kind( + MANAGER_AGENT, + hive_sh4re::ApprovalKind::UpdateMetaInputs, + &commit_ref, + description, + ) + .map_err(|e| anyhow::anyhow!("{e:#}")) + { + Ok(id) => id, + Err(e) => { + return ManagerResponse::Err { + message: format!("queue update_meta_inputs approval: {e:#}"), + }; + } + }; + tracing::info!(%id, %label, "update_meta_inputs approval queued"); + coord.emit_approval_added( + id, + MANAGER_AGENT, + "update_meta_inputs", + None, + None, + description.map(str::to_owned), + ); + ManagerResponse::Ok +} + +/// `ListSchedules` — snapshot every scheduled prompt onto the wire. +fn handle_list_schedules(coord: &Arc) -> ManagerResponse { + match coord.scheduled_prompts.list() { + Ok(schedules) => ManagerResponse::Schedules { + schedules: schedules.into_iter().map(schedule_to_wire).collect(), + }, + Err(e) => ManagerResponse::Err { + message: format!("list scheduled prompts: {e:#}"), + }, + } +} + +/// `GetLogs` — read a child container's journal via hive-priv (the +/// `-M` read needs root). `journalctl -M` wants the `h-` machine +/// name, which `container_name` derives. +async fn handle_get_logs(agent: &str, lines: Option) -> ManagerResponse { + let n = lines.unwrap_or(50); + let machine = crate::lifecycle::container_name(agent); + tracing::info!(%agent, %machine, %n, "manager: get_logs"); + match crate::priv_client::read_container_journal( + &machine, + hive_sh4re::priv_proto::JournalQuery { + lines: n, + ..Default::default() + }, + ) + .await + { + Ok((stdout, stderr)) => { + let content = if stdout.is_empty() { stderr } else { stdout }; + ManagerResponse::Logs { content } + } + Err(e) => ManagerResponse::Err { + message: format!("get_logs: {e:#}"), + }, + } +} + +/// `RequestApplyCommit` (manager socket) — queue an apply-commit +/// approval + plant the `proposal/` tag. +async fn handle_manager_apply_commit( + coord: &Arc, + agent: &str, + commit_ref: &str, + description: Option<&str>, +) -> ManagerResponse { + tracing::info!(%agent, %commit_ref, "manager: request_apply_commit"); + match submit_apply_commit(coord, agent, commit_ref, description).await { + Ok((id, sha)) => { + tracing::info!(%id, %agent, manager_ref = %commit_ref, %sha, "approval queued + proposal tag planted"); + ManagerResponse::Ok + } + Err(e) => ManagerResponse::Err { + message: format!("{e:#}"), + }, + } +} + +/// `GetLooseEnds` (manager socket) — manager's own threads by default, +/// a named agent's when given, or hive-wide for `"*"` (which requires +/// the `query_agent_state` capability). +fn handle_manager_loose_ends(coord: &Arc, agent: Option<&str>) -> ManagerResponse { + let result = match agent { + Some("*") => { + if !crate::capabilities::has_cap(MANAGER_AGENT, hive_sh4re::Capability::QueryAgentState) + { + return ManagerResponse::Err { + message: "query_agent_state capability required for hive-wide loose ends" + .into(), + }; + } + crate::loose_ends::hive_wide(coord) + } + Some(name) => crate::loose_ends::for_agent(coord, name), + None => crate::loose_ends::for_agent(coord, MANAGER_AGENT), + }; + match result { + Ok(loose_ends) => ManagerResponse::LooseEnds { loose_ends }, + Err(e) => ManagerResponse::Err { + message: format!("{e:#}"), + }, + } +} + +/// `CountPendingReminders` (manager socket) — count pending reminders +/// for the target (defaults to the manager itself). +fn handle_manager_count_pending_reminders( + coord: &Arc, + agent: Option<&str>, +) -> ManagerResponse { + let target = agent.unwrap_or(MANAGER_AGENT); + match coord.broker.count_pending_reminders_for(target) { + Ok(count) => ManagerResponse::PendingRemindersCount { count }, + Err(e) => ManagerResponse::Err { + message: format!("{e:#}"), + }, + } +} + +/// `ReminderRollup` (manager socket) — roll up reminders fired in the +/// last `since_secs` for the target (defaults to the manager itself). +fn handle_manager_reminder_rollup( + coord: &Arc, + agent: Option<&str>, + since_secs: u64, +) -> ManagerResponse { + let target = agent.unwrap_or(MANAGER_AGENT); + match coord.broker.reminder_rollup_for(target, since_secs) { + Ok(stats) => ManagerResponse::ReminderRollup(stats), + Err(e) => ManagerResponse::Err { + message: format!("{e:#}"), + }, + } +} + /// `request_apply_commit` takes a commit SHA only — not a branch or /// tag name. A branch is mutable; pinning the proposal to a concrete /// sha keeps "what the manager asked to deploy" unambiguous and means diff --git a/hive-priv/src/main.rs b/hive-priv/src/main.rs index 58733b8b..ba8e05f9 100644 --- a/hive-priv/src/main.rs +++ b/hive-priv/src/main.rs @@ -157,8 +157,10 @@ async fn write_line_event(writer: &mut OwnedWriteHalf, stream: PrivStream, data: /// the returned strings are empty. #[allow( clippy::too_many_lines, - reason = "flat dispatch table: one match arm per privileged request variant; \ - splitting it would scatter the routing logic" + reason = "flat routing table over the privileged request variants — the \ + logic-bearing arms (resource-limits / daemon-reload / \ + restart-matrix / create+update) are extracted to helpers; the \ + rest are one-line validate-then-delegate dispatches kept inline" )] async fn exec(req: PrivRequest, writer: &mut OwnedWriteHalf) -> Result<(String, String)> { match req { @@ -178,35 +180,11 @@ async fn exec(req: PrivRequest, writer: &mut OwnedWriteHalf) -> Result<(String, } PrivRequest::UpdateContainer { ref name, stream } => { - validate_container_name(name)?; - let flake_ref = agent_flake_ref(name); - let args = [ - "update", - &container_system_name(name), - "--flake", - &flake_ref, - ]; - if stream { - container_run_streaming(&args, writer).await - } else { - container_run(&args).await - } + container_flake_action("update", name, stream, writer).await } PrivRequest::CreateContainer { ref name, stream } => { - validate_container_name(name)?; - let flake_ref = agent_flake_ref(name); - let args = [ - "create", - &container_system_name(name), - "--flake", - &flake_ref, - ]; - if stream { - container_run_streaming(&args, writer).await - } else { - container_run(&args).await - } + container_flake_action("create", name, stream, writer).await } PrivRequest::DestroyContainer { ref name } => { @@ -242,15 +220,7 @@ async fn exec(req: PrivRequest, writer: &mut OwnedWriteHalf) -> Result<(String, ref container, ref memory_max, ref cpu_quota, - } => { - validate_container_system_name(container)?; - let dir = format!("/run/systemd/system/container@{container}.service.d"); - std::fs::create_dir_all(&dir).with_context(|| format!("create {dir}"))?; - let path = format!("{dir}/hyperhive-limits.conf"); - let content = format!("[Service]\nMemoryMax={memory_max}\nCPUQuota={cpu_quota}\n"); - std::fs::write(&path, content).with_context(|| format!("write {path}"))?; - Ok((String::new(), String::new())) - } + } => write_resource_limits(container, memory_max, cpu_quota), PrivRequest::RemoveServiceDropin { ref container } => { validate_container_system_name(container)?; @@ -261,21 +231,7 @@ async fn exec(req: PrivRequest, writer: &mut OwnedWriteHalf) -> Result<(String, Ok((String::new(), String::new())) } - PrivRequest::DaemonReload => { - let out = Command::new("systemctl") - .arg("daemon-reload") - .output() - .await - .context("invoke systemctl daemon-reload")?; - if !out.status.success() { - bail!( - "systemctl daemon-reload failed ({}): {}", - out.status, - String::from_utf8_lossy(&out.stderr).trim() - ); - } - Ok((String::new(), String::new())) - } + PrivRequest::DaemonReload => daemon_reload().await, PrivRequest::ReloadGatewayNginx => sync_gateway_nginx().await, @@ -327,29 +283,87 @@ async fn exec(req: PrivRequest, writer: &mut OwnedWriteHalf) -> Result<(String, } PrivRequest::RestartMatrixDaemon { ref agent_name } => { - validate_agent_name(agent_name)?; - let machine = format!("--machine=h-{agent_name}"); - let unit = "hive-matrix-daemon.service"; - let out = Command::new("systemctl") - .args([&machine, "restart", unit]) - .output() - .await - .with_context(|| format!("systemctl restart {unit} in container h-{agent_name}"))?; - if !out.status.success() { - bail!( - "systemctl restart {unit} in h-{agent_name} exited {}: {}", - out.status, - String::from_utf8_lossy(&out.stderr).trim() - ); - } - Ok(( - String::from_utf8_lossy(&out.stdout).into_owned(), - String::from_utf8_lossy(&out.stderr).into_owned(), - )) + restart_matrix_daemon(agent_name).await } } } +/// Shared body for `CreateContainer` / `UpdateContainer`: validate the +/// name, build the `nixos-container … --flake ` argv, and +/// run it (streaming line events to `writer` when `stream` is set). +async fn container_flake_action( + verb: &str, + name: &str, + stream: bool, + writer: &mut OwnedWriteHalf, +) -> Result<(String, String)> { + validate_container_name(name)?; + let flake_ref = agent_flake_ref(name); + let args = [verb, &container_system_name(name), "--flake", &flake_ref]; + if stream { + container_run_streaming(&args, writer).await + } else { + container_run(&args).await + } +} + +/// `WriteResourceLimits` — drop a systemd `MemoryMax`/`CPUQuota` +/// override into the container service's drop-in dir. +fn write_resource_limits( + container: &str, + memory_max: &str, + cpu_quota: &str, +) -> Result<(String, String)> { + validate_container_system_name(container)?; + let dir = format!("/run/systemd/system/container@{container}.service.d"); + std::fs::create_dir_all(&dir).with_context(|| format!("create {dir}"))?; + let path = format!("{dir}/hyperhive-limits.conf"); + let content = format!("[Service]\nMemoryMax={memory_max}\nCPUQuota={cpu_quota}\n"); + std::fs::write(&path, content).with_context(|| format!("write {path}"))?; + Ok((String::new(), String::new())) +} + +/// `DaemonReload` — `systemctl daemon-reload` on the host. +async fn daemon_reload() -> Result<(String, String)> { + let out = Command::new("systemctl") + .arg("daemon-reload") + .output() + .await + .context("invoke systemctl daemon-reload")?; + if !out.status.success() { + bail!( + "systemctl daemon-reload failed ({}): {}", + out.status, + String::from_utf8_lossy(&out.stderr).trim() + ); + } + Ok((String::new(), String::new())) +} + +/// `RestartMatrixDaemon` — restart the matrix daemon unit inside the +/// agent's container. +async fn restart_matrix_daemon(agent_name: &str) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + let machine = format!("--machine=h-{agent_name}"); + let unit = "hive-matrix-daemon.service"; + let out = Command::new("systemctl") + .args([&machine, "restart", unit]) + .output() + .await + .with_context(|| format!("systemctl restart {unit} in container h-{agent_name}"))?; + if !out.status.success() { + bail!( + "systemctl restart {unit} in h-{agent_name} exited {}: {}", + out.status, + String::from_utf8_lossy(&out.stderr).trim() + ); + } + Ok(( + String::from_utf8_lossy(&out.stdout).into_owned(), + String::from_utf8_lossy(&out.stderr).into_owned(), + )) +} + /// Shared helper for `WriteAgentForgeToken` and `WriteAgentMatrixToken`. /// Writes `content` to `AGENT_STATE_ROOT//state/`, /// chowns to the agent user (derived from the state dir's existing owner),