From 51d5ea52f7ad43e4c51e7cfba8c13eba7c55b1fb Mon Sep 17 00:00:00 2001 From: Yaroslav Bolyukin Date: Wed, 01 Jul 2026 01:52:45 +0000 Subject: [PATCH] console --- --- /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 { - 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>, -} - -impl TreeLayer { - 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(&self, scope: Option>, 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::() 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 Layer for TreeLayer -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::() - .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 { + writer: W, + color: bool, + last: Mutex>, +} + +impl TreeLayer { + 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(&self, scope: Option>, 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::() 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 Layer for TreeLayer +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::() + .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 -- gitstuff