From 12b2723646a201b205cd4dc0535c17fbd7df7933 Mon Sep 17 00:00:00 2001 From: Aly Raffauf Date: Wed, 5 Aug 2026 19:26:35 -0400 Subject: [PATCH] add new ui; split nix module --- Cargo.lock | 107 +++++++++++ Cargo.toml | 3 + src/main.rs | 167 +++-------------- src/models.rs | 8 - src/nix.rs | 426 +------------------------------------------- src/nix/deploy.rs | 311 ++++++++++++++++++++++++++++++++ src/nix/evaluate.rs | 203 +++++++++++++++++++++ src/process.rs | 144 +++++++++------ src/ssh.rs | 17 +- src/ui.rs | 263 +++++++++++++++++++++++++++ src/workflow.rs | 163 +++++++++++++++++ 11 files changed, 1174 insertions(+), 638 deletions(-) create mode 100644 src/nix/deploy.rs create mode 100644 src/nix/evaluate.rs create mode 100644 src/ui.rs create mode 100644 src/workflow.rs diff --git a/Cargo.lock b/Cargo.lock index a70e07d..2feb6ce 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -87,8 +87,10 @@ dependencies = [ "clap", "env_logger", "futures", + "indicatif", "log", "openssh", + "owo-colors", "serde", "serde_json", "tokio", @@ -159,6 +161,18 @@ version = "1.0.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1d07550c9036bf2ae0c684c4297d503f838287c83c53686d05370d0e139ae570" +[[package]] +name = "console" +version = "0.16.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4fe5f465a4f6fee88fad41b85d990f84c835335e85b5d9e6e63e0d06d28cba7c" +dependencies = [ + "encode_unicode", + "libc", + "unicode-width", + "windows-sys", +] + [[package]] name = "defmt" version = "1.1.1" @@ -190,6 +204,12 @@ dependencies = [ "thiserror", ] +[[package]] +name = "encode_unicode" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "34aa73646ffb006b8f5147f3dc182bd4bcb190227ce861fc4a4844bf8e3cb2c0" + [[package]] name = "env_filter" version = "2.0.0" @@ -334,6 +354,42 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" +[[package]] +name = "hermit-abi" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fc0fef456e4baa96da950455cd02c081ca953b141298e41db3fc7e36b1da849c" + +[[package]] +name = "indicatif" +version = "0.18.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9433806cd6b4ec1aba79c021c7e4c58fb4c3b9977c085062e611ac929998fb0c" +dependencies = [ + "console", + "portable-atomic", + "unicode-width", + "unit-prefix", + "web-time", +] + +[[package]] +name = "is-terminal" +version = "0.4.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46" +dependencies = [ + "hermit-abi", + "libc", + "windows-sys", +] + +[[package]] +name = "is_ci" +version = "1.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7655c9839580ee829dfacba1d1278c2b7883e50a277ff7541299489d6bdfdc45" + [[package]] name = "is_terminal_polyfill" version = "1.70.2" @@ -454,6 +510,16 @@ dependencies = [ "tokio", ] +[[package]] +name = "owo-colors" +version = "4.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d211803b9b6b570f68772237e415a029d5a50c65d382910b879fb19d3271f94d" +dependencies = [ + "supports-color 2.1.0", + "supports-color 3.0.2", +] + [[package]] name = "pin-project-lite" version = "0.2.17" @@ -628,6 +694,25 @@ version = "0.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f" +[[package]] +name = "supports-color" +version = "2.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d6398cde53adc3c4557306a96ce67b302968513830a77a95b2b17305d9719a89" +dependencies = [ + "is-terminal", + "is_ci", +] + +[[package]] +name = "supports-color" +version = "3.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c64fc7232dd8d2e4ac5ce4ef302b1d81e0b80d055b9d77c7c4f51f6aa4c867d6" +dependencies = [ + "is_ci", +] + [[package]] name = "syn" version = "2.0.119" @@ -716,6 +801,18 @@ version = "1.0.24" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" +[[package]] +name = "unicode-width" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b4ac048d71ede7ee76d585517add45da530660ef4390e49b098733c6e897f254" + +[[package]] +name = "unit-prefix" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "81e544489bf3d8ef66c953931f56617f423cd4b5494be343d9b9d3dda037b9a3" + [[package]] name = "utf8parse" version = "0.2.2" @@ -784,6 +881,16 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "web-time" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a6580f308b1fad9207618087a65c04e7a10bc77e02c8e84e9b00dd4b12fa0bb" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + [[package]] name = "windows-link" version = "0.2.1" diff --git a/Cargo.toml b/Cargo.toml index 59edd2f..022893f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -14,6 +14,7 @@ tokio = { version = "1", features = [ "rt-multi-thread", "macros", "process", + "io-util", "time", ] } clap = { version = "4", features = ["derive", "env"] } @@ -21,7 +22,9 @@ serde = { version = "1", features = ["derive"] } serde_json = "1" anyhow = "1" futures = "0.3" +indicatif = "0.18" log = "0.4" env_logger = "0.11" openssh = "0.11" +owo-colors = { version = "4", features = ["supports-colors"] } uuid = { version = "1", features = ["v4"] } diff --git a/src/main.rs b/src/main.rs index 945e4c9..a9dd0b9 100644 --- a/src/main.rs +++ b/src/main.rs @@ -4,13 +4,11 @@ mod nix; mod op; mod process; mod ssh; +mod ui; +mod workflow; -use std::collections::HashSet; - -use futures::stream::{self, StreamExt}; - -use crate::cli::{Command, CommonArgs}; -use crate::ssh::HostKeyPolicy; +use crate::cli::Command; +use crate::nix::EvaluationProgress; #[tokio::main] async fn main() -> anyhow::Result<()> { @@ -23,148 +21,41 @@ async fn main() -> anyhow::Result<()> { .format_target(false) .init(); - let start = std::time::Instant::now(); + let ui = ui::Ui::new(); + + let mut evaluation = ui.start_section("Evaluating"); + let eval_result = nix::eval_deployments(&args.flake, |progress| match progress { + EvaluationProgress::JobEvaluated { name } => { + ui.add_job(&mut evaluation, &name, &name); + ui.start_job(&evaluation, &name, &name); + } + EvaluationProgress::MetadataResolved { name } => { + ui.finish_job_success(ui.job_progress(&evaluation, &name), &name); + } + }) + .await; - log::info!("Reading nodes from {}", args.flake); - let (jobs, _debug) = nix::eval_deployments(&args.flake).await?; + let jobs = match eval_result { + Ok(result) => { + ui.finish_section_success(evaluation); + result + } + Err(error) => { + ui.finish_section_failure(evaluation); + return Err(error); + } + }; match args.command { Command::List => { - print_nodes(&jobs); + ui.print_nodes(&jobs); return Ok(()); } cmd @ (Command::Switch(_) | Command::Boot(_)) => { let (op, common) = cmd.into_deploy().expect("deploy command"); - run_deploy(op, common, jobs, start).await?; - } - } - - Ok(()) -} - -/// Apply the `--skip` and `nodes` filters to the evaluated job map. -fn filter_jobs( - mut jobs: std::collections::HashMap, - common: &CommonArgs, -) -> anyhow::Result> { - let all_names: HashSet<&String> = jobs.keys().collect(); - - for s in &common.skip { - if !all_names.contains(&s) { - log::warn!("ignoring unknown node '{s}' in --skip"); + workflow::run_deploy(op, common, jobs, &ui).await?; } } - if !common.nodes.is_empty() { - for n in &common.nodes { - if !all_names.contains(&n) { - anyhow::bail!("node '{n}' not found in flake"); - } - } - } - - jobs.retain(|name, _| !common.skip.iter().any(|s| s == name)); - if !common.nodes.is_empty() { - jobs.retain(|name, _| common.nodes.iter().any(|n| n == name)); - } - - Ok(jobs) -} - -/// Print the resolved node list with their system/user/host. -fn print_nodes(jobs: &std::collections::HashMap) { - log::info!("Found {} node(s):", jobs.len()); - for (name, spec) in jobs { - log::info!( - " {name} -> system={} user={} host={}", - spec.system, - spec.user, - spec.hostname, - ); - } -} - -/// Filter, validate, build, and deploy the jobs for the given operation. -async fn run_deploy( - op: crate::op::Operation, - common: CommonArgs, - jobs: std::collections::HashMap, - start: std::time::Instant, -) -> anyhow::Result<()> { - let jobs = filter_jobs(jobs, &common)?; - print_nodes(&jobs); - - op::validate(&jobs, op)?; - - log::info!("Operation {op} is valid for all nodes"); - - let host_key_policy = if common.accept_new_host_keys { - HostKeyPolicy::AcceptNew - } else { - HostKeyPolicy::Strict - }; - - // Build each node's closure locally or on a remote builder. - log::info!("Building {} output(s)...", jobs.len()); - let mut outs: std::collections::HashMap = - std::collections::HashMap::with_capacity(jobs.len()); - for (name, spec) in &jobs { - let (out, _debug) = nix::build_closure(spec, &common.build_host, host_key_policy).await?; - log::info!(" ✔ {name} ({})", spec.system); - outs.insert(name.clone(), out); - } - - log::info!("Deploying {} output(s)...", jobs.len()); - - // One future per node, run concurrently. - let tasks: Vec<_> = jobs - .iter() - .map(|(name, spec)| { - let out = outs[name].clone(); - let name = name.clone(); - let spec = spec.clone(); - async move { - let result = nix::deploy_closure(&spec, &out, op, host_key_policy).await; - (name, spec, result) - } - }) - .collect(); - - let results = stream::iter(tasks) - .buffer_unordered(common.parallel) - .collect::>() - .await; - - let errors: Vec<_> = results - .into_iter() - .filter_map(|(name, spec, result)| match result { - Ok(_) => { - log::info!( - " ✔ {name} ({}) -> {}@{}", - spec.system, - spec.user, - spec.hostname - ); - None - } - Err(e) => { - let target = spec.target(); - log::warn!("Failed to deploy to {target}: {e}"); - Some(e) - } - }) - .collect(); - - let duration = start.elapsed(); - - if !errors.is_empty() { - log::info!( - "Deployment failed with {} error(s) ({duration:?})", - errors.len() - ); - anyhow::bail!("deployment failed"); - } - - log::info!("Completed successfully in {duration:?}"); Ok(()) } diff --git a/src/models.rs b/src/models.rs index c63e2e4..7b8a0f3 100644 --- a/src/models.rs +++ b/src/models.rs @@ -27,14 +27,6 @@ impl fmt::Display for SystemType { } } -/// Captured details of a subprocess invocation, for debug logging. -#[derive(Debug, Clone, Default)] -pub struct DebugInfo { - pub command: String, - pub std_out: String, - pub std_err: String, -} - #[derive(Debug, Clone)] pub struct JobSpec { pub hostname: String, diff --git a/src/nix.rs b/src/nix.rs index 0c8901c..219db86 100644 --- a/src/nix.rs +++ b/src/nix.rs @@ -1,423 +1,5 @@ -use std::collections::HashMap; -use std::path::PathBuf; -use std::time::Duration; +mod deploy; +mod evaluate; -use anyhow::{Context, Result}; -use tokio::time::{sleep, Instant}; -use uuid::Uuid; - -use crate::models::{BuildResult, DebugInfo, JobSpec, NixEvalJobsResult, SystemType}; -use crate::op::Operation; -use crate::process::{get_config_attr, run, run_json, run_json_with_env, run_with_env}; -use crate::ssh::{run as run_ssh, HostKeyPolicy}; - -/// Evaluate the flake's `blzrd.nodes` and return enriched `JobSpec`s. -pub async fn eval_deployments(cfg: &str) -> Result<(HashMap, Vec)> { - let (text, mut debug_infos) = run_eval_jobs(cfg).await?; - let results = parse_json_lines(&text)?; - let basics = collect_basics(&results)?; - let jobs = build_job_specs(cfg, &results, &basics, &mut debug_infos).await?; - Ok((jobs, debug_infos)) -} - -/// Run `nix-eval-jobs` on the flake's `blzrd.nodes` and return it + `DebugInfo`. -async fn run_eval_jobs(cfg: &str) -> Result<(String, Vec)> { - let flake_reference = format!("{cfg}#blzrd.nodes"); - - let cache_home = std::env::var("XDG_CACHE_HOME") - .ok() - .filter(|s| !s.is_empty()) - .unwrap_or_else(|| { - let home = - std::env::var("HOME").expect("HOME must be set if XDG_CACHE_HOME is not set"); - format!("{home}/.cache") - }); - let gc_roots_dir: PathBuf = PathBuf::from(cache_home).join("blzrd"); - - std::fs::create_dir_all(&gc_roots_dir) - .with_context(|| format!("failed to create gc-roots dir {}", gc_roots_dir.display()))?; - - let (out, debug) = run( - "nix-eval-jobs", - &[ - "--gc-roots-dir", - &gc_roots_dir.to_string_lossy(), - "--force-recurse", - "--flake", - &flake_reference, - ], - ) - .await?; - - Ok((String::from_utf8_lossy(&out).into_owned(), vec![debug])) -} - -/// Parse the JSON-Lines text emitted by `nix-eval-jobs` into typed records. -fn parse_json_lines(text: &str) -> Result> { - text.lines() - .filter(|line| !line.is_empty()) - .map(serde_json::from_str) - .collect::, _>>() - .context("failed to parse JSON from nix-eval-jobs") -} - -/// First pass: extract the basic `(output, drv_path)` pair for each job. -fn collect_basics(results: &[NixEvalJobsResult]) -> Result> { - let mut basics: HashMap = HashMap::with_capacity(results.len()); - for result in results { - // Error records from nix-eval-jobs carry `attr` and `error`. - if let Some(err) = &result.error { - anyhow::bail!("job {}: {err}", result.attr); - } - - let attr_path = result - .attr_path - .as_ref() - .with_context(|| format!("job {}: missing attr_path", result.attr))?; - if attr_path.len() < 2 { - anyhow::bail!("job {}: malformed attr_path: {:?}", result.attr, attr_path); - } - let job_name = attr_path[0].clone(); - - let outputs = result - .outputs - .as_ref() - .with_context(|| format!("job {job_name}: missing outputs"))?; - let output_path = outputs - .get("out") - .with_context(|| format!("job {job_name}: missing 'out' output"))?; - - let drv_path = result.drv_path.clone().unwrap_or_default(); - - basics.insert(job_name, (output_path.clone(), drv_path)); - } - Ok(basics) -} - -/// Second pass: enrich each job's basic data with hostname, user, and system type. -async fn build_job_specs( - cfg: &str, - results: &[NixEvalJobsResult], - basics: &HashMap, - debug_infos: &mut Vec, -) -> Result> { - let mut jobs: HashMap = HashMap::with_capacity(basics.len()); - - for name in basics.keys() { - let (output, drv_path) = basics.get(name).cloned().unwrap_or_default(); - - // hostname: defaults to the job name, override if flake provides one. - let mut hostname = name.clone(); - if let Ok((h, debug)) = get_config_attr(cfg, name, "hostname").await { - debug_infos.push(debug); - if !h.is_empty() { - hostname = h; - } - } - - let mut user = String::new(); - if let Ok((u, debug)) = get_config_attr(cfg, name, "user").await { - debug_infos.push(debug); - user = u; - } - - // system type: from flake, or inferred from the build system string. - let mut type_str = String::new(); - if let Ok((t, debug)) = get_config_attr(cfg, name, "type").await { - debug_infos.push(debug); - type_str = t; - } - - if type_str.is_empty() { - let system = results - .iter() - .find(|r| r.attr_path.as_deref().and_then(|p| p.first()) == Some(name)) - .and_then(|r| r.system.as_deref()) - .unwrap_or_default(); - - type_str = if system.contains("darwin") { - "darwin".to_string() - } else if system.contains("linux") { - "nixos".to_string() - } else { - anyhow::bail!("job {name}: unknown system type: {system}"); - }; - } - - let system = - SystemType::parse(&type_str).map_err(|e| anyhow::anyhow!("job {name}: {e}"))?; - - if output.is_empty() { - anyhow::bail!("job {name}: missing output path"); - } - if user.is_empty() { - anyhow::bail!("job {name}: missing user"); - } - jobs.insert( - name.clone(), - JobSpec { - hostname, - system, - user, - drv_path, - }, - ); - } - - Ok(jobs) -} - -/// Build a job's derivation and return the `out` store path. -pub async fn build_closure( - spec: &JobSpec, - build_host: &str, - host_key_policy: HostKeyPolicy, -) -> Result<(String, DebugInfo)> { - let drv = format!("{}^*", spec.drv_path); - - if build_host == "localhost" { - let (results, debug): (Vec, _) = - run_json("nix", &["build", "--no-link", "--json", &drv]).await?; - let out = results - .into_iter() - .next() - .and_then(|r| r.outputs.get("out").cloned()) - .context("build result missing 'out' output")?; - return Ok((out, debug)); - } - - // Remote builder branch. - let store = format!("ssh-ng://{build_host}"); - let nix_ssh_options = nix_ssh_options(host_key_policy); - let nix_env = [("NIX_SSHOPTS", nix_ssh_options.as_str())]; - - // 1. Copy the derivation to the builder. - let (_out, _debug) = run_with_env("nix", &["copy", "--to", &store, &spec.drv_path], &nix_env) - .await - .with_context(|| format!("copy to {build_host}"))?; - - // 2. Build on the builder. - let (results, debug): (Vec, _) = run_json_with_env( - "nix", - &["build", "--no-link", "--json", "--store", &store, &drv], - &nix_env, - ) - .await - .with_context(|| format!("build on {build_host}"))?; - - let out = results - .into_iter() - .next() - .and_then(|r| r.outputs.get("out").cloned()) - .context("build result missing 'out' output")?; - - // 3. Copy the out path back. - let (_out, _debug2) = run_with_env( - "nix", - &["copy", "--from", &store, &out, "--no-check-sigs"], - &nix_env, - ) - .await - .with_context(|| format!("copy from {build_host}"))?; - - Ok((out, debug)) -} - -/// Poll a remote `blzrd-activate-*` transient systemd unit until it reaches a -/// terminal state, then return. Exponential backoff 5 -> 60 seconds, 5-minute deadline. -async fn poll_activation(target: &str, unit: &str, host_key_policy: HostKeyPolicy) -> Result<()> { - let mut sleep_val = Duration::from_secs(5); - let deadline = Instant::now() + Duration::from_mins(5); - - loop { - if Instant::now() >= deadline { - anyhow::bail!("activation on {target}: timed out waiting for {unit} unit"); - } - - let result = run_ssh(target, "systemctl", &["show", unit], host_key_policy) - .await - .ok(); - - let mut sub_state: Option = None; - let mut exec_status: Option = None; - if let Some((out, debug_info)) = result { - let text = String::from_utf8_lossy(&out); - log::debug!( - "{unit} on {target}:\n{}stderr: {}", - text.trim(), - debug_info.std_err.trim() - ); - for line in text.lines() { - if let Some(v) = line.strip_prefix("SubState=") { - sub_state = Some(v.trim().to_string()); - } else if let Some(v) = line.strip_prefix("ExecMainStatus=") { - exec_status = v.trim().parse::().ok(); - } - } - } - - // With --remain-after-exit, the unit stays active but transitions to - // SubState=exited once the main process terminates. - match (sub_state.as_deref(), exec_status.unwrap_or(0)) { - (Some("exited"), 0) => break, - (Some("exited"), code) => { - anyhow::bail!( - "activation on {target}: switch-to-configuration exited with code {code}" - ); - } - (Some("failed"), _) => { - anyhow::bail!( - "activation on {target}: unit failed (ExecMainStatus={})", - exec_status.unwrap_or(-1) - ); - } - _ => { - sleep(sleep_val).await; - sleep_val = (sleep_val * 2).min(Duration::from_mins(1)); - } - } - - if sleep_val > Duration::from_mins(1) { - anyhow::bail!("activation on {target}: gave up after backoff cap exceeded"); - } - } - - Ok(()) -} - -pub async fn deploy_closure( - spec: &JobSpec, - out_path: &str, - op: Operation, - host_key_policy: HostKeyPolicy, -) -> Result { - let target = spec.target(); - let path = out_path.to_string(); - - let unit = format!("blzrd-activate-{}", Uuid::new_v4().simple()); - - let mut cmds: Vec> = Vec::new(); - - let sys = spec.system; - - match (sys, op) { - (SystemType::Darwin, Operation::Switch) => { - cmds.push(root_command( - spec, - "/run/current-system/sw/bin/nix-env".into(), - vec![ - "-p".into(), - "/nix/var/nix/profiles/system".into(), - "--set".into(), - path.clone(), - ], - )); - cmds.push(root_command(spec, format!("{path}/activate"), vec![])); - } - - (SystemType::Darwin, Operation::Boot) => { - anyhow::bail!( - "job {}: 'boot' is not a valid darwin operation", - spec.hostname - ); - } - - (SystemType::Nixos, Operation::Switch) => { - cmds.push(root_command( - spec, - "/run/current-system/sw/bin/nix-env".into(), - vec![ - "-p".into(), - "/nix/var/nix/profiles/system".into(), - "--set".into(), - path.clone(), - ], - )); - cmds.push(root_command( - spec, - "/run/current-system/sw/bin/systemd-run".into(), - vec![ - "--unit".into(), - unit.clone(), - "--remain-after-exit".into(), - "--no-block".into(), - "--".into(), - format!("{path}/bin/switch-to-configuration"), - op.to_string(), - ], - )); - } - - (SystemType::Nixos, Operation::Boot) => { - cmds.push(root_command( - spec, - "/run/current-system/sw/bin/nix-env".into(), - vec![ - "-p".into(), - "/nix/var/nix/profiles/system".into(), - "--set".into(), - path.clone(), - ], - )); - cmds.push(root_command( - spec, - format!("{path}/bin/switch-to-configuration"), - vec![op.to_string()], - )); - } - } - - // 1. Copy the closure to the target. - let nix_ssh_options = nix_ssh_options(host_key_policy); - let nix_env = [("NIX_SSHOPTS", nix_ssh_options.as_str())]; - let (_out, debug) = run_with_env( - "nix", - &[ - "copy", - "--to", - &format!("ssh-ng://{target}"), - &path, - "--no-check-sigs", - ], - &nix_env, - ) - .await - .with_context(|| format!("copy to {target}"))?; - - // 2. Run each activation command in order. - for cmd in &cmds { - let args: Vec<&str> = cmd[1..].iter().map(String::as_str).collect(); - let (_out, _d) = run_ssh(&target, &cmd[0], &args, host_key_policy) - .await - .with_context(|| format!("activation on {target}"))?; - } - - if matches!((sys, op), (SystemType::Nixos, Operation::Switch)) { - poll_activation(&target, &unit, host_key_policy) - .await - .with_context(|| format!("activation on {target}"))?; - } - - Ok(debug) -} - -fn nix_ssh_options(host_key_policy: HostKeyPolicy) -> String { - let host_key_option = host_key_policy.nix_ssh_option(); - match std::env::var("NIX_SSHOPTS") { - Ok(existing) if !existing.trim().is_empty() => format!("{existing} {host_key_option}"), - _ => host_key_option.to_string(), - } -} - -fn root_command(spec: &JobSpec, program: String, args: Vec) -> Vec { - let mut command = Vec::with_capacity(args.len() + 3); - - if spec.user != "root" { - command.push("sudo".into()); - command.push("-n".into()); - } - - command.push(program); - command.extend(args); - command -} +pub use deploy::{build_closure, deploy_closure}; +pub use evaluate::{eval_deployments, EvaluationProgress}; diff --git a/src/nix/deploy.rs b/src/nix/deploy.rs new file mode 100644 index 0000000..3bb31ca --- /dev/null +++ b/src/nix/deploy.rs @@ -0,0 +1,311 @@ +use std::time::Duration; + +use anyhow::{Context, Result}; +use tokio::time::{sleep, Instant}; +use uuid::Uuid; + +use crate::models::{BuildResult, JobSpec, SystemType}; +use crate::op::Operation; +use crate::process::{run_json, run_json_with_env, run_with_env}; +use crate::ssh::{run as run_ssh, HostKeyPolicy}; + +/// Build a job's derivation and return the `out` store path. +pub async fn build_closure( + spec: &JobSpec, + build_host: &str, + host_key_policy: HostKeyPolicy, +) -> Result { + let derivation = format!("{}^*", spec.drv_path); + + if build_host == "localhost" { + return build_local_closure(&derivation).await; + } + + build_remote_closure(&derivation, build_host, host_key_policy).await +} + +async fn build_local_closure(derivation: &str) -> Result { + let results: Vec = + run_json("nix", &["build", "--no-link", "--json", derivation]).await?; + build_output_path(results) +} + +async fn build_remote_closure( + derivation: &str, + build_host: &str, + host_key_policy: HostKeyPolicy, +) -> Result { + let store = format!("ssh-ng://{build_host}"); + let nix_ssh_options = nix_ssh_options(host_key_policy); + let nix_env = [("NIX_SSHOPTS", nix_ssh_options.as_str())]; + + run_with_env("nix", &["copy", "--to", &store, derivation], &nix_env) + .await + .with_context(|| format!("copy to {build_host}"))?; + + let results: Vec = run_json_with_env( + "nix", + &[ + "build", + "--no-link", + "--json", + "--store", + &store, + derivation, + ], + &nix_env, + ) + .await + .with_context(|| format!("build on {build_host}"))?; + let output = build_output_path(results)?; + + run_with_env( + "nix", + &["copy", "--from", &store, &output, "--no-check-sigs"], + &nix_env, + ) + .await + .with_context(|| format!("copy from {build_host}"))?; + + Ok(output) +} + +fn build_output_path(results: Vec) -> Result { + results + .into_iter() + .next() + .and_then(|result| result.outputs.get("out").cloned()) + .context("build result missing 'out' output") +} + +pub async fn deploy_closure( + spec: &JobSpec, + output_path: &str, + operation: Operation, + host_key_policy: HostKeyPolicy, +) -> Result<()> { + let target = spec.target(); + copy_closure(output_path, &target, host_key_policy).await?; + + let activation_unit = format!("blzrd-activate-{}", Uuid::new_v4().simple()); + let commands = activation_commands(spec, output_path, operation, &activation_unit)?; + run_activation_commands(&target, &commands, host_key_policy).await?; + + if matches!( + (spec.system, operation), + (SystemType::Nixos, Operation::Switch) + ) { + poll_activation(&target, &activation_unit, host_key_policy) + .await + .with_context(|| format!("activation on {target}"))?; + } + + Ok(()) +} + +async fn copy_closure( + output_path: &str, + target: &str, + host_key_policy: HostKeyPolicy, +) -> Result<()> { + let nix_ssh_options = nix_ssh_options(host_key_policy); + let nix_env = [("NIX_SSHOPTS", nix_ssh_options.as_str())]; + let remote_store = format!("ssh-ng://{target}"); + + run_with_env( + "nix", + &[ + "copy", + "--to", + &remote_store, + output_path, + "--no-check-sigs", + ], + &nix_env, + ) + .await + .with_context(|| format!("copy to {target}"))?; + + Ok(()) +} + +async fn run_activation_commands( + target: &str, + commands: &[RemoteCommand], + host_key_policy: HostKeyPolicy, +) -> Result<()> { + for command in commands { + let args: Vec<_> = command.arguments.iter().map(String::as_str).collect(); + run_ssh(target, &command.program, &args, host_key_policy) + .await + .with_context(|| format!("activation on {target}"))?; + } + + Ok(()) +} + +fn activation_commands( + spec: &JobSpec, + output_path: &str, + operation: Operation, + activation_unit: &str, +) -> Result> { + let update_profile = root_command( + spec, + "/run/current-system/sw/bin/nix-env", + vec![ + "-p".to_owned(), + "/nix/var/nix/profiles/system".to_owned(), + "--set".to_owned(), + output_path.to_owned(), + ], + ); + + let activate = match (spec.system, operation) { + (SystemType::Darwin, Operation::Switch) => { + root_command(spec, &format!("{output_path}/activate"), Vec::new()) + } + (SystemType::Darwin, Operation::Boot) => { + anyhow::bail!( + "job {}: 'boot' is not a valid darwin operation", + spec.hostname + ); + } + (SystemType::Nixos, Operation::Switch) => root_command( + spec, + "/run/current-system/sw/bin/systemd-run", + vec![ + "--unit".to_owned(), + activation_unit.to_owned(), + "--remain-after-exit".to_owned(), + "--no-block".to_owned(), + "--".to_owned(), + format!("{output_path}/bin/switch-to-configuration"), + operation.to_string(), + ], + ), + (SystemType::Nixos, Operation::Boot) => root_command( + spec, + &format!("{output_path}/bin/switch-to-configuration"), + vec![operation.to_string()], + ), + }; + + Ok(vec![update_profile, activate]) +} + +async fn poll_activation(target: &str, unit: &str, host_key_policy: HostKeyPolicy) -> Result<()> { + let mut sleep_duration = Duration::from_secs(5); + let deadline = Instant::now() + Duration::from_mins(5); + + loop { + if Instant::now() >= deadline { + anyhow::bail!("activation on {target}: timed out waiting for {unit} unit"); + } + + let result = run_ssh(target, "systemctl", &["show", unit], host_key_policy) + .await + .ok(); + let (sub_state, exit_status) = result + .as_deref() + .map(parse_systemd_state) + .unwrap_or_default(); + + match (sub_state.as_deref(), exit_status.unwrap_or(0)) { + (Some("exited"), 0) => break, + (Some("exited"), status) => { + anyhow::bail!( + "activation on {target}: switch-to-configuration exited with code {status}" + ); + } + (Some("failed"), _) => { + anyhow::bail!( + "activation on {target}: unit failed (ExecMainStatus={})", + exit_status.unwrap_or(-1) + ); + } + _ => { + sleep(sleep_duration).await; + sleep_duration = (sleep_duration * 2).min(Duration::from_mins(1)); + } + } + } + + Ok(()) +} + +fn parse_systemd_state(output: &[u8]) -> (Option, Option) { + let mut sub_state = None; + let mut exit_status = None; + + for line in String::from_utf8_lossy(output).lines() { + if let Some(value) = line.strip_prefix("SubState=") { + sub_state = Some(value.trim().to_owned()); + } else if let Some(value) = line.strip_prefix("ExecMainStatus=") { + exit_status = value.trim().parse().ok(); + } + } + + (sub_state, exit_status) +} + +fn nix_ssh_options(host_key_policy: HostKeyPolicy) -> String { + let host_key_option = host_key_policy.nix_ssh_option(); + match std::env::var("NIX_SSHOPTS") { + Ok(existing) if !existing.trim().is_empty() => format!("{existing} {host_key_option}"), + _ => host_key_option.to_owned(), + } +} + +struct RemoteCommand { + program: String, + arguments: Vec, +} + +fn root_command(spec: &JobSpec, program: &str, arguments: Vec) -> RemoteCommand { + if spec.user == "root" { + return RemoteCommand { + program: program.to_owned(), + arguments, + }; + } + + let mut sudo_arguments = Vec::with_capacity(arguments.len() + 2); + sudo_arguments.push("-n".to_owned()); + sudo_arguments.push(program.to_owned()); + sudo_arguments.extend(arguments); + RemoteCommand { + program: "sudo".to_owned(), + arguments: sudo_arguments, + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn job(user: &str) -> JobSpec { + JobSpec { + hostname: "host".to_owned(), + system: SystemType::Nixos, + user: user.to_owned(), + drv_path: "/nix/store/example.drv".to_owned(), + } + } + + #[test] + fn root_command_runs_directly_for_root() { + let command = root_command(&job("root"), "activate", vec!["switch".to_owned()]); + + assert_eq!(command.program, "activate"); + assert_eq!(command.arguments, ["switch"]); + } + + #[test] + fn root_command_uses_noninteractive_sudo_for_other_users() { + let command = root_command(&job("deploy"), "activate", vec!["switch".to_owned()]); + + assert_eq!(command.program, "sudo"); + assert_eq!(command.arguments, ["-n", "activate", "switch"]); + } +} diff --git a/src/nix/evaluate.rs b/src/nix/evaluate.rs new file mode 100644 index 0000000..0150de7 --- /dev/null +++ b/src/nix/evaluate.rs @@ -0,0 +1,203 @@ +use std::collections::HashMap; +use std::path::PathBuf; + +use anyhow::{Context, Result}; + +use crate::models::{JobSpec, NixEvalJobsResult, SystemType}; +use crate::process::{get_config_attr, run_json_lines}; + +#[derive(Debug, Clone)] +pub enum EvaluationProgress { + JobEvaluated { name: String }, + MetadataResolved { name: String }, +} + +/// Evaluate the flake's `blzrd.nodes` and return enriched `JobSpec`s. +pub async fn eval_deployments( + cfg: &str, + mut on_progress: OnProgress, +) -> Result> +where + OnProgress: FnMut(EvaluationProgress), +{ + let results = run_eval_jobs(cfg, &mut on_progress).await?; + let basics = collect_basics(&results)?; + build_job_specs(cfg, &results, &basics, &mut on_progress).await +} + +/// Run `nix-eval-jobs` on the flake's `blzrd.nodes`. +async fn run_eval_jobs( + cfg: &str, + on_progress: &mut OnProgress, +) -> Result> +where + OnProgress: FnMut(EvaluationProgress), +{ + let flake_reference = format!("{cfg}#blzrd.nodes"); + let gc_roots_dir = gc_roots_dir()?; + + run_json_lines( + "nix-eval-jobs", + &[ + "--gc-roots-dir", + &gc_roots_dir.to_string_lossy(), + "--force-recurse", + "--flake", + &flake_reference, + ], + |result: &NixEvalJobsResult| { + let name = result + .attr_path + .as_ref() + .and_then(|path| path.first()) + .cloned() + .unwrap_or_else(|| result.attr.clone()); + on_progress(EvaluationProgress::JobEvaluated { name }); + }, + ) + .await +} + +fn gc_roots_dir() -> Result { + let cache_home = std::env::var("XDG_CACHE_HOME") + .ok() + .filter(|path| !path.is_empty()) + .unwrap_or_else(|| { + let home = + std::env::var("HOME").expect("HOME must be set if XDG_CACHE_HOME is not set"); + format!("{home}/.cache") + }); + let gc_roots_dir = PathBuf::from(cache_home).join("blzrd"); + + std::fs::create_dir_all(&gc_roots_dir) + .with_context(|| format!("failed to create gc-roots dir {}", gc_roots_dir.display()))?; + + Ok(gc_roots_dir) +} + +/// First pass: extract the basic `(output, drv_path)` pair for each job. +fn collect_basics(results: &[NixEvalJobsResult]) -> Result> { + let mut basics = HashMap::with_capacity(results.len()); + for result in results { + if let Some(error) = &result.error { + anyhow::bail!("job {}: {error}", result.attr); + } + + let attr_path = result + .attr_path + .as_ref() + .with_context(|| format!("job {}: missing attr_path", result.attr))?; + if attr_path.len() < 2 { + anyhow::bail!("job {}: malformed attr_path: {attr_path:?}", result.attr); + } + let job_name = attr_path[0].clone(); + + let outputs = result + .outputs + .as_ref() + .with_context(|| format!("job {job_name}: missing outputs"))?; + let output_path = outputs + .get("out") + .with_context(|| format!("job {job_name}: missing 'out' output"))?; + + basics.insert( + job_name, + ( + output_path.clone(), + result.drv_path.clone().unwrap_or_default(), + ), + ); + } + Ok(basics) +} + +/// Second pass: enrich each job's basic data with hostname, user, and system type. +async fn build_job_specs( + cfg: &str, + results: &[NixEvalJobsResult], + basics: &HashMap, + on_progress: &mut impl FnMut(EvaluationProgress), +) -> Result> { + let mut jobs = HashMap::with_capacity(basics.len()); + let mut names: Vec<_> = basics.keys().cloned().collect(); + names.sort(); + + for name in names { + let (output, drv_path) = basics.get(&name).cloned().unwrap_or_default(); + let hostname = configured_hostname(cfg, &name).await; + let user = configured_value(cfg, &name, "user").await; + let system = configured_system(cfg, &name, results).await?; + + if output.is_empty() { + anyhow::bail!("job {name}: missing output path"); + } + if user.is_empty() { + anyhow::bail!("job {name}: missing user"); + } + + jobs.insert( + name.clone(), + JobSpec { + hostname, + system, + user, + drv_path, + }, + ); + on_progress(EvaluationProgress::MetadataResolved { name }); + } + + Ok(jobs) +} + +async fn configured_hostname(cfg: &str, name: &str) -> String { + let hostname = configured_value(cfg, name, "hostname").await; + if hostname.is_empty() { + name.to_owned() + } else { + hostname + } +} + +async fn configured_value(cfg: &str, name: &str, attribute: &str) -> String { + get_config_attr(cfg, name, attribute) + .await + .unwrap_or_default() +} + +async fn configured_system( + cfg: &str, + name: &str, + results: &[NixEvalJobsResult], +) -> Result { + let configured_type = configured_value(cfg, name, "type").await; + let type_name = if configured_type.is_empty() { + infer_system_type(name, results)? + } else { + configured_type + }; + + SystemType::parse(&type_name).map_err(|error| anyhow::anyhow!("job {name}: {error}")) +} + +fn infer_system_type(name: &str, results: &[NixEvalJobsResult]) -> Result { + let system = results + .iter() + .find(|result| { + result + .attr_path + .as_deref() + .and_then(|path| path.first()) + .is_some_and(|job_name| job_name == name) + }) + .and_then(|result| result.system.as_deref()) + .unwrap_or_default(); + + if system.contains("darwin") { + Ok("darwin".to_string()) + } else if system.contains("linux") { + Ok("nixos".to_string()) + } else { + anyhow::bail!("job {name}: unknown system type: {system}"); + } +} diff --git a/src/process.rs b/src/process.rs index 6d4c1ff..ca7ed84 100644 --- a/src/process.rs +++ b/src/process.rs @@ -1,18 +1,10 @@ use anyhow::{Context, Result}; use serde::de::DeserializeOwned; +use std::process::Stdio; +use tokio::io::{AsyncBufReadExt, AsyncReadExt, BufReader}; use tokio::process::Command; -use crate::models::DebugInfo; - -pub async fn run(cmd: &str, args: &[&str]) -> Result<(Vec, DebugInfo)> { - run_with_env(cmd, args, &[]).await -} - -pub async fn run_with_env( - cmd: &str, - args: &[&str], - envs: &[(&str, &str)], -) -> Result<(Vec, DebugInfo)> { +pub async fn run_with_env(cmd: &str, args: &[&str], envs: &[(&str, &str)]) -> Result> { let display = format!("{} {}", cmd, args.join(" ")); let mut command = Command::new(cmd); @@ -30,13 +22,7 @@ pub async fn run_with_env( let std_err = String::from_utf8_lossy(&output.stderr).into_owned(); let was_success = output.status.success(); - let debug = DebugInfo { - command: display.clone(), - std_out: std_out.clone(), - std_err: std_err.clone(), - }; - - log_debug(&debug); + log_debug(&display, &std_out, &std_err); if !was_success { let detail = if std_err.trim().is_empty() { @@ -47,73 +33,115 @@ pub async fn run_with_env( anyhow::bail!("{display} failed: {}", detail.trim()); } - // Return stdout only — stderr stays in `debug` for logging. - Ok((output.stdout, debug)) + Ok(output.stdout) } /// Run `cmd args...` and deserialize the JSON output into a struct of type `T`. -pub async fn run_json(cmd: &str, args: &[&str]) -> Result<(T, DebugInfo)> { +pub async fn run_json(cmd: &str, args: &[&str]) -> Result { run_json_with_env(cmd, args, &[]).await } -pub async fn run_json_with_env( +/// Run a command that emits one JSON value per line, notifying the caller as +/// each value becomes available. +pub async fn run_json_lines( cmd: &str, args: &[&str], - envs: &[(&str, &str)], -) -> Result<(T, DebugInfo)> { + mut on_line: OnLine, +) -> Result> +where + T: DeserializeOwned, + OnLine: FnMut(&T), +{ let display = format!("{} {}", cmd, args.join(" ")); let mut command = Command::new(cmd); - command.args(args); - for (key, value) in envs { - command.env(key, value); + command + .args(args) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()); + let mut child = command + .spawn() + .with_context(|| format!("failed to spawn {display}"))?; + let stdout = child + .stdout + .take() + .with_context(|| format!("failed to capture stdout from {display}"))?; + let stderr = child + .stderr + .take() + .with_context(|| format!("failed to capture stderr from {display}"))?; + + let stderr_task = tokio::spawn(async move { + let mut text = String::new(); + BufReader::new(stderr).read_to_string(&mut text).await?; + Ok::(text) + }); + + let mut stdout_lines = BufReader::new(stdout).lines(); + let mut stdout_text = String::new(); + let mut values = Vec::new(); + while let Some(line) = stdout_lines + .next_line() + .await + .with_context(|| format!("failed to read stdout from {display}"))? + { + if line.trim().is_empty() { + continue; + } + + stdout_text.push_str(&line); + stdout_text.push('\n'); + let value: T = serde_json::from_str(&line) + .with_context(|| format!("invalid JSON from {display}: {line}"))?; + on_line(&value); + values.push(value); } - let output = command - .output() + let status = child + .wait() .await - .with_context(|| format!("failed to spawn {display}"))?; - - let std_out = String::from_utf8_lossy(&output.stdout).into_owned(); - let std_err = String::from_utf8_lossy(&output.stderr).into_owned(); - let was_success = output.status.success(); - - let debug = DebugInfo { - command: display.clone(), - std_out: std_out.clone(), - std_err: std_err.clone(), - }; + .with_context(|| format!("failed waiting for {display}"))?; + let stderr_text = stderr_task + .await + .context("stderr reader task failed")? + .with_context(|| format!("failed to read stderr from {display}"))?; - log_debug(&debug); + log_debug(&display, &stdout_text, &stderr_text); - if !was_success { - let detail = if std_err.trim().is_empty() { - &std_out + if !status.success() { + let detail = if stderr_text.trim().is_empty() { + &stdout_text } else { - &std_err + &stderr_text }; anyhow::bail!("{display} failed: {}", detail.trim()); } - let value: T = - serde_json::from_str(&std_out).with_context(|| format!("invalid JSON from {display}"))?; + Ok(values) +} - Ok((value, debug)) +pub async fn run_json_with_env( + cmd: &str, + args: &[&str], + envs: &[(&str, &str)], +) -> Result { + let display = format!("{} {}", cmd, args.join(" ")); + let output = run_with_env(cmd, args, envs).await?; + serde_json::from_slice(&output).with_context(|| format!("invalid JSON from {display}")) } -/// Emit a `DebugInfo` at the debug log level (only visible with `--debug`). -pub(crate) fn log_debug(debug: &DebugInfo) { - log::debug!("$ {}", debug.command); - if !debug.std_out.trim().is_empty() { - log::debug!("stdout:\n{}", debug.std_out.trim()); +/// Emit subprocess details at the debug log level (only visible with `--debug`). +pub(crate) fn log_debug(command: &str, std_out: &str, std_err: &str) { + log::debug!("$ {command}"); + if !std_out.trim().is_empty() { + log::debug!("stdout:\n{}", std_out.trim()); } - if !debug.std_err.trim().is_empty() { - log::debug!("stderr:\n{}", debug.std_err.trim()); + if !std_err.trim().is_empty() { + log::debug!("stderr:\n{}", std_err.trim()); } } -pub async fn get_config_attr(cfg: &str, job: &str, attr: &str) -> Result<(String, DebugInfo)> { +pub async fn get_config_attr(cfg: &str, job: &str, attr: &str) -> Result { let attr_path = format!("{cfg}#blzrd.nodes.{job}.{attr}"); - let (value, debug): (String, _) = run_json("nix", &["eval", "--json", &attr_path]).await?; - Ok((value, debug)) + run_json("nix", &["eval", "--json", &attr_path]).await } diff --git a/src/ssh.rs b/src/ssh.rs index c9728e0..ed75ecb 100644 --- a/src/ssh.rs +++ b/src/ssh.rs @@ -3,7 +3,6 @@ use std::time::Duration; use anyhow::{Context, Result}; use openssh::{KnownHosts, SessionBuilder}; -use crate::models::DebugInfo; use crate::process::log_debug; const CONNECT_TIMEOUT: Duration = Duration::from_secs(5); @@ -36,7 +35,7 @@ pub async fn run( program: &str, args: &[&str], host_key_policy: HostKeyPolicy, -) -> Result<(Vec, DebugInfo)> { +) -> Result> { let display = format!("ssh {target} {program} {}", args.join(" ")); let mut builder = SessionBuilder::default(); @@ -64,22 +63,16 @@ pub async fn run( let std_out = String::from_utf8_lossy(&output.stdout).into_owned(); let std_err = String::from_utf8_lossy(&output.stderr).into_owned(); - let debug = DebugInfo { - command: display.clone(), - std_out, - std_err: std_err.clone(), - }; - - log_debug(&debug); + log_debug(&display, &std_out, &std_err); if !output.status.success() { let detail = if std_err.trim().is_empty() { - &debug.std_out + &std_out } else { - &debug.std_err + &std_err }; anyhow::bail!("{display} failed: {}", detail.trim()); } - Ok((output.stdout, debug)) + Ok(output.stdout) } diff --git a/src/ui.rs b/src/ui.rs new file mode 100644 index 0000000..e1c11fc --- /dev/null +++ b/src/ui.rs @@ -0,0 +1,263 @@ +use std::collections::HashMap; +use std::fmt::Display; +use std::io::{self, IsTerminal}; +use std::sync::Mutex; +use std::time::Duration; + +use indicatif::{MultiProgress, ProgressBar, ProgressStyle}; +use owo_colors::{OwoColorize, Stream}; + +use crate::models::JobSpec; + +pub struct Ui { + progress: MultiProgress, + is_interactive: bool, + retained_progress: Mutex>, + needs_section_separator: Mutex, +} + +pub struct SectionProgress { + name: String, + progress: ProgressBar, + jobs: HashMap, +} + +impl Ui { + pub fn new() -> Self { + Self { + progress: MultiProgress::new(), + is_interactive: io::stderr().is_terminal(), + retained_progress: Mutex::new(Vec::new()), + needs_section_separator: Mutex::new(false), + } + } + + pub fn print_nodes(&self, jobs: &HashMap) { + let mut nodes: Vec<_> = jobs.iter().collect(); + nodes.sort_by(|(left_name, _), (right_name, _)| left_name.cmp(right_name)); + + for (name, spec) in nodes { + eprintln!("{name} · {} · {}", spec.system, spec.target()); + } + } + + pub fn print_warning(&self, message: impl Display) { + eprintln!("{}", warning(format!("! {message}"))); + } + + pub fn start_section(&self, name: &str) -> SectionProgress { + self.add_section_separator(); + + let heading = phase_heading(name); + if !self.is_interactive { + eprintln!("⠋ {heading}"); + } + + SectionProgress { + name: heading.clone(), + progress: self.new_spinner(heading, true), + jobs: HashMap::new(), + } + } + + pub fn add_job(&self, section: &mut SectionProgress, name: &str, display_name: &str) { + if section.jobs.contains_key(name) { + return; + } + + let progress = self.new_spinner(display_name, false); + progress.set_style(job_spinner_style()); + section.jobs.insert(name.to_owned(), progress); + } + + pub fn start_job( + &self, + section: &SectionProgress, + name: &str, + display_name: &str, + ) -> ProgressBar { + let progress = section + .jobs + .get(name) + .expect("job progress must be registered") + .clone(); + progress.enable_steady_tick(Duration::from_millis(100)); + progress.tick(); + + if !self.is_interactive { + eprintln!(" ⠋ {display_name}"); + } + + progress + } + + pub fn finish_job_success(&self, progress: ProgressBar, name: &str) { + self.finish_job(progress, format!(" {} {name}", success_marker())); + } + + pub fn finish_job_failure(&self, progress: ProgressBar, name: &str, error: impl Display) { + self.finish_job(progress, format!(" {} {name}: {error}", failure_marker())); + } + + pub fn job_progress(&self, section: &SectionProgress, name: &str) -> ProgressBar { + section + .jobs + .get(name) + .expect("job progress must be registered") + .clone() + } + + pub fn finish_section_success(&self, section: SectionProgress) { + self.finish_section(section, true); + } + + pub fn finish_section_failure(&self, section: SectionProgress) { + self.finish_section(section, false); + } + + pub fn print_completion(&self) { + self.add_section_separator(); + let message = format!( + "{} {}", + success_marker(), + phase_heading("Deployment complete") + ); + if self.is_interactive { + let progress = self.progress.add(ProgressBar::new_spinner()); + progress.set_style(completed_style()); + progress.finish_with_message(message); + self.retain_progress(progress); + } else { + eprintln!("{message}"); + } + } + + fn new_spinner(&self, message: impl Into, should_animate: bool) -> ProgressBar { + let progress = self.progress.add(ProgressBar::new_spinner()); + progress.set_style(spinner_style()); + progress.set_message(message.into()); + progress.tick(); + if should_animate { + progress.enable_steady_tick(Duration::from_millis(100)); + } + progress + } + + fn finish_job(&self, progress: ProgressBar, message: String) { + if self.is_interactive { + progress.set_style(completed_style()); + progress.finish_with_message(message); + } else { + eprintln!("{message}"); + progress.finish_and_clear(); + } + } + + fn finish_section(&self, section: SectionProgress, is_success: bool) { + let message = if is_success { + format!("{} {}", success_marker(), section.name) + } else { + format!("{} {}", failure_marker(), section.name) + }; + + if !is_success { + for progress in section.jobs.values() { + if !progress.is_finished() { + progress.finish_and_clear(); + } + } + } + + if self.is_interactive { + section.progress.set_style(completed_style()); + section.progress.finish_with_message(message); + } else { + eprintln!("{message}"); + section.progress.finish_and_clear(); + } + + self.retain(section); + *self + .needs_section_separator + .lock() + .expect("section separator mutex is not poisoned") = true; + } + + fn add_section_separator(&self) { + let mut needs_separator = self + .needs_section_separator + .lock() + .expect("section separator mutex is not poisoned"); + if !*needs_separator { + return; + } + *needs_separator = false; + drop(needs_separator); + + if self.is_interactive { + let separator = self.progress.add(ProgressBar::new_spinner()); + separator.set_style(completed_style()); + separator.finish_with_message(" ".to_owned()); + self.retain_progress(separator); + } else { + eprintln!(); + } + } + + fn retain(&self, section: SectionProgress) { + let mut retained = self + .retained_progress + .lock() + .expect("progress retention mutex is not poisoned"); + retained.push(section.progress); + retained.extend(section.jobs.into_values()); + } + + fn retain_progress(&self, progress: ProgressBar) { + self.retained_progress + .lock() + .expect("progress retention mutex is not poisoned") + .push(progress); + } +} + +fn spinner_style() -> ProgressStyle { + ProgressStyle::with_template("{spinner} {msg}").expect("spinner template is valid") +} + +fn job_spinner_style() -> ProgressStyle { + ProgressStyle::with_template(" {spinner} {msg}").expect("spinner template is valid") +} + +fn completed_style() -> ProgressStyle { + ProgressStyle::with_template("{msg}").expect("completed template is valid") +} + +fn success_marker() -> String { + format!( + "{}", + "✓".if_supports_color(Stream::Stderr, |marker| marker.green()) + ) +} + +fn failure_marker() -> String { + format!( + "{}", + "✗".if_supports_color(Stream::Stderr, |marker| marker.red()) + ) +} + +fn phase_heading(name: &str) -> String { + let uppercase_name = name.to_uppercase(); + format!( + "{}", + uppercase_name.if_supports_color(Stream::Stderr, |heading| heading.bold()) + ) +} + +fn warning(message: impl Display) -> String { + format!( + "{}", + message.if_supports_color(Stream::Stderr, |message| message.yellow()) + ) +} diff --git a/src/workflow.rs b/src/workflow.rs new file mode 100644 index 0000000..a865a64 --- /dev/null +++ b/src/workflow.rs @@ -0,0 +1,163 @@ +use std::collections::HashMap; + +use anyhow::Context; +use futures::stream::{self, StreamExt}; + +use crate::cli::CommonArgs; +use crate::models::JobSpec; +use crate::nix; +use crate::op::{self, Operation}; +use crate::ssh::HostKeyPolicy; +use crate::ui::Ui; + +/// Filter, validate, build, and deploy the jobs for the given operation. +pub async fn run_deploy( + operation: Operation, + common: CommonArgs, + jobs: HashMap, + ui: &Ui, +) -> anyhow::Result<()> { + let (jobs, unknown_skips) = filter_jobs(jobs, &common)?; + for name in unknown_skips { + ui.print_warning(format!("ignoring unknown node '{name}' in --skip")); + } + + op::validate(&jobs, operation)?; + let host_key_policy = host_key_policy(&common); + let closures = build_closures(&jobs, &common, host_key_policy, ui).await?; + deploy_closures( + jobs, + closures, + common.parallel, + operation, + host_key_policy, + ui, + ) + .await +} + +/// Apply the `--skip` and `nodes` filters to the evaluated job map. +fn filter_jobs( + mut jobs: HashMap, + common: &CommonArgs, +) -> anyhow::Result<(HashMap, Vec)> { + let unknown_skips = common + .skip + .iter() + .filter(|skip_name| !jobs.contains_key(*skip_name)) + .cloned() + .collect(); + + for node_name in &common.nodes { + if !jobs.contains_key(node_name) { + anyhow::bail!("node '{node_name}' not found in flake"); + } + } + + jobs.retain(|name, _| !common.skip.iter().any(|skip_name| skip_name == name)); + if !common.nodes.is_empty() { + jobs.retain(|name, _| common.nodes.iter().any(|node_name| node_name == name)); + } + + Ok((jobs, unknown_skips)) +} + +fn host_key_policy(common: &CommonArgs) -> HostKeyPolicy { + if common.accept_new_host_keys { + HostKeyPolicy::AcceptNew + } else { + HostKeyPolicy::Strict + } +} + +async fn build_closures( + jobs: &HashMap, + common: &CommonArgs, + host_key_policy: HostKeyPolicy, + ui: &Ui, +) -> anyhow::Result> { + let mut progress = ui.start_section("Building"); + let names = sorted_job_names(jobs); + for name in &names { + ui.add_job(&mut progress, name, name); + } + + let mut closures = HashMap::with_capacity(jobs.len()); + for name in names { + let job_progress = ui.start_job(&progress, &name, &name); + let result = nix::build_closure(&jobs[&name], &common.build_host, host_key_policy).await; + match result { + Ok(closure) => { + ui.finish_job_success(job_progress, &name); + closures.insert(name, closure); + } + Err(error) => { + ui.finish_job_failure(job_progress, &name, &error); + ui.finish_section_failure(progress); + return Err(error).context(format!("building {name}")); + } + } + } + + ui.finish_section_success(progress); + Ok(closures) +} + +async fn deploy_closures( + jobs: HashMap, + closures: HashMap, + parallelism: usize, + operation: Operation, + host_key_policy: HostKeyPolicy, + ui: &Ui, +) -> anyhow::Result<()> { + let mut progress = ui.start_section("Deploying"); + let jobs = sorted_jobs(jobs); + for (name, spec) in &jobs { + ui.add_job(&mut progress, name, &deployment_label(name, spec)); + } + + let deployment_section = &progress; + let results = stream::iter(jobs) + .map(|(name, spec)| { + let closure = closures[&name].clone(); + async move { + let label = deployment_label(&name, &spec); + let job_progress = ui.start_job(deployment_section, &name, &label); + let result = nix::deploy_closure(&spec, &closure, operation, host_key_policy).await; + match &result { + Ok(()) => ui.finish_job_success(job_progress, &label), + Err(error) => ui.finish_job_failure(job_progress, &label, error), + } + result + } + }) + .buffer_unordered(parallelism) + .collect::>() + .await; + + if results.iter().any(Result::is_err) { + ui.finish_section_failure(progress); + anyhow::bail!("deployment failed"); + } + + ui.finish_section_success(progress); + ui.print_completion(); + Ok(()) +} + +fn sorted_job_names(jobs: &HashMap) -> Vec { + let mut names: Vec<_> = jobs.keys().cloned().collect(); + names.sort(); + names +} + +fn sorted_jobs(jobs: HashMap) -> Vec<(String, JobSpec)> { + let mut jobs: Vec<_> = jobs.into_iter().collect(); + jobs.sort_by(|(left_name, _), (right_name, _)| left_name.cmp(right_name)); + jobs +} + +fn deployment_label(name: &str, spec: &JobSpec) -> String { + format!("{name} -> {}", spec.target()) +} -- 2.51.2