use std::collections::HashSet;
use std::net::SocketAddr;
use std::time::Duration;
use chrono::{DateTime, Utc};
use hydroflow::scheduled::graph::Hydroflow;
use hydroflow::util::{bind_udp_bytes, ipv4_resolve};
use hydroflow_macro::hydroflow_syntax;
use rand::seq::SliceRandom;
use rand::thread_rng;
use serde::{Deserialize, Serialize};
use crate::protocol::{Message, MessageWithAddr};
use crate::Role::{
Client, GossipingServer1, GossipingServer2, GossipingServer3, GossipingServer4,
GossipingServer5, Server,
};
use crate::{default_server_address, Opts, Role};
#[derive(PartialEq, Eq, Clone, Serialize, Deserialize, Debug, Hash)]
pub struct ChatMessage {
nickname: String,
message: String,
ts: DateTime<Utc>,
}
enum InfectionOperation {
InfectWithMessage { msg: ChatMessage },
RemoveForMessage { msg: ChatMessage },
}
pub const REMOVAL_PROBABILITY: f32 = 1.0 / 4.0;
pub(crate) async fn run_gossiping_server(opts: Opts) {
let server_address = opts.address.unwrap_or_else(default_server_address);
let all_members = [
GossipingServer1,
GossipingServer2,
GossipingServer3,
GossipingServer4,
GossipingServer5,
];
let other_members: Vec<Role> = all_members
.into_iter()
.filter(|role| *role != opts.role)
.collect();
let gossip_listening_addr = gossip_address(&opts.role);
println!("Starting server on {:?}", server_address);
let (client_outbound, client_inbound, actual_server_addr) =
bind_udp_bytes(server_address).await;
let (gossip_outbound, gossip_inbound, _) = bind_udp_bytes(gossip_listening_addr).await;
println!(
"Server is live! Listening on {:?}. Gossiping On: {:?}",
actual_server_addr, gossip_listening_addr
);
let mut hf: Hydroflow = hydroflow_syntax! {
client_out = union() -> dest_sink_serde(client_outbound);
client_in = source_stream_serde(client_inbound)
-> map(Result::unwrap)
-> map(|(msg, addr)| MessageWithAddr::from_message(msg, addr))
-> demux_enum::<MessageWithAddr>();
clients = client_in[ConnectRequest] -> map(|(addr,)| addr) -> tee();
client_in[ConnectResponse] -> for_each(|(addr,)| println!("Received unexpected `ConnectResponse` as server from addr {}.", addr));
clients[0] -> map(|addr| (Message::ConnectResponse, addr)) -> [0]client_out;
messages_from_connected_client = client_in[ChatMsg]
-> map(|(_addr, nickname, message, ts)| ChatMessage { nickname, message, ts })
-> maybe_new_messages;
clients[1] -> [1]broadcast;
broadcast = cross_join::<'tick, 'static>() -> [1]client_out;
gossip_out = dest_sink_serde(gossip_outbound);
gossip_in = source_stream_serde(gossip_inbound)
-> map(Result::unwrap)
-> map(|(message, _)| message)
-> maybe_new_messages;
maybe_new_messages = union();
actually_new_messages = difference() -> tee();
maybe_new_messages -> [pos]actually_new_messages;
all_messages -> [neg]actually_new_messages;
actually_new_messages -> defer_tick() -> all_messages; actually_new_messages
-> map(|chat_msg: ChatMessage| Message::ChatMsg {
nickname: chat_msg.nickname,
message: chat_msg.message,
ts: chat_msg.ts})
-> [0]broadcast; actually_new_messages
-> map(|msg: ChatMessage| InfectionOperation::InfectWithMessage { msg })
-> infecting_messages;
all_messages = fold::<'static>(HashSet::<ChatMessage>::new, |accum, message| {
accum.insert(message);
}) -> flatten();
infecting_messages = union() -> fold::<'static>(HashSet::<ChatMessage>::new, |accum, op| {
match op {
InfectionOperation::InfectWithMessage{ msg } => {accum.insert(msg);},
InfectionOperation::RemoveForMessage{ msg } => { accum.remove(&msg);}
}
});
source_interval(Duration::from_secs(1)) -> [0]triggered_messages; triggered_messages = cross_join()
-> map(|(_, message)| {
// Choose a random peer
let random_peer = other_members.choose(&mut thread_rng()).unwrap();
(message, gossip_address(random_peer))
})
-> tee();
infecting_messages -> flatten() -> [1]triggered_messages;
triggered_messages
-> inspect(|(msg, addr)| println!("Gossiping {:?} to {:?}", msg, addr))
-> gossip_out;
triggered_messages
-> filter_map(|(msg, _addr)| {
if rand::random::<f32>() < REMOVAL_PROBABILITY{
println!("Dropping Message {:?}", msg);
Some(InfectionOperation::RemoveForMessage{ msg })
} else {
None
}
})
-> defer_tick()
-> infecting_messages;
};
#[cfg(feature = "debugging")]
if let Some(graph) = opts.graph {
let serde_graph = hf
.meta_graph()
.expect("No graph found, maybe failed to parse.");
serde_graph.open_graph(graph, opts.write_config).unwrap();
}
hf.run_async().await.unwrap();
}
fn gossip_address(role: &Role) -> SocketAddr {
match role {
Client | Server => {
panic!("Incorrect role {:?} for gossip server.", role)
}
GossipingServer1 => ipv4_resolve("localhost:54322"),
GossipingServer2 => ipv4_resolve("localhost:54323"),
GossipingServer3 => ipv4_resolve("localhost:54324"),
GossipingServer4 => ipv4_resolve("localhost:54325"),
GossipingServer5 => ipv4_resolve("localhost:54326"),
}
.unwrap()
}