diff --git a/tools/internal/mixnet-check-all-gateways/Cargo.toml b/tools/internal/mixnet-check-all-gateways/Cargo.toml index 18bab516f8..7c6ffc5f54 100644 --- a/tools/internal/mixnet-check-all-gateways/Cargo.toml +++ b/tools/internal/mixnet-check-all-gateways/Cargo.toml @@ -31,3 +31,4 @@ uuid = { version = "1", features = ["v4", "serde"] } bincode = "1.0" dirs.workspace = true tokio-util = { workspace = true } +tempfile.workspace = true diff --git a/tools/internal/mixnet-check-all-gateways/src/main.rs b/tools/internal/mixnet-check-all-gateways/src/main.rs index ec57bfa199..b87090bcbb 100644 --- a/tools/internal/mixnet-check-all-gateways/src/main.rs +++ b/tools/internal/mixnet-check-all-gateways/src/main.rs @@ -47,7 +47,7 @@ enum TestError { #[tokio::main] async fn main() -> Result<()> { - setup_logging(); // TODO think about parsing and noise here, could just parse on errors from all libs and then have info from here? echo server metrics + info logging from this code + error logs from elsewhere should be ok + setup_logging(); // TODO think about parsing and noise here if make it concurrent. Could just parse on errors from all libs and then have info from here? echo server metrics + info logging from this code + error logs from elsewhere should be ok let entry_gw_keys = reqwest_and_parse( Url::parse("https://validator.nymtech.net/api/v1/unstable/nym-nodes/skimmed/entry-gateways/all?no_legacy=true").unwrap(), "EntryGateway", @@ -75,14 +75,13 @@ async fn main() -> Result<()> { ); info!("got {} exit gws", exit_gw_keys.len(),); - let mut port_range: u64 = 9000; // port that we start iterating upwards from, will go from port_range to (port_range + exit_gws.len()) by the end of the run. This was made configurable presuming at some point we'd try make this run concurrently for speedup. + let mut port_range: u64 = 9000; // Port that we start iterating upwards from, will go from port_range to (port_range + exit_gws.len()) by the end of the run. This was made configurable presuming at some point we'd try make this run concurrently for speedup. let cancel_token = CancellationToken::new(); let watcher_token = cancel_token.clone(); tokio::spawn(async move { signal::ctrl_c().await?; println!("CTRL_C received"); - // TODO shutdown logic for loops watcher_token.cancel(); Ok::<(), anyhow::Error>(()) }); @@ -91,8 +90,13 @@ async fn main() -> Result<()> { let time_now = start.duration_since(UNIX_EPOCH).unwrap().as_secs(); for exit_gw in exit_gw_keys.clone() { - let echo_server_token = CancellationToken::new(); - let thread_echo_server_token = echo_server_token.clone(); + let loop_token = cancel_token.child_token(); + let inner_loop_token = cancel_token.child_token(); + if loop_token.is_cancelled() { + break; + } + let thread_token = cancel_token.child_token(); + let last_check_token = thread_token.clone(); if !fs::metadata(format!("./src/results/{}", time_now)) .map(|metadata| metadata.is_dir()) @@ -122,7 +126,7 @@ async fn main() -> Result<()> { ) .as_str(), ), - "../../../envs/mainnet.env".to_string(), // TODO make configurable + "../../../envs/mainnet.env".to_string(), // TODO replace with None port_range.to_string().as_str(), ) .await @@ -143,16 +147,16 @@ async fn main() -> Result<()> { }; port_range += 1; + let echo_disconnect_signal = echo_server.disconnect_signal(); let echo_addr = echo_server.nym_address().await; debug!("echo addr: {echo_addr}"); tokio::task::spawn(async move { loop { tokio::select! { - _ = thread_echo_server_token.cancelled() => { + _ = thread_token.cancelled() => { info!("loop over; disconnecting echo server {}", echo_addr.clone()); - echo_server.disconnect().await; // TODO fix this failing lock: can't aquire lock on the ProxyServer to kill it - info!("disconnected {}", echo_addr.clone()); + echo_disconnect_signal.send(()).await?; break; } _ = echo_server.run() => {} @@ -165,97 +169,104 @@ async fn main() -> Result<()> { time::sleep(Duration::from_secs(5)).await; for entry_gw in entry_gw_keys.clone() { - // let builder = mixnet::MixnetClientBuilder::new_ephemeral() - // .request_gateway(entry_gw.clone()) - // .build()?; + if inner_loop_token.is_cancelled() { + info!("Inner loop cancelled"); + break; + } + let builder = mixnet::MixnetClientBuilder::new_ephemeral() + .request_gateway(entry_gw.clone()) + .build()?; - // let mut client = match builder.connect_to_mixnet().await { - // Ok(client) => client, - // Err(err) => { - // let res = TestResult { - // entry_gw: entry_gw.clone(), - // exit_gw: exit_gw.clone(), - // error: TestError::Other(err.to_string()), - // }; - // info!("{res:#?}"); - // results_vec.push(json!(res)); - // continue; - // } - // }; + let mut client = match builder.connect_to_mixnet().await { + Ok(client) => client, + Err(err) => { + let res = TestResult { + entry_gw: entry_gw.clone(), + exit_gw: exit_gw.clone(), + error: TestError::Other(err.to_string()), + }; + info!("{res:#?}"); + results_vec.push(json!(res)); + continue; + } + }; - // let test_address = client.nym_address(); - // info!("currently testing entry gateway: {test_address}"); + let test_address = client.nym_address(); + info!("currently testing entry gateway: {test_address}"); - // // Has to be ProxiedMessage for the moment which is slightly annoying until I - // // modify the ProxyServer to just stupidly echo back whatever it gets in a - // // ReconstructedMessage format if it can't deseralise incoming traffic to a ProxiedMessage - // let session_id = uuid::Uuid::new_v4(); - // let message_id = 0; - // let outgoing = ProxiedMessage::new( - // Payload::Data(MESSAGE.as_bytes().to_vec()), - // session_id, - // message_id, - // ); - // let coded_message = bincode::serialize(&outgoing).unwrap(); + // Has to be ProxiedMessage for the moment which is slightly annoying until I + // modify the ProxyServer to just stupidly echo back whatever it gets in a + // ReconstructedMessage format if it can't deseralise incoming traffic to a ProxiedMessage + let session_id = uuid::Uuid::new_v4(); + let message_id = 0; + let outgoing = ProxiedMessage::new( + Payload::Data(MESSAGE.as_bytes().to_vec()), + session_id, + message_id, + ); + let coded_message = bincode::serialize(&outgoing).unwrap(); - // match client - // .send_message(echo_addr, &coded_message, IncludedSurbs::Amount(30)) - // .await - // { - // Ok(_) => { - // debug!("Message sent"); - // } - // Err(err) => { - // let res = TestResult { - // entry_gw: entry_gw.clone(), - // exit_gw: exit_gw.clone(), - // error: TestError::Other(err.to_string()), - // }; - // info!("{res:#?}"); - // results_vec.push(json!(res)); - // continue; - // } - // }; + match client + .send_message(echo_addr, &coded_message, IncludedSurbs::Amount(30)) + .await + { + Ok(_) => { + debug!("Message sent"); + } + Err(err) => { + let res = TestResult { + entry_gw: entry_gw.clone(), + exit_gw: exit_gw.clone(), + error: TestError::Other(err.to_string()), + }; + info!("{res:#?}"); + results_vec.push(json!(res)); + continue; + } + }; - // let res = match timeout(Duration::from_secs(TIMEOUT), client.next()).await { - // Err(_timeout) => { - // warn!("timed out"); - // TestResult { - // entry_gw: entry_gw.clone(), - // exit_gw: exit_gw.clone(), - // error: TestError::Timeout, - // } - // } - // Ok(Some(received)) => { - // let incoming: ProxiedMessage = bincode::deserialize(&received.message).unwrap(); - // info!("got echo: {:?}", incoming); - // debug!( - // "sent message as lazy ref until I properly sort the utils for comparison: {:?}", - // MESSAGE.as_bytes() - // ); - // debug!("incoming message: {:?}", incoming.message); - // TestResult { - // entry_gw: entry_gw.clone(), - // exit_gw: exit_gw.clone(), - // error: TestError::None, - // } - // } - // Ok(None) => { - // info!("failed to receive any message back..."); - // TestResult { - // entry_gw: entry_gw.clone(), - // exit_gw: exit_gw.clone(), - // error: TestError::NoMessage, - // } - // } - // }; - // debug!("{res:#?}"); - // results_vec.push(json!(res)); - // debug!("disconnecting the client before shutting down..."); - // client.disconnect().await; + let res = match timeout(Duration::from_secs(TIMEOUT), client.next()).await { + Err(_timeout) => { + warn!("timed out"); + TestResult { + entry_gw: entry_gw.clone(), + exit_gw: exit_gw.clone(), + error: TestError::Timeout, + } + } + Ok(Some(received)) => { + let incoming: ProxiedMessage = bincode::deserialize(&received.message).unwrap(); + info!("got echo: {:?}", incoming); + debug!( + "sent message as lazy ref until I properly sort the utils for comparison: {:?}", + MESSAGE.as_bytes() + ); + debug!("incoming message: {:?}", incoming.message); + TestResult { + entry_gw: entry_gw.clone(), + exit_gw: exit_gw.clone(), + error: TestError::None, + } + } + Ok(None) => { + info!("failed to receive any message back..."); + TestResult { + entry_gw: entry_gw.clone(), + exit_gw: exit_gw.clone(), + error: TestError::NoMessage, + } + } + }; + debug!("{res:#?}"); + results_vec.push(json!(res)); + debug!("disconnecting the client before shutting down..."); + client.disconnect().await; + + tokio::time::sleep(tokio::time::Duration::from_millis(500)).await; + info!("{}", &entry_gw); if Some(&entry_gw) == entry_gw_keys.last() { - echo_server_token.cancel(); + last_check_token.cancel(); } } let json_array = json!(results_vec); @@ -278,7 +289,7 @@ fn filter_gateway_keys(json: &Value, key: &str) -> Result> { if let Some(nodes) = json["nodes"]["data"].as_array() { for node in nodes { if let Some(performance) = node.get("performance").and_then(|v| v.as_str()) { - let performance_value: f64 = performance.parse().unwrap_or(0.0); // TODO make this configurable + let performance_value: f64 = performance.parse().unwrap_or(0.0); // TODO make this configurable? let inactive = node.get("role").and_then(|v| v.as_str()) == Some("Inactive"); @@ -299,7 +310,6 @@ fn filter_gateway_keys(json: &Value, key: &str) -> Result> { } } else { // TODO make error and return / break - // println!("No nodes found "); return Err(anyhow!("Could not parse any gateways")); } Ok(filtered_keys)