Proof of Concept
--- a/miner-apps/translator/src/lib/io_task.rs
+++ b/miner-apps/translator/src/lib/io_task.rs
@@ -158,3 +158,107 @@
);
}
}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use std::{convert::TryInto as _, net::SocketAddr, time::Duration};
+
+ use async_channel::unbounded;
+ use stratum_apps::{
+ key_utils::{Secp256k1PublicKey, Secp256k1SecretKey},
+ network_helpers::{accept_noise_connection, connect_with_noise},
+ stratum_core::common_messages_sv2::{Protocol, SetupConnection},
+ };
+ use tokio::net::TcpSocket;
+
+ const TEST_PUBLIC_KEY: &str = "9auqWEzQDVyd2oe1JVGFLMLHZtCo2FFqZwtKA5gd9xbuEu7PH72";
+ const TEST_SECRET_KEY: &str = "mkDLTBBRxdBv998612qipDYoTK3YUrqLe8uWw7gu3iXbSrn2n";
+
+ fn test_keys() -> (Secp256k1PublicKey, Secp256k1SecretKey) {
+ (
+ TEST_PUBLIC_KEY.parse().unwrap(),
+ TEST_SECRET_KEY.parse().unwrap(),
+ )
+ }
+
+ fn large_setup_frame() -> Sv2Frame {
+ let field = "x".repeat(255);
+ let setup = SetupConnection {
+ protocol: Protocol::MiningProtocol,
+ min_version: 2,
+ max_version: 2,
+ flags: 0,
+ endpoint_host: field.clone().try_into().unwrap(),
+ endpoint_port: 0,
+ vendor: field.clone().try_into().unwrap(),
+ hardware_version: field.clone().try_into().unwrap(),
+ firmware: field.clone().try_into().unwrap(),
+ device_id: field.try_into().unwrap(),
+ };
+
+ Message::Common(setup.into()).try_into().unwrap()
+ }
+
+ #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
+ async fn fallback_does_not_wait_forever_on_blocked_writer() {
+ let listener_socket = TcpSocket::new_v4().unwrap();
+ listener_socket.set_send_buffer_size(1024).unwrap();
+ listener_socket
+ .bind(SocketAddr::from(([127, 0, 0, 1], 0)))
+ .unwrap();
+ let listener = listener_socket.listen(1).unwrap();
+ let address = listener.local_addr().unwrap();
+ let (authority_pubkey, authority_secret_key) = test_keys();
+
+ let server = tokio::spawn(async move {
+ let (stream, _) = listener.accept().await.unwrap();
+ accept_noise_connection::<Message>(
+ stream,
+ authority_pubkey,
+ authority_secret_key,
+ 60,
+ )
+ .await
+ .unwrap()
+ });
+
+ let client_socket = TcpSocket::new_v4().unwrap();
+ client_socket.set_recv_buffer_size(1024).unwrap();
+ let client_stream = client_socket.connect(address).await.unwrap();
+ let client = connect_with_noise::<Message>(client_stream, Some(authority_pubkey))
+ .await
+ .unwrap();
+ let server = server.await.unwrap();
+
+ let (_peer_reader, _peer_writer) = client.into_split();
+ let (reader, writer) = server.into_split();
+ let task_manager = Arc::new(TaskManager::new());
+ let (outbound_tx, outbound_rx) = unbounded::<Sv2Frame>();
+ let (inbound_tx, _inbound_rx) = unbounded::<Sv2Frame>();
+ let fallback_coordinator = FallbackCoordinator::new();
+
+ spawn_io_tasks(
+ task_manager,
+ reader,
+ writer,
+ outbound_rx,
+ inbound_tx,
+ CancellationToken::new(),
+ fallback_coordinator.clone(),
+ );
+
+ for _ in 0..16_384 {
+ outbound_tx.send(large_setup_frame()).await.unwrap();
+ }
+
+ tokio::time::sleep(Duration::from_millis(250)).await;
+
+ tokio::time::timeout(
+ Duration::from_secs(2),
+ fallback_coordinator.trigger_fallback_and_wait(),
+ )
+ .await
+ .expect("fallback cleanup hung while a peer kept the writer backpressured");
+ }
+}
Suggested Fix
diff --git a/miner-apps/translator/src/lib/io_task.rs b/miner-apps/translator/src/lib/io_task.rs
--- a/miner-apps/translator/src/lib/io_task.rs
+++ b/miner-apps/translator/src/lib/io_task.rs
@@ -127,10 +127,25 @@
match res {
Ok(frame) => {
trace!("Sending outbound frame");
- if let Err(e) = writer.write_frame(frame.into()).await {
- error!(error=?e, "Writer error");
- outbound_rx.close_and_drain();
- break;
+ tokio::select! {
+ biased;
+ _ = cancellation_token.cancelled() => {
+ trace!("Received app shutdown signal");
+ inbound_tx_clone.close();
+ break;
+ }
+ _ = fallback_token.cancelled() => {
+ trace!("Received fallback signal");
+ inbound_tx_clone.close();
+ break;
+ }
+ res = writer.write_frame(frame.into()) => {
+ if let Err(e) = res {
+ error!(error=?e, "Writer error");
+ outbound_rx.close_and_drain();
+ break;
+ }
+ }
}
}
Err(_) => {
Proof of Concept
Suggested Fix