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