diff --git a/Cargo.toml b/Cargo.toml index 44d3f67..7217af9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -5,7 +5,7 @@ authors = ["Thilo Behnke "] edition = "2018" [workspace] -members = ["pong"] +members = ["pong", "server"] [lib] crate-type = ["cdylib", "rlib"] diff --git a/pong/Cargo.toml b/pong/Cargo.toml index 9ff1639..2919920 100644 --- a/pong/Cargo.toml +++ b/pong/Cargo.toml @@ -6,6 +6,8 @@ edition = "2021" # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html [dependencies] +serde = { version = "1.0", features = ["derive"] } +serde_json = { version = "1.0.79" } [dev-dependencies] rstest = "0.12.0" diff --git a/pong/src/event.rs b/pong/src/event.rs new file mode 100644 index 0000000..48f8769 --- /dev/null +++ b/pong/src/event.rs @@ -0,0 +1,88 @@ + +pub mod event { + use std::fmt::Debug; + use std::fs::OpenOptions; + use std::io::Write; + use serde::{Deserialize, Serialize}; + + #[derive(Debug, Deserialize, Serialize)] + pub struct Event { + pub topic: String, + pub key: Option, + pub msg: String + } + + pub trait EventWriterImpl : Send + Sync { + fn write(&mut self, event: Event) -> Result<(), String>; + } + + pub struct FileEventWriterImpl {} + impl EventWriterImpl for FileEventWriterImpl { + fn write(&mut self, event: Event) -> Result<(), String> { + let options = OpenOptions::new().read(true).create(true).write(true).open("events.log"); + if let Err(e) = options { + return Err(format!("{}", e)); + } + let mut file = options.unwrap(); + match file.write(event.msg.as_bytes()) { + Ok(_) => Ok(()), + Err(e) => Err(format!("{}", e)) + } + } + } + + pub struct NoopEventWriterImpl {} + impl EventWriterImpl for NoopEventWriterImpl { + fn write(&mut self, event: Event) -> Result<(), String> { + todo!() + } + } + + pub struct EventWriter { + writer_impl: Box + } + + impl EventWriter { + pub fn new(writer_impl: Box) -> EventWriter { + EventWriter { + writer_impl + } + } + + pub fn noop() -> EventWriter { + EventWriter { + writer_impl: Box::new(NoopEventWriterImpl {}) + } + } + + pub fn file() -> EventWriter { + EventWriter { + writer_impl: Box::new(FileEventWriterImpl {}) + } + } + + pub fn write(&mut self, event: Event) -> Result<(), String> { + self.writer_impl.write(event) + } + } + + pub trait EventReaderImpl : Send + Sync { + fn read(&mut self) -> Result, String>; + } + + pub struct EventReader { + reader_impl: Box + } + + impl EventReader { + pub fn new(reader_impl: Box) -> EventReader { + EventReader { + reader_impl + } + } + + pub fn read(&mut self) -> Result, String> { + self.reader_impl.read() + } + } +} diff --git a/pong/src/game_field.rs b/pong/src/game_field.rs index 6a2b5d1..3313965 100644 --- a/pong/src/game_field.rs +++ b/pong/src/game_field.rs @@ -4,6 +4,7 @@ use crate::game_object::game_object::{DefaultGameObject, GameObject}; use crate::geom::geom::Vector; use crate::geom::shape::{Shape, ShapeType}; use crate::pong::pong_collisions::{handle_ball_bounds_collision, handle_player_ball_collision, handle_player_bound_collision}; +use crate::pong::pong_events::{PongEventWriter, DefaultPongEventWriter, NoopPongEventWriter, PongEventType, GameObjUpdate}; use crate::utils::utils::{DefaultLoggerFactory, Logger, LoggerFactory, NoopLogger}; use std::borrow::{Borrow, BorrowMut}; use std::cell::{Cell, Ref, RefCell, RefMut}; @@ -30,12 +31,13 @@ pub struct Field { pub height: u16, pub collisions: Box, objs: Vec>>>, + event_writer: Box, collision_detector: CollisionDetector, collision_handler: CollisionHandler, } impl Field { - pub fn new(logger_factory: Box) -> Field { + pub fn new(logger_factory: Box, event_writer: Box) -> Field { let width = 800; let height = 600; @@ -50,6 +52,7 @@ impl Field { collisions: Box::new(Collisions::new(vec![])), collision_detector: CollisionDetector::new(&logger_factory), collision_handler: CollisionHandler::new(&logger_factory), + event_writer, logger_factory }; @@ -85,6 +88,7 @@ impl Field { pub fn mock(width: u16, height: u16) -> Field { let logger_factory = DefaultLoggerFactory::new(Box::new(NoopLogger{})); + let event_writer = NoopPongEventWriter::new(); Field { logger: logger_factory.get("game_field"), width, @@ -96,6 +100,7 @@ impl Field { collisions: Box::new(Collisions::new(vec![])), collision_detector: CollisionDetector::new(&logger_factory), collision_handler: CollisionHandler::new(&logger_factory), + event_writer, logger_factory, } } @@ -110,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" { @@ -164,6 +169,14 @@ impl Field { let obj_b = objs.iter().find(|o| RefCell::borrow(o).id() == collision.1).unwrap().clone(); collision_handler.handle(obj_a, obj_b); } + + { + for obj in self.objs.iter().filter(|o| RefCell::borrow(o).is_dirty()) { + let mut obj = RefCell::borrow_mut(obj); + self.event_writer.write(PongEventType::GameObjUpdate(GameObjUpdate{obj_id: &obj.id().to_string(), vel: obj.vel(), orientation: obj.orientation(), pos: obj.pos()})); + obj.set_dirty(false); + } + } } fn get_collisions(&self) -> Box { diff --git a/pong/src/game_object.rs b/pong/src/game_object.rs index 632568a..b2f2157 100644 --- a/pong/src/game_object.rs +++ b/pong/src/game_object.rs @@ -17,6 +17,8 @@ pub mod game_object { fn vel(&self) -> &Vector; fn vel_mut(&mut self) -> &mut Vector; fn is_static(&self) -> bool; + fn is_dirty(&self) -> bool; + fn set_dirty(&mut self, is_dirty: bool); } // #[derive(Clone, Debug, PartialEq)] @@ -26,6 +28,7 @@ pub mod game_object { pub obj_type: String, geom: Box, physics: Box, + dirty: bool } impl DefaultGameObject { @@ -40,6 +43,7 @@ pub mod game_object { obj_type, geom, physics, + dirty: false } } } @@ -74,18 +78,20 @@ pub mod game_object { } fn update_pos(&mut self) { + // Keep last orientation if vel is now zero. + if self.vel() == &Vector::zero() { + return; + } let vel = self.vel().clone(); let center = self.geom.center_mut(); center.add(&vel); - // Keep last orientation if vel is now zero. - if vel == Vector::zero() { - return; - } let mut updated_orientation = vel.clone(); updated_orientation.normalize(); let orientation = self.geom.orientation_mut(); orientation.x = updated_orientation.x; orientation.y = updated_orientation.y; + + self.dirty = true; } fn bounding_box(&self) -> BoundingBox { @@ -103,6 +109,14 @@ pub mod game_object { fn is_static(&self) -> bool { self.physics.is_static() } + + fn is_dirty(&self) -> bool { + return self.dirty + } + + fn set_dirty(&mut self, is_dirty: bool) { + self.dirty = is_dirty; + } } } diff --git a/pong/src/geom.rs b/pong/src/geom.rs index 6393101..711def8 100644 --- a/pong/src/geom.rs +++ b/pong/src/geom.rs @@ -1,5 +1,7 @@ pub mod geom { - #[derive(Debug, Clone)] + use serde::{Serialize}; + + #[derive(Debug, Clone, Serialize)] pub struct Vector { pub x: f64, pub y: f64, diff --git a/pong/src/lib.rs b/pong/src/lib.rs index c1c5724..0ba1216 100644 --- a/pong/src/lib.rs +++ b/pong/src/lib.rs @@ -4,3 +4,4 @@ pub mod game_object; pub mod geom; pub mod pong; pub mod utils; +pub mod event; diff --git a/pong/src/pong.rs b/pong/src/pong.rs index 4ca8472..34bff2b 100644 --- a/pong/src/pong.rs +++ b/pong/src/pong.rs @@ -25,6 +25,8 @@ pub mod pong_collisions { b_to_a.sub(&player.pos()); b_to_a.normalize(); ball.pos_mut().add(&b_to_a); + + ball.set_dirty(true); } pub fn handle_ball_bounds_collision( @@ -34,6 +36,8 @@ pub mod pong_collisions { let mut ball = RefCell::borrow_mut(&ball); let bound = RefCell::borrow(&bound); ball.vel_mut().reflect(&bound.orientation()); + + ball.set_dirty(true); } pub fn handle_player_bound_collision( @@ -54,5 +58,67 @@ pub mod pong_collisions { new_pos.add(&perpendicular); let player_pos = player.pos_mut(); player_pos.y = new_pos.y; + + player.set_dirty(true); + } +} + +pub mod pong_events { + use serde_json::json; + use serde::{Serialize}; + use crate::event::event::{Event, EventWriter}; + use crate::geom::geom::Vector; + + #[derive(Serialize)] + pub enum PongEventType<'a> { + GameObjUpdate(GameObjUpdate<'a>) + } + + #[derive(Serialize)] + pub struct GameObjUpdate<'a> { + pub obj_id: &'a str, + pub pos: &'a Vector, + pub vel: &'a Vector, + pub orientation: &'a Vector + } + + pub trait PongEventWriter { + fn write(&mut self, event: PongEventType) -> Result<(), String>; + } + + pub struct DefaultPongEventWriter { + writer: EventWriter + } + + impl PongEventWriter for DefaultPongEventWriter { + fn write(&mut self, event: PongEventType) -> Result<(), String> { + let out_event = match event { + PongEventType::GameObjUpdate(ref update) => { + Event { + topic: String::from("obj_update"), + key: Some(update.obj_id.clone().to_string()), + msg: serde_json::to_string(&event).unwrap() + } + } + }; + self.writer.write(out_event) + } + } + + pub struct NoopPongEventWriter {} + impl NoopPongEventWriter { + pub fn new() -> Box { + Box::new(DefaultPongEventWriter { + writer: EventWriter::noop() + }) + } + } + + impl DefaultPongEventWriter { + pub fn new() -> Box { + Box::new(DefaultPongEventWriter { + writer: EventWriter::file() + }) + } } } diff --git a/server/.env b/server/.env new file mode 100644 index 0000000..6d7b559 --- /dev/null +++ b/server/.env @@ -0,0 +1,2 @@ +KAFKA_HOST=localhost +KAFKA_PORT=9092 diff --git a/server/Cargo.toml b/server/Cargo.toml new file mode 100644 index 0000000..a7f8428 --- /dev/null +++ b/server/Cargo.toml @@ -0,0 +1,21 @@ +[package] +name = "server" +version = "0.1.0" +edition = "2021" + +# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html + +[dependencies] +kafka = { version = "0.8.0" } +hyper = {version = "0.14.18", features = ["full"]} +tokio = { version = "1", features = ["full"] } +tokio-stream = {version = "0.1" } +serde = { version = "1.0", features = ["derive"] } +serde_json = { version = "1.0.79" } +md5 = { version = "0.7.0" } +pong = { path = "../pong", version = "0.1.0" } +hyper-tungstenite = "0.8.0" +futures = { version = "0.3.12" } + +[dev-dependencies] +rstest = "0.12.0" diff --git a/server/docker-compose.yml b/server/docker-compose.yml new file mode 100644 index 0000000..3217f09 --- /dev/null +++ b/server/docker-compose.yml @@ -0,0 +1,37 @@ +version: "3" +services: + zookeeper: + container_name: pong_server_zookeeper + image: 'bitnami/zookeeper:latest' + ports: + - '2181:2181' + environment: + - ALLOW_ANONYMOUS_LOGIN=yes + kafka: + container_name: pong_server_kafka + build: + context: kafka + ports: + - '9092:9092' + - '9093:9093' + - '7243:7243' + environment: + - KAFKA_BROKER_ID=1 + - KAFKA_CFG_INTER_BROKER_LISTENER_NAME=DOCKER + - KAFKA_CFG_LISTENERS=LOCAL://:9093,DOCKER://:9092 + - KAFKA_CFG_ADVERTISED_LISTENERS=LOCAL://127.0.0.1:9093,DOCKER://kafka:9092 + - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=LOCAL:PLAINTEXT,DOCKER:PLAINTEXT + - KAFKA_CFG_ZOOKEEPER_CONNECT=zookeeper:2181 + - ALLOW_PLAINTEXT_LISTENER=yes + depends_on: + - zookeeper +# kaka-rest-proxy: +# container_name: pong_server_kafka_rest_proxy +# image: 'confluentinc/cp-kafka-rest' +# ports: +# - '8082:8082' +# environment: +# - KAFKA_REST_BOOTSTRAP_SERVERS=kafka:9092 +# depends_on: +# - kafka + diff --git a/server/kafka/Dockerfile b/server/kafka/Dockerfile new file mode 100644 index 0000000..7ff014d --- /dev/null +++ b/server/kafka/Dockerfile @@ -0,0 +1,17 @@ +FROM rust:1.54 as build-stage + +RUN mkdir -p /opt/pong/kafka-script-proxy +WORKDIR /opt/pong/kafka-script-proxy +ADD kafka-script-proxy ./ +RUN cargo build --release + +FROM bitnami/kafka:latest + +ADD custom-entrypoint.sh / +COPY --from=build-stage /opt/pong/kafka-script-proxy/target/release/kafka-script-proxy /bin/ + +USER root +RUN mkdir -p /var/log/kafka-script-proxy + +ENTRYPOINT [] +CMD [ "/custom-entrypoint.sh" ] diff --git a/server/kafka/custom-entrypoint.sh b/server/kafka/custom-entrypoint.sh new file mode 100755 index 0000000..ff95afa --- /dev/null +++ b/server/kafka/custom-entrypoint.sh @@ -0,0 +1,4 @@ +#!/usr/bin/env bash + +/bin/kafka-script-proxy & +/opt/bitnami/scripts/kafka/entrypoint.sh "/opt/bitnami/scripts/kafka/run.sh" diff --git a/server/kafka/kafka-script-proxy/.dockerignore b/server/kafka/kafka-script-proxy/.dockerignore new file mode 100644 index 0000000..a09c56d --- /dev/null +++ b/server/kafka/kafka-script-proxy/.dockerignore @@ -0,0 +1 @@ +/.idea diff --git a/server/kafka/kafka-script-proxy/.gitignore b/server/kafka/kafka-script-proxy/.gitignore new file mode 100644 index 0000000..6db043d --- /dev/null +++ b/server/kafka/kafka-script-proxy/.gitignore @@ -0,0 +1,2 @@ +.idea +/target diff --git a/server/kafka/kafka-script-proxy/Cargo.toml b/server/kafka/kafka-script-proxy/Cargo.toml new file mode 100644 index 0000000..5e03ca8 --- /dev/null +++ b/server/kafka/kafka-script-proxy/Cargo.toml @@ -0,0 +1,13 @@ + +[package] +name = "kafka-script-proxy" +version = "0.1.0" +edition = "2018" + +[workspace] + +# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html + +[dependencies] +hyper = {version = "0.14.18", features = ["full"]} +tokio = { version = "1", features = ["full"] } diff --git a/server/kafka/kafka-script-proxy/src/main.rs b/server/kafka/kafka-script-proxy/src/main.rs new file mode 100644 index 0000000..c0ddf71 --- /dev/null +++ b/server/kafka/kafka-script-proxy/src/main.rs @@ -0,0 +1,151 @@ +extern crate core; + +use hyper::service::{make_service_fn, service_fn}; +use hyper::{Body, Method, Request, Response, Server, StatusCode}; +use std::convert::Infallible; +use std::process::Output; +use tokio::fs::OpenOptions; +use tokio::io::AsyncWriteExt; +use tokio::process::Command; + +const TOPICS: [&str; 3] = ["move", "status", "input"]; + +#[tokio::main] +pub async fn main() { + run().await.unwrap() +} + +pub async fn run() -> Result<(), ()> { + let make_svc = make_service_fn(|_| async { + Ok::<_, Infallible>(service_fn(|req: Request| async { + handle_request(req).await + })) + }); + + let host = ([0, 0, 0, 0], 7243).into(); + let server = Server::bind(&host).serve(make_svc); + write_to_log(&format!("Listening on http://{}", host)).await; + let graceful = server.with_graceful_shutdown(shutdown_signal()); + graceful.await.unwrap(); + Ok(()) +} + +async fn handle_request(req: Request) -> Result, Infallible> { + write_to_log(&format!( + "req to {} with method {}", + req.uri().path(), + req.method() + )) + .await; + match (req.method(), req.uri().path()) { + (&Method::POST, "/add_partition") => handle_add_partition().await, + _ => build_error_res("not found", StatusCode::NOT_FOUND), + } +} + +async fn handle_add_partition() -> Result, Infallible> { + write_to_log("Called to add partition.").await; + let current_count = get_highest_partition_count().await; + if let Err(_) = current_count { + let err = "Failed to retrieve max partition count."; + write_to_log(err).await; + return build_error_res(err, StatusCode::INTERNAL_SERVER_ERROR); + } + let current_count = current_count.unwrap(); + write_to_log(&format!("Successfully retrieved current max partition count: {}.", current_count)).await; + + let next_partition = current_count + 1; + write_to_log(&format!("Updating partition count to {} for the following topics: {}", next_partition, TOPICS.join(","))).await; + for topic in TOPICS { + let output = run_command(&format!("/opt/bitnami/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --alter --topic {} --partitions {}", topic, next_partition)).await; + if let Err(e) = output { + let error = format!("Failed to update the partition count: {}", e); + write_to_log(&error).await; + return build_error_res(&error, StatusCode::INTERNAL_SERVER_ERROR); + } + let output = output.unwrap(); + if !output.status.success() { + let error = format!("Failed to update the partition count: {:?}", std::str::from_utf8(&*output.stderr).unwrap()); + write_to_log(&error).await; + return build_error_res(&error, StatusCode::INTERNAL_SERVER_ERROR); + } + } + write_to_log(&format!("Successfully updated partition count to {}", next_partition)).await; + return build_success_res(¤t_count.to_string()); // current_count because count - 1 = max partition +} + +async fn get_highest_partition_count() -> Result { + let output = run_command("/opt/bitnami/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe | grep -Po '(?<=PartitionCount: )(\\d+)' | sort -r | head -1").await.unwrap(); + let output_str = std::str::from_utf8(&*output.stdout); + if let Err(e) = output_str { + let message = format!("Failed to convert command output to string: {:?}", e); + write_to_log(&message).await; + return Err(PartitionCountQueryError { message }); + } + + let output_str = output_str.unwrap().trim().replace("\n", ""); + let parse_res = output_str.parse::(); + if let Err(e) = parse_res { + let message = format!("Failed to parse partition count for output {}: {:?}", output_str, e); + write_to_log(&message).await; + return Err(PartitionCountQueryError {message}) + } + return Ok(parse_res.unwrap()) +} + +#[derive(Debug)] +struct PartitionCountQueryError { + message: String +} + +async fn run_command(command: &str) -> std::io::Result { + write_to_log(&format!("Running command: {}", command)).await; + let output = Command::new("/bin/bash") + .arg("-c") + .arg(command) + .output() + .await.unwrap(); + let stdout = std::str::from_utf8(&*output.stdout).unwrap(); + let stderr = std::str::from_utf8(&*output.stderr).unwrap(); + write_to_log(&format!("Command returned stdout: {}", stdout)).await; + write_to_log(&format!("Command returned stderr: {}", stderr)).await; + Ok(output) +} + +pub fn build_success_res(value: &str) -> Result, Infallible> { + let json = format!("{{\"data\": {}}}", value); + return Ok(Response::new(Body::from(json))); +} + +pub fn build_error_res(error: &str, status: StatusCode) -> Result, Infallible> { + let json = format!("{{\"error\": \"{}\"}}", error); + let mut res = Response::new(Body::from(json)); + *res.status_mut() = status; + return Ok(res); +} + +pub async fn shutdown_signal() { + // Wait for the CTRL+C signal + tokio::signal::ctrl_c() + .await + .expect("failed to install CTRL+C signal handler"); +} + +async fn write_to_log(value: &str) { + let content = format!("{}\n", value); + let mut options = OpenOptions::new(); + let file = options + .create(true) + .write(true) + .append(true) + .open("/var/log/kafka-script-proxy/log.txt") + .await; + if let Err(e) = file { + println!("Failed to open file: {}", e); + return; + } + let write_res = file.unwrap().write_all(content.as_bytes()).await; + if let Err(e) = write_res { + println!("Failed to write to file: {:?}", e); + } +} diff --git a/server/run.sh b/server/run.sh new file mode 100755 index 0000000..5d00b19 --- /dev/null +++ b/server/run.sh @@ -0,0 +1,10 @@ +#!/usr/bin/env bash + +source .env + +docker-compose down +docker-compose up -d --build --force-recreate +docker exec pong_server_kafka /opt/bitnami/kafka/bin/kafka-topics.sh --create --topic session --bootstrap-server "$KAFKA_HOST:$KAFKA_PORT" +docker exec pong_server_kafka /opt/bitnami/kafka/bin/kafka-topics.sh --create --topic move --bootstrap-server "$KAFKA_HOST:$KAFKA_PORT" +docker exec pong_server_kafka /opt/bitnami/kafka/bin/kafka-topics.sh --create --topic status --bootstrap-server "$KAFKA_HOST:$KAFKA_PORT" +docker exec pong_server_kafka /opt/bitnami/kafka/bin/kafka-topics.sh --create --topic input --bootstrap-server "$KAFKA_HOST:$KAFKA_PORT" diff --git a/server/src/hash.rs b/server/src/hash.rs new file mode 100644 index 0000000..7cf4c37 --- /dev/null +++ b/server/src/hash.rs @@ -0,0 +1,9 @@ +use std::ptr::hash; + +pub struct Hasher {} +impl Hasher { + pub fn hash(n: u16) -> String { + let digest = md5::compute(format!("{}", n)); + format!("{:x}", digest) + } +} diff --git a/server/src/http.rs b/server/src/http.rs new file mode 100644 index 0000000..d5d307f --- /dev/null +++ b/server/src/http.rs @@ -0,0 +1,250 @@ +use std::convert::Infallible; +use std::fs::read; +use std::io::ErrorKind::NotFound; +use std::net::SocketAddr; +use std::sync::Arc; +use hyper::{Body, body, Method, Request, Response, Server, StatusCode}; +use hyper::server::conn::AddrStream; +use hyper::service::{make_service_fn, service_fn}; +use hyper_tungstenite::HyperWebsocket; +use hyper_tungstenite::tungstenite::{Error, Message}; +use kafka::producer::Producer; +use serde_json::json; +use serde::{Deserialize, Serialize}; +use tokio::sync::Mutex; +use pong::event::event::{Event, EventReader, EventWriter}; +use futures::{sink::SinkExt, stream::StreamExt}; +use crate::kafka::{KafkaEventReaderImpl, KafkaSessionEventWriterImpl}; +use crate::player::Player; +use crate::session::{CachingSessionManager, SessionManager}; +use crate::utils::http_utils::{get_query_params, read_json_body}; + +pub struct HttpServer { + addr: [u8; 4], + port: u16, + session_manager: Arc> +} +impl HttpServer { + pub fn new(addr: [u8; 4], port: u16, kafka_host: &str) -> HttpServer { + let session_manager = Arc::new(Mutex::new(CachingSessionManager::new(kafka_host))); + HttpServer {addr, port, session_manager} + } + + pub async fn run(self) -> Result<(), Box> { + let make_svc = make_service_fn(|socket: &AddrStream| { + let mut session_manager = Arc::clone(&self.session_manager); + let addr = socket.remote_addr(); + async move { + Ok::<_, Infallible>(service_fn(move |req: Request| { + let mut session_manager = Arc::clone(&session_manager); + async move { + if hyper_tungstenite::is_upgrade_request(&req) { + let (response, websocket) = hyper_tungstenite::upgrade(req, None).unwrap(); + + // Spawn a task to handle the websocket connection. + tokio::spawn(async move { + if let Err(e) = serve_websocket(websocket, session_manager).await { + eprintln!("Error in websocket connection: {}", e); + } + }); + + // Return the response so the spawned future can continue. + return Ok(response) + } + + return handle_request(&session_manager, req, addr).await; + } + })) + } + }); + + let host = (self.addr, self.port).into(); + let server = Server::bind(&host).serve(make_svc); + println!("Listening on http://{}", host); + let graceful = server.with_graceful_shutdown(shutdown_signal()); + graceful.await?; + Ok(()) + } +} + +/// Handle a websocket connection. +async fn serve_websocket(websocket: HyperWebsocket, session_manager: Arc>) -> Result<(), Error> { + let mut websocket = websocket.await?; + while let Some(message) = websocket.next().await { + match message? { + Message::Text(msg) => { + let event = serde_json::from_str::(&msg); + if let Err(e) = event { + eprintln!("Failed to deserialize ws message to event {}: {}", msg, e); + continue; + } + let event = event.unwrap(); + let mut locked = session_manager.lock().await; + let writer = locked.get_session_writer(&event.session_id); + if let Err(e) = writer { + eprintln!("Failed to retrieve session writer for session {}: {}", event.session_id, e); + continue; + } + let writer = writer.unwrap(); + let mut writer = writer.lock().await; + let write_res = writer.write_to_session(&event.topic, &event.msg); + if let Err(e) = write_res { + eprintln!("Failed to write event {:?}: {}", event, e); + websocket.send(Message::text("Failed to write message")).await?; + continue; + } + websocket.send(Message::text("Message wrote successfully")).await?; + }, + Message::Close(msg) => { + // No need to send a reply: tungstenite takes care of this for you. + if let Some(msg) = &msg { + println!("Received close message with code {} and message: {}", msg.code, msg.reason); + } else { + println!("Received close message"); + } + }, + _ => {} + } + } + + Ok(()) +} + +// TODO: How to handle event writes/reads? This must be a websocket, but how to implement in hyper (if possible)? +// https://github.com/de-vri-es/hyper-tungstenite-rs +async fn handle_request(session_manager: &Arc>, req: Request, addr: SocketAddr) -> Result, Infallible> { + println!("req to {} with method {}", req.uri().path(), req.method()); + match (req.method(), req.uri().path()) { + (&Method::GET, "/session") => handle_get_session(session_manager, req).await, + (&Method::POST, "/create_session") => handle_session_create(session_manager, req, addr).await, + (&Method::POST, "/join_session") => handle_session_join(session_manager, req, addr).await, + (&Method::POST, "/write") => handle_event_write(session_manager, req).await, + (&Method::POST, "/read") => handle_event_read(session_manager, req).await, + _ => Ok(Response::new("unknown".into())) + } +} + +async fn handle_get_session(session_manager: &Arc>, req: Request) -> Result, Infallible> { + let mut locked = session_manager.lock().await; + let query_params = get_query_params(&req); + let session_id = query_params.get("session_id"); + if let None = session_id { + return build_error_res("Please provide a valid session id", StatusCode::BAD_REQUEST); + } + let session_id = session_id.unwrap(); + let session = locked.get_session(session_id); + if let None = session { + return build_error_res("Unable to find session for given id", StatusCode::NOT_FOUND); + } + return build_success_res(&serde_json::to_string(&session.unwrap()).unwrap()); +} + +async fn handle_session_create(session_manager: &Arc>, req: Request, addr: SocketAddr) -> Result, Infallible> { + let mut locked = session_manager.lock().await; + let session_create_res = locked.create_session(Player {id: addr.to_string()}).await; + if let Err(e) = session_create_res { + return Ok(Response::builder().status(StatusCode::INTERNAL_SERVER_ERROR).body(Body::from(e)).unwrap()); + } + let serialized = json!(session_create_res.unwrap()); + return Ok(Response::new(Body::from(serialized.to_string()))) +} + +async fn handle_session_join(session_manager: &Arc>, mut req: Request, addr: SocketAddr) -> Result, Infallible> { + let mut locked = session_manager.lock().await; + let body = read_json_body::(&mut req).await; + let session_join_res = locked.join_session(body.session_id, Player {id: addr.to_string()}).await; + if let Err(e) = session_join_res { + return Ok(Response::builder().status(StatusCode::INTERNAL_SERVER_ERROR).body(Body::from(e)).unwrap()); + } + let serialized = json!(session_join_res.unwrap()); + return Ok(Response::new(Body::from(serialized.to_string()))) +} + +async fn handle_event_write(session_manager: &Arc>, mut req: Request) -> Result, Infallible> { + let mut locked = session_manager.lock().await; + let event = read_json_body::(&mut req).await; + let writer = locked.get_session_writer(&event.session_id); + if let Err(e) = writer { + let err = format!("Failed to write event: {}", e); + println!("{}", err); + let mut res = Response::new(Body::from(err)); + *res.status_mut() = StatusCode::NOT_FOUND; + return Ok(res); + } + let mut writer = writer.unwrap(); + let mut writer_locked = writer.lock().await; + println!("Writing session event to kafka: {:?}", event); + let write_res = writer_locked.write_to_session(&event.topic, &event.msg); + if let Err(e) = write_res { + let err = format!("Failed to write event: {}", e); + println!("{}", err); + let mut res = Response::new(Body::from(err)); + *res.status_mut() = StatusCode::INTERNAL_SERVER_ERROR; + return Ok(res); + } + println!("Successfully wrote event to kafka."); + build_success_res(&serde_json::to_string(&event).unwrap()) +} + +async fn handle_event_read(session_manager: &Arc>, mut req: Request) -> Result, Infallible> { + let mut locked = session_manager.lock().await; + let read_payload = read_json_body::(&mut req).await; + let reader = locked.get_session_reader(&read_payload.session_id); + if let Err(e) = reader { + let err = format!("Failed to read events: {}", e); + println!("{}", err); + let mut res = Response::new(Body::from(err)); + *res.status_mut() = StatusCode::NOT_FOUND; + return Ok(res); + } + let mut reader = reader.unwrap(); + let mut reader_locked = reader.lock().await; + println!("Reading session events from kafka for session: {}", read_payload.session_id); + let events = reader_locked.read_from_session(); + if let Err(e) = events { + let err = format!("Failed to read events: {}", e); + println!("{}", err); + let mut res = Response::new(Body::from(err)); + *res.status_mut() = StatusCode::INTERNAL_SERVER_ERROR; + return Ok(res); + } + println!("Successfully read session events from kafka."); + let json = serde_json::to_string(&events.unwrap()).unwrap(); + build_success_res(&json) +} + +pub fn build_success_res(value: &str) -> Result, Infallible> { + let json = format!("{{\"data\": {}}}", value); + return Ok(Response::new(Body::from(json))); +} + +pub fn build_error_res(error: &str, status: StatusCode) -> Result, Infallible> { + let json = format!("{{\"error\": \"{}\"}}", error); + let mut res = Response::new(Body::from(json)); + *res.status_mut() = status; + return Ok(res); +} + +async fn shutdown_signal() { + // Wait for the CTRL+C signal + tokio::signal::ctrl_c() + .await + .expect("failed to install CTRL+C signal handler"); +} + +#[derive(Debug, Deserialize, Serialize)] +struct SessionEventWriteDTO { + session_id: String, + topic: String, + msg: String +} + +#[derive(Debug, Serialize, Deserialize)] +struct SessionReadDTO { + session_id: String +} + +#[derive(Debug, Serialize, Deserialize)] +struct SessionJoinDto { + session_id: String +} diff --git a/server/src/kafka.rs b/server/src/kafka.rs new file mode 100644 index 0000000..a7dfa73 --- /dev/null +++ b/server/src/kafka.rs @@ -0,0 +1,262 @@ +use std::hash::{BuildHasher, Hash}; +use std::process::ExitStatus; +use std::str::FromStr; +use std::time::Duration; +use hyper::{Body, Client, Method, Request, Uri}; +use serde::{Deserialize}; +use kafka::client::{KafkaClient, ProduceMessage}; +use kafka::client::metadata::Topic; +use tokio::process::Command; +use kafka::consumer::{Consumer, FetchOffset, GroupOffsetStorage, MessageSet}; +use kafka::producer::{DefaultPartitioner, Partitioner, Producer, Record, RequiredAcks, Topics}; +use pong::event::event::{Event, EventReaderImpl, EventWriter, EventWriterImpl}; +use crate::hash::Hasher; +use crate::session::Session; + +pub struct KafkaSessionEventWriterImpl { + producer: Producer +} +impl KafkaSessionEventWriterImpl { + pub fn new(host: &str) -> KafkaSessionEventWriterImpl { + println!("Connecting session_writer producer to kafka host: {}", host); + let mut producer = Producer::from_hosts(vec![host.to_owned()]) + .with_ack_timeout(Duration::from_secs(1)) + .with_required_acks(RequiredAcks::One) + .with_partitioner(SessionPartitioner {}) + .create() + .unwrap(); + KafkaSessionEventWriterImpl { + producer + } + } +} + +pub struct KafkaDefaultEventWriterImpl { + producer: Producer +} +impl KafkaDefaultEventWriterImpl { + pub fn new(host: &str) -> KafkaDefaultEventWriterImpl { + println!("Connecting default producer to kafka host: {}", host); + let mut producer = Producer::from_hosts(vec![host.to_owned()]) + .with_ack_timeout(Duration::from_secs(1)) + .with_required_acks(RequiredAcks::One) + .create() + .unwrap(); + KafkaDefaultEventWriterImpl { + producer + } + } +} + + +impl EventWriterImpl for KafkaSessionEventWriterImpl { + fn write(&mut self, event: Event) -> Result<(), String> { + match event.key { + Some(key) => { + let record = Record::from_key_value(event.topic.as_str(), key, event.msg.as_str()); + match self.producer.send(&record) { + Ok(()) => Ok(()), + Err(e) => Err(format!("{}", e)) + } + }, + None => { + let record = Record::from_value(event.topic.as_str(), event.msg.as_str()); + match self.producer.send(&record) { + Ok(()) => Ok(()), + Err(e) => Err(format!("{}", e)) + } + } + } + } +} + +impl EventWriterImpl for KafkaDefaultEventWriterImpl { + fn write(&mut self, event: Event) -> Result<(), String> { + match event.key { + Some(key) => { + let record = Record::from_key_value(event.topic.as_str(), key, event.msg.as_str()); + match self.producer.send(&record) { + Ok(()) => Ok(()), + Err(e) => Err(format!("{}", e)) + } + }, + None => { + let record = Record::from_value(event.topic.as_str(), event.msg.as_str()); + match self.producer.send(&record) { + Ok(()) => Ok(()), + Err(e) => Err(format!("{}", e)) + } + } + } + } +} + +pub struct KafkaEventReaderImpl { + consumer: Consumer +} +impl KafkaEventReaderImpl { + pub fn default() -> KafkaEventReaderImpl { + KafkaEventReaderImpl::new("localhost:9093") + } + + pub fn from(host: &str) -> KafkaEventReaderImpl { + KafkaEventReaderImpl::new(host) + } + + pub fn new(host: &str) -> KafkaEventReaderImpl { + println!("Connecting consumer to kafka host: {}", host); + let mut consumer = Consumer::from_hosts(vec!(host.to_owned())) + .with_topic("move".to_owned()) + .with_topic("status".to_owned()) + .with_topic("input".to_owned()) + .with_fallback_offset(FetchOffset::Earliest) + .with_group("group".to_owned()) + .with_offset_storage(GroupOffsetStorage::Kafka) + .create() + .unwrap(); + KafkaEventReaderImpl { + consumer + } + } + + pub fn for_partitions(host: &str, partitions: &[i32], topics: &[&str]) -> KafkaEventReaderImpl { + println!("Connecting partition specific consumer to kafka host: {}", host); + let mut builder = Consumer::from_hosts(vec!(host.to_owned())); + for topic in topics.iter() { + builder = builder.with_topic_partitions(topic.parse().unwrap(), partitions); + } + builder = builder + .with_fallback_offset(FetchOffset::Earliest) + .with_group("group".to_owned()) + .with_offset_storage(GroupOffsetStorage::Kafka); + + let consumer = builder + .create() + .unwrap(); + KafkaEventReaderImpl { + consumer + } + } +} +impl EventReaderImpl for KafkaEventReaderImpl { + fn read(&mut self) -> Result, String> { + self.consume() + } +} + +impl KafkaEventReaderImpl { + fn consume(&mut self) -> Result, String> { + // TODO: How to best filter messages by key (= game session id?) + // E.g. https://docs.rs/kafka/latest/kafka/producer/struct.DefaultPartitioner.html - is it possible to read from partition by retrieving the hash of the key? + // Does it even make sense to hash the key if it already is a hash? Custom partitioner? + let polled = self.consumer.poll().unwrap(); + let message_sets: Vec> = polled.iter().collect(); + let mut events = vec![]; + for ms in message_sets { + let topic = ms.topic(); + let partition = ms.partition(); + println!("querying topic={} partition={}", topic, partition); + for m in ms.messages() { + let event = Event {topic: String::from(topic), key: Some(std::str::from_utf8(m.key).unwrap().parse().unwrap()), msg: std::str::from_utf8(m.value).unwrap().parse().unwrap() }; + events.push(event); + } + self.consumer.consume_messageset(ms).unwrap(); + } + self.consumer.commit_consumed().unwrap(); + Ok(events) + } +} + +pub struct KafkaSessionEventReaderImpl { + inner: KafkaEventReaderImpl +} + +impl KafkaSessionEventReaderImpl { + pub fn new(host: &str, session: &Session, topics: &[&str]) -> KafkaSessionEventReaderImpl { + let partitions = [session.id as i32]; + KafkaSessionEventReaderImpl { + inner: KafkaEventReaderImpl::for_partitions(host, &partitions, topics) + } + } +} + +impl EventReaderImpl for KafkaSessionEventReaderImpl { + fn read(&mut self) -> Result, String> { + self.inner.read() + } +} + + +#[derive(Debug)] +pub struct KafkaTopicManager { + partition_management_endpoint: String +} +impl KafkaTopicManager { + + pub fn default() -> KafkaTopicManager { + KafkaTopicManager {partition_management_endpoint: "http://localhost:7243/add_partition".to_owned()} + } + + pub fn from(topic_manager_host: &str) -> KafkaTopicManager { + KafkaTopicManager {partition_management_endpoint: format!("http://{}/add_partition", topic_manager_host).to_owned()} + } + + pub async fn add_partition(&self) -> Result { + let mut client = Client::new(); + let request = Request::builder().method(Method::POST).uri(Uri::from_str(&self.partition_management_endpoint).unwrap()).body(Body::empty()).unwrap(); + let res = client.request(request).await; + if let Err(e) = res { + let error = format!("Failed to add partition: {:?}", e); + println!("{}", error); + return Err(error); + } + let status = res.as_ref().unwrap().status(); + let bytes = hyper::body::to_bytes(res.unwrap()).await; + if let Err(e) = bytes { + let error = format!("Failed to read bytes from response: {:?}", e); + println!("{}", error); + return Err(error); + } + let bytes = bytes.unwrap().to_vec(); + let res_str = std::str::from_utf8(&*bytes); + if let Err(e) = res_str { + let error = format!("Failed to deserialize bytes to string: {:?}", e); + println!("{}", error); + return Err(error); + } + if status != 200 { + let error = format!("Failed to add partition: {}", res_str.unwrap()); + println!("{}", error); + return Err(error); + } + let json = serde_json::from_str::(res_str.unwrap()); + if let Err(e) = json { + let error = format!("Failed to convert string {} to json: {:?}", res_str.unwrap(), e); + println!("{}", error); + return Err(error); + } + let updated_partition_count = json.unwrap().data; + println!("Successfully created partition: {}", updated_partition_count); + Ok(updated_partition_count) + } +} + +#[derive(Deserialize)] +struct PartitionApiDTO { + data: u16 +} + +pub struct SessionPartitioner {} + +impl Partitioner for SessionPartitioner { + fn partition(&mut self, topics: Topics, msg: &mut ProduceMessage) { + match msg.key { + Some(key) => { + let key = std::str::from_utf8(key).unwrap(); + msg.partition = key.parse::().unwrap(); + println!("Overriding message partition with key: {}", msg.partition); + }, + None => panic!("Producing message without key not allowed!") + } + } +} diff --git a/server/src/main.rs b/server/src/main.rs new file mode 100644 index 0000000..e4fb33c --- /dev/null +++ b/server/src/main.rs @@ -0,0 +1,15 @@ +extern crate core; + +use crate::http::HttpServer; + +pub mod http; +pub mod kafka; +pub mod utils; +mod hash; +mod session; +mod player; + +#[tokio::main] +pub async fn main() { + HttpServer::new([127, 0, 0, 1], 4000, "localhost:9093").run().await.expect("failed to run server"); +} diff --git a/server/src/player.rs b/server/src/player.rs new file mode 100644 index 0000000..2fce89f --- /dev/null +++ b/server/src/player.rs @@ -0,0 +1,7 @@ +use hyper::{Body, Request}; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Player { + pub id: String +} diff --git a/server/src/session.rs b/server/src/session.rs new file mode 100644 index 0000000..f283e4d --- /dev/null +++ b/server/src/session.rs @@ -0,0 +1,279 @@ +use std::collections::HashMap; +use std::rc::Rc; +use std::sync::Arc; +use kafka::producer::Producer; +use crate::hash::Hasher; +use crate::kafka::{KafkaDefaultEventWriterImpl, KafkaEventReaderImpl, KafkaSessionEventReaderImpl, KafkaSessionEventWriterImpl, KafkaTopicManager}; +use serde::{Serialize, Deserialize}; +use serde::de::DeserializeOwned; +use serde_json::json; +use tokio::sync::Mutex; +use pong::event::event::{Event, EventReader, EventWriter}; +use crate::player::Player; + +pub struct SessionManager { + kafka_host: String, + sessions: Vec, + session_producer: EventWriter, + topic_manager: KafkaTopicManager +} + +// TODO: On startup read the session events from kafka to restore the session id <-> hash mappings. +impl SessionManager { + pub fn new(kafka_host: &str) -> SessionManager { + SessionManager { + kafka_host: kafka_host.to_owned(), + sessions: vec![], + topic_manager: KafkaTopicManager::from("localhost:7243"), + session_producer: EventWriter::new(Box::new(KafkaDefaultEventWriterImpl::new(kafka_host))) + } + } + + pub fn get_session(&self, session_id: &str) -> Option { + self.sessions.iter().find(|s| s.hash == session_id).map_or_else(|| None, |s| Some(s.clone())) + } + + pub async fn create_session(&mut self, player: Player) -> Result { + let add_partition_res = self.topic_manager.add_partition().await; + if let Err(e) = add_partition_res { + println!("Failed to create partition: {}", e); + return Err(e); + } + let session_id = add_partition_res.unwrap(); + let session_hash = Hasher::hash(session_id); + let session = Session::new (session_id, session_hash, player.clone()); + println!("Successfully created session: {:?}", session); + self.write_to_producer(session_created(session.clone(), player.clone())); + self.sessions.push(session.clone()); + Ok(session) + } + + pub async fn join_session(&mut self, session_id: String, player: Player) -> Result { + let updated_session = { + let session = self.sessions.iter_mut().find(|s| s.hash == session_id); + if let None = session { + let error = format!("Can't join session that does not exist: {}", session_id); + return Err(error); + } + let mut session = session.unwrap(); + if session.state != SessionState::PENDING { + let error = format!("Can't join session that is not PENDING: {}", session_id); + return Err(error); + } + if session.players.len() > 1 { + let error = format!("Can't join session with more than 1 player: {}", session_id); + return Err(error); + } + session.players.push(player.clone()); + session.state = SessionState::RUNNING; + session.clone() + }; + { + self.write_to_producer(session_joined(updated_session.clone(), player.clone())); + }; + println!("sessions = {:?}", self.sessions); + Ok(updated_session.clone()) + } + + fn write_to_producer(&mut self, session_event: T) -> Result<(), String> where T : Serialize { + let json_event = serde_json::to_string(&session_event).unwrap(); + let session_event_write = self.session_producer.write(Event {topic: "session".to_owned(), key: None, msg: json_event}); + if let Err(e) = session_event_write { + let message = format!("Failed to write session create event to kafka: {:?}", e); + println!("{}", e); + return Err(message.to_owned()); + } + println!("Successfully produced session event."); + return Ok(()); + } + + pub fn get_session_reader(&self, session_id: &str) -> Result { + let session = self.find_session(&session_id); + if let None = session { + return Err(format!("Unable to find session with hash {}", session_id)) + } + let session = session.unwrap(); + let event_reader = EventReader::new(Box::new(KafkaSessionEventReaderImpl::new(&self.kafka_host, &session, &["move", "status", "input"]))); + Ok(SessionReader {reader: event_reader, session}) + } + + pub fn get_session_writer(&self, session_id: &str) -> Result { + let session = self.find_session(&session_id); + if let None = session { + return Err(format!("Unable to find session with hash {}", session_id)) + } + let session = session.unwrap(); + let event_writer = EventWriter::new(Box::new(KafkaSessionEventWriterImpl::new(&self.kafka_host))); + Ok(SessionWriter {writer: event_writer, session}) + } + + fn find_session(&self, session_id: &str) -> Option { + self.sessions.iter().find(|s| session_id == s.hash).map(|s| s.clone()) + } +} + +pub struct CachingSessionManager { + inner: SessionManager, + reader_cache: HashMap>>, + writer_cache: HashMap>>, +} + +impl CachingSessionManager { + pub fn new(kafka_host: &str) -> CachingSessionManager { + CachingSessionManager { + inner: SessionManager::new(kafka_host), + reader_cache: HashMap::new(), + writer_cache: HashMap::new(), + } + } + + pub fn get_session(&self, session_id: &str) -> Option { + self.inner.get_session(session_id) + } + + pub async fn create_session(&mut self, player: Player) -> Result { + self.inner.create_session(player).await + } + + pub async fn join_session(&mut self, session_id: String, player: Player) -> Result { + self.inner.join_session(session_id, player).await + } + + pub fn get_session_reader(&mut self, session_id: &str) -> Result>, String> { + let cached = self.reader_cache.get(session_id); + if let Some(reader) = cached { + println!("Reusing existing reader for session: {:?}", session_id); + return Ok(Arc::clone(reader)); + } + let reader = self.inner.get_session_reader(session_id); + if let Err(e) = reader { + return Err(e); + } + let reader = Arc::new(Mutex::new(reader.unwrap())); + self.reader_cache.insert(session_id.to_string(), Arc::clone(&reader)); + return Ok(Arc::clone(&reader)); + } + + pub fn get_session_writer(&mut self, session_id: &str) -> Result>, String> { + let cached = self.writer_cache.get(session_id); + if let Some(writer) = cached { + println!("Reusing existing writer for session: {:?}", session_id); + return Ok(Arc::clone(writer)); + } + let writer = self.inner.get_session_writer(session_id); + if let Err(e) = writer { + return Err(e); + } + let writer = Arc::new(Mutex::new(writer.unwrap())); + self.writer_cache.insert(session_id.to_string(), Arc::clone(&writer)); + return Ok(Arc::clone(&writer)); + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Session { + pub id: u16, + pub hash: String, + pub state: SessionState, + players: Vec +} + +impl Session { + pub fn new(id: u16, hash: String, player: Player) -> Session { + Session { + players: vec!(player), + id, + hash, + state: SessionState::PENDING + } + } + + pub fn can_be_joined(&self) -> bool { + self.players.len() == 1 + } + + pub fn join(&mut self, player: Player) -> bool { + if !self.can_be_joined() { + return false + } + self.players.push(player); + return true; + } +} + +#[derive(Clone, Debug, Serialize, Deserialize, Eq, PartialEq)] +pub enum SessionState { + PENDING, // 1 player is missing + RUNNING, // game is playing + CLOSED // game is over +} + +pub struct SessionWriter { + session: Session, + writer: EventWriter +} + +impl SessionWriter { + pub fn write_to_session(&mut self, topic: &str, msg: &str) -> Result<(), String> { + let event = Event {msg: msg.to_owned(), key: Some(self.session.id.to_string()), topic: topic.to_owned()}; + self.writer.write(event) + } +} + +pub struct SessionReader { + session: Session, + reader: EventReader +} + +impl SessionReader { + pub fn read_from_session(&mut self) -> Result, String> { + self.reader.read() + } +} + +#[derive(Deserialize, Serialize)] +struct SessionCreatedEvent { + event_type: SessionEventType, + session: Session, + player: Player +} + +impl SessionCreatedEvent { + pub fn new(session: Session, player: Player) -> SessionCreatedEvent { + SessionCreatedEvent { + event_type: SessionEventType::CREATED, + session, + player + } + } +} + +#[derive(Deserialize, Serialize)] +struct SessionJoinedEvent { + event_type: SessionEventType, + session: Session, + player: Player +} + +impl SessionJoinedEvent { + pub fn new(session: Session, player: Player) -> SessionJoinedEvent { + SessionJoinedEvent { + event_type: SessionEventType::JOINED, + session, + player + } + } +} + +fn session_created(session: Session, player: Player) -> SessionCreatedEvent { + SessionCreatedEvent::new(session, player) +} + +fn session_joined(session: Session, player: Player) -> SessionJoinedEvent { + SessionJoinedEvent::new(session, player) +} + +#[derive(Deserialize, Serialize)] +enum SessionEventType { + CREATED, JOINED +} diff --git a/server/src/utils.rs b/server/src/utils.rs new file mode 100644 index 0000000..4bfdb7a --- /dev/null +++ b/server/src/utils.rs @@ -0,0 +1,58 @@ +pub mod http_utils { + use std::borrow::BorrowMut; + use std::collections::HashMap; + use std::io::Read; + use hyper::{Body, body, Request}; + use hyper::body::Buf; + use serde::de::DeserializeOwned; + use serde::Deserialize; + + pub fn get_query_params(req: &Request) -> HashMap<&str, &str> { + let uri = req.uri(); + let query = uri.query(); + println!("uri={:?}, query={:?}", uri, query); + match query { + None => HashMap::new(), + Some(query) => { + query.split("&").map(|s| s.split_at(s.find("=").unwrap())).map(|(key, value)| (key, &value[1..])).collect() + } + } + } + + pub async fn read_json_body(req: &mut Request) -> T where T : DeserializeOwned { + let mut body = req.body_mut(); + let bytes = body::to_bytes(body).await.unwrap(); + let body_str = std::str::from_utf8(&*bytes).unwrap(); + serde_json::from_str::(body_str).unwrap() + } +} + +#[cfg(test)] +pub mod http_utils_tests { + use rstest::rstest; + use std::collections::HashMap; + use hyper::{Body, Request, Uri}; + use hyper::http::uri::{Builder, Parts}; + use crate::utils::http_utils::get_query_params; + use super::*; + + #[rstest] + #[case( + "?test=abc", + HashMap::from([("test", "abc")]) + )] + #[case( + "?test=abc&help=123", + HashMap::from([("test", "abc"), ("help", "123")]) + )] + #[case( + "show?topic=status&key=abc", + HashMap::from([("topic", "status"), ("key", "abc")]) + )] + fn get_query_params_tests(#[case] query_str: &str, #[case] expected: HashMap<&str, &str>) { + let uri = Builder::new().scheme("https").authority("behnke.rs").path_and_query(query_str).build().unwrap(); + let req = Request::get(uri).body(Body::empty()).unwrap(); + let res = get_query_params(&req); + assert_eq!(res, expected) + } +} diff --git a/server/tests/utils_tests.rs b/server/tests/utils_tests.rs new file mode 100644 index 0000000..7eef968 --- /dev/null +++ b/server/tests/utils_tests.rs @@ -0,0 +1,14 @@ +pub mod tests { + use std::collections::HashMap; + use rstest::rstest; + + #[rstest] + #[case( + "test=abc", + HashMap::from([("test", "abc")]) + )] + fn get_query_params_tests(#[case] query_str: &str, #[case] expected: HashMap<&str, &str>) { + let res = get_query_params(query_str); + assert_eq!(res, expected) + } +} diff --git a/src/lib.rs b/src/lib.rs index 8a87d42..429b559 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -3,7 +3,7 @@ mod utils; use std::cell::RefCell; use pong::collision::collision::{Collision, CollisionDetector}; use pong::game_field::{Field, Input, InputType}; -use pong::game_object::game_object::GameObject; +use pong::game_object::game_object::{DefaultGameObject, GameObject}; use pong::geom::geom::Vector; use pong::geom::shape::ShapeType; use pong::utils::utils::{DefaultLoggerFactory, Logger}; @@ -12,6 +12,7 @@ use serde_json::json; use std::cmp::{max, min}; use std::rc::Rc; use wasm_bindgen::prelude::*; +use pong::pong::pong_events::DefaultPongEventWriter; extern crate serde_json; extern crate web_sys; @@ -111,7 +112,7 @@ pub struct FieldWrapper { #[wasm_bindgen] impl FieldWrapper { pub fn new() -> FieldWrapper { - let field = Field::new(DefaultLoggerFactory::new(Box::new(WasmLogger::root()))); + let field = Field::new(DefaultLoggerFactory::new(Box::new(WasmLogger::root())), DefaultPongEventWriter::new()); FieldWrapper { field } }