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",
1[workspace]2members = ["crates/*", "cmds/*", "remowt/crates/*", "remowt/cmds/*"]3resolver = "3"4package.version = "0.1.9"5package.edition = "2024"6package.rust-version = "1.95.0"7package.license = "MIT"89[workspace.dependencies]10fleet-base = { path = "./crates/fleet-base" }11fleet-shared = { path = "./crates/fleet-shared" }12nix-eval = { path = "./crates/nix-eval" }13nixlike = { path = "./crates/nixlike" }14opentelemetry-exporter-env = { path = "./crates/opentelemetry-exporter-env" }15remowt-fleet = { path = "./crates/remowt-fleet" }1617remowt-client = { version = "0.1.9", path = "remowt/crates/remowt-client" }18remowt-endpoints = { version = "0.1.9", path = "remowt/crates/remowt-endpoints" }19remowt-link-shared = { version = "0.1.9", path = "remowt/crates/remowt-link-shared" }20remowt-plugin = { version = "0.1.9", path = "remowt/crates/remowt-plugin" }21remowt-polkit-shared = { version = "0.1.9", path = "remowt/crates/polkit-shared" }22remowt-ui-prompt = { version = "0.1.9", path = "remowt/crates/remowt-ui-prompt" }2324bifrostlink = "0.2.0"25bifrostlink-macros = "0.2.0"26bifrostlink-ports = "0.2.0"2728iroh = { version = "1.0.0", features = ["unstable-custom-transports"] }29iroh-base = "1.0.0"30n0-watcher = "1.0.0"31noq-udp = { version = "1.0.0", default-features = false }3233age = { version = "0.11", features = ["plugin", "ssh"] }34anyhow = "1.0"35base64 = "0.22.1"36bindgen = "0.72.0"37bytes = "1.11.0"38camino = "1.2.2"39chrono = { version = "0.4.41", features = ["serde"] }40clap = { version = "4.5", features = ["derive", "env", "unicode", "wrap_help"] }41clap_complete = "4.5"42cxx = "1.0.168"43cxx-build = "1.0.168"44ed25519-dalek = "3.0.0-rc.0"45futures = "0.3.31"46hex = "0.4.3"47hmac = "0.13.0"48hostname = "0.4.1"49human-repr = "1.1"50indicatif = "0.18"51indoc = "2.0.6"52itertools = "0.15.0"53linked-hash-map = "0.5.6"54nix = { version = "0.31.2", features = ["fs", "user"] }55nom = "8.0.0"56opentelemetry = "0.32.0"57opentelemetry-appender-tracing = "0.32.0"58opentelemetry-otlp = { version = "0.32.0", features = ["grpc-tonic", "gzip-tonic", "http-json", "reqwest-rustls"] }59opentelemetry_sdk = "0.32.0"60pbkdf2 = "0.13"61peg = "0.8.5"62pkg-config = "0.3.30"63rand = "0.10.0"64russh = { version = "0.61.2", default-features = false, features = ["flate2", "ring", "rsa"] }65russh-config = "0.58.0"66serde = { version = "1.0", features = ["derive"] }67serde_json = "1.0"68sha2 = "0.11"69shlex = "2.0.1"70tabled = "0.21.0"71tempfile = "3.20"72test-log = { version = "0.2.19", features = ["trace"] }73thiserror = "2.0.12"74time = "0.3.41"75tokio = { version = "1.45.1", features = ["fs", "macros", "rt", "rt-multi-thread", "sync", "time"] }76tracing = "0.1"77tracing-indicatif = "0.3.13"78tracing-journald = "0.3.2"79tracing-opentelemetry = "0.33.0"80uuid = { version = "1", features = ["v4"] }8182tokio-util = "0.7.11"83tracing-subscriber = { version = "0.3.19", features = ["env-filter", "fmt"] }84unicode_categories = "0.1.1"85vte = { version = "0.15.0", features = ["ansi"] }86x25519-dalek = { version = "2.0.1", features = ["getrandom"] }87zbus = "5.16.0"88zbus_polkit = "5.0.0"8990[profile.dev]91panic = "abort"92[profile.release]93panic = "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