Skip to main content

self_hosted_node/engine_backend/
java_backend.rs

1#![allow(dead_code)]
2
3#[cfg(forge_backend)]
4use std::collections::HashMap;
5use std::env;
6#[cfg(feature = "java-forge")]
7use std::io::{BufRead, BufReader, BufWriter, Write};
8use std::path::{Path, PathBuf};
9#[cfg(feature = "java-forge")]
10use std::process::{Child, ChildStdin, Command, Stdio};
11use std::sync::atomic::AtomicBool;
12use std::sync::mpsc as std_mpsc;
13#[cfg(feature = "java-forge")]
14use std::sync::mpsc::RecvTimeoutError;
15#[cfg(forge_backend)]
16use std::sync::mpsc::TryRecvError;
17use std::sync::Arc;
18#[cfg(feature = "java-forge")]
19use std::sync::Mutex;
20#[cfg(forge_backend)]
21use std::time::Duration;
22#[cfg(feature = "java-forge")]
23use std::time::Instant;
24
25use manabrew_protocol::deck_dto::{Deck, DeckCardIdentity};
26
27use crate::config::DeckSelection;
28#[cfg(feature = "java-forge")]
29use manabot::{BotAgent, SimpleAi};
30#[cfg(forge_backend)]
31use manabrew_agent_interface::game_view_dto::GameViewDto;
32use manabrew_agent_interface::prompt::{AgentMessage, ClientToServerMessage};
33#[cfg(forge_backend)]
34use manabrew_agent_interface::prompt::{
35    AgentPrompt, ChooseActionOutput, DiceRolledOutput, DirectiveInput, GameOverInput, PromptInput,
36    PromptOutput, ProtocolError, ProtocolErrorCode, ResponseViolation, StateUpdate,
37};
38#[cfg(feature = "java-forge")]
39use manabrew_agent_interface::prompt::{MulliganOutput, MulliganPutBackOutput};
40use serde::Serialize;
41#[cfg(feature = "java-forge")]
42use serde_json::json;
43#[cfg(feature = "java-forge")]
44use serde_json::Value;
45#[cfg(forge_backend)]
46use tracing::warn;
47#[cfg(forge_backend)]
48use tracing::{debug, info};
49
50use super::HostedGameOver;
51use crate::config::workspace_root;
52
53pub fn unsupported_message() -> &'static str {
54    "hosted java-forge backend is unavailable; rebuild self-hosted-node with --features java-forge"
55}
56
57#[cfg(feature = "java-forge")]
58pub fn run_smoke_game(max_prompts: usize) -> Result<(), String> {
59    let config = JavaRuntimeConfig::from_env();
60    let assets_dir = config.assets_dir.to_string_lossy().to_string();
61    let bridge = SubprocessBridge::spawn(&config)?;
62    let mut session = JavaForgeSession::new(bridge);
63    session.initialize(&assets_dir)?;
64
65    let deck_a = smoke_deck("Mountain", "Lightning Bolt");
66    let deck_b = smoke_deck("Forest", "Grizzly Bears");
67    let request = StartGameRequest::new(
68        "self-hosted-java-smoke".to_string(),
69        String::new(),
70        20,
71        42,
72        vec![
73            PlayerConfig::new("Smoke A".to_string(), &deck_a, Vec::new()),
74            PlayerConfig::new("Smoke B".to_string(), &deck_b, Vec::new()),
75        ],
76    );
77    let session_id = session.start_game(&request)?;
78    info!(session_id, "java-forge smoke session started");
79
80    let mut prompts_seen = 0usize;
81    while prompts_seen < max_prompts {
82        let Some(prompt_json) = wait_for_prompt(&mut session, 600)? else {
83            session.end_game()?;
84            return Err("timed out waiting for java-forge smoke prompt".to_string());
85        };
86        let prompt: AgentPrompt = serde_json::from_str(&prompt_json)
87            .map_err(|err| format!("failed to parse java-forge smoke prompt: {err}"))?;
88        let player = player_index(&prompt.deciding_player_id);
89        info!(prompts_seen, player, "java-forge smoke prompt");
90        let pass = PromptOutput::ChooseAction(ChooseActionOutput::Pass {
91            until: None,
92            exhaust_stack: false,
93        });
94        session.submit_action(&serde_json::to_string(&pass).map_err(|err| err.to_string())?)?;
95        prompts_seen += 1;
96    }
97
98    let snapshot_json = session.get_snapshot(Some(0))?;
99    let snapshot: Value = serde_json::from_str(&snapshot_json)
100        .map_err(|err| format!("failed to parse java-forge smoke snapshot: {err}"))?;
101    info!(
102        turn = snapshot
103            .get("turn")
104            .and_then(|value| value.as_i64())
105            .unwrap_or_default(),
106        phase = snapshot
107            .get("phase")
108            .and_then(|value| value.as_str())
109            .unwrap_or("<missing>"),
110        "java-forge smoke snapshot"
111    );
112    session.end_game()?;
113    Ok(())
114}
115
116#[cfg(not(feature = "java-forge"))]
117pub fn run_smoke_game(_max_prompts: usize) -> Result<(), String> {
118    Err(
119        "java-forge smoke requires building self-hosted-node with --features java-forge"
120            .to_string(),
121    )
122}
123
124// Unlike run_smoke_game, which spawns the JAR in a subprocess, this drives the
125// GraalVM native library in-process: isolate creation, forge_initialize (which
126// loads the card database) and a real game start. CI runs it after restoring
127// forge-harness/native/build from cache, where nothing else would notice a
128// stale or broken libforgeharness.
129#[cfg(feature = "graal-forge")]
130pub fn run_graal_smoke() -> Result<(), String> {
131    let config = JavaRuntimeConfig::from_env();
132    let engine = GraalEngineHandle::create(&config.assets_dir)?;
133
134    let deck_a = smoke_deck("Mountain", "Lightning Bolt");
135    let deck_b = smoke_deck("Forest", "Grizzly Bears");
136    let request = StartGameRequest::new(
137        "self-hosted-graal-smoke".to_string(),
138        String::new(),
139        20,
140        42,
141        vec![
142            PlayerConfig::new("Smoke A".to_string(), &deck_a, Vec::new()),
143            PlayerConfig::new("Smoke B".to_string(), &deck_b, Vec::new()),
144        ],
145    );
146
147    let session_id = engine.start_game(&request.to_json().map_err(|err| err.to_string())?)?;
148    info!(session_id, "graal-forge smoke session started");
149    engine.end_game(&session_id)?;
150    Ok(())
151}
152
153#[cfg(not(feature = "graal-forge"))]
154pub fn run_graal_smoke() -> Result<(), String> {
155    Err(
156        "graal-forge smoke requires building self-hosted-node with --features graal-forge"
157            .to_string(),
158    )
159}
160
161#[cfg(feature = "java-forge")]
162pub fn run_scenario(name: &str, max_prompts: usize) -> Result<(), String> {
163    let scenario = JavaScenario::from_name(name)?;
164    let config = JavaRuntimeConfig::from_env();
165    let assets_dir = config.assets_dir.to_string_lossy().to_string();
166    let bridge = SubprocessBridge::spawn(&config)?;
167    let mut session = JavaForgeSession::new(bridge);
168    session.initialize(&assets_dir)?;
169
170    let request = StartGameRequest::new(
171        format!("self-hosted-java-scenario-{}", scenario.name()),
172        String::new(),
173        20,
174        42,
175        vec![
176            PlayerConfig::new(
177                "Scenario A".to_string(),
178                &scenario_deck("Swamp"),
179                Vec::new(),
180            ),
181            PlayerConfig::new(
182                "Scenario B".to_string(),
183                &scenario_deck("Forest"),
184                Vec::new(),
185            ),
186        ],
187    );
188    let session_id = session.start_game(&request)?;
189    info!(
190        session_id,
191        scenario = scenario.name(),
192        "java-forge scenario started"
193    );
194
195    let result = run_scenario_loop(&mut session, scenario, max_prompts);
196    let end_result = session.end_game();
197    result.and(end_result)
198}
199
200#[cfg(not(feature = "java-forge"))]
201pub fn run_scenario(_name: &str, _max_prompts: usize) -> Result<(), String> {
202    Err(
203        "java-forge scenarios require building self-hosted-node with --features java-forge"
204            .to_string(),
205    )
206}
207
208#[cfg(feature = "java-forge")]
209pub fn run_self_play(
210    seats: &[DeckSelection],
211    starting_life: i32,
212    seed: u64,
213    max_prompts: usize,
214    games: usize,
215) -> Result<(), String> {
216    let config = JavaRuntimeConfig::from_env();
217    let assets_dir = config.assets_dir.to_string_lossy().to_string();
218    let bridge = SubprocessBridge::spawn(&config)?;
219    let mut session = JavaForgeSession::new(bridge);
220    session.initialize(&assets_dir)?;
221
222    let mut players = Vec::with_capacity(seats.len());
223    for (i, seat) in seats.iter().enumerate() {
224        let identities = deck_card_identities(&seat.deck);
225        players.push(PlayerConfig::new(
226            format!("Self-Play {}", i + 1),
227            &identities,
228            commander_names_for_java(&seat.deck, seat.commander_name.as_deref()),
229        ));
230    }
231
232    for game_index in 0..games.max(1) {
233        let request = StartGameRequest::new(
234            format!("self-hosted-java-self-play-{game_index}"),
235            String::new(),
236            starting_life,
237            seed.wrapping_add(game_index as u64),
238            players.clone(),
239        );
240        let session_id = session.start_game(&request)?;
241        info!(
242            session_id,
243            game_index,
244            games,
245            players = seats.len(),
246            starting_life,
247            max_prompts,
248            "java-forge self-play game started"
249        );
250        let result = run_self_play_loop(&mut session, max_prompts);
251        let end_result = session.end_game();
252        result.and(end_result)?;
253    }
254    Ok(())
255}
256
257#[cfg(not(feature = "java-forge"))]
258pub fn run_self_play(
259    _seats: &[DeckSelection],
260    _starting_life: i32,
261    _seed: u64,
262    _max_prompts: usize,
263    _games: usize,
264) -> Result<(), String> {
265    Err(
266        "java-forge self-play requires building self-hosted-node with --features java-forge"
267            .to_string(),
268    )
269}
270
271#[cfg(feature = "java-forge")]
272type SharedBridge = Arc<Mutex<SubprocessBridge>>;
273
274#[cfg(feature = "java-forge")]
275struct PoolSlot {
276    bridge: SharedBridge,
277    active: usize,
278}
279
280#[cfg(feature = "java-forge")]
281pub struct JavaEnginePool {
282    config: JavaRuntimeConfig,
283    max_sessions: usize,
284    sessions_per_process: usize,
285    slots: Mutex<Vec<PoolSlot>>,
286    in_use: Mutex<HashMap<String, SharedBridge>>,
287}
288
289#[cfg(feature = "java-forge")]
290#[derive(Clone)]
291pub struct JavaEngineHandle {
292    pool: Arc<JavaEnginePool>,
293}
294
295#[cfg(feature = "java-forge")]
296impl JavaEnginePool {
297    pub fn start(
298        config: &JavaRuntimeConfig,
299        max_sessions: usize,
300        sessions_per_process: usize,
301    ) -> Result<Arc<Self>, String> {
302        let max_sessions = max_sessions.max(1);
303        let sessions_per_process = sessions_per_process.max(1);
304        let processes = max_sessions.div_ceil(sessions_per_process);
305        let mut slots = Vec::with_capacity(processes);
306        for slot in 0..processes {
307            info!(
308                slot,
309                processes, sessions_per_process, "pre-warming java subprocess"
310            );
311            let bridge = SubprocessBridge::spawn(config)?;
312            slots.push(PoolSlot {
313                bridge: Arc::new(Mutex::new(bridge)),
314                active: 0,
315            });
316        }
317        Ok(Arc::new(Self {
318            config: config.clone(),
319            max_sessions,
320            sessions_per_process,
321            slots: Mutex::new(slots),
322            in_use: Mutex::new(HashMap::new()),
323        }))
324    }
325
326    pub fn handle(self: &Arc<Self>) -> JavaEngineHandle {
327        JavaEngineHandle {
328            pool: Arc::clone(self),
329        }
330    }
331}
332
333#[cfg(feature = "java-forge")]
334impl Drop for JavaEnginePool {
335    fn drop(&mut self) {
336        let slots = self.slots.get_mut().map(std::mem::take).unwrap_or_default();
337        for slot in slots {
338            if let Ok(mutex) = Arc::try_unwrap(slot.bridge) {
339                if let Ok(inner) = mutex.into_inner() {
340                    inner.shutdown();
341                }
342            }
343        }
344    }
345}
346
347#[cfg(feature = "java-forge")]
348impl JavaEnginePool {
349    fn acquire(&self) -> Result<SharedBridge, String> {
350        let deadline = Instant::now() + Duration::from_secs(60);
351        loop {
352            let claimed = {
353                let mut slots = self
354                    .slots
355                    .lock()
356                    .map_err(|_| "java engine slots poisoned".to_string())?;
357                let mut claimed = None;
358                for slot in slots.iter_mut() {
359                    if slot.active < self.sessions_per_process {
360                        slot.active += 1;
361                        claimed = Some(Arc::clone(&slot.bridge));
362                        break;
363                    }
364                }
365                claimed
366            };
367            if let Some(bridge) = claimed {
368                let alive = bridge
369                    .lock()
370                    .ok()
371                    .map(|mut guard| guard.is_alive())
372                    .unwrap_or(false);
373                if alive {
374                    return Ok(bridge);
375                }
376                warn!("discarding dead java subprocess from pool");
377                self.replace_slot(&bridge);
378                continue;
379            }
380            if Instant::now() >= deadline {
381                return Err(format!(
382                    "java engine pool exhausted (max_sessions={}); no free session slot after 60s",
383                    self.max_sessions
384                ));
385            }
386            std::thread::sleep(Duration::from_millis(50));
387        }
388    }
389
390    fn release(&self, bridge: SharedBridge) {
391        let now_idle = {
392            let mut slots = match self.slots.lock() {
393                Ok(slots) => slots,
394                Err(_) => return,
395            };
396            match slots.iter_mut().find(|s| Arc::ptr_eq(&s.bridge, &bridge)) {
397                Some(slot) => {
398                    slot.active = slot.active.saturating_sub(1);
399                    slot.active == 0
400                }
401                None => return,
402            }
403        };
404        if !now_idle {
405            return;
406        }
407        let healthy = {
408            let mut guard = match bridge.lock() {
409                Ok(guard) => guard,
410                Err(_) => return,
411            };
412            guard.is_alive() && guard.reset().is_ok()
413        };
414        if !healthy {
415            warn!("java subprocess unhealthy at idle; respawning");
416            self.replace_slot(&bridge);
417        }
418    }
419
420    fn replace_slot(&self, dead: &SharedBridge) {
421        let Ok(mut slots) = self.slots.lock() else {
422            return;
423        };
424        let Some(index) = slots.iter().position(|s| Arc::ptr_eq(&s.bridge, dead)) else {
425            return;
426        };
427        match SubprocessBridge::spawn(&self.config) {
428            Ok(replacement) => {
429                slots[index] = PoolSlot {
430                    bridge: Arc::new(Mutex::new(replacement)),
431                    active: 0,
432                };
433            }
434            Err(error) => {
435                warn!(%error, "failed to respawn java subprocess; retiring pool slot");
436                slots.remove(index);
437            }
438        }
439    }
440}
441
442#[cfg(feature = "java-forge")]
443impl JavaEngineHandle {
444    fn bridge_for(&self, session_id: &str) -> Result<SharedBridge, String> {
445        let in_use = self
446            .pool
447            .in_use
448            .lock()
449            .map_err(|_| "java engine in_use map poisoned".to_string())?;
450        in_use
451            .get(session_id)
452            .cloned()
453            .ok_or_else(|| format!("unknown java session: {session_id}"))
454    }
455
456    pub fn start_game(&self, request_json: &str) -> Result<String, String> {
457        let bridge = self.pool.acquire()?;
458        let response = {
459            let mut guard = bridge
460                .lock()
461                .map_err(|_| "java subprocess mutex poisoned".to_string())?;
462            guard.start_game_json(request_json)
463        };
464        let response = match response {
465            Ok(response) => response,
466            Err(error) => {
467                self.pool.release(bridge);
468                return Err(error);
469            }
470        };
471        let parsed: StartGameResponse = match serde_json::from_str(&response) {
472            Ok(parsed) => parsed,
473            Err(error) => {
474                self.pool.release(bridge);
475                return Err(format!("malformed startGame response: {error}"));
476            }
477        };
478        let session_id = parsed.session_id.clone();
479        let displaced = {
480            let mut in_use = self
481                .pool
482                .in_use
483                .lock()
484                .map_err(|_| "java engine in_use map poisoned".to_string())?;
485            in_use.insert(session_id.clone(), bridge)
486        };
487        if let Some(displaced) = displaced {
488            warn!(
489                session_id,
490                "session_id collision; releasing displaced java subprocess"
491            );
492            self.pool.release(displaced);
493        }
494        Ok(session_id)
495    }
496
497    pub fn submit_action(&self, session_id: &str, action_json: &str) -> Result<String, String> {
498        let bridge = self.bridge_for(session_id)?;
499        let mutex_started = Instant::now();
500        let mut guard = bridge
501            .lock()
502            .map_err(|_| "java subprocess mutex poisoned".to_string())?;
503        crate::metrics::record_forge_decision_stage("bridge_mutex", mutex_started.elapsed());
504        let call_started = Instant::now();
505        let result = guard.submit_action(session_id, action_json);
506        crate::metrics::record_forge_decision_stage("submit_action", call_started.elapsed());
507        result
508    }
509
510    pub fn get_prompt(
511        &self,
512        session_id: &str,
513        player_index: usize,
514    ) -> Result<Option<String>, String> {
515        let bridge = self.bridge_for(session_id)?;
516        let mut guard = bridge
517            .lock()
518            .map_err(|_| "java subprocess mutex poisoned".to_string())?;
519        guard.get_prompt(session_id, player_index)
520    }
521
522    pub fn is_game_over(&self, session_id: &str) -> Result<bool, String> {
523        let bridge = self.bridge_for(session_id)?;
524        let mut guard = bridge
525            .lock()
526            .map_err(|_| "java subprocess mutex poisoned".to_string())?;
527        guard.is_game_over(session_id)
528    }
529
530    pub fn get_snapshot(&self, session_id: &str, viewer: Option<usize>) -> Result<String, String> {
531        let bridge = self.bridge_for(session_id)?;
532        let mut guard = bridge
533            .lock()
534            .map_err(|_| "java subprocess mutex poisoned".to_string())?;
535        guard.get_snapshot(session_id, viewer)
536    }
537
538    pub fn end_game(&self, session_id: &str) -> Result<(), String> {
539        let bridge = {
540            let mut in_use = self
541                .pool
542                .in_use
543                .lock()
544                .map_err(|_| "java engine in_use map poisoned".to_string())?;
545            in_use.remove(session_id)
546        };
547        let Some(bridge) = bridge else {
548            return Ok(());
549        };
550        let result = {
551            let mut guard = bridge
552                .lock()
553                .map_err(|_| "java subprocess mutex poisoned".to_string())?;
554            guard.end_game(session_id)
555        };
556        self.pool.release(bridge);
557        result
558    }
559
560    pub fn abort_game(&self, session_id: &str) -> Result<(), String> {
561        let bridge = {
562            let mut in_use = self
563                .pool
564                .in_use
565                .lock()
566                .map_err(|_| "java engine in_use map poisoned".to_string())?;
567            in_use.remove(session_id)
568        };
569        let Some(bridge) = bridge else {
570            return Ok(());
571        };
572        let result = {
573            let mut guard = bridge
574                .lock()
575                .map_err(|_| "java subprocess mutex poisoned".to_string())?;
576            guard.abort_game(session_id)
577        };
578        self.pool.release(bridge);
579        result
580    }
581}
582
583#[cfg(feature = "java-forge")]
584static JAVA_ENGINE: std::sync::OnceLock<Arc<JavaEnginePool>> = std::sync::OnceLock::new();
585
586#[cfg(feature = "java-forge")]
587pub fn init_engine() -> Result<(), String> {
588    if JAVA_ENGINE.get().is_some() {
589        return Ok(());
590    }
591    let config = JavaRuntimeConfig::from_env();
592    // SELF_HOSTED_NODE_MAX_GAMES is the concurrent-session ceiling for this node.
593    // SELF_HOSTED_NODE_GAMES_PER_JVM multiplexes that many sessions into each
594    // subprocess (the engine is concurrency-safe since the endstep patches), so
595    // the pool spawns ceil(max_games / games_per_jvm) processes. Default 1 keeps
596    // the historical one-subprocess-per-game shape.
597    let max_sessions = env::var("SELF_HOSTED_NODE_MAX_GAMES")
598        .ok()
599        .and_then(|value| value.parse::<usize>().ok())
600        .filter(|n| *n >= 1)
601        .unwrap_or(1);
602    let sessions_per_process = env::var("SELF_HOSTED_NODE_GAMES_PER_JVM")
603        .ok()
604        .and_then(|value| value.parse::<usize>().ok())
605        .filter(|n| *n >= 1)
606        .unwrap_or(1);
607    let pool = JavaEnginePool::start(&config, max_sessions, sessions_per_process)?;
608    JAVA_ENGINE
609        .set(pool)
610        .map_err(|_| "java engine already initialized".to_string())
611}
612
613#[cfg(all(feature = "graal-forge", not(feature = "java-forge")))]
614pub fn init_engine() -> Result<(), String> {
615    // Loading the card database into the isolate takes tens of seconds on a
616    // small box. Pay it here, before any room is advertised, rather than
617    // leaving it for whoever starts the first game.
618    if shared_isolate_enabled() {
619        drop(GraalEngineHandle::create(
620            &JavaRuntimeConfig::from_env().assets_dir,
621        )?);
622    }
623    Ok(())
624}
625
626#[cfg(not(forge_backend))]
627pub fn init_engine() -> Result<(), String> {
628    Err(
629        "forge engine requires building self-hosted-node with --features java-forge or graal-forge"
630            .to_string(),
631    )
632}
633
634#[cfg(feature = "java-forge")]
635fn engine_handle() -> Result<JavaEngineHandle, String> {
636    JAVA_ENGINE
637        .get()
638        .map(JavaEnginePool::handle)
639        .ok_or_else(|| "java engine is not initialized".to_string())
640}
641
642#[cfg(feature = "java-forge")]
643type ForgeEngine = JavaEngineHandle;
644
645#[cfg(feature = "java-forge")]
646fn obtain_engine() -> Result<ForgeEngine, String> {
647    engine_handle()
648}
649
650#[cfg(all(feature = "graal-forge", not(feature = "java-forge")))]
651type ForgeEngine = GraalEngineHandle;
652
653#[cfg(all(feature = "graal-forge", not(feature = "java-forge")))]
654fn obtain_engine() -> Result<ForgeEngine, String> {
655    GraalEngineHandle::create(&JavaRuntimeConfig::from_env().assets_dir)
656}
657
658#[cfg(feature = "graal-forge")]
659mod graal_ffi {
660    use std::os::raw::{c_char, c_int};
661
662    #[allow(non_camel_case_types)]
663    pub type graal_isolate_t = std::ffi::c_void;
664    #[allow(non_camel_case_types)]
665    pub type graal_isolatethread_t = std::ffi::c_void;
666
667    extern "C" {
668        pub fn graal_create_isolate(
669            params: *mut std::ffi::c_void,
670            isolate: *mut *mut graal_isolate_t,
671            thread: *mut *mut graal_isolatethread_t,
672        ) -> c_int;
673        pub fn graal_tear_down_isolate(thread: *mut graal_isolatethread_t) -> c_int;
674        pub fn graal_attach_thread(
675            isolate: *mut graal_isolate_t,
676            thread: *mut *mut graal_isolatethread_t,
677        ) -> c_int;
678        pub fn graal_detach_thread(thread: *mut graal_isolatethread_t) -> c_int;
679        pub fn forge_initialize(
680            thread: *mut graal_isolatethread_t,
681            assets_dir: *const c_char,
682        ) -> *mut c_char;
683        pub fn forge_start_game(
684            thread: *mut graal_isolatethread_t,
685            request_json: *const c_char,
686        ) -> *mut c_char;
687        pub fn forge_submit_action(
688            thread: *mut graal_isolatethread_t,
689            session_id: *const c_char,
690            action_json: *const c_char,
691        ) -> *mut c_char;
692        pub fn forge_get_prompt(
693            thread: *mut graal_isolatethread_t,
694            session_id: *const c_char,
695            player_index: c_int,
696        ) -> *mut c_char;
697        pub fn forge_get_snapshot(
698            thread: *mut graal_isolatethread_t,
699            session_id: *const c_char,
700            viewer: c_int,
701        ) -> *mut c_char;
702        pub fn forge_get_game_over(
703            thread: *mut graal_isolatethread_t,
704            session_id: *const c_char,
705        ) -> *mut c_char;
706        pub fn forge_end_game(
707            thread: *mut graal_isolatethread_t,
708            session_id: *const c_char,
709        ) -> *mut c_char;
710        pub fn forge_abort_game(
711            thread: *mut graal_isolatethread_t,
712            session_id: *const c_char,
713        ) -> *mut c_char;
714        pub fn forge_free_string(thread: *mut graal_isolatethread_t, ptr: *mut c_char);
715    }
716}
717
718#[cfg(feature = "graal-forge")]
719#[derive(serde::Deserialize)]
720struct ForgeReply {
721    ok: bool,
722    #[serde(default)]
723    result: String,
724    #[serde(default)]
725    error: Option<String>,
726}
727
728// A GraalVM isolate hosts the in-process Forge engine. The `thread` handle is
729// bound to the hosted-engine thread that created or attached it (isolate
730// threads are not portable), so the handle is Rc/!Send by design. In shared
731// mode one isolate serves all rooms: the first creator initializes Forge, later
732// rooms attach their own thread to it and sessions coexist in the adapter.
733#[cfg(feature = "graal-forge")]
734struct GraalBridge {
735    thread: *mut graal_ffi::graal_isolatethread_t,
736    attached: bool,
737}
738
739#[cfg(feature = "graal-forge")]
740struct SharedIsolate(*mut graal_ffi::graal_isolate_t);
741#[cfg(feature = "graal-forge")]
742unsafe impl Send for SharedIsolate {}
743
744#[cfg(feature = "graal-forge")]
745static SHARED_GRAAL_ISOLATE: std::sync::Mutex<Option<SharedIsolate>> = std::sync::Mutex::new(None);
746
747// One isolate hosts every room's game (safe since the endstep concurrency
748// patches), so the card database loads once per node. Off keeps the historical
749// isolate-per-game shape.
750#[cfg(feature = "graal-forge")]
751fn shared_isolate_enabled() -> bool {
752    env::var("SELF_HOSTED_NODE_SHARED_ISOLATE")
753        .map(|value| value == "1" || value.eq_ignore_ascii_case("true"))
754        .unwrap_or(false)
755}
756
757#[cfg(feature = "graal-forge")]
758impl GraalBridge {
759    fn create() -> Result<Self, String> {
760        let mut isolate: *mut graal_ffi::graal_isolate_t = std::ptr::null_mut();
761        let mut thread: *mut graal_ffi::graal_isolatethread_t = std::ptr::null_mut();
762        let rc = unsafe {
763            graal_ffi::graal_create_isolate(std::ptr::null_mut(), &mut isolate, &mut thread)
764        };
765        if rc != 0 {
766            return Err(format!("graal_create_isolate failed with code {rc}"));
767        }
768        Ok(Self {
769            thread,
770            attached: false,
771        })
772    }
773
774    fn create_in_shared_isolate(assets_dir: &Path) -> Result<Self, String> {
775        let mut guard = SHARED_GRAAL_ISOLATE
776            .lock()
777            .map_err(|_| "shared graal isolate poisoned".to_string())?;
778        if let Some(shared) = guard.as_ref() {
779            let mut thread: *mut graal_ffi::graal_isolatethread_t = std::ptr::null_mut();
780            let rc = unsafe { graal_ffi::graal_attach_thread(shared.0, &mut thread) };
781            if rc != 0 {
782                return Err(format!("graal_attach_thread failed with code {rc}"));
783            }
784            return Ok(Self {
785                thread,
786                attached: true,
787            });
788        }
789        let mut isolate: *mut graal_ffi::graal_isolate_t = std::ptr::null_mut();
790        let mut thread: *mut graal_ffi::graal_isolatethread_t = std::ptr::null_mut();
791        let rc = unsafe {
792            graal_ffi::graal_create_isolate(std::ptr::null_mut(), &mut isolate, &mut thread)
793        };
794        if rc != 0 {
795            return Err(format!("graal_create_isolate failed with code {rc}"));
796        }
797        let mut bridge = Self {
798            thread,
799            attached: false,
800        };
801        let assets = cstring(&assets_dir.to_string_lossy())?;
802        bridge.decode(unsafe { graal_ffi::forge_initialize(bridge.thread, assets.as_ptr()) })?;
803        bridge.attached = true;
804        *guard = Some(SharedIsolate(isolate));
805        info!("shared graal isolate initialized");
806        Ok(bridge)
807    }
808
809    fn decode(&self, raw: *mut std::os::raw::c_char) -> Result<String, String> {
810        if raw.is_null() {
811            return Err("forge native lib returned null".to_string());
812        }
813        let envelope = unsafe { std::ffi::CStr::from_ptr(raw) }
814            .to_string_lossy()
815            .into_owned();
816        unsafe { graal_ffi::forge_free_string(self.thread, raw) };
817        let reply: ForgeReply = serde_json::from_str(&envelope)
818            .map_err(|err| format!("malformed forge envelope: {err}"))?;
819        if reply.ok {
820            Ok(reply.result)
821        } else {
822            Err(reply
823                .error
824                .unwrap_or_else(|| "unknown forge error".to_string()))
825        }
826    }
827}
828
829#[cfg(feature = "graal-forge")]
830impl Drop for GraalBridge {
831    fn drop(&mut self) {
832        if self.attached {
833            unsafe { graal_ffi::graal_detach_thread(self.thread) };
834        } else {
835            unsafe { graal_ffi::graal_tear_down_isolate(self.thread) };
836        }
837    }
838}
839
840#[cfg(feature = "graal-forge")]
841#[derive(Clone)]
842struct GraalEngineHandle {
843    bridge: std::rc::Rc<GraalBridge>,
844}
845
846#[cfg(feature = "graal-forge")]
847impl GraalEngineHandle {
848    fn create(assets_dir: &Path) -> Result<Self, String> {
849        let bridge = if shared_isolate_enabled() {
850            GraalBridge::create_in_shared_isolate(assets_dir)?
851        } else {
852            let bridge = GraalBridge::create()?;
853            let assets = cstring(&assets_dir.to_string_lossy())?;
854            bridge
855                .decode(unsafe { graal_ffi::forge_initialize(bridge.thread, assets.as_ptr()) })?;
856            bridge
857        };
858        Ok(Self {
859            bridge: std::rc::Rc::new(bridge),
860        })
861    }
862
863    fn start_game(&self, request_json: &str) -> Result<String, String> {
864        let request = cstring(request_json)?;
865        let response = self
866            .bridge
867            .decode(unsafe { graal_ffi::forge_start_game(self.bridge.thread, request.as_ptr()) })?;
868        let parsed: StartGameResponse = serde_json::from_str(&response)
869            .map_err(|err| format!("malformed startGame response: {err}"))?;
870        Ok(parsed.session_id)
871    }
872
873    fn submit_action(&self, session_id: &str, action_json: &str) -> Result<String, String> {
874        let session = cstring(session_id)?;
875        let action = cstring(action_json)?;
876        self.bridge.decode(unsafe {
877            graal_ffi::forge_submit_action(self.bridge.thread, session.as_ptr(), action.as_ptr())
878        })
879    }
880
881    fn get_prompt(&self, session_id: &str, player_index: usize) -> Result<Option<String>, String> {
882        let session = cstring(session_id)?;
883        let prompt = self.bridge.decode(unsafe {
884            graal_ffi::forge_get_prompt(
885                self.bridge.thread,
886                session.as_ptr(),
887                player_index as std::os::raw::c_int,
888            )
889        })?;
890        Ok((!prompt.is_empty()).then_some(prompt))
891    }
892
893    fn is_game_over(&self, session_id: &str) -> Result<bool, String> {
894        let session = cstring(session_id)?;
895        let value = self.bridge.decode(unsafe {
896            graal_ffi::forge_get_game_over(self.bridge.thread, session.as_ptr())
897        })?;
898        Ok(value.trim() == "true")
899    }
900
901    fn get_snapshot(&self, session_id: &str, viewer: Option<usize>) -> Result<String, String> {
902        let session = cstring(session_id)?;
903        let viewer = viewer.map_or(-1, |v| v as std::os::raw::c_int);
904        self.bridge.decode(unsafe {
905            graal_ffi::forge_get_snapshot(self.bridge.thread, session.as_ptr(), viewer)
906        })
907    }
908
909    fn end_game(&self, session_id: &str) -> Result<(), String> {
910        let session = cstring(session_id)?;
911        self.bridge
912            .decode(unsafe { graal_ffi::forge_end_game(self.bridge.thread, session.as_ptr()) })
913            .map(|_| ())
914    }
915
916    fn abort_game(&self, session_id: &str) -> Result<(), String> {
917        let session = cstring(session_id)?;
918        self.bridge
919            .decode(unsafe { graal_ffi::forge_abort_game(self.bridge.thread, session.as_ptr()) })
920            .map(|_| ())
921    }
922}
923
924#[cfg(feature = "graal-forge")]
925fn cstring(value: &str) -> Result<std::ffi::CString, String> {
926    std::ffi::CString::new(value).map_err(|_| "string contained interior NUL".to_string())
927}
928
929#[cfg(feature = "java-forge")]
930pub fn run_concurrent_self_play(
931    seats: &[DeckSelection],
932    starting_life: i32,
933    seed: u64,
934    max_prompts: usize,
935    concurrency: usize,
936) -> Result<(), String> {
937    let config = JavaRuntimeConfig::from_env();
938    let games_per_process = env::var("SELF_HOSTED_NODE_GAMES_PER_JVM")
939        .ok()
940        .and_then(|value| value.parse::<usize>().ok())
941        .filter(|n| *n >= 1)
942        .unwrap_or(1);
943    let pool = JavaEnginePool::start(&config, concurrency.max(1), games_per_process)?;
944    info!(
945        concurrency,
946        "java-engine started; launching concurrent games"
947    );
948
949    let mut players = Vec::with_capacity(seats.len());
950    for (i, seat) in seats.iter().enumerate() {
951        let identities = deck_card_identities(&seat.deck);
952        players.push(PlayerConfig::new(
953            format!("Self-Play {}", i + 1),
954            &identities,
955            commander_names_for_java(&seat.deck, seat.commander_name.as_deref()),
956        ));
957    }
958
959    let mut joins = Vec::with_capacity(concurrency.max(1));
960    for game_index in 0..concurrency.max(1) {
961        let handle = pool.handle();
962        let request = StartGameRequest::new(
963            format!("self-hosted-java-concurrent-{game_index}"),
964            String::new(),
965            starting_life,
966            seed.wrapping_add(game_index as u64),
967            players.clone(),
968        );
969        joins.push(std::thread::spawn(move || -> Result<(), String> {
970            let request_json = request.to_json().map_err(|error| error.to_string())?;
971            let session_id = handle.start_game(&request_json)?;
972            info!(session_id, game_index, "concurrent java game started");
973            let result = drive_game_via_handle(&handle, &session_id, max_prompts);
974            let _ = handle.end_game(&session_id);
975            result
976        }));
977    }
978
979    let mut outcome = Ok(());
980    for join in joins {
981        match join.join() {
982            Ok(Ok(())) => {}
983            Ok(Err(error)) => outcome = Err(error),
984            Err(_) => outcome = Err("concurrent game thread panicked".to_string()),
985        }
986    }
987    outcome
988}
989
990#[cfg(not(feature = "java-forge"))]
991pub fn run_concurrent_self_play(
992    _seats: &[DeckSelection],
993    _starting_life: i32,
994    _seed: u64,
995    _max_prompts: usize,
996    _concurrency: usize,
997) -> Result<(), String> {
998    Err(
999        "java-forge concurrent self-play requires building self-hosted-node with --features java-forge"
1000            .to_string(),
1001    )
1002}
1003
1004#[cfg(feature = "java-forge")]
1005fn drive_game_via_handle(
1006    handle: &JavaEngineHandle,
1007    session_id: &str,
1008    max_prompts: usize,
1009) -> Result<(), String> {
1010    let mut bots: HashMap<usize, SimpleAi> = HashMap::new();
1011    let mut last_prompt: Option<String> = None;
1012    let mut acted = 0usize;
1013    let mut seen_prompt = false;
1014    let max_iterations = max_prompts.saturating_mul(200).max(2_000);
1015
1016    for _ in 0..max_iterations {
1017        if let Some(prompt_json) = handle.get_prompt(session_id, 0)? {
1018            seen_prompt = true;
1019            if last_prompt.as_deref() == Some(prompt_json.as_str()) {
1020                if handle.is_game_over(session_id)? {
1021                    return Ok(());
1022                }
1023                std::thread::sleep(Duration::from_millis(20));
1024                continue;
1025            }
1026            let prompt: AgentPrompt = serde_json::from_str(&prompt_json)
1027                .map_err(|error| format!("failed to parse concurrent prompt: {error}"))?;
1028            let player = player_index(&prompt.deciding_player_id);
1029            if let Some(action) = bots.entry(player).or_default().decide(prompt) {
1030                let action_json = serde_json::to_string(&action).map_err(|err| err.to_string())?;
1031                handle.submit_action(session_id, &action_json)?;
1032                acted += 1;
1033                if acted >= max_prompts {
1034                    return Err(format!(
1035                        "concurrent game {session_id} did not finish within {max_prompts} decisions"
1036                    ));
1037                }
1038            }
1039            last_prompt = Some(prompt_json);
1040            continue;
1041        }
1042        if seen_prompt && handle.is_game_over(session_id)? {
1043            return Ok(());
1044        }
1045        std::thread::sleep(Duration::from_millis(20));
1046    }
1047    Err(format!(
1048        "concurrent game {session_id} exceeded its iteration cap"
1049    ))
1050}
1051
1052/// G1 keeps pauses roughly bounded as the live set grows; Serial does not.
1053const DEFAULT_JAVA_COLLECTOR: &str = "G1";
1054const DEFAULT_JAVA_GC_LOG: &str = "stderr";
1055const DEFAULT_JAVA_HEAP_MB: u64 = 1024;
1056const DEFAULT_JAVA_ACTIVE_PROCESSORS: u64 = 2;
1057
1058#[derive(Debug, Clone)]
1059pub struct JavaRuntimeConfig {
1060    pub assets_dir: PathBuf,
1061    pub harness_jar: PathBuf,
1062    pub java_home: Option<PathBuf>,
1063    pub extra_classpath: Vec<PathBuf>,
1064    pub heap_mb: Option<u64>,
1065    pub active_processor_count: Option<u64>,
1066    pub gc_log: Option<String>,
1067    pub collector: Option<String>,
1068    pub extra_jvm_args: Vec<String>,
1069}
1070
1071impl JavaRuntimeConfig {
1072    pub fn from_env() -> Self {
1073        let root = workspace_root();
1074        Self {
1075            assets_dir: env_path("SELF_HOSTED_NODE_FORGE_ASSETS_DIR")
1076                .or_else(|| env_path("MANA_BREW_FORGE_ASSETS_DIR"))
1077                .unwrap_or_else(|| root.join("forge/forge-gui")),
1078            harness_jar: env_path("SELF_HOSTED_NODE_FORGE_HARNESS_JAR")
1079                .or_else(|| env_path("MANA_BREW_FORGE_HARNESS_JAR"))
1080                .unwrap_or_else(|| {
1081                    root.join("forge-harness/target/forge-harness-jar-with-dependencies.jar")
1082                }),
1083            java_home: env_path("SELF_HOSTED_NODE_JAVA_HOME")
1084                .or_else(|| env_path("MANA_BREW_JAVA_HOME"))
1085                .or_else(|| env_path("JAVA_HOME")),
1086            extra_classpath: env_classpath("SELF_HOSTED_NODE_FORGE_EXTRA_CLASSPATH")
1087                .into_iter()
1088                .chain(env_classpath("MANA_BREW_FORGE_EXTRA_CLASSPATH"))
1089                .collect(),
1090            heap_mb: env_sizing("SELF_HOSTED_NODE_JAVA_HEAP_MB", DEFAULT_JAVA_HEAP_MB),
1091            active_processor_count: env_sizing(
1092                "SELF_HOSTED_NODE_JAVA_ACTIVE_PROCESSORS",
1093                DEFAULT_JAVA_ACTIVE_PROCESSORS,
1094            ),
1095            collector: env::var("SELF_HOSTED_NODE_JAVA_COLLECTOR")
1096                .ok()
1097                .map(|value| value.trim().to_string())
1098                .filter(|value| !value.is_empty())
1099                .or_else(|| Some(DEFAULT_JAVA_COLLECTOR.to_string())),
1100            // Defaults on: the fleet ships no logs, so this log exists to be
1101            // turned into metrics by the stderr pump (see metrics.rs).
1102            gc_log: env::var("SELF_HOSTED_NODE_JAVA_GC_LOG")
1103                .ok()
1104                .map(|value| value.trim().to_string())
1105                .filter(|value| !value.is_empty())
1106                .or_else(|| Some(DEFAULT_JAVA_GC_LOG.to_string())),
1107            extra_jvm_args: env::var("SELF_HOSTED_NODE_JAVA_OPTS")
1108                .unwrap_or_default()
1109                .split_whitespace()
1110                .map(str::to_string)
1111                .collect(),
1112        }
1113    }
1114
1115    #[cfg(feature = "java-forge")]
1116    fn jvm_args(&self) -> Vec<String> {
1117        let mut args = vec![
1118            "-Dfile.encoding=UTF-8".to_string(),
1119            "-Dsun.stdout.encoding=UTF-8".to_string(),
1120            "-Dsun.stderr.encoding=UTF-8".to_string(),
1121            "-Djava.awt.headless=true".to_string(),
1122            // Forge's endstep concurrency patches assume a heap exhaustion kills the process
1123            // instead of thrashing; a supervised node is restarted, a thrashing one is not.
1124            "-XX:+ExitOnOutOfMemoryError".to_string(),
1125            // Forge calls System.gc() on match boundaries (Match.java,
1126            // HostedMatch.java) and so does the harness. Each one is a full
1127            // collection, and a full collection costs about 1.3ms per MB of
1128            // live set, so on a large board they are seconds of stop-the-world
1129            // for no benefit the collector would not have reached anyway.
1130            "-XX:+DisableExplicitGC".to_string(),
1131        ];
1132        // Pause time is linear in the live set and Serial collects on one
1133        // thread, so a bigger heap under Serial means longer freezes, not
1134        // shorter ones. Name the collector rather than inheriting whatever the
1135        // host's environment happens to set.
1136        if let Some(collector) = &self.collector {
1137            args.push(format!("-XX:+Use{collector}GC"));
1138        }
1139        // A fleet runs many of these side by side, and every JVM sizes its heap and its GC
1140        // thread count from the whole machine unless told otherwise.
1141        if let Some(heap_mb) = self.heap_mb {
1142            args.push(format!("-Xmx{heap_mb}m"));
1143        }
1144        if let Some(processors) = self.active_processor_count {
1145            args.push(format!("-XX:ActiveProcessorCount={processors}"));
1146        }
1147        // In a container the only log that reaches Loki is the node's own, and the
1148        // subprocess owns stdout for the protocol, so "stderr" routes GC through the
1149        // stderr pump instead of a file nothing ships.
1150        match self.gc_log.as_deref() {
1151            None => {}
1152            Some("stderr") => args.push("-Xlog:gc*:stderr:time,uptime,level,tags".to_string()),
1153            Some(dir) => args.push(format!(
1154                "-Xlog:gc*:file={}:time,uptime,level,tags:filecount=5,filesize=20M",
1155                Path::new(dir).join("engine-gc-%p.log").display()
1156            )),
1157        }
1158        args.extend(self.extra_jvm_args.iter().cloned());
1159        args
1160    }
1161
1162    pub fn validate(&self) -> Result<(), String> {
1163        require_dir(&self.assets_dir, "Forge assets directory")?;
1164        require_file(&self.harness_jar, "Forge harness jar")?;
1165        if let Some(java_home) = &self.java_home {
1166            require_dir(java_home, "Java home")?;
1167        }
1168        for entry in &self.extra_classpath {
1169            if !entry.exists() {
1170                return Err(format!(
1171                    "Classpath entry does not exist: {}",
1172                    entry.display()
1173                ));
1174            }
1175        }
1176        Ok(())
1177    }
1178
1179    pub fn classpath_entries(&self) -> Vec<PathBuf> {
1180        let mut entries = Vec::with_capacity(1 + self.extra_classpath.len());
1181        entries.push(self.harness_jar.clone());
1182        entries.extend(self.extra_classpath.iter().cloned());
1183        entries
1184    }
1185}
1186
1187#[cfg(forge_backend)]
1188#[allow(clippy::too_many_arguments)]
1189pub fn run_hosted_engine_game(
1190    game_id: String,
1191    player_names: Vec<String>,
1192    decks: Vec<Deck>,
1193    commander_names: Vec<Option<String>>,
1194    commander_variant: bool,
1195    game_variant: String,
1196    local_player_index: Option<usize>,
1197    ai_player_indices: Vec<usize>,
1198    starting_life: i32,
1199    remote_prompt_tx: std_mpsc::Sender<(usize, AgentMessage)>,
1200    remote_response_rxs: Vec<(usize, std_mpsc::Receiver<ClientToServerMessage>)>,
1201    game_over_tx: std_mpsc::Sender<HostedGameOver>,
1202    cancel: Arc<AtomicBool>,
1203) -> Result<(), String> {
1204    run_hosted_engine_game_inner(
1205        game_id,
1206        player_names,
1207        decks,
1208        commander_names,
1209        commander_variant,
1210        game_variant,
1211        local_player_index,
1212        ai_player_indices,
1213        starting_life,
1214        remote_prompt_tx,
1215        remote_response_rxs,
1216        game_over_tx,
1217        cancel,
1218    )
1219}
1220
1221#[cfg(not(forge_backend))]
1222#[allow(clippy::too_many_arguments)]
1223pub fn run_hosted_engine_game(
1224    _game_id: String,
1225    _player_names: Vec<String>,
1226    _decks: Vec<Deck>,
1227    _commander_names: Vec<Option<String>>,
1228    _commander_variant: bool,
1229    _game_variant: String,
1230    _local_player_index: Option<usize>,
1231    _ai_player_indices: Vec<usize>,
1232    _starting_life: i32,
1233    _remote_prompt_tx: std_mpsc::Sender<(usize, AgentMessage)>,
1234    _remote_response_rxs: Vec<(usize, std_mpsc::Receiver<ClientToServerMessage>)>,
1235    _game_over_tx: std_mpsc::Sender<HostedGameOver>,
1236    _cancel: Arc<AtomicBool>,
1237) -> Result<(), String> {
1238    Err(unsupported_message().to_string())
1239}
1240
1241#[cfg(forge_backend)]
1242#[allow(clippy::too_many_arguments)]
1243fn run_hosted_engine_game_inner(
1244    game_id: String,
1245    player_names: Vec<String>,
1246    decks: Vec<Deck>,
1247    commander_names: Vec<Option<String>>,
1248    commander_variant: bool,
1249    game_variant: String,
1250    local_player_index: Option<usize>,
1251    ai_player_indices: Vec<usize>,
1252    starting_life: i32,
1253    remote_prompt_tx: std_mpsc::Sender<(usize, AgentMessage)>,
1254    remote_response_rxs: Vec<(usize, std_mpsc::Receiver<ClientToServerMessage>)>,
1255    game_over_tx: std_mpsc::Sender<HostedGameOver>,
1256    cancel: Arc<AtomicBool>,
1257) -> Result<(), String> {
1258    let engine = obtain_engine()?;
1259
1260    let mut players = Vec::with_capacity(player_names.len());
1261    for (index, name) in player_names.iter().enumerate() {
1262        let identities = deck_card_identities(&decks[index]);
1263        let seat_commander_names = if commander_variant {
1264            commander_names_for_java(&decks[index], commander_names[index].as_deref())
1265        } else {
1266            Vec::new()
1267        };
1268        players.push(PlayerConfig::new(
1269            name.clone(),
1270            &identities,
1271            seat_commander_names,
1272        ));
1273    }
1274    for &idx in &ai_player_indices {
1275        if let Some(player) = players.get_mut(idx) {
1276            player.ai = true;
1277        }
1278    }
1279    let request = StartGameRequest::new(
1280        game_id.clone(),
1281        game_variant,
1282        starting_life,
1283        rand::random(),
1284        players,
1285    );
1286    let session_id = engine.start_game(&request.to_json().map_err(|err| err.to_string())?)?;
1287    info!(game_id, session_id, "hosted java-forge session started");
1288
1289    struct SessionGuard {
1290        engine: ForgeEngine,
1291        session_id: String,
1292        armed: std::cell::Cell<bool>,
1293    }
1294    impl Drop for SessionGuard {
1295        fn drop(&mut self) {
1296            if self.armed.get() {
1297                if let Err(error) = self.engine.abort_game(&self.session_id) {
1298                    warn!(session_id = %self.session_id, %error, "failed to abort java session; context may leak");
1299                }
1300            }
1301        }
1302    }
1303    let guard = SessionGuard {
1304        engine: engine.clone(),
1305        session_id: session_id.clone(),
1306        armed: std::cell::Cell::new(true),
1307    };
1308
1309    let mut remote_response_rxs: HashMap<usize, std_mpsc::Receiver<ClientToServerMessage>> =
1310        remote_response_rxs.into_iter().collect();
1311    let mut last_prompt: Option<AgentPrompt> = None;
1312    let mut pending_roll_acks: usize = 0;
1313    let mut decision_received: Option<Instant> = None;
1314    let mut decision_submitted: Option<Instant> = None;
1315
1316    loop {
1317        if cancel.load(std::sync::atomic::Ordering::Relaxed) {
1318            info!(
1319                session_id,
1320                "hosted java-forge session cancelled; player left the game"
1321            );
1322            return Ok(());
1323        }
1324        for (player_index, rx) in &mut remote_response_rxs {
1325            loop {
1326                match rx.try_recv() {
1327                    Ok(ClientToServerMessage::Response {
1328                        action: PromptOutput::DiceRolled(DiceRolledOutput::DiceRolledAcknowledged),
1329                        ..
1330                    }) => {
1331                        if pending_roll_acks > 0 {
1332                            pending_roll_acks -= 1;
1333                            if pending_roll_acks == 0 {
1334                                let ack = serde_json::to_string(&PromptOutput::DiceRolled(
1335                                    DiceRolledOutput::DiceRolledAcknowledged,
1336                                ))
1337                                .map_err(|err| format!("failed to serialize roll ack: {err}"))?;
1338                                engine.submit_action(&session_id, &ack)?;
1339                            }
1340                        }
1341                    }
1342                    Ok(ClientToServerMessage::Response { prompt_id, action }) => {
1343                        // prompt_id 0 is transport-synthesized: the
1344                        // absent-player default, exempt from validation.
1345                        if prompt_id != 0 {
1346                            let Some(prompt) =
1347                                last_prompt.as_ref().filter(|p| p.prompt_id == prompt_id)
1348                            else {
1349                                reject_response(
1350                                    &remote_prompt_tx,
1351                                    *player_index,
1352                                    last_prompt.as_ref().filter(|p| {
1353                                        self::player_index(&p.deciding_player_id) == *player_index
1354                                    }),
1355                                    ProtocolErrorCode::StalePrompt,
1356                                    format!("response for prompt {prompt_id} is not open"),
1357                                );
1358                                continue;
1359                            };
1360                            if self::player_index(&prompt.deciding_player_id) != *player_index {
1361                                reject_response(
1362                                    &remote_prompt_tx,
1363                                    *player_index,
1364                                    None,
1365                                    ProtocolErrorCode::WrongPlayer,
1366                                    format!(
1367                                        "prompt {prompt_id} is for {}",
1368                                        prompt.deciding_player_id
1369                                    ),
1370                                );
1371                                continue;
1372                            }
1373                            match prompt.input.validate_response(&action) {
1374                                Ok(()) => {}
1375                                Err(ResponseViolation::WrongPromptType) => {
1376                                    reject_response(
1377                                        &remote_prompt_tx,
1378                                        *player_index,
1379                                        Some(prompt),
1380                                        ProtocolErrorCode::WrongPromptType,
1381                                        "response output does not match the prompt type"
1382                                            .to_string(),
1383                                    );
1384                                    continue;
1385                                }
1386                                Err(ResponseViolation::UnknownActionId(id)) => {
1387                                    reject_response(
1388                                        &remote_prompt_tx,
1389                                        *player_index,
1390                                        Some(prompt),
1391                                        ProtocolErrorCode::UnknownActionId,
1392                                        format!(
1393                                            "action id {id:?} was not advertised by the prompt"
1394                                        ),
1395                                    );
1396                                    continue;
1397                                }
1398                                Err(ResponseViolation::CancelNotAllowed) => {
1399                                    reject_response(
1400                                        &remote_prompt_tx,
1401                                        *player_index,
1402                                        Some(prompt),
1403                                        ProtocolErrorCode::CancelNotAllowed,
1404                                        "this prompt is not cancellable".to_string(),
1405                                    );
1406                                    continue;
1407                                }
1408                            }
1409                        }
1410                        decision_received = Some(Instant::now());
1411                        let action_json = serde_json::to_string(&action).map_err(|err| {
1412                            format!(
1413                                "failed to serialize prompt output for player {player_index}: {err}"
1414                            )
1415                        })?;
1416                        debug!(player_index, %action_json, "submitting remote response to java");
1417                        engine.submit_action(&session_id, &action_json)?;
1418                        decision_submitted = Some(Instant::now());
1419                    }
1420                    Ok(ClientToServerMessage::Directive {
1421                        directive: DirectiveInput::Concede,
1422                    }) => {
1423                        // A directive can arrive while another player's prompt
1424                        // is open (this loop drains every seat), so it names
1425                        // its seat.
1426                        let directive_json = directive_concede_json(*player_index);
1427                        debug!(player_index, %directive_json, "submitting concede directive to java");
1428                        engine.submit_action(&session_id, &directive_json)?;
1429                    }
1430                    Err(TryRecvError::Empty) => break,
1431                    Err(TryRecvError::Disconnected) => {
1432                        debug!(player_index, "java-forge response channel disconnected");
1433                        break;
1434                    }
1435                }
1436            }
1437        }
1438
1439        if let Some(prompt_json) = engine.get_prompt(&session_id, 0)? {
1440            let prompt: AgentPrompt = serde_json::from_str(&prompt_json)
1441                .map_err(|err| format!("failed to parse java prompt: {err}"))?;
1442            if last_prompt.as_ref().map(|p| p.prompt_id) != Some(prompt.prompt_id) {
1443                if let Some(started) = decision_submitted.take() {
1444                    crate::metrics::record_forge_decision_stage("next_prompt", started.elapsed());
1445                }
1446                if let Some(started) = decision_received.take() {
1447                    crate::metrics::record_forge_decision_stage(
1448                        "decision_total",
1449                        started.elapsed(),
1450                    );
1451                }
1452                last_prompt = Some(prompt.clone());
1453                let player = player_index(&prompt.deciding_player_id);
1454                debug!(player, "forwarding java prompt to remote");
1455                if matches!(prompt.input, PromptInput::DiceRolled(_)) {
1456                    let snapshots_started = Instant::now();
1457                    let prompt_msg = AgentMessage::Prompt(prompt);
1458                    for &agent_index in remote_response_rxs.keys() {
1459                        let state = AgentMessage::State(state_via_handle(
1460                            &engine,
1461                            &session_id,
1462                            Some(agent_index),
1463                        )?);
1464                        let _ = remote_prompt_tx.send((agent_index, state));
1465                        let _ = remote_prompt_tx.send((agent_index, prompt_msg.clone()));
1466                    }
1467                    send_observer_state(&engine, &session_id, &remote_prompt_tx);
1468                    pending_roll_acks = remote_response_rxs.len();
1469                    if pending_roll_acks == 0 {
1470                        let ack = serde_json::to_string(&PromptOutput::DiceRolled(
1471                            DiceRolledOutput::DiceRolledAcknowledged,
1472                        ))
1473                        .map_err(|err| format!("failed to serialize roll ack: {err}"))?;
1474                        engine.submit_action(&session_id, &ack)?;
1475                    }
1476                    crate::metrics::record_forge_decision_stage(
1477                        "snapshots",
1478                        snapshots_started.elapsed(),
1479                    );
1480                } else if Some(player) == local_player_index {
1481                    if let Some(output) = auto_action(&prompt) {
1482                        let action_json = serde_json::to_string(&output)
1483                            .map_err(|err| format!("failed to serialize auto action: {err}"))?;
1484                        engine.submit_action(&session_id, &action_json)?;
1485                    }
1486                } else {
1487                    let snapshots_started = Instant::now();
1488                    for &agent_index in remote_response_rxs.keys() {
1489                        let state = AgentMessage::State(state_via_handle(
1490                            &engine,
1491                            &session_id,
1492                            Some(agent_index),
1493                        )?);
1494                        if remote_prompt_tx.send((agent_index, state)).is_err() {
1495                            return Ok(());
1496                        }
1497                    }
1498                    send_observer_state(&engine, &session_id, &remote_prompt_tx);
1499                    let prompt_msg = AgentMessage::Prompt(prompt);
1500                    if remote_prompt_tx.send((player, prompt_msg)).is_err() {
1501                        return Ok(());
1502                    }
1503                    crate::metrics::record_forge_decision_stage(
1504                        "snapshots",
1505                        snapshots_started.elapsed(),
1506                    );
1507                }
1508            }
1509        }
1510
1511        if engine.is_game_over(&session_id)? {
1512            info!("hosted java-forge session reached game over");
1513            let mut final_messages = Vec::new();
1514            for &agent_index in remote_response_rxs.keys() {
1515                match state_via_handle(&engine, &session_id, Some(agent_index)) {
1516                    Ok(state_update) => {
1517                        final_messages.push((agent_index, AgentMessage::State(state_update)));
1518                    }
1519                    Err(error) => {
1520                        warn!(%error, agent_index, "game over: final snapshot unavailable; sending game-over prompt only");
1521                    }
1522                }
1523            }
1524            if let Ok(state_update) = state_via_handle(&engine, &session_id, None) {
1525                final_messages.push((
1526                    crate::host::OBSERVER_SEAT,
1527                    AgentMessage::State(state_update),
1528                ));
1529            }
1530            let game_over = AgentMessage::Prompt(game_over_prompt());
1531            for &agent_index in remote_response_rxs.keys() {
1532                final_messages.push((agent_index, game_over.clone()));
1533            }
1534            let _ = game_over_tx.send(HostedGameOver {
1535                game_id: game_id.clone(),
1536                messages: final_messages,
1537            });
1538            engine.end_game(&session_id)?;
1539            guard.armed.set(false);
1540            return Ok(());
1541        }
1542
1543        std::thread::sleep(Duration::from_millis(50));
1544    }
1545}
1546
1547#[cfg(feature = "java-forge")]
1548fn wait_for_prompt<B: JavaBridge>(
1549    session: &mut JavaForgeSession<B>,
1550    max_polls: usize,
1551) -> Result<Option<String>, String> {
1552    for _ in 0..max_polls {
1553        if let Some(prompt) = session.get_prompt(0)? {
1554            return Ok(Some(prompt));
1555        }
1556        std::thread::sleep(Duration::from_millis(50));
1557    }
1558    Ok(None)
1559}
1560
1561#[cfg(forge_backend)]
1562fn send_observer_state(
1563    engine: &ForgeEngine,
1564    session_id: &str,
1565    remote_prompt_tx: &std_mpsc::Sender<(usize, AgentMessage)>,
1566) {
1567    match state_via_handle(engine, session_id, None) {
1568        Ok(state_update) => {
1569            let _ = remote_prompt_tx.send((
1570                crate::host::OBSERVER_SEAT,
1571                AgentMessage::State(state_update),
1572            ));
1573        }
1574        Err(error) => warn!(%error, "observer snapshot unavailable"),
1575    }
1576}
1577
1578#[cfg(forge_backend)]
1579fn player_index(deciding_player_id: &str) -> usize {
1580    deciding_player_id
1581        .strip_prefix("player-")
1582        .and_then(|n| n.parse().ok())
1583        .unwrap_or(0)
1584}
1585
1586#[cfg(forge_backend)]
1587fn reject_response(
1588    remote_prompt_tx: &std_mpsc::Sender<(usize, AgentMessage)>,
1589    seat: usize,
1590    reopen_prompt: Option<&AgentPrompt>,
1591    code: ProtocolErrorCode,
1592    message: String,
1593) {
1594    let _ = remote_prompt_tx.send((
1595        seat,
1596        AgentMessage::Error(ProtocolError {
1597            code,
1598            message,
1599            prompt_id: reopen_prompt.map(|p| p.prompt_id),
1600        }),
1601    ));
1602    if let Some(prompt) = reopen_prompt {
1603        let _ = remote_prompt_tx.send((seat, AgentMessage::Prompt(prompt.clone())));
1604    }
1605}
1606
1607#[cfg(forge_backend)]
1608fn auto_action(prompt: &AgentPrompt) -> Option<PromptOutput> {
1609    match prompt.input {
1610        PromptInput::ChooseAction(_) => {
1611            Some(PromptOutput::ChooseAction(ChooseActionOutput::Pass {
1612                until: None,
1613                exhaust_stack: false,
1614            }))
1615        }
1616        _ => None,
1617    }
1618}
1619
1620#[cfg(forge_backend)]
1621fn game_over_prompt() -> AgentPrompt {
1622    AgentPrompt {
1623        prompt_id: u32::MAX,
1624        deciding_player_id: "player-0".to_string(),
1625        source_card: None,
1626        input: PromptInput::GameOver(GameOverInput {}),
1627    }
1628}
1629
1630#[cfg(forge_backend)]
1631fn state_via_handle(
1632    engine: &ForgeEngine,
1633    session_id: &str,
1634    viewer: Option<usize>,
1635) -> Result<StateUpdate, String> {
1636    let game_view: GameViewDto = serde_json::from_str(&engine.get_snapshot(session_id, viewer)?)
1637        .map_err(|err| format!("failed to parse java snapshot: {err}"))?;
1638    Ok(StateUpdate { game_view })
1639}
1640
1641#[cfg(forge_backend)]
1642fn deck_card_identities(deck: &Deck) -> Vec<DeckCardIdentity> {
1643    deck.cards
1644        .iter()
1645        .chain(deck.commanders.iter().flatten())
1646        .map(|card| card.identity.clone())
1647        .collect()
1648}
1649
1650#[cfg(forge_backend)]
1651fn commander_names_for_java(deck: &Deck, fallback: Option<&str>) -> Vec<String> {
1652    let names: Vec<String> = deck
1653        .commanders
1654        .iter()
1655        .flatten()
1656        .map(|card| java_card_name(&card.identity.name))
1657        .collect();
1658    if !names.is_empty() {
1659        return names;
1660    }
1661    fallback
1662        .filter(|name| !name.is_empty())
1663        .map(|name| vec![java_card_name(name)])
1664        .unwrap_or_default()
1665}
1666
1667#[cfg(forge_backend)]
1668fn smoke_deck(land_name: &str, spell_name: &str) -> Vec<DeckCardIdentity> {
1669    (0..24)
1670        .map(|_| DeckCardIdentity {
1671            name: land_name.to_string(),
1672            ..Default::default()
1673        })
1674        .chain((0..36).map(|_| DeckCardIdentity {
1675            name: spell_name.to_string(),
1676            ..Default::default()
1677        }))
1678        .collect()
1679}
1680
1681#[cfg(feature = "java-forge")]
1682fn scenario_deck(land_name: &str) -> Vec<DeckCardIdentity> {
1683    (0..60)
1684        .map(|_| DeckCardIdentity {
1685            name: land_name.to_string(),
1686            ..Default::default()
1687        })
1688        .collect()
1689}
1690
1691#[cfg(feature = "java-forge")]
1692enum JavaScenario {
1693    KeepAndPlayLand {
1694        played_land: bool,
1695    },
1696    MulliganOncePlayLand {
1697        mulliganed: bool,
1698        kept_second_hand: bool,
1699        put_back_done: bool,
1700        played_land: bool,
1701    },
1702}
1703
1704#[cfg(feature = "java-forge")]
1705impl JavaScenario {
1706    fn from_name(name: &str) -> Result<Self, String> {
1707        match name {
1708            "keep-and-play-land" => Ok(Self::KeepAndPlayLand { played_land: false }),
1709            "mulligan-once-play-land" => Ok(Self::MulliganOncePlayLand {
1710                mulliganed: false,
1711                kept_second_hand: false,
1712                put_back_done: false,
1713                played_land: false,
1714            }),
1715            _ => Err(format!(
1716                "unknown java-forge scenario '{name}'. Supported scenarios: keep-and-play-land, mulligan-once-play-land"
1717            )),
1718        }
1719    }
1720
1721    fn name(&self) -> &'static str {
1722        match self {
1723            Self::KeepAndPlayLand { .. } => "keep-and-play-land",
1724            Self::MulliganOncePlayLand { .. } => "mulligan-once-play-land",
1725        }
1726    }
1727
1728    fn next_action(
1729        &mut self,
1730        prompt: &Value,
1731        game_view: &GameViewDto,
1732    ) -> Result<Option<PromptOutput>, String> {
1733        match self {
1734            Self::KeepAndPlayLand { played_land } => {
1735                if *played_land && battlefield_contains(game_view, "Swamp") {
1736                    return Ok(None);
1737                }
1738                match prompt_type(prompt) {
1739                    Some("mulligan") => Ok(Some(PromptOutput::Mulligan(
1740                        MulliganOutput::MulliganDecision { keep: true },
1741                    ))),
1742                    Some("chooseAction") => {
1743                        if let Some(action) = play_first_card_action(prompt, "Swamp")? {
1744                            *played_land = true;
1745                            Ok(Some(action))
1746                        } else {
1747                            Ok(Some(PromptOutput::ChooseAction(ChooseActionOutput::Pass {
1748                                until: None,
1749                                exhaust_stack: false,
1750                            })))
1751                        }
1752                    }
1753                    other => Err(format!(
1754                        "scenario '{}' expected mulligan or chooseAction, got {:?}",
1755                        self.name(),
1756                        other
1757                    )),
1758                }
1759            }
1760            Self::MulliganOncePlayLand {
1761                mulliganed,
1762                kept_second_hand,
1763                put_back_done,
1764                played_land,
1765            } => {
1766                if *played_land && battlefield_contains(game_view, "Swamp") {
1767                    return Ok(None);
1768                }
1769                match prompt_type(prompt) {
1770                    Some("mulligan") if !*mulliganed => {
1771                        *mulliganed = true;
1772                        Ok(Some(PromptOutput::Mulligan(MulliganOutput::MulliganDecision { keep: false })))
1773                    }
1774                    Some("mulligan") if !*kept_second_hand => {
1775                        *kept_second_hand = true;
1776                        Ok(Some(PromptOutput::Mulligan(MulliganOutput::MulliganDecision { keep: true })))
1777                    }
1778                    Some("mulliganPutBack") if !*put_back_done => {
1779                        let count = prompt
1780                            .get("input")
1781                            .and_then(|input| input.get("count"))
1782                            .and_then(Value::as_u64)
1783                            .unwrap_or(1) as usize;
1784                        let card_ids = prompt_card_ids(prompt, "handCardIds", count)?;
1785                        *put_back_done = true;
1786                        Ok(Some(PromptOutput::MulliganPutBack(MulliganPutBackOutput::MulliganPutBackDecision { card_ids })))
1787                    }
1788                    Some("chooseAction") => {
1789                        if let Some(action) = play_first_card_action(prompt, "Swamp")? {
1790                            *played_land = true;
1791                            Ok(Some(action))
1792                        } else {
1793                            Ok(Some(PromptOutput::ChooseAction(ChooseActionOutput::Pass { until: None, exhaust_stack: false })))
1794                        }
1795                    }
1796                    other => Err(format!(
1797                        "scenario '{}' expected mulligan, mulliganPutBack, or chooseAction, got {:?}",
1798                        self.name(),
1799                        other
1800                    )),
1801                }
1802            }
1803        }
1804    }
1805}
1806
1807#[cfg(feature = "java-forge")]
1808fn run_scenario_loop<B: JavaBridge>(
1809    session: &mut JavaForgeSession<B>,
1810    mut scenario: JavaScenario,
1811    max_prompts: usize,
1812) -> Result<(), String> {
1813    let mut prompts_seen = 0usize;
1814    let mut last_prompt_json: Option<String> = None;
1815    while prompts_seen < max_prompts {
1816        let Some(prompt_json) = wait_for_prompt(session, 600)? else {
1817            return Err(format!(
1818                "timed out waiting for java-forge scenario '{}' prompt",
1819                scenario.name()
1820            ));
1821        };
1822        if last_prompt_json.as_deref() == Some(prompt_json.as_str()) {
1823            std::thread::sleep(Duration::from_millis(50));
1824            continue;
1825        }
1826        last_prompt_json = Some(prompt_json.clone());
1827        prompts_seen += 1;
1828
1829        let prompt: AgentPrompt = serde_json::from_str(&prompt_json)
1830            .map_err(|err| format!("failed to parse java scenario prompt: {err}"))?;
1831        let player = player_index(&prompt.deciding_player_id);
1832        if player != 0 {
1833            if let Some(output) = auto_action(&prompt) {
1834                session.submit_action(
1835                    &serde_json::to_string(&output).map_err(|err| err.to_string())?,
1836                )?;
1837            }
1838            continue;
1839        }
1840
1841        let game_view: GameViewDto = serde_json::from_str(&session.get_snapshot(Some(0))?)
1842            .map_err(|err| format!("failed to parse java scenario snapshot: {err}"))?;
1843        let normalized_prompt = serde_json::to_value(&prompt).map_err(|err| err.to_string())?;
1844        info!(
1845            scenario = scenario.name(),
1846            prompts_seen,
1847            prompt_type = prompt_type(&normalized_prompt).unwrap_or("<missing>"),
1848            "java-forge scenario prompt"
1849        );
1850        let Some(action) = scenario.next_action(&normalized_prompt, &game_view)? else {
1851            info!(
1852                scenario = scenario.name(),
1853                prompts_seen, "java-forge scenario assertions satisfied"
1854            );
1855            return Ok(());
1856        };
1857        submit_player_action(session, &action)?;
1858    }
1859    Err(format!(
1860        "java-forge scenario '{}' did not complete within {max_prompts} prompts",
1861        scenario.name()
1862    ))
1863}
1864
1865#[cfg(feature = "java-forge")]
1866pub fn run_concede_smoke() -> Result<(), String> {
1867    let config = JavaRuntimeConfig::from_env();
1868    let assets_dir = config.assets_dir.to_string_lossy().to_string();
1869    let bridge = SubprocessBridge::spawn(&config)?;
1870    let mut session = JavaForgeSession::new(bridge);
1871    session.initialize(&assets_dir)?;
1872    run_concede_game(&mut session, 3, 2, 12)?;
1873    run_concede_game(&mut session, 2, 1, 8)?;
1874    Ok(())
1875}
1876
1877#[cfg(not(feature = "java-forge"))]
1878pub fn run_concede_smoke() -> Result<(), String> {
1879    Err(unsupported_message().to_string())
1880}
1881
1882#[cfg(forge_backend)]
1883fn directive_concede_json(player: usize) -> String {
1884    format!(r#"{{"type":"directive","directive":{{"type":"concede"}},"player":{player}}}"#)
1885}
1886
1887#[cfg(feature = "java-forge")]
1888fn run_concede_game<B: JavaBridge>(
1889    session: &mut JavaForgeSession<B>,
1890    seats: usize,
1891    conceder: usize,
1892    concede_after: usize,
1893) -> Result<(), String> {
1894    const POST_CONCEDE_DECISIONS: usize = 12;
1895    const STALL_REPEATS: usize = 300;
1896
1897    let cards: Vec<manabrew_protocol::deck_dto::DeckCard> = (0..60)
1898        .map(|_| manabrew_protocol::deck_dto::DeckCard {
1899            identity: DeckCardIdentity {
1900                name: "Mountain".to_string(),
1901                set_code: "M20".to_string(),
1902                ..Default::default()
1903            },
1904            ..Default::default()
1905        })
1906        .collect();
1907    let deck = Deck {
1908        name: "concede-smoke".to_string(),
1909        cards,
1910        ..Default::default()
1911    };
1912    let identities = deck_card_identities(&deck);
1913    let players: Vec<PlayerConfig> = (0..seats)
1914        .map(|i| PlayerConfig::new(format!("Concede {}", i + 1), &identities, Vec::new()))
1915        .collect();
1916    let request = StartGameRequest::new(
1917        format!("concede-smoke-{seats}p"),
1918        String::new(),
1919        20,
1920        7,
1921        players,
1922    );
1923    let session_id = session.start_game(&request)?;
1924    info!(
1925        session_id,
1926        seats, conceder, concede_after, "concede smoke game started"
1927    );
1928
1929    let mut bots: HashMap<usize, SimpleAi> = HashMap::new();
1930    let mut last_prompt_json: Option<String> = None;
1931    let mut acted = 0usize;
1932    let mut acted_after_concede = 0usize;
1933    let mut conceded = false;
1934    let mut repeat_count = 0usize;
1935
1936    for _ in 0..40_000 {
1937        if session.is_game_over()? {
1938            if !conceded {
1939                return Err("concede smoke: game ended before the concede fired".to_string());
1940            }
1941            break;
1942        }
1943        let Some(prompt_json) = session.get_prompt(0)? else {
1944            std::thread::sleep(Duration::from_millis(10));
1945            continue;
1946        };
1947        if last_prompt_json.as_deref() == Some(prompt_json.as_str()) {
1948            repeat_count += 1;
1949            if repeat_count > STALL_REPEATS {
1950                return Err(format!(
1951                    "concede smoke stalled on the same prompt (seats={seats} acted={acted} conceded={conceded} after={acted_after_concede}): {prompt_json}"
1952                ));
1953            }
1954            std::thread::sleep(Duration::from_millis(10));
1955            continue;
1956        }
1957        repeat_count = 0;
1958
1959        let prompt: AgentPrompt = serde_json::from_str(&prompt_json)
1960            .map_err(|err| format!("concede smoke: bad prompt: {err}"))?;
1961        let player = player_index(&prompt.deciding_player_id);
1962
1963        if !conceded && acted >= concede_after {
1964            info!(conceder, acted, "concede smoke: injecting concede");
1965            session.submit_action(&directive_concede_json(conceder))?;
1966            conceded = true;
1967            last_prompt_json = None;
1968            continue;
1969        }
1970        if conceded && player == conceder {
1971            session.submit_action(&directive_concede_json(conceder))?;
1972            last_prompt_json = Some(prompt_json);
1973            continue;
1974        }
1975
1976        if let Some(action) = bots.entry(player).or_default().decide(prompt) {
1977            submit_player_action(session, &action)?;
1978            acted += 1;
1979            if conceded {
1980                acted_after_concede += 1;
1981            }
1982        }
1983        last_prompt_json = Some(prompt_json);
1984        if conceded && seats > 2 && acted_after_concede >= POST_CONCEDE_DECISIONS {
1985            break;
1986        }
1987    }
1988
1989    if !conceded {
1990        return Err("concede smoke: never reached the injection point".to_string());
1991    }
1992    if seats == 2 {
1993        if !session.is_game_over()? {
1994            return Err("concede smoke: 2p concession did not end the game".to_string());
1995        }
1996        info!(acted, "concede smoke: 2p concession ended the game");
1997    } else {
1998        if acted_after_concede < POST_CONCEDE_DECISIONS && !session.is_game_over()? {
1999            return Err(format!(
2000                "concede smoke: game did not progress after the concession (only {acted_after_concede} decisions)"
2001            ));
2002        }
2003        let snapshot = parse_snapshot(session)?;
2004        let status = snapshot
2005            .pointer(&format!("/players/{conceder}/status"))
2006            .and_then(Value::as_str)
2007            .unwrap_or("");
2008        if status != "conceded" {
2009            return Err(format!(
2010                "concede smoke: seat {conceder} has status '{status}', expected 'conceded'"
2011            ));
2012        }
2013        info!(
2014            acted,
2015            acted_after_concede, "concede smoke: game continued past the concession"
2016        );
2017    }
2018    session.end_game()?;
2019    Ok(())
2020}
2021
2022#[cfg(feature = "java-forge")]
2023fn run_self_play_loop<B: JavaBridge>(
2024    session: &mut JavaForgeSession<B>,
2025    max_prompts: usize,
2026) -> Result<(), String> {
2027    const STALL_REPEATS: usize = 100;
2028
2029    let mut bots: HashMap<usize, SimpleAi> = HashMap::new();
2030    let mut last_prompt_json: Option<String> = None;
2031    let mut acted = 0usize;
2032    let mut repeat_count = 0usize;
2033    let mut seen_prompt = false;
2034    let max_iterations = max_prompts.saturating_mul(200).max(2_000);
2035
2036    for _ in 0..max_iterations {
2037        if let Some(prompt_json) = session.get_prompt(0)? {
2038            seen_prompt = true;
2039            if last_prompt_json.as_deref() == Some(prompt_json.as_str()) {
2040                if session.is_game_over()? {
2041                    info!(acted, "java-forge self-play reached game over");
2042                    return Ok(());
2043                }
2044                repeat_count += 1;
2045                if repeat_count > STALL_REPEATS {
2046                    let raw_value: Value =
2047                        serde_json::from_str(&prompt_json).unwrap_or(Value::Null);
2048                    dump_stuck(
2049                        "java re-emitted the same prompt after the bot acted (stall)",
2050                        &raw_value,
2051                        None,
2052                        session,
2053                    );
2054                    return Err(
2055                        "self-play stalled: java re-emitted the same prompt after the bot's action"
2056                            .to_string(),
2057                    );
2058                }
2059                std::thread::sleep(Duration::from_millis(20));
2060                continue;
2061            }
2062            repeat_count = 0;
2063
2064            let prompt: AgentPrompt = serde_json::from_str(&prompt_json)
2065                .map_err(|err| format!("failed to parse java self-play prompt: {err}"))?;
2066            let player = player_index(&prompt.deciding_player_id);
2067            let raw_value: Value = serde_json::from_str(&prompt_json).unwrap_or(Value::Null);
2068            let normalized = raw_value.clone();
2069
2070            match bots.entry(player).or_default().decide(prompt) {
2071                Some(action) => {
2072                    if let Err(err) = submit_player_action(session, &action) {
2073                        dump_stuck(
2074                            "java rejected the bot action",
2075                            &raw_value,
2076                            Some(&normalized),
2077                            session,
2078                        );
2079                        return Err(format!(
2080                            "self-play: java rejected action for player {player}: {err}"
2081                        ));
2082                    }
2083                    acted += 1;
2084                    if acted >= max_prompts {
2085                        dump_stuck(
2086                            "did not reach game over within max prompts",
2087                            &raw_value,
2088                            Some(&normalized),
2089                            session,
2090                        );
2091                        return Err(format!(
2092                            "self-play did not reach game over within {max_prompts} decisions"
2093                        ));
2094                    }
2095                }
2096                None => debug!(
2097                    player,
2098                    prompt_type = prompt_type(&normalized).unwrap_or("<missing>"),
2099                    "self-play: no action for prompt (display-only)"
2100                ),
2101            }
2102            last_prompt_json = Some(prompt_json);
2103            continue;
2104        }
2105
2106        if seen_prompt && session.is_game_over()? {
2107            info!(acted, "java-forge self-play reached game over");
2108            return Ok(());
2109        }
2110        std::thread::sleep(Duration::from_millis(20));
2111    }
2112
2113    dump_stuck(
2114        "self-play exceeded its iteration cap without game over",
2115        &Value::Null,
2116        None,
2117        session,
2118    );
2119    Err("self-play exceeded its iteration cap without reaching game over".to_string())
2120}
2121
2122#[cfg(feature = "java-forge")]
2123fn parse_snapshot<B: JavaBridge>(session: &mut JavaForgeSession<B>) -> Result<Value, String> {
2124    let snapshot_json = session.get_snapshot(Some(0))?;
2125    serde_json::from_str(&snapshot_json)
2126        .map_err(|err| format!("failed to parse java self-play snapshot: {err}"))
2127}
2128
2129#[cfg(feature = "java-forge")]
2130fn dump_stuck<B: JavaBridge>(
2131    reason: &str,
2132    prompt: &Value,
2133    normalized: Option<&Value>,
2134    session: &mut JavaForgeSession<B>,
2135) {
2136    let snapshot = parse_snapshot(session).unwrap_or(Value::Null);
2137    let artifact = json!({
2138        "reason": reason,
2139        "rawPrompt": prompt,
2140        "normalizedPrompt": normalized,
2141        "snapshot": snapshot,
2142    });
2143    let ts = std::time::SystemTime::now()
2144        .duration_since(std::time::UNIX_EPOCH)
2145        .map(|d| d.as_millis())
2146        .unwrap_or(0);
2147    let path = workspace_root().join(format!("target/self-play-stuck-{ts}.json"));
2148    match serde_json::to_string_pretty(&artifact) {
2149        Ok(body) => {
2150            if let Err(error) = std::fs::write(&path, body) {
2151                warn!(%error, reason, "self-play stuck; failed to write artifact");
2152            } else {
2153                warn!(path = %path.display(), reason, "self-play stuck; wrote artifact");
2154            }
2155        }
2156        Err(error) => warn!(%error, reason, "self-play stuck; failed to serialize artifact"),
2157    }
2158}
2159
2160#[cfg(feature = "java-forge")]
2161fn submit_player_action<B: JavaBridge>(
2162    session: &mut JavaForgeSession<B>,
2163    action: &PromptOutput,
2164) -> Result<(), String> {
2165    let action_json = serde_json::to_string(action)
2166        .map_err(|err| format!("failed to serialize scenario action: {err}"))?;
2167    session.submit_action(&action_json)?;
2168    Ok(())
2169}
2170
2171#[cfg(feature = "java-forge")]
2172fn prompt_type(prompt: &Value) -> Option<&str> {
2173    prompt
2174        .get("input")
2175        .and_then(|input| input.get("type"))
2176        .and_then(Value::as_str)
2177}
2178
2179#[cfg(feature = "java-forge")]
2180fn play_first_card_action(prompt: &Value, card_name: &str) -> Result<Option<PromptOutput>, String> {
2181    let Some(action) = prompt
2182        .get("input")
2183        .and_then(|input| input.get("actions"))
2184        .and_then(Value::as_array)
2185        .and_then(|actions| {
2186            actions.iter().find(|action| {
2187                action
2188                    .get("modeLabel")
2189                    .and_then(Value::as_str)
2190                    .is_some_and(|label| label.contains(card_name))
2191            })
2192        })
2193    else {
2194        return Ok(None);
2195    };
2196    let action_id = action
2197        .get("id")
2198        .and_then(Value::as_str)
2199        .ok_or_else(|| format!("playable action for '{card_name}' is missing id"))?;
2200    Ok(Some(PromptOutput::ChooseAction(ChooseActionOutput::Act {
2201        action_id: action_id.to_string(),
2202    })))
2203}
2204
2205#[cfg(feature = "java-forge")]
2206fn prompt_card_ids(prompt: &Value, field: &str, count: usize) -> Result<Vec<String>, String> {
2207    let card_ids = prompt
2208        .get("input")
2209        .and_then(|input| input.get(field))
2210        .and_then(Value::as_array)
2211        .ok_or_else(|| format!("prompt is missing {field}"))?;
2212    if card_ids.len() < count {
2213        return Err(format!(
2214            "prompt field {field} has {} cards, need {count}",
2215            card_ids.len()
2216        ));
2217    }
2218    Ok(card_ids
2219        .iter()
2220        .take(count)
2221        .filter_map(Value::as_str)
2222        .map(str::to_string)
2223        .collect())
2224}
2225
2226#[cfg(feature = "java-forge")]
2227fn battlefield_contains(game_view: &GameViewDto, card_name: &str) -> bool {
2228    use manabrew_agent_interface::game_view_dto::{CardView, ZoneKind};
2229    game_view.zones.iter().any(|zone| {
2230        zone.zone == ZoneKind::Battlefield
2231            && zone.owner_id == "player-0"
2232            && zone.cards.iter().any(|card| match card {
2233                CardView::Visible(dto) => dto.identity.name == card_name,
2234                CardView::Hidden { .. } => false,
2235            })
2236    })
2237}
2238
2239pub trait JavaBridge {
2240    fn initialize(&mut self, assets_dir: &str) -> Result<(), String>;
2241    fn start_game_json(&mut self, request_json: &str) -> Result<String, String>;
2242    fn submit_action(&mut self, session_id: &str, action_json: &str) -> Result<String, String>;
2243    fn get_prompt(
2244        &mut self,
2245        session_id: &str,
2246        player_index: usize,
2247    ) -> Result<Option<String>, String>;
2248    fn get_snapshot(&mut self, session_id: &str, viewer: Option<usize>) -> Result<String, String>;
2249    fn is_game_over(&mut self, session_id: &str) -> Result<bool, String>;
2250    fn end_game(&mut self, session_id: &str) -> Result<(), String>;
2251    fn abort_game(&mut self, session_id: &str) -> Result<(), String>;
2252}
2253
2254pub struct JavaForgeSession<B> {
2255    bridge: B,
2256    session_id: Option<String>,
2257}
2258
2259impl<B: JavaBridge> JavaForgeSession<B> {
2260    pub fn new(bridge: B) -> Self {
2261        Self {
2262            bridge,
2263            session_id: None,
2264        }
2265    }
2266
2267    pub fn initialize(&mut self, assets_dir: &str) -> Result<(), String> {
2268        self.bridge.initialize(assets_dir)
2269    }
2270
2271    pub fn start_game(&mut self, request: &StartGameRequest) -> Result<String, String> {
2272        let request_json = request.to_json().map_err(|err| err.to_string())?;
2273        let response_json = self.bridge.start_game_json(&request_json)?;
2274        let response: StartGameResponse =
2275            serde_json::from_str(&response_json).map_err(|err| err.to_string())?;
2276        self.session_id = Some(response.session_id.clone());
2277        Ok(response.session_id)
2278    }
2279
2280    pub fn submit_action(&mut self, action_json: &str) -> Result<String, String> {
2281        let session_id = self.require_session_id()?.to_string();
2282        self.bridge.submit_action(&session_id, action_json)
2283    }
2284
2285    pub fn get_prompt(&mut self, player_index: usize) -> Result<Option<String>, String> {
2286        let session_id = self.require_session_id()?.to_string();
2287        self.bridge.get_prompt(&session_id, player_index)
2288    }
2289
2290    pub fn get_snapshot(&mut self, viewer: Option<usize>) -> Result<String, String> {
2291        let session_id = self.require_session_id()?.to_string();
2292        self.bridge.get_snapshot(&session_id, viewer)
2293    }
2294
2295    pub fn is_game_over(&mut self) -> Result<bool, String> {
2296        let session_id = self.require_session_id()?.to_string();
2297        self.bridge.is_game_over(&session_id)
2298    }
2299
2300    pub fn end_game(&mut self) -> Result<(), String> {
2301        let Some(session_id) = self.session_id.take() else {
2302            return Ok(());
2303        };
2304        self.bridge.end_game(&session_id)
2305    }
2306
2307    fn require_session_id(&self) -> Result<&str, String> {
2308        self.session_id
2309            .as_deref()
2310            .ok_or_else(|| "java-forge session has not started".to_string())
2311    }
2312}
2313
2314pub struct UnavailableJavaBridge;
2315
2316impl JavaBridge for UnavailableJavaBridge {
2317    fn initialize(&mut self, _assets_dir: &str) -> Result<(), String> {
2318        Err(unsupported_message().to_string())
2319    }
2320
2321    fn start_game_json(&mut self, _request_json: &str) -> Result<String, String> {
2322        Err(unsupported_message().to_string())
2323    }
2324
2325    fn submit_action(&mut self, _session_id: &str, _action_json: &str) -> Result<String, String> {
2326        Err(unsupported_message().to_string())
2327    }
2328
2329    fn get_prompt(
2330        &mut self,
2331        _session_id: &str,
2332        _player_index: usize,
2333    ) -> Result<Option<String>, String> {
2334        Err(unsupported_message().to_string())
2335    }
2336
2337    fn get_snapshot(
2338        &mut self,
2339        _session_id: &str,
2340        _viewer: Option<usize>,
2341    ) -> Result<String, String> {
2342        Err(unsupported_message().to_string())
2343    }
2344
2345    fn is_game_over(&mut self, _session_id: &str) -> Result<bool, String> {
2346        Err(unsupported_message().to_string())
2347    }
2348
2349    fn end_game(&mut self, _session_id: &str) -> Result<(), String> {
2350        Err(unsupported_message().to_string())
2351    }
2352
2353    fn abort_game(&mut self, _session_id: &str) -> Result<(), String> {
2354        Err(unsupported_message().to_string())
2355    }
2356}
2357
2358#[cfg(feature = "java-forge")]
2359#[derive(serde::Deserialize)]
2360struct SubprocessReply {
2361    ok: bool,
2362    #[serde(default)]
2363    result: String,
2364    #[serde(default)]
2365    error: Option<String>,
2366}
2367
2368#[cfg(feature = "java-forge")]
2369const CALL_TIMEOUT: Duration = Duration::from_secs(60);
2370#[cfg(feature = "java-forge")]
2371const SHUTDOWN_GRACE: Duration = Duration::from_secs(5);
2372
2373#[cfg(feature = "java-forge")]
2374pub struct SubprocessBridge {
2375    child: Child,
2376    stdin: BufWriter<ChildStdin>,
2377    stdout_rx: std_mpsc::Receiver<String>,
2378    stdout_handle: Option<std::thread::JoinHandle<()>>,
2379    stderr_handle: Option<std::thread::JoinHandle<()>>,
2380}
2381
2382#[cfg(feature = "java-forge")]
2383impl SubprocessBridge {
2384    fn spawn(config: &JavaRuntimeConfig) -> Result<Self, String> {
2385        config.validate()?;
2386
2387        let java_bin = resolve_java_bin(config);
2388        let jvm_args = config.jvm_args();
2389        info!(target: "self_hosted_node::java", args = %jvm_args.join(" "), "spawning java engine");
2390        let mut cmd = Command::new(&java_bin);
2391        // The JVM applies JAVA_TOOL_OPTIONS before anything on the command line,
2392        // so whatever is set on the host silently governs every flag we do not
2393        // name. Production ran -Xmx512m -XX:+UseSerialGC that way for months
2394        // with nothing in the repo to show it. Use SELF_HOSTED_NODE_JAVA_OPTS
2395        // instead: it is passed here, logged above, and lives in config.
2396        cmd.env_remove("JAVA_TOOL_OPTIONS");
2397        cmd.args(&jvm_args);
2398        cmd.arg("-jar").arg(&config.harness_jar);
2399        cmd.arg("--interactive-server");
2400        cmd.arg("--forge-home")
2401            .arg(format!("{}/", config.assets_dir.display()));
2402
2403        let mut child = cmd
2404            .stdin(Stdio::piped())
2405            .stdout(Stdio::piped())
2406            .stderr(Stdio::piped())
2407            .spawn()
2408            .map_err(|err| format!("failed to spawn java subprocess: {err}"))?;
2409
2410        let stdin = child
2411            .stdin
2412            .take()
2413            .ok_or_else(|| "java subprocess has no stdin".to_string())?;
2414        let stdout = child
2415            .stdout
2416            .take()
2417            .ok_or_else(|| "java subprocess has no stdout".to_string())?;
2418        let stderr = child.stderr.take();
2419
2420        // Bounded so a chatty Java side can't grow the queue without bound.
2421        // Cap is generous — protocol replies are one line per request.
2422        let (stdout_tx, stdout_rx) = std_mpsc::sync_channel::<String>(1024);
2423        let stdout_handle = std::thread::spawn(move || {
2424            let reader = BufReader::new(stdout);
2425            for line in reader.lines().map_while(Result::ok) {
2426                if stdout_tx.send(line).is_err() {
2427                    break;
2428                }
2429            }
2430        });
2431
2432        let stderr_handle = std::thread::spawn(move || {
2433            if let Some(stderr) = stderr {
2434                let reader = BufReader::new(stderr);
2435                for line in reader.lines().map_while(Result::ok) {
2436                    if let Some(pause) = parse_gc_pause(&line) {
2437                        crate::metrics::record_jvm_gc(
2438                            pause.kind,
2439                            Duration::from_secs_f64(pause.millis / 1000.0),
2440                            pause.heap_after_mb,
2441                        );
2442                    }
2443                    if line.contains("Exception") || line.contains("ERROR") {
2444                        warn!(target: "self_hosted_node::java", "[java] {line}");
2445                    } else if line.contains("][gc") {
2446                        // -Xlog:gc*:stderr is opt-in, and debug would drop it before Loki.
2447                        info!(target: "self_hosted_node::java", "[java] {line}");
2448                    } else {
2449                        debug!(target: "self_hosted_node::java", "[java] {line}");
2450                    }
2451                }
2452            }
2453        });
2454
2455        Ok(Self {
2456            child,
2457            stdin: BufWriter::new(stdin),
2458            stdout_rx,
2459            stdout_handle: Some(stdout_handle),
2460            stderr_handle: Some(stderr_handle),
2461        })
2462    }
2463
2464    fn call(&mut self, request_json: &str) -> Result<String, String> {
2465        // Drain anything still queued from a prior request — a previous call()
2466        // that timed out may have left its reply in the channel, and consuming
2467        // it now would shift every subsequent call off-by-one.
2468        loop {
2469            match self.stdout_rx.try_recv() {
2470                Ok(stale) => {
2471                    debug!(target: "self_hosted_node::java", line = %stale, "discarding stale stdout line");
2472                }
2473                Err(TryRecvError::Empty) => break,
2474                Err(TryRecvError::Disconnected) => {
2475                    return Err("java subprocess closed stdout (crashed?)".to_string());
2476                }
2477            }
2478        }
2479
2480        self.stdin
2481            .write_all(request_json.as_bytes())
2482            .map_err(|err| format!("failed to write subprocess stdin: {err}"))?;
2483        self.stdin
2484            .write_all(b"\n")
2485            .map_err(|err| format!("failed to write subprocess newline: {err}"))?;
2486        self.stdin
2487            .flush()
2488            .map_err(|err| format!("failed to flush subprocess stdin: {err}"))?;
2489
2490        let deadline = Instant::now() + CALL_TIMEOUT;
2491        loop {
2492            let remaining = deadline.saturating_duration_since(Instant::now());
2493            if remaining.is_zero() {
2494                return Err(format!(
2495                    "java subprocess timed out after {}s",
2496                    CALL_TIMEOUT.as_secs()
2497                ));
2498            }
2499            match self.stdout_rx.recv_timeout(remaining) {
2500                Ok(line) => {
2501                    let trimmed = line.trim();
2502                    if trimmed.is_empty() {
2503                        continue;
2504                    }
2505                    match serde_json::from_str::<SubprocessReply>(trimmed) {
2506                        Ok(reply) if reply.ok => return Ok(reply.result),
2507                        Ok(reply) => {
2508                            return Err(reply.error.unwrap_or_else(|| "unknown java error".into()));
2509                        }
2510                        Err(_) => {
2511                            debug!(target: "self_hosted_node::java", line = trimmed, "non-protocol stdout line");
2512                        }
2513                    }
2514                }
2515                Err(RecvTimeoutError::Timeout) => {
2516                    return Err(format!(
2517                        "java subprocess timed out after {}s",
2518                        CALL_TIMEOUT.as_secs()
2519                    ));
2520                }
2521                Err(RecvTimeoutError::Disconnected) => {
2522                    return Err("java subprocess closed stdout (crashed?)".to_string());
2523                }
2524            }
2525        }
2526    }
2527
2528    fn reset(&mut self) -> Result<(), String> {
2529        self.call("{\"command\":\"reset\"}").map(|_| ())
2530    }
2531
2532    fn is_alive(&mut self) -> bool {
2533        matches!(self.child.try_wait(), Ok(None))
2534    }
2535
2536    fn shutdown(mut self) {
2537        let _ = self.stdin.write_all(b"{\"command\":\"quit\"}\n");
2538        let _ = self.stdin.flush();
2539        let deadline = Instant::now() + SHUTDOWN_GRACE;
2540        loop {
2541            match self.child.try_wait() {
2542                Ok(Some(_)) => break,
2543                Ok(None) if Instant::now() >= deadline => {
2544                    let _ = self.child.kill();
2545                    let _ = self.child.wait();
2546                    break;
2547                }
2548                Ok(None) => std::thread::sleep(Duration::from_millis(100)),
2549                Err(_) => break,
2550            }
2551        }
2552        if let Some(handle) = self.stdout_handle.take() {
2553            let _ = handle.join();
2554        }
2555        if let Some(handle) = self.stderr_handle.take() {
2556            let _ = handle.join();
2557        }
2558    }
2559}
2560
2561#[cfg(feature = "java-forge")]
2562impl JavaBridge for SubprocessBridge {
2563    fn initialize(&mut self, _assets_dir: &str) -> Result<(), String> {
2564        Ok(())
2565    }
2566
2567    fn start_game_json(&mut self, request_json: &str) -> Result<String, String> {
2568        let body = json!({ "command": "startGame", "payload": request_json });
2569        self.call(&body.to_string())
2570    }
2571
2572    fn submit_action(&mut self, session_id: &str, action_json: &str) -> Result<String, String> {
2573        let body = json!({
2574            "command": "submitAction",
2575            "sessionId": session_id,
2576            "payload": action_json,
2577        });
2578        self.call(&body.to_string())
2579    }
2580
2581    fn get_prompt(
2582        &mut self,
2583        session_id: &str,
2584        player_index: usize,
2585    ) -> Result<Option<String>, String> {
2586        let body = json!({
2587            "command": "getPrompt",
2588            "sessionId": session_id,
2589            "playerIndex": player_index,
2590        });
2591        let prompt = self.call(&body.to_string())?;
2592        Ok((!prompt.is_empty()).then_some(prompt))
2593    }
2594
2595    fn get_snapshot(&mut self, session_id: &str, viewer: Option<usize>) -> Result<String, String> {
2596        let mut body = json!({ "command": "getSnapshot", "sessionId": session_id });
2597        if let Some(viewer) = viewer {
2598            body["viewer"] = json!(viewer);
2599        }
2600        self.call(&body.to_string())
2601    }
2602
2603    fn is_game_over(&mut self, session_id: &str) -> Result<bool, String> {
2604        let body = json!({ "command": "getGameOver", "sessionId": session_id });
2605        let value = self.call(&body.to_string())?;
2606        Ok(value.trim() == "true")
2607    }
2608
2609    fn end_game(&mut self, session_id: &str) -> Result<(), String> {
2610        let body = json!({ "command": "endGame", "sessionId": session_id });
2611        self.call(&body.to_string()).map(|_| ())
2612    }
2613
2614    fn abort_game(&mut self, session_id: &str) -> Result<(), String> {
2615        let body = json!({ "command": "abortGame", "sessionId": session_id });
2616        self.call(&body.to_string()).map(|_| ())
2617    }
2618}
2619
2620#[cfg(feature = "java-forge")]
2621impl Drop for SubprocessBridge {
2622    fn drop(&mut self) {
2623        let _ = self.child.kill();
2624        let _ = self.child.wait();
2625        if let Some(handle) = self.stdout_handle.take() {
2626            let _ = handle.join();
2627        }
2628        if let Some(handle) = self.stderr_handle.take() {
2629            let _ = handle.join();
2630        }
2631    }
2632}
2633
2634#[cfg(feature = "java-forge")]
2635fn resolve_java_bin(config: &JavaRuntimeConfig) -> String {
2636    if let Some(home) = &config.java_home {
2637        let bin = home.join("bin").join("java");
2638        if bin.is_file() {
2639            return bin.to_string_lossy().to_string();
2640        }
2641    }
2642    if let Ok(home) = env::var("JAVA_HOME") {
2643        let bin = PathBuf::from(home).join("bin").join("java");
2644        if bin.is_file() {
2645            return bin.to_string_lossy().to_string();
2646        }
2647    }
2648    "java".to_string()
2649}
2650
2651#[derive(Debug, Clone, Serialize)]
2652#[serde(rename_all = "camelCase")]
2653pub struct StartGameRequest {
2654    game_id: String,
2655    variant: String,
2656    starting_life: i32,
2657    seed: u64,
2658    players: Vec<PlayerConfig>,
2659}
2660
2661#[derive(Debug, Clone, Serialize)]
2662#[serde(rename_all = "camelCase")]
2663pub struct PlayerConfig {
2664    name: String,
2665    deck: Vec<CardIdentityForJava>,
2666    commander_names: Vec<String>,
2667    ai: bool,
2668}
2669
2670#[derive(Debug, Clone, Serialize)]
2671#[serde(rename_all = "camelCase")]
2672pub struct CardIdentityForJava {
2673    name: String,
2674    set_code: Option<String>,
2675}
2676
2677#[derive(Debug, serde::Deserialize)]
2678#[serde(rename_all = "camelCase")]
2679struct StartGameResponse {
2680    session_id: String,
2681    #[allow(dead_code)]
2682    player_indexes: Vec<usize>,
2683}
2684
2685impl StartGameRequest {
2686    pub fn new(
2687        game_id: String,
2688        variant: String,
2689        starting_life: i32,
2690        seed: u64,
2691        players: Vec<PlayerConfig>,
2692    ) -> Self {
2693        Self {
2694            game_id,
2695            variant,
2696            starting_life,
2697            seed,
2698            players,
2699        }
2700    }
2701
2702    pub fn to_json(&self) -> Result<String, serde_json::Error> {
2703        serde_json::to_string(self)
2704    }
2705}
2706
2707impl PlayerConfig {
2708    pub fn new(name: String, deck: &[DeckCardIdentity], commander_names: Vec<String>) -> Self {
2709        Self {
2710            name,
2711            deck: deck.iter().map(CardIdentityForJava::from).collect(),
2712            commander_names,
2713            ai: false,
2714        }
2715    }
2716}
2717
2718impl From<&DeckCardIdentity> for CardIdentityForJava {
2719    fn from(identity: &DeckCardIdentity) -> Self {
2720        Self {
2721            name: java_card_name(&identity.name),
2722            set_code: (!identity.set_code.is_empty()).then(|| identity.set_code.clone()),
2723        }
2724    }
2725}
2726
2727fn java_card_name(name: &str) -> String {
2728    name.split_once(" // ")
2729        .map(|(front, _)| front.to_string())
2730        .unwrap_or_else(|| name.to_string())
2731}
2732
2733fn env_path(key: &str) -> Option<PathBuf> {
2734    env::var_os(key)
2735        .filter(|value| !value.is_empty())
2736        .map(PathBuf::from)
2737}
2738
2739fn env_sizing(key: &str, default: u64) -> Option<u64> {
2740    let value = match env::var(key) {
2741        Ok(raw) if !raw.trim().is_empty() => match raw.trim().parse::<u64>() {
2742            Ok(parsed) => parsed,
2743            Err(_) => {
2744                tracing::warn!(target: "self_hosted_node::java", key, raw, "ignoring unparseable jvm sizing");
2745                default
2746            }
2747        },
2748        _ => default,
2749    };
2750    (value > 0).then_some(value)
2751}
2752
2753fn env_classpath(key: &str) -> Vec<PathBuf> {
2754    let Some(value) = env::var_os(key) else {
2755        return Vec::new();
2756    };
2757    env::split_paths(&value).collect()
2758}
2759
2760fn require_dir(path: &Path, label: &str) -> Result<(), String> {
2761    if path.is_dir() {
2762        Ok(())
2763    } else {
2764        Err(format!("{label} does not exist: {}", path.display()))
2765    }
2766}
2767
2768fn require_file(path: &Path, label: &str) -> Result<(), String> {
2769    if path.is_file() {
2770        Ok(())
2771    } else {
2772        Err(format!("{label} does not exist: {}", path.display()))
2773    }
2774}
2775
2776/// One `Pause` line from the JVM's unified GC log.
2777struct GcPause {
2778    kind: &'static str,
2779    millis: f64,
2780    heap_after_mb: Option<u64>,
2781}
2782
2783/// Parse a unified-log GC pause, e.g.
2784/// `[12.3s][info][gc] GC(7) Pause Full (Allocation Failure) 240M->171M(247M) 226.856ms`
2785fn parse_gc_pause(line: &str) -> Option<GcPause> {
2786    let kind = if line.contains("Pause Full") {
2787        "full"
2788    } else if line.contains("Pause Young") {
2789        "young"
2790    } else {
2791        return None;
2792    };
2793    let millis = line
2794        .rsplit(' ')
2795        .find_map(|token| token.strip_suffix("ms")?.parse::<f64>().ok())?;
2796    // "240M->171M(247M)": the figure after the arrow is what survived.
2797    let heap_after_mb = line.split_once("->").and_then(|(_, rest)| {
2798        let digits: String = rest.chars().take_while(char::is_ascii_digit).collect();
2799        digits.parse::<u64>().ok()
2800    });
2801    Some(GcPause {
2802        kind,
2803        millis,
2804        heap_after_mb,
2805    })
2806}
2807
2808#[cfg(test)]
2809mod gc_log_tests {
2810    use super::parse_gc_pause;
2811
2812    #[test]
2813    fn reads_a_full_collection() {
2814        let pause = parse_gc_pause(
2815            "[117.6s][info][gc] GC(21) Pause Full (Allocation Failure) 240M->171M(247M) 226.856ms",
2816        )
2817        .expect("parsed");
2818        assert_eq!(pause.kind, "full");
2819        assert_eq!(pause.heap_after_mb, Some(171));
2820        assert!((pause.millis - 226.856).abs() < f64::EPSILON);
2821    }
2822
2823    #[test]
2824    fn reads_a_young_collection() {
2825        let pause = parse_gc_pause(
2826            "[3.1s][info][gc] GC(2) Pause Young (Allocation Failure) 216M->91M(494M) 17.254ms",
2827        )
2828        .expect("parsed");
2829        assert_eq!(pause.kind, "young");
2830        assert_eq!(pause.heap_after_mb, Some(91));
2831    }
2832
2833    /// The node asks for `-Xlog:gc*:stderr:time,uptime,level,tags`, which is a
2834    /// richer line than plain `-Xlog:gc`. Real production output.
2835    #[test]
2836    fn reads_the_format_the_node_actually_requests() {
2837        let pause = parse_gc_pause(
2838            "[2026-08-18T20:41:42.907+0000][120025.819s][info][gc] GC(1202) Pause Full (Allocation Failure) 494M->487M(494M) 545.344ms",
2839        )
2840        .expect("parsed");
2841        assert_eq!(pause.kind, "full");
2842        assert_eq!(pause.heap_after_mb, Some(487));
2843        assert!((pause.millis - 545.344).abs() < f64::EPSILON);
2844    }
2845
2846    /// `gc*` also emits region and metaspace lines that carry `->`; they must
2847    /// not be mistaken for pauses.
2848    #[test]
2849    fn ignores_the_extra_lines_gc_star_adds() {
2850        assert!(
2851            parse_gc_pause("[120025.819s][info][gc,heap] GC(1202) Eden regions: 10->0(12)")
2852                .is_none()
2853        );
2854        assert!(
2855            parse_gc_pause("[120025.819s][info][gc,metaspace] Metaspace: 48M->48M(1088M)")
2856                .is_none()
2857        );
2858    }
2859
2860    #[test]
2861    fn ignores_other_output() {
2862        assert!(parse_gc_pause("[java] LOGGER ERROR: something").is_none());
2863        assert!(parse_gc_pause("[1.0s][info][gc,init] CardTable entry size: 512").is_none());
2864    }
2865}