diff --git a/pong/src/game_field.rs b/pong/src/game_field.rs index 6aa58ff..c6e99b4 100644 --- a/pong/src/game_field.rs +++ b/pong/src/game_field.rs @@ -19,10 +19,11 @@ pub enum InputType { DOWN, } -#[derive(Clone, Copy, Debug, PartialEq, Eq)] +#[derive(Clone, Debug, PartialEq, Eq)] pub struct Input { pub input: InputType, pub obj_id: u16, + pub player: u16 } pub struct Field { diff --git a/server/src/http.rs b/server/src/http.rs index 02a7859..e58df8f 100644 --- a/server/src/http.rs +++ b/server/src/http.rs @@ -4,6 +4,7 @@ use std::io::ErrorKind::NotFound; use std::net::SocketAddr; use std::str::FromStr; use std::sync::Arc; +use std::time::Duration; use hyper::{Body, body, Method, Request, Response, Server, StatusCode}; use hyper::server::conn::AddrStream; use hyper::service::{make_service_fn, service_fn}; @@ -15,6 +16,7 @@ use serde::{Deserialize, Serialize}; use tokio::sync::Mutex; use pong::event::event::{Event, EventReader, EventWriter}; use futures::{sink::SinkExt, stream::StreamExt}; +use tokio::time::sleep; use crate::kafka::{KafkaEventReaderImpl, KafkaSessionEventWriterImpl}; use crate::player::Player; use crate::session::{Session, SessionManager}; @@ -116,6 +118,7 @@ async fn serve_websocket(websocket_session: WebSocketSession, websocket: HyperWe match message.unwrap() { Message::Text(msg) => { let events = serde_json::from_str::(&msg); + println!("Received ws message to persist events to kafka"); if let Err(e) = events { eprintln!("Failed to deserialize ws message to event {}: {}", msg, e); continue; @@ -137,7 +140,7 @@ async fn serve_websocket(websocket_session: WebSocketSession, websocket: HyperWe if any_error { eprintln!("Failed to write at least one message for session {}", event_wrapper.session_id); } else { - // println!("Successfully wrote {} messages to kafka for session {:?}", event_count, websocket_session_read_copy) + println!("Successfully wrote {} messages to kafka for session {:?}", event_count, websocket_session_read_copy) } }, Message::Close(msg) => { @@ -151,11 +154,13 @@ async fn serve_websocket(websocket_session: WebSocketSession, websocket: HyperWe _ => {} } } + println!("!!!! Exit websocket receiver !!!!") }); let websocket_session_write_copy = websocket_session.clone(); tokio::spawn(async move { println!("Ready to read messages from kafka: {:?}", websocket_session_write_copy); loop { + println!("Reading messages from kafka."); let messages = event_reader.read_from_session(); if let Err(_) = messages { eprintln!("Failed to read messages from kafka for session: {:?}", websocket_session_write_copy); @@ -164,14 +169,18 @@ async fn serve_websocket(websocket_session: WebSocketSession, websocket: HyperWe // println!("Read messages for websocket_session {:?} from consumer: {:?}", websocket_session_write_copy, messages); let messages = messages.unwrap(); if messages.len() == 0 { + println!("No new messages from kafka."); continue; } + println!("{} new messages from kafka.", messages.len()); let json = serde_json::to_string(&messages).unwrap(); let message = Message::from(json); + println!("Sending kafka messages through websocket."); let send_res = websocket_writer.send(message).await; if let Err(e) = send_res { eprintln!("Failed to send message to websocket for session {:?}: {:?}", websocket_session_write_copy, e) } + sleep(Duration::from_millis(2000)); } }); Ok(()) diff --git a/server/src/kafka.rs b/server/src/kafka.rs index 41b553a..2185be2 100644 --- a/server/src/kafka.rs +++ b/server/src/kafka.rs @@ -157,13 +157,16 @@ impl KafkaEventReaderImpl { let message_sets: Vec> = polled.iter().collect(); let mut events = vec![]; for ms in message_sets { + let mut topic_event_count = 0; 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() }; + topic_event_count += 1; events.push(event); } + println!("returned {:?} events for topic={} partition={}", topic_event_count, topic, partition); self.consumer.consume_messageset(ms).unwrap(); } self.consumer.commit_consumed().unwrap(); diff --git a/src/lib.rs b/src/lib.rs index e520a2e..a4f6aef 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -93,6 +93,7 @@ impl InputTypeDTO { pub struct InputDTO { pub input: InputTypeDTO, pub obj_id: u16, + pub player: u16 } impl InputDTO { @@ -100,6 +101,7 @@ impl InputDTO { return Input { input: self.input.to_input_type(), obj_id: self.obj_id, + player: self.player }; } } diff --git a/www/index.js b/www/index.js index d874c5b..dd58841 100644 --- a/www/index.js +++ b/www/index.js @@ -38,7 +38,15 @@ const renderLoop = () => { reset(); return; } - actions = getInputActions(); + actions = getInputActions().filter(it => { + if (!networkSession) { + return it + } + if (networkSession && isHost) { + return it.player === 1; + } + return it.player === 2; + }); if (paused) { requestAnimationFrame(renderLoop); return; @@ -72,7 +80,14 @@ const tick = () => { field.tick(actions, update); objects = JSON.parse(field.objects()); } else if (isHost) { - field.tick(actions, update); + // Would mean that input events would get lost if latency is higher than 100 ms. + const peerInputEvents = events.filter(e => e.topic === "input").filter(it => it.msg.player !== player.id) + const lastPeerInputEvents = peerInputEvents.length ? [peerInputEvents[peerInputEvents - 1]] : [] + const allActions = [ + ...actions, ...lastPeerInputEvents + ]; + console.warn({allActions}) + field.tick(allActions, update); objects = JSON.parse(field.objects()); sendEvents([...getInputEvents(), ...getMoveEvents(objects)]) } else { @@ -87,22 +102,29 @@ const tick = () => { const latestMoveEvents = Object.entries(moveEventsByObj) .map(([_, moveEvents]) => moveEvents[moveEvents.length - 1]); objects = latestMoveEvents.map(({msg}) => msg); - sendEvents(getInputEvents()) + // sendEvents(getInputEvents()) } render(objects); } const getMoveEvents = objects => { - return objects.map(o => ({session_id: networkSession.hash, topic: 'move', msg: JSON.stringify({...o, session_id: networkSession.hash})})); + return objects.map(o => ({session_id: networkSession.hash, topic: 'move', msg: JSON.stringify({...o, session_id: networkSession.hash, ts: Date.now()})})); } const getInputEvents = () => { - return actions.map(({input}) => ({msg: JSON.stringify({input, player: player.id, session_id: networkSession.hash}), session_id: networkSession.hash, topic: 'input'})); + const inputEvents = actions.map(({input}) => ({msg: JSON.stringify({inputs: [input], player: player.id, session_id: networkSession.hash, ts: Date.now()}), session_id: networkSession.hash, topic: 'input'})); + if (inputEvents.length) { + return inputEvents; + } + const noInputs = {inputs: [], obj_id: isHost ? 0 : 1, player: isHost ? 1 : 2 } + return [{msg: JSON.stringify({input: noInputs, player: player.id, session_id: networkSession.hash, ts: Date.now()}), session_id: networkSession.hash, topic: 'input'}]; } const sendEvents = events => { const eventWrapper = {session_id: networkSession.hash, events}; - websocket.send(JSON.stringify(eventWrapper)); + if (eventWrapper.events.length) { + websocket.send(JSON.stringify(eventWrapper)); + } } const render = objects => { @@ -316,13 +338,13 @@ const getInputActions = () => { return [...keysDown].map(key => { switch(key) { case 'KeyW': - return {input: 'UP', obj_id: 0} + return {input: 'UP', obj_id: 0, player: 1} case 'KeyS': - return {input: 'DOWN', obj_id: 0} + return {input: 'DOWN', obj_id: 0, player: 1} case 'ArrowUp': - return {input: 'UP', obj_id: 1} + return {input: 'UP', obj_id: 1, player: 2} case 'ArrowDown': - return {input: 'DOWN', obj_id: 1} + return {input: 'DOWN', obj_id: 1, player: 2} default: return null }