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