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