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