From f0ad944a4437ceca72fcef278cd83f3dd5ef4e64 Mon Sep 17 00:00:00 2001 From: Hansie Odendaal <39146854+hansieodendaal@users.noreply.github.com> Date: Tue, 21 Jul 2026 12:45:14 +0200 Subject: [PATCH] feat: update overwatch for tokio-console (#3148) --- Cargo.lock | 40 +++++++++- Cargo.toml | 4 +- README.md | 58 ++++++++++++++ nodes/node/binary/Cargo.toml | 2 +- nodes/node/binary/src/config/tests.rs | 28 +++++++ .../src/config/tracing/serde/console.rs | 4 + services/tracing/src/console.rs | 78 +++++++++++++++++-- services/tracing/src/lib.rs | 22 +++++- .../features/manual_control.feature | 6 +- .../src/cucumber/steps/manual_nodes/utils.rs | 3 + .../cucumber/steps/tokio_console/profile.rs | 1 + .../src/cucumber/steps/tokio_console/steps.rs | 19 +++-- .../src/framework/local/provisioning.rs | 18 +++++ 13 files changed, 261 insertions(+), 22 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 53a107930..be89d12c0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5383,6 +5383,29 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" +[[package]] +name = "manyhow" +version = "0.11.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b33efb3ca6d3b07393750d4030418d594ab1139cee518f0dc88db70fec873587" +dependencies = [ + "manyhow-macros", + "proc-macro2", + "quote", + "syn 2.0.118", +] + +[[package]] +name = "manyhow-macros" +version = "0.11.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "46fce34d199b78b6e6073abf984c9cf5fd3e9330145a93ee0738a7443e371495" +dependencies = [ + "proc-macro-utils", + "proc-macro2", + "quote", +] + [[package]] name = "match-lookup" version = "0.1.2" @@ -6174,7 +6197,7 @@ dependencies = [ [[package]] name = "overwatch" version = "0.1.0" -source = "git+https://github.com/logos-co/Overwatch?rev=448c192#448c192895b8311c742b1726a1bb12ee314ad95c" +source = "git+https://github.com/logos-co/Overwatch?rev=ef70b69fb82b0b585cfef1427972c4e97ea88c84#ef70b69fb82b0b585cfef1427972c4e97ea88c84" dependencies = [ "async-trait", "futures", @@ -6189,10 +6212,10 @@ dependencies = [ [[package]] name = "overwatch-derive" version = "0.1.0" -source = "git+https://github.com/logos-co/Overwatch?rev=448c192#448c192895b8311c742b1726a1bb12ee314ad95c" +source = "git+https://github.com/logos-co/Overwatch?rev=ef70b69fb82b0b585cfef1427972c4e97ea88c84#ef70b69fb82b0b585cfef1427972c4e97ea88c84" dependencies = [ "convert_case 0.8.0", - "proc-macro-error2", + "manyhow", "proc-macro2", "quote", "syn 2.0.118", @@ -6584,6 +6607,17 @@ dependencies = [ "syn 2.0.118", ] +[[package]] +name = "proc-macro-utils" +version = "0.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eeaf08a13de400bc215877b5bdc088f241b12eb42f0a548d3390dc1c56bb7071" +dependencies = [ + "proc-macro2", + "quote", + "smallvec", +] + [[package]] name = "proc-macro2" version = "1.0.106" diff --git a/Cargo.toml b/Cargo.toml index 87f931fbb..cddc99cbb 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -245,8 +245,8 @@ opentelemetry-http = { default-features = false, version = "0.31 opentelemetry-otlp = { default-features = false, version = "0.31" } opentelemetry-semantic-conventions = { default-features = false, version = "0.31" } opentelemetry_sdk = { default-features = false, version = "0.31" } -overwatch = { default-features = false, git = "https://github.com/logos-co/Overwatch", rev = "448c192" } -overwatch-derive = { default-features = false, git = "https://github.com/logos-co/Overwatch", rev = "448c192" } +overwatch = { default-features = false, git = "https://github.com/logos-co/Overwatch", rev = "ef70b69fb82b0b585cfef1427972c4e97ea88c84" } +overwatch-derive = { default-features = false, git = "https://github.com/logos-co/Overwatch", rev = "ef70b69fb82b0b585cfef1427972c4e97ea88c84" } owo-colors = { default-features = false, version = "4.2.3" } pprof = { default-features = false, version = "0.15" } proc-macro2 = { default-features = false, version = "1" } diff --git a/README.md b/README.md index b89c412b7..d888aa741 100644 --- a/README.md +++ b/README.md @@ -234,8 +234,66 @@ tracing: console: !Console bind_address: 127.0.0.1 port: 6669 + recording_path: /absolute/path/to/node-1-tokio-console.jsonl ``` +`recording_path` enables subscriber-side raw Tokio Console recording. Omitting it +leaves raw recording disabled while the live Tokio Console endpoint remains +available. Absolute paths are recommended; the node process must be able to +create and write the file. Use a unique path for each node and profiling run. +Recordings can grow substantially during long runs, and recording adds +instrumentation and disk-I/O overhead. Stop the node gracefully so the recorder +can flush as much telemetry as possible. The file contains raw subscriber +telemetry for offline analysis and is different from Tokio Console client +diagnostic logs. + +The verified `console-subscriber 0.5.0` format is newline-delimited JSON: the +first line is a version header (`{"v":1}`), followed by raw `Spawn`, `Enter`, +`Exit`, `Close`, and `Waker` event records. This recording format is currently +experimental and may change between subscriber versions. The Tokio Console +client connects to a live endpoint; it does not currently replay these raw +files directly, so offline analysis requires a compatible parser or tooling. + +For multiple nodes, use unique ports and recording paths: + +```yaml +# NODE_1 +tracing: + console: !Console + bind_address: 127.0.0.1 + port: 6669 + recording_path: /profiles/run-001/node-1.jsonl +``` + +```yaml +# NODE_2 +tracing: + console: !Console + bind_address: 127.0.0.1 + port: 6670 + recording_path: /profiles/run-001/node-2.jsonl +``` + +When using Cucumber, select recording independently for each profiled node: + +```gherkin +And I will have tokio console profile nodes: + | node_name | record_raw | + | NODE_1 | true | + | NODE_2 | false | +``` + +`true` enables the live endpoint and raw recording; `false` enables only the +live endpoint. Nodes omitted from the table are unaffected. For an enabled +node, Cucumber stores the recording at: + +```text +//tokio-console-raw.jsonl +``` + +The path is resolved after the node runtime directory is created and remains +with the scenario and node artifacts after shutdown. + **Note:** Port `6669` is the default port for the console, but you can change it in your config and use the corresponding value in the runtime if needed, for example, use a different port when multiple instrumented node processes are running on the same host. diff --git a/nodes/node/binary/Cargo.toml b/nodes/node/binary/Cargo.toml index 2ee93ee05..acc4cd3b6 100644 --- a/nodes/node/binary/Cargo.toml +++ b/nodes/node/binary/Cargo.toml @@ -100,5 +100,5 @@ testing-disable-proposal-publish = ["lb-chain-leader-service/testing-disable-pro # # Builds using this feature must also be compiled with: # RUSTFLAGS="--cfg tokio_unstable" -tokio-console = ["lb-tracing-service/tokio-console", "tokio/tracing"] +tokio-console = ["lb-tracing-service/tokio-console", "overwatch/tokio-task-names", "tokio/tracing"] tracing = [] diff --git a/nodes/node/binary/src/config/tests.rs b/nodes/node/binary/src/config/tests.rs index faf78b205..bf9c4ecca 100644 --- a/nodes/node/binary/src/config/tests.rs +++ b/nodes/node/binary/src/config/tests.rs @@ -47,6 +47,7 @@ fn tokio_console_config_defaults_to_loopback_default_console_port() { assert_eq!(config.bind_address, IpAddr::V4(Ipv4Addr::LOCALHOST)); assert_eq!(config.port, 6_669); + assert_eq!(config.recording_path, None); } #[test] @@ -66,6 +67,33 @@ fn tokio_console_config_deserializes() { assert_eq!(config.bind_address, IpAddr::V4(Ipv4Addr::LOCALHOST)); assert_eq!(config.port, 6_669); + assert_eq!(config.recording_path, None); +} + +#[test] +fn tokio_console_config_deserializes_recording_path_and_converts() { + let layer: ConsoleLayer = serde_yaml::from_str( + " + !Console + bind_address: 127.0.0.1 + port: 6669 + recording_path: /tmp/node-tokio-console.jsonl + ", + ) + .expect("tokio-console config should deserialize"); + + let ConsoleLayer::Console(config) = layer else { + panic!("expected console layer"); + }; + let tracing_config: lb_tracing_service::ConsoleLayerSettings = + ConsoleLayer::Console(config.clone()).into(); + let lb_tracing_service::ConsoleLayerSettings::Console(tracing_config) = tracing_config else { + panic!("expected console layer"); + }; + + assert_eq!(tracing_config.bind_address, "127.0.0.1"); + assert_eq!(tracing_config.port, 6_669); + assert_eq!(tracing_config.recording_path, config.recording_path); } #[test] diff --git a/nodes/node/binary/src/config/tracing/serde/console.rs b/nodes/node/binary/src/config/tracing/serde/console.rs index 7e12a39c2..193d6a3bd 100644 --- a/nodes/node/binary/src/config/tracing/serde/console.rs +++ b/nodes/node/binary/src/config/tracing/serde/console.rs @@ -1,4 +1,5 @@ use core::net::{IpAddr, Ipv4Addr}; +use std::path::PathBuf; use lb_tracing_service::{ConsoleLayerSettings, TokioConsoleConfig}; use serde::{Deserialize, Serialize}; @@ -16,6 +17,7 @@ impl From for ConsoleLayerSettings { Layer::Console(config) => Self::Console(TokioConsoleConfig { bind_address: config.bind_address.to_string(), port: config.port, + recording_path: config.recording_path, }), Layer::None => Self::None, } @@ -27,6 +29,7 @@ impl From for ConsoleLayerSettings { pub struct TokioConfig { pub bind_address: IpAddr, pub port: u16, + pub recording_path: Option, } impl Default for TokioConfig { @@ -34,6 +37,7 @@ impl Default for TokioConfig { Self { bind_address: Ipv4Addr::LOCALHOST.into(), port: 6_669, + recording_path: None, } } } diff --git a/services/tracing/src/console.rs b/services/tracing/src/console.rs index a51e89f4c..45d2ae74f 100644 --- a/services/tracing/src/console.rs +++ b/services/tracing/src/console.rs @@ -1,4 +1,4 @@ -use std::net::SocketAddr; +use std::{fs, io, net::SocketAddr, path::Path}; use console_subscriber::ConsoleLayer; use tracing_subscriber::Layer; @@ -7,7 +7,7 @@ use crate::TokioConsoleConfig; pub fn create_console_layer( config: &TokioConsoleConfig, -) -> Option + Send + Sync>> +) -> Result + Send + Sync>>, io::Error> where S: tracing::Subscriber + for<'span> tracing_subscriber::registry::LookupSpan<'span> @@ -17,10 +17,78 @@ where let bind_addr = format!("{}:{}", config.bind_address, config.port); bind_addr.parse::().map_or_else( - |_| None, + |_| Ok(None), |socket_addr| { - let console_layer = ConsoleLayer::builder().server_addr(socket_addr).spawn(); - Some(Box::new(console_layer)) + let mut builder = ConsoleLayer::builder().server_addr(socket_addr); + if let Some(recording_path) = &config.recording_path { + prepare_recording_path(recording_path)?; + builder = builder.recording_path(recording_path); + } + Ok(Some( + Box::new(builder.spawn()) as Box + Send + Sync> + )) }, ) } + +fn prepare_recording_path(path: &Path) -> io::Result<()> { + if let Some(parent) = path.parent() { + fs::create_dir_all(parent).map_err(|source| { + io::Error::new( + source.kind(), + format!( + "failed to create Tokio console recording directory `{}`: {source}", + path.display() + ), + ) + })?; + } + + fs::OpenOptions::new() + .create(true) + .append(true) + .open(path) + .map(|_| ()) + .map_err(|source| { + io::Error::new( + source.kind(), + format!( + "failed to open Tokio console recording path `{}`: {source}", + path.display() + ), + ) + }) +} + +#[cfg(test)] +mod tests { + use std::{fs, thread, time::Duration}; + + use console_subscriber::ConsoleLayer; + use tracing::Level; + use tracing_subscriber::layer::SubscriberExt as _; + + // Running this requires the same `--cfg tokio_unstable` flag as a node + // binary built with the Tokio Console feature. + #[ignore = "requires RUSTFLAGS=--cfg tokio_unstable"] + #[test] + fn configured_recorder_creates_non_empty_raw_recording() { + let path = + std::env::temp_dir().join(format!("logos-tokio-console-{}.jsonl", std::process::id())); + drop(fs::remove_file(&path)); + + let (layer, _server) = ConsoleLayer::builder().recording_path(&path).build(); + let subscriber = tracing_subscriber::registry().with(layer); + tracing::subscriber::with_default(subscriber, || { + let span = tracing::span!(target: "tokio", Level::INFO, "runtime.test"); + let _entered = span.enter(); + }); + + let non_empty = (0..20).any(|_| { + thread::sleep(Duration::from_millis(50)); + fs::metadata(&path).is_ok_and(|metadata| metadata.len() > 0) + }); + assert!(non_empty); + fs::remove_file(path).expect("test recording should be removable"); + } +} diff --git a/services/tracing/src/lib.rs b/services/tracing/src/lib.rs index e593f2eea..b1d226123 100644 --- a/services/tracing/src/lib.rs +++ b/services/tracing/src/lib.rs @@ -205,6 +205,7 @@ pub enum MetricsLayerSettings { pub struct TokioConsoleConfig { pub bind_address: String, pub port: u16, + pub recording_path: Option, } #[derive(Clone, Debug, Serialize, Deserialize)] @@ -357,15 +358,30 @@ where .. } = logger_layers; + #[cfg(feature = "tokio-console")] + let console_layer = match &config.console { + ConsoleLayerSettings::Console(console_config) => { + console::create_console_layer::(console_config)? + } + ConsoleLayerSettings::None => None, + }; + ONCE_INIT.call_once(move || { let mut layers: Vec + Send + Sync>> = vec![]; let level_filter = { #[cfg(feature = "tokio-console")] { - if let ConsoleLayerSettings::Console(console_config) = &config.console - && let Some(console_layer) = console::create_console_layer(console_config) - { + if let Some(console_layer) = console_layer { + if let ConsoleLayerSettings::Console(console_config) = &config.console + && let Some(recording_path) = &console_config.recording_path + { + tracing::info!( + target: LOG_TARGET, + "Tokio console raw recording is enabled at `{}`", + recording_path.display() + ); + } layers.push(console_layer); LevelFilter::TRACE } else { diff --git a/tests/cucumber_tests/features/manual_control.feature b/tests/cucumber_tests/features/manual_control.feature index 91500e549..0b163b746 100644 --- a/tests/cucumber_tests/features/manual_control.feature +++ b/tests/cucumber_tests/features/manual_control.feature @@ -75,9 +75,9 @@ Feature: Manual control of transactions | 1 | 3 | 1000000 | | 2 | 3 | 1000000 | And I will have tokio console profile nodes: - | node_name | - | NODE_1 | - | NODE_2 | + | node_name | record_raw | + | NODE_1 | true | + | NODE_2 | false | And I have a cluster with capacity of 2 nodes # And we use IBD peers And all peers must be mode online after startup in 30 seconds diff --git a/tests/src/cucumber/steps/manual_nodes/utils.rs b/tests/src/cucumber/steps/manual_nodes/utils.rs index a0482f0c2..8c62ee3a6 100644 --- a/tests/src/cucumber/steps/manual_nodes/utils.rs +++ b/tests/src/cucumber/steps/manual_nodes/utils.rs @@ -1064,6 +1064,9 @@ fn prepare_config_patch( config.user.tracing.console = ConsoleLayer::Console(TokioConfig { bind_address: IpAddr::V4(Ipv4Addr::LOCALHOST), port: node.port, + recording_path: node + .record_raw + .then(|| PathBuf::from("tokio-console-raw.jsonl")), }); } diff --git a/tests/src/cucumber/steps/tokio_console/profile.rs b/tests/src/cucumber/steps/tokio_console/profile.rs index 38054cb9f..c47845cb2 100644 --- a/tests/src/cucumber/steps/tokio_console/profile.rs +++ b/tests/src/cucumber/steps/tokio_console/profile.rs @@ -3,6 +3,7 @@ use std::collections::BTreeMap; #[derive(Debug, Clone)] pub struct TokioConsoleProfileNode { pub port: u16, + pub record_raw: bool, } #[derive(Debug, Clone, Default)] diff --git a/tests/src/cucumber/steps/tokio_console/steps.rs b/tests/src/cucumber/steps/tokio_console/steps.rs index 8c91a5644..63e369bab 100644 --- a/tests/src/cucumber/steps/tokio_console/steps.rs +++ b/tests/src/cucumber/steps/tokio_console/steps.rs @@ -6,7 +6,10 @@ use lb_testing_framework::get_reserved_available_tcp_port; use crate::cucumber::{ error::{StepError, StepResult}, steps::{ - parse_steps::{invalid_table_row, parse_required_string_cell, parse_table_rows}, + parse_steps::{ + invalid_table_row, parse_required_bool_cell, parse_required_string_cell, + parse_table_rows, + }, tokio_console::profile::{TokioConsoleProfile, TokioConsoleProfileNode}, }, world::CucumberWorld, @@ -16,6 +19,7 @@ const DEFAULT_BASE_PORT: u16 = 16_669; pub(super) struct TokioConsoleResourcesRow { pub node_name: String, + pub record_raw: bool, } #[given("I will have tokio console profile nodes:")] @@ -36,6 +40,7 @@ fn step_tokio_console_profile_nodes(world: &mut CucumberWorld, step: &Step) -> S row.node_name, TokioConsoleProfileNode { port: get_reserved_available_tcp_port().unwrap_or(DEFAULT_BASE_PORT), + record_raw: row.record_raw, }, ) }) @@ -48,7 +53,7 @@ fn step_tokio_console_profile_nodes(world: &mut CucumberWorld, step: &Step) -> S pub(super) fn tokio_console_rows(step: &Step) -> Result, StepError> { let rows = parse_table_rows( step, - &["node_name"], + &["node_name", "record_raw"], "Tokio console", parse_node_resource_row, )?; @@ -58,11 +63,15 @@ pub(super) fn tokio_console_rows(step: &Step) -> Result Result { match row { - [node_name] => { + [node_name, record_raw] => { let node_name = parse_required_string_cell(node_name, "node_name", "Tokio console")?; + let record_raw = parse_required_bool_cell(record_raw, "record_raw", "Tokio console")?; - Ok(TokioConsoleResourcesRow { node_name }) + Ok(TokioConsoleResourcesRow { + node_name, + record_raw, + }) } - _ => invalid_table_row("Tokio console", &["node_name"], row.len()), + _ => invalid_table_row("Tokio console", &["node_name", "record_raw"], row.len()), } } diff --git a/tests/testing_framework/src/framework/local/provisioning.rs b/tests/testing_framework/src/framework/local/provisioning.rs index e32908c8b..c9b00d7eb 100644 --- a/tests/testing_framework/src/framework/local/provisioning.rs +++ b/tests/testing_framework/src/framework/local/provisioning.rs @@ -170,6 +170,24 @@ impl LocalDeployerEnv for LbcEnv { config.user.state.base_folder = dir.to_path_buf(); "db".clone_into(&mut config.user.storage.backend.folder_name); + if let config::tracing::serde::console::Layer::Console(console) = + &mut config.user.tracing.console + && let Some(recording_path) = &mut console.recording_path + { + if recording_path.is_relative() { + let relative_path = recording_path.clone(); + *recording_path = dir.join(relative_path); + } + if let Some(parent) = recording_path.parent() { + fs::create_dir_all(parent).map_err(|source| { + io::Error::other(format!( + "failed to prepare Tokio console recording path `{}`: {source}", + recording_path.display() + )) + })?; + } + } + let user_yaml = serde_yaml::to_string(&config.user).map_err(io::Error::other)?; let deployment_yaml = serde_yaml::to_string(&config.deployment).map_err(io::Error::other)?;