Skip to main content

self_hosted_node/
host.rs

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