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