10 files changed
--- /dev/null
+++ b/.cargo/config.toml
@@ -0,0 +1,2 @@
+[build]
+rustflags = ["--cfg", "tokio_unstable"]
445 "fs_extra",445 "fs_extra",
446]446]
447
448[[package]]
449name = "axum"
450version = "0.8.9"
451source = "registry+https://github.com/rust-lang/crates.io-index"
452checksum = "31b698c5f9a010f6573133b09e0de5408834d0c82f8d7475a89fc1867a71cd90"
453dependencies = [
454 "axum-core",
455 "bytes",
456 "futures-util",
457 "http",
458 "http-body",
459 "http-body-util",
460 "itoa",
461 "matchit",
462 "memchr",
463 "mime",
464 "percent-encoding",
465 "pin-project-lite",
466 "serde_core",
467 "sync_wrapper",
468 "tower",
469 "tower-layer",
470 "tower-service",
471]
472
473[[package]]
474name = "axum-core"
475version = "0.5.6"
476source = "registry+https://github.com/rust-lang/crates.io-index"
477checksum = "08c78f31d7b1291f7ee735c1c6780ccde7785daae9a9206026862dab7d8792d1"
478dependencies = [
479 "bytes",
480 "futures-core",
481 "http",
482 "http-body",
483 "http-body-util",
484 "mime",
485 "pin-project-lite",
486 "sync_wrapper",
487 "tower-layer",
488 "tower-service",
489]
447490
448[[package]]491[[package]]
449name = "backon"492name = "backon"
995 "windows-sys 0.61.2",1038 "windows-sys 0.61.2",
996]1039]
1040
1041[[package]]
1042name = "console-api"
1043version = "0.9.0"
1044source = "registry+https://github.com/rust-lang/crates.io-index"
1045checksum = "e8599749b6667e2f0c910c1d0dff6901163ff698a52d5a39720f61b5be4b20d3"
1046dependencies = [
1047 "futures-core",
1048 "prost",
1049 "prost-types",
1050 "tonic",
1051 "tonic-prost",
1052 "tracing-core",
1053]
1054
1055[[package]]
1056name = "console-subscriber"
1057version = "0.5.0"
1058source = "registry+https://github.com/rust-lang/crates.io-index"
1059checksum = "fb4915b7d8dd960457a1b6c380114c2944f728e7c65294ab247ae6b6f1f37592"
1060dependencies = [
1061 "console-api",
1062 "crossbeam-channel",
1063 "crossbeam-utils",
1064 "futures-task",
1065 "hdrhistogram",
1066 "humantime",
1067 "hyper-util",
1068 "prost",
1069 "prost-types",
1070 "serde",
1071 "serde_json",
1072 "thread_local",
1073 "tokio",
1074 "tokio-stream",
1075 "tonic",
1076 "tracing",
1077 "tracing-core",
1078 "tracing-subscriber",
1079]
9971080
998[[package]]1081[[package]]
999name = "const-hex"1082name = "const-hex"
1834 "clap_complete",1917 "clap_complete",
1835 "fleet-base",1918 "fleet-base",
1836 "futures",1919 "futures",
1837 "human-repr",1920 "goodlog-subscriber",
1838 "indicatif",
1839 "itertools 0.15.0",1921 "itertools 0.15.0",
1840 "nix-eval",1922 "nix-eval",
1841 "opentelemetry",
1842 "opentelemetry-appender-tracing",
1843 "opentelemetry-exporter-env",
1844 "opentelemetry_sdk",
1845 "remowt-fleet",1923 "remowt-fleet",
1846 "serde",1924 "serde",
1847 "serde_json",1925 "serde_json",
1850 "tempfile",1928 "tempfile",
1851 "tokio",1929 "tokio",
1852 "tracing",1930 "tracing",
1853 "tracing-indicatif",
1854 "tracing-opentelemetry",
1855 "tracing-subscriber",
1856]1931]
18571932
1858[[package]]1933[[package]]
2242 "wasm-bindgen",2317 "wasm-bindgen",
2243]2318]
2319
2320[[package]]
2321name = "goodlog-subscriber"
2322version = "0.1.9"
2323dependencies = [
2324 "anyhow",
2325 "clap",
2326 "console-subscriber",
2327 "human-repr",
2328 "indicatif",
2329 "opentelemetry",
2330 "opentelemetry-appender-tracing",
2331 "opentelemetry-exporter-env",
2332 "opentelemetry_sdk",
2333 "tracing",
2334 "tracing-indicatif",
2335 "tracing-opentelemetry",
2336 "tracing-subscriber",
2337]
22442338
2245[[package]]2339[[package]]
2246name = "group"2340name = "group"
2283 "foldhash",2377 "foldhash",
2284]2378]
2379
2380[[package]]
2381name = "hdrhistogram"
2382version = "7.5.4"
2383source = "registry+https://github.com/rust-lang/crates.io-index"
2384checksum = "765c9198f173dd59ce26ff9f95ef0aafd0a0fe01fb9d72841bc5066a4c06511d"
2385dependencies = [
2386 "base64 0.21.7",
2387 "byteorder",
2388 "flate2",
2389 "nom 7.1.3",
2390 "num-traits",
2391]
22852392
2286[[package]]2393[[package]]
2287name = "heck"2394name = "heck"
2491source = "registry+https://github.com/rust-lang/crates.io-index"2598source = "registry+https://github.com/rust-lang/crates.io-index"
2492checksum = "f58b778a5761513caf593693f8951c97a5b610841e754788400f32102eefdff1"2599checksum = "f58b778a5761513caf593693f8951c97a5b610841e754788400f32102eefdff1"
2600
2601[[package]]
2602name = "humantime"
2603version = "2.3.0"
2604source = "registry+https://github.com/rust-lang/crates.io-index"
2605checksum = "135b12329e5e3ce057a9f972339ea52bc954fe1e9358ef27f95e89716fbc5424"
24932606
2494[[package]]2607[[package]]
2495name = "hybrid-array"2608name = "hybrid-array"
3368 "regex-automata",3481 "regex-automata",
3369]3482]
3483
3484[[package]]
3485name = "matchit"
3486version = "0.8.4"
3487source = "registry+https://github.com/rust-lang/crates.io-index"
3488checksum = "47e1ffaa40ddd1f3ed91f717a33c8c0ee23fff369e3aa8772b9605cc1d22f4c3"
33703489
3371[[package]]3490[[package]]
3372name = "md5"3491name = "md5"
3389 "autocfg",3508 "autocfg",
3390]3509]
3510
3511[[package]]
3512name = "mime"
3513version = "0.3.17"
3514source = "registry+https://github.com/rust-lang/crates.io-index"
3515checksum = "6877bb514081ee2a7ff5ef9de3281f14a4dd4bceac4c09388074a6b5df8a139a"
33913516
3392[[package]]3517[[package]]
3393name = "minimal-lexical"3518name = "minimal-lexical"
6440checksum = "ac2a5518c70fa84342385732db33fb3f44bc4cc748936eb5833d2df34d6445ef"6565checksum = "ac2a5518c70fa84342385732db33fb3f44bc4cc748936eb5833d2df34d6445ef"
6441dependencies = [6566dependencies = [
6442 "async-trait",6567 "async-trait",
6568 "axum",
6443 "base64 0.22.1",6569 "base64 0.22.1",
6444 "bytes",6570 "bytes",
6445 "flate2",6571 "flate2",
6572 "h2",
6446 "http",6573 "http",
6447 "http-body",6574 "http-body",
6448 "http-body-util",6575 "http-body-util",
6451 "hyper-util",6578 "hyper-util",
6452 "percent-encoding",6579 "percent-encoding",
6453 "pin-project",6580 "pin-project",
6581 "socket2",
6454 "sync_wrapper",6582 "sync_wrapper",
6455 "tokio",6583 "tokio",
6456 "tokio-stream",6584 "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());
- }
-}
--- a/cmds/fleet/src/main.rs
+++ b/cmds/fleet/src/main.rs
@@ -1,7 +1,6 @@
#![recursion_limit = "512"]
pub(crate) mod cmds;
-mod log_tree;
use std::{process::ExitCode, sync::Arc};
@@ -18,25 +17,12 @@
};
use fleet_base::{host::Config, opts::FleetOpts};
use futures::{TryStreamExt, stream::FuturesUnordered};
-#[cfg(feature = "indicatif")]
-use human_repr::HumanCount;
-#[cfg(feature = "indicatif")]
-use indicatif::{ProgressState, ProgressStyle};
-use log_tree::TreeLayer;
+use goodlog_subscriber::{LogOpts, setup_logging};
use nix_eval::{
eval_store, gc_register_my_thread, gc_unregister_my_thread, init_libraries, init_tokio_for_nix,
-};
-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};
use tokio::task::spawn_blocking;
use tracing::{Instrument, error, info, info_span};
-#[cfg(feature = "indicatif")]
-use tracing_indicatif::IndicatifLayer;
-use tracing_subscriber::{EnvFilter, prelude::*};
#[derive(Parser)]
struct Prefetch {}
@@ -102,14 +88,9 @@
fleet_opts: FleetOpts,
#[clap(subcommand)]
command: Opts,
- #[clap(long, next_help_heading = "Telemetry", env = "OTEL_FLEET")]
- otel: bool,
- #[clap(flatten)]
- otlp_base: OtlpBaseSettings,
- #[clap(flatten)]
- otel_logs: OtlpLogsSettings,
+
#[clap(flatten)]
- otel_traces: OtlpTracesSettings,
+ log: LogOpts,
}
async fn run_command(config: &Config, opts: FleetOpts, command: Opts) -> Result<()> {
@@ -127,100 +108,6 @@
Ok(())
}
-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(
- "{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");
- }
- },
- ),
- )
- };
-
- 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().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(())
-}
-
fn main() -> ExitCode {
let opts = RootOpts::parse();
if let Opts::Complete(c) = &opts.command {
@@ -228,7 +115,7 @@
return ExitCode::SUCCESS;
}
- if let Err(e) = setup_logging(&opts) {
+ if let Err(e) = setup_logging(&opts.log) {
eprintln!("{e:#}");
return ExitCode::FAILURE;
}
@@ -250,16 +137,12 @@
init_tokio_for_nix(runtime.clone());
runtime.block_on(async {
- tokio::task::spawn(async move {
- if let Err(e) = main_real(opts).await {
- error!("{e:#}");
- ExitCode::FAILURE
- } else {
- ExitCode::SUCCESS
- }
- })
- .await
- .expect("primary task panicked")
+ if let Err(e) = main_real(opts).await {
+ error!("{e:#}");
+ ExitCode::FAILURE
+ } else {
+ ExitCode::SUCCESS
+ }
})
}
--- /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