single-node broadcast works
This commit is contained in:
parent
3186c08f32
commit
8f869ff0d5
5 changed files with 94 additions and 2 deletions
10
Cargo.lock
generated
10
Cargo.lock
generated
|
@ -184,6 +184,16 @@ dependencies = [
|
||||||
"slab",
|
"slab",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "gg-broadcast"
|
||||||
|
version = "0.0.1"
|
||||||
|
dependencies = [
|
||||||
|
"async-trait",
|
||||||
|
"maelstrom-node",
|
||||||
|
"serde_json",
|
||||||
|
"tokio",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "gg-echo"
|
name = "gg-echo"
|
||||||
version = "0.0.1"
|
version = "0.0.1"
|
||||||
|
|
|
@ -1,5 +1,5 @@
|
||||||
[workspace]
|
[workspace]
|
||||||
members = ["gg-echo", "gg-uid"]
|
members = ["gg-echo", "gg-uid", "gg-broadcast"]
|
||||||
resolver = "2"
|
resolver = "2"
|
||||||
|
|
||||||
[workspace.package]
|
[workspace.package]
|
||||||
|
@ -11,3 +11,4 @@ license-file = "LICENSE.md"
|
||||||
async-trait = "0.1"
|
async-trait = "0.1"
|
||||||
maelstrom-node = "0.1"
|
maelstrom-node = "0.1"
|
||||||
tokio = { version = "1", default-features = false, features = ["io-util", "io-std", "rt-multi-thread", "macros"] }
|
tokio = { version = "1", default-features = false, features = ["io-util", "io-std", "rt-multi-thread", "macros"] }
|
||||||
|
serde_json = "1"
|
||||||
|
|
12
gg-broadcast/Cargo.toml
Normal file
12
gg-broadcast/Cargo.toml
Normal file
|
@ -0,0 +1,12 @@
|
||||||
|
[package]
|
||||||
|
name = "gg-broadcast"
|
||||||
|
edition = "2021"
|
||||||
|
version.workspace = true
|
||||||
|
authors.workspace = true
|
||||||
|
license-file.workspace = true
|
||||||
|
|
||||||
|
[dependencies]
|
||||||
|
async-trait.workspace = true
|
||||||
|
maelstrom-node.workspace = true
|
||||||
|
serde_json.workspace = true
|
||||||
|
tokio.workspace = true
|
69
gg-broadcast/src/main.rs
Normal file
69
gg-broadcast/src/main.rs
Normal file
|
@ -0,0 +1,69 @@
|
||||||
|
use async_trait::async_trait;
|
||||||
|
use maelstrom::protocol::Message;
|
||||||
|
use maelstrom::{done, protocol::MessageBody, Node, Result, Runtime};
|
||||||
|
//use serde_json::{Map, Value};
|
||||||
|
use std::borrow::BorrowMut;
|
||||||
|
use std::collections::HashMap;
|
||||||
|
use std::sync::{Arc, Mutex};
|
||||||
|
|
||||||
|
#[tokio::main]
|
||||||
|
async fn main() {
|
||||||
|
let handler = Arc::new(Handler::default());
|
||||||
|
let _ = Runtime::new()
|
||||||
|
.with_handler(handler)
|
||||||
|
.run()
|
||||||
|
.await
|
||||||
|
.unwrap_or_default();
|
||||||
|
}
|
||||||
|
|
||||||
|
type Store = Arc<Mutex<HashMap<String, i64>>>;
|
||||||
|
|
||||||
|
#[derive(Clone, Default)]
|
||||||
|
struct Handler {
|
||||||
|
pub store: Store,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[async_trait]
|
||||||
|
impl Node for Handler {
|
||||||
|
async fn process(&self, runtime: Runtime, req: Message) -> Result<()> {
|
||||||
|
let typ = req.get_type();
|
||||||
|
|
||||||
|
match typ {
|
||||||
|
"broadcast" => {
|
||||||
|
let val = req.body.extra.get("message").and_then(|v| v.as_i64());
|
||||||
|
if let Some(val) = val {
|
||||||
|
let frm = req.src.as_str();
|
||||||
|
let mid = req.body.msg_id;
|
||||||
|
let key = format!("{frm}{mid}");
|
||||||
|
{
|
||||||
|
// there's only one thread, safe to unwrap
|
||||||
|
let mut guard = self.store.lock().unwrap();
|
||||||
|
guard.borrow_mut().entry(key).or_insert(val);
|
||||||
|
}
|
||||||
|
let body = MessageBody::new().with_type("broadcast_ok");
|
||||||
|
return runtime.reply(req, body).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
"read" => {
|
||||||
|
let vals = {
|
||||||
|
let g = self.store.lock().unwrap();
|
||||||
|
g.values().cloned().collect::<Vec<_>>()
|
||||||
|
};
|
||||||
|
let body = MessageBody::from_extra(
|
||||||
|
[("messages".to_string(), vals.into())]
|
||||||
|
.into_iter()
|
||||||
|
.collect(),
|
||||||
|
)
|
||||||
|
.with_type("read_ok");
|
||||||
|
return runtime.reply(req, body).await;
|
||||||
|
}
|
||||||
|
"topology" => {
|
||||||
|
let body = MessageBody::new().with_type("topology_ok");
|
||||||
|
return runtime.reply(req, body).await;
|
||||||
|
}
|
||||||
|
_ => {}
|
||||||
|
}
|
||||||
|
|
||||||
|
done(runtime, req)
|
||||||
|
}
|
||||||
|
}
|
|
@ -8,5 +8,5 @@ license-file.workspace = true
|
||||||
[dependencies]
|
[dependencies]
|
||||||
async-trait.workspace = true
|
async-trait.workspace = true
|
||||||
maelstrom-node.workspace = true
|
maelstrom-node.workspace = true
|
||||||
serde_json = "1.0.117"
|
serde_json.workspace = true
|
||||||
tokio.workspace = true
|
tokio.workspace = true
|
||||||
|
|
Loading…
Reference in a new issue