22 files changed
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -5084,6 +5084,7 @@
"remowt-ui-prompt",
"tokio",
"tracing",
+ "tracing-journald",
"tracing-subscriber",
]
@@ -6576,6 +6577,17 @@
]
[[package]]
+name = "tracing-journald"
+version = "0.3.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "2d3a81ed245bfb62592b1e2bc153e77656d94ee6a0497683a65a12ccaf2438d0"
+dependencies = [
+ "libc",
+ "tracing-core",
+ "tracing-subscriber",
+]
+
+[[package]]
name = "tracing-log"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -75,6 +75,7 @@
tokio = { version = "1.45.1", features = ["fs", "macros", "rt", "rt-multi-thread", "sync", "time"] }
tracing = "0.1"
tracing-indicatif = "0.3.13"
+tracing-journald = "0.3.2"
tracing-opentelemetry = "0.33.0"
uuid = { version = "1", features = ["v4"] }
# For fixed coloring needs to be updated to https://github.com/tokio-rs/tracing/pull/3484
--- a/cmds/fleet/src/cmds/build_systems.rs
+++ b/cmds/fleet/src/cmds/build_systems.rs
@@ -14,6 +14,8 @@
use tokio::task::spawn_blocking;
use tracing::{Instrument, error, field, info, info_span, warn};
+const TARGET: &str = "build-systems";
+
#[derive(Parser)]
pub struct Deploy {
/// Disable automatic rollback
@@ -32,7 +34,7 @@
}
async fn build_task(config: Config, hostname: String, build_attr: &str) -> Result<Utf8PathBuf> {
- info!("building");
+ info!(target: TARGET, "building");
let host = config.host(&hostname)?;
// let action = Action::from(self.subcommand.clone());
let nixos = host.nixos_config()?;
@@ -43,7 +45,7 @@
// We already have system profiles for backups.
if !host.local {
- info!("adding gc root");
+ info!(target: TARGET, "adding gc root");
let local = config.local_host();
let plugin_id = local.ensure_nix_plugin().await?;
let nix = local
@@ -70,7 +72,7 @@
let build_attr = self.build_attr.clone();
for host in hosts {
let config = config.clone();
- let span = info_span!("build", host = field::display(&host.name));
+ let span = info_span!(target: TARGET, "build", host = field::display(&host.name));
let hostname = host.name;
let build_attr = build_attr.clone();
tasks.push(
@@ -78,7 +80,7 @@
let built = match build_task(config, hostname.clone(), &build_attr).await {
Ok(path) => path,
Err(e) => {
- error!("failed to deploy host: {:?}", e);
+ error!(target: TARGET, "failed to deploy host: {:?}", e);
return;
}
};
@@ -86,9 +88,9 @@
let mut out = current_dir().expect("cwd exists");
out.push(format!("built-{hostname}"));
- info!("linking iso image to {:?}", out);
+ info!(target: TARGET, "linking iso image to {:?}", out);
if let Err(e) = symlink(built, out) {
- error!("failed to symlink: {e}")
+ error!(target: TARGET, "failed to symlink: {e}")
}
})
.instrument(span),
@@ -105,7 +107,7 @@
let tasks = FuturesUnordered::new();
for host in hosts.into_iter() {
let config = config.clone();
- let span = info_span!("deploy", host = field::display(&host.name));
+ let span = info_span!(target: TARGET, "deploy", host = field::display(&host.name));
let hostname = host.name.clone();
let opts = opts.clone();
if let Some(deploy_kind) = opts.action_attr::<DeployKind>(&host, "deploy_kind")? {
@@ -125,7 +127,7 @@
{
Ok(path) => path,
Err(e) => {
- error!("failed to build host system closure: {:?}", e);
+ error!(target: TARGET, "failed to build host system closure: {:?}", e);
return;
}
};
@@ -133,7 +135,7 @@
let deploy_kind = match host.deploy_kind().await {
Ok(v) => v,
Err(e) => {
- error!("failed to query target deploy kind: {e}");
+ error!(target: TARGET, "failed to query target deploy kind: {e}");
return;
}
};
@@ -141,7 +143,7 @@
// TODO: Make disable_rollback a host attribute instead
let mut disable_rollback = self.disable_rollback;
if !disable_rollback && deploy_kind != DeployKind::Fleet {
- warn!("disabling rollback, as not supported by non-fleet deployment kinds");
+ warn!(target: TARGET, "disabling rollback, as not supported by non-fleet deployment kinds");
disable_rollback = true;
}
@@ -149,7 +151,7 @@
match upload_task(&host, GenerationStorage::Deployer, built).await {
Ok(v) => v,
Err(e) => {
- error!("upload failed: {e}");
+ error!(target: TARGET, "upload failed: {e}");
return;
}
};
@@ -169,7 +171,7 @@
)
.await
{
- error!("activation failed: {e}");
+ error!(target: TARGET, "activation failed: {e}");
}
})
.instrument(span),
--- /dev/null
+++ b/cmds/fleet/src/log_tree.rs
@@ -0,0 +1,243 @@
+//! Streaming, indented tree formatter.
+//!
+//! Unlike `tracing-tree` it keeps no per-thread depth state, so interleaved
+//! parallel output stays correct; unlike `tracing-forest` it doesn't buffer
+//! whole trees, so logs stream live. Visual style mirrors `@meteor-it/logger`'s
+//! node receiver: a fixed-width, right-aligned, level-colored target column,
+//! then a magenta `>`/`<` gutter for span enter/exit (with duration on exit),
+//! and 2-space indentation per depth for events.
+
+use std::{
+ fmt::Write as _,
+ io::{IsTerminal as _, Write as _},
+ sync::{
+ Mutex,
+ atomic::{AtomicBool, Ordering},
+ },
+ time::{Duration, Instant},
+};
+
+use tracing::{Event, Level, Subscriber, field::Field, field::Visit, span::Attributes, span::Id};
+use tracing_subscriber::{
+ fmt::MakeWriter,
+ layer::{Context, Layer},
+ registry::{LookupSpan, Scope},
+};
+
+/// Width of the right-aligned target column.
+const NAME_LIMIT: usize = 18;
+
+#[derive(Default)]
+struct Fields {
+ message: String,
+ rest: String,
+}
+
+impl Visit for Fields {
+ fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
+ if field.name() == "message" {
+ let _ = write!(self.message, "{value:?}");
+ } else {
+ let _ = write!(self.rest, " {}={:?}", field.name(), value);
+ }
+ }
+}
+
+impl Fields {
+ fn render(&self) -> String {
+ format!("{}{}", self.message, self.rest).trim().to_owned()
+ }
+}
+
+/// Per-span state stashed in the registry's extensions.
+struct SpanData {
+ fields: String,
+ started: Instant,
+ // Whether the span's enter line was written, so we only emit an exit line
+ // for spans that actually appeared.
+ opened: AtomicBool,
+}
+
+/// SGR code → escape for the level-colored target column.
+fn level_code(level: &Level) -> &'static str {
+ match *level {
+ Level::ERROR => "31m", // red
+ Level::WARN => "33m", // yellow
+ Level::INFO => "34m", // blue
+ Level::DEBUG => "39m", // default fg
+ Level::TRACE => "90m", // gray
+ }
+}
+
+/// 2-space indent per depth, then a gutter symbol.
+fn gutter(depth: usize, sym: char) -> String {
+ let mut s = " ".repeat(depth);
+ s.push(sym);
+ s
+}
+
+fn fmt_dur(d: Duration) -> String {
+ if d.as_secs_f64() >= 1.0 {
+ format!("{:.2}s", d.as_secs_f64())
+ } else if d.as_millis() >= 1 {
+ format!("{}ms", d.as_millis())
+ } else {
+ format!("{}µs", d.as_micros())
+ }
+}
+
+pub struct TreeLayer<W> {
+ writer: W,
+ color: bool,
+ // Serializes the whole compute-and-write so enter/exit decisions, output
+ // ordering, and the writes stay consistent under parallel emission.
+ last: Mutex<Vec<Id>>,
+}
+
+impl<W> TreeLayer<W> {
+ pub fn new(writer: W) -> Self {
+ Self {
+ writer,
+ color: std::io::stderr().is_terminal() && std::env::var_os("NO_COLOR").is_none(),
+ last: Mutex::new(vec![]),
+ }
+ }
+
+ /// Raw SGR escape, or empty when color is disabled.
+ fn sgr(&self, code: &str) -> String {
+ if self.color {
+ format!("\x1b[{code}")
+ } else {
+ String::new()
+ }
+ }
+
+ fn name_col(&self, mut target: &str, code: &str) -> String {
+ let w = NAME_LIMIT;
+ if target.len() > w {
+ target = &target[..w];
+ }
+ if self.color {
+ format!("\x1b[{code}\x1b[1m{target:>w$}\x1b[0m")
+ } else {
+ format!("{target:>w$}")
+ }
+ }
+
+ fn print_ancestors<S>(&self, scope: Option<Scope<'_, S>>, out: &mut String) -> usize
+ where
+ S: Subscriber + for<'a> LookupSpan<'a>,
+ {
+ let scope: Vec<_> = scope.map(|s| s.from_root().collect()).unwrap_or_default();
+
+ let mut _last = self.last.lock().expect("not poisoned");
+
+ let skip = scope
+ .iter()
+ .enumerate()
+ .take_while(|(i, s)| _last.get(*i) == Some(&s.id()))
+ .count();
+
+ for (depth, span) in scope.iter().enumerate().skip(skip) {
+ let ext = span.extensions();
+ let Some(data) = ext.get::<SpanData>() else {
+ continue;
+ };
+ let newly_open = !data.opened.swap(true, Ordering::Relaxed);
+ let _ = writeln!(
+ out,
+ "{} {}{}{} {}{}{}",
+ self.name_col(span.metadata().target(), "44m"),
+ self.sgr("35m"),
+ gutter(depth, if newly_open { '>' } else { '.' }),
+ self.sgr("1m"),
+ span.name(),
+ data.fields,
+ self.sgr("0m"),
+ );
+ }
+
+ *_last = scope.iter().map(|v| v.id()).collect();
+
+ scope.len()
+ }
+}
+
+impl<S, W> Layer<S> for TreeLayer<W>
+where
+ S: Subscriber + for<'a> LookupSpan<'a>,
+ W: for<'a> MakeWriter<'a> + 'static,
+{
+ fn on_new_span(&self, attrs: &Attributes<'_>, id: &Id, ctx: Context<'_, S>) {
+ let span = ctx.span(id).expect("new span exists");
+ let mut fields = Fields::default();
+ attrs.record(&mut fields);
+ span.extensions_mut().insert(SpanData {
+ fields: fields.rest,
+ started: Instant::now(),
+ opened: AtomicBool::new(false),
+ });
+ }
+
+ fn on_event(&self, event: &Event<'_>, ctx: Context<'_, S>) {
+ let mut out = String::new();
+
+ let scope = ctx.event_scope(event);
+ let pad = self.print_ancestors(scope, &mut out);
+
+ let mut fields = Fields::default();
+ event.record(&mut fields);
+ let message = fields.render();
+ let meta = event.metadata();
+ for (i, message) in message.split('\n').enumerate() {
+ // Event line: black bg + colored target column, then indent + message.
+ let _ = writeln!(
+ out,
+ "{}{}{}{}{}",
+ self.sgr("40m"),
+ self.name_col(
+ if i == 0 { meta.target() } else { "|" },
+ level_code(meta.level())
+ ),
+ self.sgr("0m"),
+ gutter(pad, ' '),
+ message,
+ );
+ }
+
+ let _ = self.writer.make_writer().write_all(out.as_bytes());
+ }
+
+ fn on_close(&self, id: Id, ctx: Context<'_, S>) {
+ let Some(span) = ctx.span(&id) else { return };
+ let (opened, dur) = span
+ .extensions()
+ .get::<SpanData>()
+ .map_or((false, Duration::ZERO), |d| {
+ (d.opened.load(Ordering::Relaxed), d.started.elapsed())
+ });
+ if !opened {
+ return;
+ }
+
+ let mut out = String::new();
+ let _pad = self.print_ancestors(ctx.span_scope(&id), &mut out);
+
+ let depth = span.scope().skip(1).count(); // ancestors == depth
+ let _ = writeln!(
+ out,
+ "{} {}{}{} {}{} (Done in {}){}",
+ self.name_col(span.metadata().target(), "44m"),
+ self.sgr("35m"),
+ gutter(depth, '<'),
+ self.sgr("1m"),
+ span.name(),
+ self.sgr("22m"),
+ fmt_dur(dur),
+ self.sgr("0m"),
+ );
+
+ let _guard = self.last.lock().expect("not poisoned");
+ let _ = self.writer.make_writer().write_all(out.as_bytes());
+ }
+}
--- a/cmds/fleet/src/main.rs
+++ b/cmds/fleet/src/main.rs
@@ -1,6 +1,7 @@
#![recursion_limit = "512"]
pub(crate) mod cmds;
+mod log_tree;
use std::{process::ExitCode, sync::Arc};
@@ -21,6 +22,7 @@
use human_repr::HumanCount;
#[cfg(feature = "indicatif")]
use indicatif::{ProgressState, ProgressStyle};
+use log_tree::TreeLayer;
use nix_eval::{
eval_store, gc_register_my_thread, gc_unregister_my_thread, init_libraries, init_tokio_for_nix,
};
@@ -128,15 +130,21 @@
fn setup_logging(opts: &RootOpts) -> Result<()> {
#[cfg(feature = "indicatif")]
let indicatif_layer = {
+ use std::fmt;
use std::time::Duration;
IndicatifLayer::new().with_max_progress_bars(10, Some(ProgressStyle::default_spinner()))
+ .with_span_child_prefix_indent(" ")
+ .with_span_child_prefix_symbol("")
.with_progress_style(
ProgressStyle::with_template(
- "{color_start}{span_child_prefix} {span_name}{{{span_fields}}}{color_end} {wide_msg} {color_start}{download_progress} {elapsed}{color_end}",
+ "{span_child_prefix:.magenta}{chevron:.magenta}{span_name:<18.blue.bold} {wide_msg} {span_fields:.dim} {color_start}{download_progress} {elapsed}{color_end}",
)
.unwrap()
- .with_key("download_progress", |state: &ProgressState, writer: &mut dyn std::fmt::Write| {
+ .with_key("chevron", |_: &ProgressState, writer: &mut dyn fmt::Write| {
+ let _ = write!(writer, "> ");
+ })
+ .with_key("download_progress", |state: &ProgressState, writer: &mut dyn fmt::Write| {
let Some(len) = state.len() else {
return;
};
@@ -149,7 +157,7 @@
})
.with_key(
"color_start",
- |state: &ProgressState, writer: &mut dyn std::fmt::Write| {
+ |state: &ProgressState, writer: &mut dyn fmt::Write| {
let elapsed = state.elapsed();
if elapsed > Duration::from_secs(60) {
@@ -163,7 +171,7 @@
)
.with_key(
"color_end",
- |state: &ProgressState, writer: &mut dyn std::fmt::Write| {
+ |state: &ProgressState, writer: &mut dyn fmt::Write| {
if state.elapsed() > Duration::from_secs(30) {
let _ = write!(writer, "\x1b[0m");
}
@@ -174,12 +182,15 @@
let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info"));
- let reg = tracing_subscriber::registry().with({
- let sub = tracing_subscriber::fmt::layer().without_time();
+ let tree = {
#[cfg(feature = "indicatif")]
- let sub = sub.with_writer(indicatif_layer.get_stderr_writer());
- sub.with_filter(filter) // .without,
- });
+ let writer = indicatif_layer.get_stderr_writer();
+ #[cfg(not(feature = "indicatif"))]
+ let writer = || std::io::stderr();
+ TreeLayer::new(writer)
+ };
+
+ let reg = tracing_subscriber::registry().with(tree.with_filter(filter));
#[cfg(feature = "indicatif")]
let reg = reg.with(indicatif_layer);
@@ -200,14 +211,12 @@
let logger = OpenTelemetryTracingBridge::new(&log_provider);
let tracer = span_provider.tracer("fleet");
- let reg = reg
- .with(tracing_opentelemetry::layer().with_tracer(tracer))
- .with(logger);
-
- reg.init();
+ reg.with(tracing_opentelemetry::layer().with_tracer(tracer))
+ .with(logger)
+ .init();
} else {
reg.init();
- };
+ }
Ok(())
}
@@ -252,7 +261,6 @@
.await
.expect("primary task panicked")
})
- // async_main(opts)
}
async fn main_real(opts: RootOpts) -> Result<()> {
1use std::time::Duration;23use anyhow::{Context as _, Result, anyhow, bail};4use camino::Utf8PathBuf;5use clap::ValueEnum;6use itertools::Itertools;7use remowt_endpoints::fs::FsClient;8use remowt_endpoints::systemd::SystemdClient;9use remowt_fleet::NixClient;10use remowt_link_shared::BifConfig;11use tokio::time::sleep;12use tracing::{Instrument as _, error, info, info_span, warn};1314use crate::host::{ConfigHost, DeployKind};15use crate::pins::{Generation, GenerationStorage};1617#[derive(ValueEnum, Clone, Copy)]18pub enum DeployAction {19 20 Upload,21 22 Test,23 24 Boot,25 26 Switch,27}2829impl DeployAction {30 pub(crate) fn name(&self) -> Option<&'static str> {31 match self {32 Self::Upload => None,33 Self::Test => Some("test"),34 Self::Boot => Some("boot"),35 Self::Switch => Some("switch"),36 }37 }38 pub(crate) fn should_switch_profile(&self) -> bool {39 matches!(self, Self::Switch | Self::Boot)40 }41 pub(crate) fn should_activate(&self) -> bool {42 matches!(self, Self::Switch | Self::Test | Self::Boot)43 }44 pub(crate) fn should_create_rollback_marker(&self) -> bool {45 46 47 !matches!(self, Self::Upload)48 }49 pub(crate) fn should_schedule_rollback_run(&self) -> bool {50 matches!(self, Self::Switch | Self::Test)51 }52}5354async fn get_current_generation(host: &ConfigHost) -> Result<Generation> {55 let generations = host.list_generations("system").await?;56 let current = generations57 .into_iter()58 .filter(|g| g.current)59 .at_most_one()60 .map_err(|_e| anyhow!("bad list-generations output"))?61 .ok_or_else(|| anyhow!("failed to find generation"))?;62 Ok(current)63}6465pub async fn deploy_task(66 action: DeployAction,67 host: &ConfigHost,68 built: Utf8PathBuf,69 specialisation: Option<String>,70 disable_rollback: bool,71) -> Result<()> {72 let remowt = host.remowt().await?;73 let deploy_kind = host.deploy_kind().await?;74 if (deploy_kind == DeployKind::NixosInstall || deploy_kind == DeployKind::NixosLustrate)75 && !matches!(action, DeployAction::Boot | DeployAction::Upload)76 {77 bail!("{deploy_kind:?} deploy kind only supports boot and upload actions");78 }7980 let mut failed = false;8182 83 84 85 86 87 if !disable_rollback && action.should_create_rollback_marker() {88 89 info!("preparing for rollback");90 let generation = get_current_generation(host).await?;91 info!(92 "rollback target would be {} {}",93 generation.id, generation.datetime94 );95 {96 let mut cmd = remowt.cmd("sh");97 cmd.arg("-c").arg(format!("mark=$(mktemp -p /etc -t fleet_rollback_marker.XXXXX) && echo -n {} > $mark && mv --no-clobber $mark /etc/fleet_rollback_marker", generation.id));98 if let Err(e) = cmd.sudo().run().await {99 error!("failed to set rollback marker: {e}");100 failed = true;101 }102 }103 104 105 106 107 108109 110 111 112 113 if action.should_schedule_rollback_run() {114 let mut cmd = remowt.cmd("systemd-run");115 cmd.comparg("--on-active", "3min")116 .comparg("--unit", "rollback-watchdog-run")117 .arg("systemctl")118 .arg("start")119 .arg("rollback-watchdog.service");120 if let Err(e) = cmd.sudo().run().await {121 error!("failed to schedule rollback run: {e}");122 failed = true;123 }124 }125 }126127 let remowt = host.remowt().await?;128 let fs = remowt.endpoints::<FsClient<_>>();129 if deploy_kind == DeployKind::NixosLustrate {130 131 132 if !fs133 .file_exists(Utf8PathBuf::from("/etc/NIXOS_LUSTRATE"))134 .await?135 {136 bail!("/etc/NIXOS_LUSTRATE should be created on remote host");137 }138 139 let mut cmd = remowt.cmd("touch");140 cmd.arg("/etc/NIXOS");141 cmd.sudo().run().await.context("creating /etc/NIXOS")?;142 }143 if deploy_kind == DeployKind::NixosInstall {144 info!(145 "running nixos-install to switch profile, install bootloader, and perform activation"146 );147 let mut cmd = remowt.cmd("nixos-install");148 cmd.arg("--system").arg(&built).args([149 150 151 "--no-channel-copy",152 "--root",153 "/mnt",154 ]);155 if let Err(e) = cmd.sudo().run().await {156 error!("failed to execute nixos-install: {e}");157 failed = true;158 }159 } else {160 if action.should_switch_profile() && !failed {161 info!("switching system profile generation");162163 match host.ensure_nix_plugin().await {164 Ok(plugin_id) => {165 let nix_elevated = remowt.plugin_endpoints::<NixClient<_>>(plugin_id);166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 if let Err(e) = nix_elevated184 .switch_profile("/nix/var/nix/profiles/system".to_owned(), built.clone())185 .await186 {187 error!("failed to switch system profile generation: {e}");188 failed = true;189 }190 }191 Err(e) => {192 failed = true;193 error!("failed to enable nix plugin: {e:?}");194 }195 }196 }197198 199200 if action.should_activate() && !failed {201 202 info!("executing activation script");203 let specialised = if let Some(specialisation) = specialisation {204 let mut specialised = built.join("specialisation");205 specialised.push(specialisation);206 specialised207 } else {208 built.clone()209 };210 let switch_script = specialised.join("bin/switch-to-configuration");211 let mut cmd = remowt.cmd("systemd-run");212 if deploy_kind == DeployKind::NixosLustrate {213 cmd.arg("--setenv=NIXOS_INSTALL_BOOTLOADER=1");214 }215 cmd.arg("--setenv=FLEET_ONLINE_ACTIVATION=1")216 .arg("--collect")217 .arg("--no-ask-password")218 .arg("--pipe")219 .arg("--quiet")220 .arg("--service-type=exec")221 .arg("--unit=fleet-switch-to-configuration")222 .arg(switch_script)223 .arg(action.name().expect("upload.should_activate == false"));224 if let Err(e) = cmd.sudo().run().in_current_span().await {225 error!("failed to activate: {e}");226 failed = true;227 }228 }229 }230 if action.should_create_rollback_marker() {231 let elevated_systemd = remowt.run0_endpoints::<SystemdClient<BifConfig>>().await?;232 let elevated_fs = remowt.run0_endpoints::<FsClient<BifConfig>>().await?;233 if !disable_rollback {234 if failed {235 if action.should_schedule_rollback_run() {236 info!("executing rollback");237 if let Err(e) = elevated_systemd238 .start("rollback-watchdog.service".to_owned())239 .instrument(info_span!("rollback"))240 .await241 {242 error!("failed to trigger rollback: {e}")243 }244 }245 } else {246 info!("trying to mark upgrade as successful");247 if let Err(e) = elevated_fs248 .rm_file(Utf8PathBuf::from("/etc/fleet_rollback_marker"))249 .in_current_span()250 .await251 {252 error!(253 "failed to remove rollback marker. This is bad, as the system will be rolled back by watchdog: {e}"254 )255 }256 }257 info!("disarming watchdog, just in case");258 if let Err(_e) = elevated_systemd259 .stop("rollback-watchdog.timer".to_owned())260 .await261 {262 263 }264 if action.should_schedule_rollback_run()265 && let Err(e) = elevated_systemd266 .stop("rollback-watchdog-run.timer".to_owned())267 .await268 {269 error!("failed to disarm rollback run: {e}");270 }271 } else if let Err(_e) = elevated_fs272 .rm_file(Utf8PathBuf::from("/etc/fleet_rollback_marker"))273 .in_current_span()274 .await275 {276 277 }278 }279 Ok(())280}281282pub async fn upload_task(283 host: &ConfigHost,284 location: GenerationStorage,285 generation: Utf8PathBuf,286) -> Result<Utf8PathBuf> {287 if matches!(location, GenerationStorage::Pusher) {288 bail!("pusher is not enabled in this version of fleet");289 }290 if !host.local {291 info!("uploading system closure");292 let mut tries = 0;293 loop {294 match host.remote_derivation(&generation).await {295 Ok(remote) => {296 assert!(remote == generation, "CA derivations aren't implemented");297 return Ok(remote);298 }299 Err(e) if tries < 3 => {300 tries += 1;301 warn!("copy failure ({}/3): {:#}", tries, e);302 sleep(Duration::from_millis(5000)).await;303 }304 Err(e) => {305 bail!("upload failed: {e:#}");306 }307 }308 }309 }310 Ok(generation)311}
1use std::time::Duration;23use anyhow::{Context as _, Result, anyhow, bail};4use camino::Utf8PathBuf;5use clap::ValueEnum;6use itertools::Itertools;7use remowt_endpoints::fs::FsClient;8use remowt_endpoints::systemd::SystemdClient;9use remowt_fleet::NixClient;10use remowt_link_shared::BifConfig;11use tokio::time::sleep;12use tracing::{Instrument as _, error, info, info_span, warn};1314use crate::host::{ConfigHost, DeployKind};15use crate::pins::{Generation, GenerationStorage};1617#[derive(ValueEnum, Clone, Copy)]18pub enum DeployAction {19 20 Upload,21 22 Test,23 24 Boot,25 26 Switch,27}2829impl DeployAction {30 pub(crate) fn name(&self) -> Option<&'static str> {31 match self {32 Self::Upload => None,33 Self::Test => Some("test"),34 Self::Boot => Some("boot"),35 Self::Switch => Some("switch"),36 }37 }38 pub(crate) fn should_switch_profile(&self) -> bool {39 matches!(self, Self::Switch | Self::Boot)40 }41 pub(crate) fn should_activate(&self) -> bool {42 matches!(self, Self::Switch | Self::Test | Self::Boot)43 }44 pub(crate) fn should_create_rollback_marker(&self) -> bool {45 46 47 !matches!(self, Self::Upload)48 }49 pub(crate) fn should_schedule_rollback_run(&self) -> bool {50 matches!(self, Self::Switch | Self::Test)51 }52}5354async fn get_current_generation(host: &ConfigHost) -> Result<Generation> {55 let generations = host.list_generations("system").await?;56 let current = generations57 .into_iter()58 .filter(|g| g.current)59 .at_most_one()60 .map_err(|_e| anyhow!("bad list-generations output"))?61 .ok_or_else(|| anyhow!("failed to find generation"))?;62 Ok(current)63}6465pub async fn deploy_task(66 action: DeployAction,67 host: &ConfigHost,68 built: Utf8PathBuf,69 specialisation: Option<String>,70 disable_rollback: bool,71) -> Result<()> {72 let remowt = host.remowt().await?;73 let deploy_kind = host.deploy_kind().await?;74 if (deploy_kind == DeployKind::NixosInstall || deploy_kind == DeployKind::NixosLustrate)75 && !matches!(action, DeployAction::Boot | DeployAction::Upload)76 {77 bail!("{deploy_kind:?} deploy kind only supports boot and upload actions");78 }7980 let mut failed = false;8182 83 84 85 86 87 if !disable_rollback && action.should_create_rollback_marker() {88 89 info!("preparing for rollback");90 let generation = get_current_generation(host).await?;91 info!(92 "rollback target would be {} {}",93 generation.id, generation.datetime94 );95 {96 let mut cmd = remowt.cmd("sh");97 cmd.arg("-c").arg(format!("mark=$(mktemp -p /etc -t fleet_rollback_marker.XXXXX) && echo -n {} > $mark && mv --no-clobber $mark /etc/fleet_rollback_marker", generation.id));98 if let Err(e) = cmd.sudo().run().await {99 error!("failed to set rollback marker: {e}");100 failed = true;101 }102 }103 104 105 106 107 108109 110 111 112 113 if action.should_schedule_rollback_run() {114 let mut cmd = remowt.cmd("systemd-run");115 cmd.comparg("--on-active", "3min")116 .comparg("--unit", "rollback-watchdog-run")117 .arg("systemctl")118 .arg("start")119 .arg("rollback-watchdog.service");120 if let Err(e) = cmd.sudo().run().await {121 error!("failed to schedule rollback run: {e}");122 failed = true;123 }124 }125 }126127 let remowt = host.remowt().await?;128 let fs = remowt.endpoints::<FsClient<_>>();129 if deploy_kind == DeployKind::NixosLustrate {130 131 132 if !fs133 .file_exists(Utf8PathBuf::from("/etc/NIXOS_LUSTRATE"))134 .await?135 {136 bail!("/etc/NIXOS_LUSTRATE should be created on remote host");137 }138 139 let mut cmd = remowt.cmd("touch");140 cmd.arg("/etc/NIXOS");141 cmd.sudo().run().await.context("creating /etc/NIXOS")?;142 }143 if deploy_kind == DeployKind::NixosInstall {144 info!(145 "running nixos-install to switch profile, install bootloader, and perform activation"146 );147 let mut cmd = remowt.cmd("nixos-install");148 cmd.arg("--system").arg(&built).args([149 150 151 "--no-channel-copy",152 "--root",153 "/mnt",154 ]);155 if let Err(e) = cmd.sudo().run().await {156 error!("failed to execute nixos-install: {e}");157 failed = true;158 }159 } else {160 if action.should_switch_profile() && !failed {161 info!("switching system profile generation");162163 match host.ensure_nix_plugin().await {164 Ok(plugin_id) => {165 let nix_elevated = remowt.plugin_endpoints::<NixClient<_>>(plugin_id);166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 if let Err(e) = nix_elevated184 .switch_profile("/nix/var/nix/profiles/system".to_owned(), built.clone())185 .await186 {187 error!("failed to switch system profile generation: {e}");188 failed = true;189 }190 }191 Err(e) => {192 failed = true;193 error!("failed to enable nix plugin: {e:?}");194 }195 }196 }197198 199200 if action.should_activate() && !failed {201 202 info!("executing activation script");203 let specialised = if let Some(specialisation) = specialisation {204 let mut specialised = built.join("specialisation");205 specialised.push(specialisation);206 specialised207 } else {208 built.clone()209 };210 let switch_script = specialised.join("bin/switch-to-configuration");211 let mut cmd = remowt.cmd("systemd-run");212 if deploy_kind == DeployKind::NixosLustrate {213 cmd.arg("--setenv=NIXOS_INSTALL_BOOTLOADER=1");214 }215 cmd.arg("--setenv=FLEET_ONLINE_ACTIVATION=1")216 .arg("--collect")217 .arg("--no-ask-password")218 .arg("--pipe")219 .arg("--quiet")220 .arg("--service-type=exec")221 .arg("--unit=fleet-switch-to-configuration")222 .arg(switch_script)223 .arg(action.name().expect("upload.should_activate == false"));224 if let Err(e) = cmd.sudo().run().in_current_span().await {225 error!("failed to activate: {e}");226 failed = true;227 }228 }229 }230 if action.should_create_rollback_marker() {231 let elevated_systemd = remowt.run0_endpoints::<SystemdClient<BifConfig>>().await?;232 let elevated_fs = remowt.run0_endpoints::<FsClient<BifConfig>>().await?;233 if !disable_rollback {234 if failed {235 if action.should_schedule_rollback_run() {236 info!("executing rollback");237 if let Err(e) = elevated_systemd238 .start("rollback-watchdog.service".to_owned())239 .instrument(info_span!("rollback"))240 .await241 {242 error!("failed to trigger rollback: {e}")243 }244 }245 } else {246 info!("trying to mark upgrade as successful");247 if let Err(e) = elevated_fs248 .rm_file(Utf8PathBuf::from("/etc/fleet_rollback_marker"))249 .in_current_span()250 .await251 {252 error!(253 "failed to remove rollback marker. This is bad, as the system will be rolled back by watchdog: {e}"254 )255 }256 }257 info!("disarming watchdog, just in case");258 if let Err(_e) = elevated_systemd259 .stop("rollback-watchdog.timer".to_owned())260 .await261 {262 263 }264 if action.should_schedule_rollback_run()265 && let Err(e) = elevated_systemd266 .stop("rollback-watchdog-run.timer".to_owned())267 .await268 {269 error!("failed to disarm rollback run: {e}");270 }271 } else if let Err(_e) = elevated_fs272 .rm_file(Utf8PathBuf::from("/etc/fleet_rollback_marker"))273 .in_current_span()274 .await275 {276 277 }278 }279 Ok(())280}281282pub async fn upload_task(283 host: &ConfigHost,284 location: GenerationStorage,285 generation: Utf8PathBuf,286) -> Result<Utf8PathBuf> {287 if matches!(location, GenerationStorage::Pusher) {288 bail!("pusher is not enabled in this version of fleet");289 }290 if !host.local {291 info!("uploading system closure");292 let mut tries = 0;293 294 loop {295 let name = format!("{}-system", host.gc_root_prefix());296 match host.remote_derivation(&generation, Some(&name)).await {297 Ok(remote) => {298 assert!(remote == generation, "CA derivations aren't implemented");299 return Ok(remote);300 }301 Err(e) if tries < 3 => {302 tries += 1;303 warn!("copy failure ({}/3): {:#}", tries, e);304 sleep(Duration::from_millis(5000)).await;305 }306 Err(e) => {307 bail!("upload failed: {e:#}");308 }309 }310 }311 }312 Ok(generation)313}
--- a/crates/fleet-base/src/host.rs
+++ b/crates/fleet-base/src/host.rs
@@ -13,16 +13,16 @@
use camino::{Utf8Path, Utf8PathBuf};
use chrono::{DateTime, Utc};
use fleet_shared::SecretData;
-use nix_eval::{Store, Value, drv::DrvGraph, eval_store, nix_go, nix_go_json, util::assert_warn};
+use nix_eval::{Store, Value, eval_store, nix_go, nix_go_json, util::assert_warn};
use remowt_client::{AgentBundle, Remowt};
use remowt_endpoints::fs::FsClient;
use remowt_fleet::NixClient;
-use remowt_link_shared::BifConfig;
+use remowt_link_shared::{Address, BifConfig};
use remowt_ui_prompt::auto::AutoPrompter;
use remowt_ui_prompt::bifrost::PromptEndpoints;
use remowt_ui_prompt::{PrependSourcePrompter, Source};
use tempfile::NamedTempFile;
-use tokio::{sync::OnceCell, task::spawn_blocking};
+use tokio::{process::Command, sync::OnceCell, task::spawn_blocking};
use tracing::{info, instrument, warn};
use crate::fleetdata::{
@@ -104,11 +104,14 @@
session_destination: OnceLock<String>,
legacy_ssh_store: OnceLock<bool>,
+ /// fleetConfiguration.hosts.host, None for local (as local is only used for local tool overrides)
pub host_config: Option<Value>,
pub nixos_config: OnceLock<Value>,
pub nixos_unchecked_config: OnceLock<Value>,
pub pkgs_override: Option<Value>,
+ binary_cache: OnceCell<Option<BinaryCache>>,
+
// TODO: Move command helpers away with connectivity refactor
pub local: bool,
pub remowt: OnceCell<Remowt>,
@@ -131,13 +134,35 @@
AgentBundle::from_dir(agents_dir()?)
}
-enum RemoteDerivationMode {
- /// Closure is (pre)signed, remote host is expected to trust our signatures.
- CopySigned,
- /// Closure is not signed, using escalated nix daemon to bypass signature verification.
- CopyPrivileged,
- /// Closure is uploaded to attic, remote host is expected to trust its signature.
- Attic { name: String },
+#[derive(Clone)]
+pub enum BinaryCache {
+ Attic { name: String, pin: bool },
+}
+impl BinaryCache {
+ async fn upload(&self, local: &Remowt, path: &Utf8Path, pin_name: Option<&str>) -> Result<()> {
+ match self {
+ BinaryCache::Attic { name, pin } => {
+ let mut cmd = local.cmd("attic");
+ cmd.arg("push").arg(name).arg(path.as_str());
+ cmd.run()
+ .await
+ .context("failed to push path to attic (is attic-cli on PATH?)")?;
+
+ if *pin && let Some(pin_name) = pin_name {
+ let mut cmd = local.cmd("attic");
+ cmd.arg("pin")
+ .arg("create")
+ .arg(name)
+ .arg(pin_name)
+ .arg(path.as_str());
+ cmd.run()
+ .await
+ .context("failed to pin path in attic (is your attic-cli compatible?)")?;
+ }
+ Ok(())
+ }
+ }
+ }
}
impl ConfigHost {
@@ -156,6 +181,32 @@
.set(legacy)
.expect("legacy ssh store is already set")
}
+ pub fn binary_cache(&self) -> Result<Option<BinaryCache>> {
+ if let Some(v) = self.binary_cache.get() {
+ return Ok(v.clone());
+ }
+ let cache = if let Some(host_config) = &self.host_config {
+ let cache = nix_go!(host_config.cache);
+ let enable: bool = nix_go_json!(cache.enable);
+ if enable {
+ let kind: String = nix_go_json!(cache.kind);
+ let name: String = nix_go_json!(cache.name);
+ match kind.as_str() {
+ "attic" => {
+ let pin: bool = nix_go_json!(cache.pin);
+ Some(BinaryCache::Attic { name, pin })
+ }
+ v => bail!("unknown binary cache kind: {v}"),
+ }
+ } else {
+ None
+ }
+ } else {
+ None
+ };
+ let _ = self.binary_cache.set(cache.clone());
+ Ok(cache)
+ }
pub async fn deploy_kind(&self) -> Result<DeployKind> {
self.deploy_kind.get_or_try_init(|| async {
let remowt = self.remowt().await?;
@@ -238,8 +289,10 @@
let built = plugin
.build("out")
.context("failed to build the fleet nix plugin")?;
+ let pin = format!("{}-remowt-plugin-fleet", self.gc_root_prefix());
let copied = self
- .remote_derivation(&built)
+ // TODO: Pin? But this pin will be host-specific, and heavy because of nix dependency...
+ .remote_derivation(&built, Some(&pin))
.await
.context("failed to copy the fleet nix plugin to the host store")?;
let bin = copied.join("bin/remowt-plugin-fleet");
@@ -248,6 +301,13 @@
.run0_load_plugin_path(NIX_PLUGIN_ID, bin.as_str())
.await
.context("failed to load the fleet nix plugin")?;
+ // Should not fail, as plugin loading waits for connection on remote, and we only wait for forward... On the other hand, we should monitor connection to the remote itself somewhere too.
+ self.remowt()
+ .await?
+ .rpc()
+ .wait_for_connection_to(Address::Plugin(NIX_PLUGIN_ID))
+ .await
+ .map_err(|_| anyhow!("failed to wait"))?;
anyhow::Ok(())
})
.await?;
@@ -339,17 +399,31 @@
}
/// Returns path for futureproofing, as path might change i.e on conversion to CA
#[instrument(skip(self))]
- pub async fn remote_derivation(&self, path: &Utf8Path) -> Result<Utf8PathBuf> {
+ pub async fn remote_derivation(
+ &self,
+ path: &Utf8Path,
+ pin_name: Option<&str>,
+ ) -> Result<Utf8PathBuf> {
let path = path.to_owned();
if self.local {
// Path is located locally, thus already trusted.
return Ok(path);
}
+ let cache = self.binary_cache()?;
+ let mut cache = if let Some(cache) = cache {
+ let local = self.config.local_host().remowt().await?;
+ let path = path.to_owned();
+ let pin_name = pin_name.map(str::to_owned);
+ Some(tokio::spawn(async move {
+ if let Err(e) = cache.upload(&local, &path, pin_name.as_deref()).await {
+ warn!("failed to upload closure {path} to cache: {e:?}");
+ }
+ }))
+ } else {
+ None
+ };
let sign: Pin<Box<dyn Future<Output = Result<()>> + Send>> = {
let path = path.clone();
- let graph = DrvGraph::resolve(&eval_store(), &path)
- .context("failed to resolve graph to be uploaded")?;
- info!("signing {} paths", graph.nodes.len());
Box::pin(async move {
let local = self.config.local_host();
@@ -377,13 +451,62 @@
.expect("copy_to panicked")
.context("copying closure to remote store")?;
}
+
+ if let Some(cache) = cache.take() {
+ let _ = cache.await;
+ }
+
Ok(path)
}
+
+ /// Like [`Self::remote_derivation`], but instead of copying the closure over
+ /// the nix daemon tunnel, uploads it to an attic binary cache and has the
+ /// remote host substitute it from there. The remote is expected to already
+ /// trust the cache's signing key and have it configured as a substituter.
+ #[instrument(skip(self))]
+ pub async fn remote_derivation_attic(
+ &self,
+ path: &Utf8Path,
+ cache_name: &str,
+ ) -> Result<Utf8PathBuf> {
+ let path = path.to_owned();
+ if self.local {
+ // Path is located locally, thus already trusted.
+ return Ok(path);
+ }
+
+ info!("uploading {path} closure to attic cache {cache_name}");
+ let status = Command::new("attic")
+ .arg("push")
+ .arg(cache_name)
+ .arg(path.as_str())
+ .status()
+ .await
+ .context("failed to spawn attic-cli, is it on PATH?")?;
+ ensure!(status.success(), "attic push failed: {status}");
+
+ info!("substituting {path} on the remote host");
+ let nix = self.nix_client().await?;
+ let substituted = nix
+ .substitute(vec![path.clone()])
+ .await
+ .map_err(|e| anyhow!("{e:?}"))?
+ .map_err(|e| anyhow!("{e}"))?;
+ ensure!(
+ substituted.contains(&path),
+ "remote host failed to substitute {path} from attic cache {cache_name}"
+ );
+ Ok(path)
+ }
}
struct HostSecretDefinition(Value);
impl ConfigHost {
+ pub fn gc_root_prefix(&self) -> String {
+ format!("{}-{}", self.config.gc_root_prefix(), self.name)
+ }
+
// TOCTOU is possible here in case if config is changed, but this case is not handled anywhere anyway,
// assuming getting tags always returns the same value.
pub fn tags(&self) -> Result<Vec<String>> {
@@ -504,6 +627,7 @@
cell
},
pkgs_override: Some(self.default_pkgs.clone()),
+ binary_cache: OnceCell::new(),
local: true,
remowt: OnceCell::new(),
@@ -546,6 +670,7 @@
nixos_unchecked_config: OnceLock::new(),
groups: OnceLock::new(),
pkgs_override: None,
+ binary_cache: OnceCell::new(),
// TODO: Remove with connectivit refactor
local: self.localhost == name,
--- a/crates/fleet-base/src/keys.rs
+++ b/crates/fleet-base/src/keys.rs
@@ -9,6 +9,9 @@
use crate::{fleetdata::SecretOwner, host::Config};
impl Config {
+ pub fn gc_root_prefix(&self) -> String {
+ self.data.gc_root_prefix.clone()
+ }
fn cached_host_key(&self, host: &str) -> Option<String> {
let hosts = self.data.hosts.read().expect("no poisoning");
let key = hosts.get(host).map(|h| &h.encryption_key);
--- a/crates/fleet-base/src/primops.rs
+++ b/crates/fleet-base/src/primops.rs
@@ -119,7 +119,7 @@
.context("failed to build generator for target host")?;
let generator = host_on
- .remote_derivation(&generator)
+ .remote_derivation(&generator, None)
.await
.context("failed to copy generator to target host")?;
--- a/crates/nix-eval/src/lib.rs
+++ b/crates/nix-eval/src/lib.rs
@@ -39,7 +39,7 @@
make_attrs, make_bindings_builder, make_list, make_list_builder, realised_string,
realised_string_free, realised_string_get_buffer_size, realised_string_get_buffer_start,
realised_string_get_store_path, realised_string_get_store_path_count, register_primop,
- set_err_msg, setting_set, state_free, store_copy_closure, store_free, store_open,
+ set_err_msg, setting_get, setting_set, state_free, store_copy_closure, store_free, store_open,
store_parse_path, store_path_free, store_path_name, string_realise, value, value_call,
value_decref, value_force, value_incref,
};
@@ -390,6 +390,14 @@
with_default_context(|c, _| unsafe { setting_set(c, s.as_ptr(), v.as_ptr()) }).map(|_| ())
}
+pub fn get_setting(s: &CStr) -> Result<String> {
+ let mut out = String::new();
+ with_default_context(|c, _| unsafe {
+ setting_get(c, s.as_ptr(), Some(copy_nix_str), (&raw mut out).cast())
+ })?;
+ Ok(out)
+}
+
#[derive(Debug)]
pub struct AddedFile {
pub store_path: Utf8PathBuf,
--- a/crates/nix-eval/src/logging.rs
+++ b/crates/nix-eval/src/logging.rs
@@ -109,6 +109,12 @@
(ActivityType::CopyPaths, []) => {
debug_span!(target: "nix::copy-paths", "copying paths")
}
+ // Umbrella progress activity. Nix emits it at lvlError ("always show"),
+ // which would otherwise become a noisy ERROR `action` span — give it an
+ // explicit debug level like the other container activities.
+ (ActivityType::Builds, []) => {
+ debug_span!(target: "nix::builds", "building")
+ }
(ActivityType::Unknown, [])
if s.starts_with("copying \"") && s.ends_with("\" to the store") =>
{
@@ -365,6 +371,16 @@
BuildGraphGuard { paths }
}
+/// Drop the lazily-created `building` span for a derivation. Used when a build
+/// fails: its dependents will never build, so their speculatively-created spans
+/// (made when a dependency started building) would otherwise linger until the
+/// whole build graph is torn down.
+pub fn remove_drv_span(drv_path: &Utf8Path) {
+ if let Some(entry) = DRV_GRAPH.lock().expect("not poisoned").get_mut(drv_path) {
+ entry.span = None;
+ }
+}
+
fn ensure_drv_span(drv_path: &Utf8Path) -> Option<Span> {
let mut drv_graph = DRV_GRAPH.lock().expect("not poisoned");
@@ -524,7 +540,7 @@
// ResultType::FileLinked => todo!(),
(ResultType::BuildLogLine, [Str(s)]) => {
let s = ansi_filter(s);
- info!("{s}");
+ info!(target: "nix::build", "{s}");
}
// ResultType::UntrustedPath => todo!(),
// ResultType::CorruptedPath => todo!(),
--- a/crates/nix-eval/src/scheduler.rs
+++ b/crates/nix-eval/src/scheduler.rs
@@ -45,9 +45,36 @@
},
}
+// Derivations whose names contain any of these are "heavy": memory-hungry builds that already
+// saturate all cores internally (e.g. composable-kernels), so running anything alongside them
+// just thrashes RAM. A heavy build grabs every semaphore permit and runs exclusively.
+const DEFAULT_HEAVY_DRVS: &[&str] = &[
+ "composable-kernel",
+ "composable_kernel",
+ "ghc",
+ "gfortran",
+ "llvm",
+ "clang",
+];
+
+fn heavy_patterns_from_env() -> Vec<String> {
+ let mut pats: Vec<String> = DEFAULT_HEAVY_DRVS.iter().map(|s| (*s).to_owned()).collect();
+ if let Ok(extra) = std::env::var("FLEET_HEAVY_DRVS") {
+ pats.extend(
+ extra
+ .split(',')
+ .map(str::trim)
+ .filter(|s| !s.is_empty())
+ .map(str::to_owned),
+ );
+ }
+ pats
+}
+
pub struct Scheduler {
store: Arc<Store>,
parallelism: usize,
+ heavy_patterns: Vec<String>,
events: broadcast::Sender<BuildEvent>,
}
@@ -58,10 +85,20 @@
Self {
store: eval_store(),
parallelism,
+ heavy_patterns: heavy_patterns_from_env(),
events,
}
}
+ // Permits a drv must hold to run. Heavy drvs take all of them, pausing everything else.
+ fn permits_for(&self, name: &str) -> u32 {
+ if self.heavy_patterns.iter().any(|p| name.contains(p)) {
+ self.parallelism as u32
+ } else {
+ 1
+ }
+ }
+
pub fn subscribe(&self) -> broadcast::Receiver<BuildEvent> {
self.events.subscribe()
}
@@ -139,6 +176,7 @@
.get(&path)
.map(|n| n.name.clone())
.unwrap_or_default();
+ crate::logging::remove_drv_span(&path);
let _ = self.events.send(BuildEvent::DrvCancelled {
drv_path: path.clone(),
name,
@@ -153,8 +191,21 @@
let graph = graph.clone();
let wanted_here = wanted.get(&path).cloned().unwrap_or_default();
let store = self.store.clone();
+ let permits = graph
+ .nodes
+ .get(&path)
+ .map(|n| self.permits_for(&n.name))
+ .unwrap_or(1);
in_flight.push(tokio::spawn(async move {
- let _permit = sem.acquire_owned().await.expect("semaphore not closed");
+ if permits > 1 {
+ debug!(permits, "heavy drv: waiting for exclusive build slot");
+ }
+ // Heavy drvs acquire every permit, so they only start once all in-flight
+ // builds drain and nothing new can slip in (the semaphore is FIFO-fair).
+ let _permit = sem
+ .acquire_many_owned(permits)
+ .await
+ .expect("semaphore not closed");
let node = graph
.nodes
.get(&path)
@@ -224,6 +275,7 @@
propagate_done(&dependents, &mut indeg, &mut ready, &finished);
}
Err(e) => {
+ crate::logging::remove_drv_span(&finished);
failed.insert(finished.clone(), format!("{e:#}"));
mark_tainted(&dependents, &finished, &mut tainted);
propagate_done(&dependents, &mut indeg, &mut ready, &finished);
@@ -368,11 +420,23 @@
// TODO: Parallelism as a metric works poorly with multiple machines, but I haven't thought about bringing
// hercy here yet. In case of remote machines - they will handle parallelism on their own, and this one
// will work as a hard cap.
+/// Number of derivations to build concurrently, taken from Nix's `max-jobs`
+/// setting. `auto` resolves to the CPU count, matching Nix itself.
+fn max_jobs() -> usize {
+ let cores = || {
+ std::thread::available_parallelism()
+ .map(|p| p.get())
+ .unwrap_or(4)
+ };
+ match crate::get_setting(c"max-jobs") {
+ Ok(v) if v.trim() == "auto" => cores(),
+ Ok(v) => v.trim().parse().unwrap_or_else(|_| cores()),
+ Err(_) => cores(),
+ }
+}
+
pub fn build_graph_sync(graph: Arc<DrvGraph>, root_outputs: Vec<String>) -> Result<()> {
- let parallelism = std::thread::available_parallelism()
- .map(|p| p.get())
- .unwrap_or(4);
- let scheduler = Scheduler::new(parallelism);
+ let scheduler = Scheduler::new(max_jobs());
crate::await_in_nix(async move { scheduler.run(graph, root_outputs).await })
.context("scheduler run")
}
--- a/crates/remowt-fleet/src/lib.rs
+++ b/crates/remowt-fleet/src/lib.rs
@@ -28,6 +28,8 @@
Sign(String),
#[error("listing generations failed: {0}")]
ListGenerations(String),
+ #[error("substitution failed: {0}")]
+ Substitute(String),
}
#[endpoints(ns = 91)]
@@ -60,6 +62,20 @@
.map_err(|e| NixError::Sign(e.to_string()))
}
+ /// Download the given paths (and their closures) into the local store using
+ /// the configured substituters. Used to fetch closures that were uploaded to
+ /// a binary cache (e.g. attic) instead of copied over the nix daemon tunnel.
+ #[endpoints(id = 6)]
+ async fn substitute(&self, paths: Vec<Utf8PathBuf>) -> Result<Vec<Utf8PathBuf>, NixError> {
+ spawn_blocking(move || {
+ let store = eval_store();
+ store.substitute_paths(&paths)
+ })
+ .await
+ .expect("substitution panicked")
+ .map_err(|e| NixError::Substitute(e.to_string()))
+ }
+
#[endpoints(id = 5)]
async fn list_generations(
&self,
--- /dev/null
+++ b/modules/cache.nix
@@ -0,0 +1,61 @@
+{
+ lib,
+ fleetLib,
+ config,
+ ...
+}:
+let
+ inherit (lib.options) mkOption mkEnableOption literalExpression;
+ inherit (lib.types)
+ enum
+ str
+ bool
+ submodule
+ ;
+ inherit (fleetLib.options) mkHostsOption;
+
+ _file = ./cache.nix;
+
+ cacheOptions = {
+ enable = mkEnableOption "fleet binary cache";
+ kind = mkOption {
+ description = ''
+ Kind of the binary cache to use.
+ '';
+ type = enum [ "attic" ];
+ default = "attic";
+ };
+ name = mkOption {
+ description = ''
+ Name of the binary cache.
+ '';
+ type = str;
+ };
+ pin = mkOption {
+ description = ''
+ Whether to pin pushed paths in the attic cache, preventing them from
+ being garbage collected.
+
+ Requires compatible fork of attic server+client: https://github.com/zhaofengli/attic/pull/227#issuecomment-4759154796
+ '';
+ type = bool;
+ default = false;
+ };
+ };
+in
+{
+ options = {
+ cache = cacheOptions;
+ hosts = mkHostsOption {
+ inherit _file;
+ options.cache = mkOption {
+ description = ''
+ Binary cache configuration for this host.
+ '';
+ type = submodule { options = cacheOptions; };
+ default = config.cache;
+ defaultText = literalExpression "fleetConfiguration.cache";
+ };
+ };
+ };
+}
--- a/modules/module-list.nix
+++ b/modules/module-list.nix
@@ -1,5 +1,6 @@
[
./assertions.nix
+ ./cache.nix
./fleetLib.nix
./hosts.nix
./nixos.nix
--- a/pkgs/default.nix
+++ b/pkgs/default.nix
@@ -12,8 +12,6 @@
fleet-generator-helper = callPackage ./fleet-generator-helper.nix { inherit craneLib; };
inherit remowt-agents-bundle;
- remowt-plugin-fleet = callPackage ./remowt-plugin-fleet.nix {
- inherit craneLib inputs remowt-agents-bundle;
- };
+ remowt-plugin-fleet = callPackage ./remowt-plugin-fleet.nix { inherit craneLib inputs; };
remowt-ssh = callPackage ./remowt-ssh.nix { inherit craneLib remowt-agents-bundle; };
}
--- a/remowt/cmds/remowt-agent/src/main.rs
+++ b/remowt/cmds/remowt-agent/src/main.rs
@@ -273,10 +273,11 @@
let helper = SocketHelper {
fallback: SuidHelper,
};
- register_auth_agent(&system_conn, Agent::new(helper, RofiPrompter)).await?;
+ let rofi = RofiPrompter::default();
+ register_auth_agent(&system_conn, Agent::new(helper, rofi.clone())).await?;
let mut rpc = Rpc::<BifConfig>::new(Address::User);
- serve_prompts(&mut rpc, RofiPrompter);
+ serve_prompts(&mut rpc, rofi);
gateway::serve(rpc.clone(), &gateway::local_socket()?).await?;
--- a/remowt/cmds/remowt-ssh/Cargo.toml
+++ b/remowt/cmds/remowt-ssh/Cargo.toml
@@ -8,6 +8,7 @@
[dependencies]
clap = { workspace = true, features = ["derive"] }
tracing-subscriber.workspace = true
+tracing-journald.workspace = true
remowt-link-shared.workspace = true
remowt-client.workspace = true
tokio = { workspace = true, features = [
--- a/remowt/cmds/remowt-ssh/src/main.rs
+++ b/remowt/cmds/remowt-ssh/src/main.rs
@@ -20,6 +20,8 @@
use tokio::io::{AsyncRead, ReadBuf};
use tokio::signal::unix::{signal, SignalKind};
use tracing::debug;
+use tracing_subscriber::prelude::*;
+use tracing_subscriber::EnvFilter;
#[derive(Parser)]
enum Opts {
@@ -45,10 +47,16 @@
#[tokio::main(flavor = "current_thread")]
async fn main() -> anyhow::Result<()> {
- tracing_subscriber::fmt()
- .with_writer(std::io::stderr)
- .without_time()
- .init();
+ // Log to the journal, not stderr: this is an interactive client that puts
+ // the terminal in raw mode, so anything on stderr corrupts the session
+ // output. If the journal socket is unavailable, drop logs rather than fall
+ // back to stderr.
+ if let Ok(journald) = tracing_journald::layer() {
+ tracing_subscriber::registry()
+ .with(EnvFilter::from_default_env())
+ .with(journald.with_syslog_identifier("remowt-ssh".to_owned()))
+ .init();
+ }
let opts = Opts::parse();
let bundle = AgentBundle::from_dir(agents_dir()?)?;
--- a/remowt/crates/remowt-ui-prompt/Cargo.toml
+++ b/remowt/crates/remowt-ui-prompt/Cargo.toml
@@ -12,5 +12,5 @@
remowt-link-shared.workspace = true
serde.workspace = true
thiserror.workspace = true
-tokio = { workspace = true, features = ["io-util", "macros", "process", "rt"] }
+tokio = { workspace = true, features = ["io-util", "macros", "process", "rt", "sync"] }
tracing.workspace = true
--- a/remowt/crates/remowt-ui-prompt/src/auto.rs
+++ b/remowt/crates/remowt-ui-prompt/src/auto.rs
@@ -24,7 +24,7 @@
};
Self {
remote,
- fallback: RofiPrompter,
+ fallback: RofiPrompter::default(),
}
}
--- a/remowt/crates/remowt-ui-prompt/src/rofi.rs
+++ b/remowt/crates/remowt-ui-prompt/src/rofi.rs
@@ -1,13 +1,18 @@
use std::process::Stdio;
+use std::sync::Arc;
use tokio::io::AsyncWriteExt;
use tokio::process::Command;
+use tokio::sync::Mutex;
use tracing::trace;
use crate::{Error, Prompter, Result, Source};
-#[derive(Clone)]
-pub struct RofiPrompter;
+#[derive(Clone, Default)]
+pub struct RofiPrompter {
+ // Rofi can't run concurrently; serialize invocations.
+ lock: Arc<Mutex<()>>,
+}
fn fixup_prompt(prompt: &str) -> &str {
// Rofi always appends such suffix
@@ -27,6 +32,7 @@
source: &[Source],
) -> Result<u32> {
trace!("rofi radio");
+ let _guard = self.lock.lock().await;
let mut cmd = rofi_command();
let mesg = if source.is_empty() {
description.to_owned()
@@ -113,6 +119,7 @@
source: &[Source],
) -> Result<String> {
trace!("rofi text");
+ let _guard = self.lock.lock().await;
let mut cmd = rofi_command();
let mesg = if source.is_empty() {
description.to_owned()
@@ -162,6 +169,7 @@
async fn display_text(&self, error: bool, description: &str, source: &[Source]) -> Result<()> {
trace!("rofi display");
+ let _guard = self.lock.lock().await;
let mut cmd = rofi_command();
let mut mesg = if source.is_empty() {
description.to_owned()
@@ -206,7 +214,7 @@
#[ignore = "interactive"]
async fn test() {
let prompter = PrependSourcePrompter {
- prompter: RofiPrompter,
+ prompter: RofiPrompter::default(),
description: "test".to_owned(),
source: vec![Source(Cow::Borrowed("ssh"))],
};