diff --git a/network-runner/Cargo.lock b/network-runner/Cargo.lock index e133e41..c6800a3 100644 --- a/network-runner/Cargo.lock +++ b/network-runner/Cargo.lock @@ -1865,7 +1865,7 @@ dependencies = [ [[package]] name = "nomos-mix" version = "0.1.0" -source = "git+https://github.com/logos-co/nomos-node?rev=9b29c17#9b29c17e2f7e5da7f17c4d7a2b8e8749c0b32617" +source = "git+https://github.com/logos-co/nomos-node?rev=e095964#e0959644a9927d3930197df281684615dee7a6ef" dependencies = [ "cached", "futures", @@ -1882,7 +1882,7 @@ dependencies = [ [[package]] name = "nomos-mix-message" version = "0.1.0" -source = "git+https://github.com/logos-co/nomos-node?rev=9b29c17#9b29c17e2f7e5da7f17c4d7a2b8e8749c0b32617" +source = "git+https://github.com/logos-co/nomos-node?rev=e095964#e0959644a9927d3930197df281684615dee7a6ef" dependencies = [ "serde", "sphinx-packet", diff --git a/network-runner/Cargo.toml b/network-runner/Cargo.toml index e11de03..f9e76d1 100644 --- a/network-runner/Cargo.toml +++ b/network-runner/Cargo.toml @@ -26,14 +26,7 @@ humantime = "2.1" humantime-serde = "1" once_cell = "1.17" parking_lot = "0.12" -polars = { version = "0.27", features = [ - "serde", - "object", - "json", - "csv-file", - "parquet", - "dtype-struct", -], optional = true } +polars = { version = "0.27", features = ["serde", "object", "json", "csv-file", "parquet", "dtype-struct"], optional = true } rand = { version = "0.8", features = ["small_rng"] } rayon = "1.7" scopeguard = "1" @@ -41,19 +34,12 @@ serde = { version = "1.0", features = ["derive", "rc"] } serde_with = "2.3" serde_json = "1.0" thiserror = "1" -tracing = { version = "0.1", default-features = false, features = [ - "log", - "attributes", -] } -tracing-subscriber = { version = "0.3", features = [ - "json", - "env-filter", - "tracing-log", -] } -nomos-mix = { git = "https://github.com/logos-co/nomos-node", rev = "9b29c17", package = "nomos-mix" } -nomos-mix-message = { git = "https://github.com/logos-co/nomos-node", rev = "9b29c17", package = "nomos-mix-message" } -rand_chacha = "0.3.1" -multiaddr = "0.18.2" +tracing = { version = "0.1", default-features = false, features = ["log", "attributes"] } +tracing-subscriber = { version = "0.3", features = ["json", "env-filter", "tracing-log"]} +nomos-mix = { git = "https://github.com/logos-co/nomos-node", rev = "e095964", package = "nomos-mix" } +nomos-mix-message = { git = "https://github.com/logos-co/nomos-node", rev = "e095964", package = "nomos-mix-message" } +rand_chacha = "0.3" +multiaddr = "0.18" [target.'cfg(target_arch = "wasm32")'.dependencies] getrandom = { version = "0.2", features = ["js"] } diff --git a/network-runner/src/node/mix/mod.rs b/network-runner/src/node/mix/mod.rs index 52c6033..9e70381 100644 --- a/network-runner/src/node/mix/mod.rs +++ b/network-runner/src/node/mix/mod.rs @@ -29,9 +29,7 @@ use std::{ use stream_wrapper::CrossbeamReceiverStream; #[derive(Debug, Clone)] -pub enum MixMessage { - Dummy(Vec), -} +pub struct MixMessage(Vec); impl PayloadSize for MixMessage { fn size_bytes(&self) -> u32 { @@ -79,12 +77,6 @@ impl MixNode { settings: MixnodeSettings, network_interface: InMemoryNetworkInterface, ) -> Self { - let state = MixnodeState { - node_id: id, - mock_counter: 0, - step_id: 0, - }; - let mut rng_generator = ChaCha12Rng::seed_from_u64(settings.seed); // Init Tier-1: Persistent transmission @@ -102,7 +94,7 @@ impl MixNode { ), ); - // Init Tier-2: Temporal processor + // Init Tier-2: message blend let (blend_sender, blend_receiver) = channel::unbounded(); let (blend_update_time_sender, blend_update_time_receiver) = channel::unbounded(); let nodes: Vec< @@ -137,7 +129,11 @@ impl MixNode { id, network_interface, settings, - state, + state: MixnodeState { + node_id: id, + mock_counter: 0, + step_id: 0, + }, persistent_sender, persistent_update_time_sender, persistent_transmission_messages, @@ -187,15 +183,12 @@ impl Node for MixNode { let messages = self.network_interface.receive_messages(); for message in messages { println!(">>>>> Node {}, message: {message:?}", self.id); - let MixMessage::Dummy(msg) = message.into_payload(); - blend_sender.send(msg).unwrap(); + blend_sender.send(message.into_payload().0).unwrap(); } - self.state.step_id += 1; - self.state.mock_counter += 1; - let waker = futures::task::noop_waker(); let mut cx = futures::task::Context::from_waker(&waker); + // Proceed message blend if let Poll::Ready(Some(msg)) = pin::pin!(blend_messages).poll_next(&mut cx) { match msg { MixOutgoingMessage::Outbound(msg) => { @@ -206,11 +199,14 @@ impl Node for MixNode { } } } - + // Proceed persistent transmission if let Poll::Ready(Some(msg)) = pin::pin!(persistent_transmission_messages).poll_next(&mut cx) { - self.forward(MixMessage::Dummy(msg)); + self.forward(MixMessage(msg)); } + + self.state.step_id += 1; + self.state.mock_counter += 1; } }