feature/event-writer

This commit is contained in:
Thilo Behnke
2022-06-01 20:44:22 +02:00
parent 97ac6b6f55
commit 9aca924037
27 changed files with 1349 additions and 10 deletions

View File

@@ -5,7 +5,7 @@ authors = ["Thilo Behnke <thilo.behnke@gmx.net>"]
edition = "2018"
[workspace]
members = ["pong"]
members = ["pong", "server"]
[lib]
crate-type = ["cdylib", "rlib"]

View File

@@ -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"

88
pong/src/event.rs Normal file
View File

@@ -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<String>,
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<dyn EventWriterImpl>
}
impl EventWriter {
pub fn new(writer_impl: Box<dyn EventWriterImpl>) -> 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<Vec<Event>, String>;
}
pub struct EventReader {
reader_impl: Box<dyn EventReaderImpl>
}
impl EventReader {
pub fn new(reader_impl: Box<dyn EventReaderImpl>) -> EventReader {
EventReader {
reader_impl
}
}
pub fn read(&mut self) -> Result<Vec<Event>, String> {
self.reader_impl.read()
}
}
}

View File

@@ -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<dyn CollisionRegistry>,
objs: Vec<Rc<RefCell<Box<dyn GameObject>>>>,
event_writer: Box<dyn PongEventWriter>,
collision_detector: CollisionDetector,
collision_handler: CollisionHandler,
}
impl Field {
pub fn new(logger_factory: Box<dyn LoggerFactory>) -> Field {
pub fn new(logger_factory: Box<dyn LoggerFactory>, event_writer: Box<dyn PongEventWriter>) -> 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<Input>) {
pub fn tick(&mut self, inputs: Vec<Input>) {
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<dyn CollisionRegistry> {

View File

@@ -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<dyn GeomComp>,
physics: Box<dyn PhysicsComp>,
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;
}
}
}

View File

@@ -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,

View File

@@ -4,3 +4,4 @@ pub mod game_object;
pub mod geom;
pub mod pong;
pub mod utils;
pub mod event;

View File

@@ -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<dyn PongEventWriter> {
Box::new(DefaultPongEventWriter {
writer: EventWriter::noop()
})
}
}
impl DefaultPongEventWriter {
pub fn new() -> Box<dyn PongEventWriter> {
Box::new(DefaultPongEventWriter {
writer: EventWriter::file()
})
}
}
}

2
server/.env Normal file
View File

@@ -0,0 +1,2 @@
KAFKA_HOST=localhost
KAFKA_PORT=9092

21
server/Cargo.toml Normal file
View File

@@ -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"

37
server/docker-compose.yml Normal file
View File

@@ -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

17
server/kafka/Dockerfile Normal file
View File

@@ -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" ]

View File

@@ -0,0 +1,4 @@
#!/usr/bin/env bash
/bin/kafka-script-proxy &
/opt/bitnami/scripts/kafka/entrypoint.sh "/opt/bitnami/scripts/kafka/run.sh"

View File

@@ -0,0 +1 @@
/.idea

View File

@@ -0,0 +1,2 @@
.idea
/target

View File

@@ -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"] }

View File

@@ -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<Body>| 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<Body>) -> Result<Response<Body>, 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<Response<Body>, 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(&current_count.to_string()); // current_count because count - 1 = max partition
}
async fn get_highest_partition_count() -> Result<u32, PartitionCountQueryError> {
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::<u32>();
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<Output> {
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<Response<Body>, Infallible> {
let json = format!("{{\"data\": {}}}", value);
return Ok(Response::new(Body::from(json)));
}
pub fn build_error_res(error: &str, status: StatusCode) -> Result<Response<Body>, 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);
}
}

10
server/run.sh Executable file
View File

@@ -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"

9
server/src/hash.rs Normal file
View File

@@ -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)
}
}

250
server/src/http.rs Normal file
View File

@@ -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<Mutex<CachingSessionManager>>
}
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<dyn std::error::Error + Send + Sync>> {
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<Body>| {
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<Mutex<CachingSessionManager>>) -> 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::<SessionEventWriteDTO>(&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<Mutex<CachingSessionManager>>, req: Request<Body>, addr: SocketAddr) -> Result<Response<Body>, 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<Mutex<CachingSessionManager>>, req: Request<Body>) -> Result<Response<Body>, 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<Mutex<CachingSessionManager>>, req: Request<Body>, addr: SocketAddr) -> Result<Response<Body>, 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<Mutex<CachingSessionManager>>, mut req: Request<Body>, addr: SocketAddr) -> Result<Response<Body>, Infallible> {
let mut locked = session_manager.lock().await;
let body = read_json_body::<SessionJoinDto>(&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<Mutex<CachingSessionManager>>, mut req: Request<Body>) -> Result<Response<Body>, Infallible> {
let mut locked = session_manager.lock().await;
let event = read_json_body::<SessionEventWriteDTO>(&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<Mutex<CachingSessionManager>>, mut req: Request<Body>) -> Result<Response<Body>, Infallible> {
let mut locked = session_manager.lock().await;
let read_payload = read_json_body::<SessionReadDTO>(&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<Response<Body>, Infallible> {
let json = format!("{{\"data\": {}}}", value);
return Ok(Response::new(Body::from(json)));
}
pub fn build_error_res(error: &str, status: StatusCode) -> Result<Response<Body>, 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
}

262
server/src/kafka.rs Normal file
View File

@@ -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<SessionPartitioner>
}
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<Vec<Event>, String> {
self.consume()
}
}
impl KafkaEventReaderImpl {
fn consume(&mut self) -> Result<Vec<Event>, 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<MessageSet<'_>> = 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<Vec<Event>, 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<u16, String> {
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::<PartitionApiDTO>(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::<i32>().unwrap();
println!("Overriding message partition with key: {}", msg.partition);
},
None => panic!("Producing message without key not allowed!")
}
}
}

15
server/src/main.rs Normal file
View File

@@ -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");
}

7
server/src/player.rs Normal file
View File

@@ -0,0 +1,7 @@
use hyper::{Body, Request};
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Player {
pub id: String
}

279
server/src/session.rs Normal file
View File

@@ -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>,
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<Session> {
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<Session, String> {
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<Session, String> {
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<T>(&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<SessionReader, String> {
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<SessionWriter, String> {
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<Session> {
self.sessions.iter().find(|s| session_id == s.hash).map(|s| s.clone())
}
}
pub struct CachingSessionManager {
inner: SessionManager,
reader_cache: HashMap<String, Arc<Mutex<SessionReader>>>,
writer_cache: HashMap<String, Arc<Mutex<SessionWriter>>>,
}
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<Session> {
self.inner.get_session(session_id)
}
pub async fn create_session(&mut self, player: Player) -> Result<Session, String> {
self.inner.create_session(player).await
}
pub async fn join_session(&mut self, session_id: String, player: Player) -> Result<Session, String> {
self.inner.join_session(session_id, player).await
}
pub fn get_session_reader(&mut self, session_id: &str) -> Result<Arc<Mutex<SessionReader>>, 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<Arc<Mutex<SessionWriter>>, 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<Player>
}
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<Vec<Event>, 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
}

58
server/src/utils.rs Normal file
View File

@@ -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<Body>) -> 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<T>(req: &mut Request<Body>) -> 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::<T>(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)
}
}

View File

@@ -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)
}
}

View File

@@ -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 }
}