diff --git a/init-kafka.sh b/init-kafka.sh index b6b69f2..3350566 100755 --- a/init-kafka.sh +++ b/init-kafka.sh @@ -5,7 +5,6 @@ set -e source .env 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 host_tick --bootstrap-server "$KAFKA_HOST:$KAFKA_PORT" +docker exec pong_server_kafka /opt/bitnami/kafka/bin/kafka-topics.sh --create --topic peer_tick --bootstrap-server "$KAFKA_HOST:$KAFKA_PORT" docker exec pong_server_kafka /opt/bitnami/kafka/bin/kafka-topics.sh --create --topic heart_beat --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/kafka/kafka-script-proxy/src/main.rs b/kafka/kafka-script-proxy/src/main.rs index b6ee521..18413ca 100644 --- a/kafka/kafka-script-proxy/src/main.rs +++ b/kafka/kafka-script-proxy/src/main.rs @@ -8,7 +8,7 @@ use tokio::fs::OpenOptions; use tokio::io::AsyncWriteExt; use tokio::process::Command; -const TOPICS: [&str; 5] = ["move", "status", "input", "heart_beat", "session"]; +const TOPICS: [&str; 5] = ["host_tick", "peer_tick", "heart_beat", "session"]; #[tokio::main] pub async fn main() { diff --git a/server/src/event.rs b/server/src/event.rs index f959ab8..5ce14e0 100644 --- a/server/src/event.rs +++ b/server/src/event.rs @@ -126,20 +126,6 @@ impl FromStr for SessionEventType { } } -pub fn deserialize(event: &str) -> Option { - let wrapper = serde_json::from_str::(event); - wrapper.ok().and_then(|w| { - match w.topic.as_str() { - "move" => serde_json::from_str::(&w.event).ok().map(|e| PongEvent::Move(w.session_id, e)), - "input" => serde_json::from_str::(&w.event).ok().map(|e| PongEvent::Input(w.session_id, e)), - "status" => serde_json::from_str::(&w.event).ok().map(|e| PongEvent::Status(w.session_id, e)), - "heart_beat" => serde_json::from_str::(&w.event).ok().map(|e| PongEvent::HeartBeat(w.session_id, e)), - "session" => serde_json::from_str::(&w.event).ok().map(|e| PongEvent::Session(w.session_id, e)), - _ => None - } - }) -} - #[cfg(test)] mod tests { use crate::event::{SessionEvent, SessionEventPayload}; diff --git a/server/src/websocket_handler.rs b/server/src/websocket_handler.rs index 0a2e599..5f78e72 100644 --- a/server/src/websocket_handler.rs +++ b/server/src/websocket_handler.rs @@ -341,9 +341,9 @@ impl FromStr for WebSocketConnectionType { impl WebSocketConnectionType { pub fn get_topics(&self) -> &[&str] { match self { - WebSocketConnectionType::HOST => &["input", "session"], + WebSocketConnectionType::HOST => &["peer_tick", "session"], WebSocketConnectionType::PEER | WebSocketConnectionType::OBSERVER => { - &["move", "input", "status", "session"] + &["host_tick", "session"] } } }