client/server
This commit is contained in:
@@ -0,0 +1,12 @@
|
||||
[package]
|
||||
name = "audiopoker_server"
|
||||
version = "0.1.0"
|
||||
edition = "2021"
|
||||
|
||||
[dependencies]
|
||||
audiopoker_core = { path = "../core" }
|
||||
serde = { version = "1.0", features = ["derive"] }
|
||||
serde_json = "1.0"
|
||||
tokio = { version = "1", features = ["rt-multi-thread", "net", "macros", "sync", "io-util"] }
|
||||
tokio-tungstenite = "0.24"
|
||||
futures-util = "0.3"
|
||||
@@ -0,0 +1,36 @@
|
||||
//! Verwaltet alle aktiven Tische. Einfache Umsetzung der "Tischsuche" aus
|
||||
//! PLAN.md: Der Tischname selbst ist der Suchbegriff - unbekannte Namen
|
||||
//! legen automatisch einen neuen Tisch (mit eigenem Tokio-Task) an.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Mutex;
|
||||
|
||||
use tokio::sync::mpsc;
|
||||
|
||||
use crate::table::{TableActor, TableCommand};
|
||||
|
||||
#[derive(Default)]
|
||||
pub struct Lobby {
|
||||
tables: Mutex<HashMap<String, mpsc::UnboundedSender<TableCommand>>>,
|
||||
}
|
||||
|
||||
impl Lobby {
|
||||
pub fn new() -> Self {
|
||||
Self::default()
|
||||
}
|
||||
|
||||
/// Gibt den Befehls-Sender für den genannten Tisch zurück; legt bei
|
||||
/// Bedarf einen neuen Tisch samt eigenem Hintergrund-Task an.
|
||||
pub fn table_sender(&self, table_name: &str) -> mpsc::UnboundedSender<TableCommand> {
|
||||
let mut tables = self.tables.lock().expect("Lobby-Mutex vergiftet");
|
||||
if let Some(tx) = tables.get(table_name) {
|
||||
return tx.clone();
|
||||
}
|
||||
|
||||
let (tx, rx) = mpsc::unbounded_channel();
|
||||
let actor = TableActor::new(table_name.to_string());
|
||||
tokio::spawn(actor.run(rx));
|
||||
tables.insert(table_name.to_string(), tx.clone());
|
||||
tx
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,155 @@
|
||||
//! Nativer WebSocket-Server für Audiopoker (siehe PLAN.md, Phase 2).
|
||||
//! Führt die Spiellogik autoritativ aus (`audiopoker_core::game`) und
|
||||
//! synchronisiert alle verbundenen Clients über `audiopoker_core::network`.
|
||||
|
||||
mod lobby;
|
||||
mod table;
|
||||
|
||||
use std::net::SocketAddr;
|
||||
use std::sync::Arc;
|
||||
|
||||
use futures_util::stream::{SplitSink, SplitStream};
|
||||
use futures_util::{SinkExt, StreamExt};
|
||||
use tokio::net::{TcpListener, TcpStream};
|
||||
use tokio::sync::{mpsc, oneshot};
|
||||
use tokio_tungstenite::tungstenite::Message;
|
||||
use tokio_tungstenite::WebSocketStream;
|
||||
|
||||
use audiopoker_core::network::{ClientMessage, ServerMessage};
|
||||
use lobby::Lobby;
|
||||
use table::TableCommand;
|
||||
|
||||
type WsSink = SplitSink<WebSocketStream<TcpStream>, Message>;
|
||||
type WsStream = SplitStream<WebSocketStream<TcpStream>>;
|
||||
|
||||
const DEFAULT_ADDR: &str = "0.0.0.0:9001";
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() {
|
||||
let addr = std::env::var("AUDIOPOKER_ADDR").unwrap_or_else(|_| DEFAULT_ADDR.to_string());
|
||||
let listener = TcpListener::bind(&addr)
|
||||
.await
|
||||
.unwrap_or_else(|err| panic!("Konnte nicht auf {addr} binden: {err}"));
|
||||
println!("Audiopoker-Server lauscht auf ws://{addr}");
|
||||
|
||||
let lobby = Arc::new(Lobby::new());
|
||||
|
||||
loop {
|
||||
let (stream, peer) = match listener.accept().await {
|
||||
Ok(conn) => conn,
|
||||
Err(err) => {
|
||||
eprintln!("Fehler beim Annehmen einer Verbindung: {err}");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
let lobby = Arc::clone(&lobby);
|
||||
tokio::spawn(async move {
|
||||
if let Err(err) = handle_connection(stream, peer, lobby).await {
|
||||
eprintln!("Verbindung zu {peer} beendet: {err}");
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_connection(
|
||||
stream: TcpStream,
|
||||
peer: SocketAddr,
|
||||
lobby: Arc<Lobby>,
|
||||
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
let ws_stream = tokio_tungstenite::accept_async(stream).await?;
|
||||
println!("Neue Verbindung von {peer}");
|
||||
let (mut ws_sink, mut ws_stream) = ws_stream.split();
|
||||
|
||||
let Some((table_name, player_name)) = await_join(&mut ws_sink, &mut ws_stream).await? else {
|
||||
return Ok(()); // Verbindung wurde beendet, bevor ein Tisch gewählt wurde.
|
||||
};
|
||||
|
||||
let table_tx = lobby.table_sender(&table_name);
|
||||
let (reply_tx, mut reply_rx) = mpsc::unbounded_channel::<ServerMessage>();
|
||||
let (assigned_id_tx, assigned_id_rx) = oneshot::channel();
|
||||
|
||||
table_tx.send(TableCommand::Join {
|
||||
player_name,
|
||||
reply_tx,
|
||||
assigned_id_tx,
|
||||
})?;
|
||||
|
||||
let player_id = match assigned_id_rx.await? {
|
||||
Ok(id) => id,
|
||||
Err(reason) => {
|
||||
send_error(&mut ws_sink, &reason).await;
|
||||
return Ok(());
|
||||
}
|
||||
};
|
||||
|
||||
loop {
|
||||
tokio::select! {
|
||||
outgoing = reply_rx.recv() => {
|
||||
match outgoing {
|
||||
Some(message) => send_message(&mut ws_sink, &message).await?,
|
||||
None => break,
|
||||
}
|
||||
}
|
||||
incoming = ws_stream.next() => {
|
||||
match incoming {
|
||||
Some(Ok(Message::Text(text))) => {
|
||||
match serde_json::from_str::<ClientMessage>(&text) {
|
||||
Ok(client_msg) => {
|
||||
table_tx.send(TableCommand::Action { player_id, message: client_msg })?;
|
||||
}
|
||||
Err(err) => send_error(&mut ws_sink, &format!("Ungültige Nachricht: {err}")).await,
|
||||
}
|
||||
}
|
||||
Some(Ok(Message::Close(_))) | None => break,
|
||||
Some(Ok(_)) => {} // Ping/Pong/Binary werden ignoriert.
|
||||
Some(Err(err)) => {
|
||||
eprintln!("WebSocket-Fehler von {peer}: {err}");
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let _ = table_tx.send(TableCommand::Disconnected { player_id });
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Wartet auf die erste Client-Nachricht, die ein `JoinTable` sein muss.
|
||||
/// Alles andere wird mit einer Fehlermeldung beantwortet, ohne die
|
||||
/// Verbindung zu schließen (der Client darf es erneut versuchen).
|
||||
async fn await_join(
|
||||
ws_sink: &mut WsSink,
|
||||
ws_stream: &mut WsStream,
|
||||
) -> Result<Option<(String, String)>, Box<dyn std::error::Error + Send + Sync>> {
|
||||
loop {
|
||||
let Some(msg) = ws_stream.next().await else {
|
||||
return Ok(None);
|
||||
};
|
||||
match msg? {
|
||||
Message::Text(text) => match serde_json::from_str::<ClientMessage>(&text) {
|
||||
Ok(ClientMessage::JoinTable { table_name, player_name }) => {
|
||||
return Ok(Some((table_name, player_name)));
|
||||
}
|
||||
Ok(_) => send_error(ws_sink, "Bitte zuerst mit JoinTable einem Tisch beitreten.").await,
|
||||
Err(err) => send_error(ws_sink, &format!("Ungültige Nachricht: {err}")).await,
|
||||
},
|
||||
Message::Close(_) => return Ok(None),
|
||||
_ => {} // Ping/Pong/Binary werden ignoriert.
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn send_message(
|
||||
ws_sink: &mut WsSink,
|
||||
message: &ServerMessage,
|
||||
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
let text = serde_json::to_string(message)?;
|
||||
ws_sink.send(Message::Text(text)).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn send_error(ws_sink: &mut WsSink, reason: &str) {
|
||||
let _ = send_message(ws_sink, &ServerMessage::Error(reason.to_string())).await;
|
||||
}
|
||||
@@ -0,0 +1,192 @@
|
||||
//! Ein `TableActor` läuft als eigener Tokio-Task und besitzt exklusiv den
|
||||
//! `GameState` für genau einen Tisch (Actor-Modell: keine geteilten Mutexe
|
||||
//! nötig, alle Zugriffe laufen sequenziell über einen mpsc-Channel).
|
||||
|
||||
use std::collections::HashMap;
|
||||
|
||||
use tokio::sync::{mpsc, oneshot};
|
||||
|
||||
use audiopoker_core::game::{BettingAction, GameManager, GameState, Round, RAISE_INCREMENT};
|
||||
use audiopoker_core::logic::Deck;
|
||||
use audiopoker_core::network::{ClientMessage, ServerMessage};
|
||||
|
||||
/// Maximale Spielerzahl pro Tisch. Sobald diese Zahl an `JoinTable`-Anfragen
|
||||
/// eingetroffen ist, wird automatisch eine neue Hand gestartet.
|
||||
pub const MAX_PLAYERS: usize = 2;
|
||||
|
||||
/// Befehle, die eine WebSocket-Verbindung an den Tisch-Aktor schickt.
|
||||
pub enum TableCommand {
|
||||
Join {
|
||||
player_name: String,
|
||||
reply_tx: mpsc::UnboundedSender<ServerMessage>,
|
||||
assigned_id_tx: oneshot::Sender<Result<u32, String>>,
|
||||
},
|
||||
Action {
|
||||
player_id: u32,
|
||||
message: ClientMessage,
|
||||
},
|
||||
/// Verbindung wurde geschlossen (aktuell nur zur Kenntnisnahme geloggt;
|
||||
/// ein vorzeitig aussteigender Spieler mitten in einer Hand ist noch
|
||||
/// nicht behandelt - siehe PLAN.md, offene Punkte).
|
||||
Disconnected {
|
||||
player_id: u32,
|
||||
},
|
||||
}
|
||||
|
||||
pub struct TableActor {
|
||||
name: String,
|
||||
state: Option<GameState>,
|
||||
pending_names: Vec<String>,
|
||||
senders: HashMap<u32, mpsc::UnboundedSender<ServerMessage>>,
|
||||
}
|
||||
|
||||
impl TableActor {
|
||||
pub fn new(name: String) -> Self {
|
||||
Self {
|
||||
name,
|
||||
state: None,
|
||||
pending_names: Vec::new(),
|
||||
senders: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn run(mut self, mut rx: mpsc::UnboundedReceiver<TableCommand>) {
|
||||
while let Some(cmd) = rx.recv().await {
|
||||
match cmd {
|
||||
TableCommand::Join { player_name, reply_tx, assigned_id_tx } => {
|
||||
self.handle_join(player_name, reply_tx, assigned_id_tx);
|
||||
}
|
||||
TableCommand::Action { player_id, message } => {
|
||||
self.handle_action(player_id, message);
|
||||
}
|
||||
TableCommand::Disconnected { player_id } => {
|
||||
println!("[Tisch {}] Spieler {player_id} getrennt.", self.name);
|
||||
self.senders.remove(&player_id);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn broadcast(&self, message: ServerMessage) {
|
||||
for tx in self.senders.values() {
|
||||
let _ = tx.send(message.clone());
|
||||
}
|
||||
}
|
||||
|
||||
fn send_to(&self, player_id: u32, message: ServerMessage) {
|
||||
if let Some(tx) = self.senders.get(&player_id) {
|
||||
let _ = tx.send(message);
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_join(
|
||||
&mut self,
|
||||
player_name: String,
|
||||
reply_tx: mpsc::UnboundedSender<ServerMessage>,
|
||||
assigned_id_tx: oneshot::Sender<Result<u32, String>>,
|
||||
) {
|
||||
if self.state.is_some() {
|
||||
let _ = assigned_id_tx.send(Err("Tisch ist bereits voll, das Spiel läuft schon.".into()));
|
||||
return;
|
||||
}
|
||||
|
||||
let player_id = self.pending_names.len() as u32;
|
||||
self.pending_names.push(player_name.clone());
|
||||
self.senders.insert(player_id, reply_tx);
|
||||
let _ = assigned_id_tx.send(Ok(player_id));
|
||||
self.send_to(player_id, ServerMessage::Welcome { player_id });
|
||||
|
||||
if self.pending_names.len() < MAX_PLAYERS {
|
||||
self.broadcast(ServerMessage::Announce(format!(
|
||||
"{player_name} ist dem Tisch beigetreten. Warte auf weitere Spieler ({}/{}).",
|
||||
self.pending_names.len(),
|
||||
MAX_PLAYERS
|
||||
)));
|
||||
return;
|
||||
}
|
||||
|
||||
self.start_hand();
|
||||
}
|
||||
|
||||
/// Teilt Karten aus und startet eine neue Hand, sobald genug Spieler da sind.
|
||||
fn start_hand(&mut self) {
|
||||
let mut game = GameManager::start_new_game(std::mem::take(&mut self.pending_names));
|
||||
|
||||
let mut deck = Deck::new();
|
||||
deck.shuffle();
|
||||
for player in game.players.iter_mut() {
|
||||
player.hole_cards = deck.deal_hand(2);
|
||||
}
|
||||
game.community_cards = deck.deal_hand(5);
|
||||
|
||||
self.broadcast(ServerMessage::Announce(format!(
|
||||
"Tisch ist voll, Spiel startet mit {} Spielern.",
|
||||
game.players.len()
|
||||
)));
|
||||
self.broadcast(ServerMessage::PlaySfx("shuffle".into()));
|
||||
|
||||
for player in &game.players {
|
||||
self.send_to(player.id, ServerMessage::HoleCards(player.hole_cards.clone()));
|
||||
}
|
||||
|
||||
self.broadcast(ServerMessage::TurnToAct {
|
||||
player_id: game.players[game.to_act].id,
|
||||
});
|
||||
|
||||
self.state = Some(game);
|
||||
}
|
||||
|
||||
fn handle_action(&mut self, player_id: u32, message: ClientMessage) {
|
||||
let Some(state) = self.state.as_mut() else {
|
||||
self.send_to(player_id, ServerMessage::Error("Das Spiel hat noch nicht begonnen.".into()));
|
||||
return;
|
||||
};
|
||||
|
||||
if state.to_act as u32 != player_id {
|
||||
self.send_to(player_id, ServerMessage::Error("Du bist gerade nicht am Zug.".into()));
|
||||
return;
|
||||
}
|
||||
|
||||
let to_call = GameManager::amount_to_call(state, player_id as usize);
|
||||
let action = match message {
|
||||
ClientMessage::Fold => BettingAction::Fold,
|
||||
ClientMessage::Check if to_call == 0 => BettingAction::Check,
|
||||
// Client wollte checken, es ist aber ein Call fällig - wir werten das
|
||||
// grosszügig als Call, statt die Aktion abzulehnen.
|
||||
ClientMessage::Check => BettingAction::Call(to_call),
|
||||
ClientMessage::Call => BettingAction::Call(to_call),
|
||||
// Der Erhöhungsbetrag wird bewusst serverseitig festgelegt (siehe
|
||||
// Kommentar in `core::network`), damit Clients ihn nicht manipulieren können.
|
||||
ClientMessage::Raise => BettingAction::Raise(to_call + RAISE_INCREMENT),
|
||||
ClientMessage::JoinTable { .. } => {
|
||||
self.send_to(player_id, ServerMessage::Error("Du bist diesem Tisch bereits beigetreten.".into()));
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let round_before = state.current_round;
|
||||
let outcome = GameManager::take_action(state, player_id as usize, action);
|
||||
|
||||
for line in outcome.announcements {
|
||||
self.broadcast(ServerMessage::Announce(line));
|
||||
}
|
||||
for cue in outcome.sfx_cues {
|
||||
self.broadcast(ServerMessage::PlaySfx(cue.to_string()));
|
||||
}
|
||||
|
||||
let state = self.state.as_ref().expect("state wurde oben bereits entpackt");
|
||||
|
||||
if state.current_round != round_before && state.current_round != Round::Showdown {
|
||||
let revealed_count = state.current_round.revealed_community_count();
|
||||
self.broadcast(ServerMessage::CommunityCards(
|
||||
state.community_cards[..revealed_count].to_vec(),
|
||||
));
|
||||
}
|
||||
|
||||
if state.current_round != Round::Showdown {
|
||||
self.broadcast(ServerMessage::TurnToAct {
|
||||
player_id: state.players[state.to_act].id,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user