This commit is contained in:
Youngjoon Lee 2024-11-06 20:27:58 +07:00
parent b5b491021a
commit 65dc154664
No known key found for this signature in database
GPG Key ID: 25CA11F37F095E5D
3 changed files with 23 additions and 41 deletions

View File

@ -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",

View File

@ -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"] }

View File

@ -29,9 +29,7 @@ use std::{
use stream_wrapper::CrossbeamReceiverStream;
#[derive(Debug, Clone)]
pub enum MixMessage {
Dummy(Vec<u8>),
}
pub struct MixMessage(Vec<u8>);
impl PayloadSize for MixMessage {
fn size_bytes(&self) -> u32 {
@ -79,12 +77,6 @@ impl MixNode {
settings: MixnodeSettings,
network_interface: InMemoryNetworkInterface<MixMessage>,
) -> 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;
}
}