Skip to main content

self_hosted_node/
host.rs

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