diff --git a/hive-priv/src/main.rs b/hive-priv/src/main.rs index 43b4fc8e..dbfb4ee6 100644 --- a/hive-priv/src/main.rs +++ b/hive-priv/src/main.rs @@ -350,19 +350,9 @@ async fn exec( match req { PrivRequest::StartContainer { ref name } => start_container(name).await, - PrivRequest::StopContainer { ref name } => { - validate_container_name(name)?; - stop_and_release(&container_system_name(name)).await - } + PrivRequest::StopContainer { ref name } => stop_container(name).await, - PrivRequest::KillContainer { ref name } => { - validate_container_name(name)?; - // nixos-container has no kill verb. Use machinectl to send SIGKILL - // to all processes in the container — the right semantics for a - // forced shutdown after a graceful stop has already been attempted. - let machine = container_system_name(name); - machinectl_run(&["kill", &machine, "--signal=SIGKILL"]).await - } + PrivRequest::KillContainer { ref name } => kill_container(name).await, PrivRequest::UpdateContainer { ref name, stream } => { container_flake_action("update", name, stream, writer).await @@ -372,20 +362,14 @@ async fn exec( container_flake_action("create", name, stream, writer).await } - PrivRequest::DestroyContainer { ref name } => { - validate_container_name(name)?; - container_run(&["destroy", &container_system_name(name)]).await - } + PrivRequest::DestroyContainer { ref name } => destroy_container(name).await, PrivRequest::ListContainers => container_run(&["list"]).await, PrivRequest::ReadContainerJournal { ref container, ref query, - } => { - validate_container_system_name(container)?; - read_container_journal(container, query).await - } + } => exec_read_container_journal(container, query).await, PrivRequest::WriteNspawnFlags { ref container, @@ -413,18 +397,12 @@ async fn exec( PrivRequest::SetAgentPaused { ref agent_name, paused, - } => { - validate_agent_name(agent_name)?; - set_agent_paused(agent_name, paused) - } + } => exec_set_agent_paused(agent_name, paused), PrivRequest::WriteAgentForgeToken { ref agent_name, ref token, - } => { - validate_agent_name(agent_name)?; - write_agent_state_file(agent_name, "forge-token", &format!("{token}\n")) - } + } => write_forge_token(agent_name, token), PrivRequest::WriteAgentMatrixToken { ref agent_name, @@ -436,10 +414,7 @@ async fn exec( PrivRequest::WriteAgentGithubToken { ref agent_name, ref token, - } => { - validate_agent_name(agent_name)?; - write_agent_state_file(agent_name, "github-token", &format!("{token}\n")) - } + } => write_github_token(agent_name, token), PrivRequest::WriteAgentExtraForgeAccount { ref agent_name, @@ -451,12 +426,7 @@ async fn exec( PrivRequest::DeleteAgentExtraForgeAccount { ref agent_name, ref label, - } => { - validate_agent_name(agent_name)?; - validate_name_chars(label)?; - delete_agent_state_file(agent_name, &format!("forge-{label}-token"))?; - delete_agent_state_file(agent_name, &format!("forge-{label}.json")) - } + } => delete_extra_forge_account(agent_name, label), PrivRequest::RestartMatrixDaemon { ref agent_name } => { restart_matrix_daemon(agent_name).await @@ -469,52 +439,37 @@ async fn exec( } PrivRequest::EnsureAgentSubvolume { ref agent_name } => { - validate_agent_name(agent_name)?; - ensure_agent_subvolume(agent_name).await + exec_ensure_agent_subvolume(agent_name).await } PrivRequest::DeleteAgentSubvolume { ref agent_name } => { - validate_agent_name(agent_name)?; - delete_agent_subvolume(agent_name).await + exec_delete_agent_subvolume(agent_name).await } PrivRequest::EnsureBtrfsQuota => ensure_btrfs_quota().await, PrivRequest::ReadSubvolumeUsage { ref agent_name } => { - validate_agent_name(agent_name)?; - read_subvolume_usage(agent_name).await + exec_read_subvolume_usage(agent_name).await } PrivRequest::SetSubvolumeQuota { ref agent_name, limit_bytes, - } => { - validate_agent_name(agent_name)?; - set_subvolume_quota(agent_name, limit_bytes).await - } + } => exec_set_subvolume_quota(agent_name, limit_bytes).await, PrivRequest::UpgradeAgentSubvolume { ref agent_name } => { - validate_agent_name(agent_name)?; - upgrade_agent_subvolume(agent_name).await + exec_upgrade_agent_subvolume(agent_name).await } PrivRequest::SnapshotAgentSubvolume { ref agent_name, ref snapshot_name, - } => { - validate_agent_name(agent_name)?; - validate_snapshot_name(snapshot_name)?; - snapshot_agent_subvolume(agent_name, snapshot_name).await - } + } => exec_snapshot_agent_subvolume(agent_name, snapshot_name).await, PrivRequest::DeleteAgentSnapshot { ref agent_name, ref snapshot_name, - } => { - validate_agent_name(agent_name)?; - validate_snapshot_name(snapshot_name)?; - delete_agent_snapshot(agent_name, snapshot_name).await - } + } => exec_delete_agent_snapshot(agent_name, snapshot_name).await, PrivRequest::SendAgentSnapshotToFile { ref agent_name, @@ -522,13 +477,7 @@ async fn exec( ref parent_snapshot_name, ref dest_file_name, } => { - validate_agent_name(agent_name)?; - validate_snapshot_name(snapshot_name)?; - if let Some(parent) = parent_snapshot_name { - validate_snapshot_name(parent)?; - } - validate_credential_name(dest_file_name)?; - send_agent_snapshot_to_file( + exec_send_agent_snapshot_to_file( agent_name, snapshot_name, parent_snapshot_name.as_deref(), @@ -542,17 +491,11 @@ async fn exec( ref snapshot_name, ref parent_snapshot_name, } => { - validate_agent_name(agent_name)?; - validate_snapshot_name(snapshot_name)?; - if let Some(parent) = parent_snapshot_name { - validate_snapshot_name(parent)?; - } - let dest = fd.context("no descriptor to stream into")?; - send_agent_snapshot_to_fd( + exec_send_agent_snapshot_to_fd( agent_name, snapshot_name, parent_snapshot_name.as_deref(), - dest, + fd, ) .await } @@ -638,6 +581,155 @@ fn write_extra_forge_account( Ok(res) } +/// `StopContainer`. +async fn stop_container(name: &str) -> Result<(String, String)> { + validate_container_name(name)?; + stop_and_release(&container_system_name(name)).await +} + +/// `KillContainer`: `nixos-container` has no kill verb, so this uses +/// `machinectl` to send `SIGKILL` to every process in the container — the +/// right semantics for a forced shutdown after a graceful stop has already +/// been attempted. +async fn kill_container(name: &str) -> Result<(String, String)> { + validate_container_name(name)?; + let machine = container_system_name(name); + machinectl_run(&["kill", &machine, "--signal=SIGKILL"]).await +} + +/// `DestroyContainer`. +async fn destroy_container(name: &str) -> Result<(String, String)> { + validate_container_name(name)?; + container_run(&["destroy", &container_system_name(name)]).await +} + +/// `ReadContainerJournal`. +async fn exec_read_container_journal( + container: &str, + query: &JournalQuery, +) -> Result<(String, String)> { + validate_container_system_name(container)?; + read_container_journal(container, query).await +} + +/// `SetAgentPaused`. +fn exec_set_agent_paused(agent_name: &str, paused: bool) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + set_agent_paused(agent_name, paused) +} + +/// `WriteAgentForgeToken`. +fn write_forge_token(agent_name: &str, token: &str) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + write_agent_state_file(agent_name, "forge-token", &format!("{token}\n")) +} + +/// `WriteAgentGithubToken`. +fn write_github_token(agent_name: &str, token: &str) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + write_agent_state_file(agent_name, "github-token", &format!("{token}\n")) +} + +/// `DeleteAgentExtraForgeAccount`. Missing files are not an error +/// (idempotent revoke). +fn delete_extra_forge_account(agent_name: &str, label: &str) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + validate_name_chars(label)?; + delete_agent_state_file(agent_name, &format!("forge-{label}-token"))?; + delete_agent_state_file(agent_name, &format!("forge-{label}.json")) +} + +/// `EnsureAgentSubvolume`. +async fn exec_ensure_agent_subvolume(agent_name: &str) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + ensure_agent_subvolume(agent_name).await +} + +/// `DeleteAgentSubvolume`. +async fn exec_delete_agent_subvolume(agent_name: &str) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + delete_agent_subvolume(agent_name).await +} + +/// `ReadSubvolumeUsage`. +async fn exec_read_subvolume_usage(agent_name: &str) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + read_subvolume_usage(agent_name).await +} + +/// `SetSubvolumeQuota`. +async fn exec_set_subvolume_quota( + agent_name: &str, + limit_bytes: Option, +) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + set_subvolume_quota(agent_name, limit_bytes).await +} + +/// `UpgradeAgentSubvolume`. +async fn exec_upgrade_agent_subvolume(agent_name: &str) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + upgrade_agent_subvolume(agent_name).await +} + +/// `SnapshotAgentSubvolume`. +async fn exec_snapshot_agent_subvolume( + agent_name: &str, + snapshot_name: &str, +) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + validate_snapshot_name(snapshot_name)?; + snapshot_agent_subvolume(agent_name, snapshot_name).await +} + +/// `DeleteAgentSnapshot`. +async fn exec_delete_agent_snapshot( + agent_name: &str, + snapshot_name: &str, +) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + validate_snapshot_name(snapshot_name)?; + delete_agent_snapshot(agent_name, snapshot_name).await +} + +/// `SendAgentSnapshotToFile`. +async fn exec_send_agent_snapshot_to_file( + agent_name: &str, + snapshot_name: &str, + parent_snapshot_name: Option<&str>, + dest_file_name: &str, +) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + validate_snapshot_name(snapshot_name)?; + if let Some(parent) = parent_snapshot_name { + validate_snapshot_name(parent)?; + } + validate_credential_name(dest_file_name)?; + send_agent_snapshot_to_file( + agent_name, + snapshot_name, + parent_snapshot_name, + dest_file_name, + ) + .await +} + +/// `SendAgentSnapshotToFd`. +async fn exec_send_agent_snapshot_to_fd( + agent_name: &str, + snapshot_name: &str, + parent_snapshot_name: Option<&str>, + fd: Option, +) -> Result<(String, String)> { + validate_agent_name(agent_name)?; + validate_snapshot_name(snapshot_name)?; + if let Some(parent) = parent_snapshot_name { + validate_snapshot_name(parent)?; + } + let dest = fd.context("no descriptor to stream into")?; + send_agent_snapshot_to_fd(agent_name, snapshot_name, parent_snapshot_name, dest).await +} + /// Shared body for `CreateContainer` / `UpdateContainer`: validate the /// name, build the toplevel ourselves, and run /// `nixos-container … --system-path ` (streaming line