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