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