diff --git a/Cargo.lock b/Cargo.lock index 6148d82b7..910d1b9d2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4675,6 +4675,7 @@ dependencies = [ "logos-blockchain-key-management-system-keys", "logos-blockchain-log-targets", "logos-blockchain-poseidon2", + "logos-blockchain-utils", "logos-blockchain-utxotree", "tokio", "tracing", @@ -5000,6 +5001,7 @@ dependencies = [ "logos-blockchain-log-targets", "logos-blockchain-services-utils", "logos-blockchain-tracing", + "logos-blockchain-utils", "overwatch", "rocksdb", "serde", diff --git a/blend/core/Cargo.toml b/blend/core/Cargo.toml index 473487a5a..0023653c8 100644 --- a/blend/core/Cargo.toml +++ b/blend/core/Cargo.toml @@ -20,6 +20,7 @@ lb-blend-proofs = { workspace = true } lb-blend-scheduling = { workspace = true } [features] +tokio-task-names = ["lb-blend-scheduling/tokio-task-names"] unsafe-test-functions = [ "lb-blend-message/unsafe-test-functions", "lb-blend-network/unsafe-test-functions", diff --git a/blend/scheduling/Cargo.toml b/blend/scheduling/Cargo.toml index 10b3cac47..cd5b93c9a 100644 --- a/blend/scheduling/Cargo.toml +++ b/blend/scheduling/Cargo.toml @@ -42,4 +42,5 @@ libp2p = { workspace = true } test-log = { features = ["trace"], workspace = true } [features] +tokio-task-names = ["lb-utils/tokio-task-names"] unsafe-test-functions = [] diff --git a/blend/scheduling/src/message_blend/provers/core/mod.rs b/blend/scheduling/src/message_blend/provers/core/mod.rs index 93e039742..c30fdf310 100644 --- a/blend/scheduling/src/message_blend/provers/core/mod.rs +++ b/blend/scheduling/src/message_blend/provers/core/mod.rs @@ -9,8 +9,8 @@ use lb_blend_proofs::{quota::inputs::prove::PublicInputs, selection::VerifiedPro use lb_groth16::fr_to_bytes; use lb_key_management_system_keys::keys::UnsecuredEd25519Key; use lb_log_targets::blend; -use lb_utils::tokio::stream::Buffered; -use tokio::{task::spawn, time::Instant}; +use lb_utils::tokio::{stream::Buffered, task::spawn}; +use tokio::time::Instant; use crate::message_blend::{ CoreProofOfQuotaGenerator, buffer_size, @@ -102,7 +102,7 @@ where // Without this, `generate_poq` would only begin when `FuturesOrdered` first // polls the future — which only happens when the consumer polls the stream — // causing avoidable latency when the consumer is idle. - let task = spawn(async move { + let task = spawn("logos/blend/core-proof-generator", async move { let (proof_of_quota, secret_selection_randomness) = proof_of_quota_generator .generate_poq( &PublicInputs { diff --git a/blend/scheduling/src/message_blend/provers/leader/mod.rs b/blend/scheduling/src/message_blend/provers/leader/mod.rs index fb90ac0d1..fd3c6cc4f 100644 --- a/blend/scheduling/src/message_blend/provers/leader/mod.rs +++ b/blend/scheduling/src/message_blend/provers/leader/mod.rs @@ -16,8 +16,8 @@ use lb_cryptarchia_engine::Epoch; use lb_groth16::fr_to_bytes; use lb_key_management_system_keys::keys::UnsecuredEd25519Key; use lb_log_targets::blend; -use lb_utils::tokio::stream::Buffered; -use tokio::{task::spawn_blocking, time::Instant}; +use lb_utils::tokio::{stream::Buffered, task::spawn_blocking}; +use tokio::time::Instant; use crate::message_blend::{ buffer_size, @@ -108,7 +108,7 @@ fn create_proof_stream( // Without this, `spawn_blocking` would only be called when `FuturesOrdered` // first polls the future — which only happens when the consumer polls the // stream — causing avoidable latency when the consumer is idle. - let task = spawn_blocking(move || { + let task = spawn_blocking("logos/blend/leader-poq-blocking", move || { let ephemeral_signing_key = UnsecuredEd25519Key::generate_with_blake_rng(); let (proof_of_quota, secret_selection_randomness) = VerifiedProofOfQuota::new( &PublicInputs { diff --git a/kms/operators/Cargo.toml b/kms/operators/Cargo.toml index 0f96d614a..2dd521b48 100644 --- a/kms/operators/Cargo.toml +++ b/kms/operators/Cargo.toml @@ -20,9 +20,11 @@ lb-groth16 = { workspace = true } lb-key-management-system-keys = { features = ["unsafe"], workspace = true } lb-log-targets = { workspace = true } lb-poseidon2 = { workspace = true } +lb-utils = { features = ["tokio"], workspace = true } lb-utxotree = { workspace = true } tokio = { features = ["rt"], workspace = true } tracing = { workspace = true } [features] -unsafe = [] +tokio-task-names = ["lb-utils/tokio-task-names"] +unsafe = [] diff --git a/kms/operators/src/blend/poq.rs b/kms/operators/src/blend/poq.rs index 93fc86998..314ad934d 100644 --- a/kms/operators/src/blend/poq.rs +++ b/kms/operators/src/blend/poq.rs @@ -13,7 +13,8 @@ use lb_key_management_system_keys::keys::{ secured_key::{SecureKeyOperator, SecuredKey}, }; use lb_log_targets::kms; -use tokio::{sync::oneshot, task::spawn_blocking}; +use lb_utils::tokio::task::spawn_blocking; +use tokio::sync::oneshot; use tracing::trace; const LOG_TARGET: &str = kms::operators::BLEND_POQ; @@ -64,10 +65,11 @@ impl SecureKeyOperator for PoQOperator { let public_inputs = self.public_inputs; // spawn a blocking task as this computation is heavy atm because it needs of an // external binary. - let poq_result = - spawn_blocking(move || VerifiedProofOfQuota::new(&public_inputs, private_inputs)) - .await - .map_err(Self::Error::FailedOperatorCall)?; + let poq_result = spawn_blocking("logos/blend/core-poq-blocking", move || { + VerifiedProofOfQuota::new(&public_inputs, private_inputs) + }) + .await + .map_err(Self::Error::FailedOperatorCall)?; drop(self.response_channel.send(poq_result).inspect_err(|_| { trace!(target: LOG_TARGET, "Error sending generated proof of quota, most likely due to an epoch rotation that discarded the receiver side of the channel."); })); diff --git a/nodes/node/binary/Cargo.toml b/nodes/node/binary/Cargo.toml index aaa8714f5..4ed88bd0b 100644 --- a/nodes/node/binary/Cargo.toml +++ b/nodes/node/binary/Cargo.toml @@ -101,5 +101,19 @@ 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", "overwatch/tokio-task-names", "tokio/tracing"] -tracing = [] +tokio-console = [ + "lb-blend-service/tokio-task-names", + "lb-blend/tokio-task-names", + "lb-chain-leader-service/tokio-task-names", + "lb-chain-network-service/tokio-task-names", + "lb-key-management-system-service/tokio-task-names", + "lb-network-service/tokio-task-names", + "lb-storage-service/tokio-task-names", + "lb-tracing-service/tokio-console", + "lb-tx-service/tokio-task-names", + "lb-utils/tokio-task-names", + "lb-wallet-service/tokio-task-names", + "overwatch/tokio-task-names", + "tokio/tracing", +] +tracing = [] diff --git a/services/blend/Cargo.toml b/services/blend/Cargo.toml index 73b4f9c6e..49e463807 100644 --- a/services/blend/Cargo.toml +++ b/services/blend/Cargo.toml @@ -51,4 +51,5 @@ libp2p-swarm-test = { workspace = true } test-log = { features = ["trace"], workspace = true } [features] -default = [] +default = [] +tokio-task-names = ["lb-utils/tokio-task-names"] diff --git a/services/blend/src/core/backends/libp2p/mod.rs b/services/blend/src/core/backends/libp2p/mod.rs index 6fb752864..f6156e0ee 100644 --- a/services/blend/src/core/backends/libp2p/mod.rs +++ b/services/blend/src/core/backends/libp2p/mod.rs @@ -10,6 +10,7 @@ use lb_blend::message::encap::validated::{ }; use lb_chain_service::Epoch; use lb_log_targets::blend; +use lb_utils::tokio::task::spawn_on; use libp2p::PeerId; use overwatch::overwatch::handle::OverwatchHandle; use rand::RngCore; @@ -80,9 +81,11 @@ where }); let (swarm_task_abort_handle, swarm_task_abort_registration) = AbortHandle::new_pair(); - overwatch_handle - .runtime() - .spawn(Abortable::new(swarm.run(), swarm_task_abort_registration)); + spawn_on( + overwatch_handle.runtime(), + "logos/blend/libp2p-swarm", + Abortable::new(swarm.run(), swarm_task_abort_registration), + ); Self { swarm_task_abort_handle, diff --git a/services/chain/chain-leader/Cargo.toml b/services/chain/chain-leader/Cargo.toml index 53eff8a31..d9e1d66e3 100644 --- a/services/chain/chain-leader/Cargo.toml +++ b/services/chain/chain-leader/Cargo.toml @@ -28,7 +28,7 @@ lb-storage-service = { workspace = true } lb-time-service = { workspace = true } lb-tracing = { workspace = true } lb-tx-service = { workspace = true } -lb-utils = { workspace = true } +lb-utils = { features = ["tokio"], workspace = true } lb-wallet-service = { workspace = true } overwatch = { workspace = true } rand = { workspace = true } @@ -47,3 +47,4 @@ default = [] # Suppresses publishing of locally-proposed blocks over the blend network. # The block is still self-applied to this node's chain. Test-only build flag. testing-disable-proposal-publish = [] +tokio-task-names = ["lb-utils/tokio-task-names"] diff --git a/services/chain/chain-leader/src/leadership.rs b/services/chain/chain-leader/src/leadership.rs index a29ab67ed..30804114e 100644 --- a/services/chain/chain-leader/src/leadership.rs +++ b/services/chain/chain-leader/src/leadership.rs @@ -19,6 +19,7 @@ use lb_key_management_system_service::{ }; use lb_ledger::{EpochState, UtxoTree}; use lb_time_service::{EpochSlotTickStream, SlotTick, TimeServiceMessage}; +use lb_utils::tokio::task::spawn_blocking; use lb_wallet_service::{ UtxoWithKeyId, api::{WalletApi, WalletApiError, WalletServiceData}, @@ -117,7 +118,7 @@ where } }; - let res = tokio::task::spawn_blocking(move || { + let res = spawn_blocking("logos/chain/leader-proof-blocking", move || { Groth16LeaderProof::prove(private_inputs, voucher_cm) }) .await; diff --git a/services/chain/chain-leader/src/lib.rs b/services/chain/chain-leader/src/lib.rs index 6929a56ba..ddb1447a2 100644 --- a/services/chain/chain-leader/src/lib.rs +++ b/services/chain/chain-leader/src/lib.rs @@ -39,6 +39,7 @@ use lb_tx_service::{ network::NetworkAdapter as MempoolNetworkAdapter, storage::MempoolStorageAdapter, }; +use lb_utils::tokio::task::spawn; use lb_wallet_service::api::{WalletApi, WalletApiError}; use overwatch::{ DynError, OpaqueServiceResourcesHandle, @@ -732,14 +733,17 @@ where // channel providing backpressure. let (epoch_handoff_sender, epoch_handoff_receiver) = mpsc::channel(WINNING_POL_EPOCH_HANDOFF_BUFFER_SIZE); - tokio::spawn(search_for_winning_slots( - (*cryptarchia).clone(), - (*wallet).clone(), - (*kms).clone(), - (*time_relay).clone(), - (*ledger_config).clone(), - epoch_handoff_sender, - )); + spawn( + "logos/chain/winning-slot-scanner", + search_for_winning_slots( + (*cryptarchia).clone(), + (*wallet).clone(), + (*kms).clone(), + (*time_relay).clone(), + (*ledger_config).clone(), + epoch_handoff_sender, + ), + ); let stream: WinningPolEpochSlotsStream = Box::pin(ReceiverStream::new(epoch_handoff_receiver)); if sender.send(stream).is_err() { diff --git a/services/chain/chain-network/Cargo.toml b/services/chain/chain-network/Cargo.toml index 60241815d..e59a36171 100644 --- a/services/chain/chain-network/Cargo.toml +++ b/services/chain/chain-network/Cargo.toml @@ -28,7 +28,7 @@ lb-storage-service = { workspace = true } lb-time-service = { workspace = true } lb-tracing = { workspace = true } lb-tx-service = { workspace = true } -lb-utils = { workspace = true } +lb-utils = { features = ["tokio"], workspace = true } lru = { workspace = true } overwatch = { workspace = true } rand = { workspace = true } @@ -46,4 +46,5 @@ lb-network-service = { workspace = true } lb-utils = { workspace = true } [features] -default = [] +default = [] +tokio-task-names = ["lb-utils/tokio-task-names"] diff --git a/services/chain/chain-network/src/lib.rs b/services/chain/chain-network/src/lib.rs index 539a4ad3b..890542b57 100644 --- a/services/chain/chain-network/src/lib.rs +++ b/services/chain/chain-network/src/lib.rs @@ -34,7 +34,7 @@ use lb_tx_service::{ TxMempoolService, backend::RecoverableMempool, network::NetworkAdapter as MempoolNetworkAdapter, storage::MempoolStorageAdapter, }; -use lb_utils::bounded::BoundedError; +use lb_utils::{bounded::BoundedError, tokio::task::spawn}; use network::NetworkAdapter; use overwatch::{ DynError, OpaqueServiceResourcesHandle, @@ -462,7 +462,7 @@ where let adapter = tip_poll_adapter.clone(); let cryptarchia = relays.cryptarchia().clone(); let tx = polled_tip_tx.clone(); - tip_poll_task = Some(tokio::spawn(async move { + tip_poll_task = Some(spawn("logos/chain/tip-poll", async move { if let Some(polled) = poll_peer_tips_if_behind( &adapter, &cryptarchia, diff --git a/services/key-management-system/Cargo.toml b/services/key-management-system/Cargo.toml index ac579d51d..d5525caa0 100644 --- a/services/key-management-system/Cargo.toml +++ b/services/key-management-system/Cargo.toml @@ -33,4 +33,5 @@ serde = { features = ["std"], workspace = true } serde_yaml = { workspace = true } [features] -unsafe = ["lb-key-management-system-keys/unsafe", "lb-key-management-system-operators/unsafe"] +tokio-task-names = ["lb-key-management-system-operators/tokio-task-names"] +unsafe = ["lb-key-management-system-keys/unsafe", "lb-key-management-system-operators/unsafe"] diff --git a/services/network/Cargo.toml b/services/network/Cargo.toml index 7cace874b..e3540b161 100644 --- a/services/network/Cargo.toml +++ b/services/network/Cargo.toml @@ -20,6 +20,7 @@ lb-cryptarchia-sync = { workspace = true } lb-libp2p = { workspace = true } lb-log-targets = { workspace = true } lb-tracing = { workspace = true } +lb-utils = { features = ["tokio"], workspace = true } overwatch = { workspace = true } rand = { features = ["std"], workspace = true } rand_chacha = { workspace = true } @@ -35,5 +36,6 @@ lb-utils = { workspace = true } tracing-subscriber = { features = ["env-filter", "fmt", "std"], workspace = true } [features] -default = [] -openapi = ["dep:utoipa"] +default = [] +openapi = ["dep:utoipa"] +tokio-task-names = ["lb-utils/tokio-task-names"] diff --git a/services/network/src/backends/libp2p/mod.rs b/services/network/src/backends/libp2p/mod.rs index d054ad459..400668ba8 100644 --- a/services/network/src/backends/libp2p/mod.rs +++ b/services/network/src/backends/libp2p/mod.rs @@ -7,6 +7,7 @@ pub use lb_libp2p::{ libp2p::gossipsub::{Message, TopicHash}, }; use lb_log_targets::network_service; +use lb_utils::tokio::task::spawn_on; use overwatch::overwatch::handle::OverwatchHandle; use rand::SeedableRng as _; use rand_chacha::ChaCha20Rng; @@ -58,9 +59,13 @@ impl NetworkBackend for Libp2p { rng, ); - overwatch_handle.runtime().spawn(async move { - swarm_handler.run(initial_peers).await; - }); + spawn_on( + overwatch_handle.runtime(), + "logos/network/libp2p-swarm", + async move { + swarm_handler.run(initial_peers).await; + }, + ); Self { pubsub_events_tx, diff --git a/services/network/src/backends/libp2p/swarm/gossipsub.rs b/services/network/src/backends/libp2p/swarm/gossipsub.rs index de4e66242..9717833f6 100644 --- a/services/network/src/backends/libp2p/swarm/gossipsub.rs +++ b/services/network/src/backends/libp2p/swarm/gossipsub.rs @@ -1,5 +1,6 @@ use lb_libp2p::{behaviour::gossipsub::swarm_ext::topic_hash, gossipsub}; use lb_log_targets::network_service; +use lb_utils::tokio::task::spawn; use rand::RngCore; use crate::backends::libp2p::{ @@ -88,7 +89,7 @@ impl SwarmHandler { ); let commands_tx = self.commands_tx.clone(); - tokio::spawn(async move { + spawn("logos/network/gossipsub-retry", async move { tokio::time::sleep(wait).await; let Some(new_retry_count) = retry_count.checked_add(1) else { tracing::error!(target: LOG_TARGET, "retry count overflow."); diff --git a/services/network/src/backends/libp2p/swarm/mod.rs b/services/network/src/backends/libp2p/swarm/mod.rs index 06d310184..af96c1518 100644 --- a/services/network/src/backends/libp2p/swarm/mod.rs +++ b/services/network/src/backends/libp2p/swarm/mod.rs @@ -28,6 +28,7 @@ use lb_libp2p::{ }, }; use lb_log_targets::network_service; +use lb_utils::tokio::task::spawn; use rand::RngCore; use tokio::sync::{broadcast, mpsc, oneshot}; use tokio_stream::StreamExt as _; @@ -355,7 +356,7 @@ impl SwarmHandler { tracing::debug!(target: LOG_TARGET, "Retry dialing in {wait:?}: {dial:?}"); let commands_tx = self.commands_tx.clone(); - tokio::spawn(async move { + spawn("logos/network/dial-retry", async move { tokio::time::sleep(wait).await; Self::schedule_connect(dial, commands_tx).await; }); diff --git a/services/storage/Cargo.toml b/services/storage/Cargo.toml index 9ef00e2a7..7675dd6ed 100644 --- a/services/storage/Cargo.toml +++ b/services/storage/Cargo.toml @@ -21,6 +21,7 @@ lb-cryptarchia-engine = { workspace = true } lb-log-targets = { workspace = true } lb-services-utils = { workspace = true } lb-tracing = { workspace = true } +lb-utils = { features = ["tokio"], workspace = true } overwatch = { workspace = true } rocksdb = { features = ["bindgen-runtime"], optional = true, workspace = true } serde = { workspace = true } @@ -32,8 +33,9 @@ tracing = { workspace = true } tempfile = { workspace = true } [features] -default = [] -rocksdb-backend = ["dep:rocksdb"] +default = [] +rocksdb-backend = ["dep:rocksdb"] +tokio-task-names = ["lb-utils/tokio-task-names"] [[bin]] name = "logos-blockchain-rocksdb" diff --git a/services/storage/src/backends/rocksdb.rs b/services/storage/src/backends/rocksdb.rs index 3b958797c..870d766b2 100644 --- a/services/storage/src/backends/rocksdb.rs +++ b/services/storage/src/backends/rocksdb.rs @@ -2,6 +2,7 @@ use std::{collections::HashMap, num::NonZeroUsize, path::PathBuf, sync::Arc}; use async_trait::async_trait; use bytes::Bytes; +use lb_utils::tokio::task::spawn_blocking; use rocksdb::{DB, Direction, Error, IteratorMode, Options}; use serde::{Deserialize, Serialize}; @@ -142,7 +143,7 @@ impl StorageBackend for RocksBackend { // Use spawn_blocking to avoid blocking the async runtime during the bulk // operation - tokio::task::spawn_blocking(move || { + spawn_blocking("logos/storage/rocksdb-bulk-store-blocking", move || { let mut batch = rocksdb::WriteBatch::default(); let mut has_items = false; diff --git a/services/tracing/src/lib.rs b/services/tracing/src/lib.rs index b1d226123..63cb9dee3 100644 --- a/services/tracing/src/lib.rs +++ b/services/tracing/src/lib.rs @@ -286,6 +286,7 @@ impl ServiceCore for Tracing + Display + Send, { + #[expect(clippy::too_many_lines, reason = "TODO: Address this at some point.")] fn init( service_resources_handle: OpaqueServiceResourcesHandle, _initial_state: Self::State, @@ -369,6 +370,8 @@ where ONCE_INIT.call_once(move || { let mut layers: Vec + Send + Sync>> = vec![]; + #[cfg(feature = "tokio-console")] + let mut display_tokio_console_msg = None; let level_filter = { #[cfg(feature = "tokio-console")] { @@ -376,11 +379,10 @@ where if let ConsoleLayerSettings::Console(console_config) = &config.console && let Some(recording_path) = &console_config.recording_path { - tracing::info!( - target: LOG_TARGET, + display_tokio_console_msg = Some(format!( "Tokio console raw recording is enabled at `{}`", recording_path.display() - ); + )); } layers.push(console_layer); LevelFilter::TRACE @@ -403,6 +405,11 @@ where .with(level_filter) .with(layers) .init(); + + #[cfg(feature = "tokio-console")] + if let Some(msg) = display_tokio_console_msg { + tracing::info!(target: LOG_TARGET, "{msg}"); + } }); Ok(Self { diff --git a/services/tx-service/Cargo.toml b/services/tx-service/Cargo.toml index cccc59759..43c2cab36 100644 --- a/services/tx-service/Cargo.toml +++ b/services/tx-service/Cargo.toml @@ -29,7 +29,7 @@ lb-network-service = { workspace = true } lb-services-utils = { workspace = true } lb-storage-service = { workspace = true } lb-tracing = { workspace = true } -lb-utils = { workspace = true } +lb-utils = { features = ["tokio"], workspace = true } overwatch = { workspace = true } serde = { workspace = true } serde_json = { optional = true, workspace = true } @@ -50,4 +50,5 @@ default = [] rocksdb-backend = ["lb-storage-service/rocksdb-backend"] # enable to help generate OpenAPI -openapi = ["dep:serde_json", "dep:utoipa"] +openapi = ["dep:serde_json", "dep:utoipa"] +tokio-task-names = ["lb-utils/tokio-task-names"] diff --git a/services/tx-service/src/tx/service.rs b/services/tx-service/src/tx/service.rs index 3552a0545..852418095 100644 --- a/services/tx-service/src/tx/service.rs +++ b/services/tx-service/src/tx/service.rs @@ -24,6 +24,7 @@ use lb_services_utils::{ wait_until_services_are_ready, }; use lb_storage_service::{StorageService, recovery::StorageRecoveryBackend}; +use lb_utils::tokio::task::spawn; use overwatch::{ OpaqueServiceResourcesHandle, services::{AsServiceId, ServiceCore, ServiceData, relay::OutboundRelay}, @@ -395,7 +396,7 @@ where Err(MempoolError::ExistingItem) => { // Tx already in pool, but since this came from a local submission // (not gossip), re-gossip it so leader nodes can pick it up. - tokio::spawn(async move { + spawn("logos/mempool/transaction-regossip", async move { let adapter = NetworkAdapter::new(settings, network_relay).await; adapter.send(item_for_broadcast).await; }); @@ -496,7 +497,7 @@ where ) { state_updater.update(Some(::save(pool).into())); - tokio::spawn(async move { + spawn("logos/mempool/transaction-broadcast", async move { let adapter = NetworkAdapter::new(settings, network_relay).await; adapter.send(item_for_broadcast).await; }); diff --git a/services/wallet/Cargo.toml b/services/wallet/Cargo.toml index f32fdbc03..e6691e696 100644 --- a/services/wallet/Cargo.toml +++ b/services/wallet/Cargo.toml @@ -29,7 +29,7 @@ lb-log-targets = { workspace = true } lb-mmr = { workspace = true } lb-services-utils = { workspace = true } lb-storage-service = { workspace = true } -lb-utils = { workspace = true } +lb-utils = { features = ["tokio"], workspace = true } lb-wallet = { workspace = true } overwatch = { workspace = true } @@ -42,3 +42,6 @@ ignored = [ # Required by `wait_until_services_are_ready` macro. "futures", ] + +[features] +tokio-task-names = ["lb-utils/tokio-task-names"] diff --git a/services/wallet/src/lib.rs b/services/wallet/src/lib.rs index 9d1452304..779ecb8d1 100644 --- a/services/wallet/src/lib.rs +++ b/services/wallet/src/lib.rs @@ -52,7 +52,7 @@ use lb_services_utils::{ use lb_storage_service::{ api::chain::StorageChainApi, backends::StorageBackend, recovery::StorageRecoveryBackend, }; -use lb_utils::bounded::BoundedError; +use lb_utils::{bounded::BoundedError, tokio::task::spawn_blocking}; use lb_wallet::{WalletBalance, WalletBlock, WalletError}; use overwatch::{ DynError, OpaqueServiceResourcesHandle, @@ -882,7 +882,7 @@ where let rewards_root = leader_claim_op.rewards_root; // TODO: This should happen in KMS - let poc = tokio::task::spawn_blocking(move || { + let poc = spawn_blocking("logos/wallet/leader-claim-proof-blocking", move || { Self::generate_poc(voucher_secret, &path, rewards_root, tx_hash) }) .await??; diff --git a/tests/cucumber_tests/features/manual_control.feature b/tests/cucumber_tests/features/manual_control.feature index 0b163b746..f1bb787bd 100644 --- a/tests/cucumber_tests/features/manual_control.feature +++ b/tests/cucumber_tests/features/manual_control.feature @@ -85,7 +85,7 @@ Feature: Manual control of transactions | node_name | account_index | wallet_name | connected_to | | NODE_1 | 1 | WALLET_1A | | | NODE_2 | 2 | WALLET_2A | NODE_1 | - When all nodes have at least 2 blocks and converged to within 1 blocks in 300 seconds + When all nodes have at least 2 blocks and converged to within 0 blocks in 300 seconds When I perform manual control of transactions for all wallets no time-out Then I stop all nodes diff --git a/utils/Cargo.toml b/utils/Cargo.toml index 1b6644a2e..645ef9378 100644 --- a/utils/Cargo.toml +++ b/utils/Cargo.toml @@ -33,10 +33,11 @@ tokio = { optional = true, workspace = true } tracing = { workspace = true } [features] -rng = ["dep:blake2"] -serde = ["const-hex/alloc", "serde/alloc"] -time = ["dep:humantime", "dep:serde_with", "dep:time"] -tokio = ["dep:futures", "dep:tokio"] +rng = ["dep:blake2"] +serde = ["const-hex/alloc", "serde/alloc"] +time = ["dep:humantime", "dep:serde_with", "dep:time"] +tokio = ["dep:futures", "dep:tokio"] +tokio-task-names = ["tokio", "tokio/rt", "tokio/tracing"] [dev-dependencies] bincode = { workspace = true } diff --git a/utils/src/tokio/mod.rs b/utils/src/tokio/mod.rs index baf29e06a..277af4c9c 100644 --- a/utils/src/tokio/mod.rs +++ b/utils/src/tokio/mod.rs @@ -1 +1,115 @@ pub mod stream; + +pub mod task { + use tokio::{runtime::Handle, task::JoinHandle}; + + #[expect( + unexpected_cfgs, + reason = "tokio_unstable is supplied externally through RUSTFLAGS" + )] + #[expect(clippy::allow_attributes, reason = "cfg-selected spawn implementation")] + #[allow(clippy::needless_return, reason = "cfg-selected spawn implementation")] + pub fn spawn( + name: &'static str, + future: impl Future + Send + 'static, + ) -> JoinHandle + where + T: Send + 'static, + { + #[cfg(all(feature = "tokio-task-names", tokio_unstable))] + { + return tokio::task::Builder::new() + .name(name) + .spawn(future) + .unwrap_or_else(|_| panic!("failed to spawn named Tokio task `{name}`")); + } + + #[cfg(not(all(feature = "tokio-task-names", tokio_unstable)))] + { + let _ = name; + tokio::spawn(future) + } + } + + #[expect( + unexpected_cfgs, + reason = "tokio_unstable is supplied externally through RUSTFLAGS" + )] + #[expect(clippy::allow_attributes, reason = "cfg-selected spawn implementation")] + #[allow(clippy::needless_return, reason = "cfg-selected spawn implementation")] + pub fn spawn_blocking( + name: &'static str, + function: impl FnOnce() -> T + Send + 'static, + ) -> JoinHandle + where + T: Send + 'static, + { + #[cfg(all(feature = "tokio-task-names", tokio_unstable))] + { + return tokio::task::Builder::new() + .name(name) + .spawn_blocking(function) + .unwrap_or_else(|_| panic!("failed to spawn named Tokio blocking task `{name}`")); + } + + #[cfg(not(all(feature = "tokio-task-names", tokio_unstable)))] + { + let _ = name; + tokio::task::spawn_blocking(function) + } + } + + #[expect( + unexpected_cfgs, + reason = "tokio_unstable is supplied externally through RUSTFLAGS" + )] + #[expect(clippy::allow_attributes, reason = "cfg-selected spawn implementation")] + #[allow(clippy::needless_return, reason = "cfg-selected spawn implementation")] + pub fn spawn_on( + runtime: &Handle, + name: &'static str, + future: impl Future + Send + 'static, + ) -> JoinHandle + where + T: Send + 'static, + { + #[cfg(all(feature = "tokio-task-names", tokio_unstable))] + { + return tokio::task::Builder::new() + .name(name) + .spawn_on(future, runtime) + .unwrap_or_else(|_| panic!("failed to spawn named Tokio task `{name}`")); + } + + #[cfg(not(all(feature = "tokio-task-names", tokio_unstable)))] + { + let _ = name; + runtime.spawn(future) + } + } + + #[cfg(test)] + mod tests { + use super::{spawn, spawn_blocking, spawn_on}; + + #[test] + fn spawn_forms_preserve_join_handle_results() { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("test runtime should build"); + let handle = runtime.handle().clone(); + + runtime.block_on(async move { + assert_eq!(spawn("test/async", async { 1 }).await.unwrap(), 1); + assert_eq!(spawn_blocking("test/blocking", || 2).await.unwrap(), 2); + assert_eq!( + spawn_on(&handle, "test/explicit", async { 3 }) + .await + .unwrap(), + 3 + ); + }); + } + } +}