Skip to content

Commit 2a98bea

Browse files
committed
IT adaptations for extensions configs in tProxy, JDC, and Pool
1 parent 44ee0e6 commit 2a98bea

14 files changed

Lines changed: 103 additions & 59 deletions

integration-tests/Cargo.lock

Lines changed: 17 additions & 17 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

integration-tests/lib/mod.rs

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -67,7 +67,11 @@ pub fn start_sniffer(
6767
(sniffer, listening_address)
6868
}
6969

70-
pub async fn start_pool(template_provider_address: Option<SocketAddr>) -> (PoolSv2, SocketAddr) {
70+
pub async fn start_pool(
71+
template_provider_address: Option<SocketAddr>,
72+
supported_extensions: Vec<u16>,
73+
required_extensions: Vec<u16>,
74+
) -> (PoolSv2, SocketAddr) {
7175
use pool_sv2::config::PoolConfig;
7276
let listening_address = get_available_address();
7377
let authority_public_key = Secp256k1PublicKey::try_from(
@@ -106,6 +110,8 @@ pub async fn start_pool(template_provider_address: Option<SocketAddr>) -> (PoolS
106110
SHARES_PER_MINUTE,
107111
share_batch_size,
108112
1,
113+
supported_extensions,
114+
required_extensions,
109115
);
110116
let pool = PoolSv2::new(config);
111117
let pool_clone = pool.clone();
@@ -130,6 +136,8 @@ pub fn start_template_provider(
130136
pub fn start_jdc(
131137
pool: &[(SocketAddr, SocketAddr)], // (pool_address, jds_address)
132138
tp_address: SocketAddr,
139+
supported_extensions: Vec<u16>,
140+
required_extensions: Vec<u16>,
133141
) -> (JobDeclaratorClient, SocketAddr) {
134142
use jd_client_sv2::config::{
135143
JobDeclaratorClientConfig, PoolConfig, ProtocolConfig, TPConfig, Upstream,
@@ -187,6 +195,8 @@ pub fn start_jdc(
187195
upstreams,
188196
jdc_signature,
189197
None,
198+
supported_extensions,
199+
required_extensions,
190200
);
191201
let ret = jd_client_sv2::JobDeclaratorClient::new(jd_client_proxy);
192202
let ret_clone = ret.clone();
@@ -247,6 +257,8 @@ pub fn start_jds(tp_rpc_connection: &ConnectParams) -> (JobDeclaratorServer, Soc
247257
pub async fn start_sv2_translator(
248258
upstream: SocketAddr,
249259
aggregate_channels: bool,
260+
supported_extensions: Vec<u16>,
261+
required_extensions: Vec<u16>,
250262
) -> (TranslatorSv2, SocketAddr) {
251263
let upstream_address = upstream.ip().to_string();
252264
let upstream_port = upstream.port();
@@ -284,6 +296,8 @@ pub async fn start_sv2_translator(
284296
downstream_extranonce2_size,
285297
"user_identity".to_string(),
286298
aggregate_channels,
299+
supported_extensions,
300+
required_extensions,
287301
);
288302
let translator_v2 = translator_sv2::TranslatorSv2::new(config);
289303
let clone_translator_v2 = translator_v2.clone();

integration-tests/lib/utils.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -334,6 +334,7 @@ pub fn into_static(m: AnyMessage<'_>) -> AnyMessage<'static> {
334334
TemplateDistribution::SubmitSolution(m.into_static()),
335335
),
336336
},
337+
AnyMessage::Extensions(extensions) => AnyMessage::Extensions(extensions.into_static()),
337338
}
338339
}
339340

integration-tests/tests/jd_integration.rs

Lines changed: 17 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -19,10 +19,10 @@ use stratum_apps::stratum_core::{
1919
async fn jds_should_not_panic_if_jdc_shutsdown() {
2020
start_tracing();
2121
let (tp, tp_addr) = start_template_provider(None, DifficultyLevel::Low);
22-
let (_pool, pool_addr) = start_pool(Some(tp_addr)).await;
22+
let (_pool, pool_addr) = start_pool(Some(tp_addr), vec![], vec![]).await;
2323
let (_jds, jds_addr) = start_jds(tp.rpc_info());
2424
let (sniffer_a, sniffer_addr_a) = start_sniffer("0", jds_addr, false, vec![], None);
25-
let (jdc, jdc_addr) = start_jdc(&[(pool_addr, sniffer_addr_a)], tp_addr);
25+
let (jdc, jdc_addr) = start_jdc(&[(pool_addr, sniffer_addr_a)], tp_addr, vec![], vec![]);
2626
sniffer_a
2727
.wait_for_message_type(MessageDirection::ToUpstream, MESSAGE_TYPE_SETUP_CONNECTION)
2828
.await;
@@ -36,7 +36,7 @@ async fn jds_should_not_panic_if_jdc_shutsdown() {
3636
tokio::time::sleep(tokio::time::Duration::from_millis(2000)).await;
3737
assert!(tokio::net::TcpListener::bind(jdc_addr).await.is_ok());
3838
let (sniffer, sniffer_addr) = start_sniffer("0", jds_addr, false, vec![], None);
39-
let (_jdc_1, _jdc_addr_1) = start_jdc(&[(pool_addr, sniffer_addr)], tp_addr);
39+
let (_jdc_1, _jdc_addr_1) = start_jdc(&[(pool_addr, sniffer_addr)], tp_addr, vec![], vec![]);
4040
sniffer
4141
.wait_for_message_type(MessageDirection::ToUpstream, MESSAGE_TYPE_SETUP_CONNECTION)
4242
.await;
@@ -50,13 +50,18 @@ async fn jds_should_not_panic_if_jdc_shutsdown() {
5050
async fn jdc_tp_success_setup() {
5151
start_tracing();
5252
let (tp, tp_addr) = start_template_provider(None, DifficultyLevel::Low);
53-
let (_pool, pool_addr) = start_pool(Some(tp_addr)).await;
53+
let (_pool, pool_addr) = start_pool(Some(tp_addr), vec![], vec![]).await;
5454
let (_jds, jds_addr) = start_jds(tp.rpc_info());
5555
let (tp_jdc_sniffer, tp_jdc_sniffer_addr) = start_sniffer("0", tp_addr, false, vec![], None);
56-
let (_jdc, jdc_addr) = start_jdc(&[(pool_addr, jds_addr)], tp_jdc_sniffer_addr);
56+
let (_jdc, jdc_addr) = start_jdc(
57+
&[(pool_addr, jds_addr)],
58+
tp_jdc_sniffer_addr,
59+
vec![],
60+
vec![],
61+
);
5762
// This is needed because jd-client waits for a downstream connection before it starts
5863
// exchanging messages with the Template Provider.
59-
start_sv2_translator(jdc_addr, false).await;
64+
start_sv2_translator(jdc_addr, false, vec![], vec![]).await;
6065
tp_jdc_sniffer
6166
.wait_for_message_type(MessageDirection::ToUpstream, MESSAGE_TYPE_SETUP_CONNECTION)
6267
.await;
@@ -78,7 +83,7 @@ async fn jds_receive_solution_while_processing_declared_job_test() {
7883
start_tracing();
7984
let (tp_1, tp_addr_1) = start_template_provider(None, DifficultyLevel::Low);
8085
let (tp_2, tp_addr_2) = start_template_provider(None, DifficultyLevel::Low);
81-
let (_pool, pool_addr) = start_pool(Some(tp_addr_1)).await;
86+
let (_pool, pool_addr) = start_pool(Some(tp_addr_1), vec![], vec![]).await;
8287
let (_jds, jds_addr) = start_jds(tp_1.rpc_info());
8388

8489
let prev_hash = U256::Owned(vec![
@@ -111,8 +116,8 @@ async fn jds_receive_solution_while_processing_declared_job_test() {
111116
vec![submit_solution_replace.into()],
112117
None,
113118
);
114-
let (_jdc, jdc_addr) = start_jdc(&[(pool_addr, sniffer_a_addr)], tp_addr_2);
115-
let (_translator, tproxy_addr) = start_sv2_translator(jdc_addr, false).await;
119+
let (_jdc, jdc_addr) = start_jdc(&[(pool_addr, sniffer_a_addr)], tp_addr_2, vec![], vec![]);
120+
let (_translator, tproxy_addr) = start_sv2_translator(jdc_addr, false, vec![], vec![]).await;
116121
let (_minerd_process, _minerd_addr) = start_minerd(tproxy_addr, None, None, false).await;
117122
assert!(tp_2.fund_wallet().is_ok());
118123
assert!(tp_2.create_mempool_transaction().is_ok());
@@ -170,7 +175,7 @@ async fn jds_wont_exit_upon_receiving_unexpected_txids_in_provide_missing_transa
170175
assert!(tp_2.fund_wallet().is_ok());
171176
assert!(tp_2.create_mempool_transaction().is_ok());
172177

173-
let (_pool, pool_addr) = start_pool(Some(tp_addr_1)).await;
178+
let (_pool, pool_addr) = start_pool(Some(tp_addr_1), vec![], vec![]).await;
174179
let (_jds, jds_addr) = start_jds(tp_1.rpc_info());
175180

176181
let provide_missing_transaction_success_replace = ReplaceMessage::new(
@@ -196,8 +201,8 @@ async fn jds_wont_exit_upon_receiving_unexpected_txids_in_provide_missing_transa
196201
None,
197202
);
198203

199-
let (_, jdc_addr_1) = start_jdc(&[(pool_addr, sniffer_addr)], tp_addr_2);
200-
let (_translator, tproxy_addr) = start_sv2_translator(jdc_addr_1, false).await;
204+
let (_, jdc_addr_1) = start_jdc(&[(pool_addr, sniffer_addr)], tp_addr_2, vec![], vec![]);
205+
let (_translator, tproxy_addr) = start_sv2_translator(jdc_addr_1, false, vec![], vec![]).await;
201206
let (_minerd_process, _minerd_addr) = start_minerd(tproxy_addr, None, None, false).await;
202207

203208
sniffer

integration-tests/tests/jd_provide_missing_transaction.rs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,11 +6,11 @@ async fn jds_ask_for_missing_transactions() {
66
start_tracing();
77
let (tp_1, tp_addr_1) = start_template_provider(None, DifficultyLevel::Low);
88
let (tp_2, tp_addr_2) = start_template_provider(None, DifficultyLevel::Low);
9-
let (_pool, pool_addr) = start_pool(Some(tp_addr_1)).await;
9+
let (_pool, pool_addr) = start_pool(Some(tp_addr_1), vec![], vec![]).await;
1010
let (_jds, jds_addr) = start_jds(tp_1.rpc_info());
1111
let (sniffer, sniffer_addr) = start_sniffer("A", jds_addr, false, vec![], None);
12-
let (_jdc, jdc_addr) = start_jdc(&[(pool_addr, sniffer_addr)], tp_addr_2);
13-
let (_translator, tproxy_addr) = start_sv2_translator(jdc_addr, false).await;
12+
let (_jdc, jdc_addr) = start_jdc(&[(pool_addr, sniffer_addr)], tp_addr_2, vec![], vec![]);
13+
let (_translator, tproxy_addr) = start_sv2_translator(jdc_addr, false, vec![], vec![]).await;
1414
let (_minerd_process, _minerd_addr) = start_minerd(tproxy_addr, None, None, false).await;
1515
assert!(tp_2.fund_wallet().is_ok());
1616
assert!(tp_2.create_mempool_transaction().is_ok());

integration-tests/tests/jd_tproxy_integration.rs

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -5,14 +5,20 @@ use stratum_apps::stratum_core::{common_messages_sv2::*, mining_sv2::*};
55
async fn jd_non_aggregated_tproxy_integration() {
66
start_tracing();
77
let (tp, tp_addr) = start_template_provider(None, DifficultyLevel::Low);
8-
let (_pool, pool_addr) = start_pool(Some(tp_addr)).await;
8+
let (_pool, pool_addr) = start_pool(Some(tp_addr), vec![], vec![]).await;
99
let (jdc_pool_sniffer, jdc_pool_sniffer_addr) =
1010
start_sniffer("0", pool_addr, false, vec![], None);
1111
let (_jds, jds_addr) = start_jds(tp.rpc_info());
12-
let (_jdc, jdc_addr) = start_jdc(&[(jdc_pool_sniffer_addr, jds_addr)], tp_addr);
12+
let (_jdc, jdc_addr) = start_jdc(
13+
&[(jdc_pool_sniffer_addr, jds_addr)],
14+
tp_addr,
15+
vec![],
16+
vec![],
17+
);
1318
let (tproxy_jdc_sniffer, tproxy_jdc_sniffer_addr) =
1419
start_sniffer("1", jdc_addr, false, vec![], None);
15-
let (_translator, tproxy_addr) = start_sv2_translator(tproxy_jdc_sniffer_addr, false).await;
20+
let (_translator, tproxy_addr) =
21+
start_sv2_translator(tproxy_jdc_sniffer_addr, false, vec![], vec![]).await;
1622

1723
// start two minerd processes
1824
let (_minerd_process, _minerd_addr) = start_minerd(tproxy_addr, None, None, false).await;
@@ -74,14 +80,20 @@ async fn jd_non_aggregated_tproxy_integration() {
7480
async fn jd_aggregated_tproxy_integration() {
7581
start_tracing();
7682
let (tp, tp_addr) = start_template_provider(None, DifficultyLevel::Low);
77-
let (_pool, pool_addr) = start_pool(Some(tp_addr)).await;
83+
let (_pool, pool_addr) = start_pool(Some(tp_addr), vec![], vec![]).await;
7884
let (jdc_pool_sniffer, jdc_pool_sniffer_addr) =
7985
start_sniffer("0", pool_addr, false, vec![], None);
8086
let (_jds, jds_addr) = start_jds(tp.rpc_info());
81-
let (_jdc, jdc_addr) = start_jdc(&[(jdc_pool_sniffer_addr, jds_addr)], tp_addr);
87+
let (_jdc, jdc_addr) = start_jdc(
88+
&[(jdc_pool_sniffer_addr, jds_addr)],
89+
tp_addr,
90+
vec![],
91+
vec![],
92+
);
8293
let (tproxy_jdc_sniffer, tproxy_jdc_sniffer_addr) =
8394
start_sniffer("1", jdc_addr, false, vec![], None);
84-
let (_translator, tproxy_addr) = start_sv2_translator(tproxy_jdc_sniffer_addr, true).await;
95+
let (_translator, tproxy_addr) =
96+
start_sv2_translator(tproxy_jdc_sniffer_addr, true, vec![], vec![]).await;
8597

8698
// start two minerd processes
8799
let (_minerd_process, _minerd_addr) = start_minerd(tproxy_addr, None, None, false).await;

integration-tests/tests/jdc_block_propagation.rs

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ async fn propagated_from_jdc_to_tp() {
1111
start_tracing();
1212
let (tp, tp_addr) = start_template_provider(None, DifficultyLevel::Low);
1313
let current_block_hash = tp.get_best_block_hash().unwrap();
14-
let (_pool, pool_addr) = start_pool(Some(tp_addr)).await;
14+
let (_pool, pool_addr) = start_pool(Some(tp_addr), vec![], vec![]).await;
1515
let (_jds, jds_addr) = start_jds(tp.rpc_info());
1616
let ignore_push_solution =
1717
IgnoreMessage::new(MessageDirection::ToUpstream, MESSAGE_TYPE_PUSH_SOLUTION);
@@ -23,8 +23,13 @@ async fn propagated_from_jdc_to_tp() {
2323
None,
2424
);
2525
let (jdc_tp_sniffer, jdc_tp_sniffer_addr) = start_sniffer("1", tp_addr, false, vec![], None);
26-
let (_jdc, jdc_addr) = start_jdc(&[(pool_addr, jdc_jds_sniffer_addr)], jdc_tp_sniffer_addr);
27-
let (_translator, tproxy_addr) = start_sv2_translator(jdc_addr, false).await;
26+
let (_jdc, jdc_addr) = start_jdc(
27+
&[(pool_addr, jdc_jds_sniffer_addr)],
28+
jdc_tp_sniffer_addr,
29+
vec![],
30+
vec![],
31+
);
32+
let (_translator, tproxy_addr) = start_sv2_translator(jdc_addr, false, vec![], vec![]).await;
2833
let (_minerd_process, _minerd_addr) = start_minerd(tproxy_addr, None, None, false).await;
2934
jdc_tp_sniffer
3035
.wait_for_message_type(MessageDirection::ToUpstream, MESSAGE_TYPE_SUBMIT_SOLUTION)

0 commit comments

Comments
 (0)