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