git.delta.rocks / fleet / refs/commits / 78196e185cf1

difftreelog

feat attic

ukmvyzywYaroslav Bolyukin2026-06-19parent: #2438e2b.patch.diff

22 files changed

modifiedCargo.lockdiffbeforeafterboth
--- 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"
modifiedCargo.tomldiffbeforeafterboth
--- 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
modifiedcmds/fleet/src/cmds/build_systems.rsdiffbeforeafterboth
--- 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),
addedcmds/fleet/src/log_tree.rsdiffbeforeafterboth
--- /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());
+	}
+}
modifiedcmds/fleet/src/main.rsdiffbeforeafterboth
--- 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<()> {
modifiedcrates/fleet-base/src/deploy.rsdiffbeforeafterboth
--- a/crates/fleet-base/src/deploy.rs
+++ b/crates/fleet-base/src/deploy.rs
@@ -290,8 +290,10 @@
 	if !host.local {
 		info!("uploading system closure");
 		let mut tries = 0;
+		// TODO: Use gc prefix?..
 		loop {
-			match host.remote_derivation(&generation).await {
+			let name = format!("{}-system", host.gc_root_prefix());
+			match host.remote_derivation(&generation, Some(&name)).await {
 				Ok(remote) => {
 					assert!(remote == generation, "CA derivations aren't implemented");
 					return Ok(remote);
modifiedcrates/fleet-base/src/host.rsdiffbeforeafterboth
--- 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,
modifiedcrates/fleet-base/src/keys.rsdiffbeforeafterboth
--- 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);
modifiedcrates/fleet-base/src/primops.rsdiffbeforeafterboth
--- 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")?;
 
modifiedcrates/nix-eval/src/lib.rsdiffbeforeafterboth
--- 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,
modifiedcrates/nix-eval/src/logging.rsdiffbeforeafterboth
--- 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!(),
modifiedcrates/nix-eval/src/scheduler.rsdiffbeforeafterboth
--- 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")
 }
modifiedcrates/remowt-fleet/src/lib.rsdiffbeforeafterboth
--- 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,
addedmodules/cache.nixdiffbeforeafterboth
--- /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";
+      };
+    };
+  };
+}
modifiedmodules/module-list.nixdiffbeforeafterboth
--- a/modules/module-list.nix
+++ b/modules/module-list.nix
@@ -1,5 +1,6 @@
 [
   ./assertions.nix
+  ./cache.nix
   ./fleetLib.nix
   ./hosts.nix
   ./nixos.nix
modifiedpkgs/default.nixdiffbeforeafterboth
--- 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; };
 }
modifiedremowt/cmds/remowt-agent/src/main.rsdiffbeforeafterboth
--- 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?;
 
modifiedremowt/cmds/remowt-ssh/Cargo.tomldiffbeforeafterboth
--- 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 = [
modifiedremowt/cmds/remowt-ssh/src/main.rsdiffbeforeafterboth
--- 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()?)?;
modifiedremowt/crates/remowt-ui-prompt/Cargo.tomldiffbeforeafterboth
--- 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
modifiedremowt/crates/remowt-ui-prompt/src/auto.rsdiffbeforeafterboth
--- 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(),
 		}
 	}
 
modifiedremowt/crates/remowt-ui-prompt/src/rofi.rsdiffbeforeafterboth
before · remowt/crates/remowt-ui-prompt/src/rofi.rs
1use std::process::Stdio;23use tokio::io::AsyncWriteExt;4use tokio::process::Command;5use tracing::trace;67use crate::{Error, Prompter, Result, Source};89#[derive(Clone)]10pub struct RofiPrompter;1112fn fixup_prompt(prompt: &str) -> &str {13	// Rofi always appends such suffix14	prompt.strip_suffix(": ").unwrap_or(prompt)15}1617fn rofi_command() -> Command {18	Command::new(option_env!("ROFI").unwrap_or("rofi"))19}2021impl Prompter for RofiPrompter {22	async fn prompt_enum(23		&self,24		prompt: &str,25		description: &str,26		variants: &[&str],27		source: &[Source],28	) -> Result<u32> {29		trace!("rofi radio");30		let mut cmd = rofi_command();31		let mesg = if source.is_empty() {32			description.to_owned()33		} else {34			let mut out = format!("{description}\n\n<b>Requested on ",);35			for (i, s) in source.iter().enumerate() {36				if i != 0 {37					out.push_str(" -> ");38				}39				out.push_str(&s.to_string());40			}41			out.push_str("</b>");42			out43		};44		cmd.args([45			"-dmenu",46			"-mesg",47			&mesg,48			"-sync",49			"-no-custom",50			"-p",51			fixup_prompt(prompt),52			"-format",53			"i",54			"-markup-rows",55		]);56		cmd.stdin(Stdio::piped());57		cmd.stdout(Stdio::piped());58		cmd.kill_on_drop(true);59		let mut child = cmd60			.spawn()61			.map_err(|e| Error::InputError(format!("failed to spawn rofi: {e}")))?;6263		let mut stdin = child.stdin.take().expect("stdin is piped");64		for var in variants {65			stdin66				.write_all(var.replace('\n', " ").as_bytes())67				.await68				.map_err(|e| Error::InputError(format!("failed to write rofi variants: {e}")))?;69			stdin70				.write_all(b"\n")71				.await72				.map_err(|e| Error::InputError(format!("failed to write rofi variants: {e}")))?;73		}74		// write_all already flushes, just to be sure.75		let _ = stdin.flush().await;76		drop(stdin);7778		let out = child79			.wait_with_output()80			.await81			.map_err(|e| Error::InputError(format!("failed to wait for rofi: {e}")))?;82		match out.status.code() {83			Some(0) => {}84			Some(1) => return Err(Error::Cancel),85			other => {86				return Err(Error::InputError(format!(87					"rofi exited with status {other:?}"88				)));89			}90		}91		let stdout = out92			.stdout93			.strip_suffix(b"\n")94			.unwrap_or(&out.stdout)95			.to_owned();9697		let id: u32 = String::from_utf8(stdout)98			.map_err(|e| Error::InputError(format!("rofi produced invalid output: {e}")))?99			.parse()100			.map_err(|e| Error::InputError(format!("rofi produced invalid output: {e}")))?;101		if id as usize >= variants.len() {102			return Err(Error::InputError("invalid rofi response".to_owned()));103		}104105		Ok(id)106	}107108	async fn prompt_text(109		&self,110		echo: bool,111		prompt: &str,112		description: &str,113		source: &[Source],114	) -> Result<String> {115		trace!("rofi text");116		let mut cmd = rofi_command();117		let mesg = if source.is_empty() {118			description.to_owned()119		} else {120			let mut out = format!("{description}\n\n<b>Requested on ",);121			for (i, s) in source.iter().enumerate() {122				if i != 0 {123					out.push_str(" -> ");124				}125				out.push_str(&s.to_string());126			}127			out.push_str("</b>");128			out129		};130		cmd.args(["-dmenu", "-mesg", &mesg, "-p", fixup_prompt(prompt)]);131		if !echo {132			cmd.arg("-password");133		}134		cmd.stdin(Stdio::null());135		cmd.stdout(Stdio::piped());136		cmd.kill_on_drop(true);137		let child = cmd138			.spawn()139			.map_err(|e| Error::InputError(format!("failed to spawn rofi: {e}")))?;140141		let out = child142			.wait_with_output()143			.await144			.map_err(|e| Error::InputError(format!("failed to wait for rofi: {e}")))?;145		match out.status.code() {146			Some(0) => {}147			Some(1) => return Err(Error::Cancel),148			other => {149				return Err(Error::InputError(format!(150					"rofi exited with status {other:?}"151				)));152			}153		}154		let stdout = out155			.stdout156			.strip_suffix(b"\n")157			.unwrap_or(&out.stdout)158			.to_owned();159160		Ok(String::from_utf8_lossy(&stdout).to_string())161	}162163	async fn display_text(&self, error: bool, description: &str, source: &[Source]) -> Result<()> {164		trace!("rofi display");165		let mut cmd = rofi_command();166		let mut mesg = if source.is_empty() {167			description.to_owned()168		} else {169			let mut out = format!("{description}\n\n<b>Coming from ",);170			for s in source.iter() {171				out.push_str(&s.to_string());172			}173			out.push_str("</b>");174			out175		};176		if error {177			mesg.insert_str(0, "<span color=\"red\">");178			mesg.push_str("</span>");179		}180		cmd.args(["-e", &mesg, "-markup"]);181		cmd.stdin(Stdio::null());182		cmd.stdout(Stdio::null());183		cmd.kill_on_drop(true);184		let mut child = cmd185			.spawn()186			.map_err(|e| Error::InputError(format!("failed to spawn rofi: {e}")))?;187188		child189			.wait()190			.await191			.map_err(|e| Error::InputError(format!("failed to wait for rofi: {e}")))?;192193		Ok(())194	}195}196197#[cfg(test)]198mod tests {199	use std::borrow::Cow;200201	use crate::rofi::RofiPrompter;202	use crate::{PrependSourcePrompter, Prompter as _, Source};203204	// #[tokio::test]205	#[tokio::test]206	#[ignore = "interactive"]207	async fn test() {208		let prompter = PrependSourcePrompter {209			prompter: RofiPrompter,210			description: "test".to_owned(),211			source: vec![Source(Cow::Borrowed("ssh"))],212		};213		prompter214			.prompt_radio("Enable", "Polkit needs access", &[])215			.await216			.expect("rofi");217		prompter218			.prompt_text(false, "Password", "Polkit needs access", &[])219			.await220			.expect("rofi");221		prompter222			.display_text(true, "Polkit needs access", &[])223			.await224			.expect("rofi");225	}226}
after · remowt/crates/remowt-ui-prompt/src/rofi.rs
1use std::process::Stdio;2use std::sync::Arc;34use tokio::io::AsyncWriteExt;5use tokio::process::Command;6use tokio::sync::Mutex;7use tracing::trace;89use crate::{Error, Prompter, Result, Source};1011#[derive(Clone, Default)]12pub struct RofiPrompter {13	// Rofi can't run concurrently; serialize invocations.14	lock: Arc<Mutex<()>>,15}1617fn fixup_prompt(prompt: &str) -> &str {18	// Rofi always appends such suffix19	prompt.strip_suffix(": ").unwrap_or(prompt)20}2122fn rofi_command() -> Command {23	Command::new(option_env!("ROFI").unwrap_or("rofi"))24}2526impl Prompter for RofiPrompter {27	async fn prompt_enum(28		&self,29		prompt: &str,30		description: &str,31		variants: &[&str],32		source: &[Source],33	) -> Result<u32> {34		trace!("rofi radio");35		let _guard = self.lock.lock().await;36		let mut cmd = rofi_command();37		let mesg = if source.is_empty() {38			description.to_owned()39		} else {40			let mut out = format!("{description}\n\n<b>Requested on ",);41			for (i, s) in source.iter().enumerate() {42				if i != 0 {43					out.push_str(" -> ");44				}45				out.push_str(&s.to_string());46			}47			out.push_str("</b>");48			out49		};50		cmd.args([51			"-dmenu",52			"-mesg",53			&mesg,54			"-sync",55			"-no-custom",56			"-p",57			fixup_prompt(prompt),58			"-format",59			"i",60			"-markup-rows",61		]);62		cmd.stdin(Stdio::piped());63		cmd.stdout(Stdio::piped());64		cmd.kill_on_drop(true);65		let mut child = cmd66			.spawn()67			.map_err(|e| Error::InputError(format!("failed to spawn rofi: {e}")))?;6869		let mut stdin = child.stdin.take().expect("stdin is piped");70		for var in variants {71			stdin72				.write_all(var.replace('\n', " ").as_bytes())73				.await74				.map_err(|e| Error::InputError(format!("failed to write rofi variants: {e}")))?;75			stdin76				.write_all(b"\n")77				.await78				.map_err(|e| Error::InputError(format!("failed to write rofi variants: {e}")))?;79		}80		// write_all already flushes, just to be sure.81		let _ = stdin.flush().await;82		drop(stdin);8384		let out = child85			.wait_with_output()86			.await87			.map_err(|e| Error::InputError(format!("failed to wait for rofi: {e}")))?;88		match out.status.code() {89			Some(0) => {}90			Some(1) => return Err(Error::Cancel),91			other => {92				return Err(Error::InputError(format!(93					"rofi exited with status {other:?}"94				)));95			}96		}97		let stdout = out98			.stdout99			.strip_suffix(b"\n")100			.unwrap_or(&out.stdout)101			.to_owned();102103		let id: u32 = String::from_utf8(stdout)104			.map_err(|e| Error::InputError(format!("rofi produced invalid output: {e}")))?105			.parse()106			.map_err(|e| Error::InputError(format!("rofi produced invalid output: {e}")))?;107		if id as usize >= variants.len() {108			return Err(Error::InputError("invalid rofi response".to_owned()));109		}110111		Ok(id)112	}113114	async fn prompt_text(115		&self,116		echo: bool,117		prompt: &str,118		description: &str,119		source: &[Source],120	) -> Result<String> {121		trace!("rofi text");122		let _guard = self.lock.lock().await;123		let mut cmd = rofi_command();124		let mesg = if source.is_empty() {125			description.to_owned()126		} else {127			let mut out = format!("{description}\n\n<b>Requested on ",);128			for (i, s) in source.iter().enumerate() {129				if i != 0 {130					out.push_str(" -> ");131				}132				out.push_str(&s.to_string());133			}134			out.push_str("</b>");135			out136		};137		cmd.args(["-dmenu", "-mesg", &mesg, "-p", fixup_prompt(prompt)]);138		if !echo {139			cmd.arg("-password");140		}141		cmd.stdin(Stdio::null());142		cmd.stdout(Stdio::piped());143		cmd.kill_on_drop(true);144		let child = cmd145			.spawn()146			.map_err(|e| Error::InputError(format!("failed to spawn rofi: {e}")))?;147148		let out = child149			.wait_with_output()150			.await151			.map_err(|e| Error::InputError(format!("failed to wait for rofi: {e}")))?;152		match out.status.code() {153			Some(0) => {}154			Some(1) => return Err(Error::Cancel),155			other => {156				return Err(Error::InputError(format!(157					"rofi exited with status {other:?}"158				)));159			}160		}161		let stdout = out162			.stdout163			.strip_suffix(b"\n")164			.unwrap_or(&out.stdout)165			.to_owned();166167		Ok(String::from_utf8_lossy(&stdout).to_string())168	}169170	async fn display_text(&self, error: bool, description: &str, source: &[Source]) -> Result<()> {171		trace!("rofi display");172		let _guard = self.lock.lock().await;173		let mut cmd = rofi_command();174		let mut mesg = if source.is_empty() {175			description.to_owned()176		} else {177			let mut out = format!("{description}\n\n<b>Coming from ",);178			for s in source.iter() {179				out.push_str(&s.to_string());180			}181			out.push_str("</b>");182			out183		};184		if error {185			mesg.insert_str(0, "<span color=\"red\">");186			mesg.push_str("</span>");187		}188		cmd.args(["-e", &mesg, "-markup"]);189		cmd.stdin(Stdio::null());190		cmd.stdout(Stdio::null());191		cmd.kill_on_drop(true);192		let mut child = cmd193			.spawn()194			.map_err(|e| Error::InputError(format!("failed to spawn rofi: {e}")))?;195196		child197			.wait()198			.await199			.map_err(|e| Error::InputError(format!("failed to wait for rofi: {e}")))?;200201		Ok(())202	}203}204205#[cfg(test)]206mod tests {207	use std::borrow::Cow;208209	use crate::rofi::RofiPrompter;210	use crate::{PrependSourcePrompter, Prompter as _, Source};211212	// #[tokio::test]213	#[tokio::test]214	#[ignore = "interactive"]215	async fn test() {216		let prompter = PrependSourcePrompter {217			prompter: RofiPrompter::default(),218			description: "test".to_owned(),219			source: vec![Source(Cow::Borrowed("ssh"))],220		};221		prompter222			.prompt_radio("Enable", "Polkit needs access", &[])223			.await224			.expect("rofi");225		prompter226			.prompt_text(false, "Password", "Polkit needs access", &[])227			.await228			.expect("rofi");229		prompter230			.display_text(true, "Polkit needs access", &[])231			.await232			.expect("rofi");233	}234}