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