10 files changed
--- /dev/null
+++ b/.cargo/config.toml
@@ -0,0 +1,2 @@
+[build]
+rustflags = ["--cfg", "tokio_unstable"]
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -446,6 +446,49 @@
]
[[package]]
+name = "axum"
+version = "0.8.9"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "31b698c5f9a010f6573133b09e0de5408834d0c82f8d7475a89fc1867a71cd90"
+dependencies = [
+ "axum-core",
+ "bytes",
+ "futures-util",
+ "http",
+ "http-body",
+ "http-body-util",
+ "itoa",
+ "matchit",
+ "memchr",
+ "mime",
+ "percent-encoding",
+ "pin-project-lite",
+ "serde_core",
+ "sync_wrapper",
+ "tower",
+ "tower-layer",
+ "tower-service",
+]
+
+[[package]]
+name = "axum-core"
+version = "0.5.6"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "08c78f31d7b1291f7ee735c1c6780ccde7785daae9a9206026862dab7d8792d1"
+dependencies = [
+ "bytes",
+ "futures-core",
+ "http",
+ "http-body",
+ "http-body-util",
+ "mime",
+ "pin-project-lite",
+ "sync_wrapper",
+ "tower-layer",
+ "tower-service",
+]
+
+[[package]]
name = "backon"
version = "1.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -996,6 +1039,46 @@
]
[[package]]
+name = "console-api"
+version = "0.9.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "e8599749b6667e2f0c910c1d0dff6901163ff698a52d5a39720f61b5be4b20d3"
+dependencies = [
+ "futures-core",
+ "prost",
+ "prost-types",
+ "tonic",
+ "tonic-prost",
+ "tracing-core",
+]
+
+[[package]]
+name = "console-subscriber"
+version = "0.5.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "fb4915b7d8dd960457a1b6c380114c2944f728e7c65294ab247ae6b6f1f37592"
+dependencies = [
+ "console-api",
+ "crossbeam-channel",
+ "crossbeam-utils",
+ "futures-task",
+ "hdrhistogram",
+ "humantime",
+ "hyper-util",
+ "prost",
+ "prost-types",
+ "serde",
+ "serde_json",
+ "thread_local",
+ "tokio",
+ "tokio-stream",
+ "tonic",
+ "tracing",
+ "tracing-core",
+ "tracing-subscriber",
+]
+
+[[package]]
name = "const-hex"
version = "1.19.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -1834,14 +1917,9 @@
"clap_complete",
"fleet-base",
"futures",
- "human-repr",
- "indicatif",
+ "goodlog-subscriber",
"itertools 0.15.0",
"nix-eval",
- "opentelemetry",
- "opentelemetry-appender-tracing",
- "opentelemetry-exporter-env",
- "opentelemetry_sdk",
"remowt-fleet",
"serde",
"serde_json",
@@ -1850,9 +1928,6 @@
"tempfile",
"tokio",
"tracing",
- "tracing-indicatif",
- "tracing-opentelemetry",
- "tracing-subscriber",
]
[[package]]
@@ -2243,6 +2318,25 @@
]
[[package]]
+name = "goodlog-subscriber"
+version = "0.1.9"
+dependencies = [
+ "anyhow",
+ "clap",
+ "console-subscriber",
+ "human-repr",
+ "indicatif",
+ "opentelemetry",
+ "opentelemetry-appender-tracing",
+ "opentelemetry-exporter-env",
+ "opentelemetry_sdk",
+ "tracing",
+ "tracing-indicatif",
+ "tracing-opentelemetry",
+ "tracing-subscriber",
+]
+
+[[package]]
name = "group"
version = "0.14.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2284,6 +2378,19 @@
]
[[package]]
+name = "hdrhistogram"
+version = "7.5.4"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "765c9198f173dd59ce26ff9f95ef0aafd0a0fe01fb9d72841bc5066a4c06511d"
+dependencies = [
+ "base64 0.21.7",
+ "byteorder",
+ "flate2",
+ "nom 7.1.3",
+ "num-traits",
+]
+
+[[package]]
name = "heck"
version = "0.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2492,6 +2599,12 @@
checksum = "f58b778a5761513caf593693f8951c97a5b610841e754788400f32102eefdff1"
[[package]]
+name = "humantime"
+version = "2.3.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "135b12329e5e3ce057a9f972339ea52bc954fe1e9358ef27f95e89716fbc5424"
+
+[[package]]
name = "hybrid-array"
version = "0.4.12"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -3369,6 +3482,12 @@
]
[[package]]
+name = "matchit"
+version = "0.8.4"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "47e1ffaa40ddd1f3ed91f717a33c8c0ee23fff369e3aa8772b9605cc1d22f4c3"
+
+[[package]]
name = "md5"
version = "0.8.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -3390,6 +3509,12 @@
]
[[package]]
+name = "mime"
+version = "0.3.17"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "6877bb514081ee2a7ff5ef9de3281f14a4dd4bceac4c09388074a6b5df8a139a"
+
+[[package]]
name = "minimal-lexical"
version = "0.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -6440,9 +6565,11 @@
checksum = "ac2a5518c70fa84342385732db33fb3f44bc4cc748936eb5833d2df34d6445ef"
dependencies = [
"async-trait",
+ "axum",
"base64 0.22.1",
"bytes",
"flate2",
+ "h2",
"http",
"http-body",
"http-body-util",
@@ -6451,6 +6578,7 @@
"hyper-util",
"percent-encoding",
"pin-project",
+ "socket2",
"sync_wrapper",
"tokio",
"tokio-stream",
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -39,6 +39,7 @@
chrono = { version = "0.4.41", features = ["serde"] }
clap = { version = "4.5", features = ["derive", "env", "unicode", "wrap_help"] }
clap_complete = "4.5"
+console-subscriber = "0.5.0"
cxx = "1.0.168"
cxx-build = "1.0.168"
ed25519-dalek = "3.0.0-rc.0"
@@ -72,7 +73,7 @@
test-log = { version = "0.2.19", features = ["trace"] }
thiserror = "2.0.12"
time = "0.3.41"
-tokio = { version = "1.45.1", features = ["fs", "macros", "rt", "rt-multi-thread", "sync", "time"] }
+tokio = { version = "1.45.1", features = ["fs", "macros", "rt", "rt-multi-thread", "sync", "time", "tracing"] }
tracing = "0.1"
tracing-indicatif = "0.3.13"
tracing-journald = "0.3.2"
@@ -86,6 +87,7 @@
x25519-dalek = { version = "2.0.1", features = ["getrandom"] }
zbus = "5.16.0"
zbus_polkit = "5.0.0"
+goodlog-subscriber = { version = "0.1.9", path = "crates/goodlog-subscriber" }
[profile.dev]
panic = "abort"
--- a/cmds/fleet/Cargo.toml
+++ b/cmds/fleet/Cargo.toml
@@ -20,28 +20,15 @@
tempfile.workspace = true
tokio.workspace = true
tracing.workspace = true
-tracing-subscriber.workspace = true
futures.workspace = true
itertools.workspace = true
shlex.workspace = true
tabled.workspace = true
-human-repr = { workspace = true, optional = true }
-indicatif = { workspace = true, optional = true }
-opentelemetry.workspace = true
-opentelemetry-appender-tracing.workspace = true
-opentelemetry-exporter-env.workspace = true
-opentelemetry_sdk.workspace = true
-tracing-indicatif = { workspace = true, optional = true }
-tracing-opentelemetry.workspace = true
+goodlog-subscriber.workspace = true
[features]
default = ["indicatif"]
# Not quite stable
-indicatif = [
- "dep:tracing-indicatif",
- "dep:indicatif",
- "dep:human-repr",
- "nix-eval/indicatif",
-]
+indicatif = ["nix-eval/indicatif", "goodlog-subscriber/indicatif"]
--- a/cmds/fleet/src/log_tree.rs
+++ /dev/null
@@ -1,243 +0,0 @@
-//! 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());
- }
-}
1#![recursion_limit = "512"]23pub(crate) mod cmds;4mod log_tree;56use std::{process::ExitCode, sync::Arc};78use anyhow::{Context as _, Result, bail};9use camino::Utf8PathBuf;10use clap::{CommandFactory, Parser};11use cmds::{12 build_systems::{BuildSystems, Deploy},13 complete::Complete,14 info::Info,15 rollback::RollbackSingle,16 secrets::Secret,17 tf::Tf,18};19use fleet_base::{host::Config, opts::FleetOpts};20use futures::{TryStreamExt, stream::FuturesUnordered};21#[cfg(feature = "indicatif")]22use human_repr::HumanCount;23#[cfg(feature = "indicatif")]24use indicatif::{ProgressState, ProgressStyle};25use log_tree::TreeLayer;26use nix_eval::{27 eval_store, gc_register_my_thread, gc_unregister_my_thread, init_libraries, init_tokio_for_nix,28};29use opentelemetry::trace::TracerProvider;30use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;31use opentelemetry_exporter_env::{32 OtlpBaseSettings, OtlpLogsSettings, OtlpTracesSettings, ResolvedOtlpSettings,33};34use opentelemetry_sdk::{logs::SdkLoggerProvider, trace::SdkTracerProvider};35use tokio::task::spawn_blocking;36use tracing::{Instrument, error, info, info_span};37#[cfg(feature = "indicatif")]38use tracing_indicatif::IndicatifLayer;39use tracing_subscriber::{EnvFilter, prelude::*};4041#[derive(Parser)]42struct Prefetch {}43impl Prefetch {44 async fn run(&self, config: &Config) -> Result<()> {45 let mut prefetch_dir = config.directory.to_path_buf();46 prefetch_dir.push("prefetch");47 if !prefetch_dir.is_dir() {48 info!("nothing to prefetch: no prefetch directory");49 return Ok(());50 }51 let tasks = FuturesUnordered::new();52 for entry in std::fs::read_dir(&prefetch_dir)? {53 let entry = entry?;54 if !entry.metadata()?.is_file() {55 bail!("only files should exist in prefetch directory");56 }57 let name = entry.file_name().to_string_lossy().into_owned();58 let path =59 Utf8PathBuf::try_from(entry.path()).context("prefetch path should be utf8")?;60 let span = info_span!("prefetching", name = %name);61 tasks.push(async move {62 let store = eval_store();63 let added = spawn_blocking(move || store.add_file(&name, &path))64 .instrument(span.clone())65 .await??;66 let _g = span.enter();67 info!("{} -> {}", added.hash, added.store_path);68 anyhow::Ok(())69 });70 }71 tasks.try_collect::<Vec<()>>().await?;72 Ok(())73 }74}7576#[derive(Parser)]77enum Opts {78 79 BuildSystems(BuildSystems),80 81 Deploy(Deploy),82 83 RollbackSingle(RollbackSingle),84 85 #[clap(subcommand)]86 Secret(Secret),87 88 Prefetch(Prefetch),89 90 Info(Info),91 92 #[clap(hide(true))]93 Complete(Complete),94 95 Tf(Tf),96}9798#[derive(Parser)]99#[clap(version, author)]100struct RootOpts {101 #[clap(flatten)]102 fleet_opts: FleetOpts,103 #[clap(subcommand)]104 command: Opts,105 #[clap(long, next_help_heading = "Telemetry", env = "OTEL_FLEET")]106 otel: bool,107 #[clap(flatten)]108 otlp_base: OtlpBaseSettings,109 #[clap(flatten)]110 otel_logs: OtlpLogsSettings,111 #[clap(flatten)]112 otel_traces: OtlpTracesSettings,113}114115async fn run_command(config: &Config, opts: FleetOpts, command: Opts) -> Result<()> {116 match command {117 Opts::BuildSystems(c) => c.run(config, &opts).await?,118 Opts::Deploy(d) => d.run(config, &opts).await?,119 Opts::RollbackSingle(r) => r.run(config, &opts).await?,120 Opts::Secret(s) => s.run(config, &opts).await?,121 Opts::Info(i) => i.run(config).await?,122 Opts::Prefetch(p) => p.run(config).await?,123 Opts::Tf(t) => t.run(config).await?,124 125 Opts::Complete(c) => spawn_blocking(move || c.run(RootOpts::command())).await?,126 };127 Ok(())128}129130fn setup_logging(opts: &RootOpts) -> Result<()> {131 #[cfg(feature = "indicatif")]132 let indicatif_layer = {133 use std::fmt;134 use std::time::Duration;135136 IndicatifLayer::new().with_max_progress_bars(10, Some(ProgressStyle::default_spinner()))137 .with_span_child_prefix_indent(" ")138 .with_span_child_prefix_symbol("")139 .with_progress_style(140 ProgressStyle::with_template(141 "{span_child_prefix:.magenta}{chevron:.magenta}{span_name:<18.blue.bold} {wide_msg} {span_fields:.dim} {color_start}{download_progress} {elapsed}{color_end}",142 )143 .unwrap()144 .with_key("chevron", |_: &ProgressState, writer: &mut dyn fmt::Write| {145 let _ = write!(writer, "> ");146 })147 .with_key("download_progress", |state: &ProgressState, writer: &mut dyn fmt::Write| {148 let Some(len) = state.len() else {149 return;150 };151 let pos = state.pos();152 if pos > len {153 let _ = write!(writer, "{}", pos.human_count_bare());154 } else {155 let _ = write!(writer, "{} / {}", pos.human_count_bare(), len.human_count_bare());156 }157 })158 .with_key(159 "color_start",160 |state: &ProgressState, writer: &mut dyn fmt::Write| {161 let elapsed = state.elapsed();162163 if elapsed > Duration::from_secs(60) {164 165 let _ = write!(writer, "\x1b[{}m", 1 + 30);166 } else if elapsed > Duration::from_secs(30) {167 168 let _ = write!(writer, "\x1b[{}m", 3 + 30);169 }170 },171 )172 .with_key(173 "color_end",174 |state: &ProgressState, writer: &mut dyn fmt::Write| {175 if state.elapsed() > Duration::from_secs(30) {176 let _ = write!(writer, "\x1b[0m");177 }178 },179 ),180 )181 };182183 let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("noq_udp=warn,info"));184185 let tree = {186 #[cfg(feature = "indicatif")]187 let writer = indicatif_layer.get_stderr_writer();188 #[cfg(not(feature = "indicatif"))]189 let writer = || std::io::stderr();190 TreeLayer::new(writer)191 };192193 let reg = tracing_subscriber::registry().with(tree.with_filter(filter));194195 #[cfg(feature = "indicatif")]196 let reg = reg.with(indicatif_layer);197198 if opts.otel {199 let traces = ResolvedOtlpSettings::traces(&opts.otlp_base, &opts.otel_traces)?;200 let span_exporter = traces.span_exporter()?;201 let logs = ResolvedOtlpSettings::logs(&opts.otlp_base, &opts.otel_logs)?;202 let log_exporter = logs.log_exporter()?;203204 let span_provider = SdkTracerProvider::builder()205 .with_batch_exporter(span_exporter)206 .build();207 let log_provider = SdkLoggerProvider::builder()208 .with_batch_exporter(log_exporter)209 .build();210211 let logger = OpenTelemetryTracingBridge::new(&log_provider);212 let tracer = span_provider.tracer("fleet");213214 reg.with(tracing_opentelemetry::layer().with_tracer(tracer))215 .with(logger)216 .init();217 } else {218 reg.init();219 }220221 Ok(())222}223224fn main() -> ExitCode {225 let opts = RootOpts::parse();226 if let Opts::Complete(c) = &opts.command {227 c.run(RootOpts::command());228 return ExitCode::SUCCESS;229 }230231 if let Err(e) = setup_logging(&opts) {232 eprintln!("{e:#}");233 return ExitCode::FAILURE;234 }235236 init_libraries();237238 let runtime = tokio::runtime::Builder::new_multi_thread()239 .enable_all()240 .on_thread_start(|| {241 gc_register_my_thread();242 })243 .on_thread_stop(|| {244 gc_unregister_my_thread();245 })246 .build()247 .expect("failed to build runtime");248 let runtime = Arc::new(runtime);249250 init_tokio_for_nix(runtime.clone());251252 runtime.block_on(async {253 tokio::task::spawn(async move {254 if let Err(e) = main_real(opts).await {255 error!("{e:#}");256 ExitCode::FAILURE257 } else {258 ExitCode::SUCCESS259 }260 })261 .await262 .expect("primary task panicked")263 })264}265266async fn main_real(opts: RootOpts) -> Result<()> {267 let config = opts.fleet_opts.build(matches!(268 opts.command,269 Opts::Deploy(_) | Opts::BuildSystems(_)270 ))?;271272 match run_command(&config, opts.fleet_opts, opts.command).await {273 Ok(()) => {274 config.save()?;275 Ok(())276 }277 Err(e) => {278 let _ = config.save();279 Err(e)280 }281 }282}283284#[cfg(test)]285mod tests {286 use super::*;287288 #[test]289 fn verify_command() {290 use clap::CommandFactory;291 RootOpts::command().debug_assert();292 }293}
--- /dev/null
+++ b/crates/goodlog-subscriber/Cargo.toml
@@ -0,0 +1,31 @@
+[package]
+name = "goodlog-subscriber"
+version.workspace = true
+edition.workspace = true
+rust-version.workspace = true
+license.workspace = true
+
+[dependencies]
+tracing-subscriber.workspace = true
+
+# Tokio-console
+console-subscriber = { workspace = true, optional = true }
+
+# Indicatif
+human-repr = { workspace = true, optional = true }
+indicatif = { workspace = true, optional = true }
+tracing-indicatif = { workspace = true, optional = true }
+
+# OTEL
+opentelemetry.workspace = true
+opentelemetry-appender-tracing.workspace = true
+opentelemetry-exporter-env.workspace = true
+opentelemetry_sdk.workspace = true
+tracing-opentelemetry.workspace = true
+clap = { workspace = true, features = ["derive"] }
+anyhow.workspace = true
+tracing.workspace = true
+
+[features]
+indicatif = ["dep:tracing-indicatif", "dep:indicatif", "dep:human-repr"]
+tokio-console = ["dep:console-subscriber"]
--- /dev/null
+++ b/crates/goodlog-subscriber/src/lib.rs
@@ -0,0 +1,135 @@
+mod log_tree;
+
+use anyhow::Result;
+use clap::Parser;
+#[cfg(feature = "indicatif")]
+use human_repr::HumanCount;
+#[cfg(feature = "indicatif")]
+use indicatif::{ProgressState, ProgressStyle};
+use log_tree::TreeLayer;
+use opentelemetry::trace::TracerProvider;
+use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;
+use opentelemetry_exporter_env::{
+ OtlpBaseSettings, OtlpLogsSettings, OtlpTracesSettings, ResolvedOtlpSettings,
+};
+use opentelemetry_sdk::{logs::SdkLoggerProvider, trace::SdkTracerProvider};
+#[cfg(feature = "indicatif")]
+use tracing_indicatif::IndicatifLayer;
+use tracing_subscriber::layer::SubscriberExt as _;
+use tracing_subscriber::util::SubscriberInitExt as _;
+use tracing_subscriber::{EnvFilter, Layer as _};
+
+#[derive(Parser)]
+#[clap(next_help_heading = "Telemetry")]
+pub struct LogOpts {
+ /// Whatever opentelemetry logging should be enabled
+ #[clap(long, env = "OTEL_FLEET")]
+ otel: bool,
+ #[clap(flatten)]
+ otlp_base: OtlpBaseSettings,
+ #[clap(flatten)]
+ otel_logs: OtlpLogsSettings,
+ #[clap(flatten)]
+ otel_traces: OtlpTracesSettings,
+}
+
+pub fn setup_logging(opts: &LogOpts) -> 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(
+ "{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("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;
+ };
+ let pos = state.pos();
+ if pos > len {
+ let _ = write!(writer, "{}", pos.human_count_bare());
+ } else {
+ let _ = write!(writer, "{} / {}", pos.human_count_bare(), len.human_count_bare());
+ }
+ })
+ .with_key(
+ "color_start",
+ |state: &ProgressState, writer: &mut dyn fmt::Write| {
+ let elapsed = state.elapsed();
+
+ if elapsed > Duration::from_secs(60) {
+ // Red
+ let _ = write!(writer, "\x1b[{}m", 1 + 30);
+ } else if elapsed > Duration::from_secs(30) {
+ // Yellow
+ let _ = write!(writer, "\x1b[{}m", 3 + 30);
+ }
+ },
+ )
+ .with_key(
+ "color_end",
+ |state: &ProgressState, writer: &mut dyn fmt::Write| {
+ if state.elapsed() > Duration::from_secs(30) {
+ let _ = write!(writer, "\x1b[0m");
+ }
+ },
+ ),
+ )
+ };
+
+ // TODO: Default filter should be configurable
+ let filter =
+ EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("noq_udp=warn,info"));
+
+ let tree = {
+ #[cfg(feature = "indicatif")]
+ let writer = indicatif_layer.get_stderr_writer();
+ #[cfg(not(feature = "indicatif"))]
+ let writer = || std::io::stderr();
+ TreeLayer::new(writer)
+ };
+
+ let reg = tracing_subscriber::registry();
+
+ #[cfg(feature = "tokio-console")]
+ let reg = reg.with(console_subscriber::spawn());
+
+ let reg = reg.with(tree.with_filter(filter));
+
+ #[cfg(feature = "indicatif")]
+ let reg = reg.with(indicatif_layer);
+
+ if opts.otel {
+ let traces = ResolvedOtlpSettings::traces(&opts.otlp_base, &opts.otel_traces)?;
+ let span_exporter = traces.span_exporter()?;
+ let logs = ResolvedOtlpSettings::logs(&opts.otlp_base, &opts.otel_logs)?;
+ let log_exporter = logs.log_exporter()?;
+
+ let span_provider = SdkTracerProvider::builder()
+ .with_batch_exporter(span_exporter)
+ .build();
+ let log_provider = SdkLoggerProvider::builder()
+ .with_batch_exporter(log_exporter)
+ .build();
+
+ let logger = OpenTelemetryTracingBridge::new(&log_provider);
+ let tracer = span_provider.tracer("fleet");
+
+ reg.with(tracing_opentelemetry::layer().with_tracer(tracer))
+ .with(logger)
+ .init();
+ } else {
+ reg.init();
+ }
+
+ Ok(())
+}
--- /dev/null
+++ b/crates/goodlog-subscriber/src/log_tree.rs
@@ -0,0 +1,233 @@
+//! Streaming, indented tree formatter.
+
+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()
+ }
+}
+
+struct SpanData {
+ fields: String,
+ started: Instant,
+ opened: AtomicBool,
+}
+
+fn level_code(level: &Level) -> &'static str {
+ match *level {
+ // Red
+ Level::ERROR => "31m",
+ // Yellow
+ Level::WARN => "33m",
+ // Blue
+ Level::INFO => "34m",
+ // Default fg
+ Level::DEBUG => "39m",
+ // Gray
+ Level::TRACE => "90m",
+ }
+}
+
+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,
+ 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![]),
+ }
+ }
+
+ 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/pkgs/fleet.nix
+++ b/pkgs/fleet.nix
@@ -31,6 +31,7 @@
# TODO: built-in fleet prompter should be a prodash widget, or it should require
# tty remowt prompter running on host machine idk.
ROFI = "${rofi}/bin/rofi";
+ RUSTFLAGS = "--cfg tokio_unstable";
buildInputs = [
inputs.nix.packages.${system}.nix-expr-c