diff --git a/Cargo.lock b/Cargo.lock index 232a0a9bc3..b37bef0b75 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4948,10 +4948,11 @@ dependencies = [ [[package]] name = "js-sys" -version = "0.3.61" +version = "0.3.89" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "445dde2150c55e483f3d8416706b97ec8e8237c307e5b7b4b8dd15e6af2a0730" +checksum = "f4eacb0641a310445a4c513f2a5e23e19952e269c6a38887254d5f837a305506" dependencies = [ + "once_cell", "wasm-bindgen", ] @@ -5132,9 +5133,9 @@ dependencies = [ [[package]] name = "keccak" -version = "0.1.6" +version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cb26cec98cce3a3d96cbb7bced3c4b16e3d13f27ec56dbd62cbc8f39cfb9d653" +checksum = "3afef3b6eff9ce9d8ff9b3601125eec7f0c8cbac7abd14f355d053fa56c98768" dependencies = [ "cpufeatures", ] @@ -5191,9 +5192,9 @@ dependencies = [ [[package]] name = "libc" -version = "0.2.178" +version = "0.2.182" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "37c93d8daa9d8a012fd8ab92f088405fb202ea0b6ab73ee2482ae66af4f42091" +checksum = "6800badb6cb2082ffd7b6a67e6125bb39f18782f793520caee8cb8846be06112" [[package]] name = "libloading" @@ -5599,13 +5600,13 @@ dependencies = [ [[package]] name = "libredox" -version = "0.1.9" +version = "0.1.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "391290121bad3d37fbddad76d8f5d1c1c314cfc646d143d7e07a3086ddff0ce3" +checksum = "3d0b95e02c851351f877147b7deea7b1afb1df71b63aa5f8270716e0c5720616" dependencies = [ "bitflags 2.11.0", "libc", - "redox_syscall 0.5.17", + "redox_syscall 0.7.1", ] [[package]] @@ -5725,9 +5726,9 @@ checksum = "f051f77a7c8e6957c0696eac88f26b0117e54f52d3fc682ab19397a8812846a4" [[package]] name = "linux-raw-sys" -version = "0.11.0" +version = "0.12.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "df1d3c3b53da64cf5760482273a98e575c651a67eec7f77df96b5b642de8f039" +checksum = "32a66949e030da00e8c7d4434b251670a91556f4144941d37452769c25d58a53" [[package]] name = "lock_api" @@ -9863,6 +9864,15 @@ dependencies = [ "bitflags 2.11.0", ] +[[package]] +name = "redox_syscall" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "35985aa610addc02e24fc232012c86fd11f14111180f902b67e2d5331f8ebf2b" +dependencies = [ + "bitflags 2.11.0", +] + [[package]] name = "redox_users" version = "0.3.5" @@ -10338,14 +10348,14 @@ dependencies = [ [[package]] name = "rustix" -version = "1.1.3" +version = "1.1.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "146c9e247ccc180c1f61615433868c99f3de3ae256a30a43b49f67c2d9171f34" +checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190" dependencies = [ "bitflags 2.11.0", "errno", "libc", - "linux-raw-sys 0.11.0", + "linux-raw-sys 0.12.1", "windows-sys 0.60.2", ] @@ -12479,7 +12489,10 @@ dependencies = [ "futures-channel", "futures-timer", "hex", + "jsonrpc-core", + "jsonrpc-ws-server", "libloading", + "once_cell", "parking_lot 0.12.1", "rand 0.9.2", "rand_core 0.6.4", @@ -12495,7 +12508,6 @@ dependencies = [ "starcoin-rpc-api", "starcoin-rpc-client", "starcoin-service-registry", - "starcoin-stratum", "starcoin-time-service", "starcoin-vm-types", "starcoin-vm1-types", @@ -12791,7 +12803,6 @@ dependencies = [ "starcoin-state-service", "starcoin-statedb", "starcoin-storage", - "starcoin-stratum", "starcoin-sync", "starcoin-sync-api", "starcoin-txpool", @@ -13293,30 +13304,6 @@ dependencies = [ "thiserror", ] -[[package]] -name = "starcoin-stratum" -version = "2.1.1" -dependencies = [ - "anyhow", - "byteorder", - "bytes 1.6.1", - "futures 0.3.31", - "hex", - "once_cell", - "serde 1.0.228", - "serde_json", - "starcoin-config", - "starcoin-crypto 1.10.0-rc.2 (git+https://github.com/starcoinorg/starcoin-crypto?rev=bc3ca2138cf66d646f14a0f673ae5d853743f50a)", - "starcoin-logger", - "starcoin-miner", - "starcoin-service-registry", - "starcoin-vm-types", - "starcoin-vm1-types", - "stest", - "tokio", - "tokio-util 0.7.18", -] - [[package]] name = "starcoin-sync" version = "2.1.1" @@ -14649,14 +14636,14 @@ dependencies = [ [[package]] name = "tempfile" -version = "3.25.0" +version = "3.26.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0136791f7c95b1f6dd99f9cc786b91bb81c3800b639b3478e561ddb7be95e5f1" +checksum = "82a72c767771b47409d2345987fda8628641887d5466101319899796367354a0" dependencies = [ "fastrand 2.3.0", "getrandom 0.3.3", "once_cell", - "rustix 1.1.3", + "rustix 1.1.4", "windows-sys 0.60.2", ] @@ -15952,26 +15939,14 @@ dependencies = [ [[package]] name = "wasm-bindgen" -version = "0.2.84" +version = "0.2.112" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "31f8dcbc21f30d9b8f2ea926ecb58f6b91192c17e9d33594b3df58b2007ca53b" +checksum = "05d7d0fce354c88b7982aec4400b3e7fcf723c32737cef571bd165f7613557ee" dependencies = [ "cfg-if 1.0.0", - "wasm-bindgen-macro", -] - -[[package]] -name = "wasm-bindgen-backend" -version = "0.2.84" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "95ce90fd5bcc06af55a641a86428ee4229e44e07033963a2290a8e241607ccb9" -dependencies = [ - "bumpalo", - "log 0.4.27", "once_cell", - "proc-macro2 1.0.101", - "quote 1.0.36", - "syn 1.0.107", + "rustversion", + "wasm-bindgen-macro", "wasm-bindgen-shared", ] @@ -15989,9 +15964,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-macro" -version = "0.2.84" +version = "0.2.112" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4c21f77c0bedc37fd5dc21f897894a5ca01e7bb159884559461862ae90c0b4c5" +checksum = "55839b71ba921e4f75b674cb16f843f4b1f3b26ddfcb3454de1cf65cc021ec0f" dependencies = [ "quote 1.0.36", "wasm-bindgen-macro-support", @@ -15999,22 +15974,25 @@ dependencies = [ [[package]] name = "wasm-bindgen-macro-support" -version = "0.2.84" +version = "0.2.112" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2aff81306fcac3c7515ad4e177f521b5c9a15f2b08f4e32d823066102f35a5f6" +checksum = "caf2e969c2d60ff52e7e98b7392ff1588bffdd1ccd4769eba27222fd3d621571" dependencies = [ + "bumpalo", "proc-macro2 1.0.101", "quote 1.0.36", - "syn 1.0.107", - "wasm-bindgen-backend", + "syn 2.0.114", "wasm-bindgen-shared", ] [[package]] name = "wasm-bindgen-shared" -version = "0.2.84" +version = "0.2.112" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0046fef7e28c3804e5e38bfa31ea2a0f73905319b677e57ebe37e49358989b5d" +checksum = "0861f0dcdf46ea819407495634953cdcc8a8c7215ab799a7a7ce366be71c7b30" +dependencies = [ + "unicode-ident", +] [[package]] name = "wasm-timer" @@ -16033,9 +16011,9 @@ dependencies = [ [[package]] name = "web-sys" -version = "0.3.61" +version = "0.3.89" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e33b99f4b23ba3eec1a53ac264e35a755f00e966e0065077d6027c0f575b0b97" +checksum = "10053fbf9a374174094915bbce141e87a6bf32ecd9a002980db4b638405e8962" dependencies = [ "js-sys", "wasm-bindgen", @@ -16883,7 +16861,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "af3a19837351dc82ba89f8a125e22a3c475f05aba604acc023d62b2739ae2909" dependencies = [ "libc", - "rustix 1.1.3", + "rustix 1.1.4", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 33da873521..e5af84e83e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -130,7 +130,6 @@ members = [ #"cmd/peer-watcher", "cmd/airdrop", #"cmd/replay", - "stratum", "cmd/miner_client/api", #"cmd/db-exporter", "cmd/genesis-nft-miner", @@ -248,7 +247,6 @@ default-members = [ "cmd/airdrop", #"cmd/replay", "cmd/genesis-nft-miner", - "stratum", "cmd/miner_client/api", #"cmd/db-exporter", "cmd/cmd-utils", @@ -256,6 +254,8 @@ default-members = [ "dataformat-generator-vm2", ] +exclude = ["cmd/stratumd"] + [profile.dev] panic = "unwind" @@ -467,6 +467,7 @@ pin-utils = "0.1.0" pretty = "0.12.5" proc-macro2 = "1.0" prometheus = "0.13.0" +postgres = "0.19.12" proptest = "1.10.0" proptest-derive = "0.8.0" quote = "1.0.16" @@ -559,7 +560,6 @@ starcoin-state-store-api = { path = "state/state-store-api" } starcoin-state-tree = { path = "state/state-tree" } starcoin-statedb = { path = "state/statedb" } starcoin-storage = { path = "storage" } -starcoin-stratum = { path = "stratum" } starcoin-sync = { path = "sync" } starcoin-sync-api = { path = "sync/api" } starcoin-system = { path = "commons/system", package = "starcoin-system" } diff --git a/cmd/miner_client/Cargo.toml b/cmd/miner_client/Cargo.toml index 180bdda498..1d33800fec 100644 --- a/cmd/miner_client/Cargo.toml +++ b/cmd/miner_client/Cargo.toml @@ -36,12 +36,14 @@ starcoin-miner-client-api = { workspace = true } starcoin-rpc-api = { workspace = true } starcoin-rpc-client = { workspace = true } starcoin-service-registry = { workspace = true } -starcoin-stratum = { workspace = true } stest = { workspace = true } thiserror = { workspace = true } [dev-dependencies] +jsonrpc-core = { workspace = true } +jsonrpc-ws-server = { workspace = true } starcoin-miner = { workspace = true } +once_cell = { workspace = true } [package] authors = { workspace = true } diff --git a/cmd/miner_client/src/lib.rs b/cmd/miner_client/src/lib.rs index 4d4149a990..af1501671f 100644 --- a/cmd/miner_client/src/lib.rs +++ b/cmd/miner_client/src/lib.rs @@ -7,6 +7,7 @@ pub mod miner; mod solver; pub mod stratum_client; pub mod stratum_client_service; +pub mod stratum_compat; use anyhow::Result; use futures::stream::BoxStream; use starcoin_config::TimeService; diff --git a/cmd/miner_client/src/main.rs b/cmd/miner_client/src/main.rs index 52ff4a3106..4b0d64499e 100644 --- a/cmd/miner_client/src/main.rs +++ b/cmd/miner_client/src/main.rs @@ -9,8 +9,8 @@ use starcoin_miner_client::stratum_client::StratumJobClient; use starcoin_miner_client::stratum_client_service::{ StratumClientService, StratumClientServiceServiceFactory, }; +use starcoin_miner_client::stratum_compat::LoginRequest; use starcoin_service_registry::{RegistryAsyncService, RegistryService}; -use starcoin_stratum::rpc::LoginRequest; use starcoin_time_service::RealTimeService; use std::sync::Arc; diff --git a/cmd/miner_client/src/stratum_client.rs b/cmd/miner_client/src/stratum_client.rs index 11294aad43..f17fca3547 100644 --- a/cmd/miner_client/src/stratum_client.rs +++ b/cmd/miner_client/src/stratum_client.rs @@ -1,4 +1,7 @@ -use crate::stratum_client_service::{ShareRequest, StratumClientService, SubmitSealRequest}; +use crate::stratum_client_service::{ + LoginServiceRequest, ShareRequest, StratumClientService, SubmitSealRequest, +}; +use crate::stratum_compat::{target_hex_to_difficulty, LoginRequest}; use crate::{ConsensusStrategy, JobClient, SealEvent}; use anyhow::Result; use byteorder::{LittleEndian, WriteBytesExt}; @@ -6,8 +9,6 @@ use futures::future; use futures::stream::{BoxStream, StreamExt}; use starcoin_logger::prelude::error; use starcoin_service_registry::ServiceRef; -use starcoin_stratum::rpc::LoginRequest; -use starcoin_stratum::target_hex_to_difficulty; use starcoin_time_service::TimeService; use starcoin_types::system_events::{MintBlockEvent, MintEventExtra}; use std::sync::Arc; @@ -39,7 +40,7 @@ impl JobClient for StratumJobClient { let login = self.login.clone(); let fut = async move { let stream = srv - .send(login) + .send(LoginServiceRequest(login)) .await? .await .map_err(|e| anyhow::anyhow!(format!("{}", e))) diff --git a/cmd/miner_client/src/stratum_client_service.rs b/cmd/miner_client/src/stratum_client_service.rs index ecfd95bff3..14e69c16e5 100644 --- a/cmd/miner_client/src/stratum_client_service.rs +++ b/cmd/miner_client/src/stratum_client_service.rs @@ -1,3 +1,7 @@ +use crate::stratum_compat::JsonStreamCodec; +pub use crate::stratum_compat::{ + LoginRequest, ShareRequest, Status, StratumJob, StratumJobResponse, +}; use anyhow::anyhow; use anyhow::Result; use futures::{select, Sink, SinkExt, Stream, StreamExt, TryStreamExt}; @@ -9,10 +13,6 @@ use starcoin_logger::prelude::*; use starcoin_service_registry::{ ActorService, ServiceContext, ServiceFactory, ServiceHandler, ServiceRequest, }; -use starcoin_stratum::codec::JsonStreamCodec; -pub use starcoin_stratum::rpc::{ - LoginRequest, ShareRequest, Status, StratumJob, StratumJobResponse, -}; use std::collections::HashMap; use std::convert::TryFrom; use std::pin::Pin; @@ -127,6 +127,13 @@ impl ServiceRequest for SubmitSealRequest { #[derive(Clone, Debug, Deserialize, Serialize)] pub struct SubmitSealRequest(pub ShareRequest); +#[derive(Clone, Debug)] +pub struct LoginServiceRequest(pub LoginRequest); + +impl ServiceRequest for LoginServiceRequest { + type Response = oneshot::Receiver>; +} + fn build_request_string( method: &str, argument: &T, @@ -277,16 +284,16 @@ impl ActorService for StratumClientService { } } -impl ServiceHandler for StratumClientService { +impl ServiceHandler for StratumClientService { fn handle( &mut self, - msg: LoginRequest, + msg: LoginServiceRequest, _ctx: &mut ServiceContext, - ) -> ::Response { + ) -> ::Response { match self.sender.clone() { Some(sender) => { let (s, r) = futures::channel::oneshot::channel(); - if let Err(err) = sender.unbounded_send(Request::LoginRequest(msg, s)) { + if let Err(err) = sender.unbounded_send(Request::LoginRequest(msg.0, s)) { error!("stratum handle login_request failed: {}", err); } r diff --git a/stratum/src/codec.rs b/cmd/miner_client/src/stratum_compat.rs similarity index 58% rename from stratum/src/codec.rs rename to cmd/miner_client/src/stratum_compat.rs index 28ec1b61b3..4c60f9fad3 100644 --- a/stratum/src/codec.rs +++ b/cmd/miner_client/src/stratum_compat.rs @@ -1,15 +1,74 @@ -use bytes::BytesMut; -use std::{io, str}; +use serde::{Deserialize, Serialize}; +use starcoin_types::block::BlockHeaderExtra; +use starcoin_types::U256; +use std::{convert::TryInto, io, str}; +use tokio_util::bytes::BytesMut; use tokio_util::codec::{Decoder, Encoder}; -/// Separator for enveloping messages in streaming codecs. +const MAX_INBOUND_BYTES: usize = 256 * 1024; + +pub fn target_hex_to_difficulty(target: &str) -> anyhow::Result { + let mut temp = hex::decode(target)?; + temp.reverse(); + let temp = hex::encode(temp); + let temp = U256::from_str_radix(&temp, 16)?; + Ok(U256::from(u64::MAX) / temp) +} + +#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)] +pub struct LoginRequest { + pub login: String, + pub pass: String, + pub agent: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub algo: Option>, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ShareRequest { + pub id: String, + pub job_id: String, + pub nonce: String, + pub result: String, +} + +#[derive(Debug, PartialEq, Eq, Clone, Deserialize, Serialize)] +pub struct Status { + pub status: String, +} + +#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)] +pub struct StratumJobResponse { + #[serde(skip_serializing_if = "Option::is_none")] + pub login: Option, + pub id: String, + pub status: String, + pub job: StratumJob, +} + +#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)] +pub struct StratumJob { + pub height: u64, + pub id: String, + pub target: String, + pub job_id: String, + pub blob: String, +} + +impl StratumJob { + pub fn get_extra(&self) -> anyhow::Result { + let blob = hex::decode(&self.blob)?; + if blob.len() != 76 { + return Err(anyhow::anyhow!("Invalid stratum job")); + } + let extra: [u8; 4] = blob[35..39].try_into()?; + Ok(BlockHeaderExtra::new(extra)) + } +} + #[derive(Debug, Clone)] pub enum Separator { - /// No envelope is expected between messages. Decoder will try to figure out - /// message boundaries by accumulating incoming bytes until valid JSON is formed. - /// Encoder will send messages without any boundaries between requests. Empty, - /// Byte is used as a sentinel between messages. Byte(u8), } @@ -19,7 +78,6 @@ impl Default for Separator { } } -/// Stream codec for streaming JSON-RPC over TCP. #[derive(Debug, Default)] pub struct JsonStreamCodec { incoming_separator: Separator, @@ -27,12 +85,10 @@ pub struct JsonStreamCodec { } impl JsonStreamCodec { - /// Default codec with streaming input data. Input can be both enveloped and not. pub fn stream_incoming() -> Self { Self::new(Separator::Empty, Default::default()) } - /// New custom stream codec. pub fn new(incoming_separator: Separator, outgoing_separator: Separator) -> Self { Self { incoming_separator, @@ -50,6 +106,12 @@ impl Decoder for JsonStreamCodec { type Error = io::Error; fn decode(&mut self, buf: &mut BytesMut) -> io::Result> { + if buf.len() > MAX_INBOUND_BYTES { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "jsonrpc message too large", + )); + } if let Separator::Byte(separator) = self.incoming_separator { if let Some(i) = buf.as_ref().iter().position(|&b| b == separator) { let line = buf.split_to(i); @@ -57,7 +119,7 @@ impl Decoder for JsonStreamCodec { match str::from_utf8(line.as_ref()) { Ok(s) => Ok(Some(s.to_string())), - Err(_) => Err(io::Error::new(io::ErrorKind::Other, "invalid UTF-8")), + Err(_) => Err(io::Error::other("invalid UTF-8")), } } else { Ok(None) diff --git a/cmd/miner_client/tests/stratum_compat.rs b/cmd/miner_client/tests/stratum_compat.rs new file mode 100644 index 0000000000..8228a568de --- /dev/null +++ b/cmd/miner_client/tests/stratum_compat.rs @@ -0,0 +1,258 @@ +use anyhow::{Context, Result}; +use futures::StreamExt; +use jsonrpc_core::{Error as JsonRpcError, IoHandler, Params}; +use jsonrpc_ws_server::{Server as WsServer, ServerBuilder}; +use once_cell::sync::Lazy; +use serde_json::json; +use starcoin_config::MinerClientConfig; +use starcoin_miner_client::stratum_client::StratumJobClient; +use starcoin_miner_client::stratum_client_service::{ + StratumClientService, StratumClientServiceServiceFactory, +}; +use starcoin_miner_client::stratum_compat::LoginRequest; +use starcoin_miner_client::JobClient; +use starcoin_service_registry::{RegistryAsyncService, RegistryService}; +use starcoin_time_service::RealTimeService; +use starcoin_types::genesis_config::ConsensusStrategy; +use starcoin_types::system_events::{MintBlockEvent, SealEvent}; +use starcoin_types::U256; +use std::net::{Ipv4Addr, SocketAddr}; +use std::path::PathBuf; +use std::process::{Child, Command, Stdio}; +use std::sync::{Arc, Mutex as StdMutex}; +use std::time::Duration; +use tokio::net::TcpStream; +use tokio::sync::Mutex; +use tokio::time::Instant; + +static TEST_MUTEX: Lazy> = Lazy::new(|| Mutex::new(())); + +#[derive(Clone)] +struct MockRpcState { + current_job: MintBlockEvent, + submit_count: usize, +} + +struct MockMiningRpc { + state: Arc>, + server: Option, + addr: SocketAddr, +} + +impl MockMiningRpc { + fn start(initial_job: MintBlockEvent) -> Result { + let state = Arc::new(StdMutex::new(MockRpcState { + current_job: initial_job, + submit_count: 0, + })); + let mut io = IoHandler::default(); + + let state_for_get_job = state.clone(); + io.add_sync_method("mining.get_job", move |_params: Params| { + let guard = state_for_get_job + .lock() + .map_err(|_| JsonRpcError::internal_error())?; + serde_json::to_value(Some(guard.current_job.clone())) + .map_err(|_| JsonRpcError::internal_error()) + }); + + let state_for_submit = state.clone(); + io.add_sync_method("mining.submit", move |params: Params| { + let (_minting_blob, _nonce, _extra): (String, u32, String) = params + .parse() + .map_err(|_| JsonRpcError::invalid_params("invalid submit params"))?; + let mut guard = state_for_submit + .lock() + .map_err(|_| JsonRpcError::internal_error())?; + guard.submit_count = guard.submit_count.saturating_add(1); + let block_hash = guard.current_job.parent_hash; + Ok(json!({ "block_hash": block_hash })) + }); + + let addr = SocketAddr::new(Ipv4Addr::LOCALHOST.into(), pick_free_port()?); + let server = ServerBuilder::new(io) + .start(&addr) + .context("start mock mining rpc ws server failed")?; + Ok(Self { + state, + server: Some(server), + addr, + }) + } + + fn ws_url(&self) -> String { + format!("ws://{}", self.addr) + } + + fn submit_count(&self) -> Result { + let guard = self + .state + .lock() + .map_err(|_| anyhow::anyhow!("mock rpc mutex poisoned"))?; + Ok(guard.submit_count) + } +} + +impl Drop for MockMiningRpc { + fn drop(&mut self) { + if let Some(server) = self.server.take() { + std::mem::forget(server); + } + } +} + +struct StratumdProcess { + child: Child, +} + +impl StratumdProcess { + async fn spawn(listen: SocketAddr, node_rpc: &str) -> Result { + let bin = resolve_stratumd_bin()?; + let mut cmd = Command::new(&bin); + cmd.arg("--listen") + .arg(listen.to_string()) + .arg("--node-rpc") + .arg(node_rpc) + .arg("--job-poll-ms") + .arg("50") + .stdout(Stdio::null()) + .stderr(Stdio::null()); + + let mut child = cmd.spawn().context("spawn stratumd failed")?; + wait_for_server_ready(&mut child, listen, Duration::from_secs(6)).await?; + Ok(Self { child }) + } +} + +impl Drop for StratumdProcess { + fn drop(&mut self) { + let _ = self.child.kill(); + let _ = self.child.wait(); + } +} + +fn resolve_stratumd_bin() -> Result { + let bin = std::env::var("STRATUMD_BIN").context( + "STRATUMD_BIN is not set. Point it to standalone starcoin_stratumd binary path.", + )?; + let path = PathBuf::from(bin); + if !path.exists() { + return Err(anyhow::anyhow!( + "STRATUMD_BIN path does not exist: {:?}", + path + )); + } + Ok(path) +} + +fn pick_free_port() -> Result { + let listener = std::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))?; + Ok(listener.local_addr()?.port()) +} + +fn build_mint_event(number: u64) -> MintBlockEvent { + let mut minting_blob = vec![0u8; 76]; + minting_blob[0..8].copy_from_slice(&number.to_le_bytes()); + MintBlockEvent::new( + starcoin_crypto::HashValue::random(), + ConsensusStrategy::Dummy, + minting_blob, + U256::from(1u64), + number, + None, + ) +} + +async fn wait_for_server_ready( + child: &mut Child, + addr: SocketAddr, + timeout: Duration, +) -> Result<()> { + let start = Instant::now(); + loop { + if let Ok(stream) = TcpStream::connect(addr).await { + drop(stream); + return Ok(()); + } + if let Some(status) = child.try_wait()? { + return Err(anyhow::anyhow!( + "stratumd exited before ready, status: {}", + status + )); + } + if start.elapsed() > timeout { + return Err(anyhow::anyhow!("wait stratumd ready timeout: {}", addr)); + } + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + +async fn wait_submit_count(mock: &MockMiningRpc, expected: usize, timeout: Duration) -> Result<()> { + let start = Instant::now(); + loop { + if mock.submit_count()? >= expected { + return Ok(()); + } + if start.elapsed() >= timeout { + return Err(anyhow::anyhow!( + "submit count timeout, expected >= {}, got {}", + expected, + mock.submit_count()? + )); + } + tokio::time::sleep(Duration::from_millis(20)).await; + } +} + +#[stest::test] +#[ignore = "requires external standalone stratumd binary via STRATUMD_BIN"] +async fn test_miner_client_stratum_compat() -> Result<()> { + let _guard = TEST_MUTEX.lock().await; + + let mock = MockMiningRpc::start(build_mint_event(1))?; + let listen = SocketAddr::new(Ipv4Addr::LOCALHOST.into(), pick_free_port()?); + let _stratumd = StratumdProcess::spawn(listen, &mock.ws_url()).await?; + + let registry = RegistryService::launch(); + let client_config = MinerClientConfig { + server: Some(listen.to_string()), + plugin_path: None, + miner_thread: 1, + enable_stderr: false, + }; + registry.put_shared(client_config).await?; + + let stratum_cli_srv = registry + .register_by_factory::() + .await?; + let time_srv = Arc::new(RealTimeService::new()); + let login = LoginRequest { + login: "test".into(), + pass: "x".into(), + agent: "test-client".into(), + algo: None, + }; + let job_client = StratumJobClient::new(stratum_cli_srv, time_srv, login); + + let mut jobs = job_client.subscribe().await?; + let job = tokio::time::timeout(Duration::from_secs(5), jobs.next()) + .await + .map_err(|_| anyhow::anyhow!("job notification timeout"))? + .ok_or_else(|| anyhow::anyhow!("job stream closed"))?; + let extra = job + .extra + .clone() + .ok_or_else(|| anyhow::anyhow!("missing mint extra"))?; + + let seal = SealEvent { + minting_blob: job.minting_blob.clone(), + nonce: 0, + extra: Some(extra), + hash_result: "00".into(), + }; + job_client.submit_seal(seal).await?; + wait_submit_count(&mock, 1, Duration::from_secs(3)).await?; + + let _ = registry.shutdown_system().await; + Ok(()) +} diff --git a/code_layout.md b/code_layout.md index 639cde5c5a..a961f7fece 100644 --- a/code_layout.md +++ b/code_layout.md @@ -33,7 +33,7 @@ All codes directories are listed here. (in alphabetical order) | [scripts](scripts) | scripts helping for building, testing etc., | | | [state](state) | maintain the chain state | | | [storage](storage) | backend storage of chain | | -| [stratum](stratum) | stratum mining | [Stratum Mining Protocol](stratum/stratum_mining_protocol.md) | +| [cmd/stratumd](cmd/stratumd) | standalone stratum gateway | [Stratum Mining Protocol](cmd/stratumd/stratum_mining_protocol.md) | | [sync](sync) | sync from network | | | [test-helper](test-helper) | test helper functions | | | [testsuite](testsuite) | codes & scripts for setting up test | | diff --git a/config/example/barnard/config.toml b/config/example/barnard/config.toml index a0707a1966..32bdffa622 100644 --- a/config/example/barnard/config.toml +++ b/config/example/barnard/config.toml @@ -37,9 +37,6 @@ apis = "chain,miner,node,pubsub,state,txpool,contract" [storage] max_open_files = 40960 -[stratum] -port = 8090 - [sync] [tx_pool] diff --git a/config/example/main/config.toml b/config/example/main/config.toml index a0707a1966..32bdffa622 100644 --- a/config/example/main/config.toml +++ b/config/example/main/config.toml @@ -37,9 +37,6 @@ apis = "chain,miner,node,pubsub,state,txpool,contract" [storage] max_open_files = 40960 -[stratum] -port = 8090 - [sync] [tx_pool] diff --git a/config/example/proxima/config.toml b/config/example/proxima/config.toml index a0707a1966..32bdffa622 100644 --- a/config/example/proxima/config.toml +++ b/config/example/proxima/config.toml @@ -37,9 +37,6 @@ apis = "chain,miner,node,pubsub,state,txpool,contract" [storage] max_open_files = 40960 -[stratum] -port = 8090 - [sync] [tx_pool] diff --git a/config/src/lib.rs b/config/src/lib.rs index 239cbfd640..1cdd7a374c 100644 --- a/config/src/lib.rs +++ b/config/src/lib.rs @@ -37,7 +37,6 @@ mod miner_config; mod network_config; mod rpc_config; mod storage_config; -mod stratum_config; mod sync_config; #[cfg(test)] mod tests; @@ -48,7 +47,6 @@ use thiserror::Error; use crate::account_provider_config::AccountProviderConfig; use crate::genesis_config::vm2::GenesisConfig as GenesisConfig2; -use crate::stratum_config::StratumConfig; pub use api_config::{Api, ApiSet}; pub use api_quota::{ApiQuotaConfig, QuotaDuration}; pub use available_port::{ @@ -224,9 +222,6 @@ pub struct StarcoinOpt { pub sync: SyncConfig, #[clap(flatten)] pub vault: AccountVaultConfig, - #[serde(default)] - #[clap(flatten)] - pub stratum: StratumConfig, #[clap(flatten)] pub account_provider: AccountProviderConfig, } @@ -481,8 +476,9 @@ pub struct NodeConfig { pub metrics: MetricsConfig, #[serde(default)] pub logger: LoggerConfig, - #[serde(default)] - pub stratum: StratumConfig, + // Compatibility shim: keep accepting legacy `[stratum]` table from old config files. + #[serde(default, skip_serializing)] + pub stratum: Option, #[serde(default)] pub account_provider: AccountProviderConfig, } @@ -572,7 +568,6 @@ impl NodeConfig { self.vault.merge_with_opt(opt, base.clone())?; self.metrics.merge_with_opt(opt, base.clone())?; self.logger.merge_with_opt(opt, base.clone())?; - self.stratum.merge_with_opt(opt, base.clone())?; self.account_provider.merge_with_opt(opt, base)?; Ok(()) } diff --git a/config/src/stratum_config.rs b/config/src/stratum_config.rs deleted file mode 100644 index 59558a0ae3..0000000000 --- a/config/src/stratum_config.rs +++ /dev/null @@ -1,77 +0,0 @@ -use crate::{ - get_available_port_from, get_random_available_port, BaseConfig, ConfigModule, Parser, - StarcoinOpt, -}; -use anyhow::Result; -use serde::{Deserialize, Serialize}; -use starcoin_logger::prelude::*; -use std::net::{IpAddr, Ipv4Addr, SocketAddr}; -use std::sync::Arc; - -const DEFAULT_STRATUM_PORT: u16 = 9880; -// UNSPECIFIED is 0.0.0.0 -const DEFAULT_STRATUM_ADDRESS: IpAddr = IpAddr::V4(Ipv4Addr::UNSPECIFIED); - -#[derive(Debug, Default, Clone, PartialEq, Deserialize, Serialize, Parser)] -pub struct StratumConfig { - #[serde(skip)] - #[clap(name = "disable-stratum", long, help = "disable stratum")] - pub disable: bool, - - #[serde(skip_serializing_if = "Option::is_none")] - #[clap(name = "stratum-port", long)] - /// Default tcp port is 9880 - pub port: Option, - - #[serde(skip_serializing_if = "Option::is_none")] - #[clap(long = "stratum-address")] - /// Stratum address, default is 0.0.0.0 - pub address: Option, - - #[clap(skip)] - #[serde(skip)] - base: Option>, -} - -impl StratumConfig { - fn base(&self) -> &BaseConfig { - self.base.as_ref().expect("Config should init.") - } - pub fn get_address(&self) -> Option { - if self.disable { - return None; - } - let base = self.base(); - let address = self.address.unwrap_or(DEFAULT_STRATUM_ADDRESS).to_string(); - let port = self.port.unwrap_or_else(|| { - if base.net().is_test() { - get_random_available_port() - } else if base.net().is_dev() { - get_available_port_from(DEFAULT_STRATUM_PORT) - } else { - DEFAULT_STRATUM_PORT - } - }); - format!("{}:{}", address, port).parse::().ok() - } -} - -impl ConfigModule for StratumConfig { - fn merge_with_opt(&mut self, opt: &StarcoinOpt, base: Arc) -> Result<()> { - self.base = Some(base); - if opt.stratum.address.is_some() { - self.address = opt.rpc.rpc_address; - } - if opt.stratum.disable { - self.disable = true; - } - if opt.stratum.port.is_some() { - self.port = opt.stratum.port; - } - info!( - "Stratum listen address: {:?}, port:{:?}", - self.address, self.port - ); - Ok(()) - } -} diff --git a/config/src/tests.rs b/config/src/tests.rs index 5f3120774e..e96ec3e572 100644 --- a/config/src/tests.rs +++ b/config/src/tests.rs @@ -131,11 +131,6 @@ fn test_example_config_compact() -> Result<()> { //Vault "--vault-dir", "/data/my_starcoin_vault", - //Stratum - "--stratum-port", - "8090", - "--stratum-address", - "127.0.0.1", ]; let opt = StarcoinOpt::try_parse_from(args)?; let config = NodeConfig::load_with_opt(&opt)?; diff --git a/node/Cargo.toml b/node/Cargo.toml index b04bcbca21..04c0831ec5 100644 --- a/node/Cargo.toml +++ b/node/Cargo.toml @@ -36,7 +36,6 @@ starcoin-state-api = { workspace = true } starcoin-state-service = { workspace = true } starcoin-statedb = { workspace = true } starcoin-storage = { workspace = true } -starcoin-stratum = { workspace = true } starcoin-sync = { workspace = true } starcoin-sync-api = { workspace = true } starcoin-txpool = { workspace = true } diff --git a/node/src/node.rs b/node/src/node.rs index 8e310c3d1e..8a47bdba0d 100644 --- a/node/src/node.rs +++ b/node/src/node.rs @@ -44,8 +44,6 @@ use starcoin_storage::{ errors::StorageInitError, metrics::StorageMetrics, storage::StorageInstance, BlockStore, Storage, Storage2, }; -use starcoin_stratum::service::{StratumService, StratumServiceFactory}; -use starcoin_stratum::stratum::{Stratum, StratumFactory}; use starcoin_sync::announcement::AnnouncementService; use starcoin_sync::block_connector::{BlockConnectorService, ExecuteService, ResetRequest}; use starcoin_sync::sync::SyncService; @@ -390,6 +388,7 @@ impl NodeService { info!("Self peer_id is: {}", peer_id.to_base58()); info!("Self address is: {}", config.network.self_address()); + warn!("In-process stratum is removed from node; use standalone `starcoin_stratumd`."); // Create NewHeaderChannel for miner services use starcoin_miner::{NewHeaderChannel, NewHeaderService}; @@ -411,10 +410,6 @@ impl NodeService { info!("Config.miner.enable_miner_client is false, No in process MinerClient."); } - registry - .register_by_factory::() - .await?; - registry.register::().await?; // wait for service init. @@ -428,10 +423,6 @@ impl NodeService { registry .register_by_factory::() .await?; - registry - .register_by_factory::() - .await?; - // start metrics server if !config.metrics.disable_metrics() { registry.register::().await?; diff --git a/rpc/client/src/async_client.rs b/rpc/client/src/async_client.rs index c056a12630..0aad55be34 100644 --- a/rpc/client/src/async_client.rs +++ b/rpc/client/src/async_client.rs @@ -19,9 +19,12 @@ use parking_lot::Mutex; // Starcoin crates use starcoin_crypto::HashValue; use starcoin_rpc_api::{ - chain::GetBlockOption, + chain::{GetBlockOption, GetEventOption}, node::NodeInfo, - types::{BlockView, MintedBlockView, MultiStateView}, + types::{ + BlockView, ChainInfoView, MintedBlockView, MultiStateView, TransactionEventResponse, + TransactionInfoView as TransactionInfoViewRpc, + }, }; use starcoin_types::system_events::MintBlockEvent; use starcoin_vm2_account_api::AccountInfo; @@ -227,6 +230,32 @@ impl AsyncRpcClient { .await .map_err(map_err) } + + pub async fn chain_info(&self) -> anyhow::Result { + self.call_rpc_async(|inner| inner.chain_client.info()) + .await + .map_err(map_err) + } + + pub async fn chain_get_block_txn_infos( + &self, + block_hash: HashValue, + ) -> anyhow::Result> { + self.call_rpc_async(|inner| inner.chain_client.get_block_txn_infos(block_hash)) + .await + .map_err(map_err) + } + + pub async fn chain_get_events_by_txn_hash( + &self, + txn_hash: HashValue, + option: Option, + ) -> anyhow::Result> { + self.call_rpc_async(|inner| inner.chain_client.get_events_by_txn_hash(txn_hash, option)) + .await + .map_err(map_err) + } + pub async fn node_info(&self) -> anyhow::Result { self.call_rpc_async(|inner| inner.node_client.info()) .await @@ -244,6 +273,12 @@ impl AsyncRpcClient { .map_err(map_err) } + pub async fn miner_get_job(&self) -> anyhow::Result> { + self.call_rpc_async(|inner| inner.miner_client.get_job()) + .await + .map_err(map_err) + } + pub async fn subscribe_new_mint_blocks( &self, ) -> anyhow::Result> { diff --git a/rpc/client/src/lib.rs b/rpc/client/src/lib.rs index 4c0980ca2c..bba8dce0db 100644 --- a/rpc/client/src/lib.rs +++ b/rpc/client/src/lib.rs @@ -1029,6 +1029,11 @@ impl RpcClient { .map_err(map_err) } + pub fn miner_get_job(&self) -> anyhow::Result> { + self.call_rpc_blocking(|inner| inner.miner_client.get_job()) + .map_err(map_err) + } + pub fn txpool_status(&self) -> anyhow::Result { self.call_rpc_blocking(|inner| inner.txpool_client.state()) .map_err(map_err) diff --git a/simnet/src/lib.rs b/simnet/src/lib.rs index 9c661609ad..f153a83127 100644 --- a/simnet/src/lib.rs +++ b/simnet/src/lib.rs @@ -2,7 +2,7 @@ // SPDX-License-Identifier: Apache-2.0 use rand_chacha::{ - rand_core::{RngCore, SeedableRng}, + rand_core::{Rng, SeedableRng}, ChaCha8Rng, }; use std::cmp::Ordering; @@ -182,12 +182,12 @@ impl BlkStream { } } -fn exp_sample(rng: &mut R, mean_interval: u64) -> u64 { +fn exp_sample(rng: &mut R, mean_interval: u64) -> u64 { let u = next_unit_interval(rng); (-u.ln() * mean_interval as f64) as u64 } -fn next_unit_interval(rng: &mut R) -> f64 { +fn next_unit_interval(rng: &mut R) -> f64 { // Map u64 samples to (0, 1]; clamp to avoid ln(0) during exponential sampling. const SCALE: f64 = 1.0 / (u64::MAX as f64 + 1.0); (rng.next_u64() as f64 * SCALE).clamp(f64::EPSILON, 1.0) diff --git a/stratum/Cargo.toml b/stratum/Cargo.toml deleted file mode 100644 index a308872a6f..0000000000 --- a/stratum/Cargo.toml +++ /dev/null @@ -1,32 +0,0 @@ -[dependencies] -anyhow = { workspace = true } -byteorder = { workspace = true } -bytes = { workspace = true } -futures = { workspace = true } -hex = { workspace = true } -serde = { workspace = true } -serde_json = { features = ["arbitrary_precision"], workspace = true } -tokio = { features = ["full"], workspace = true } -tokio-util = { workspace = true } -starcoin-config = { workspace = true } -starcoin-crypto = { workspace = true } -starcoin-logger = { workspace = true } -starcoin-miner = { workspace = true } -starcoin-service-registry = { workspace = true } -starcoin-types = { workspace = true } - -[dev-dependencies] -stest = { workspace = true } -starcoin-vm2-vm-types = { workspace = true } -once_cell = { workspace = true } - -[package] -authors = { workspace = true } -edition = { workspace = true } -name = "starcoin-stratum" -version = { workspace = true } -homepage = { workspace = true } -license = { workspace = true } -publish = { workspace = true } -repository = { workspace = true } -rust-version = { workspace = true } diff --git a/stratum/src/diff_manager.rs b/stratum/src/diff_manager.rs deleted file mode 100644 index c1b966d98a..0000000000 --- a/stratum/src/diff_manager.rs +++ /dev/null @@ -1,73 +0,0 @@ -use crate::difficulty_to_target_hex; -use starcoin_logger::prelude::*; -use starcoin_types::U256; -pub const SHARE_SUBMIT_PERIOD: u64 = 2; -pub const INIT_HASH_RATE: u64 = 10000; -pub const MINI_UPDATE_PERIOD: u64 = 20; - -pub struct DifficultyManager { - pub timestamp_since_last_update: u64, - pub submits_since_last_update: u32, - pub hash_rate: u64, - pub difficulty: U256, -} -impl Default for DifficultyManager { - fn default() -> Self { - Self::new() - } -} -impl DifficultyManager { - pub fn get_target(&self) -> String { - difficulty_to_target_hex(self.difficulty) - } - - pub fn new() -> Self { - Self { - timestamp_since_last_update: Self::current_timestamp(), - submits_since_last_update: 0, - hash_rate: INIT_HASH_RATE, - difficulty: Self::get_difficulty_from_hashrate(INIT_HASH_RATE, SHARE_SUBMIT_PERIOD), - } - } - - pub fn find_seal(&mut self) { - self.submits_since_last_update += 1; - } - - pub fn try_update(&mut self, worker: String) -> bool { - self.find_seal(); - let current_timestamp = Self::current_timestamp(); - - let pass_time = current_timestamp - self.timestamp_since_last_update; - if pass_time < MINI_UPDATE_PERIOD { - return false; - } - - if self.submits_since_last_update == 0 { - self.hash_rate /= 2 - } else { - // hash_rate = difficulty / avg_time = difficulty / (pass_time / submits_of_share) - self.hash_rate = - (self.difficulty * self.submits_since_last_update / pass_time).as_u64(); - } - info!("Miner:{} hash rate is:{}", worker, self.hash_rate); - self.timestamp_since_last_update = current_timestamp; - self.difficulty = Self::get_difficulty_from_hashrate(self.hash_rate, SHARE_SUBMIT_PERIOD); - self.submits_since_last_update = 0; - true - } - - fn get_difficulty_from_hashrate(hash_rate: u64, share_submit_period: u64) -> U256 { - if hash_rate == 0 { - return 1.into(); - } - U256::from(hash_rate * share_submit_period) - } - - fn current_timestamp() -> u64 { - std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .expect("time went backwards") - .as_secs() - } -} diff --git a/stratum/src/lib.rs b/stratum/src/lib.rs deleted file mode 100644 index 169b2858ba..0000000000 --- a/stratum/src/lib.rs +++ /dev/null @@ -1,33 +0,0 @@ -use starcoin_types::U256; - -pub mod codec; -pub mod diff_manager; -pub mod rpc; -pub mod service; -pub mod stratum; -pub use anyhow::Result; - -pub fn difficulty_to_target_hex(difficulty: U256) -> String { - let target = format!("{:x}", U256::from(u64::MAX) / difficulty); - let mut temp = "0".repeat(16 - target.len()); - temp.push_str(&target); - let mut t = hex::decode(temp).expect("Decode target never failed"); - t.reverse(); - hex::encode(&t) -} - -pub fn target_hex_to_difficulty(target: &str) -> Result { - let mut temp = hex::decode(target)?; - temp.reverse(); - let temp = hex::encode(temp); - let temp = U256::from_str_radix(&temp, 16)?; - Ok(U256::from(u64::MAX) / temp) -} - -#[test] -fn test() { - let target = difficulty_to_target_hex(U256::from(35652289346123_u64)); - println!("{}", target); - let diff = target_hex_to_difficulty(&target).unwrap(); - println!("{}", diff); -} diff --git a/stratum/src/rpc.rs b/stratum/src/rpc.rs deleted file mode 100644 index bdeb4034ab..0000000000 --- a/stratum/src/rpc.rs +++ /dev/null @@ -1,212 +0,0 @@ -use crate::diff_manager::DifficultyManager; -use byteorder::{ByteOrder, LittleEndian, WriteBytesExt}; -use futures::channel::mpsc; -use serde::{Deserialize, Serialize}; -use starcoin_crypto::hash::DefaultHasher; -use starcoin_miner::SubmitSealRequest as MinerSubmitSealRequest; -use starcoin_service_registry::ServiceRequest; -use starcoin_types::block::BlockHeaderExtra; -use starcoin_types::system_events::MintBlockEvent; -use std::convert::TryInto; -use std::sync::{Arc, RwLock}; - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct ShareRequest { - pub id: String, - pub job_id: String, - pub nonce: String, - pub result: String, -} - -impl TryInto for ShareRequest { - type Error = anyhow::Error; - fn try_into(self) -> anyhow::Result { - let nonce_temp = u32::from_str_radix(self.nonce.as_str(), 16)?; - let mut n = Vec::new(); - let _ = n.write_u32::(nonce_temp); - let nonce = byteorder::BigEndian::read_u32(&n); - let extra = hex::decode(self.id)?; - let extra: [u8; 4] = extra - .try_into() - .map_err(|_| anyhow::anyhow!("Failed to parse extra"))?; - Ok(MinerSubmitSealRequest { - nonce, - extra: BlockHeaderExtra::new(extra), - minting_blob: vec![], - }) - } -} - -#[derive(Clone, Debug, Serialize, Deserialize)] -pub struct SubmitResult { - pub result: Status, -} - -#[derive(Debug, PartialEq, Eq, Clone, Deserialize, Serialize)] -pub struct KeepalivedResult { - pub result: Status, -} - -#[derive(Debug, PartialEq, Eq, Clone, Deserialize, Serialize)] -pub struct Status { - pub status: String, -} - -#[derive(Debug, Clone)] -pub struct SubscribeJobEvent(pub LoginRequest); - -impl ServiceRequest for SubscribeJobEvent { - type Response = anyhow::Result>; -} - -#[derive(Debug, Clone)] -pub struct SubmitShareEvent(pub ShareRequest); - -impl ServiceRequest for SubmitShareEvent { - type Response = anyhow::Result<()>; -} - -#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)] -pub struct LoginRequest { - pub login: String, - pub pass: String, - pub agent: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub algo: Option>, -} - -impl ServiceRequest for LoginRequest { - type Response = - futures::channel::oneshot::Receiver>; -} - -#[derive(Debug, PartialEq, Eq, Hash, Clone, Copy)] -pub struct WorkerId { - buff: [u8; 4], -} -impl WorkerId { - pub fn from_hex(input: String) -> anyhow::Result { - let worker_id: [u8; 4] = hex::decode(input) - .map_err(|_| anyhow::anyhow!("Decode worker id failed"))? - .try_into() - .map_err(|_| anyhow::anyhow!("Invalid length of worker id"))?; - Ok(WorkerId { buff: worker_id }) - } - pub fn to_hex(&self) -> String { - hex::encode(self.buff) - } -} -pub struct MinerWorker { - pub base_info: LoginRequest, - pub sub_id: u32, - pub worker_id: WorkerId, - pub diff_manager: Arc>, -} -impl MinerWorker { - fn generate_worker_id(login_name: String, sub_id: u32) -> WorkerId { - let mut hash = DefaultHasher::new(b""); - hash.update(login_name.as_bytes()); - let mut output: [u8; 4] = hash.finish().to_vec()[0..4] - .try_into() - .expect("Hash len should have 8 bytes"); - output - .iter_mut() - .zip(u32::to_le_bytes(sub_id).iter()) - .for_each(|(x1, x2)| *x1 ^= *x2); - WorkerId { buff: output } - } - - pub fn new(sub_id: u32, base_info: LoginRequest) -> Self { - let worker_id = Self::generate_worker_id(base_info.login.clone(), sub_id); - let diff_manager = Arc::new(RwLock::new(DifficultyManager::new())); - Self { - base_info, - sub_id, - worker_id, - diff_manager, - } - } - pub fn diff_manager(&self) -> Arc> { - self.diff_manager.clone() - } -} - -#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)] -pub struct StratumJobResponse { - #[serde(skip_serializing_if = "Option::is_none")] - pub login: Option, - pub id: String, - pub status: String, - pub job: StratumJob, -} - -#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)] -pub struct StratumJob { - pub height: u64, - pub id: String, - pub target: String, - pub job_id: String, - pub blob: String, -} - -impl StratumJob { - pub fn get_extra(&self) -> anyhow::Result { - let blob = hex::decode(&self.blob)?; - if blob.len() != 76 { - return Err(anyhow::anyhow!("Invalid stratum job")); - } - let extra: [u8; 4] = blob[35..39].try_into()?; - - Ok(BlockHeaderExtra::new(extra)) - } -} -#[derive(Debug, PartialEq, Eq)] -pub struct JobId { - pub job_id: [u8; 8], -} -impl JobId { - pub fn from_bob(minting_bob: &[u8]) -> JobId { - let mut job_id = [0u8; 8]; - job_id.copy_from_slice(&minting_bob[0..8]); - Self { job_id } - } - pub fn encode(&self) -> String { - hex::encode(self.job_id) - } - pub fn equal_with(&self, minting_bob: &[u8]) -> bool { - self.job_id[..] == minting_bob[0..8] - } - pub fn new(job_id: &String) -> anyhow::Result { - let job_id: [u8; 8] = hex::decode(job_id) - .map_err(|_| anyhow::anyhow!("Decode job_id failed"))? - .try_into() - .map_err(|_| anyhow::anyhow!("Invalid job id with bad length"))?; - Ok(Self { job_id }) - } -} - -impl StratumJobResponse { - pub fn from( - e: &MintBlockEvent, - login: Option, - worker_id: WorkerId, - target: String, - ) -> Self { - let mut minting_blob = e.minting_blob.clone(); - minting_blob[35..39].copy_from_slice(&worker_id.buff); - - let job_id = JobId::from_bob(&e.minting_blob).encode(); - Self { - login, - id: worker_id.to_hex(), - status: "OK".into(), - job: StratumJob { - height: 0, - id: worker_id.to_hex(), - target, - job_id, - blob: hex::encode(&minting_blob), - }, - } - } -} diff --git a/stratum/src/service.rs b/stratum/src/service.rs deleted file mode 100644 index 5d8eb79c2d..0000000000 --- a/stratum/src/service.rs +++ /dev/null @@ -1,377 +0,0 @@ -use crate::codec::JsonStreamCodec; -use crate::rpc::{LoginRequest, ShareRequest, Status, SubmitShareEvent, SubscribeJobEvent}; -use crate::stratum::Stratum; -use anyhow::Result; -use futures::{SinkExt, StreamExt}; -use serde::de::DeserializeOwned; -use serde::{Deserialize, Serialize}; -use starcoin_config::NodeConfig; -use starcoin_logger::prelude::*; -use starcoin_service_registry::{ActorService, ServiceContext, ServiceFactory, ServiceRef}; -use std::sync::Arc; -use std::thread::JoinHandle; -use tokio::net::{TcpListener, TcpStream}; -use tokio::runtime::Runtime; -use tokio::sync::oneshot; -use tokio_util::codec::Framed; - -pub struct StratumService { - config: Arc, - shutdown_tx: Option>, - join_handle: Option>, -} - -impl ActorService for StratumService { - fn started(&mut self, ctx: &mut ServiceContext) -> Result<()> { - if let Some(address) = self.config.stratum.get_address() { - let stratum = ctx.service_ref::()?.clone(); - let (shutdown_tx, shutdown_rx) = oneshot::channel(); - let join_handle = std::thread::spawn(move || { - let runtime = Runtime::new().expect("create stratum tokio runtime"); - let result = runtime.block_on(run_stratum_server(address, stratum, shutdown_rx)); - if let Err(err) = result { - error!(target: "stratum", "stratum server stopped with error: {}", err); - } - }); - self.shutdown_tx = Some(shutdown_tx); - self.join_handle = Some(join_handle); - info!(target: "stratum", "Stratum tcp server start at: {}", address); - } - Ok(()) - } - - fn stopped(&mut self, _ctx: &mut ServiceContext) -> Result<()> { - if let Some(shutdown_tx) = self.shutdown_tx.take() { - let _ = shutdown_tx.send(()); - } - if let Some(join_handle) = self.join_handle.take() { - let _ = join_handle.join(); - } - Ok(()) - } -} - -pub struct StratumServiceFactory; - -impl ServiceFactory for StratumServiceFactory { - fn create(ctx: &mut ServiceContext) -> Result { - let config = ctx.get_shared::>()?; - Ok(StratumService { - config, - shutdown_tx: None, - join_handle: None, - }) - } -} - -#[derive(Debug, Deserialize)] -struct JsonRpcRequest { - #[allow(dead_code)] - #[serde(default)] - jsonrpc: Option, - #[serde(default)] - id: Option, - method: String, - #[serde(default)] - params: serde_json::Value, -} - -#[derive(Debug, Deserialize)] -#[serde(untagged)] -enum JsonRpcId { - Number(u64), - String(String), -} - -#[derive(Debug, Serialize)] -struct JsonRpcOutput { - #[serde(skip_serializing_if = "Option::is_none")] - jsonrpc: Option<&'static str>, - result: T, - id: u32, - error: Option, -} - -#[derive(Debug, Serialize)] -struct JsonRpcFailure { - #[serde(skip_serializing_if = "Option::is_none")] - jsonrpc: Option<&'static str>, - id: u32, - error: JsonRpcError, -} - -#[derive(Debug, Serialize)] -struct JsonRpcError { - code: i32, - message: String, -} - -#[derive(Debug, Serialize)] -struct JsonRpcNotification { - #[serde(skip_serializing_if = "Option::is_none")] - jsonrpc: Option<&'static str>, - method: &'static str, - params: T, -} - -async fn run_stratum_server( - address: std::net::SocketAddr, - stratum: ServiceRef, - mut shutdown_rx: oneshot::Receiver<()>, -) -> Result<()> { - let listener = TcpListener::bind(address).await?; - loop { - tokio::select! { - _ = &mut shutdown_rx => { - break; - } - accept_result = listener.accept() => { - match accept_result { - Ok((stream, peer_addr)) => { - info!(target: "stratum", "stratum client connected: {}", peer_addr); - tokio::spawn(handle_connection(stream, stratum.clone())); - } - Err(err) => { - error!(target: "stratum", "accept connection failed: {}", err); - } - } - } - } - } - Ok(()) -} - -async fn handle_connection(stream: TcpStream, stratum: ServiceRef) { - let framed = Framed::new(stream, JsonStreamCodec::stream_incoming()); - let (mut sink, mut stream) = framed.split(); - let (out_tx, mut out_rx) = futures::channel::mpsc::unbounded::(); - - let writer = tokio::spawn(async move { - while let Some(msg) = out_rx.next().await { - if sink.send(msg).await.is_err() { - break; - } - } - }); - - while let Some(item) = stream.next().await { - let line = match item { - Ok(line) => line, - Err(err) => { - debug!(target: "stratum", "stratum read error: {}", err); - break; - } - }; - if line.trim().is_empty() { - continue; - } - let request: JsonRpcRequest = match serde_json::from_str(&line) { - Ok(req) => req, - Err(err) => { - debug!(target: "stratum", "invalid jsonrpc request: {}", err); - continue; - } - }; - let request_id = parse_request_id(request.id); - match request.method.as_str() { - "login" => { - if let Err(err) = handle_login(request_id, request.params, &stratum, &out_tx).await - { - debug!(target: "stratum", "handle login failed: {}", err); - } - } - "submit" => { - if let Err(err) = handle_submit(request_id, request.params, &stratum, &out_tx).await - { - debug!(target: "stratum", "handle submit failed: {}", err); - } - } - "keepalived" => { - if let Some(id) = request_id { - let status = Status { - status: "KEEPALIVED".to_string(), - }; - if let Err(err) = send_output(&out_tx, id, status) { - debug!(target: "stratum", "send keepalived response failed: {}", err); - } - } - } - "logout" => { - if let Err(err) = handle_logout(request_id, request.params, &out_tx).await { - debug!(target: "stratum", "handle logout failed: {}", err); - } - } - other => { - if let Some(id) = request_id { - let _ = send_failure(&out_tx, id, -1, format!("unknown method {}", other)); - } - } - } - } - - writer.abort(); -} - -async fn handle_login( - request_id: Option, - params: serde_json::Value, - stratum: &ServiceRef, - out_tx: &futures::channel::mpsc::UnboundedSender, -) -> Result<()> { - let login: LoginRequest = match parse_params(params) { - Ok(login) => login, - Err(err) => { - if let Some(id) = request_id { - let _ = send_failure(out_tx, id, -1, err.to_string()); - } - return Ok(()); - } - }; - let mut job_rx = match stratum.send(SubscribeJobEvent(login)).await { - Ok(Ok(rx)) => rx, - Ok(Err(err)) => { - if let Some(id) = request_id { - let _ = send_failure(out_tx, id, -1, err.to_string()); - } - return Ok(()); - } - Err(err) => { - if let Some(id) = request_id { - let _ = send_failure(out_tx, id, -1, err.to_string()); - } - return Ok(()); - } - }; - - if let Some(id) = request_id { - let mut first_job = match job_rx.next().await { - Some(job) => job, - None => { - let _ = send_failure(out_tx, id, -1, "no job".to_string()); - return Ok(()); - } - }; - first_job.login = None; - send_output(out_tx, id, first_job)?; - } - - let out_tx = out_tx.clone(); - tokio::spawn(async move { - while let Some(job_resp) = job_rx.next().await { - let notif = JsonRpcNotification { - jsonrpc: Some("2.0"), - method: "job", - params: job_resp.job, - }; - let msg = match serde_json::to_string(¬if) { - Ok(msg) => msg, - Err(err) => { - debug!(target: "stratum", "serialize job notification failed: {}", err); - continue; - } - }; - if out_tx.unbounded_send(msg).is_err() { - break; - } - } - }); - - Ok(()) -} - -async fn handle_submit( - request_id: Option, - params: serde_json::Value, - stratum: &ServiceRef, - out_tx: &futures::channel::mpsc::UnboundedSender, -) -> Result<()> { - let share: ShareRequest = match parse_params(params) { - Ok(share) => share, - Err(err) => { - if let Some(id) = request_id { - let _ = send_failure(out_tx, id, -1, err.to_string()); - } - return Ok(()); - } - }; - let submit_result = stratum.send(SubmitShareEvent(share)).await; - match submit_result { - Ok(Ok(())) => { - if let Some(id) = request_id { - let status = Status { - status: "OK".to_string(), - }; - let _ = send_output(out_tx, id, status); - } - } - Ok(Err(err)) => { - if let Some(id) = request_id { - let _ = send_failure(out_tx, id, -1, err.to_string()); - } - } - Err(err) => { - if let Some(id) = request_id { - let _ = send_failure(out_tx, id, -1, err.to_string()); - } - } - } - Ok(()) -} - -async fn handle_logout( - request_id: Option, - params: serde_json::Value, - out_tx: &futures::channel::mpsc::UnboundedSender, -) -> Result<()> { - info!(target: "stratum", "receive logout request params: {}", params); - if let Some(id) = request_id { - let _ = send_output(out_tx, id, false); - } - Ok(()) -} - -fn parse_request_id(id: Option) -> Option { - match id { - Some(JsonRpcId::Number(num)) => u32::try_from(num).ok(), - Some(JsonRpcId::String(s)) => s.parse::().ok(), - None => None, - } -} - -fn parse_params(params: serde_json::Value) -> Result { - serde_json::from_value(params).map_err(|err| anyhow::anyhow!("invalid params: {}", err)) -} - -fn send_output( - out_tx: &futures::channel::mpsc::UnboundedSender, - id: u32, - result: T, -) -> Result<()> { - let output = JsonRpcOutput { - jsonrpc: Some("2.0"), - result, - id, - error: None, - }; - let msg = serde_json::to_string(&output)?; - out_tx - .unbounded_send(msg) - .map_err(|_| anyhow::anyhow!("send response failed")) -} - -fn send_failure( - out_tx: &futures::channel::mpsc::UnboundedSender, - id: u32, - code: i32, - message: String, -) -> Result<()> { - let failure = JsonRpcFailure { - jsonrpc: Some("2.0"), - id, - error: JsonRpcError { code, message }, - }; - let msg = serde_json::to_string(&failure)?; - out_tx - .unbounded_send(msg) - .map_err(|_| anyhow::anyhow!("send response failed")) -} diff --git a/stratum/src/stratum.rs b/stratum/src/stratum.rs deleted file mode 100644 index 081ab2fee0..0000000000 --- a/stratum/src/stratum.rs +++ /dev/null @@ -1,170 +0,0 @@ -use crate::{rpc::*, target_hex_to_difficulty}; -use anyhow::Result; -use futures::channel::mpsc; -use starcoin_logger::prelude::*; -use starcoin_miner::{ - MinerService, SubmitSealRequest as MinerSubmitSealRequest, UpdateSubscriberNumRequest, -}; -use starcoin_service_registry::{ - ActorService, EventHandler, ServiceContext, ServiceFactory, ServiceHandler, ServiceRef, -}; -use starcoin_types::system_events::MintBlockEvent; -use std::collections::HashMap; -use std::convert::TryInto; -use std::sync::atomic; - -pub struct Stratum { - uid: atomic::AtomicU32, - mint_block_subscribers: - HashMap, MinerWorker)>, - miner_service: ServiceRef, -} - -impl Stratum { - fn new(miner_service: ServiceRef) -> Self { - Self { - miner_service, - uid: atomic::AtomicU32::new(1), - mint_block_subscribers: Default::default(), - } - } - - fn next_id(&self) -> u32 { - self.uid.fetch_add(1, atomic::Ordering::SeqCst) - } - - fn sync_upstream_job(&mut self) -> Result> { - let service = self.miner_service.clone(); - let subscribers_num = self.mint_block_subscribers.len() as u32; - futures::executor::block_on(service.send(UpdateSubscriberNumRequest { - number: Some(subscribers_num), - })) - } - - fn get_downstream_job( - miner: &MinerWorker, - set_login: bool, - upstreaum_event: &MintBlockEvent, - ) -> StratumJobResponse { - let login = miner.base_info.clone(); - - let target = miner.diff_manager.read().unwrap().get_target(); - info!( - "set downstream job diff:{:?}", - target_hex_to_difficulty(&target).unwrap() - ); - StratumJobResponse::from( - upstreaum_event, - if set_login { Some(login) } else { None }, - miner.worker_id, - target, - ) - } - - fn dispatch_job_to_clients(&mut self, event: MintBlockEvent) { - let mut remove_outdated = vec![]; - for (id, (ch, worker)) in self.mint_block_subscribers.iter() { - let job = Self::get_downstream_job(worker, false, &event); - info!(target: "stratum", "dispatch startum job:{:?}", job); - if let Err(err) = ch.unbounded_send(job) { - if err.is_disconnected() { - warn!("stratum disconnect worker:{:?}", err); - remove_outdated.push(*id); - } else if err.is_full() { - error!(target: "stratum", "subscription {:?} fail to new messages, channel is full", id); - } - } - } - for id in remove_outdated { - self.mint_block_subscribers.remove(&id); - } - } -} - -impl ActorService for Stratum { - fn started(&mut self, ctx: &mut ServiceContext) -> Result<()> { - ctx.set_mailbox_capacity(1024); - ctx.subscribe::(); - Ok(()) - } - - fn stopped(&mut self, ctx: &mut ServiceContext) -> Result<()> { - ctx.unsubscribe::(); - Ok(()) - } -} - -impl EventHandler for Stratum { - fn handle_event(&mut self, event: MintBlockEvent, _ctx: &mut ServiceContext) { - self.dispatch_job_to_clients(event); - } -} - -impl ServiceHandler for Stratum { - fn handle( - &mut self, - msg: SubscribeJobEvent, - _ctx: &mut ServiceContext, - ) -> anyhow::Result> { - let SubscribeJobEvent(login) = msg; - let (sender, receiver) = mpsc::unbounded(); - let sub_id = self.next_id(); - info!(target: "stratum", "receive subscribe event {:?},sub_id:{}", login, sub_id); - let miner_worker = MinerWorker::new(sub_id, login); - let worker_id = miner_worker.worker_id; - self.mint_block_subscribers - .insert(worker_id, (sender.clone(), miner_worker)); - let event = self.sync_upstream_job()?; - let downstream_job = event.as_ref().and_then(|event| { - self.mint_block_subscribers - .get(&worker_id) - .map(|(_, worker)| Self::get_downstream_job(worker, true, event)) - }); - if let Some(downstream_job) = downstream_job { - info!(target:"stratum", "Respond to stratum subscribe:{:?}", downstream_job); - if let Err(err) = sender.unbounded_send(downstream_job) { - error!(target: "stratum", "Failed to send MintBlockEvent: {}", err); - } - } else { - warn!(target: "stratum", "current mint job is empty"); - } - Ok(receiver) - } -} - -impl ServiceHandler for Stratum { - fn handle(&mut self, msg: SubmitShareEvent, _ctx: &mut ServiceContext) -> Result<()> { - info!(target: "stratum", "received submit share event:{:?}", &msg.0); - if let Some(current_mint_event) = self.sync_upstream_job()? { - let worker_id = WorkerId::from_hex(msg.0.id.clone())?; - if let Some((_job_sender, worker)) = self.mint_block_subscribers.get(&worker_id) { - let _updated_diff = worker - .diff_manager() - .write() - .unwrap() - .try_update(worker.base_info.login.clone()); - }; - let job_id = JobId::new(&msg.0.job_id)?; - let submit_job_id = JobId::from_bob(¤t_mint_event.minting_blob); - if job_id != submit_job_id { - warn!(target: "stratum", "received job mismatch with current job,{:?},{:?}",job_id, submit_job_id); - return Ok(()); - }; - - let mut seal: MinerSubmitSealRequest = msg.0.try_into()?; - - seal.minting_blob = current_mint_event.minting_blob; - self.miner_service.try_send(seal)?; - } - Ok(()) - } -} - -pub struct StratumFactory; - -impl ServiceFactory for StratumFactory { - fn create(ctx: &mut ServiceContext) -> Result { - let miner_service = ctx.service_ref::()?.clone(); - Ok(Stratum::new(miner_service)) - } -} diff --git a/stratum/stratum_mining_protocol.md b/stratum/stratum_mining_protocol.md deleted file mode 100644 index 47be40f456..0000000000 --- a/stratum/stratum_mining_protocol.md +++ /dev/null @@ -1,134 +0,0 @@ -# Stratum mining protocol -## login - -Miner send `login` request after connection successfully established for authorization on pool. - -#### Example request: -```json -{ - "id": 1, - "jsonrpc": "2.0", - "method": "login", - "params": { - "login": "48edfHu7V9Z84YzzMa6fUueoELZ9ZRXq9VetWzYGzKt52XU5xvqgzYnDK9URnRoJMk1j8nLwEVsaSWJ4fhdUyZijBGUicoD", - "pass": "x", - "agent": "Ibctminer/1.0.0" - } -} -``` - -#### Example success reply: -```json -{ - "id": 1, - "jsonrpc": "2.0", - "error": null, - "result": { - "id": "1be0b7b6-b15a-47be-a17d-46b2911cf7d0", // the id of the working miner - "job": { - "blob": "070780e6b9d60586ba419a0c224e3c6c3e134cc45c4fa04d8ee2d91c2595463c57eef0a4f0796c000000002fcc4d62fa6c77e76c30017c768be5c61d83ec9d3a085d524ba8053ecc3224660d", - "job_id": "q7PLUPL25UV0z5Ij14IyMk8htXbj", - "id": "1be0b7b6-b15a-47be-a17d-46b2911cf7d0", //the id of the working miner - "target": "b88d0600", - "height": 0 // height must always be 0 - }, - "status": "OK" - } -} -``` - -#### Example error reply: -```json -{ - "id": 1, - "jsonrpc": "2.0", - "error": { - "code": -1, - "message": "Invalid payment address provided" - } -} -``` - -## job -Pool send new job to miner. Miner should switch to new job as fast as possible. - -#### Example notification: -```json -{ - "jsonrpc": "2.0", - "method": "job", - "params": { - "blob": "0707d5efb9d6057e95a35f868231780b3a8649c4e57f3c77eaf437329243eef0b9f4b6987d05b900000000cae7754cb85a0ad8eebf3e0bf55f3ec5e754a1d6b05d46e5c358f907dbcbb72b01", - "job_id": "4BiGm3/RgGQzgkTI/xV0smdA+EGZ", - "target": "b88d0600", - "height": 0 - } -} -``` - -## submit -Miner send `submit` request after share was found. - -#### Example request: -```json -{ - "id": 2, - "jsonrpc": "2.0", - "method": "submit", - "params": { - "id": "1be0b7b6-b15a-47be-a17d-46b2911cf7d0", - "job_id": "4BiGm3/RgGQzgkTI/xV0smdA+EGZ", - "nonce": "d0030040", - "result": "e1364b8782719d7683e2ccd3d8f724bc59dfa780a9e960e7c0e0046acdb40100" - } -} -``` - -#### Example success reply: -```json -{ - "id": 2, - "jsonrpc": "2.0", - "error": null, - "result": { - "status": "OK" - } -} -``` - -#### Example error reply: -```json -{ - "id": 2, - "jsonrpc": "2.0", - "error": { - "code": -1, - "message": "Low difficulty share" - } -} -``` - -## keepalived -Miner send `keepalived` to prevent connection timeout. -#### Example request: -```json -{ - "id": 2, - "method": "keepalived", - "params": { - "id": "1be0b7b6-b15a-47be-a17d-46b2911cf7d0" //the id of the working miner - } -} -``` - -#### Example success reply: -```json -{ - "id": 2, - "jsonrpc": "2.0", - "error": null, - "result": { - "status": "KEEPALIVED" - } -} -``` diff --git a/stratum/tests/stratum_mining_protocol.rs b/stratum/tests/stratum_mining_protocol.rs deleted file mode 100644 index 453234f52b..0000000000 --- a/stratum/tests/stratum_mining_protocol.rs +++ /dev/null @@ -1,349 +0,0 @@ -use anyhow::Result; -use once_cell::sync::Lazy; -use serde_json::json; -use starcoin_config::{BuiltinNetworkID, NodeConfig, StarcoinOpt}; -use starcoin_crypto::HashValue; -use starcoin_miner::{DispatchMintBlockTemplate, MinerService}; -use starcoin_service_registry::{RegistryAsyncService, RegistryService}; -use starcoin_stratum::rpc::StratumJobResponse; -use starcoin_stratum::service::{StratumService, StratumServiceFactory}; -use starcoin_stratum::stratum::{Stratum, StratumFactory}; -use starcoin_types::block::{BlockBody, BlockTemplate}; -use starcoin_types::block_metadata::BlockMetadata; -use starcoin_types::genesis_config::{ChainId, ConsensusStrategy}; -use starcoin_types::U256; -use starcoin_vm2_vm_types::account_address::AccountAddress as Vm2AccountAddress; -use starcoin_vm2_vm_types::on_chain_resource::ChainId as Vm2ChainId; -use std::fs; -use std::io; -use std::net::{Ipv4Addr, SocketAddr}; -use std::sync::atomic::{AtomicUsize, Ordering}; -use std::sync::Arc; -use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; -use tokio::net::tcp::{OwnedReadHalf, OwnedWriteHalf}; -use tokio::net::TcpStream; -use tokio::sync::Mutex; -use tokio::time::{sleep, Duration, Instant}; - -static TEST_MUTEX: Lazy> = Lazy::new(|| Mutex::new(())); - -fn pick_free_port() -> std::io::Result { - let listener = std::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))?; - let port = listener.local_addr()?.port(); - Ok(port) -} - -fn prepare_config() -> Result> { - static TEST_DIR_COUNTER: AtomicUsize = AtomicUsize::new(0); - let suffix = TEST_DIR_COUNTER.fetch_add(1, Ordering::Relaxed); - let base_dir = std::env::temp_dir().join(format!( - "starcoin-stratum-test-{}-{}", - std::process::id(), - suffix, - )); - fs::create_dir_all(&base_dir)?; - let opt = StarcoinOpt { - net: Some(BuiltinNetworkID::Dev.into()), - base_data_dir: Some(base_dir), - ..StarcoinOpt::default() - }; - let mut config = NodeConfig::load_with_opt(&opt)?; - config.stratum.address = Some(Ipv4Addr::LOCALHOST.into()); - let port = match pick_free_port() { - Ok(port) => port, - Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => { - eprintln!("Skipping test: cannot bind local port in this environment."); - return Ok(None); - } - Err(err) => return Err(err.into()), - }; - config.stratum.port = Some(port); - let addr = config.stratum.get_address().expect("stratum address"); - Ok(Some((config, addr))) -} - -fn build_block_template(number: u64, timestamp: u64) -> BlockTemplate { - let parent_hash = HashValue::zero(); - let author = Vm2AccountAddress::from_hex_literal("0x1").expect("valid address"); - let chain_id_v2 = Vm2ChainId::test(); - let metadata = BlockMetadata::new( - parent_hash, - timestamp, - author, - 0, - number, - chain_id_v2, - 0, - vec![], - 0, - ); - BlockTemplate::new( - HashValue::zero(), - HashValue::zero(), - HashValue::zero(), - HashValue::zero(), - HashValue::zero(), - 0, - BlockBody::new_empty(), - ChainId::test(), - U256::from(1u64), - ConsensusStrategy::CryptoNight, - metadata, - 0, - HashValue::zero(), - vec![], - ) -} - -async fn connect_with_retry(addr: SocketAddr, timeout: Duration) -> Result { - let start = Instant::now(); - loop { - match TcpStream::connect(addr).await { - Ok(stream) => return Ok(stream), - Err(err) if err.kind() == io::ErrorKind::ConnectionRefused => { - if start.elapsed() >= timeout { - return Err(anyhow::anyhow!( - "connect timeout after {:?}: {}", - timeout, - err - )); - } - sleep(Duration::from_millis(50)).await; - } - Err(err) => return Err(err.into()), - } - } -} - -async fn read_json_line(reader: &mut BufReader) -> Result { - let mut line = String::new(); - let read = tokio::time::timeout(Duration::from_secs(5), reader.read_line(&mut line)) - .await - .map_err(|_| anyhow::anyhow!("response timeout"))??; - if read == 0 { - return Err(anyhow::anyhow!("connection closed")); - } - let value: serde_json::Value = serde_json::from_str(line.trim())?; - Ok(value) -} - -fn extract_job(value: &serde_json::Value) -> Option { - if let Some(result) = value.get("result") { - if result.get("job").is_some() { - return serde_json::from_value(result.clone()).ok(); - } - } - if value.get("method").and_then(|m| m.as_str()) == Some("job") { - if let Some(params) = value.get("params") { - if let Some(result) = params.get("result") { - if let Ok(job) = serde_json::from_value(result.clone()) { - return Some(job); - } - } - if let Ok(job) = serde_json::from_value(params.clone()) { - return Some(job); - } - } - } - None -} - -async fn wait_for_job(reader: &mut BufReader) -> Result { - let start = Instant::now(); - loop { - if start.elapsed() > Duration::from_secs(5) { - return Err(anyhow::anyhow!("job notification timeout")); - } - let value = read_json_line(reader).await?; - if let Some(job) = extract_job(&value) { - return Ok(job); - } - } -} - -async fn write_json_line(writer: &mut OwnedWriteHalf, value: serde_json::Value) -> Result<()> { - let payload = format!("{}\n", value); - writer.write_all(payload.as_bytes()).await?; - writer.flush().await?; - Ok(()) -} - -#[stest::test] -async fn test_login_request() -> Result<()> { - let _guard = TEST_MUTEX.lock().await; - let Some((config, addr)) = prepare_config()? else { - return Ok(()); - }; - - let registry = RegistryService::launch(); - registry.put_shared(Arc::new(config)).await?; - - let result = tokio::time::timeout(Duration::from_secs(20), async { - registry.register::().await?; - registry - .register_by_factory::() - .await?; - registry - .register_by_factory::() - .await?; - - let miner = registry.service_ref::().await?; - miner - .send(DispatchMintBlockTemplate { - block_template: build_block_template(1, 0), - }) - .await?; - - sleep(Duration::from_millis(100)).await; - - let stream = connect_with_retry(addr, Duration::from_secs(5)).await?; - let (reader, mut writer) = stream.into_split(); - let mut reader = BufReader::new(reader); - - let login_req = json!({ - "id": 1, - "jsonrpc": "2.0", - "method": "login", - "params": { - "login": "test", - "pass": "x", - "agent": "test-client" - } - }); - write_json_line(&mut writer, login_req).await?; - - let job = wait_for_job(&mut reader).await?; - assert_eq!(job.status, "OK"); - assert!(!job.id.is_empty()); - assert!(!job.job.job_id.is_empty()); - assert!(!job.job.blob.is_empty()); - - Ok::<(), anyhow::Error>(()) - }) - .await - .map_err(|_| anyhow::anyhow!("test timeout"))?; - - let _ = registry.shutdown_system().await; - result -} - -#[stest::test] -async fn test_submit_request() -> Result<()> { - let _guard = TEST_MUTEX.lock().await; - let Some((config, addr)) = prepare_config()? else { - return Ok(()); - }; - - let registry = RegistryService::launch(); - registry.put_shared(Arc::new(config)).await?; - - let result = tokio::time::timeout(Duration::from_secs(20), async { - registry.register::().await?; - registry - .register_by_factory::() - .await?; - registry - .register_by_factory::() - .await?; - - let miner = registry.service_ref::().await?; - miner - .send(DispatchMintBlockTemplate { - block_template: build_block_template(1, 0), - }) - .await?; - - sleep(Duration::from_millis(100)).await; - - let stream = connect_with_retry(addr, Duration::from_secs(5)).await?; - let (reader, mut writer) = stream.into_split(); - let mut reader = BufReader::new(reader); - - let login_req = json!({ - "id": 1, - "jsonrpc": "2.0", - "method": "login", - "params": { - "login": "test", - "pass": "x", - "agent": "test-client" - } - }); - write_json_line(&mut writer, login_req).await?; - let job = wait_for_job(&mut reader).await?; - - let submit_req = json!({ - "id": 2, - "jsonrpc": "2.0", - "method": "submit", - "params": { - "id": job.id, - "job_id": job.job.job_id, - "nonce": "00000000", - "result": "00" - } - }); - write_json_line(&mut writer, submit_req).await?; - - let response = read_json_line(&mut reader).await?; - let status = response["result"]["status"] - .as_str() - .ok_or_else(|| anyhow::anyhow!("missing submit status"))?; - assert_eq!(status, "OK"); - - Ok::<(), anyhow::Error>(()) - }) - .await - .map_err(|_| anyhow::anyhow!("test timeout"))?; - - let _ = registry.shutdown_system().await; - result -} - -#[stest::test] -async fn test_keepalived_request() -> Result<()> { - let _guard = TEST_MUTEX.lock().await; - let Some((config, addr)) = prepare_config()? else { - return Ok(()); - }; - - let registry = RegistryService::launch(); - registry.put_shared(Arc::new(config)).await?; - - let result = tokio::time::timeout(Duration::from_secs(20), async { - registry.register::().await?; - registry - .register_by_factory::() - .await?; - registry - .register_by_factory::() - .await?; - - let stream = connect_with_retry(addr, Duration::from_secs(5)).await?; - let (reader, mut writer) = stream.into_split(); - let mut reader = BufReader::new(reader); - - let keep_req = json!({ - "id": 2, - "jsonrpc": "2.0", - "method": "keepalived", - "params": { - "id": "test" - } - }); - write_json_line(&mut writer, keep_req).await?; - - let response = read_json_line(&mut reader).await?; - let status = response["result"]["status"] - .as_str() - .ok_or_else(|| anyhow::anyhow!("missing keepalived status"))?; - assert_eq!(status, "KEEPALIVED"); - - Ok::<(), anyhow::Error>(()) - }) - .await - .map_err(|_| anyhow::anyhow!("test timeout"))?; - - let _ = registry.shutdown_system().await; - result -} diff --git a/vm2/vm-runtime/src/starcoin_vm.rs b/vm2/vm-runtime/src/starcoin_vm.rs index cadc813a36..1884e9210b 100644 --- a/vm2/vm-runtime/src/starcoin_vm.rs +++ b/vm2/vm-runtime/src/starcoin_vm.rs @@ -242,7 +242,7 @@ impl StarcoinVM { bail!("failed to load the gas schedule when trying to print its info"); } Some(gs) => { - gs.info(_message); + gs.info("gas schedule from VMConfig"); } } }