wip - issues with sending ws messages to kafka

This commit is contained in:
Thilo Behnke
2022-06-06 17:21:50 +02:00
parent a508c051a1
commit 2f57c831a1
5 changed files with 49 additions and 12 deletions
+2 -1
View File
@@ -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 {
+10 -1
View File
@@ -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::<SessionEventListDTO>(&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(())
+3
View File
@@ -157,13 +157,16 @@ impl KafkaEventReaderImpl {
let message_sets: Vec<MessageSet<'_>> = 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();
+2
View File
@@ -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
};
}
}
+32 -10
View File
@@ -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
}