From c0d72a631cd66a2c1fc308561a5a9b0588283a2b Mon Sep 17 00:00:00 2001 From: Thilo Behnke Date: Thu, 12 May 2022 22:55:29 +0200 Subject: [PATCH] issue with sending event writer between threads --- pong/src/event.rs | 11 +++++------ pong/src/game_field.rs | 2 +- pong/src/pong.rs | 4 ++-- server/Cargo.toml | 1 + server/src/http.rs | 28 +++++++++++++++++++++++----- server/src/kafka.rs | 32 ++++++++++++++++++++++++++++++++ server/src/main.rs | 1 + 7 files changed, 65 insertions(+), 14 deletions(-) create mode 100644 server/src/kafka.rs diff --git a/pong/src/event.rs b/pong/src/event.rs index 6745e45..9494a48 100644 --- a/pong/src/event.rs +++ b/pong/src/event.rs @@ -10,12 +10,12 @@ pub mod event { } pub trait EventWriterImpl { - fn write(&self, event: Event) -> Result<(), ()>; + fn write(&mut self, event: Event) -> Result<(), ()>; } pub struct FileEventWriterImpl {} impl EventWriterImpl for FileEventWriterImpl { - fn write(&self, event: Event) -> Result<(), ()> { + fn write(&mut self, event: Event) -> Result<(), ()> { let options = OpenOptions::new().read(true).create(true).write(true).open("events.log"); if let Err(_) = options { return Err(()); @@ -30,7 +30,7 @@ pub mod event { pub struct NoopEventWriterImpl {} impl EventWriterImpl for NoopEventWriterImpl { - fn write(&self, event: Event) -> Result<(), ()> { + fn write(&mut self, event: Event) -> Result<(), ()> { todo!() } } @@ -40,7 +40,7 @@ pub mod event { } impl EventWriter { - fn new(writer_impl: Box) -> EventWriter { + pub fn new(writer_impl: Box) -> EventWriter { EventWriter { writer_impl } @@ -57,9 +57,8 @@ pub mod event { writer_impl: Box::new(FileEventWriterImpl {}) } } - // TODO: Kafka - pub fn write(&self, event: Event) -> Result<(), ()> { + pub fn write(&mut self, event: Event) -> Result<(), ()> { self.writer_impl.write(event) } } diff --git a/pong/src/game_field.rs b/pong/src/game_field.rs index dcd850f..3313965 100644 --- a/pong/src/game_field.rs +++ b/pong/src/game_field.rs @@ -115,7 +115,7 @@ impl Field { self.objs.push(Rc::new(RefCell::new(ball))); } - pub fn tick(&self, inputs: Vec) { + pub fn tick(&mut self, inputs: Vec) { for obj in self.objs.iter() { let mut obj_mut = RefCell::borrow_mut(obj); if obj_mut.obj_type() != "ball" { diff --git a/pong/src/pong.rs b/pong/src/pong.rs index 5162d2b..4b0b575 100644 --- a/pong/src/pong.rs +++ b/pong/src/pong.rs @@ -83,7 +83,7 @@ pub mod pong_events { } pub trait PongEventWriter { - fn write(&self, event: PongEventType) -> Result<(), ()>; + fn write(&mut self, event: PongEventType) -> Result<(), ()>; } pub struct DefaultPongEventWriter { @@ -91,7 +91,7 @@ pub mod pong_events { } impl PongEventWriter for DefaultPongEventWriter { - fn write(&self, event: PongEventType) -> Result<(), ()> { + fn write(&mut self, event: PongEventType) -> Result<(), ()> { let out_event = match event { PongEventType::GameObjUpdate(ref update) => { Event { diff --git a/server/Cargo.toml b/server/Cargo.toml index 986f069..f1041d1 100644 --- a/server/Cargo.toml +++ b/server/Cargo.toml @@ -10,3 +10,4 @@ kafka = { version = "0.8.0" } hyper = {version = "0.14.18", features = ["full"]} tokio = { version = "1", features = ["full"] } tokio-stream = {version = "0.1" } +pong = { path = "../pong", version = "0.1.0" } diff --git a/server/src/http.rs b/server/src/http.rs index b473dec..708065d 100644 --- a/server/src/http.rs +++ b/server/src/http.rs @@ -1,20 +1,33 @@ use std::convert::Infallible; +use std::sync::Arc; use hyper::{Body, Request, Response, Server}; use hyper::server::conn::AddrStream; use hyper::service::{make_service_fn, service_fn}; +use kafka::producer::Producer; +use tokio::sync::Mutex; +use pong::event::event::{Event, EventWriter}; +use crate::kafka::KafkaEventWriterImpl; pub struct HttpServer { addr: [u8; 4], - port: u16 + port: u16, + event_writer: Arc> } impl HttpServer { pub fn new(addr: [u8; 4], port: u16) -> HttpServer { - HttpServer {addr, port} + let event_writer = Arc::new(Mutex::new(EventWriter::new(Box::new(KafkaEventWriterImpl::default())))); + HttpServer {addr, port, event_writer} } pub async fn run(&self) -> Result<(), Box> { + let mut event_writer = Arc::clone(&self.event_writer); let make_svc = make_service_fn(|socket: &AddrStream| async { - Ok::<_, Infallible>(service_fn(handle_request)) + Ok::<_, Infallible>(service_fn(move |req: Request| { + async move { + let mut event_writer = Arc::clone(&event_writer); + handle_request(&event_writer, req).await + } + })) }); let host = (self.addr, self.port).into(); @@ -26,10 +39,15 @@ impl HttpServer { } } -async fn handle_request(req: Request) -> Result, Infallible> { - Ok(Response::new("hello".into())) +async fn handle_request(event_writer: &Arc>, req: Request) -> Result, Infallible> { + let mut locked = event_writer.lock().await; + if let err = locked.write(Event {topic: "topic".into(), key: "key".into(), msg: "msg".into()}) { + println!("Failed to write to kafka! {:?}", err); + } + Ok(Response::new("response".into())) } + async fn shutdown_signal() { // Wait for the CTRL+C signal tokio::signal::ctrl_c() diff --git a/server/src/kafka.rs b/server/src/kafka.rs new file mode 100644 index 0000000..016c63d --- /dev/null +++ b/server/src/kafka.rs @@ -0,0 +1,32 @@ +use std::time::Duration; +use kafka::producer::{Producer, Record, RequiredAcks}; +use pong::event::event::{Event, EventWriter, EventWriterImpl}; + +pub struct KafkaEventWriterImpl { + producer: Producer +} +impl KafkaEventWriterImpl { + pub fn default() -> KafkaEventWriterImpl { + KafkaEventWriterImpl::new("localhost:9092") + } + + pub fn new(host: &str) -> KafkaEventWriterImpl { + let mut producer = Producer::from_hosts(vec![host.to_owned()]) + .with_ack_timeout(Duration::from_secs(1)) + .with_required_acks(RequiredAcks::One) + .create() + .unwrap(); + KafkaEventWriterImpl { + producer + } + } +} +impl EventWriterImpl for KafkaEventWriterImpl { + fn write(&mut self, event: Event) -> Result<(), ()> { + let record = Record::from_key_value(event.topic.as_str(), event.key.as_str(), event.msg.as_str()); + match self.producer.send(&record) { + Ok(()) => Ok(()), + Err(_) => Err(()) + } + } +} diff --git a/server/src/main.rs b/server/src/main.rs index 81ae76a..6cfed44 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -1,6 +1,7 @@ use crate::http::HttpServer; mod http; +mod kafka; #[tokio::main] pub async fn main() {