diff --git a/.env.sample b/.env.sample index 6fefb3d5..a00b8a67 100644 --- a/.env.sample +++ b/.env.sample @@ -4,10 +4,12 @@ PRECONF_ROUTER_ADDRESS=0x8dDEA87FfA296B951881aDacbAe8011090C54cAC TAIKO_INBOX_ADDRESS=0xc2DD6e8DC8d0558F00Cc1FA6A16FFF1A62Cc436B ANCHOR_ADDRESS=0x1670100000000000000000000000000000010001 VALIDATOR_INDEX=1 -L2_RPC_URL=ws://127.0.0.1:8546 +# Ordinary execution RPC; Shasta and Realtime subscriptions use L2_WS_RPC_URL. +L2_RPC_URL=http://127.0.0.1:8545 +L2_WS_RPC_URL=ws://127.0.0.1:8546 L2_AUTH_RPC_URL=http://127.0.0.1:8551 L2_DRIVER_URL=http://127.0.0.1:1235 L1_RPC_URLS=ws://127.0.0.1:32003,wss://123.123.123.123:32001 L1_BEACON_URL=http://127.0.0.1:33001 RUST_LOG=debug,reqwest=info,hyper=info,alloy_transport=info,alloy_rpc_client=info,alloy_provider=info -JWT_SECRET_FILE_PATH=/some/path/jwtsecret.hex \ No newline at end of file +JWT_SECRET_FILE_PATH=/some/path/jwtsecret.hex diff --git a/CHANGELOG.md b/CHANGELOG.md index 7dde9972..f11da8da 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,13 @@ All notable changes to Catalyst are documented here, organized by release version. +## [Unreleased] + +### Features +- Add `L2_WS_RPC_URL` to split WebSocket subscriptions from ordinary L2 execution RPC requests (#2) + +--- + ## [v1.41.0] — 2026-06-25 ### Features diff --git a/common/src/config/mod.rs b/common/src/config/mod.rs index 2a0ab348..59b5b8fd 100644 --- a/common/src/config/mod.rs +++ b/common/src/config/mod.rs @@ -1,6 +1,7 @@ mod config_trait; pub use config_trait::ConfigTrait; +use crate::shared::alloy_tools::{RpcTransport, parse_rpc_transport}; use alloy::primitives::Address; use anyhow::Error; use std::str::FromStr; @@ -27,6 +28,7 @@ pub struct Config { pub preconf_heartbeat_ms: u64, // L2 pub l2_rpc_url: String, + pub l2_ws_rpc_url: String, pub l2_auth_rpc_url: String, pub l2_driver_url: String, /// jwt secret file path for L2 EL and L2 driver @@ -121,6 +123,35 @@ fn get_env_with_deprecation(new_key: &str, deprecated_key: &str) -> Option, +) -> Result { + let l2_rpc_transport = parse_rpc_transport("L2_RPC_URL", l2_rpc_url)?; + + let configured_l2_ws_rpc_url = configured_l2_ws_rpc_url.filter(|url| !url.trim().is_empty()); + if let Some(url) = configured_l2_ws_rpc_url { + let parsed = reqwest::Url::parse(&url) + .map_err(|error| anyhow::anyhow!("L2_WS_RPC_URL must be a valid URL: {error}"))?; + match parsed.scheme() { + "ws" | "wss" => return Ok(parsed.to_string()), + scheme => { + return Err(anyhow::anyhow!( + "L2_WS_RPC_URL must use the ws or wss scheme, got '{scheme}'" + )); + } + } + } + + if let RpcTransport::WebSocket(parsed) = l2_rpc_transport { + return Ok(parsed.to_string()); + } + + Err(anyhow::anyhow!( + "L2_WS_RPC_URL must be set to a ws:// or wss:// URL when L2_RPC_URL does not use WebSocket" + )) +} + impl Config { pub fn read_env_variables() -> Result { // Load environment variables from .env file @@ -501,6 +532,9 @@ impl Config { "ws://127.0.0.1:1234".to_string() }); + let l2_ws_rpc_url = + resolve_l2_ws_rpc_url(&l2_rpc_url, std::env::var("L2_WS_RPC_URL").ok())?; + let l2_auth_rpc_url = get_env_with_deprecation("L2_AUTH_RPC_URL", "TAIKO_GETH_AUTH_RPC_URL").unwrap_or_else( || { @@ -518,6 +552,7 @@ impl Config { let config = Self { preconfer_address, l2_rpc_url, + l2_ws_rpc_url, l2_auth_rpc_url, l2_driver_url, catalyst_node_ecdsa_private_key, @@ -580,6 +615,7 @@ impl Config { r#" Configuration:{} L2 RPC URL: {}, +L2 WebSocket RPC URL: {}, L2 auth RPC URL: {}, L2 driver URL: {}, L1 RPC URL: {}, @@ -636,6 +672,7 @@ internal server port: {} "".to_string() }, config.l2_rpc_url, + config.l2_ws_rpc_url, config.l2_auth_rpc_url, config.l2_driver_url, match config.l1_rpc_urls.split_first() { @@ -697,3 +734,150 @@ internal server port: {} Ok(config) } } + +#[cfg(test)] +mod tests { + use super::resolve_l2_ws_rpc_url; + + #[test] + fn explicit_l2_ws_rpc_url_is_used_with_http_rpc() { + let resolved = resolve_l2_ws_rpc_url( + "http://127.0.0.1:8545", + Some("ws://127.0.0.1:8546".to_string()), + ) + .unwrap(); + + assert_eq!(resolved, "ws://127.0.0.1:8546/"); + } + + #[test] + fn explicit_l2_ws_rpc_url_trims_surrounding_whitespace() { + let resolved = resolve_l2_ws_rpc_url( + "http://127.0.0.1:8545", + Some(" wss://l2.example ".to_string()), + ) + .unwrap(); + + assert_eq!(resolved, "wss://l2.example/"); + } + + #[test] + fn explicit_l2_ws_rpc_url_uses_canonical_serialization() { + let resolved = resolve_l2_ws_rpc_url( + "http://127.0.0.1:8545", + Some("wss://l2.example/pa th".to_string()), + ) + .unwrap(); + + assert_eq!(resolved, "wss://l2.example/pa%20th"); + } + + #[test] + fn legacy_websocket_l2_rpc_url_is_reused() { + let resolved = resolve_l2_ws_rpc_url("wss://l2.example", None).unwrap(); + + assert_eq!(resolved, "wss://l2.example/"); + } + + #[test] + fn legacy_websocket_l2_rpc_url_trims_surrounding_whitespace() { + let resolved = resolve_l2_ws_rpc_url(" wss://l2.example ", None).unwrap(); + + assert_eq!(resolved, "wss://l2.example/"); + } + + #[test] + fn legacy_websocket_l2_rpc_url_uses_canonical_serialization() { + let resolved = resolve_l2_ws_rpc_url("wss://l2.example/pa th", None).unwrap(); + + assert_eq!(resolved, "wss://l2.example/pa%20th"); + } + + #[test] + fn empty_l2_ws_rpc_url_reuses_legacy_websocket_l2_rpc_url() { + let resolved = resolve_l2_ws_rpc_url("wss://l2.example", Some(String::new())).unwrap(); + + assert_eq!(resolved, "wss://l2.example/"); + } + + #[test] + fn whitespace_only_l2_ws_rpc_url_reuses_legacy_websocket_l2_rpc_url() { + let resolved = + resolve_l2_ws_rpc_url("wss://l2.example", Some(" \t ".to_string())).unwrap(); + + assert_eq!(resolved, "wss://l2.example/"); + } + + #[test] + fn http_l2_rpc_url_requires_explicit_websocket_url() { + let error = resolve_l2_ws_rpc_url("https://l2.example", None) + .unwrap_err() + .to_string(); + + assert!(error.contains("L2_WS_RPC_URL must be set")); + } + + #[test] + fn malformed_l2_rpc_url_is_rejected_with_explicit_websocket_url() { + let error = resolve_l2_ws_rpc_url("not a valid URL", Some("wss://l2.example".to_string())) + .unwrap_err() + .to_string(); + + assert!(error.contains("L2_RPC_URL must be a valid URL")); + } + + #[test] + fn empty_l2_rpc_url_has_a_specific_error() { + let error = resolve_l2_ws_rpc_url("", Some("wss://l2.example".to_string())) + .unwrap_err() + .to_string(); + + assert_eq!(error, "L2_RPC_URL must not be empty"); + } + + #[test] + fn l2_rpc_url_rejects_unsupported_scheme() { + let error = resolve_l2_ws_rpc_url("ftp://l2.example", Some("wss://l2.example".to_string())) + .unwrap_err() + .to_string(); + + assert!(error.contains("L2_RPC_URL must use the http, https, ws or wss scheme")); + } + + #[test] + fn explicit_l2_ws_rpc_url_rejects_non_websocket_scheme() { + let error = resolve_l2_ws_rpc_url( + "http://127.0.0.1:8545", + Some("http://127.0.0.1:8546".to_string()), + ) + .unwrap_err() + .to_string(); + + assert!(error.contains("L2_WS_RPC_URL must use the ws or wss scheme")); + } + + #[test] + fn explicit_l2_ws_rpc_url_rejects_unsupported_scheme_as_websocket_only() { + let error = resolve_l2_ws_rpc_url( + "http://127.0.0.1:8545", + Some("ftp://l2.example".to_string()), + ) + .unwrap_err() + .to_string(); + + assert_eq!( + error, + "L2_WS_RPC_URL must use the ws or wss scheme, got 'ftp'" + ); + } + + #[test] + fn explicit_l2_ws_rpc_url_rejects_malformed_url() { + let error = + resolve_l2_ws_rpc_url("http://127.0.0.1:8545", Some("not a valid URL".to_string())) + .unwrap_err() + .to_string(); + + assert!(error.contains("L2_WS_RPC_URL must be a valid URL")); + } +} diff --git a/common/src/shared/alloy_tools.rs b/common/src/shared/alloy_tools.rs index 3bc8df9a..8e6a6c24 100644 --- a/common/src/shared/alloy_tools.rs +++ b/common/src/shared/alloy_tools.rs @@ -75,22 +75,22 @@ fn find_errors_from_trace(trace_str: &str) -> Option { pub async fn construct_alloy_provider( signer: &Signer, - execution_ws_rpc_url: &str, + execution_rpc_url: &str, ) -> Result { match signer { Signer::PrivateKey(private_key, _) => { debug!( "Creating alloy provider with URL: {} and private key signer.", - execution_ws_rpc_url + execution_rpc_url ); let signer = PrivateKeySigner::from_str(private_key.as_str())?; - Ok(create_alloy_provider_with_wallet(signer.into(), execution_ws_rpc_url).await?) + Ok(create_alloy_provider_with_wallet(signer.into(), execution_rpc_url).await?) } Signer::Web3signer(web3signer, address) => { debug!( "Creating alloy provider with URL: {} and web3signer signer.", - execution_ws_rpc_url + execution_rpc_url ); let preconfer_address = *address; @@ -100,52 +100,103 @@ pub async fn construct_alloy_provider( )?; let wallet = EthereumWallet::new(tx_signer); - Ok(create_alloy_provider_with_wallet(wallet, execution_ws_rpc_url).await?) + Ok(create_alloy_provider_with_wallet(wallet, execution_rpc_url).await?) } } } +#[derive(Debug)] +pub(crate) enum RpcTransport { + Http(reqwest::Url), + WebSocket(reqwest::Url), +} + +pub(crate) fn parse_rpc_transport(source: &str, url: &str) -> Result { + if url.trim().is_empty() { + return Err(anyhow::anyhow!("{source} must not be empty")); + } + + let parsed = reqwest::Url::parse(url) + .map_err(|error| anyhow::anyhow!("{source} must be a valid URL: {error}"))?; + + match parsed.scheme() { + "http" | "https" => Ok(RpcTransport::Http(parsed)), + "ws" | "wss" => Ok(RpcTransport::WebSocket(parsed)), + scheme => Err(anyhow::anyhow!( + "{source} must use the http, https, ws or wss scheme, got '{scheme}'" + )), + } +} + async fn create_alloy_provider_with_wallet( wallet: EthereumWallet, url: &str, ) -> Result { - if url.contains("ws://") || url.contains("wss://") { - let ws = WsConnect::new(url); - Ok(ProviderBuilder::new() - .wallet(wallet) - .connect_ws(ws.clone()) - .await - .map_err(|e| Error::msg(format!("Execution layer: Failed to connect to WS: {e}")))? - .erased()) - } else if url.contains("http://") || url.contains("https://") { - Ok(ProviderBuilder::new() + match parse_rpc_transport("RPC URL", url)? { + RpcTransport::Http(url) => Ok(ProviderBuilder::new() .wallet(wallet) - .connect_http(url.parse::()?) - .erased()) - } else { - Err(anyhow::anyhow!( - "Invalid URL, only websocket and http are supported: {}", - url - )) + .connect_http(url) + .erased()), + RpcTransport::WebSocket(url) => { + let ws = WsConnect::new(url.as_str()); + Ok(ProviderBuilder::new() + .wallet(wallet) + .connect_ws(ws.clone()) + .await + .map_err(|e| Error::msg(format!("Execution layer: Failed to connect to WS: {e}")))? + .erased()) + } } } pub async fn create_alloy_provider_without_wallet(url: &str) -> Result { - if url.contains("ws://") || url.contains("wss://") { - let ws = WsConnect::new(url); - Ok(ProviderBuilder::new() - .connect_ws(ws.clone()) - .await - .map_err(|e| Error::msg(format!("Execution layer: Failed to connect to WS: {e}")))? - .erased()) - } else if url.contains("http://") || url.contains("https://") { - Ok(ProviderBuilder::new() - .connect_http(url.parse::()?) - .erased()) - } else { - Err(anyhow::anyhow!( - "Invalid URL, only websocket and http are supported: {}", - url - )) + match parse_rpc_transport("RPC URL", url)? { + RpcTransport::Http(url) => Ok(ProviderBuilder::new().connect_http(url).erased()), + RpcTransport::WebSocket(url) => { + let ws = WsConnect::new(url.as_str()); + Ok(ProviderBuilder::new() + .connect_ws(ws.clone()) + .await + .map_err(|e| Error::msg(format!("Execution layer: Failed to connect to WS: {e}")))? + .erased()) + } + } +} + +#[cfg(test)] +mod tests { + use super::{RpcTransport, create_alloy_provider_without_wallet, parse_rpc_transport}; + use std::time::Duration; + use tokio::{net::TcpListener, time::timeout}; + + #[test] + fn http_url_with_websocket_text_in_path_uses_http_transport() { + let transport = parse_rpc_transport("RPC URL", "https://l2.example/ws://archive").unwrap(); + + assert!(matches!(transport, RpcTransport::Http(_))); + } + + #[tokio::test] + async fn websocket_url_with_whitespace_is_normalized_before_connecting() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let accepted = tokio::spawn(async move { + matches!( + timeout(Duration::from_secs(1), listener.accept()).await, + Ok(Ok(_)) + ) + }); + + let url = format!(" ws://{address} "); + let _ = timeout( + Duration::from_secs(1), + create_alloy_provider_without_wallet(&url), + ) + .await; + + assert!( + accepted.await.unwrap(), + "normalized URL did not reach the WebSocket server" + ); } } diff --git a/realtime/src/lib.rs b/realtime/src/lib.rs index 821f8ad2..d5cea51e 100644 --- a/realtime/src/lib.rs +++ b/realtime/src/lib.rs @@ -111,7 +111,7 @@ pub async fn create_realtime_node( .first() .ok_or_else(|| anyhow::anyhow!("L1 RPC URL is required"))? .clone(), - config.l2_rpc_url.clone(), + config.l2_ws_rpc_url.clone(), realtime_config.realtime_inbox, cancel_token.clone(), "ProposedAndProved", diff --git a/shasta/src/lib.rs b/shasta/src/lib.rs index 0939d4b8..b59ac8af 100644 --- a/shasta/src/lib.rs +++ b/shasta/src/lib.rs @@ -125,7 +125,7 @@ pub async fn create_shasta_node( .first() .expect("L1 RPC URL is required") .clone(), - config.l2_rpc_url.clone(), + config.l2_ws_rpc_url.clone(), shasta_config.shasta_inbox, cancel_token.clone(), "Proposed",