From ea2c83470de86087735df2853be5d98f45d0f7ce Mon Sep 17 00:00:00 2001 From: Bohdan Ohorodnii <273991985+varex83agent@users.noreply.github.com> Date: Tue, 25 Aug 2026 17:02:01 +0200 Subject: [PATCH 1/2] refactor: qualify free-function calls and generalize public API params (#605) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Part A — qualify free-function calls: import parent modules/types instead of importing free (standalone) functions bare, and call them qualified (module::func()). Traits, types, structs, enums and constants stay imported bare. Touches ~30 production use statements across app, cli, cluster, consensus, core, dkg, eth2api, eth2util, p2p, priority, relay-server. Part B — generalize public API parameter types for consistency within modules: - eth2util::network: valid_network / network_to_genesis_time now take impl AsRef (matching sibling network_to_fork_version). - cluster::version: support_pregen_registrations / support_node_signatures now take impl AsRef. - k1util::load / save now take impl AsRef (matching the rest of the workspace). - build-proto::compile_protos now takes impl AsRef instead of &str. Part C — document both conventions in the rust-style skill: import modules/types not free functions; public APIs accept impl AsRef/impl AsRef/impl Into/impl IntoIterator where only the borrowed/converted form is needed. No behavior change. Co-Authored-By: Bohdan Ohorodnii <35969035+varex83@users.noreply.github.com> --- .claude/skills/rust-style/SKILL.md | 35 ++++++++++++- crates/app/src/monitoringapi/checker.rs | 4 +- crates/app/src/obolapi/exit.rs | 38 +++++++------- crates/build-proto/src/lib.rs | 16 ++++-- crates/cli/src/commands/create_cluster.rs | 13 +++-- crates/cli/src/commands/create_dkg.rs | 29 +++++----- crates/cli/src/commands/run.rs | 4 +- crates/cli/src/commands/test/helpers.rs | 5 +- crates/cli/src/commands/test/peers.rs | 23 ++++---- crates/cluster/src/definition.rs | 43 ++++++++------- crates/cluster/src/eip712sigs.rs | 9 ++-- crates/cluster/src/helpers.rs | 6 +-- crates/cluster/src/version.rs | 14 +++-- crates/core/src/bcast/mod.rs | 3 +- crates/core/src/signeddata.rs | 28 +++++----- crates/core/src/tracker/inclusion.rs | 4 +- crates/core/src/unsigneddata.rs | 19 +++---- crates/dkg/src/dkg.rs | 6 +-- crates/dkg/src/exchanger.rs | 4 +- crates/dkg/src/frost.rs | 6 +-- crates/dkg/src/frostp2p/behaviour.rs | 4 +- crates/dkg/src/signing.rs | 19 +++---- crates/eth2api/src/spec/serde_utils.rs | 5 +- crates/eth2util/src/network.rs | 4 +- crates/k1util/src/k1util.rs | 7 +-- crates/p2p/src/conn_logger.rs | 64 +++++++++++++---------- crates/p2p/src/force_direct.rs | 12 ++--- crates/p2p/src/k1.rs | 6 +-- crates/p2p/src/p2p.rs | 22 ++++---- crates/p2p/src/peer.rs | 6 +-- crates/p2p/src/quic_upgrade.rs | 24 ++++----- crates/p2p/src/relay/manager.rs | 4 +- crates/p2p/src/relay/manager/tests.rs | 2 +- crates/priority/src/calculate.rs | 4 +- crates/priority/src/component.rs | 6 +-- crates/priority/src/prioritiser.rs | 4 +- crates/relay-server/src/p2p.rs | 34 ++++++------ crates/relay-server/src/web.rs | 4 +- 38 files changed, 296 insertions(+), 244 deletions(-) diff --git a/.claude/skills/rust-style/SKILL.md b/.claude/skills/rust-style/SKILL.md index 46e73df3..5013dff6 100644 --- a/.claude/skills/rust-style/SKILL.md +++ b/.claude/skills/rust-style/SKILL.md @@ -127,9 +127,41 @@ Rules: - Prefer copying doc comments from Go and adapting to Rust conventions (avoid “Type is a …”). - Avoid leaving TODOs in merged code. If a short-lived internal note is necessary, use `// TODO:` and remove before PR merge. +## Imports + +Import **modules, types, traits, enums, and constants** — but **not free (standalone) functions**. Call free functions qualified through their parent module so the call site shows where the function comes from. + +```rust +// Bad — free function imported bare +use crate::name::peer_name; +let label = peer_name(&id); + +// Good — import the module, call qualified +use crate::name; +let label = name::peer_name(&id); +``` + +For a function re-exported at a crate root, qualify through the crate name (already in scope) rather than adding a bare `use`: + +```rust +// Bad +use pluto_k1util::load; +let key = load(path)?; + +// Good +let key = pluto_k1util::load(path)?; +``` + +Rules: + +- Types, structs, enums, and constants **should** be imported bare (`use foo::Bar;`, then `Bar`). +- Traits **must** be imported bare — they need to be in scope for method resolution. +- Only free functions are qualified through their module. When a `use` mixes types and a free function from the same module, add `self` and drop the function: `use foo::bar::{self, SomeType};`, then call `bar::the_fn()`. +- This mirrors how helpers such as `peer_name`, `hash_proto`, and `to_0x_hex` are called across the workspace (`module::func()`), keeping call sites self-documenting. + ## Generalized Parameter Types -Prefer generic parameters over concrete types when a function only needs the behavior of a trait. This mirrors the standard library's own conventions and makes functions callable with a wider range of inputs without extra allocations. +Prefer generic parameters over concrete types when a function only needs the behavior of a trait. This mirrors the standard library's own conventions and makes functions callable with a wider range of inputs without extra allocations. Apply this especially to **public APIs**, and keep it consistent across sibling functions in the same module — a module should not mix `&str` and `impl AsRef` for the same kind of argument. Full uniformity is not the goal: private helpers and hot internal paths may keep concrete types. | Instead of | Prefer | Accepts | | --- | --- | --- | @@ -138,6 +170,7 @@ Prefer generic parameters over concrete types when a function only needs the beh | `&[u8]` | `impl AsRef<[u8]>` | `&[u8]`, `Vec`, arrays, … | | `&Vec` | `impl AsRef<[T]>` | `Vec`, slices, arrays, … | | `String` (owned, read-only) | `impl Into` | `&str`, `String`, … | +| `&[T]` / iterator | `impl IntoIterator` | `Vec`, arrays, iterators, … | Examples: diff --git a/crates/app/src/monitoringapi/checker.rs b/crates/app/src/monitoringapi/checker.rs index 01ee4835..21c68291 100644 --- a/crates/app/src/monitoringapi/checker.rs +++ b/crates/app/src/monitoringapi/checker.rs @@ -22,7 +22,7 @@ use super::{ metrics::MONITORING_METRICS, readiness::{ReadinessError, ReadyResult, ReadyState}, }; -use crate::eth2wrap::version::check_beacon_node_version; +use crate::eth2wrap::version; /// Slots behind head after which the beacon node is considered too far behind. const BN_FAR_BEHIND_SLOTS: u64 = 320; @@ -128,7 +128,7 @@ async fn set_beacon_node_version(beacon_node: &EthBeaconNodeApiClient) { MONITORING_METRICS.beacon_node_version[&label].set(1); // The semantic compatibility check uses the FULL (untruncated) version. - check_beacon_node_version(&version); + version::check_beacon_node_version(&version); } /// Maximum length (in bytes) for the upstream-supplied beacon-node version diff --git a/crates/app/src/obolapi/exit.rs b/crates/app/src/obolapi/exit.rs index 8269d371..9861fbcb 100644 --- a/crates/app/src/obolapi/exit.rs +++ b/crates/app/src/obolapi/exit.rs @@ -9,18 +9,18 @@ use pluto_crypto::{tbls, types::Signature}; use serde::{Deserialize, Serialize}; use pluto_cluster::{ - helpers::to_0x_hex, + helpers, ssz::{SSZ_LEN_BLS_SIG, SSZ_LEN_PUB_KEY}, }; use pluto_eth2api::types::{ GetPoolVoluntaryExitsResponseResponseDatum, Phase0SignedVoluntaryExitMessage, }; -use pluto_ssz::{HashRoot, HashWalker, Hasher, put_bytes_n}; +use pluto_ssz::{HashRoot, HashWalker, Hasher}; use crate::obolapi::{ client::Client, error::{Error, Result}, - helper::{bearer_string, from_0x}, + helper, }; /// Type alias for signed voluntary exit from eth2api. @@ -44,8 +44,8 @@ impl SszHashable for SignedVoluntaryExit { let index = hh.index(); self.message.hash_with(hh)?; - let sig_bytes = from_0x(&self.signature, SSZ_LEN_BLS_SIG)?; - put_bytes_n(hh, &sig_bytes, SSZ_LEN_BLS_SIG)?; + let sig_bytes = helper::from_0x(&self.signature, SSZ_LEN_BLS_SIG)?; + pluto_ssz::put_bytes_n(hh, &sig_bytes, SSZ_LEN_BLS_SIG)?; hh.merkleize(index)?; Ok(()) @@ -88,7 +88,7 @@ impl SszHashable for ExitBlob { "missing public key".to_string(), )) })?; - let pk_bytes = from_0x(pk, SSZ_LEN_PUB_KEY)?; + let pk_bytes = helper::from_0x(pk, SSZ_LEN_PUB_KEY)?; hh.put_bytes(&pk_bytes)?; self.signed_exit_message.hash_with(hh)?; @@ -178,7 +178,7 @@ impl TryFrom for PartialExitRequest { type Error = Error; fn try_from(dto: PartialExitRequestDto) -> Result { - let signature = from_0x(&dto.signature, 65)?; + let signature = helper::from_0x(&dto.signature, 65)?; Ok(Self { unsigned: dto.unsigned, @@ -191,7 +191,7 @@ impl From for PartialExitRequestDto { fn from(req: PartialExitRequest) -> Self { Self { unsigned: req.unsigned, - signature: to_0x_hex(&req.signature), + signature: helpers::to_0x_hex(&req.signature), } } } @@ -233,7 +233,7 @@ impl SszHashable for FullExitAuthBlob { let index = hh.index(); hh.put_bytes(&self.lock_hash)?; - put_bytes_n(hh, &self.validator_pubkey, SSZ_LEN_PUB_KEY)?; + pluto_ssz::put_bytes_n(hh, &self.validator_pubkey, SSZ_LEN_PUB_KEY)?; hh.put_uint64(self.share_index)?; hh.merkleize(index)?; @@ -252,7 +252,7 @@ impl Client { identity_key: &k256::SecretKey, mut exit_blobs: Vec, ) -> Result<()> { - let lock_hash_str = to_0x_hex(lock_hash); + let lock_hash_str = helpers::to_0x_hex(lock_hash); let path = submit_partial_exit_url(&lock_hash_str); let url = self.build_url(&path)?; @@ -297,9 +297,9 @@ impl Client { identity_key: &k256::SecretKey, ) -> Result { // Validate public key is 48 bytes - let val_pubkey_bytes = from_0x(val_pubkey, 48)?; + let val_pubkey_bytes = helper::from_0x(val_pubkey, 48)?; - let path = fetch_full_exit_url(val_pubkey, &to_0x_hex(lock_hash), share_index); + let path = fetch_full_exit_url(val_pubkey, &helpers::to_0x_hex(lock_hash), share_index); let url = self.build_url(&path)?; @@ -316,7 +316,7 @@ impl Client { let headers = vec![( "Authorization".to_string(), - bearer_string(&lock_hash_signature), + helper::bearer_string(&lock_hash_signature), )]; let response_body = self.http_get(url, Some(&headers)).await?; @@ -337,7 +337,7 @@ impl Client { } // A BLS signature is 96 bytes long - let sig_bytes = from_0x(sig_str, 96)?; + let sig_bytes = helper::from_0x(sig_str, 96)?; // Convert to Signature type let mut sig = [0u8; 96]; @@ -364,7 +364,7 @@ impl Client { epoch: epoch_u64.to_string(), validator_index: exit_response.validator_index.to_string(), }, - signature: to_0x_hex(&full_sig), + signature: helpers::to_0x_hex(&full_sig), }, }) } @@ -380,9 +380,9 @@ impl Client { identity_key: &k256::SecretKey, ) -> Result<()> { // Validate public key is 48 bytes - let val_pubkey_bytes = from_0x(val_pubkey, 48)?; + let val_pubkey_bytes = helper::from_0x(val_pubkey, 48)?; - let path = delete_partial_exit_url(val_pubkey, &to_0x_hex(lock_hash), share_index); + let path = delete_partial_exit_url(val_pubkey, &helpers::to_0x_hex(lock_hash), share_index); let url = self.build_url(&path)?; @@ -398,7 +398,7 @@ impl Client { let headers = vec![( "Authorization".to_string(), - bearer_string(&lock_hash_signature), + helper::bearer_string(&lock_hash_signature), )]; self.http_delete(url, Some(&headers)).await?; @@ -475,7 +475,7 @@ mod tests { 404142434445464748494a4b4c4d4e4f505152535455565758595a5b5c5d5e5f"; let exit_blob = ExitBlob { - public_key: Some(to_0x_hex(&validator_pubkey)), + public_key: Some(helpers::to_0x_hex(&validator_pubkey)), signed_exit_message: SignedVoluntaryExit { message: Phase0SignedVoluntaryExitMessage { epoch: "194048".to_string(), diff --git a/crates/build-proto/src/lib.rs b/crates/build-proto/src/lib.rs index a10a0688..81fc313d 100644 --- a/crates/build-proto/src/lib.rs +++ b/crates/build-proto/src/lib.rs @@ -2,10 +2,15 @@ //! //! This crate compiles the protobuf files. -use std::{fs, io::Result, path::PathBuf}; +use std::{ + fs, + io::Result, + path::{Path, PathBuf}, +}; /// Compiles the protobuf files in the given directory. -pub fn compile_protos(proto_dir: &str) -> Result<()> { +pub fn compile_protos(proto_dir: impl AsRef) -> Result<()> { + let proto_dir = proto_dir.as_ref(); let proto_files: Vec = { let mut files: Vec = fs::read_dir(proto_dir)? .filter_map(|entry| entry.ok()) @@ -17,7 +22,10 @@ pub fn compile_protos(proto_dir: &str) -> Result<()> { }; if proto_files.is_empty() { - println!("cargo:warning=No .proto files found in {}", proto_dir); + println!( + "cargo:warning=No .proto files found in {}", + proto_dir.display() + ); return Ok(()); } @@ -44,7 +52,7 @@ pub fn compile_protos(proto_dir: &str) -> Result<()> { } /// Adds file attributes to the generated files. -fn add_file_attributes(proto_dir: &str) -> Result<()> { +fn add_file_attributes(proto_dir: &Path) -> Result<()> { let header = r#"// This file is @generated by prost-build. #![allow(dead_code)] #![allow(missing_docs)] diff --git a/crates/cli/src/commands/create_cluster.rs b/crates/cli/src/commands/create_cluster.rs index 68c8e68f..c7ac00e6 100644 --- a/crates/cli/src/commands/create_cluster.rs +++ b/crates/cli/src/commands/create_cluster.rs @@ -40,13 +40,12 @@ use pluto_eth2util::{ network, registration as eth2util_registration, }; use pluto_p2p::k1 as p2p_k1; -use pluto_ssz::to_0x_hex; use rand::rngs::OsRng; use tracing::{debug, info, warn}; use crate::{ commands::{ - address_validation::validate_addresses, + address_validation, constants::{MIN_NODES, MIN_THRESHOLD}, create_dkg, }, @@ -935,7 +934,7 @@ fn new_def_from_config(args: &CreateClusterArgs) -> Result { return Err(CreateClusterError::MissingNumValidatorsOrDefinitionFile); } - let (fee_recipient_addrs, withdrawal_addrs) = validate_addresses( + let (fee_recipient_addrs, withdrawal_addrs) = address_validation::validate_addresses( num_validators, &args.fee_recipient_addrs, &args.withdrawal_addrs, @@ -1201,7 +1200,7 @@ async fn load_definition( info!( url = def_file, - definition_hash = to_0x_hex(&def.definition_hash), + definition_hash = pluto_ssz::to_0x_hex(&def.definition_hash), "Cluster definition downloaded from URL" ); @@ -1213,7 +1212,7 @@ async fn load_definition( info!( path = def_file, - definition_hash = to_0x_hex(&def.definition_hash), + definition_hash = pluto_ssz::to_0x_hex(&def.definition_hash), "Cluster definition loaded from disk", ); @@ -2331,7 +2330,7 @@ mod tests { async fn multiple_addresses() { // "insufficient fee recipient addresses": 0 addrs for 4 validators → error { - let err = super::validate_addresses(4, &[], &[]).unwrap_err(); + let err = address_validation::validate_addresses(4, &[], &[]).unwrap_err(); let err_str = format!("{err}"); assert!( err_str.contains("mismatching --num-validators and --fee-recipient-addresses"), @@ -2343,7 +2342,7 @@ mod tests { // error { let fee_addr = "0x0000000000000000000000000000000000000000".to_string(); - let err = super::validate_addresses(1, &[fee_addr], &[]).unwrap_err(); + let err = address_validation::validate_addresses(1, &[fee_addr], &[]).unwrap_err(); let err_str = format!("{err}"); assert!( err_str.contains("mismatching --num-validators and --withdrawal-addresses"), diff --git a/crates/cli/src/commands/create_dkg.rs b/crates/cli/src/commands/create_dkg.rs index d39e8daf..8ef1b749 100644 --- a/crates/cli/src/commands/create_dkg.rs +++ b/crates/cli/src/commands/create_dkg.rs @@ -11,14 +11,12 @@ use pluto_cluster::{ definition::{Creator, Definition}, operator::Operator, }; -use pluto_consensus::protocols::is_supported_protocol_name; +use pluto_consensus::protocols; use pluto_eth2util::{ - deposit::{eths_to_gweis, verify_deposit_amounts}, + deposit, enr::Record, - helpers::{checksum_address, public_key_to_address}, - network::{ - GNOSIS, GOERLI, HOODI, MAINNET, PRATER, SEPOLIA, network_to_fork_version, valid_network, - }, + helpers, + network::{self, GNOSIS, GOERLI, HOODI, MAINNET, PRATER, SEPOLIA}, }; use thiserror::Error; use tracing::{info, warn}; @@ -328,7 +326,7 @@ async fn run_create_dkg(mut args: CreateDkgArgs) -> Result<(), CreateDkgError> { } for (i, addr) in args.operator_addresses.iter().enumerate() { - let checksum_addr = checksum_address(addr) + let checksum_addr = helpers::checksum_address(addr) .map_err(|source| CreateDkgError::InvalidOperatorAddress { index: i, source })?; operators.push(Operator { address: checksum_addr, @@ -351,12 +349,12 @@ async fn run_create_dkg(mut args: CreateDkgArgs) -> Result<(), CreateDkgError> { args.threshold }; - let fork_version_hex = network_to_fork_version(&args.network)?; + let fork_version_hex = network::network_to_fork_version(&args.network)?; let (priv_key, creator) = if args.publish { // Temporary creator address let key = SecretKey::random(&mut OsRng); - let addr = public_key_to_address(&key.public_key()); + let addr = helpers::public_key_to_address(&key.public_key()); ( Some(key), Creator { @@ -368,7 +366,7 @@ async fn run_create_dkg(mut args: CreateDkgArgs) -> Result<(), CreateDkgError> { (None, Creator::default()) }; - let deposit_amounts_gwei: Vec = eths_to_gweis(&args.deposit_amounts); + let deposit_amounts_gwei: Vec = deposit::eths_to_gweis(&args.deposit_amounts); let mut def = Definition::new( args.name.clone(), @@ -430,16 +428,17 @@ fn validate_dkg_config( return Err(CreateDkgError::TooFewOperators { num_operators }); } - if !valid_network(network) { + if !network::valid_network(network) { return Err(CreateDkgError::UnsupportedNetwork); } if !deposit_amounts.is_empty() { - let gweis = eths_to_gweis(deposit_amounts); - verify_deposit_amounts(&gweis, compounding)?; + let gweis = deposit::eths_to_gweis(deposit_amounts); + deposit::verify_deposit_amounts(&gweis, compounding)?; } - if !consensus_protocol.is_empty() && !is_supported_protocol_name(consensus_protocol) { + if !consensus_protocol.is_empty() && !protocols::is_supported_protocol_name(consensus_protocol) + { return Err(CreateDkgError::UnsupportedConsensusProtocol); } @@ -485,7 +484,7 @@ pub fn validate_withdrawal_addrs( network: &str, ) -> Result<(), WithdrawalValidationError> { for addr in addrs { - let checksum_addr = checksum_address(addr).map_err(|e| { + let checksum_addr = helpers::checksum_address(addr).map_err(|e| { WithdrawalValidationError::InvalidWithdrawalAddress { address: addr.clone(), reason: e.to_string(), diff --git a/crates/cli/src/commands/run.rs b/crates/cli/src/commands/run.rs index 500e71fa..30587b20 100644 --- a/crates/cli/src/commands/run.rs +++ b/crates/cli/src/commands/run.rs @@ -32,7 +32,7 @@ use std::{ time::Duration as StdDuration, }; -use pluto_eth2util::helpers::validate_http_headers; +use pluto_eth2util::helpers; use pluto_featureset::{Feature, FeaturesetError, Status}; use tokio_util::sync::CancellationToken; use tracing::{error, info, warn}; @@ -763,7 +763,7 @@ impl TryFrom for RunConfig { warn!("Jaeger flags are disabled and will be removed in a future release"); } - validate_http_headers(&general.beacon_node_headers) + helpers::validate_http_headers(&general.beacon_node_headers) .map_err(|err| CliError::Other(err.to_string()))?; let max_graffiti_bytes = if general.graffiti_disable_client_append { diff --git a/crates/cli/src/commands/test/helpers.rs b/crates/cli/src/commands/test/helpers.rs index d307d8c2..15f5fd5c 100644 --- a/crates/cli/src/commands/test/helpers.rs +++ b/crates/cli/src/commands/test/helpers.rs @@ -12,7 +12,6 @@ use crate::{ use k256::SecretKey; use pluto_app::obolapi::{Client, ClientOptions}; use pluto_eth2util::enr::Record; -use pluto_k1util::{load, sign}; use pluto_ssz::{HashRoot, HashWalker, Hasher}; use reqwest::{Method, StatusCode, header::CONTENT_TYPE}; use serde_with::{base64::Base64, serde_as}; @@ -291,7 +290,7 @@ pub(crate) async fn publish_result_to_obol_api( let enr = Record::new(&private_key, vec![])?; let sign_data_bytes = serde_json::to_vec(&data)?; let hash = hash_ssz(&sign_data_bytes)?; - let sig = sign(&private_key, &hash)?; + let sig = pluto_k1util::sign(&private_key, &hash)?; let result = ObolApiResult { enr: enr.to_string(), @@ -556,7 +555,7 @@ pub(crate) fn must_output_to_file_on_quiet( async fn load_or_generate_key(path: &Path) -> CliResult { if tokio::fs::try_exists(path).await? { - Ok(load(path)?) + Ok(pluto_k1util::load(path)?) } else { tracing::warn!( private_key_file = %path.display(), diff --git a/crates/cli/src/commands/test/peers.rs b/crates/cli/src/commands/test/peers.rs index 454e52af..e3f67394 100644 --- a/crates/cli/src/commands/test/peers.rs +++ b/crates/cli/src/commands/test/peers.rs @@ -17,17 +17,16 @@ use libp2p::{ }; use pluto_cluster::{definition::Definition, lock::Lock}; use pluto_eth2util::enr::Record; -use pluto_k1util::load as load_key; use pluto_p2p::{ behaviours::pluto::PlutoBehaviourEvent, - bootnode::new_relays, + bootnode, config::{DEFAULT_RELAYS, P2PConfig, RelayAddr}, gater::ConnGater, p2p::{Node, NodeType}, p2p_context::P2PContext, peer::{MutablePeer, Peer, peer_id_from_key, verify_p2p_key}, relay::RelayManager, - utils::is_relay_addr, + utils, }; use reqwest::Method; use sha2::{Digest, Sha256}; @@ -41,7 +40,7 @@ use super::{ write_result_to_writer, }; use crate::{ - commands::common::parse_relay_addrs, + commands::common, duration::Duration as CliDuration, error::{CliError, Result}, }; @@ -258,7 +257,7 @@ pub async fn run( tracing::debug!("enr_strings: {:?}", enr_strings); let cluster_peers = parse_peers(&enr_strings)?; - let private_key = load_key(&args.private_key_file)?; + let private_key = pluto_k1util::load(&args.private_key_file)?; verify_p2p_key(&cluster_peers, &private_key)?; @@ -280,7 +279,7 @@ pub async fn run( disable_reuse_port: args.p2p_disable_reuseport, }; - let relay_addrs = parse_relay_addrs(&args.p2p_relays)?; + let relay_addrs = common::parse_relay_addrs(&args.p2p_relays)?; let (node, relay_peers) = setup_p2p( timeout_ct.clone(), @@ -707,7 +706,7 @@ async fn run_peer_event_loop( // Once we have a relay circuit listen address our reservation is // active and other nodes can reach us. Trigger outbound dials. if let SwarmEvent::NewListenAddr { ref address, .. } = event - && is_relay_addr(address) + && utils::is_relay_addr(address) && !dialed_via_relay && !queued_tests.is_empty() { @@ -800,7 +799,7 @@ fn handle_swarm_event( } => { if let Some(state) = states.get_mut(&peer_id) { let addr = endpoint_addr(&endpoint); - let is_relay = is_relay_addr(addr); + let is_relay = utils::is_relay_addr(addr); if state.connect_time.is_none() { state.connect_time = Some(Instant::now()); @@ -855,7 +854,7 @@ fn handle_swarm_event( tracing::info!(timeout = ?direct_connection_timeout, target = %peer_target_name(peer, enr_str), "Trying to establish direct connection..."); } for addr in &info.listen_addrs { - if !is_relay_addr(addr) { + if !utils::is_relay_addr(addr) { let mut direct_addr = addr.clone(); direct_addr.push(Protocol::P2p(peer_id)); if let Err(e) = node.dial(direct_addr.clone()) { @@ -1063,7 +1062,7 @@ async fn setup_p2p( self_peer_id: PeerId, enr_hash: &str, ) -> Result<(Node, Vec)> { - let relay_peers = new_relays(cancel.clone(), relay_addrs, enr_hash).await?; + let relay_peers = bootnode::new_relays(cancel.clone(), relay_addrs, enr_hash).await?; let mut all_peer_ids: Vec = cluster_peers.iter().map(|p| p.id).collect(); all_peer_ids.push(self_peer_id); @@ -1401,7 +1400,7 @@ mod tests { // `--p2p-relays=""` parses to no relays, so there is nothing to probe. // Probing the raw flag strings instead used to key a target off the // empty string and report it as a failing relay. - let relays = parse_relay_addrs(&["".to_string()]).expect("relays"); + let relays = common::parse_relay_addrs(&["".to_string()]).expect("relays"); let queued = [TestCaseName::new("PingRelay", 1)]; let results = run_relay_http_tests(&relays, &queued, CancellationToken::new()).await; @@ -1411,7 +1410,7 @@ mod tests { #[tokio::test] async fn relay_http_tests_key_targets_by_address() { - let relays = parse_relay_addrs(&[ + let relays = common::parse_relay_addrs(&[ "http://127.0.0.1:1/enr".to_string(), "/ip4/127.0.0.1/tcp/3610/p2p/16Uiu2HAm7ULrTMdiEmQCJ2N9nsuGvfUDvfDGgHXJ4vNjrCwCzGDs" .to_string(), diff --git a/crates/cluster/src/definition.rs b/crates/cluster/src/definition.rs index 8abf21f5..29cf5bdb 100644 --- a/crates/cluster/src/definition.rs +++ b/crates/cluster/src/definition.rs @@ -3,13 +3,9 @@ use std::collections::HashSet; use pluto_ssz::{Hasher, serde_utils::HexBytes}; use crate::{ - eip712sigs::{ - EIP712Error, digest_eip712, eip712_creator_config_hash, eip712_enr, - get_operator_eip712_type, - }, - helpers::from_0x_hex_str, + eip712sigs::{self, EIP712Error}, operator::{Operator, OperatorV1X1, OperatorV1X2OrLater}, - ssz::{SSZ_MAX_VALIDATORS, SSZError, hash_definition}, + ssz::{self, SSZ_MAX_VALIDATORS, SSZError}, version::{CURRENT_VERSION, DKG_ALGO, versions::*}, }; use chrono::{DateTime, Timelike, Utc}; @@ -25,7 +21,7 @@ use serde_with::{ }; use uuid::Uuid; -use crate::helpers::{VerifySigError, verify_sig}; +use crate::helpers::{self, VerifySigError}; /// Length of the fork version in bytes. pub const FORK_VERSION_LEN: usize = 4; @@ -485,7 +481,7 @@ impl Definition { }) .collect(); - def.fork_version = from_0x_hex_str(&fork_version_hex, FORK_VERSION_LEN)?; + def.fork_version = helpers::from_0x_hex_str(&fork_version_hex, FORK_VERSION_LEN)?; for opt in opts { opt(&mut def); @@ -555,8 +551,8 @@ impl Definition { }; } - let operator_config_hash_digest = digest_eip712( - &get_operator_eip712_type(self.version.as_str())?, + let operator_config_hash_digest = eip712sigs::digest_eip712( + &eip712sigs::get_operator_eip712_type(self.version.as_str())?, self, &Operator::default(), )?; @@ -587,7 +583,7 @@ impl Definition { } // Check that we have a valid config signature for each operator. - let is_valid_operator_config_sig = verify_sig( + let is_valid_operator_config_sig = helpers::verify_sig( operator.address.as_str(), operator_config_hash_digest.as_slice(), operator.config_signature.as_slice(), @@ -608,9 +604,9 @@ impl Definition { } // Check that we have a valid enr signature for each operator. - let enr_digest = digest_eip712(&eip712_enr(), self, operator)?; + let enr_digest = eip712sigs::digest_eip712(&eip712sigs::eip712_enr(), self, operator)?; - let is_valid_operator_enr_sig = verify_sig( + let is_valid_operator_enr_sig = helpers::verify_sig( operator.address.as_str(), enr_digest.as_slice(), operator.enr_signature.as_slice(), @@ -650,10 +646,13 @@ impl Definition { return Err(DefinitionError::EmptyCreatorConfigSignature); } - let creator_config_hash_digest = - digest_eip712(&eip712_creator_config_hash(), self, &Operator::default())?; + let creator_config_hash_digest = eip712sigs::digest_eip712( + &eip712sigs::eip712_creator_config_hash(), + self, + &Operator::default(), + )?; - let is_valid_creator_sig = verify_sig( + let is_valid_creator_sig = helpers::verify_sig( self.creator.address.as_str(), creator_config_hash_digest.as_slice(), self.creator.config_signature.as_slice(), @@ -734,12 +733,12 @@ impl Definition { /// Sets the definition hashes. pub fn set_definition_hashes(&mut self) -> Result<(), DefinitionError> { let config_hash = - hash_definition(self, true).map_err(|e| DefinitionError::SSZError(Box::new(e)))?; + ssz::hash_definition(self, true).map_err(|e| DefinitionError::SSZError(Box::new(e)))?; self.config_hash = config_hash.to_vec(); - let definition_hash = - hash_definition(self, false).map_err(|e| DefinitionError::SSZError(Box::new(e)))?; + let definition_hash = ssz::hash_definition(self, false) + .map_err(|e| DefinitionError::SSZError(Box::new(e)))?; self.definition_hash = definition_hash.to_vec(); @@ -750,7 +749,7 @@ impl Definition { /// doesn't matches actual hashes. pub fn verify_hashes(&self) -> Result<(), DefinitionError> { let config_hash = - hash_definition(self, true).map_err(|e| DefinitionError::SSZError(Box::new(e)))?; + ssz::hash_definition(self, true).map_err(|e| DefinitionError::SSZError(Box::new(e)))?; if config_hash != self.config_hash.as_slice() { return Err(DefinitionError::InvalidConfigHash { @@ -759,8 +758,8 @@ impl Definition { }); } - let definition_hash = - hash_definition(self, false).map_err(|e| DefinitionError::SSZError(Box::new(e)))?; + let definition_hash = ssz::hash_definition(self, false) + .map_err(|e| DefinitionError::SSZError(Box::new(e)))?; if definition_hash != self.definition_hash.as_slice() { return Err(DefinitionError::InvalidDefinitionHash { diff --git a/crates/cluster/src/eip712sigs.rs b/crates/cluster/src/eip712sigs.rs index 916d43fe..107e68fd 100644 --- a/crates/cluster/src/eip712sigs.rs +++ b/crates/cluster/src/eip712sigs.rs @@ -2,10 +2,9 @@ use crate::{definition::Definition, operator::Operator, version::V1_3}; use k256::SecretKey; use pluto_eth2util::{ eip712::{ - Domain, Field, PRIMITIVE_STRING, PRIMITIVE_UINT256, Primitive, Type, TypedData, Value, - hash_typed_data, + self, Domain, Field, PRIMITIVE_STRING, PRIMITIVE_UINT256, Primitive, Type, TypedData, Value, }, - network::fork_version_to_chain_id, + network, }; type ValueFunc = Box Value>; @@ -139,7 +138,7 @@ pub(crate) fn digest_eip712( definition: &Definition, operator: &Operator, ) -> Result> { - let chain_id = fork_version_to_chain_id(definition.fork_version.as_ref())?; + let chain_id = network::fork_version_to_chain_id(definition.fork_version.as_ref())?; let fields = typ .fields @@ -163,7 +162,7 @@ pub(crate) fn digest_eip712( }, }; - let digest = hash_typed_data(&data).map_err(EIP712Error::FailedToHashTypedData)?; + let digest = eip712::hash_typed_data(&data).map_err(EIP712Error::FailedToHashTypedData)?; Ok(digest) } diff --git a/crates/cluster/src/helpers.rs b/crates/cluster/src/helpers.rs index 8458b227..a0c9b13a 100644 --- a/crates/cluster/src/helpers.rs +++ b/crates/cluster/src/helpers.rs @@ -1,6 +1,6 @@ use chrono::{DateTime, Utc}; use pluto_crypto::tbls; -use pluto_eth2util::helpers::{checksum_address, public_key_to_address}; +use pluto_eth2util::helpers; use pluto_k1util::K1UtilError; use serde::{Deserialize, Deserializer, Serializer}; use serde_with::{DeserializeAs, SerializeAs}; @@ -31,9 +31,9 @@ pub fn verify_sig( digest: &[u8], sig: &[u8], ) -> std::result::Result { - let expected_addr = checksum_address(expected_addr)?; + let expected_addr = helpers::checksum_address(expected_addr)?; let recovered = pluto_k1util::recover(digest, sig)?; - let actual_addr = public_key_to_address(&recovered); + let actual_addr = helpers::public_key_to_address(&recovered); Ok(expected_addr == actual_addr) } diff --git a/crates/cluster/src/version.rs b/crates/cluster/src/version.rs index 6e16617a..0bbb19ab 100644 --- a/crates/cluster/src/version.rs +++ b/crates/cluster/src/version.rs @@ -43,14 +43,20 @@ pub const SUPPORTED_VERSIONS: [&str; 11] = [ /// Returns true if the provided cluster definition version supports /// pre-generated builder registrations. #[must_use] -pub fn support_pregen_registrations(version: &str) -> bool { - !matches!(version, V1_0 | V1_1 | V1_2 | V1_3 | V1_4 | V1_5 | V1_6) +pub fn support_pregen_registrations(version: impl AsRef) -> bool { + !matches!( + version.as_ref(), + V1_0 | V1_1 | V1_2 | V1_3 | V1_4 | V1_5 | V1_6 + ) } /// Returns true if the provided cluster lock version supports node signatures. #[must_use] -pub fn support_node_signatures(version: &str) -> bool { - !matches!(version, V1_0 | V1_1 | V1_2 | V1_3 | V1_4 | V1_5 | V1_6) +pub fn support_node_signatures(version: impl AsRef) -> bool { + !matches!( + version.as_ref(), + V1_0 | V1_1 | V1_2 | V1_3 | V1_4 | V1_5 | V1_6 + ) } #[cfg(test)] diff --git a/crates/core/src/bcast/mod.rs b/crates/core/src/bcast/mod.rs index 99a3f53e..3e3651e2 100644 --- a/crates/core/src/bcast/mod.rs +++ b/crates/core/src/bcast/mod.rs @@ -18,7 +18,6 @@ use tree_hash::TreeHash; pub use recast::Recaster; use crate::{ - bcast::metrics::instrument_duty, signeddata::{ SignedSyncContributionAndProof, SignedSyncMessage, SignedVoluntaryExit, VersionedAttestation, VersionedSignedAggregateAndProof, VersionedSignedProposal, @@ -303,7 +302,7 @@ impl Broadcaster { tracing::warn!(%error, %duty, "Failed to compute broadcast delay"); }) .ok(); - instrument_duty(&duty, delay); + metrics::instrument_duty(&duty, delay); Ok(()) } diff --git a/crates/core/src/signeddata.rs b/crates/core/src/signeddata.rs index 083dc49b..c7eac58c 100644 --- a/crates/core/src/signeddata.rs +++ b/crates/core/src/signeddata.rs @@ -4,7 +4,7 @@ use alloy::primitives::U256; use serde::{Deserialize, Deserializer, Serialize, Serializer}; use tree_hash::TreeHash; -use pluto_crypto::types::sig_to_eth2; +use pluto_crypto::types; use pluto_eth2api::{ ConsensusVersion, ProduceBlockV3ResponseResponse, spec::{ @@ -236,7 +236,7 @@ impl SignedData for VersionedSignedProposal { if proposal.version == versioned::DataVersion::Unknown { return Err(SignedDataError::UnknownVersion); } - let eth2_sig = sig_to_eth2(signature); + let eth2_sig = types::sig_to_eth2(signature); proposal.block.set_signature(eth2_sig); Ok(out) @@ -383,7 +383,7 @@ impl SignedData for Attestation { fn set_signature(&self, signature: Signature) -> Result { let mut out = self.clone(); - out.0.signature = sig_to_eth2(signature); + out.0.signature = types::sig_to_eth2(signature); Ok(out) } @@ -467,7 +467,7 @@ impl SignedData for VersionedAttestation { .attestation .as_mut() .ok_or(SignedDataError::MissingAttestation(version))? - .set_signature(sig_to_eth2(signature)); + .set_signature(types::sig_to_eth2(signature)); Ok(out) } @@ -590,7 +590,7 @@ impl SignedData for SignedVoluntaryExit { fn set_signature(&self, signature: Signature) -> Result { let mut out = self.clone(); - out.0.signature = sig_to_eth2(signature); + out.0.signature = types::sig_to_eth2(signature); Ok(out) } @@ -673,7 +673,7 @@ impl SignedData for VersionedSignedValidatorRegistration { let Some(v1) = out.0.v1.as_mut() else { return Err(SignedDataError::MissingV1Registration); }; - v1.signature = sig_to_eth2(signature); + v1.signature = types::sig_to_eth2(signature); } versioned::BuilderVersion::Unknown => { return Err(SignedDataError::UnknownVersion); @@ -762,7 +762,7 @@ impl SignedData for SignedRandao { fn set_signature(&self, signature: Signature) -> Result { let mut out = self.clone(); - out.0.signature = sig_to_eth2(signature); + out.0.signature = types::sig_to_eth2(signature); Ok(out) } @@ -812,7 +812,7 @@ impl SignedData for BeaconCommitteeSelection { fn set_signature(&self, signature: Signature) -> Result { let mut out = self.clone(); - out.0.selection_proof = sig_to_eth2(signature); + out.0.selection_proof = types::sig_to_eth2(signature); Ok(out) } @@ -855,7 +855,7 @@ impl SignedData for SyncCommitteeSelection { fn set_signature(&self, signature: Signature) -> Result { let mut out = self.clone(); - out.0.selection_proof = sig_to_eth2(signature); + out.0.selection_proof = types::sig_to_eth2(signature); Ok(out) } @@ -898,7 +898,7 @@ impl SignedData for SignedAggregateAndProof { fn set_signature(&self, signature: Signature) -> Result { let mut out = self.clone(); - out.0.signature = sig_to_eth2(signature); + out.0.signature = types::sig_to_eth2(signature); Ok(out) } @@ -984,7 +984,7 @@ impl SignedData for VersionedSignedAggregateAndProof { } out.0 .aggregate_and_proof - .set_signature(sig_to_eth2(signature)); + .set_signature(types::sig_to_eth2(signature)); Ok(out) } @@ -1078,7 +1078,7 @@ impl SignedData for SignedSyncMessage { fn set_signature(&self, signature: Signature) -> Result { let mut out = self.clone(); - out.0.signature = sig_to_eth2(signature); + out.0.signature = types::sig_to_eth2(signature); Ok(out) } @@ -1121,7 +1121,7 @@ impl SignedData for SyncContributionAndProof { fn set_signature(&self, signature: Signature) -> Result { let mut out = self.clone(); - out.0.selection_proof = sig_to_eth2(signature); + out.0.selection_proof = types::sig_to_eth2(signature); Ok(out) } @@ -1164,7 +1164,7 @@ impl SignedData for SignedSyncContributionAndProof { fn set_signature(&self, signature: Signature) -> Result { let mut out = self.clone(); - out.0.signature = sig_to_eth2(signature); + out.0.signature = types::sig_to_eth2(signature); Ok(out) } diff --git a/crates/core/src/tracker/inclusion.rs b/crates/core/src/tracker/inclusion.rs index 7bb94c33..83a240bc 100644 --- a/crates/core/src/tracker/inclusion.rs +++ b/crates/core/src/tracker/inclusion.rs @@ -37,7 +37,7 @@ use crate::{ Attestation, SignedAggregateAndProof, SignedDataError, VersionedAttestation, VersionedSignedAggregateAndProof, VersionedSignedProposal, }, - tracker::{StepError, analysis::incl_supported, metrics::TRACKER_METRICS}, + tracker::{StepError, analysis, metrics::TRACKER_METRICS}, types::{Duty, DutyType, PubKey, SignedData, SignedDataSet}, }; @@ -221,7 +221,7 @@ impl InclusionCore { data: Box, delay: Duration, ) -> Result<(), InclusionError> { - if !incl_supported(&self.feature_set).contains(&duty.duty_type) { + if !analysis::incl_supported(&self.feature_set).contains(&duty.duty_type) { return Ok(()); } diff --git a/crates/core/src/unsigneddata.rs b/crates/core/src/unsigneddata.rs index 141999ce..2adeca6c 100644 --- a/crates/core/src/unsigneddata.rs +++ b/crates/core/src/unsigneddata.rs @@ -3,14 +3,14 @@ use std::collections::HashMap; use pluto_eth2api::spec::phase0; -use pluto_ssz::decode::{decode_u32, decode_u64}; +use pluto_ssz::decode; use serde::{Deserialize, Deserializer, de}; use ssz::{Decode, Encode}; use crate::{ ParSigExCodecError, corepb::v1::core as pbcore, - parsigex_codec::looks_like_json, + parsigex_codec, signeddata::{ AttestationData, AttesterDuty, SyncContribution, VersionedAggregatedAttestation, VersionedProposal, @@ -162,7 +162,7 @@ fn decode_versioned_proposal(data: &[u8]) -> Result(data) { return Ok(VersionedAggregatedAttestation(decoded.0)); @@ -222,7 +222,7 @@ fn decode_sync_contribution(data: &[u8]) -> Result Result Result Result path::PathBuf { - let mut lock_path = key_path(data_dir); + let mut lock_path = k1::key_path(data_dir); let file_name = lock_path .file_name() .and_then(OsStr::to_str) diff --git a/crates/dkg/src/exchanger.rs b/crates/dkg/src/exchanger.rs index 8f62075b..04ffa41e 100644 --- a/crates/dkg/src/exchanger.rs +++ b/crates/dkg/src/exchanger.rs @@ -55,7 +55,7 @@ use pluto_core::{ }, types::{Duty, DutyType, ParSignedData, ParSignedDataSet, PubKey, SlotNumber}, }; -use pluto_parsigex::{Handle, ReceivedSub, received_subscriber}; +use pluto_parsigex::{Handle, ReceivedSub}; /// Numeric identifier for a DKG signature exchange round, encoded as /// `Duty.slot`. @@ -215,7 +215,7 @@ impl Exchanger { { let sigdb_clone = Arc::clone(&sigdb); - let sub: ReceivedSub = received_subscriber(move |duty, set| { + let sub: ReceivedSub = pluto_parsigex::received_subscriber(move |duty, set| { let sigdb = sigdb_clone.clone(); async move { if let Err(e) = sigdb.lock().await.store_external(&duty, &set).await { diff --git a/crates/dkg/src/frost.rs b/crates/dkg/src/frost.rs index 3598ad5e..2c6c8c5a 100644 --- a/crates/dkg/src/frost.rs +++ b/crates/dkg/src/frost.rs @@ -1,7 +1,7 @@ use std::collections::{BTreeMap, HashMap}; use async_trait::async_trait; -use pluto_crypto::types::{PublicKey, privkey_from_bytes, pubkey_from_bytes}; +use pluto_crypto::types::{self, PublicKey}; use pluto_frost::{ G1Affine, G1Projective, KeyPackage, kryptology::{self, Round1Bcast, Round1Secret, Round2Bcast, ShamirShare}, @@ -437,7 +437,7 @@ fn make_shares( shares.push(Share { pub_key: point_to_pubkey(G1Affine::from(pub_key).to_compressed())?, - secret_share: privkey_from_bytes(&kryptology::scalar_to_be(&secret_share))?, + secret_share: types::privkey_from_bytes(&kryptology::scalar_to_be(&secret_share))?, public_shares: pub_shares.get(&v_idx).cloned().unwrap_or_default(), }); } @@ -449,7 +449,7 @@ fn point_to_pubkey(point: [u8; 48]) -> Result { // `pubkey_from_bytes` only checks length; transport bytes still need G1 // validation. G1Projective::from_compressed(&point).ok_or(FrostError::InvalidPublicKeyPoint)?; - Ok(pubkey_from_bytes(&point)?) + Ok(types::pubkey_from_bytes(&point)?) } fn validate_participant_inputs( diff --git a/crates/dkg/src/frostp2p/behaviour.rs b/crates/dkg/src/frostp2p/behaviour.rs index 54c44ad1..c901ff76 100644 --- a/crates/dkg/src/frostp2p/behaviour.rs +++ b/crates/dkg/src/frostp2p/behaviour.rs @@ -24,7 +24,7 @@ use tracing::{debug, warn}; use super::{ FrostP2PError, FrostP2PEvent, handler::{FrostP2PHandler, InEvent, OutEvent}, - transport::validate_round1_p2p, + transport, }; use crate::{dkgpb::v1::frost::FrostRound1P2p, frost::FrostError}; @@ -383,7 +383,7 @@ impl NetworkBehaviour for FrostP2PBehaviour { }; match event { OutEvent::Received(msg) => { - if let Err(error) = validate_round1_p2p( + if let Err(error) = transport::validate_round1_p2p( peer_id, &self.share_idx_by_peer, self.local_share_idx, diff --git a/crates/dkg/src/signing.rs b/crates/dkg/src/signing.rs index ab489a94..b7a9c621 100644 --- a/crates/dkg/src/signing.rs +++ b/crates/dkg/src/signing.rs @@ -15,17 +15,17 @@ use pluto_core::{ signeddata::VersionedSignedValidatorRegistration, types::{ParSignedData, ParSignedDataSet, PubKey}, }; -use pluto_crypto::{tbls, types::pubkey_to_eth2}; +use pluto_crypto::{tbls, types}; use pluto_eth2api::{spec::phase0, v1, versioned}; use pluto_eth2util::{deposit, network, registration}; use tracing::{info, warn}; use crate::{ - aggregate::{agg_deposit_data, agg_lock_hash_sig, agg_validator_registrations}, + aggregate, dkg::AppendConfig, exchanger::{Exchanger, SIG_DEPOSIT_DATA, SIG_LOCK, SIG_VALIDATOR_REG}, share::Share, - validators::create_dist_validators, + validators, }; /// Result type for DKG signing helpers. @@ -138,7 +138,7 @@ pub fn sign_deposit_msgs( let mut set = ParSignedDataSet::new(); for (share, withdrawal_address) in shares.iter().zip(withdrawal_addresses.iter()) { - let eth2_pubkey = pubkey_to_eth2(share.pub_key); + let eth2_pubkey = types::pubkey_to_eth2(share.pub_key); let pub_key = share_pubkey(share, "signing deposit message")?; let withdrawal_address = pluto_eth2util::helpers::checksum_address(withdrawal_address)?; @@ -177,7 +177,7 @@ pub fn sign_validator_registrations( let mut set = ParSignedDataSet::new(); for (share, fee_recipient) in shares.iter().zip(fee_recipients.iter()) { - let eth2_pubkey = pubkey_to_eth2(share.pub_key); + let eth2_pubkey = types::pubkey_to_eth2(share.pub_key); let pub_key = share_pubkey(share, "signing validator registration")?; let reg_msg = registration::new_message( @@ -233,7 +233,7 @@ pub(crate) async fn sign_and_agg_deposit_data( .checked_add(u64::try_from(i)?) .ok_or(SigningError::Overflow)?; let peer_sigs = exchanger.exchange(sig_type, set).await?; - let deposit_data = agg_deposit_data(&peer_sigs, shares, &msgs, network)?; + let deposit_data = aggregate::agg_deposit_data(&peer_sigs, shares, &msgs, network)?; result.push(deposit_data); } @@ -270,7 +270,7 @@ pub(crate) async fn sign_and_agg_validator_registrations( )?; let peer_sigs = exchanger.exchange(SIG_VALIDATOR_REG, set).await?; - Ok(agg_validator_registrations( + Ok(aggregate::agg_validator_registrations( &peer_sigs, shares, &msgs, @@ -294,7 +294,7 @@ pub(crate) async fn sign_and_aggregate_lock_hash( val_regs: Vec, append_config: Option<&AppendConfig>, ) -> Result { - let mut validators = create_dist_validators(new_shares, &deposit_datas, &val_regs)?; + let mut validators = validators::create_dist_validators(new_shares, &deposit_datas, &val_regs)?; if let Some(append) = append_config { let mut merged = append.cluster_lock.distributed_validators.clone(); @@ -362,7 +362,8 @@ pub(crate) async fn sign_and_aggregate_lock_hash( }) .collect::>()?; - let (agg_sig, all_pubshares) = agg_lock_hash_sig(&peer_sigs, &shares_map, &lock.lock_hash)?; + let (agg_sig, all_pubshares) = + aggregate::agg_lock_hash_sig(&peer_sigs, &shares_map, &lock.lock_hash)?; tbls::verify_aggregate(&all_pubshares, agg_sig, &lock.lock_hash)?; lock.signature_aggregate = agg_sig.to_vec(); diff --git a/crates/eth2api/src/spec/serde_utils.rs b/crates/eth2api/src/spec/serde_utils.rs index 781f4d14..fedb8ca7 100644 --- a/crates/eth2api/src/spec/serde_utils.rs +++ b/crates/eth2api/src/spec/serde_utils.rs @@ -1,6 +1,6 @@ //! Shared serde helpers for consensus-spec JSON encoding. -use pluto_ssz::serde_utils::trim_0x_prefix; +use pluto_ssz::serde_utils; /// Error raised while converting a loosely-typed beacon-API value (whose /// numeric and byte fields are carried as decimal / `0x`-hex strings) into a @@ -30,7 +30,8 @@ pub(crate) fn parse_u64(value: &str, field: &'static str) -> Result Result, ConversionError> { - hex::decode(trim_0x_prefix(value)).map_err(|_| ConversionError::DecodeHex { field }) + hex::decode(serde_utils::trim_0x_prefix(value)) + .map_err(|_| ConversionError::DecodeHex { field }) } /// Decodes a `0x`-prefixed (or bare) hex string into a fixed-size byte array, diff --git a/crates/eth2util/src/network.rs b/crates/eth2util/src/network.rs index a28ae7d0..ee975b3a 100644 --- a/crates/eth2util/src/network.rs +++ b/crates/eth2util/src/network.rs @@ -237,12 +237,12 @@ pub fn network_to_fork_version_bytes(network: impl AsRef) -> Result } /// Valid network. -pub fn valid_network(name: &str) -> bool { +pub fn valid_network(name: impl AsRef) -> bool { network_from_name(name).is_ok() } /// Network to genesis time. -pub fn network_to_genesis_time(name: &str) -> Result> { +pub fn network_to_genesis_time(name: impl AsRef) -> Result> { let network = network_from_name(name)?; DateTime::::from_timestamp( i64::try_from(network.genesis_timestamp).map_err(|_| { diff --git a/crates/k1util/src/k1util.rs b/crates/k1util/src/k1util.rs index 4ef26b33..7cafa2a9 100644 --- a/crates/k1util/src/k1util.rs +++ b/crates/k1util/src/k1util.rs @@ -212,8 +212,8 @@ pub fn recover(hash: &[u8], sig: &[u8]) -> Result { } /// Load loads a secret key from a file. -pub fn load(file: &Path) -> Result { - let contents = std::fs::read_to_string(file).map_err(K1UtilError::FailedToReadFile)?; +pub fn load(file: impl AsRef) -> Result { + let contents = std::fs::read_to_string(file.as_ref()).map_err(K1UtilError::FailedToReadFile)?; let decoded = hex::decode(contents.trim())?; @@ -228,7 +228,8 @@ pub fn load(file: &Path) -> Result { /// matching Charon's `app/k1util/k1util.go` `Save` which writes via /// `os.WriteFile(file, ..., 0o600)`. This prevents the private key from being /// world-readable. -pub fn save(key: &SecretKey, file: &Path) -> Result<()> { +pub fn save(key: &SecretKey, file: impl AsRef) -> Result<()> { + let file = file.as_ref(); let encoded = hex::encode(key.to_bytes()); #[cfg(unix)] diff --git a/crates/p2p/src/conn_logger.rs b/crates/p2p/src/conn_logger.rs index 088467d3..67a48493 100644 --- a/crates/p2p/src/conn_logger.rs +++ b/crates/p2p/src/conn_logger.rs @@ -8,7 +8,7 @@ use tracing::{debug, instrument}; use crate::{ metrics::{ConnectionType, P2P_METRICS, PeerConnectionLabels, Protocol, RelayConnectionLabels}, - name::peer_name, + name, p2p_context::{P2PContext, Peer}, utils, }; @@ -59,12 +59,12 @@ impl ConnectionLoggerMetrics for DefaultConnectionLoggerMetrics { } fn inc_peer_connection_total(&self, peer: &PeerId) { - P2P_METRICS.peer_connection_total[&peer_name(peer)].inc(); + P2P_METRICS.peer_connection_total[&name::peer_name(peer)].inc(); } fn set_peer_connection_type(&self, peer: &PeerId, addr: &Multiaddr, count: u64) { P2P_METRICS.peer_connection_types[&PeerConnectionLabels::new( - &peer_name(peer), + &name::peer_name(peer), utils::addr_type(addr), utils::addr_protocol(addr), )] @@ -73,7 +73,7 @@ impl ConnectionLoggerMetrics for DefaultConnectionLoggerMetrics { fn set_relay_connection_type(&self, peer: &PeerId, addr: &Multiaddr, count: u64) { P2P_METRICS.relay_connection_types[&RelayConnectionLabels::new( - &peer_name(peer), + &name::peer_name(peer), utils::addr_type(addr), utils::addr_protocol(addr), )] @@ -171,7 +171,7 @@ impl NetworkBehaviour for ConnectionLogger type ConnectionHandler = dummy::ConnectionHandler; type ToSwarm = (); - #[instrument(skip(self, _local_addr, remote_addr), fields(peer = %peer_name(&peer), addr = %remote_addr))] + #[instrument(skip(self, _local_addr, remote_addr), fields(peer = %name::peer_name(&peer), addr = %remote_addr))] fn handle_established_inbound_connection( &mut self, _connection_id: ConnectionId, @@ -180,7 +180,7 @@ impl NetworkBehaviour for ConnectionLogger remote_addr: &Multiaddr, ) -> Result, ConnectionDenied> { debug!( - peer = %peer_name(&peer), + peer = %name::peer_name(&peer), addr = %remote_addr, conn_type = ?utils::addr_type(remote_addr), protocol = ?utils::addr_protocol(remote_addr), @@ -190,7 +190,7 @@ impl NetworkBehaviour for ConnectionLogger Ok(dummy::ConnectionHandler) } - #[instrument(skip(self, _role_override, _port_use), fields(peer = %peer_name(&peer), addr = %addr))] + #[instrument(skip(self, _role_override, _port_use), fields(peer = %name::peer_name(&peer), addr = %addr))] fn handle_established_outbound_connection( &mut self, _connection_id: ConnectionId, @@ -200,7 +200,7 @@ impl NetworkBehaviour for ConnectionLogger _port_use: libp2p::core::transport::PortUse, ) -> Result, ConnectionDenied> { debug!( - peer = %peer_name(&peer), + peer = %name::peer_name(&peer), addr = %addr, conn_type = ?utils::addr_type(addr), protocol = ?utils::addr_protocol(addr), @@ -214,7 +214,7 @@ impl NetworkBehaviour for ConnectionLogger match event { libp2p::swarm::FromSwarm::ConnectionEstablished(event) => { debug!( - peer = %peer_name(&event.peer_id), + peer = %name::peer_name(&event.peer_id), endpoint = ?event.endpoint, other_established = event.other_established, "connection established" @@ -236,7 +236,7 @@ impl NetworkBehaviour for ConnectionLogger } libp2p::swarm::FromSwarm::ConnectionClosed(event) => { debug!( - peer = %peer_name(&event.peer_id), + peer = %name::peer_name(&event.peer_id), endpoint = ?event.endpoint, num_established = event.remaining_established, "connection closed" @@ -327,12 +327,12 @@ mod tests { } fn inc_peer_connection_total(&self, peer: &PeerId) { - self.metrics.peer_connection_total[&peer_name(peer)].inc(); + self.metrics.peer_connection_total[&name::peer_name(peer)].inc(); } fn set_peer_connection_type(&self, peer: &PeerId, addr: &Multiaddr, count: u64) { self.metrics.peer_connection_types[&PeerConnectionLabels::new( - &peer_name(peer), + &name::peer_name(peer), utils::addr_type(addr), utils::addr_protocol(addr), )] @@ -341,7 +341,7 @@ mod tests { fn set_relay_connection_type(&self, peer: &PeerId, addr: &Multiaddr, count: u64) { self.metrics.relay_connection_types[&RelayConnectionLabels::new( - &peer_name(peer), + &name::peer_name(peer), utils::addr_type(addr), utils::addr_protocol(addr), )] @@ -519,12 +519,16 @@ mod tests { behaviour.increment_connection(peer, &addr); // Check peer_connection_total was incremented - let total = behaviour.metrics().inner().peer_connection_total[&peer_name(&peer)].get(); + let total = + behaviour.metrics().inner().peer_connection_total[&name::peer_name(&peer)].get(); assert_eq!(total, 1); // Check peer_connection_types was set - let labels = - PeerConnectionLabels::new(&peer_name(&peer), ConnectionType::Direct, Protocol::Tcp); + let labels = PeerConnectionLabels::new( + &name::peer_name(&peer), + ConnectionType::Direct, + Protocol::Tcp, + ); let count = behaviour.metrics().inner().peer_connection_types[&labels].get(); assert_eq!(count, 1); } @@ -540,7 +544,7 @@ mod tests { // Check relay_connection_types was set (not peer_connection_types) let labels = RelayConnectionLabels::new( - &peer_name(&unknown_peer), + &name::peer_name(&unknown_peer), ConnectionType::Direct, Protocol::Tcp, ); @@ -548,8 +552,9 @@ mod tests { assert_eq!(count, 1); // peer_connection_total should not have been incremented for unknown peer - let total = - behaviour.metrics().inner().peer_connection_total[&peer_name(&unknown_peer)].get(); + let total = behaviour.metrics().inner().peer_connection_total + [&name::peer_name(&unknown_peer)] + .get(); assert_eq!(total, 0); } @@ -563,8 +568,11 @@ mod tests { behaviour.decrement_connection(peer, &addr); // After decrementing to zero, the gauge should be set to 0 - let labels = - PeerConnectionLabels::new(&peer_name(&peer), ConnectionType::Direct, Protocol::Tcp); + let labels = PeerConnectionLabels::new( + &name::peer_name(&peer), + ConnectionType::Direct, + Protocol::Tcp, + ); let count = behaviour.metrics().inner().peer_connection_types[&labels].get(); assert_eq!(count, 0); } @@ -581,7 +589,7 @@ mod tests { // After decrementing, relay_connection_types should be 0 let labels = RelayConnectionLabels::new( - &peer_name(&relay_peer), + &name::peer_name(&relay_peer), ConnectionType::Direct, Protocol::Tcp, ); @@ -647,7 +655,8 @@ mod tests { assert_eq!(behaviour.counts().get(&key), Some(&1)); // Verify metrics were updated - let total = behaviour.metrics().inner().peer_connection_total[&peer_name(&peer)].get(); + let total = + behaviour.metrics().inner().peer_connection_total[&name::peer_name(&peer)].get(); assert_eq!(total, 1); } @@ -678,7 +687,8 @@ mod tests { assert_eq!(behaviour.counts().get(&key), Some(&1)); // Verify metrics were updated - let total = behaviour.metrics().inner().peer_connection_total[&peer_name(&peer)].get(); + let total = + behaviour.metrics().inner().peer_connection_total[&name::peer_name(&peer)].get(); assert_eq!(total, 1); } @@ -916,17 +926,17 @@ mod tests { // Known peer should have peer_connection_total incremented let known_total = - behaviour.metrics().inner().peer_connection_total[&peer_name(&known_peer)].get(); + behaviour.metrics().inner().peer_connection_total[&name::peer_name(&known_peer)].get(); assert_eq!(known_total, 1); // Relay peer should NOT have peer_connection_total incremented let relay_total = - behaviour.metrics().inner().peer_connection_total[&peer_name(&relay_peer)].get(); + behaviour.metrics().inner().peer_connection_total[&name::peer_name(&relay_peer)].get(); assert_eq!(relay_total, 0); // But relay_connection_types should be set let relay_labels = RelayConnectionLabels::new( - &peer_name(&relay_peer), + &name::peer_name(&relay_peer), ConnectionType::Direct, Protocol::Tcp, ); diff --git a/crates/p2p/src/force_direct.rs b/crates/p2p/src/force_direct.rs index b8c9f49c..d905c098 100644 --- a/crates/p2p/src/force_direct.rs +++ b/crates/p2p/src/force_direct.rs @@ -19,7 +19,7 @@ use std::time::Duration; use tokio::time::Interval; use tracing::{debug, warn}; -use crate::{name::peer_name, p2p_context::P2PContext, utils}; +use crate::{name, p2p_context::P2PContext, utils}; const FORCE_DIRECT_INTERVAL: Duration = Duration::from_secs(60); @@ -123,7 +123,7 @@ impl ForceDirectBehaviour { if connections.is_empty() { warn!( - peer = %peer_name(peer), + peer = %name::peer_name(peer), "no connections to peer" ); continue; @@ -134,7 +134,7 @@ impl ForceDirectBehaviour { .any(|c| !utils::is_relay_addr(&c.remote_addr)) { debug!( - peer = %peer_name(peer), + peer = %name::peer_name(peer), "not all connections to peer are relay connections, skipping force direct" ); continue; @@ -142,7 +142,7 @@ impl ForceDirectBehaviour { let Some(addresses) = available_addresses else { warn!( - peer = %peer_name(peer), + peer = %name::peer_name(peer), "no known addresses for peer" ); continue; @@ -157,14 +157,14 @@ impl ForceDirectBehaviour { if direct_addresses.is_empty() { warn!( - peer = %peer_name(peer), + peer = %name::peer_name(peer), "no direct addresses for peer, cannot force direct connection" ); continue; } debug!( - peer = %peer_name(peer), + peer = %name::peer_name(peer), direct_addresses = ?direct_addresses, "forcing direct connection to peer using {} available addresses", direct_addresses.len() diff --git a/crates/p2p/src/k1.rs b/crates/p2p/src/k1.rs index 9665b9f3..aa12b2ce 100644 --- a/crates/p2p/src/k1.rs +++ b/crates/p2p/src/k1.rs @@ -29,7 +29,7 @@ pub fn key_path(data_dir: &Path) -> PathBuf { /// Loads the private key from the data dir. pub fn load_priv_key(data_dir: &Path) -> Result { - pluto_k1util::load(&key_path(data_dir)).map_err(K1Error::K1UtilError) + pluto_k1util::load(key_path(data_dir)).map_err(K1Error::K1UtilError) } /// Generates a new private key and saves it to the data dir. @@ -40,7 +40,7 @@ pub fn new_saved_priv_key(data_dir: &Path) -> Result { let key = SecretKey::random(&mut OsRng); - pluto_k1util::save(&key, &key_path(data_dir)).map_err(K1Error::K1UtilError)?; + pluto_k1util::save(&key, key_path(data_dir)).map_err(K1Error::K1UtilError)?; Ok(key) } @@ -89,7 +89,7 @@ mod tests { fn create_test_key_file(data_dir: &Path) -> Result { let key = SecretKey::random(&mut OsRng); - pluto_k1util::save(&key, &key_path(data_dir))?; + pluto_k1util::save(&key, key_path(data_dir))?; Ok(key) } diff --git a/crates/p2p/src/p2p.rs b/crates/p2p/src/p2p.rs index ceb4b580..afbad4a6 100644 --- a/crates/p2p/src/p2p.rs +++ b/crates/p2p/src/p2p.rs @@ -108,7 +108,7 @@ use crate::{ behaviours::pluto::{PlutoBehaviour, PlutoBehaviourBuilder, PlutoBehaviourEvent}, config::{P2PConfig, P2PConfigError}, metrics::P2P_METRICS, - name::peer_name, + name, p2p_context::P2PContext, utils, }; @@ -681,7 +681,7 @@ impl Node { // Connection errors SwarmEvent::OutgoingConnectionError { peer_id, error, .. } => { if let Some(peer) = peer_id { - warn!(peer = %peer_name(peer), %error, "outgoing connection failed"); + warn!(peer = %name::peer_name(peer), %error, "outgoing connection failed"); } else { warn!(%error, "outgoing connection failed"); } @@ -755,7 +755,7 @@ fn record_ping_metrics( return; } - let peer_label = peer_name(peer); + let peer_label = name::peer_name(peer); match result { Ok(duration) => { P2P_METRICS.ping_latency_secs[&peer_label].observe(duration.as_secs_f64()); @@ -779,7 +779,7 @@ fn init_ping_metrics(ctx: &P2PContext) { if Some(*peer) == local { continue; } - P2P_METRICS.ping_success[&peer_name(peer)].set(0); + P2P_METRICS.ping_success[&name::peer_name(peer)].set(0); } } @@ -790,7 +790,7 @@ fn clear_ping_success(ctx: &P2PContext, peer: &PeerId) { return; } - P2P_METRICS.ping_success[&peer_name(peer)].set(0); + P2P_METRICS.ping_success[&name::peer_name(peer)].set(0); } /// Stores identify-reported listen addresses for a peer, gated to known cluster @@ -854,7 +854,7 @@ mod tests { fn ping_metrics_recorded_for_known_peer() { let known = random_peer_id(); let ctx = P2PContext::new([known]); - let label = peer_name(&known); + let label = name::peer_name(&known); record_ping_metrics(&ctx, &known, &Ok(Duration::from_millis(20))); @@ -884,7 +884,7 @@ mod tests { let known = random_peer_id(); let relay = random_peer_id(); let ctx = P2PContext::new([known]); - let label = peer_name(&relay); + let label = name::peer_name(&relay); record_ping_metrics(&ctx, &relay, &Ok(Duration::from_millis(20))); record_ping_metrics(&ctx, &relay, &Err(ping::Failure::Timeout)); @@ -903,7 +903,7 @@ mod tests { fn ping_success_cleared_when_last_connection_closes() { let known = random_peer_id(); let ctx = P2PContext::new([known]); - let label = peer_name(&known); + let label = name::peer_name(&known); record_ping_metrics(&ctx, &known, &Ok(Duration::from_millis(20))); assert_eq!( @@ -924,7 +924,7 @@ mod tests { let known = random_peer_id(); let relay = random_peer_id(); let ctx = P2PContext::new([known]); - let label = peer_name(&relay); + let label = name::peer_name(&relay); clear_ping_success(&ctx, &relay); @@ -946,12 +946,12 @@ mod tests { assert_eq!( P2P_METRICS .ping_success - .get(&peer_name(&peer)) + .get(&name::peer_name(&peer)) .map(Gauge::get), Some(0) ); assert!( - !P2P_METRICS.ping_success.contains(&peer_name(&local)), + !P2P_METRICS.ping_success.contains(&name::peer_name(&local)), "Charon's ping service skips self, so no self series" ); } diff --git a/crates/p2p/src/peer.rs b/crates/p2p/src/peer.rs index 097d651f..92b3e384 100644 --- a/crates/p2p/src/peer.rs +++ b/crates/p2p/src/peer.rs @@ -9,7 +9,7 @@ use libp2p::{Multiaddr, PeerId, identity::PublicKey as Libp2pPublicKey, multiadd use pluto_eth2util::enr::Record; use tokio::sync::watch; -use crate::name::peer_name; +use crate::name; /// Peer error. #[derive(Debug, thiserror::Error)] @@ -86,7 +86,7 @@ impl Peer { id: info.id, addresses: info.addrs.clone(), index: 0, - name: peer_name(&info.id), + name: name::peer_name(&info.id), } } @@ -97,7 +97,7 @@ impl Peer { Ok(Peer { id, index, - name: peer_name(&id), + name: name::peer_name(&id), addresses: vec![], }) } diff --git a/crates/p2p/src/quic_upgrade.rs b/crates/p2p/src/quic_upgrade.rs index 0bcae00a..137d321a 100644 --- a/crates/p2p/src/quic_upgrade.rs +++ b/crates/p2p/src/quic_upgrade.rs @@ -20,7 +20,7 @@ use tokio::time::Interval; use tracing::{debug, info}; use crate::{ - name::peer_name, + name, p2p_context::P2PContext, utils::{ filter_direct_quic_addrs, has_direct_quic_conn, has_direct_tcp_conn, is_quic_addr, @@ -166,7 +166,7 @@ impl QuicUpgradeBehaviour { { backoff.tickers_remaining = backoff.tickers_remaining.saturating_sub(1); debug!( - peer = %peer_name(peer), + peer = %name::peer_name(peer), remaining = backoff.tickers_remaining, backoff_duration_minutes = backoff.backoff_duration, "skipping QUIC upgrade due to backoff" @@ -222,7 +222,7 @@ impl QuicUpgradeBehaviour { if conns.is_empty() { debug!( - peer = %peer_name(&peer_id), + peer = %name::peer_name(&peer_id), "no connection to peer" ); continue; @@ -232,7 +232,7 @@ impl QuicUpgradeBehaviour { if has_direct_quic_conn(&conn_refs) { debug!( - peer = %peer_name(&peer_id), + peer = %name::peer_name(&peer_id), "already has direct QUIC connection to peer" ); @@ -243,7 +243,7 @@ impl QuicUpgradeBehaviour { .collect(); debug!( - peer = %peer_name(&peer_id), + peer = %name::peer_name(&peer_id), "closing {} redundant TCP connections after QUIC upgrade", tcp_conn_ids.len() ); @@ -254,7 +254,7 @@ impl QuicUpgradeBehaviour { if !has_direct_tcp_conn(&conn_refs) { debug!( - peer = %peer_name(&peer_id), + peer = %name::peer_name(&peer_id), "no direct connection via TCP to peer" ); continue; @@ -269,14 +269,14 @@ impl QuicUpgradeBehaviour { if quic_addrs.is_empty() { debug!( - peer = %peer_name(&peer_id), + peer = %name::peer_name(&peer_id), "no known QUIC addresses to peer" ); continue; } info!( - peer = %peer_name(&peer_id), + peer = %name::peer_name(&peer_id), quic_addrs = ?quic_addrs, "trying to upgrade to QUIC connection with peer" ); @@ -315,13 +315,13 @@ impl QuicUpgradeBehaviour { { if is_quic_addr(addr) && !is_relay_addr(addr) { info!( - peer = %peer_name(&peer_id), + peer = %name::peer_name(&peer_id), addr = %addr, "upgraded connection to QUIC" ); debug!( - peer = %peer_name(&peer_id), + peer = %name::peer_name(&peer_id), "closing {} redundant TCP connections after QUIC upgrade", tcp_conn_ids.len() ); @@ -334,7 +334,7 @@ impl QuicUpgradeBehaviour { })); } else { debug!( - peer = %peer_name(&peer_id), + peer = %name::peer_name(&peer_id), addr = %addr, "connected via non-direct address instead of direct QUIC" ); @@ -352,7 +352,7 @@ impl QuicUpgradeBehaviour { && self.pending_upgrades.remove(&peer_id).is_some() { info!( - peer = %peer_name(&peer_id), + peer = %name::peer_name(&peer_id), "failed to connect to peer during QUIC upgrade" ); self.record_failure(peer_id, "dial failed"); diff --git a/crates/p2p/src/relay/manager.rs b/crates/p2p/src/relay/manager.rs index 7a6527d8..52571dab 100644 --- a/crates/p2p/src/relay/manager.rs +++ b/crates/p2p/src/relay/manager.rs @@ -27,7 +27,7 @@ use super::{ }; use crate::{ metrics::P2P_METRICS, - name::peer_name, + name, p2p_context::P2PContext, peer::{MutablePeer, Peer}, }; @@ -231,7 +231,7 @@ impl RelayManager { /// `p2p_relay_connections` — not how many transport connections exist, /// matching Charon's `relay.go`. fn report_relay_connection(relay_id: PeerId, reserved: bool) { - P2P_METRICS.relay_connections[&peer_name(&relay_id)].set(i64::from(reserved)); + P2P_METRICS.relay_connections[&name::peer_name(&relay_id)].set(i64::from(reserved)); } /// Polls every active dial state once, queuing a `ToSwarm::Dial` event for diff --git a/crates/p2p/src/relay/manager/tests.rs b/crates/p2p/src/relay/manager/tests.rs index 1fa16daf..2ff621a7 100644 --- a/crates/p2p/src/relay/manager/tests.rs +++ b/crates/p2p/src/relay/manager/tests.rs @@ -869,7 +869,7 @@ async fn poll_fires_swept_peer_dial_within_the_same_watchdog_pass() { fn relay_connections(relay_id: PeerId) -> Option { P2P_METRICS .relay_connections - .get(&peer_name(&relay_id)) + .get(&crate::name::peer_name(&relay_id)) .map(vise::Gauge::get) } diff --git a/crates/priority/src/calculate.rs b/crates/priority/src/calculate.rs index a5465886..9cd13495 100644 --- a/crates/priority/src/calculate.rs +++ b/crates/priority/src/calculate.rs @@ -3,7 +3,7 @@ use std::collections::{HashMap, HashSet}; -use pluto_consensus::qbft::msg::hash_proto_bytes; +use pluto_consensus::qbft::msg; use pluto_core::corepb::v1::priority::{ PriorityMsg, PriorityResult, PriorityScoredResult, PriorityTopicProposal, PriorityTopicResult, }; @@ -31,7 +31,7 @@ const COUNT_WEIGHT: i64 = MAX_PRIORITIES as i64; /// envelope bytes are hashed directly rather than the inner concrete message. fn hash_any(any: &Any) -> Result { let encoded = any.encode_to_vec(); - hash_proto_bytes(&encoded).map_err(Error::HashProto) + msg::hash_proto_bytes(&encoded).map_err(Error::HashProto) } /// Returns the cluster-wide priorities given the priorities of each peer. diff --git a/crates/priority/src/component.rs b/crates/priority/src/component.rs index 13585edd..47d634b0 100644 --- a/crates/priority/src/component.rs +++ b/crates/priority/src/component.rs @@ -9,7 +9,7 @@ use std::{collections::HashMap, sync::Arc, time::Duration}; use chrono::Utc; use k256::{PublicKey, SecretKey}; use libp2p::PeerId; -use pluto_consensus::qbft::msg::hash_proto; +use pluto_consensus::qbft::msg; use pluto_core::{ corepb::v1::priority::{PriorityMsg, PriorityTopicProposal, PriorityTopicResult}, deadline::{DeadlineCalculator, DeadlinerTask}, @@ -84,7 +84,7 @@ pub fn sign_msg(msg: &PriorityMsg, privkey: &SecretKey) -> Result { let mut clone = msg.clone(); clone.signature = Default::default(); - let hash = hash_proto(&clone).map_err(Error::HashProto)?; + let hash = msg::hash_proto(&clone).map_err(Error::HashProto)?; let sig = pluto_k1util::sign(privkey, &hash).map_err(Error::Sign)?; clone.signature = sig.to_vec().into(); @@ -104,7 +104,7 @@ pub(crate) fn verify_msg_sig(msg: &PriorityMsg, pubkey: &PublicKey) -> Result Result<()> { - let result = calculate_result(msgs, inner.min_required) + let result = calculate::calculate_result(msgs, inner.min_required) .map_err(|e| Error::CalculateResult(Box::new(e)))?; let consensus = inner.consensus.clone(); diff --git a/crates/relay-server/src/p2p.rs b/crates/relay-server/src/p2p.rs index d90a77e6..930bdcce 100644 --- a/crates/relay-server/src/p2p.rs +++ b/crates/relay-server/src/p2p.rs @@ -5,7 +5,7 @@ use std::{collections::HashSet, future::Future, net::SocketAddr, sync::Arc}; use futures::StreamExt; use k256::SecretKey; use libp2p::{Multiaddr, PeerId, core::transport::ListenerId, relay, swarm::SwarmEvent}; -use pluto_p2p::{behaviours::pluto::PlutoBehaviourEvent, name::peer_name}; +use pluto_p2p::{behaviours::pluto::PlutoBehaviourEvent, name}; use tokio::{net::TcpListener, sync::RwLock}; use tokio_util::sync::CancellationToken; use tracing::{debug, info, instrument, warn}; @@ -22,7 +22,7 @@ use pluto_p2p::{ BandwidthFactory, PeerConnectionMetrics, p2p::{Node, NodeType}, p2p_context::P2PContext, - utils::external_multiaddrs, + utils, }; /// Runs a relay P2P node: binds every listener, then serves until `ct` is @@ -269,7 +269,7 @@ pub async fn bind_relay(config: &Config, key: SecretKey) -> Result { // they're advertised on `/` and folded into ENR responses on `/enr` even // when libp2p only sees private listen addresses (e.g., K8s pods behind // NodePort). - let external_addrs = external_multiaddrs(&config.p2p_config, &bound_addrs)?; + let external_addrs = utils::external_multiaddrs(&config.p2p_config, &bound_addrs)?; let state = Arc::new(AppState::new( config.p2p_config.clone(), @@ -392,14 +392,14 @@ fn handle_swarm_event(event: &SwarmEvent>) // Track connections for metrics SwarmEvent::ConnectionEstablished { peer_id, .. } => { - debug!(peer = %peer_name(peer_id), "connection established"); + debug!(peer = %name::peer_name(peer_id), "connection established"); let labels = relay_labels(peer_id); RELAY_METRICS.connection_total[&labels].inc(); RELAY_METRICS.active_connections[&labels].inc_by(1); AddrUpdate::None } SwarmEvent::ConnectionClosed { peer_id, cause, .. } => { - debug!(peer = %peer_name(peer_id), cause = ?cause, "connection closed"); + debug!(peer = %name::peer_name(peer_id), cause = ?cause, "connection closed"); let labels = relay_labels(peer_id); RELAY_METRICS.active_connections[&labels].dec_by(1); AddrUpdate::None @@ -412,20 +412,20 @@ fn handle_swarm_event(event: &SwarmEvent>) renewed, }, )) => { - info!(peer = %peer_name(src_peer_id), renewed, "relay reservation accepted"); + info!(peer = %name::peer_name(src_peer_id), renewed, "relay reservation accepted"); AddrUpdate::None } SwarmEvent::Behaviour(PlutoBehaviourEvent::Inner(relay::Event::ReservationReqDenied { src_peer_id, status, })) => { - warn!(peer = %peer_name(src_peer_id), ?status, "relay reservation denied"); + warn!(peer = %name::peer_name(src_peer_id), ?status, "relay reservation denied"); AddrUpdate::None } SwarmEvent::Behaviour(PlutoBehaviourEvent::Inner(relay::Event::ReservationTimedOut { src_peer_id, })) => { - debug!(peer = %peer_name(src_peer_id), "relay reservation timed out"); + debug!(peer = %name::peer_name(src_peer_id), "relay reservation timed out"); AddrUpdate::None } SwarmEvent::Behaviour(PlutoBehaviourEvent::Inner(relay::Event::CircuitReqAccepted { @@ -433,8 +433,8 @@ fn handle_swarm_event(event: &SwarmEvent>) dst_peer_id, })) => { info!( - src = %peer_name(src_peer_id), - dst = %peer_name(dst_peer_id), + src = %name::peer_name(src_peer_id), + dst = %name::peer_name(dst_peer_id), "relay circuit accepted" ); AddrUpdate::None @@ -451,15 +451,15 @@ fn handle_swarm_event(event: &SwarmEvent>) // visible at warn. if matches!(status, relay::StatusCode::NoReservation) { debug!( - src = %peer_name(src_peer_id), - dst = %peer_name(dst_peer_id), + src = %name::peer_name(src_peer_id), + dst = %name::peer_name(dst_peer_id), ?status, "relay circuit denied" ); } else { warn!( - src = %peer_name(src_peer_id), - dst = %peer_name(dst_peer_id), + src = %name::peer_name(src_peer_id), + dst = %name::peer_name(dst_peer_id), ?status, "relay circuit denied" ); @@ -472,8 +472,8 @@ fn handle_swarm_event(event: &SwarmEvent>) error, })) => { debug!( - src = %peer_name(src_peer_id), - dst = %peer_name(dst_peer_id), + src = %name::peer_name(src_peer_id), + dst = %name::peer_name(dst_peer_id), error = ?error, "relay circuit closed" ); @@ -492,5 +492,5 @@ fn handle_swarm_event(event: &SwarmEvent>) /// The `peer_cluster` label is left empty since the relay server does not /// track cluster membership. fn relay_labels(peer_id: &PeerId) -> PeerWithPeerClusterLabels { - PeerWithPeerClusterLabels::new(peer_name(peer_id), "") + PeerWithPeerClusterLabels::new(name::peer_name(peer_id), "") } diff --git a/crates/relay-server/src/web.rs b/crates/relay-server/src/web.rs index dcbdd6a8..1f43a782 100644 --- a/crates/relay-server/src/web.rs +++ b/crates/relay-server/src/web.rs @@ -24,7 +24,7 @@ use crate::{ config::EXTERNAL_HOST_RESOLVE_INTERVAL, error::{RelayP2PError, Result}, }; -use pluto_p2p::{config::P2PConfig, manet::Manet, name::peer_name}; +use pluto_p2p::{config::P2PConfig, manet::Manet, name}; /// Shared application state for HTTP handlers. #[derive(Clone)] @@ -132,7 +132,7 @@ pub async fn enr_server( info!( "Relay started {peer_name} on {tcp_addrs} and {udp_addrs}", - peer_name = peer_name(&state.peer_id), + peer_name = name::peer_name(&state.peer_id), tcp_addrs = state.p2p_config.tcp_addrs.join(", "), udp_addrs = state.p2p_config.udp_addrs.join(", "), ); From f45ed70b01d545addd79f0f79570d4333eb7b994 Mon Sep 17 00:00:00 2001 From: Bohdan Ohorodnii <273991985+varex83agent@users.noreply.github.com> Date: Mon, 31 Aug 2026 12:01:20 +0200 Subject: [PATCH 2/2] style: cargo fmt after free-function qualification Co-Authored-By: varex83 --- crates/p2p/src/conn_logger.rs | 3 ++- crates/p2p/src/quic_upgrade.rs | 8 ++++++-- 2 files changed, 8 insertions(+), 3 deletions(-) diff --git a/crates/p2p/src/conn_logger.rs b/crates/p2p/src/conn_logger.rs index 1ba845f1..975b4ad1 100644 --- a/crates/p2p/src/conn_logger.rs +++ b/crates/p2p/src/conn_logger.rs @@ -550,7 +550,8 @@ mod tests { let count = behaviour.metrics().inner().relay_connection_types[&labels].get(); assert_eq!(count, 1); - // peer_connection_total should not have been incremented for unknown peer + // peer_connection_total should not have been incremented for unknown + // peer let total = behaviour.metrics().inner().peer_connection_total [&name::peer_name(&unknown_peer)] .get(); diff --git a/crates/p2p/src/quic_upgrade.rs b/crates/p2p/src/quic_upgrade.rs index f146bd62..a7b78346 100644 --- a/crates/p2p/src/quic_upgrade.rs +++ b/crates/p2p/src/quic_upgrade.rs @@ -231,7 +231,9 @@ impl QuicUpgradeBehaviour { let tcp_conn_ids: Vec<_> = conns .iter() - .filter(|c| utils::is_tcp_addr(&c.remote_addr) && !utils::is_relay_addr(&c.remote_addr)) + .filter(|c| { + utils::is_tcp_addr(&c.remote_addr) && !utils::is_relay_addr(&c.remote_addr) + }) .map(|c| c.connection_id) .collect(); @@ -276,7 +278,9 @@ impl QuicUpgradeBehaviour { let tcp_conn_ids: Vec<_> = conns .iter() - .filter(|c| utils::is_tcp_addr(&c.remote_addr) && !utils::is_relay_addr(&c.remote_addr)) + .filter(|c| { + utils::is_tcp_addr(&c.remote_addr) && !utils::is_relay_addr(&c.remote_addr) + }) .map(|c| c.connection_id) .collect();