1#![allow(clippy::too_many_arguments)]
2use std::collections::{HashMap, HashSet};
3use std::sync::atomic::{AtomicBool, Ordering};
4use std::sync::{mpsc as std_mpsc, Arc, Mutex};
5use std::thread;
6use std::time::{Duration, Instant};
7
8use crate::config::{Config, DeckSelection, SelfPlayConfig};
9use crate::engine_backend::{java_backend, rust_backend, EngineBackendKind, HostedGameOver};
10use crate::updater::{run_stale_monitor, StaleConfig};
11use futures_util::stream::{SplitSink, SplitStream};
12use futures_util::{SinkExt, StreamExt};
13use manabot::{run_bot, AgentKind, BotConfig};
14use manabrew_agent_interface::ids_codec::{parse_player_slot, player_slot};
15use manabrew_agent_interface::prompt::{AgentMessage, ClientToServerMessage, PromptOutput};
16use manabrew_agent_interface::protocol::{
17 ClientMessage, ClientPlatform, EngineKind, GameFormat, PlayerDeckInfo, ResumeRoomRequest,
18 RoomInfo, RoomStatus, ServerMessage, StateEnvelope, PROTOCOL_VERSION,
19};
20use manabrew_protocol::deck_dto::Deck;
21use manabrew_protocol::transport::DirectiveInput;
22use serde::Deserialize;
23use serde_json::{json, Value};
24use tokio::net::TcpStream;
25use tokio::sync::mpsc as tokio_mpsc;
26use tokio::task::JoinHandle;
27use tokio::time;
28use tokio_tungstenite::tungstenite::Message;
29use tokio_tungstenite::{connect_async_with_config, MaybeTlsStream, WebSocketStream};
30use tracing::{debug, error, info, warn};
31use uuid::Uuid;
32
33type WsStream = WebSocketStream<MaybeTlsStream<TcpStream>>;
34type WsWrite = SplitSink<WsStream, Message>;
35type WsRead = SplitStream<WsStream>;
36
37const SELF_HOSTED_NODE_PROTOCOL: &str = "self-hosted-node";
38const PRIVATE_MESSAGE_MISSING_PLAYER: &str =
39 "cannot forward private engine message without a mapped player";
40
41struct RelayClient {
42 username: String,
43 outbound: tokio_mpsc::UnboundedSender<Message>,
44 read: WsRead,
45 writer: JoinHandle<()>,
46}
47
48enum EngineSession {
49 Manabrew {
50 game_id: String,
51 remote_response_txs: HashMap<usize, std_mpsc::Sender<ClientToServerMessage>>,
52 },
53 Forge {
54 game_id: String,
55 remote_response_txs: HashMap<usize, std_mpsc::Sender<ClientToServerMessage>>,
56 cancel: Arc<AtomicBool>,
57 },
58}
59
60impl EngineSession {
61 fn game_id(&self) -> &str {
62 match self {
63 EngineSession::Manabrew { game_id, .. } | EngineSession::Forge { game_id, .. } => {
64 game_id
65 }
66 }
67 }
68}
69
70struct BotState {
71 handle: JoinHandle<()>,
72 shutdown: Arc<tokio::sync::Notify>,
73}
74
75#[derive(Debug, Deserialize)]
76#[serde(rename_all = "camelCase")]
77struct SpawnBotPayload {
78 #[serde(default)]
79 deck: Option<SpawnBotDeckPayload>,
80 #[serde(default)]
81 decks: Option<Vec<SpawnBotDeckPayload>>,
82}
83
84#[derive(Debug, Deserialize)]
85#[serde(rename_all = "camelCase")]
86struct SpawnBotDeckPayload {
87 deck_name: String,
88 deck: Deck,
89 #[serde(default)]
90 commander_name: Option<String>,
91}
92
93type SharedEngineSession = Arc<Mutex<Option<EngineSession>>>;
94type SessionRegistry = Arc<Mutex<Vec<SharedEngineSession>>>;
95
96fn registry_idle(registry: &SessionRegistry) -> bool {
97 registry
98 .lock()
99 .map(|sessions| {
100 sessions
101 .iter()
102 .all(|session| session.lock().map(|guard| guard.is_none()).unwrap_or(false))
103 })
104 .unwrap_or(false)
105}
106
107static DRAINING: AtomicBool = AtomicBool::new(false);
108
109fn draining() -> bool {
110 DRAINING.load(Ordering::Relaxed)
111}
112
113type SharedBotState = Arc<Mutex<Vec<BotState>>>;
114
115#[derive(Clone)]
116struct GameStart {
117 game_id: String,
118 player_order: Vec<String>,
119 player_decks: Vec<PlayerDeckInfo>,
120 starting_life: i32,
121}
122
123#[derive(Default)]
124struct HostSnapshot {
125 room_info: Option<RoomInfo>,
126 resume_token: Option<String>,
127 game: Option<GameStart>,
128 last_state: Option<Value>,
129 last_state_by_slot: HashMap<String, Value>,
130 pending_prompts: HashMap<String, Value>,
131 pending_end_game: Option<String>,
132}
133
134type SharedHostSnapshot = Arc<Mutex<HostSnapshot>>;
135
136const RECONNECT_BACKOFF_SECS: [u64; 6] = [1, 2, 4, 8, 15, 30];
137
138pub(crate) const OBSERVER_SEAT: usize = usize::MAX;
140const CLOSE_DRAIN_TIMEOUT: Duration = Duration::from_secs(3);
141const BOT_STOP_TIMEOUT: Duration = Duration::from_secs(5);
142
143enum LoopExit {
144 Cancelled,
145 Disconnected,
146}
147
148pub async fn cli_entry() {
149 tracing_subscriber::fmt()
150 .with_env_filter(
151 tracing_subscriber::EnvFilter::try_from_default_env()
152 .unwrap_or_else(|_| "self_hosted_node=info".into()),
153 )
154 .init();
155 crate::metrics::init_from_env();
156
157 if std::env::var("SELF_HOSTED_NODE_JAVA_SMOKE").is_ok() {
158 let max_prompts = std::env::var("SELF_HOSTED_NODE_JAVA_SMOKE_PROMPTS")
159 .ok()
160 .and_then(|value| value.parse().ok())
161 .unwrap_or(4);
162 if let Err(error) = java_backend::run_smoke_game(max_prompts) {
163 error!(%error, "java-forge smoke failed");
164 std::process::exit(1);
165 }
166 info!(max_prompts, "java-forge smoke completed");
167 return;
168 }
169 if std::env::var("SELF_HOSTED_NODE_GRAAL_SMOKE").is_ok() {
170 if let Err(error) = java_backend::run_graal_smoke() {
171 error!(%error, "graal-forge smoke failed");
172 std::process::exit(1);
173 }
174 info!("graal-forge smoke completed");
175 return;
176 }
177 if std::env::var("SELF_HOSTED_NODE_JAVA_CONCEDE_SMOKE").is_ok() {
178 if let Err(error) = java_backend::run_concede_smoke() {
179 error!(%error, "java-forge concede smoke failed");
180 std::process::exit(1);
181 }
182 info!("java-forge concede smoke completed");
183 return;
184 }
185 if let Ok(scenario_name) = std::env::var("SELF_HOSTED_NODE_JAVA_SCENARIO") {
186 let max_prompts = std::env::var("SELF_HOSTED_NODE_JAVA_SCENARIO_PROMPTS")
187 .ok()
188 .and_then(|value| value.parse().ok())
189 .unwrap_or(20);
190 if let Err(error) = java_backend::run_scenario(&scenario_name, max_prompts) {
191 error!(scenario_name, %error, "java-forge scenario failed");
192 std::process::exit(1);
193 }
194 info!(scenario_name, max_prompts, "java-forge scenario completed");
195 return;
196 }
197 if std::env::var("SELF_HOSTED_NODE_JAVA_SELF_PLAY").is_ok() {
198 let cfg = SelfPlayConfig::from_env();
199 let max_prompts = std::env::var("SELF_HOSTED_NODE_JAVA_SELF_PLAY_PROMPTS")
200 .ok()
201 .and_then(|value| value.parse().ok())
202 .unwrap_or(2_000);
203 let games = std::env::var("SELF_HOSTED_NODE_JAVA_SELF_PLAY_GAMES")
204 .ok()
205 .and_then(|value| value.parse().ok())
206 .unwrap_or(1);
207 if let Err(error) =
208 java_backend::run_self_play(&cfg.seats, cfg.starting_life, cfg.seed, max_prompts, games)
209 {
210 error!(%error, "java-forge self-play failed");
211 std::process::exit(1);
212 }
213 info!(
214 players = cfg.seats.len(),
215 max_prompts, games, "java-forge self-play completed"
216 );
217 return;
218 }
219 if let Ok(value) = std::env::var("SELF_HOSTED_NODE_JAVA_CONCURRENT_GAMES") {
220 let cfg = SelfPlayConfig::from_env();
221 let max_prompts = std::env::var("SELF_HOSTED_NODE_JAVA_SELF_PLAY_PROMPTS")
222 .ok()
223 .and_then(|value| value.parse().ok())
224 .unwrap_or(2_000);
225 let concurrency = value.parse().unwrap_or(2);
226 if let Err(error) = java_backend::run_concurrent_self_play(
227 &cfg.seats,
228 cfg.starting_life,
229 cfg.seed,
230 max_prompts,
231 concurrency,
232 ) {
233 error!(%error, "java-forge concurrent self-play failed");
234 std::process::exit(1);
235 }
236 info!(
237 players = cfg.seats.len(),
238 concurrency, "java-forge concurrent self-play completed"
239 );
240 return;
241 }
242 if std::env::var("SELF_HOSTED_NODE_RUST_SELF_PLAY").is_ok() {
243 let cfg = SelfPlayConfig::from_env();
244 let max_turns = std::env::var("SELF_HOSTED_NODE_RUST_SELF_PLAY_TURNS")
245 .ok()
246 .and_then(|value| value.parse().ok())
247 .unwrap_or(200);
248 if let Err(error) =
249 rust_backend::run_self_play(&cfg.seats, cfg.starting_life, cfg.seed, max_turns)
250 {
251 error!(%error, "rust self-play failed");
252 std::process::exit(1);
253 }
254 info!(
255 players = cfg.seats.len(),
256 max_turns, "rust self-play completed"
257 );
258 return;
259 }
260
261 let config = Config::from_env();
262 if let Err(error) = run(config).await {
263 error!(%error, "room node exited");
264 std::process::exit(1);
265 }
266}
267
268pub type RoomCancel = Arc<tokio::sync::Notify>;
270
271fn ensure_engine_ready(config: &Config) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
272 if config.engine_enabled
273 && config.backend.is_supported()
274 && matches!(config.backend, EngineBackendKind::Forge)
275 {
276 info!("initializing forge engine backend");
277 java_backend::init_engine()?;
278 }
279 Ok(())
280}
281
282pub async fn host_room(
285 config: Config,
286 cancel: RoomCancel,
287 ready: tokio::sync::oneshot::Sender<String>,
288) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
289 ensure_engine_ready(&config)?;
290 host_one_room(config, None, cancel, Some(ready), None).await
291}
292
293async fn run(mut config: Config) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
294 ensure_engine_ready(&config)?;
295
296 let registry: SessionRegistry = Arc::new(Mutex::new(Vec::new()));
297
298 let slots = config.max_games.max(1);
299 let single = config.room_id.is_some() || slots <= 1;
300 let room_count = if single { 1 } else { slots };
301 let cancels: Vec<RoomCancel> = (0..room_count)
302 .map(|_| Arc::new(tokio::sync::Notify::new()))
303 .collect();
304
305 let monitor_registry = registry.clone();
306 let stale_cancels = cancels.clone();
307 tokio::spawn(run_stale_monitor(
308 StaleConfig::from_env_and_args(),
309 move || registry_idle(&monitor_registry),
310 move || notify_all(&stale_cancels),
311 ));
312
313 let signal_cancels = cancels.clone();
314 tokio::spawn(async move {
315 wait_for_shutdown_signal().await;
316 info!("shutdown signal received; closing rooms");
317 notify_all(&signal_cancels);
318 });
319
320 #[cfg(unix)]
321 tokio::spawn(async move {
322 let mut usr1 =
323 match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::user_defined1()) {
324 Ok(usr1) => usr1,
325 Err(error) => {
326 warn!(%error, "failed to install SIGUSR1 handler");
327 return;
328 }
329 };
330 while usr1.recv().await.is_some() {
331 DRAINING.store(true, Ordering::Relaxed);
332 info!("drain signal received; refusing new games and closing each room once idle");
333 }
334 });
335
336 if single {
337 return host_one_room(config, None, cancels[0].clone(), None, Some(registry)).await;
338 }
339
340 config.format = GameFormat::Any;
341 let hosts: Vec<(Config, String)> = (0..slots)
342 .map(|slot| (config.clone(), (slot + 1).to_string()))
343 .collect();
344
345 info!(rooms = hosts.len(), "hosting multiple rooms on one node");
346 let mut handles = Vec::with_capacity(hosts.len());
347 for ((cfg, label), cancel) in hosts.into_iter().zip(cancels) {
348 let registry = registry.clone();
349 handles.push(tokio::spawn(async move {
350 if let Err(error) =
351 host_one_room(cfg, Some(label.clone()), cancel, None, Some(registry)).await
352 {
353 error!(%error, label, "room host exited");
354 }
355 }));
356 }
357 for handle in handles {
358 let _ = handle.await;
359 }
360 Ok(())
361}
362
363fn notify_all(cancels: &[RoomCancel]) {
364 for cancel in cancels {
365 cancel.notify_one();
366 }
367}
368
369async fn wait_for_shutdown_signal() {
370 #[cfg(unix)]
371 {
372 let mut term =
373 match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) {
374 Ok(term) => term,
375 Err(error) => {
376 warn!(%error, "failed to install SIGTERM handler");
377 let _ = tokio::signal::ctrl_c().await;
378 return;
379 }
380 };
381 tokio::select! {
382 _ = tokio::signal::ctrl_c() => {}
383 _ = term.recv() => {}
384 }
385 }
386 #[cfg(not(unix))]
387 {
388 let _ = tokio::signal::ctrl_c().await;
389 }
390}
391
392async fn host_one_room(
393 mut config: Config,
394 label: Option<String>,
395 cancel: RoomCancel,
396 ready: Option<tokio::sync::oneshot::Sender<String>>,
397 sessions: Option<SessionRegistry>,
398) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
399 if let Some(label) = &label {
400 config.username = format!("{}-{label}", config.username);
401 config.bot_username = format!("{}-{label}", config.bot_username);
402 config.room_name = format!("{} ({label})", config.room_name);
403 }
404
405 let _rooms_hosted = crate::metrics::RoomHostedGuard::new(if label.is_some() {
406 crate::metrics::PoolKind::Pod
407 } else {
408 crate::metrics::PoolKind::Solo
409 });
410
411 info!(
412 relay_url = %config.relay_url,
413 username = %config.username,
414 room_name = %config.room_name,
415 auto_start = config.auto_start,
416 engine_enabled = config.engine_enabled,
417 host_plays = config.host_plays,
418 bot_enabled = config.bot_enabled,
419 "starting room host"
420 );
421
422 let snapshot: SharedHostSnapshot = Arc::new(Mutex::new(HostSnapshot::default()));
423 let engine_session: SharedEngineSession = Arc::new(Mutex::new(None));
424 if let Some(sessions) = &sessions {
425 if let Ok(mut registered) = sessions.lock() {
426 registered.push(engine_session.clone());
427 }
428 }
429 let bot_state: SharedBotState = Arc::new(Mutex::new(Vec::new()));
430 let (outbound_tx, mut outbound_rx) = tokio_mpsc::unbounded_channel::<ClientMessage>();
431
432 let mut host =
433 RelayClient::connect(&config.relay_url, &config.username, &config.password).await?;
434 let mut room_id = establish_room(&mut host, &config, &snapshot).await?;
435
436 if let Some(ready) = ready {
437 let _ = ready.send(room_id.clone());
438 }
439
440 if config.bot_enabled {
441 spawn_bots(
442 &config,
443 std::slice::from_ref(&config.bot_deck),
444 &room_id,
445 &bot_state,
446 );
447 }
448
449 loop {
450 let exit = run_client_loop(
451 &mut host,
452 &config,
453 &room_id,
454 &engine_session,
455 &snapshot,
456 &bot_state,
457 &outbound_tx,
458 &mut outbound_rx,
459 &cancel,
460 )
461 .await;
462 if matches!(exit, LoopExit::Cancelled) {
463 stop_bots(&bot_state);
464 host.close().await;
465 return Ok(());
466 }
467
468 crate::metrics::record_relay_reconnect();
469 let mut attempt: usize = 0;
470 host = loop {
471 let delay = RECONNECT_BACKOFF_SECS[attempt.min(RECONNECT_BACKOFF_SECS.len() - 1)];
472 tokio::select! {
473 _ = cancel.notified() => {
474 info!(username = %config.username, "room host cancelled while reconnecting");
475 cancel_engine(&engine_session);
476 stop_bots(&bot_state);
477 return Ok(());
478 }
479 _ = time::sleep(Duration::from_secs(delay)) => {}
480 }
481 attempt += 1;
482 let mut client =
483 match RelayClient::connect(&config.relay_url, &config.username, &config.password)
484 .await
485 {
486 Ok(client) => client,
487 Err(error) => {
488 warn!(%error, attempt, "relay reconnect failed");
489 continue;
490 }
491 };
492 match reestablish_room(&mut client, &config, &engine_session, &snapshot).await {
493 Ok(new_room_id) => {
494 if new_room_id != room_id {
495 stop_bots(&bot_state);
496 if config.bot_enabled {
497 spawn_bots(
498 &config,
499 std::slice::from_ref(&config.bot_deck),
500 &new_room_id,
501 &bot_state,
502 );
503 }
504 room_id = new_room_id;
505 }
506 break client;
507 }
508 Err(error) => {
509 warn!(%error, attempt, "failed to re-establish room; retrying");
510 }
511 }
512 };
513 info!(username = %config.username, room_id, "relay connection re-established");
514 }
515}
516
517async fn establish_room(
518 client: &mut RelayClient,
519 config: &Config,
520 snapshot: &SharedHostSnapshot,
521) -> Result<String, Box<dyn std::error::Error + Send + Sync>> {
522 let room_id = if let Some(room_id) = &config.room_id {
523 client
524 .send(&ClientMessage::JoinRoom {
525 room_id: room_id.clone(),
526 observe: !config.host_plays,
527 as_bot: false,
528 password: config.room_password.clone(),
529 })
530 .await?;
531 info!(room_id, "joining configured room");
532 room_id.clone()
533 } else {
534 client
535 .send(&ClientMessage::CreateRoom {
536 room_name: config.room_name.clone(),
537 max_players: config.max_players,
538 format: config.format.clone(),
539 protocol_version: PROTOCOL_VERSION,
540 hosted: !config.host_plays,
541 engine: engine_kind(config),
542 draft_config: None,
543 sealed_config: None,
544 official_key: config.official_key.clone(),
545 password: config.room_password.clone(),
546 reconnect_timeout_s: config.reconnect_timeout_s,
547 })
548 .await?;
549 info!(room_name = %config.room_name, "creating room");
550 wait_for_host_room(client, config, snapshot).await?
551 };
552
553 if config.host_plays {
554 seat_client(client, &config.host_deck).await?;
555 } else {
556 info!(
557 username = %config.username,
558 "hosting room without occupying a player seat"
559 );
560 }
561 Ok(room_id)
562}
563
564async fn reestablish_room(
565 client: &mut RelayClient,
566 config: &Config,
567 engine_session: &SharedEngineSession,
568 snapshot: &SharedHostSnapshot,
569) -> Result<String, Box<dyn std::error::Error + Send + Sync>> {
570 let resume = {
571 let snap = snapshot.lock().map_err(|error| error.to_string())?;
572 match (&snap.resume_token, &snap.game, &snap.room_info) {
573 (Some(token), Some(game), Some(room_info)) => {
574 Some(resume_room_request(config, token, game, room_info))
575 }
576 _ => None,
577 }
578 };
579 let game_active = engine_session
580 .lock()
581 .map(|guard| guard.is_some())
582 .unwrap_or(false);
583
584 if let (Some(request), true) = (resume, game_active) {
585 let room_id = request.room_id.clone();
586 client.send(&ClientMessage::ResumeRoom(request)).await?;
587 match wait_for_room_resumed(client, engine_session, snapshot).await? {
588 Some(room) => {
589 let (last_state, seat_states, prompts, player_order, pending_end_game) = {
590 let mut snap = snapshot.lock().map_err(|error| error.to_string())?;
591 snap.room_info = Some(room);
592 (
593 snap.last_state.clone(),
594 snap.last_state_by_slot.clone(),
595 snap.pending_prompts.clone(),
596 snap.game
597 .as_ref()
598 .map(|game| game.player_order.clone())
599 .unwrap_or_default(),
600 snap.pending_end_game.clone(),
601 )
602 };
603 if let Some(state) = last_state {
604 client
605 .send(&ClientMessage::BroadcastState {
606 state,
607 target_player: None,
608 })
609 .await?;
610 }
611 for (slot, state) in seat_states {
612 let Some(target_player) =
613 parse_player_slot(&slot).and_then(|index| player_order.get(index).cloned())
614 else {
615 warn!(slot, "cannot replay private state without a mapped player");
616 continue;
617 };
618 client
619 .send(&ClientMessage::BroadcastState {
620 state,
621 target_player: Some(target_player),
622 })
623 .await?;
624 }
625 for (slot, prompt) in prompts {
626 let Some(target_player) =
627 parse_player_slot(&slot).and_then(|index| player_order.get(index).cloned())
628 else {
629 warn!(slot, "cannot replay prompt without a mapped player");
630 continue;
631 };
632 client
633 .send(&ClientMessage::BroadcastState {
634 state: prompt,
635 target_player: Some(target_player),
636 })
637 .await?;
638 }
639 if let Some(game_id) = pending_end_game {
640 client.send(&ClientMessage::EndGame { game_id }).await?;
641 }
642 info!(room_id, "room resumed; snapshot re-broadcast");
643 return Ok(room_id);
644 }
645 None => {
646 warn!(room_id, "room resume rejected; abandoning hosted game");
647 abort_engine_session(engine_session);
648 }
649 }
650 }
651
652 establish_room(client, config, snapshot).await
653}
654
655fn engine_kind(config: &Config) -> EngineKind {
656 if matches!(config.backend, EngineBackendKind::Forge) {
657 EngineKind::Forge
658 } else {
659 EngineKind::Manabrew
660 }
661}
662
663fn resume_room_request(
664 config: &Config,
665 token: &str,
666 game: &GameStart,
667 room_info: &RoomInfo,
668) -> ResumeRoomRequest {
669 ResumeRoomRequest {
670 room_id: room_info.room_id.clone(),
671 resume_token: token.to_string(),
672 room_name: room_info.room_name.clone(),
673 max_players: room_info.max_players,
674 format: room_info.format.clone(),
675 hosted: !config.host_plays,
676 engine: engine_kind(config),
677 official_key: config.official_key.clone(),
678 password: config.room_password.clone(),
679 reconnect_timeout_s: Some(room_info.reconnect_timeout_s),
680 draft_config: room_info.draft_config.clone(),
681 sealed_config: room_info.sealed_config.clone(),
682 player_order: game.player_order.clone(),
683 player_decks: game.player_decks.clone(),
684 starting_life: game.starting_life,
685 bot_players: room_info
686 .players
687 .iter()
688 .filter(|player| player.is_bot)
689 .map(|player| player.username.clone())
690 .collect(),
691 game_id: game.game_id.clone(),
692 }
693}
694
695async fn wait_for_room_resumed(
696 client: &mut RelayClient,
697 engine_session: &SharedEngineSession,
698 snapshot: &SharedHostSnapshot,
699) -> Result<Option<RoomInfo>, Box<dyn std::error::Error + Send + Sync>> {
700 loop {
701 match client.recv().await? {
702 Some(ServerMessage::RoomResumed { room }) => {
703 if room.status == RoomStatus::Lobby {
704 abort_stale_engine_session(
705 engine_session,
706 snapshot,
707 "resumed room is no longer in game",
708 );
709 } else {
710 concede_missing_active_seats(engine_session, snapshot, &room);
711 }
712 return Ok(Some(room));
713 }
714 Some(ServerMessage::RoomUpdate { room }) => {
715 if room.status == RoomStatus::Lobby {
716 abort_stale_engine_session(
717 engine_session,
718 snapshot,
719 "room reset while reconnecting",
720 );
721 } else {
722 concede_missing_active_seats(engine_session, snapshot, &room);
723 }
724 if let Ok(mut snap) = snapshot.lock() {
725 snap.room_info = Some(room);
726 }
727 }
728 Some(ServerMessage::StateUpdate { from_player, state }) => {
729 let Ok(envelope) = serde_json::from_value::<StateEnvelope>(state.clone()) else {
730 continue;
731 };
732 match envelope {
733 StateEnvelope::Response { .. } => {
734 route_remote_response(engine_session, snapshot, &from_player, &state);
735 }
736 StateEnvelope::Directive {
737 from_player: claimed_slot,
738 directive,
739 } => route_remote_directive(
740 engine_session,
741 snapshot,
742 &from_player,
743 &claimed_slot,
744 &directive,
745 ),
746 _ => {}
747 }
748 }
749 Some(ServerMessage::Error { code, message }) => {
750 warn!(code, message, "resume rejected by relay");
751 return Ok(None);
752 }
753 Some(other) => debug!(?other, "ignored message while waiting for room resume"),
754 None => return Err("relay closed while waiting for room resume".into()),
755 }
756 }
757}
758
759fn cancel_engine(engine_session: &SharedEngineSession) {
760 if let Ok(guard) = engine_session.lock() {
761 if let Some(EngineSession::Forge { cancel, .. }) = guard.as_ref() {
762 cancel.store(true, Ordering::Relaxed);
763 }
764 }
765}
766
767fn abort_engine_session(engine_session: &SharedEngineSession) {
768 cancel_engine(engine_session);
769 clear_engine_session(engine_session);
772}
773
774fn abort_stale_engine_session(
775 engine_session: &SharedEngineSession,
776 snapshot: &SharedHostSnapshot,
777 reason: &str,
778) {
779 let stale_game_id = engine_session
780 .lock()
781 .ok()
782 .and_then(|guard| guard.as_ref().map(|s| s.game_id().to_string()));
783 let Some(stale_game_id) = stale_game_id else {
784 return;
785 };
786 warn!(game_id = %stale_game_id, reason, "aborting stale engine session");
787 abort_engine_session(engine_session);
788 clear_game_snapshot(snapshot, &stale_game_id);
789}
790
791fn clear_game_snapshot(snapshot: &SharedHostSnapshot, game_id: &str) {
792 if let Ok(mut snap) = snapshot.lock() {
793 if snap
794 .game
795 .as_ref()
796 .is_some_and(|game| game.game_id == game_id)
797 {
798 snap.game = None;
799 snap.last_state = None;
800 snap.last_state_by_slot.clear();
801 snap.pending_prompts.clear();
802 snap.pending_end_game = None;
803 }
804 }
805}
806
807fn end_game_without_humans(
813 engine_session: &SharedEngineSession,
814 snapshot: &SharedHostSnapshot,
815 outbound_tx: &tokio_mpsc::UnboundedSender<ClientMessage>,
816 room: &RoomInfo,
817) {
818 if room.players.iter().any(|seat| !seat.is_bot) {
819 return;
820 }
821 let game_id = engine_session
822 .lock()
823 .ok()
824 .and_then(|guard| guard.as_ref().map(|s| s.game_id().to_string()));
825 let Some(game_id) = game_id else {
826 return;
827 };
828 info!(game_id, "no human seats remain; ending hosted game");
829 abort_engine_session(engine_session);
830 clear_game_snapshot(snapshot, &game_id);
831 let _ = outbound_tx.send(ClientMessage::EndGame { game_id });
832}
833
834fn spawn_bots(config: &Config, decks: &[DeckSelection], room_id: &str, bot_state: &SharedBotState) {
835 stop_bots(bot_state);
836 let mut guard = match bot_state.lock() {
837 Ok(guard) => guard,
838 Err(error) => {
839 warn!(%error, "bot state lock poisoned");
840 return;
841 }
842 };
843 let open_seats = config.max_players.saturating_sub(1) as usize;
844 for (index, deck) in decks.iter().take(open_seats).enumerate() {
845 let username = if index == 0 {
846 config.bot_username.clone()
847 } else {
848 format!("{} {}", config.bot_username, index + 1)
849 };
850 let relay_url = config.relay_url.clone();
851 let bot_config = BotConfig {
852 username,
853 password: config.password.clone(),
854 room_id: room_id.to_string(),
855 room_password: config.room_password.clone(),
856 deck_name: deck.name.clone(),
857 deck: deck.deck.clone(),
858 commander_name: deck.commander_name.clone(),
859 agent: AgentKind::Simple,
860 answer_delay_ms: None,
861 };
862 let shutdown = Arc::new(tokio::sync::Notify::new());
863 let bot_shutdown = shutdown.clone();
864 let handle = tokio::spawn(async move {
865 if let Err(error) = run_bot(relay_url, bot_config, bot_shutdown).await {
866 error!(%error, "bot task exited");
867 }
868 });
869 guard.push(BotState { handle, shutdown });
870 }
871}
872
873fn stop_bots(bot_state: &SharedBotState) {
874 match bot_state.lock() {
875 Ok(mut state) => {
876 for bot in state.drain(..) {
877 bot.shutdown.notify_one();
878 let mut handle = bot.handle;
879 tokio::spawn(async move {
880 if time::timeout(BOT_STOP_TIMEOUT, &mut handle).await.is_err() {
881 handle.abort();
882 }
883 });
884 }
885 }
886 Err(error) => warn!(%error, "bot state lock poisoned"),
887 }
888}
889
890async fn wait_for_host_room(
891 client: &mut RelayClient,
892 config: &Config,
893 snapshot: &SharedHostSnapshot,
894) -> Result<String, Box<dyn std::error::Error + Send + Sync>> {
895 loop {
896 match client.recv().await? {
897 Some(ServerMessage::RoomCreated {
898 room_id,
899 room_name,
900 room,
901 resume_token,
902 }) => {
903 info!(room_id, room_name, "room created");
904 if let Ok(mut snap) = snapshot.lock() {
905 snap.resume_token = resume_token;
906 snap.room_info = Some(room);
907 }
908 return Ok(room_id);
909 }
910 Some(ServerMessage::RoomUpdate { room }) if room.host == config.username => {
911 info!(room_id = %room.room_id, room_name = %room.room_name, "host room update");
912 return Ok(room.room_id);
913 }
914 Some(ServerMessage::AuthResult { success, error, .. }) => {
915 if !success {
916 return Err(format!("authentication failed: {error:?}").into());
917 }
918 }
919 Some(ServerMessage::Error { code, message }) => {
920 return Err(format!("server error while creating room: {code}: {message}").into());
921 }
922 Some(other) => debug!(?other, "ignored message while waiting for room creation"),
923 None => return Err("relay closed while waiting for room creation".into()),
924 }
925 }
926}
927
928async fn seat_client(
929 client: &mut RelayClient,
930 deck: &DeckSelection,
931) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
932 client
933 .send(&ClientMessage::SetDeckSelection {
934 deck_name: deck.name.clone(),
935 deck: deck.deck.clone(),
936 published_deck_id: None,
937 commander_name: deck.commander_name.clone(),
938 avatar: None,
939 })
940 .await?;
941 client
942 .send(&ClientMessage::SetReady { ready: true })
943 .await?;
944 info!(
945 username = %client.username,
946 deck = %deck.name,
947 cards = deck.deck.cards.len(),
948 "selected deck and marked ready"
949 );
950 Ok(())
951}
952
953async fn run_client_loop(
954 client: &mut RelayClient,
955 config: &Config,
956 room_id: &str,
957 engine_session: &SharedEngineSession,
958 snapshot: &SharedHostSnapshot,
959 bot_state: &SharedBotState,
960 outbound_tx: &tokio_mpsc::UnboundedSender<ClientMessage>,
961 outbound_rx: &mut tokio_mpsc::UnboundedReceiver<ClientMessage>,
962 cancel: &RoomCancel,
963) -> LoopExit {
964 let mut heartbeat = time::interval(Duration::from_secs(30));
965 heartbeat.set_missed_tick_behavior(time::MissedTickBehavior::Delay);
966 let mut drain_tick = time::interval(Duration::from_secs(5));
967 drain_tick.set_missed_tick_behavior(time::MissedTickBehavior::Delay);
968 let mut bot_usernames: HashSet<String> = HashSet::new();
969
970 loop {
971 tokio::select! {
972 _ = cancel.notified() => {
973 info!(username = %client.username, "room host cancelled; shutting down");
974 cancel_engine(engine_session);
975 return LoopExit::Cancelled;
976 }
977 _ = drain_tick.tick() => {
978 let idle = engine_session
979 .lock()
980 .map(|guard| guard.is_none())
981 .unwrap_or(false);
982 if draining() && idle {
983 info!(username = %client.username, "draining and no active game; closing room");
984 return LoopExit::Cancelled;
985 }
986 }
987 _ = heartbeat.tick() => {
988 if let Err(error) = client.broadcast_room_message(SELF_HOSTED_NODE_PROTOCOL, json!({
989 "type": "heartbeat",
990 "node": client.username,
991 "capacity": config.max_games,
992 })).await {
993 warn!(%error, username = %client.username, "relay send failed");
994 return LoopExit::Disconnected;
995 }
996 }
997 outbound = outbound_rx.recv() => {
998 let Some(outbound) = outbound else {
999 warn!(username = %client.username, "outbound channel closed");
1000 return LoopExit::Cancelled;
1001 };
1002 if let Err(error) = client.send(&outbound).await {
1003 warn!(%error, username = %client.username, "relay send failed");
1004 return LoopExit::Disconnected;
1005 }
1006 }
1007 message = client.recv() => {
1008 let message = match message {
1009 Ok(Some(message)) => message,
1010 Ok(None) => {
1011 warn!(username = %client.username, "relay websocket closed");
1012 return LoopExit::Disconnected;
1013 }
1014 Err(error) => {
1015 warn!(%error, username = %client.username, "relay receive failed");
1016 return LoopExit::Disconnected;
1017 }
1018 };
1019 if let Err(error) = handle_server_message(
1020 client,
1021 config,
1022 room_id,
1023 engine_session,
1024 snapshot,
1025 bot_state,
1026 outbound_tx,
1027 &mut bot_usernames,
1028 message,
1029 ).await {
1030 warn!(%error, username = %client.username, "relay send failed");
1031 return LoopExit::Disconnected;
1032 }
1033 }
1034 }
1035 }
1036}
1037
1038async fn handle_server_message(
1039 client: &mut RelayClient,
1040 config: &Config,
1041 room_id: &str,
1042 engine_session: &SharedEngineSession,
1043 snapshot: &SharedHostSnapshot,
1044 bot_state: &SharedBotState,
1045 outbound_tx: &tokio_mpsc::UnboundedSender<ClientMessage>,
1046 bot_usernames: &mut HashSet<String>,
1047 message: ServerMessage,
1048) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
1049 match message {
1050 ServerMessage::RoomUpdate { room } => {
1051 if let Ok(mut snap) = snapshot.lock() {
1052 snap.room_info = Some(room.clone());
1053 }
1054 *bot_usernames = room
1055 .players
1056 .iter()
1057 .filter(|p| p.is_bot)
1058 .map(|p| p.username.clone())
1059 .collect();
1060 log_room_update(&client.username, &room);
1061 if room.status == RoomStatus::Lobby {
1062 abort_stale_engine_session(engine_session, snapshot, "room reset to lobby");
1063 } else {
1064 end_game_without_humans(engine_session, snapshot, outbound_tx, &room);
1065 }
1066 maybe_auto_start_room(client, config, &room).await?;
1067 }
1068 ServerMessage::StateUpdate { from_player, state } => {
1069 handle_state_update(
1070 client,
1071 config,
1072 room_id,
1073 engine_session,
1074 snapshot,
1075 bot_state,
1076 from_player,
1077 state,
1078 )
1079 .await?;
1080 }
1081 ServerMessage::ReadyStateChanged { username, ready } => {
1082 info!(username, ready, observer = %client.username, "ready changed");
1083 }
1084 ServerMessage::PlayerJoined { username, room_id } => {
1085 info!(username, room_id, observer = %client.username, "player joined");
1086 }
1087 ServerMessage::PlayerLeft { username, room_id } => {
1088 info!(username, room_id, observer = %client.username, "player left");
1089 let leaver_still_needed =
1090 active_player_usernames(snapshot).is_some_and(|active| active.contains(&username));
1091 if leaver_still_needed {
1092 if let Some(index) = seat_index_of(snapshot, &username) {
1093 info!(
1094 username,
1095 index, "player abandoned mid-game; conceding their seat"
1096 );
1097 concede_seat(engine_session, index);
1098 }
1099 }
1100 }
1101 ServerMessage::GameStarted {
1102 room_id,
1103 game_id,
1104 player_order,
1105 player_decks,
1106 starting_life,
1107 } => {
1108 info!(room_id, game_id, ?player_order, observer = %client.username, "game started");
1109 abort_stale_engine_session(engine_session, snapshot, "relay started a new game");
1110 if let Ok(mut snap) = snapshot.lock() {
1111 snap.game = Some(GameStart {
1112 game_id: game_id.clone(),
1113 player_order: player_order.clone(),
1114 player_decks: player_decks.clone(),
1115 starting_life,
1116 });
1117 snap.last_state = None;
1118 snap.last_state_by_slot.clear();
1119 snap.pending_prompts.clear();
1120 snap.pending_end_game = None;
1121 }
1122 maybe_start_hosted_engine(
1123 config,
1124 engine_session,
1125 snapshot,
1126 outbound_tx,
1127 game_id,
1128 player_order,
1129 player_decks,
1130 starting_life,
1131 bot_usernames,
1132 );
1133 }
1134 ServerMessage::ServerShuttingDown { reconnect_in_s } => {
1135 info!(reconnect_in_s, observer = %client.username, "relay is restarting");
1136 }
1137 ServerMessage::Error { code, message } => {
1138 warn!(code, message, observer = %client.username, "server error");
1139 }
1140 other => {
1141 debug!(?other, observer = %client.username, "server message");
1142 }
1143 }
1144
1145 Ok(())
1146}
1147
1148async fn handle_state_update(
1149 client: &mut RelayClient,
1150 config: &Config,
1151 room_id: &str,
1152 engine_session: &SharedEngineSession,
1153 snapshot: &SharedHostSnapshot,
1154 bot_state: &SharedBotState,
1155 from_player: String,
1156 state: Value,
1157) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
1158 let Ok(envelope) = serde_json::from_value::<StateEnvelope>(state.clone()) else {
1159 debug!(from_player, state = %state, observer = %client.username, "state update");
1160 return Ok(());
1161 };
1162
1163 match envelope {
1164 StateEnvelope::Response { .. } => {
1165 route_remote_response(engine_session, snapshot, &from_player, &state);
1166 Ok(())
1167 }
1168 StateEnvelope::Directive {
1169 from_player: claimed_slot,
1170 directive,
1171 } => {
1172 route_remote_directive(
1173 engine_session,
1174 snapshot,
1175 &from_player,
1176 &claimed_slot,
1177 &directive,
1178 );
1179 Ok(())
1180 }
1181 StateEnvelope::RoomRelay {
1182 protocol,
1183 from_player: _,
1184 room_id: requested_room_id,
1185 payload,
1186 ..
1187 } => {
1188 info!(
1189 from_player,
1190 protocol,
1191 payload = %payload,
1192 observer = %client.username,
1193 "room relay message"
1194 );
1195 if protocol != SELF_HOSTED_NODE_PROTOCOL {
1196 return Ok(());
1197 }
1198 let payload_type = payload.get("type").and_then(Value::as_str);
1199 match payload_type {
1200 Some("removeBot") => {
1201 info!(observer = %client.username, "received removeBot request");
1202 stop_bots(bot_state);
1203 }
1204 Some("spawnBot") => {
1205 let effective_room_id = requested_room_id.as_deref().unwrap_or(room_id);
1206 if effective_room_id == room_id {
1207 info!(observer = %client.username, room_id, "received spawnBot request");
1208 let bot_decks = bot_decks_from_payload(config, &payload);
1209 spawn_bots(config, &bot_decks, room_id, bot_state);
1210 } else {
1211 debug!(
1212 observer = %client.username,
1213 current_room_id = room_id,
1214 requested_room_id = effective_room_id,
1215 "ignoring spawnBot request for another room"
1216 );
1217 }
1218 }
1219 _ => {}
1220 }
1221 Ok(())
1222 }
1223 _ => {
1224 debug!(from_player, state = %state, observer = %client.username, "state update");
1225 Ok(())
1226 }
1227 }
1228}
1229
1230fn bot_decks_from_payload(config: &Config, payload: &Value) -> Vec<DeckSelection> {
1231 let parsed = serde_json::from_value::<SpawnBotPayload>(payload.clone()).ok();
1232 let requested = match parsed {
1233 Some(SpawnBotPayload {
1234 decks: Some(decks), ..
1235 }) if !decks.is_empty() => decks,
1236 Some(SpawnBotPayload {
1237 deck: Some(deck), ..
1238 }) => vec![deck],
1239 _ => return vec![config.bot_deck.clone()],
1240 };
1241 requested
1242 .into_iter()
1243 .map(deck_selection_from_payload)
1244 .collect()
1245}
1246
1247fn deck_selection_from_payload(deck: SpawnBotDeckPayload) -> DeckSelection {
1248 info!(
1249 deck = %deck.deck_name,
1250 cards = deck.deck.cards.len(),
1251 commander = ?deck.commander_name,
1252 "using requested bot deck"
1253 );
1254 let commander_name = deck.commander_name.or_else(|| {
1255 deck.deck
1256 .commanders
1257 .as_ref()
1258 .and_then(|commanders| commanders.first())
1259 .map(|commander| commander.identity.name.clone())
1260 });
1261 DeckSelection {
1262 name: deck.deck_name,
1263 deck: deck.deck,
1264 commander_name,
1265 }
1266}
1267
1268async fn maybe_auto_start_room(
1269 client: &mut RelayClient,
1270 config: &Config,
1271 room: &RoomInfo,
1272) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
1273 if !config.auto_start || draining() {
1274 return Ok(());
1275 }
1276 if config.format == GameFormat::Any {
1277 return Ok(());
1278 }
1279 if room.host != config.username
1280 || room.status != RoomStatus::Lobby
1281 || room.players.len() < config.max_players as usize
1282 {
1283 return Ok(());
1284 }
1285 let all_ready = room
1286 .players
1287 .iter()
1288 .all(|player| player.connected && player.ready && player.selected_deck_name.is_some());
1289 if all_ready {
1290 info!(
1291 room_id = %room.room_id,
1292 players = room.players.len(),
1293 "all players ready; auto-starting hosted game"
1294 );
1295 client
1296 .send(&ClientMessage::StartGame { format: None })
1297 .await?;
1298 }
1299 Ok(())
1300}
1301
1302const ENGINE_STACK_BYTES: usize = 64 * 1024 * 1024;
1303
1304fn spawn_engine_thread<F: FnOnce() + Send + 'static>(body: F) {
1305 if let Err(error) = thread::Builder::new()
1306 .name("hosted-engine".to_string())
1307 .stack_size(ENGINE_STACK_BYTES)
1308 .spawn(body)
1309 {
1310 error!(%error, "failed to spawn hosted engine thread");
1311 }
1312}
1313
1314fn is_commander_variant(format: GameFormat) -> bool {
1317 match format {
1318 GameFormat::Any | GameFormat::Commander | GameFormat::Brawl | GameFormat::Oathbreaker => {
1319 true
1320 }
1321 GameFormat::Standard
1322 | GameFormat::Pioneer
1323 | GameFormat::Modern
1324 | GameFormat::Legacy
1325 | GameFormat::Vintage
1326 | GameFormat::Pauper
1327 | GameFormat::Draft
1328 | GameFormat::Sealed => false,
1329 }
1330}
1331
1332fn java_game_variant(format: GameFormat) -> &'static str {
1335 match format {
1336 GameFormat::Any => "",
1337 GameFormat::Commander => "Commander",
1338 GameFormat::Brawl => "Brawl",
1339 GameFormat::Oathbreaker => "Oathbreaker",
1340 GameFormat::Standard
1341 | GameFormat::Pioneer
1342 | GameFormat::Modern
1343 | GameFormat::Legacy
1344 | GameFormat::Vintage
1345 | GameFormat::Pauper
1346 | GameFormat::Draft
1347 | GameFormat::Sealed => "Constructed",
1348 }
1349}
1350
1351fn maybe_start_hosted_engine(
1352 config: &Config,
1353 engine_session: &SharedEngineSession,
1354 snapshot: &SharedHostSnapshot,
1355 outbound_tx: &tokio_mpsc::UnboundedSender<ClientMessage>,
1356 game_id: String,
1357 player_order: Vec<String>,
1358 player_decks: Vec<PlayerDeckInfo>,
1359 starting_life: i32,
1360 bot_usernames: &HashSet<String>,
1361) {
1362 if !config.engine_enabled {
1363 debug!("hosted engine disabled for this node");
1364 return;
1365 }
1366 let backend = config.backend;
1367 if !backend.is_supported() {
1368 if let Err(error) = java_backend::JavaRuntimeConfig::from_env().validate() {
1369 warn!(
1370 backend = backend.label(),
1371 %error,
1372 "hosted engine backend runtime is not ready"
1373 );
1374 return;
1375 }
1376 warn!(
1377 backend = backend.label(),
1378 message = java_backend::unsupported_message(),
1379 "hosted engine backend is not implemented yet"
1380 );
1381 return;
1382 }
1383
1384 let local_player_index = if config.host_plays {
1385 match player_order
1386 .iter()
1387 .position(|name| name == &config.username)
1388 {
1389 Some(index) => Some(index),
1390 None => {
1391 warn!(
1392 username = %config.username,
1393 ?player_order,
1394 "room node is configured as a player but is not in player order; not starting engine"
1395 );
1396 return;
1397 }
1398 }
1399 } else {
1400 None
1401 };
1402
1403 let mut guard = match engine_session.lock() {
1404 Ok(guard) => guard,
1405 Err(error) => {
1406 warn!(%error, "engine session lock poisoned");
1407 return;
1408 }
1409 };
1410 if let Some(session) = guard.as_ref() {
1411 warn!(
1412 game_id,
1413 stale_game_id = session.game_id(),
1414 "engine session still present at game start; not starting engine"
1415 );
1416 return;
1417 }
1418
1419 let session_handle = engine_session.clone();
1420 let num_players = player_order.len();
1421 if num_players < 2 {
1422 warn!(num_players, "not enough players to start hosted engine");
1423 return;
1424 }
1425
1426 let mut deck_map: HashMap<String, PlayerDeckInfo> = player_decks
1427 .into_iter()
1428 .map(|deck| (deck.username.clone(), deck))
1429 .collect();
1430 let mut ordered_decks = Vec::with_capacity(num_players);
1431 let mut commander_names = Vec::with_capacity(num_players);
1432 let mut ai_player_indices = Vec::new();
1433 for (index, username) in player_order.iter().enumerate() {
1434 let Some(deck) = deck_map.remove(username) else {
1435 warn!(username, "missing deck for player; not starting engine");
1436 return;
1437 };
1438 if config.forge_ai && bot_usernames.contains(username) {
1439 ai_player_indices.push(index);
1440 }
1441 ordered_decks.push(deck.deck);
1442 commander_names.push(deck.commander_name);
1443 }
1444
1445 let player_names = player_order;
1446 let room_format = snapshot
1447 .lock()
1448 .ok()
1449 .and_then(|snap| snap.room_info.as_ref().map(|room| room.format.clone()))
1450 .unwrap_or(GameFormat::Any);
1451 let commander_variant = is_commander_variant(room_format.clone());
1452 let game_variant = java_game_variant(room_format).to_string();
1453
1454 match backend {
1455 EngineBackendKind::Manabrew => {
1456 let (remote_prompt_tx, remote_prompt_rx) = std_mpsc::channel::<(usize, AgentMessage)>();
1457 let mut remote_response_txs = HashMap::new();
1458 let mut remote_response_rxs = Vec::new();
1459 for i in 0..num_players {
1460 if Some(i) == local_player_index {
1461 continue;
1462 }
1463 let (response_tx, response_rx) = std_mpsc::channel::<ClientToServerMessage>();
1464 remote_response_txs.insert(i, response_tx);
1465 remote_response_rxs.push((i, response_rx));
1466 }
1467 *guard = Some(EngineSession::Manabrew {
1468 game_id: game_id.clone(),
1469 remote_response_txs,
1470 });
1471 drop(guard);
1472
1473 spawn_remote_prompt_forwarder(
1474 outbound_tx.clone(),
1475 snapshot.clone(),
1476 engine_session.clone(),
1477 game_id.clone(),
1478 remote_prompt_rx,
1479 Some(player_names.clone()),
1480 config.state_delta,
1481 );
1482 let (game_over_tx, game_over_rx) = std_mpsc::channel::<HostedGameOver>();
1483 spawn_game_over_forwarder(
1484 outbound_tx.clone(),
1485 game_over_rx,
1486 engine_session.clone(),
1487 snapshot.clone(),
1488 game_id.clone(),
1489 Some(player_names.clone()),
1490 );
1491 let outbound_tx = outbound_tx.clone();
1492 let snapshot = snapshot.clone();
1493 spawn_engine_thread(move || {
1494 info!(
1495 game_id,
1496 backend = backend.label(),
1497 players = num_players,
1498 local_player_index,
1499 "starting hosted engine thread"
1500 );
1501 crate::metrics::record_engine_session_started();
1502 let started = Instant::now();
1503 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1504 rust_backend::run_hosted_engine_game(
1505 game_id.clone(),
1506 player_names,
1507 ordered_decks,
1508 commander_names,
1509 local_player_index,
1510 starting_life,
1511 remote_prompt_tx,
1512 remote_response_rxs,
1513 game_over_tx,
1514 )
1515 }));
1516 finish_hosted_engine(
1517 result,
1518 &game_id,
1519 num_players,
1520 started,
1521 &outbound_tx,
1522 &session_handle,
1523 &snapshot,
1524 );
1525 });
1526 }
1527 EngineBackendKind::Forge => {
1528 let (remote_prompt_tx, remote_prompt_rx) = std_mpsc::channel::<(usize, AgentMessage)>();
1529 let mut remote_response_txs = HashMap::new();
1530 let mut remote_response_rxs = Vec::new();
1531 for i in 0..num_players {
1532 if Some(i) == local_player_index {
1533 continue;
1534 }
1535 let (response_tx, response_rx) = std_mpsc::channel::<ClientToServerMessage>();
1536 remote_response_txs.insert(i, response_tx);
1537 remote_response_rxs.push((i, response_rx));
1538 }
1539 let cancel = Arc::new(AtomicBool::new(false));
1540 *guard = Some(EngineSession::Forge {
1541 game_id: game_id.clone(),
1542 remote_response_txs,
1543 cancel: cancel.clone(),
1544 });
1545 drop(guard);
1546
1547 spawn_remote_prompt_forwarder(
1548 outbound_tx.clone(),
1549 snapshot.clone(),
1550 engine_session.clone(),
1551 game_id.clone(),
1552 remote_prompt_rx,
1553 Some(player_names.clone()),
1554 config.state_delta,
1555 );
1556 let (game_over_tx, game_over_rx) = std_mpsc::channel::<HostedGameOver>();
1557 spawn_game_over_forwarder(
1558 outbound_tx.clone(),
1559 game_over_rx,
1560 engine_session.clone(),
1561 snapshot.clone(),
1562 game_id.clone(),
1563 Some(player_names.clone()),
1564 );
1565 let outbound_tx = outbound_tx.clone();
1566 let snapshot = snapshot.clone();
1567 spawn_engine_thread(move || {
1568 info!(
1569 game_id,
1570 backend = backend.label(),
1571 players = num_players,
1572 local_player_index,
1573 "starting hosted engine thread"
1574 );
1575 crate::metrics::record_engine_session_started();
1576 let started = Instant::now();
1577 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1578 java_backend::run_hosted_engine_game(
1579 game_id.clone(),
1580 player_names,
1581 ordered_decks,
1582 commander_names,
1583 commander_variant,
1584 game_variant,
1585 local_player_index,
1586 ai_player_indices,
1587 starting_life,
1588 remote_prompt_tx,
1589 remote_response_rxs,
1590 game_over_tx,
1591 cancel,
1592 )
1593 }));
1594 finish_hosted_engine(
1595 result,
1596 &game_id,
1597 num_players,
1598 started,
1599 &outbound_tx,
1600 &session_handle,
1601 &snapshot,
1602 );
1603 });
1604 }
1605 }
1606}
1607
1608fn clear_engine_session(engine_session: &SharedEngineSession) {
1609 match engine_session.lock() {
1610 Ok(mut guard) => *guard = None,
1611 Err(error) => warn!(%error, "engine session lock poisoned on reset"),
1612 }
1613}
1614
1615fn active_player_usernames(snapshot: &SharedHostSnapshot) -> Option<HashSet<String>> {
1620 let snap = snapshot.lock().ok()?;
1621 let game = snap.game.as_ref()?;
1622 let mut active: HashSet<String> = game.player_order.iter().cloned().collect();
1623 let players = snap
1624 .last_state
1625 .as_ref()
1626 .and_then(|state| state.pointer("/state/gameView/players"))
1627 .and_then(Value::as_array);
1628 if let Some(players) = players {
1629 for player in players {
1630 let playing = player
1631 .get("status")
1632 .and_then(Value::as_str)
1633 .is_none_or(|status| status == "playing");
1634 if playing {
1635 continue;
1636 }
1637 let index = player
1638 .get("id")
1639 .and_then(Value::as_str)
1640 .and_then(parse_player_slot);
1641 if let Some(username) = index.and_then(|i| game.player_order.get(i)) {
1642 active.remove(username);
1643 }
1644 }
1645 }
1646 Some(active)
1647}
1648
1649fn seat_index_of(snapshot: &SharedHostSnapshot, username: &str) -> Option<usize> {
1650 let snap = snapshot.lock().ok()?;
1651 let game = snap.game.as_ref()?;
1652 game.player_order.iter().position(|name| name == username)
1653}
1654
1655fn concede_missing_active_seats(
1656 engine_session: &SharedEngineSession,
1657 snapshot: &SharedHostSnapshot,
1658 room: &RoomInfo,
1659) {
1660 let Some(active) = active_player_usernames(snapshot) else {
1661 return;
1662 };
1663 let present: HashSet<&str> = room
1664 .players
1665 .iter()
1666 .map(|player| player.username.as_str())
1667 .collect();
1668 for username in active {
1669 if !present.contains(username.as_str()) {
1670 if let Some(index) = seat_index_of(snapshot, &username) {
1671 concede_seat(engine_session, index);
1672 }
1673 }
1674 }
1675}
1676
1677fn concede_seat(engine_session: &SharedEngineSession, player_index: usize) {
1678 send_seat_message(
1679 engine_session,
1680 player_index,
1681 ClientToServerMessage::Directive {
1682 directive: manabrew_protocol::transport::DirectiveInput::Concede,
1683 },
1684 );
1685}
1686
1687fn send_seat_message(
1688 engine_session: &SharedEngineSession,
1689 player_index: usize,
1690 message: ClientToServerMessage,
1691) {
1692 let guard = match engine_session.lock() {
1693 Ok(guard) => guard,
1694 Err(error) => {
1695 warn!(%error, "engine session lock poisoned");
1696 return;
1697 }
1698 };
1699 let Some(session) = guard.as_ref() else {
1700 debug!(player_index, "no engine session for seat action");
1701 return;
1702 };
1703 let txs = match session {
1704 EngineSession::Manabrew {
1705 remote_response_txs,
1706 ..
1707 }
1708 | EngineSession::Forge {
1709 remote_response_txs,
1710 ..
1711 } => remote_response_txs,
1712 };
1713 let Some(tx) = txs.get(&player_index) else {
1714 debug!(player_index, "no response channel for player");
1715 return;
1716 };
1717 if let Err(error) = tx.send(message) {
1718 warn!(player_index, %error, "failed to route seat message");
1719 }
1720}
1721
1722fn finish_hosted_engine(
1723 result: std::thread::Result<Result<(), String>>,
1724 game_id: &str,
1725 players: usize,
1726 started: Instant,
1727 outbound_tx: &tokio_mpsc::UnboundedSender<ClientMessage>,
1728 session_handle: &SharedEngineSession,
1729 snapshot: &SharedHostSnapshot,
1730) {
1731 let fatal = match result {
1732 Ok(Ok(())) => {
1733 info!(game_id, "hosted engine thread finished");
1734 None
1735 }
1736 Ok(Err(message)) => {
1737 error!(game_id, message, "hosted engine exited with a fatal error");
1738 Some(message)
1739 }
1740 Err(panic) => {
1741 let message = if let Some(message) = panic.downcast_ref::<String>() {
1742 message.clone()
1743 } else if let Some(message) = panic.downcast_ref::<&str>() {
1744 message.to_string()
1745 } else {
1746 "the host engine panicked".to_string()
1747 };
1748 error!(game_id, message, "hosted engine thread panicked");
1749 Some(message)
1750 }
1751 };
1752 crate::metrics::record_engine_session_finished(players, started, fatal.as_deref());
1753 let still_owner = session_handle
1754 .lock()
1755 .map(|guard| guard.as_ref().is_some_and(|s| s.game_id() == game_id))
1756 .unwrap_or(false);
1757 if let Some(message) = fatal {
1758 if still_owner {
1759 if let Ok(mut snap) = snapshot.lock() {
1760 if snap
1761 .game
1762 .as_ref()
1763 .is_some_and(|game| game.game_id == game_id)
1764 {
1765 snap.pending_end_game = Some(game_id.to_string());
1766 }
1767 }
1768 if let Ok(state) = serde_json::to_value(StateEnvelope::Fatal { message }) {
1769 let _ = outbound_tx.send(ClientMessage::BroadcastState {
1770 state,
1771 target_player: None,
1772 });
1773 }
1774 let _ = outbound_tx.send(ClientMessage::EndGame {
1775 game_id: game_id.to_string(),
1776 });
1777 } else {
1778 warn!(game_id, message, "stale engine session finished with error");
1779 }
1780 }
1781}
1782
1783fn route_remote_response(
1784 engine_session: &SharedEngineSession,
1785 snapshot: &SharedHostSnapshot,
1786 authenticated_username: &str,
1787 state: &Value,
1788) {
1789 let envelope: StateEnvelope = match serde_json::from_value(state.clone()) {
1790 Ok(envelope) => envelope,
1791 Err(error) => {
1792 warn!(%error, state = %state, "relay response invalid envelope");
1793 return;
1794 }
1795 };
1796 let StateEnvelope::Response {
1797 from_player,
1798 prompt_id,
1799 action: action_value,
1800 } = envelope
1801 else {
1802 warn!(state = %state, "expected response envelope");
1803 return;
1804 };
1805 let Some(player_index) =
1806 authenticated_player_index(snapshot, authenticated_username, &from_player, "response")
1807 else {
1808 return;
1809 };
1810 let action: PromptOutput = match serde_json::from_value(action_value) {
1811 Ok(action) => action,
1812 Err(error) => {
1813 warn!(from_player, %error, "relay response has invalid action");
1814 return;
1815 }
1816 };
1817
1818 let guard = match engine_session.lock() {
1819 Ok(guard) => guard,
1820 Err(error) => {
1821 warn!(%error, "engine session lock poisoned");
1822 return;
1823 }
1824 };
1825 let Some(session) = guard.as_ref() else {
1826 debug!(from_player, "no engine session for relay response");
1827 return;
1828 };
1829 let tx = match session {
1830 EngineSession::Manabrew {
1831 remote_response_txs,
1832 ..
1833 } => remote_response_txs.get(&player_index),
1834 EngineSession::Forge {
1835 remote_response_txs,
1836 ..
1837 } => {
1838 debug!(
1839 from_player,
1840 player_index, "routing relay response to java engine"
1841 );
1842 remote_response_txs.get(&player_index)
1843 }
1844 };
1845 let Some(tx) = tx else {
1846 debug!(from_player, player_index, "no response channel for player");
1847 return;
1848 };
1849 if let Err(error) = tx.send(ClientToServerMessage::Response { prompt_id, action }) {
1850 warn!(from_player, %error, "failed to route relay response");
1851 return;
1852 }
1853 drop(guard);
1854 if let Ok(mut snap) = snapshot.lock() {
1855 snap.pending_prompts.remove(&from_player);
1856 }
1857}
1858
1859fn route_remote_directive(
1860 engine_session: &SharedEngineSession,
1861 snapshot: &SharedHostSnapshot,
1862 authenticated_username: &str,
1863 claimed_slot: &str,
1864 directive: &Value,
1865) {
1866 let directive: DirectiveInput = match serde_json::from_value(directive.clone()) {
1867 Ok(directive) => directive,
1868 Err(error) => {
1869 warn!(claimed_slot, %error, "relay directive is invalid");
1870 return;
1871 }
1872 };
1873 let Some(player_index) =
1874 authenticated_player_index(snapshot, authenticated_username, claimed_slot, "directive")
1875 else {
1876 return;
1877 };
1878 info!(claimed_slot, player_index, ?directive, "routing directive");
1879 match directive {
1880 DirectiveInput::Concede => concede_seat(engine_session, player_index),
1881 }
1882}
1883
1884fn authenticated_player_index(
1885 snapshot: &SharedHostSnapshot,
1886 authenticated_username: &str,
1887 claimed_slot: &str,
1888 message_kind: &str,
1889) -> Option<usize> {
1890 let Some(player_index) = seat_index_of(snapshot, authenticated_username) else {
1891 warn!(
1892 authenticated_username,
1893 message_kind, "relay sender has no engine seat"
1894 );
1895 return None;
1896 };
1897 let authenticated_slot = player_slot(player_index);
1898 if claimed_slot != authenticated_slot {
1899 warn!(
1900 authenticated_username,
1901 claimed_slot,
1902 authenticated_slot,
1903 message_kind,
1904 "relay sender claimed another engine seat"
1905 );
1906 return None;
1907 }
1908 Some(player_index)
1909}
1910
1911fn stamp_fingerprint(mut envelope: Value) -> Value {
1914 let Some(state) = envelope.get("state") else {
1915 return envelope;
1916 };
1917 let fingerprint = manabrew_relay_protocol::state_delta::fingerprint(state);
1918 if let Some(object) = envelope.as_object_mut() {
1919 object.insert("fingerprint".to_string(), Value::String(fingerprint));
1920 }
1921 envelope
1922}
1923
1924fn patch_against_last(
1928 bases: &mut HashMap<usize, (Value, String)>,
1929 player_index: usize,
1930 per_seat: bool,
1931 slot: &str,
1932 envelope: &Value,
1933) -> Option<Value> {
1934 let next = envelope.get("state")?.clone();
1935 let fingerprint = envelope.get("fingerprint")?.as_str()?.to_string();
1936 let previous = bases.insert(player_index, (next.clone(), fingerprint.clone()));
1937 let (base_state, base) = previous?;
1938 let patch = manabrew_relay_protocol::state_delta::diff(&base_state, &next)?;
1939 serde_json::to_value(StateEnvelope::StateDelta {
1940 for_player: per_seat.then(|| slot.to_string()),
1941 base,
1942 fingerprint,
1943 patch,
1944 })
1945 .ok()
1946}
1947
1948fn spawn_remote_prompt_forwarder(
1949 outbound_tx: tokio_mpsc::UnboundedSender<ClientMessage>,
1950 snapshot: SharedHostSnapshot,
1951 engine_session: SharedEngineSession,
1952 game_id: String,
1953 remote_prompt_rx: std_mpsc::Receiver<(usize, AgentMessage)>,
1954 seat_usernames: Option<Vec<String>>,
1955 state_delta: bool,
1956) {
1957 thread::spawn(move || {
1958 let mut last_state_by_seat: HashMap<usize, Value> = HashMap::new();
1959 let mut delta_bases: HashMap<usize, (Value, String)> = HashMap::new();
1960 let mut last_state: Option<Value> = None;
1961 let mut last_display: Option<Value> = None;
1962 while let Ok((player_index, message)) = remote_prompt_rx.recv() {
1963 let Ok(session) = engine_session.lock() else {
1964 break;
1965 };
1966 if session
1967 .as_ref()
1968 .is_none_or(|session| session.game_id() != game_id)
1969 {
1970 break;
1971 }
1972 let per_seat = seat_usernames.is_some() && player_index != OBSERVER_SEAT;
1973 let slot = player_slot(player_index);
1974 let envelope = match &message {
1975 AgentMessage::State(state_update) if !per_seat => StateEnvelope::State {
1976 for_player: None,
1977 state: serde_json::to_value(state_update).unwrap_or(Value::Null),
1978 fingerprint: None,
1979 },
1980 _ => StateEnvelope::for_agent_message(slot.clone(), &message),
1981 };
1982 let Ok(state) = serde_json::to_value(envelope) else {
1983 continue;
1984 };
1985 let state = match &message {
1986 AgentMessage::State(_) if state_delta => stamp_fingerprint(state),
1987 _ => state,
1988 };
1989 match &message {
1990 AgentMessage::State(_) if per_seat => {
1991 if last_state_by_seat.get(&player_index) == Some(&state) {
1992 continue;
1993 }
1994 last_state_by_seat.insert(player_index, state.clone());
1995 }
1996 AgentMessage::State(_) if last_state.as_ref() == Some(&state) => continue,
1997 AgentMessage::State(_) => last_state = Some(state.clone()),
1998 AgentMessage::Display(_) if last_display.as_ref() == Some(&state) => continue,
1999 AgentMessage::Display(_) => last_display = Some(state.clone()),
2000 AgentMessage::Prompt(_) | AgentMessage::Error(_) => {}
2001 }
2002 if let Ok(mut snap) = snapshot.lock() {
2003 match &message {
2004 AgentMessage::State(_) if per_seat => {
2005 snap.last_state_by_slot.insert(slot.clone(), state.clone());
2006 }
2007 AgentMessage::State(_) => snap.last_state = Some(state.clone()),
2008 AgentMessage::Prompt(_) => {
2009 snap.pending_prompts.insert(slot.clone(), state.clone());
2010 }
2011 AgentMessage::Display(_) | AgentMessage::Error(_) => {}
2012 }
2013 }
2014 let state = match &message {
2015 AgentMessage::State(_) if state_delta => {
2016 patch_against_last(&mut delta_bases, player_index, per_seat, &slot, &state)
2017 .unwrap_or(state)
2018 }
2019 _ => state,
2020 };
2021 let target_player = if per_seat
2022 && matches!(
2023 message,
2024 AgentMessage::State(_) | AgentMessage::Prompt(_) | AgentMessage::Error(_)
2025 ) {
2026 let Some(target_player) = seat_usernames
2027 .as_ref()
2028 .and_then(|names| names.get(player_index).cloned())
2029 else {
2030 warn!(player_index, PRIVATE_MESSAGE_MISSING_PLAYER);
2031 continue;
2032 };
2033 Some(target_player)
2034 } else {
2035 None
2036 };
2037 if outbound_tx
2038 .send(ClientMessage::BroadcastState {
2039 state,
2040 target_player,
2041 })
2042 .is_err()
2043 {
2044 break;
2045 }
2046 }
2047 });
2048}
2049
2050fn spawn_game_over_forwarder(
2051 outbound_tx: tokio_mpsc::UnboundedSender<ClientMessage>,
2052 game_over_rx: std_mpsc::Receiver<HostedGameOver>,
2053 session_handle: SharedEngineSession,
2054 snapshot: SharedHostSnapshot,
2055 game_id: String,
2056 seat_usernames: Option<Vec<String>>,
2057) {
2058 thread::spawn(move || {
2059 while let Ok(game_over) = game_over_rx.recv() {
2060 let Ok(session) = session_handle.lock() else {
2061 return;
2062 };
2063 if session
2064 .as_ref()
2065 .is_none_or(|session| session.game_id() != game_id)
2066 {
2067 warn!(
2068 game_id,
2069 "stale engine session reached game over; not ending the relay game"
2070 );
2071 continue;
2072 }
2073 if let Ok(mut snap) = snapshot.lock() {
2074 if snap
2075 .game
2076 .as_ref()
2077 .is_some_and(|game| game.game_id == game_id)
2078 {
2079 snap.pending_end_game = Some(game_id.clone());
2080 }
2081 }
2082 let mut last_state_by_seat: HashMap<usize, Value> = HashMap::new();
2083 let mut last_state: Option<Value> = None;
2084 for (player_index, message) in game_over.messages {
2085 let per_seat = seat_usernames.is_some() && player_index != OBSERVER_SEAT;
2086 let envelope = match &message {
2087 AgentMessage::State(state_update) if !per_seat => StateEnvelope::State {
2088 for_player: None,
2089 state: serde_json::to_value(state_update).unwrap_or(Value::Null),
2090 fingerprint: None,
2091 },
2092 _ => StateEnvelope::for_agent_message(player_slot(player_index), &message),
2093 };
2094 let Ok(state) = serde_json::to_value(envelope) else {
2095 continue;
2096 };
2097 match &message {
2098 AgentMessage::State(_) if per_seat => {
2099 if last_state_by_seat.get(&player_index) == Some(&state) {
2100 continue;
2101 }
2102 last_state_by_seat.insert(player_index, state.clone());
2103 }
2104 AgentMessage::State(_) if last_state.as_ref() == Some(&state) => continue,
2105 AgentMessage::State(_) => last_state = Some(state.clone()),
2106 AgentMessage::Display(_) | AgentMessage::Prompt(_) | AgentMessage::Error(_) => {
2107 }
2108 }
2109 let target_player = if per_seat
2110 && matches!(
2111 message,
2112 AgentMessage::State(_) | AgentMessage::Prompt(_) | AgentMessage::Error(_)
2113 ) {
2114 let Some(target_player) = seat_usernames
2115 .as_ref()
2116 .and_then(|names| names.get(player_index).cloned())
2117 else {
2118 warn!(player_index, PRIVATE_MESSAGE_MISSING_PLAYER);
2119 continue;
2120 };
2121 Some(target_player)
2122 } else {
2123 None
2124 };
2125 if outbound_tx
2126 .send(ClientMessage::BroadcastState {
2127 state,
2128 target_player,
2129 })
2130 .is_err()
2131 {
2132 return;
2133 }
2134 }
2135 let Ok(state) = serde_json::to_value(StateEnvelope::RoomRelay {
2136 protocol: SELF_HOSTED_NODE_PROTOCOL.to_string(),
2137 version: 1,
2138 message_id: Uuid::new_v4().to_string(),
2139 from_player: None,
2140 target_player: None,
2141 room_id: None,
2142 payload: json!({ "type": "gameOver", "gameId": game_over.game_id }),
2143 }) else {
2144 continue;
2145 };
2146 if outbound_tx
2147 .send(ClientMessage::BroadcastState {
2148 state,
2149 target_player: None,
2150 })
2151 .is_err()
2152 {
2153 return;
2154 }
2155 if outbound_tx
2156 .send(ClientMessage::EndGame {
2157 game_id: game_id.clone(),
2158 })
2159 .is_err()
2160 {
2161 return;
2162 }
2163 }
2164 });
2165}
2166
2167async fn relay_writer(mut write: WsWrite, mut outbound: tokio_mpsc::UnboundedReceiver<Message>) {
2168 while let Some(message) = outbound.recv().await {
2169 let started = Instant::now();
2170 let closing = matches!(message, Message::Close(_));
2171 if let Err(error) = write.send(message).await {
2172 warn!(%error, "relay write failed");
2173 break;
2174 }
2175 crate::metrics::record_relay_send(started.elapsed());
2176 if closing {
2177 break;
2178 }
2179 }
2180 let _ = write.close().await;
2181}
2182
2183impl RelayClient {
2184 async fn connect(
2185 relay_url: &str,
2186 username: &str,
2187 password: &str,
2188 ) -> Result<Self, Box<dyn std::error::Error + Send + Sync>> {
2189 info!(relay_url, username, "connecting relay client");
2190 let (socket, _) = connect_async_with_config(relay_url, None, true).await?;
2191 let (write, read) = socket.split();
2192 let (outbound, outbound_rx) = tokio_mpsc::unbounded_channel();
2193 let writer = tokio::spawn(relay_writer(write, outbound_rx));
2194 let mut client = Self {
2195 username: username.to_string(),
2196 outbound,
2197 read,
2198 writer,
2199 };
2200 client
2201 .send(&ClientMessage::Authenticate {
2202 username: username.to_string(),
2203 password: password.to_string(),
2204 service: true,
2205 identity: None,
2206 client_platform: ClientPlatform::Unknown,
2207 })
2208 .await?;
2209 client.wait_for_auth().await?;
2210 Ok(client)
2211 }
2212
2213 async fn wait_for_auth(&mut self) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
2214 loop {
2215 match self.recv().await? {
2216 Some(ServerMessage::AuthResult { success, error, .. }) => {
2217 if success {
2218 info!(username = %self.username, "authenticated");
2219 return Ok(());
2220 }
2221 return Err(format!("authentication failed: {error:?}").into());
2222 }
2223 Some(other) => {
2224 debug!(?other, username = %self.username, "ignored pre-auth message")
2225 }
2226 None => return Err("relay closed before authentication".into()),
2227 }
2228 }
2229 }
2230
2231 async fn send(
2232 &mut self,
2233 message: &ClientMessage,
2234 ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
2235 self.enqueue(Message::Text(serde_json::to_string(message)?))
2236 }
2237
2238 fn enqueue(&self, message: Message) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
2239 self.outbound
2240 .send(message)
2241 .map_err(|_| "relay writer stopped".into())
2242 }
2243
2244 async fn recv(
2245 &mut self,
2246 ) -> Result<Option<ServerMessage>, Box<dyn std::error::Error + Send + Sync>> {
2247 loop {
2248 let Some(message) = self.read.next().await else {
2249 return Ok(None);
2250 };
2251 let message = message?;
2252 match message {
2253 Message::Text(text) => {
2254 return Ok(Some(serde_json::from_str(&text)?));
2255 }
2256 Message::Ping(payload) => {
2257 self.enqueue(Message::Pong(payload))?;
2258 }
2259 Message::Close(_) => return Ok(None),
2260 _ => {}
2261 }
2262 }
2263 }
2264
2265 async fn close(&mut self) {
2266 if self.enqueue(Message::Close(None)).is_err() {
2267 return;
2268 }
2269 let _ = time::timeout(CLOSE_DRAIN_TIMEOUT, async {
2270 while let Some(Ok(_)) = self.read.next().await {}
2271 })
2272 .await;
2273 self.writer.abort();
2274 }
2275
2276 async fn broadcast_room_message(
2277 &mut self,
2278 protocol: &str,
2279 payload: Value,
2280 ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
2281 let envelope = StateEnvelope::RoomRelay {
2282 protocol: protocol.to_string(),
2283 version: 1,
2284 message_id: Uuid::new_v4().to_string(),
2285 from_player: Some(self.username.clone()),
2286 target_player: None,
2287 room_id: None,
2288 payload,
2289 };
2290 self.send(&ClientMessage::BroadcastState {
2291 state: serde_json::to_value(envelope)?,
2292 target_player: None,
2293 })
2294 .await
2295 }
2296}
2297
2298fn log_room_update(observer: &str, room: &RoomInfo) {
2299 let players = room
2300 .players
2301 .iter()
2302 .map(|player| {
2303 format!(
2304 "{}{}{}",
2305 player.username,
2306 if player.ready { ":ready" } else { "" },
2307 if player.connected { "" } else { ":offline" },
2308 )
2309 })
2310 .collect::<Vec<_>>()
2311 .join(",");
2312 info!(
2313 observer,
2314 room_id = %room.room_id,
2315 room_name = %room.room_name,
2316 host = %room.host,
2317 players,
2318 "room update"
2319 );
2320}