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<()> {
--- 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);
--- 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);
1use std::collections::{BTreeMap, BTreeSet, HashMap};2use std::sync::{Arc, OnceLock};34use anyhow::{Context, bail, ensure};5use fleet_shared::SecretData;6use itertools::Itertools;7use nix_eval::{NativeFn, Value, await_in_nix, nix_go, nix_go_json};8use remowt_endpoints::fs::FsClient;9use remowt_link_shared::BifConfig;10use serde::Deserialize;11use tracing::{info, warn};1213use crate::fleetdata::{14 Expectations, FleetSecretData, FleetSecretDistribution, FleetSecretPart, GeneratorPart,15 RegenerationConstraints, SecretOwner,16};17use crate::host::{Config, ConfigHost};18use anyhow::{Result, anyhow};1920pub static PRIMOPS_DATA: OnceLock<Config> = OnceLock::new();2122#[derive(Deserialize)]23#[serde(rename_all = "camelCase")]24enum GeneratorKind {25 Impure,26 Pure,27}2829pub fn get_pkgs_and_generators(host_on: &ConfigHost, recipients: Vec<String>) -> Result<Value> {30 let pkgs = host_on.pkgs()?;31 let default_mk_secret_generators = nix_go!(pkgs.mkSecretGenerators);32 let generators = nix_go!(default_mk_secret_generators(Obj { recipients }));33 pkgs.clone().attrs_update(generators)34}35pub fn get_default_pkgs_and_generators(config: &Config) -> Result<Value> {36 let host_on = config.local_host();37 get_pkgs_and_generators(&host_on, vec![])38}39pub fn call_package(config: &Config, pkgs: &Value, package: &Value) -> Result<Value> {40 ensure!(41 package.is_function(),42 "package should be a function to be called with callPackage"43 );44 45 let nixpkgs = &config.nixpkgs;46 let call_package = nix_go!(nixpkgs.lib.callPackageWith(pkgs));47 Ok(nix_go!(call_package(package)(Obj {})))48}4950pub fn get_default_generator_drv(config: &Config, generator: &Value) -> Result<Value> {51 let default_pkgs_and_generators = get_default_pkgs_and_generators(config)?;52 let default_generator_drv = call_package(config, &default_pkgs_and_generators, generator)53 .context("failed to initialize generator to get metadata")?;5455 Ok(default_generator_drv)56}5758fn secret_to_parts(59 secret_name: &str,60 secret: &BTreeMap<String, FleetSecretPart>,61 expected: &BTreeMap<String, GeneratorPart>,62) -> Value {63 let mut out = HashMap::new();64 for (part_name, part) in secret {65 if !expected.contains_key(part_name) {66 warn!(67 "secret {secret_name} part {part_name} is stored, but not defined in nixos config, it will not be passed to nix"68 );69 continue;70 };71 out.insert(72 part_name.as_str(),73 Value::new_attrs(HashMap::from_iter([(74 "raw",75 Value::new_str(&part.raw.to_string()),76 )])),77 );78 }7980 Value::new_attrs(out)81}8283pub async fn generate(84 config: &Config,85 expectations: Expectations,86 generator: &Value,87 default_generator_drv: &Value,88) -> Result<FleetSecretDistribution> {89 let kind: GeneratorKind = nix_go_json!(default_generator_drv.generatorKind);9091 match kind {92 GeneratorKind::Impure => {93 let impure_on: Option<String> = nix_go_json!(default_generator_drv.impureOn);9495 let host_on = if let Some(on) = &impure_on {96 Arc::new(97 config98 .host(on)99 .context("failed to get secret generation target host")?,100 )101 } else {102 config.local_host()103 };104105 let remowt = host_on.remowt().await?;106 let fs = remowt.endpoints::<FsClient<BifConfig>>();107108 let mut recipients = Vec::new();109 for owner in &expectations.owners {110 recipients.push(config.key(owner).await?);111 }112 let pkgs_and_generators = get_pkgs_and_generators(&host_on, recipients)113 .context("failed to get pkgs for target host")?;114 let generator = call_package(config, &pkgs_and_generators, generator)115 .context("failed to evaluate generator for target host")?;116117 let generator = generator118 .build("out")119 .context("failed to build generator for target host")?;120121 let generator = host_on122 .remote_derivation(&generator)123 .await124 .context("failed to copy generator to target host")?;125126 127 let out_parent = fs128 .mktemp_dir()129 .await130 .map_err(|e| anyhow!("{e:?}"))131 .context("failed to prepare generator output dir on target host")?;132 let out = out_parent.join("out");133 let mut generator_cmd = remowt.cmd(generator);134 generator_cmd.env("out", &out);135 if impure_on.is_none() {136 let project_path: String = config137 .directory138 .clone()139 .into_os_string()140 .into_string()141 .map_err(|e| anyhow!("fleet project path is not utf-8: {e:?}"))?;142 generator_cmd.env("FLEET_PROJECT", project_path);143 };144 generator_cmd145 .run()146 .await147 .context("failed to run impure generator")?;148149 {150 let marker = fs.read_file_text(out.join("marker")).await?;151 ensure!(152 marker == "SUCCESS",153 "impure generator ended prematurely, secret generation failed"154 );155 }156157 let mut missing_parts = expectations.parts.clone();158 let mut parts = BTreeMap::new();159 for part in fs.read_dir(&out).await? {160 let part = part.into_string();161 if part == "created_at" || part == "expires_at" || part == "marker" {162 continue;163 }164 let Some(part_def) = missing_parts.remove(&part) else {165 bail!("secret generator has produced an unexpected part: {part}");166 };167 let contents: SecretData = fs168 .read_file_text(out.join("part"))169 .await?170 .parse()171 .map_err(|e| anyhow!("failed to decode secret {out:?} part {part:?}: {e}"))?;172173 ensure!(174 contents.encrypted == part_def.encrypted,175 "part {part} produced by generator is supposed to be {}, but it was not",176 if part_def.encrypted {177 "encrypted"178 } else {179 "plaintext"180 }181 );182183 parts.insert(part.to_owned(), FleetSecretPart { raw: contents });184 }185 if !missing_parts.is_empty() {186 bail!(187 "the secret generator has not produced the following expected parts: {missing_parts:?}"188 );189 }190191 let created_at = fs.read_file_value(out.join("created_at")).await??;192 let expires_at = match fs.read_file_value(out.join("expires_at")).await {193 Ok(v) => Some(v?),194 Err(remowt_endpoints::fs::Error::NotFound) => None,195 Err(e) => return Err(e.into()),196 };197198 let new_data = FleetSecretData {199 created_at,200 expires_at,201 parts,202 generation_data: expectations.generation_data.clone(),203 };204205 let new_data =206 FleetSecretDistribution::new(expectations.owners.clone(), new_data, config.now);207208 Ok(new_data)209 }210 GeneratorKind::Pure => {211 bail!("pure generators are disabled for now")212 }213 }214}215216pub fn init_primops() {217 NativeFn::new(218 c"__fleetEnsureHostSecrets",219 c"Ensure no extra secrets are stored for the host, pruning unknown",220 [c"host", c"expectedNonshared", c"expectedShared", c"rest"],221 |_es, [host, expected_nonshared, expected_shared, rest]| {222 let host = SecretOwner::host(host.to_string()?);223 let expected_nonshared: BTreeSet<String> = expected_nonshared.as_json()?;224 let expected_shared: BTreeSet<String> = expected_shared.as_json()?;225226 let mut expected = expected_nonshared;227 expected.extend(expected_shared);228229 let config = PRIMOPS_DATA230 .get()231 .expect("primops data should be set on init");232233 config234 .data235 .secrets236 .write()237 .expect("no poisoning")238 .prune_host(&host, expected);239240 Ok(rest.clone())241 },242 )243 .register();244 NativeFn::new(245 c"__fleetEnsureHostSecret",246 c"Ensure secret existence for a host, regenerating it in case of some mismatch",247 [c"host", c"secret", c"generator"],248 |es, [host, secret, generator]| {249 let host = SecretOwner::host(&host.to_string()?);250 let secret = secret.to_string()?;251252 let config = PRIMOPS_DATA253 .get()254 .expect("primops data should be set on init");255256 let shared_def = config.secret_definition(&secret).context("failed to get shared secret definition")?;257258 let (shared, generator, expected_owners) = if generator.is_string() {259 assert_eq!(generator.to_string()?, "shared", "asserted by nixos type system");260 let Some(shared_def) = shared_def else {261 bail!("secret {secret} is defined on host {host} as shared, but there is no shared secret with same name defined at fleetConfiguration.secrets.{secret}.generator")262 };263 let expected_owners = shared_def.expected_owners()?;264265 ensure!(expected_owners.contains(&host), "secret {secret} does not define {host} as expected owner");266267 (Some(shared_def.clone()), shared_def.generator()?, expected_owners)268 } else {269 if shared_def.is_some() {270 bail!("hosts can only have their own generators for non-shared secrets, either set host secret generator to \"shared\", or remove shared secret generator at fleetConfiguration.secrets.{secret}.generator")271 }272273 (None, generator.clone(), BTreeSet::from_iter([host.clone()]))274 };275276 let default_generator_drv = get_default_generator_drv(config, &generator)?;277 let mut expectations = Expectations {278 parts: nix_go_json!(default_generator_drv.parts),279 generation_data: nix_go_json!(default_generator_drv.generationData),280 owners: expected_owners.clone(),281 };282 let constraints = if let Some(shared) = &shared{283 RegenerationConstraints {284 allow_different: nix_go_json!(default_generator_drv.allowDifferent) && shared.allow_different()?,285 regenerate_on_owner_added: shared.regenerate_on_owner_added()?,286 regenerate_on_owner_removed: shared.regenerate_on_owner_added()?,287 }288 } else {289 RegenerationConstraints::host_personal()290 };291292 let mut secrets = config.data.secrets.write().expect("no poisoning");293 let dists = secrets.get_or_create(&secret);294295 if shared.is_some() {296 dists.prune_shared(&expected_owners, !constraints.allow_different, &expectations.parts, &expectations.generation_data, constraints.regenerate_on_owner_removed, constraints.regenerate_on_owner_added, &config.prefer_identities, config.now);297 } else {298 dists.prune_host(host.clone(), &expectations.parts, &expectations.generation_data, config.now);299 };300301 if let Some(dist) = dists.get(&host) {302 return Ok(secret_to_parts(&secret, &dist.secret.parts, &expectations.parts));303 };304305 let mut reencrypt_targets = expectations.owners.clone();306 for dist in dists.distributions() {307 for own in dist.owners() {308 reencrypt_targets.remove(own);309 }310 }311 if !constraints.regenerate_on_owner_added {312 if let Some(unpruned) = dists.try_unprune(host.clone()) {313 return Ok(secret_to_parts(&secret, &unpruned.secret.parts, &expectations.parts));314 } else if let Some(best) = dists.best_distribution_for_reencryption(&config.prefer_identities) {315 let new_owners = reencrypt_targets.clone();316 let mut reencrypt_targets = reencrypt_targets;317 reencrypt_targets.extend(best.owners().cloned());318319 let mut preferred = best.owners().collect_vec();320 preferred.sort_by_key(|v| !config.prefer_identities.contains(*v));321322 warn!("reencrypting secret {secret} as it is missing for host {host}");323324 for owner in preferred {325 if let Some(hostname) = owner.as_host() && let Ok(host) = config.host(hostname) {326 let best = best.clone();327 let reencrypt_targets = reencrypt_targets.clone();328 let reencrypted = match await_in_nix(async move {329 host.reencrypt_distribution(&best, reencrypt_targets.clone(), config.now).await330 }) {331 Ok(r) => r,332 Err(e) => {333 warn!("reencryption failed on {hostname}: {e:?}");334 continue;335 }336 };337 dists.extend(reencrypted.clone(), format!("secret was reencrypted to extend with new owners: {new_owners:?}"));338 return Ok(secret_to_parts(&secret, &reencrypted.secret.parts, &expectations.parts));339 };340 }341 warn!("failed to reencrypt using any host")342 };343 };344345 if constraints.allow_different {346 for dist in dists.distributions() {347 for own in dist.owners() {348 expectations.owners.remove(own);349 }350 }351 }352 info!("secret {secret} is being generated for {:?}", expectations.owners);353354 let expectations_ = expectations.clone();355 let generated = await_in_nix(async move {356 generate(config, expectations_, &generator, &default_generator_drv).await357 })?;358359 dists.extend(generated.clone(), "secret was generated".to_string());360361 Ok(secret_to_parts(&secret, &generated.secret.parts, &expectations.parts))362 },363 )364 .register();365}
1use std::collections::{BTreeMap, BTreeSet, HashMap};2use std::sync::{Arc, OnceLock};34use anyhow::{Context, bail, ensure};5use fleet_shared::SecretData;6use itertools::Itertools;7use nix_eval::{NativeFn, Value, await_in_nix, nix_go, nix_go_json};8use remowt_endpoints::fs::FsClient;9use remowt_link_shared::BifConfig;10use serde::Deserialize;11use tracing::{info, warn};1213use crate::fleetdata::{14 Expectations, FleetSecretData, FleetSecretDistribution, FleetSecretPart, GeneratorPart,15 RegenerationConstraints, SecretOwner,16};17use crate::host::{Config, ConfigHost};18use anyhow::{Result, anyhow};1920pub static PRIMOPS_DATA: OnceLock<Config> = OnceLock::new();2122#[derive(Deserialize)]23#[serde(rename_all = "camelCase")]24enum GeneratorKind {25 Impure,26 Pure,27}2829pub fn get_pkgs_and_generators(host_on: &ConfigHost, recipients: Vec<String>) -> Result<Value> {30 let pkgs = host_on.pkgs()?;31 let default_mk_secret_generators = nix_go!(pkgs.mkSecretGenerators);32 let generators = nix_go!(default_mk_secret_generators(Obj { recipients }));33 pkgs.clone().attrs_update(generators)34}35pub fn get_default_pkgs_and_generators(config: &Config) -> Result<Value> {36 let host_on = config.local_host();37 get_pkgs_and_generators(&host_on, vec![])38}39pub fn call_package(config: &Config, pkgs: &Value, package: &Value) -> Result<Value> {40 ensure!(41 package.is_function(),42 "package should be a function to be called with callPackage"43 );44 45 let nixpkgs = &config.nixpkgs;46 let call_package = nix_go!(nixpkgs.lib.callPackageWith(pkgs));47 Ok(nix_go!(call_package(package)(Obj {})))48}4950pub fn get_default_generator_drv(config: &Config, generator: &Value) -> Result<Value> {51 let default_pkgs_and_generators = get_default_pkgs_and_generators(config)?;52 let default_generator_drv = call_package(config, &default_pkgs_and_generators, generator)53 .context("failed to initialize generator to get metadata")?;5455 Ok(default_generator_drv)56}5758fn secret_to_parts(59 secret_name: &str,60 secret: &BTreeMap<String, FleetSecretPart>,61 expected: &BTreeMap<String, GeneratorPart>,62) -> Value {63 let mut out = HashMap::new();64 for (part_name, part) in secret {65 if !expected.contains_key(part_name) {66 warn!(67 "secret {secret_name} part {part_name} is stored, but not defined in nixos config, it will not be passed to nix"68 );69 continue;70 };71 out.insert(72 part_name.as_str(),73 Value::new_attrs(HashMap::from_iter([(74 "raw",75 Value::new_str(&part.raw.to_string()),76 )])),77 );78 }7980 Value::new_attrs(out)81}8283pub async fn generate(84 config: &Config,85 expectations: Expectations,86 generator: &Value,87 default_generator_drv: &Value,88) -> Result<FleetSecretDistribution> {89 let kind: GeneratorKind = nix_go_json!(default_generator_drv.generatorKind);9091 match kind {92 GeneratorKind::Impure => {93 let impure_on: Option<String> = nix_go_json!(default_generator_drv.impureOn);9495 let host_on = if let Some(on) = &impure_on {96 Arc::new(97 config98 .host(on)99 .context("failed to get secret generation target host")?,100 )101 } else {102 config.local_host()103 };104105 let remowt = host_on.remowt().await?;106 let fs = remowt.endpoints::<FsClient<BifConfig>>();107108 let mut recipients = Vec::new();109 for owner in &expectations.owners {110 recipients.push(config.key(owner).await?);111 }112 let pkgs_and_generators = get_pkgs_and_generators(&host_on, recipients)113 .context("failed to get pkgs for target host")?;114 let generator = call_package(config, &pkgs_and_generators, generator)115 .context("failed to evaluate generator for target host")?;116117 let generator = generator118 .build("out")119 .context("failed to build generator for target host")?;120121 let generator = host_on122 .remote_derivation(&generator, None)123 .await124 .context("failed to copy generator to target host")?;125126 127 let out_parent = fs128 .mktemp_dir()129 .await130 .map_err(|e| anyhow!("{e:?}"))131 .context("failed to prepare generator output dir on target host")?;132 let out = out_parent.join("out");133 let mut generator_cmd = remowt.cmd(generator);134 generator_cmd.env("out", &out);135 if impure_on.is_none() {136 let project_path: String = config137 .directory138 .clone()139 .into_os_string()140 .into_string()141 .map_err(|e| anyhow!("fleet project path is not utf-8: {e:?}"))?;142 generator_cmd.env("FLEET_PROJECT", project_path);143 };144 generator_cmd145 .run()146 .await147 .context("failed to run impure generator")?;148149 {150 let marker = fs.read_file_text(out.join("marker")).await?;151 ensure!(152 marker == "SUCCESS",153 "impure generator ended prematurely, secret generation failed"154 );155 }156157 let mut missing_parts = expectations.parts.clone();158 let mut parts = BTreeMap::new();159 for part in fs.read_dir(&out).await? {160 let part = part.into_string();161 if part == "created_at" || part == "expires_at" || part == "marker" {162 continue;163 }164 let Some(part_def) = missing_parts.remove(&part) else {165 bail!("secret generator has produced an unexpected part: {part}");166 };167 let contents: SecretData = fs168 .read_file_text(out.join("part"))169 .await?170 .parse()171 .map_err(|e| anyhow!("failed to decode secret {out:?} part {part:?}: {e}"))?;172173 ensure!(174 contents.encrypted == part_def.encrypted,175 "part {part} produced by generator is supposed to be {}, but it was not",176 if part_def.encrypted {177 "encrypted"178 } else {179 "plaintext"180 }181 );182183 parts.insert(part.to_owned(), FleetSecretPart { raw: contents });184 }185 if !missing_parts.is_empty() {186 bail!(187 "the secret generator has not produced the following expected parts: {missing_parts:?}"188 );189 }190191 let created_at = fs.read_file_value(out.join("created_at")).await??;192 let expires_at = match fs.read_file_value(out.join("expires_at")).await {193 Ok(v) => Some(v?),194 Err(remowt_endpoints::fs::Error::NotFound) => None,195 Err(e) => return Err(e.into()),196 };197198 let new_data = FleetSecretData {199 created_at,200 expires_at,201 parts,202 generation_data: expectations.generation_data.clone(),203 };204205 let new_data =206 FleetSecretDistribution::new(expectations.owners.clone(), new_data, config.now);207208 Ok(new_data)209 }210 GeneratorKind::Pure => {211 bail!("pure generators are disabled for now")212 }213 }214}215216pub fn init_primops() {217 NativeFn::new(218 c"__fleetEnsureHostSecrets",219 c"Ensure no extra secrets are stored for the host, pruning unknown",220 [c"host", c"expectedNonshared", c"expectedShared", c"rest"],221 |_es, [host, expected_nonshared, expected_shared, rest]| {222 let host = SecretOwner::host(host.to_string()?);223 let expected_nonshared: BTreeSet<String> = expected_nonshared.as_json()?;224 let expected_shared: BTreeSet<String> = expected_shared.as_json()?;225226 let mut expected = expected_nonshared;227 expected.extend(expected_shared);228229 let config = PRIMOPS_DATA230 .get()231 .expect("primops data should be set on init");232233 config234 .data235 .secrets236 .write()237 .expect("no poisoning")238 .prune_host(&host, expected);239240 Ok(rest.clone())241 },242 )243 .register();244 NativeFn::new(245 c"__fleetEnsureHostSecret",246 c"Ensure secret existence for a host, regenerating it in case of some mismatch",247 [c"host", c"secret", c"generator"],248 |es, [host, secret, generator]| {249 let host = SecretOwner::host(&host.to_string()?);250 let secret = secret.to_string()?;251252 let config = PRIMOPS_DATA253 .get()254 .expect("primops data should be set on init");255256 let shared_def = config.secret_definition(&secret).context("failed to get shared secret definition")?;257258 let (shared, generator, expected_owners) = if generator.is_string() {259 assert_eq!(generator.to_string()?, "shared", "asserted by nixos type system");260 let Some(shared_def) = shared_def else {261 bail!("secret {secret} is defined on host {host} as shared, but there is no shared secret with same name defined at fleetConfiguration.secrets.{secret}.generator")262 };263 let expected_owners = shared_def.expected_owners()?;264265 ensure!(expected_owners.contains(&host), "secret {secret} does not define {host} as expected owner");266267 (Some(shared_def.clone()), shared_def.generator()?, expected_owners)268 } else {269 if shared_def.is_some() {270 bail!("hosts can only have their own generators for non-shared secrets, either set host secret generator to \"shared\", or remove shared secret generator at fleetConfiguration.secrets.{secret}.generator")271 }272273 (None, generator.clone(), BTreeSet::from_iter([host.clone()]))274 };275276 let default_generator_drv = get_default_generator_drv(config, &generator)?;277 let mut expectations = Expectations {278 parts: nix_go_json!(default_generator_drv.parts),279 generation_data: nix_go_json!(default_generator_drv.generationData),280 owners: expected_owners.clone(),281 };282 let constraints = if let Some(shared) = &shared{283 RegenerationConstraints {284 allow_different: nix_go_json!(default_generator_drv.allowDifferent) && shared.allow_different()?,285 regenerate_on_owner_added: shared.regenerate_on_owner_added()?,286 regenerate_on_owner_removed: shared.regenerate_on_owner_added()?,287 }288 } else {289 RegenerationConstraints::host_personal()290 };291292 let mut secrets = config.data.secrets.write().expect("no poisoning");293 let dists = secrets.get_or_create(&secret);294295 if shared.is_some() {296 dists.prune_shared(&expected_owners, !constraints.allow_different, &expectations.parts, &expectations.generation_data, constraints.regenerate_on_owner_removed, constraints.regenerate_on_owner_added, &config.prefer_identities, config.now);297 } else {298 dists.prune_host(host.clone(), &expectations.parts, &expectations.generation_data, config.now);299 };300301 if let Some(dist) = dists.get(&host) {302 return Ok(secret_to_parts(&secret, &dist.secret.parts, &expectations.parts));303 };304305 let mut reencrypt_targets = expectations.owners.clone();306 for dist in dists.distributions() {307 for own in dist.owners() {308 reencrypt_targets.remove(own);309 }310 }311 if !constraints.regenerate_on_owner_added {312 if let Some(unpruned) = dists.try_unprune(host.clone()) {313 return Ok(secret_to_parts(&secret, &unpruned.secret.parts, &expectations.parts));314 } else if let Some(best) = dists.best_distribution_for_reencryption(&config.prefer_identities) {315 let new_owners = reencrypt_targets.clone();316 let mut reencrypt_targets = reencrypt_targets;317 reencrypt_targets.extend(best.owners().cloned());318319 let mut preferred = best.owners().collect_vec();320 preferred.sort_by_key(|v| !config.prefer_identities.contains(*v));321322 warn!("reencrypting secret {secret} as it is missing for host {host}");323324 for owner in preferred {325 if let Some(hostname) = owner.as_host() && let Ok(host) = config.host(hostname) {326 let best = best.clone();327 let reencrypt_targets = reencrypt_targets.clone();328 let reencrypted = match await_in_nix(async move {329 host.reencrypt_distribution(&best, reencrypt_targets.clone(), config.now).await330 }) {331 Ok(r) => r,332 Err(e) => {333 warn!("reencryption failed on {hostname}: {e:?}");334 continue;335 }336 };337 dists.extend(reencrypted.clone(), format!("secret was reencrypted to extend with new owners: {new_owners:?}"));338 return Ok(secret_to_parts(&secret, &reencrypted.secret.parts, &expectations.parts));339 };340 }341 warn!("failed to reencrypt using any host")342 };343 };344345 if constraints.allow_different {346 for dist in dists.distributions() {347 for own in dist.owners() {348 expectations.owners.remove(own);349 }350 }351 }352 info!("secret {secret} is being generated for {:?}", expectations.owners);353354 let expectations_ = expectations.clone();355 let generated = await_in_nix(async move {356 generate(config, expectations_, &generator, &default_generator_drv).await357 })?;358359 dists.extend(generated.clone(), "secret was generated".to_string());360361 Ok(secret_to_parts(&secret, &generated.secret.parts, &expectations.parts))362 },363 )364 .register();365}
--- 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"))],
};