use std::net::SocketAddr;
use chrono::prelude::*;
use hydroflow::hydroflow_syntax;
use hydroflow::lattices::{Max, Merge};
use hydroflow::scheduled::graph::Hydroflow;
use hydroflow::util::{UdpSink, UdpStream};
use crate::protocol::EchoMsg;
use crate::Opts;
pub(crate) async fn run_server(outbound: UdpSink, inbound: UdpStream, opts: Opts) {
let bot: Max<usize> = Max::new(0);
println!("Server live!");
let mut flow: Hydroflow = hydroflow_syntax! {
inbound_chan = source_stream_serde(inbound) -> map(Result::unwrap) -> tee();
inbound_chan[print]
-> for_each(|(msg, addr): (EchoMsg, SocketAddr)| println!("{}: Got {:?} from {:?}", Utc::now(), msg, addr));
inbound_chan[merge] -> map(|(msg, _addr): (EchoMsg, SocketAddr)| msg.lamport_clock) -> mergevc;
mergevc = fold::<'static>(
|| bot,
|old: &mut Max<usize>, lamport_clock: Max<usize>| {
let bump = Max::new(old.into_reveal() + 1);
old.merge(bump);
old.merge(lamport_clock);
}
);
inbound_chan[1] -> map(|(EchoMsg {payload, ..}, addr)| (payload, addr) )
-> [0]stamped_output;
mergevc -> [1]stamped_output;
stamped_output = cross_join::<'tick, 'tick>() -> map(|((payload, addr), the_clock): ((String, SocketAddr), Max<usize>)| (EchoMsg { payload, lamport_clock: the_clock }, addr))
-> dest_sink_serde(outbound);
};
#[cfg(feature = "debugging")]
if let Some(graph) = opts.graph {
let serde_graph = flow
.meta_graph()
.expect("No graph found, maybe failed to parse.");
serde_graph.open_graph(graph, opts.write_config).unwrap();
}
let _ = opts;
flow.run_async().await;
}