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(forge_backend)]
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        pub fn forge_dump_heap(
716            thread: *mut graal_isolatethread_t,
717            path: *const c_char,
718            live_objects_only: c_int,
719        ) -> *mut c_char;
720        pub fn forge_gc_stats(thread: *mut graal_isolatethread_t) -> *mut c_char;
721    }
722}
723
724/// Isolate-wide GC and heap counters, as reported by the engine itself.
725/// `collections` and `pause_millis` are -1 when the native image exposes no
726/// collector beans; the heap figures are always present.
727#[cfg(feature = "graal-forge")]
728#[derive(serde::Deserialize)]
729#[serde(rename_all = "camelCase")]
730pub struct GcStats {
731    pub heap_used_bytes: u64,
732    pub heap_total_bytes: u64,
733    pub heap_max_bytes: u64,
734    pub collections: i64,
735    pub pause_millis: i64,
736    /// Keyed by collector name. An incremental collection is a few milliseconds;
737    /// a complete one traces the whole live set, which here is a permanently
738    /// reachable card database, so it costs the same whether it reclaims
739    /// anything or not. The totals above hide that difference.
740    #[serde(default)]
741    pub by_collector: std::collections::HashMap<String, GcCollector>,
742}
743
744#[cfg(feature = "graal-forge")]
745#[derive(serde::Deserialize)]
746#[serde(rename_all = "camelCase")]
747pub struct GcCollector {
748    pub collections: u64,
749    pub pause_millis: u64,
750}
751
752#[cfg(feature = "graal-forge")]
753#[derive(serde::Deserialize)]
754struct ForgeReply {
755    ok: bool,
756    #[serde(default)]
757    result: String,
758    #[serde(default)]
759    error: Option<String>,
760}
761
762// A GraalVM isolate hosts the in-process Forge engine. The `thread` handle is
763// bound to the hosted-engine thread that created or attached it (isolate
764// threads are not portable), so the handle is Rc/!Send by design. In shared
765// mode one isolate serves all rooms: the first creator initializes Forge, later
766// rooms attach their own thread to it and sessions coexist in the adapter.
767#[cfg(feature = "graal-forge")]
768struct GraalBridge {
769    thread: *mut graal_ffi::graal_isolatethread_t,
770    attached: bool,
771}
772
773#[cfg(feature = "graal-forge")]
774struct SharedIsolate(*mut graal_ffi::graal_isolate_t);
775#[cfg(feature = "graal-forge")]
776unsafe impl Send for SharedIsolate {}
777
778#[cfg(feature = "graal-forge")]
779static SHARED_GRAAL_ISOLATE: std::sync::Mutex<Option<SharedIsolate>> = std::sync::Mutex::new(None);
780
781// One isolate hosts every room's game (safe since the endstep concurrency
782// patches), so the card database loads once per node. Off keeps the historical
783// isolate-per-game shape.
784#[cfg(feature = "graal-forge")]
785fn shared_isolate_enabled() -> bool {
786    env::var("SELF_HOSTED_NODE_SHARED_ISOLATE")
787        .map(|value| value == "1" || value.eq_ignore_ascii_case("true"))
788        .unwrap_or(false)
789}
790
791#[cfg(feature = "graal-forge")]
792impl GraalBridge {
793    fn create() -> Result<Self, String> {
794        let mut isolate: *mut graal_ffi::graal_isolate_t = std::ptr::null_mut();
795        let mut thread: *mut graal_ffi::graal_isolatethread_t = std::ptr::null_mut();
796        let rc = unsafe {
797            graal_ffi::graal_create_isolate(std::ptr::null_mut(), &mut isolate, &mut thread)
798        };
799        if rc != 0 {
800            return Err(format!("graal_create_isolate failed with code {rc}"));
801        }
802        Ok(Self {
803            thread,
804            attached: false,
805        })
806    }
807
808    fn create_in_shared_isolate(assets_dir: &Path) -> Result<Self, String> {
809        let mut guard = SHARED_GRAAL_ISOLATE
810            .lock()
811            .map_err(|_| "shared graal isolate poisoned".to_string())?;
812        if let Some(shared) = guard.as_ref() {
813            let mut thread: *mut graal_ffi::graal_isolatethread_t = std::ptr::null_mut();
814            let rc = unsafe { graal_ffi::graal_attach_thread(shared.0, &mut thread) };
815            if rc != 0 {
816                return Err(format!("graal_attach_thread failed with code {rc}"));
817            }
818            return Ok(Self {
819                thread,
820                attached: true,
821            });
822        }
823        let mut isolate: *mut graal_ffi::graal_isolate_t = std::ptr::null_mut();
824        let mut thread: *mut graal_ffi::graal_isolatethread_t = std::ptr::null_mut();
825        let rc = unsafe {
826            graal_ffi::graal_create_isolate(std::ptr::null_mut(), &mut isolate, &mut thread)
827        };
828        if rc != 0 {
829            return Err(format!("graal_create_isolate failed with code {rc}"));
830        }
831        let mut bridge = Self {
832            thread,
833            attached: false,
834        };
835        let assets = cstring(&assets_dir.to_string_lossy())?;
836        bridge.decode(unsafe { graal_ffi::forge_initialize(bridge.thread, assets.as_ptr()) })?;
837        bridge.attached = true;
838        *guard = Some(SharedIsolate(isolate));
839        info!("shared graal isolate initialized");
840        Ok(bridge)
841    }
842
843    fn decode(&self, raw: *mut std::os::raw::c_char) -> Result<String, String> {
844        if raw.is_null() {
845            return Err("forge native lib returned null".to_string());
846        }
847        let envelope = unsafe { std::ffi::CStr::from_ptr(raw) }
848            .to_string_lossy()
849            .into_owned();
850        unsafe { graal_ffi::forge_free_string(self.thread, raw) };
851        let reply: ForgeReply = serde_json::from_str(&envelope)
852            .map_err(|err| format!("malformed forge envelope: {err}"))?;
853        if reply.ok {
854            Ok(reply.result)
855        } else {
856            Err(reply
857                .error
858                .unwrap_or_else(|| "unknown forge error".to_string()))
859        }
860    }
861}
862
863#[cfg(feature = "graal-forge")]
864impl Drop for GraalBridge {
865    fn drop(&mut self) {
866        if self.attached {
867            unsafe { graal_ffi::graal_detach_thread(self.thread) };
868        } else {
869            unsafe { graal_ffi::graal_tear_down_isolate(self.thread) };
870        }
871    }
872}
873
874#[cfg(feature = "graal-forge")]
875#[derive(Clone)]
876struct GraalEngineHandle {
877    bridge: std::rc::Rc<GraalBridge>,
878}
879
880#[cfg(feature = "graal-forge")]
881impl GraalEngineHandle {
882    fn create(assets_dir: &Path) -> Result<Self, String> {
883        let bridge = if shared_isolate_enabled() {
884            GraalBridge::create_in_shared_isolate(assets_dir)?
885        } else {
886            let bridge = GraalBridge::create()?;
887            let assets = cstring(&assets_dir.to_string_lossy())?;
888            bridge
889                .decode(unsafe { graal_ffi::forge_initialize(bridge.thread, assets.as_ptr()) })?;
890            bridge
891        };
892        Ok(Self {
893            bridge: std::rc::Rc::new(bridge),
894        })
895    }
896
897    fn start_game(&self, request_json: &str) -> Result<String, String> {
898        let request = cstring(request_json)?;
899        let response = self
900            .bridge
901            .decode(unsafe { graal_ffi::forge_start_game(self.bridge.thread, request.as_ptr()) })?;
902        let parsed: StartGameResponse = serde_json::from_str(&response)
903            .map_err(|err| format!("malformed startGame response: {err}"))?;
904        Ok(parsed.session_id)
905    }
906
907    fn submit_action(&self, session_id: &str, action_json: &str) -> Result<String, String> {
908        let session = cstring(session_id)?;
909        let action = cstring(action_json)?;
910        let started = Instant::now();
911        let result = self.bridge.decode(unsafe {
912            graal_ffi::forge_submit_action(self.bridge.thread, session.as_ptr(), action.as_ptr())
913        });
914        crate::metrics::record_forge_decision_stage("submit_action", started.elapsed());
915        result
916    }
917
918    /// Cumulative collections, pause time and heap for the whole isolate. Every
919    /// room on this node shares it, so a collection here stops all of them.
920    fn gc_stats(&self) -> Result<GcStats, String> {
921        let raw = self
922            .bridge
923            .decode(unsafe { graal_ffi::forge_gc_stats(self.bridge.thread) })?;
924        serde_json::from_str(&raw).map_err(|err| format!("malformed gcStats response: {err}"))
925    }
926
927    fn get_prompt(&self, session_id: &str, player_index: usize) -> Result<Option<String>, String> {
928        let session = cstring(session_id)?;
929        let prompt = self.bridge.decode(unsafe {
930            graal_ffi::forge_get_prompt(
931                self.bridge.thread,
932                session.as_ptr(),
933                player_index as std::os::raw::c_int,
934            )
935        })?;
936        Ok((!prompt.is_empty()).then_some(prompt))
937    }
938
939    fn is_game_over(&self, session_id: &str) -> Result<bool, String> {
940        let session = cstring(session_id)?;
941        let value = self.bridge.decode(unsafe {
942            graal_ffi::forge_get_game_over(self.bridge.thread, session.as_ptr())
943        })?;
944        Ok(value.trim() == "true")
945    }
946
947    fn get_snapshot(&self, session_id: &str, viewer: Option<usize>) -> Result<String, String> {
948        let session = cstring(session_id)?;
949        let viewer = viewer.map_or(-1, |v| v as std::os::raw::c_int);
950        self.bridge.decode(unsafe {
951            graal_ffi::forge_get_snapshot(self.bridge.thread, session.as_ptr(), viewer)
952        })
953    }
954
955    fn end_game(&self, session_id: &str) -> Result<(), String> {
956        let session = cstring(session_id)?;
957        self.bridge
958            .decode(unsafe { graal_ffi::forge_end_game(self.bridge.thread, session.as_ptr()) })
959            .map(|_| ())
960    }
961
962    fn abort_game(&self, session_id: &str) -> Result<(), String> {
963        let session = cstring(session_id)?;
964        self.bridge
965            .decode(unsafe { graal_ffi::forge_abort_game(self.bridge.thread, session.as_ptr()) })
966            .map(|_| ())
967    }
968}
969
970/// Write an hprof dump of the shared isolate's heap. Attaches its own isolate
971/// thread, so it is callable from a signal handler task rather than the engine
972/// thread that owns the room's bridge.
973#[cfg(feature = "graal-forge")]
974pub fn dump_shared_heap(path: &str) -> Result<(), String> {
975    let isolate = SHARED_GRAAL_ISOLATE
976        .lock()
977        .map_err(|_| "shared graal isolate poisoned".to_string())?
978        .as_ref()
979        .map(|shared| shared.0)
980        .ok_or_else(|| "no shared isolate to dump".to_string())?;
981    let mut thread: *mut graal_ffi::graal_isolatethread_t = std::ptr::null_mut();
982    let rc = unsafe { graal_ffi::graal_attach_thread(isolate, &mut thread) };
983    if rc != 0 {
984        return Err(format!("graal_attach_thread failed with code {rc}"));
985    }
986    let bridge = GraalBridge {
987        thread,
988        attached: true,
989    };
990    let target = cstring(path)?;
991    bridge
992        .decode(unsafe { graal_ffi::forge_dump_heap(bridge.thread, target.as_ptr(), 1) })
993        .map(|_| ())
994}
995
996#[cfg(feature = "graal-forge")]
997fn cstring(value: &str) -> Result<std::ffi::CString, String> {
998    std::ffi::CString::new(value).map_err(|_| "string contained interior NUL".to_string())
999}
1000
1001#[cfg(feature = "java-forge")]
1002pub fn run_concurrent_self_play(
1003    seats: &[DeckSelection],
1004    starting_life: i32,
1005    seed: u64,
1006    max_prompts: usize,
1007    concurrency: usize,
1008) -> Result<(), String> {
1009    let config = JavaRuntimeConfig::from_env();
1010    let games_per_process = env::var("SELF_HOSTED_NODE_GAMES_PER_JVM")
1011        .ok()
1012        .and_then(|value| value.parse::<usize>().ok())
1013        .filter(|n| *n >= 1)
1014        .unwrap_or(1);
1015    let pool = JavaEnginePool::start(&config, concurrency.max(1), games_per_process)?;
1016    info!(
1017        concurrency,
1018        "java-engine started; launching concurrent games"
1019    );
1020
1021    let mut players = Vec::with_capacity(seats.len());
1022    for (i, seat) in seats.iter().enumerate() {
1023        let identities = deck_card_identities(&seat.deck);
1024        players.push(PlayerConfig::new(
1025            format!("Self-Play {}", i + 1),
1026            &identities,
1027            commander_names_for_java(&seat.deck, seat.commander_name.as_deref()),
1028        ));
1029    }
1030
1031    let mut joins = Vec::with_capacity(concurrency.max(1));
1032    for game_index in 0..concurrency.max(1) {
1033        let handle = pool.handle();
1034        let request = StartGameRequest::new(
1035            format!("self-hosted-java-concurrent-{game_index}"),
1036            String::new(),
1037            starting_life,
1038            seed.wrapping_add(game_index as u64),
1039            players.clone(),
1040        );
1041        joins.push(std::thread::spawn(move || -> Result<(), String> {
1042            let request_json = request.to_json().map_err(|error| error.to_string())?;
1043            let session_id = handle.start_game(&request_json)?;
1044            info!(session_id, game_index, "concurrent java game started");
1045            let result = drive_game_via_handle(&handle, &session_id, max_prompts);
1046            let _ = handle.end_game(&session_id);
1047            result
1048        }));
1049    }
1050
1051    let mut outcome = Ok(());
1052    for join in joins {
1053        match join.join() {
1054            Ok(Ok(())) => {}
1055            Ok(Err(error)) => outcome = Err(error),
1056            Err(_) => outcome = Err("concurrent game thread panicked".to_string()),
1057        }
1058    }
1059    outcome
1060}
1061
1062#[cfg(not(feature = "java-forge"))]
1063pub fn run_concurrent_self_play(
1064    _seats: &[DeckSelection],
1065    _starting_life: i32,
1066    _seed: u64,
1067    _max_prompts: usize,
1068    _concurrency: usize,
1069) -> Result<(), String> {
1070    Err(
1071        "java-forge concurrent self-play requires building self-hosted-node with --features java-forge"
1072            .to_string(),
1073    )
1074}
1075
1076/// Sampled once per completed decision, which is the cadence the stalls happen
1077/// at and cheap enough not to need its own timer.
1078#[cfg(all(feature = "graal-forge", not(feature = "java-forge")))]
1079fn report_engine_gc(engine: &ForgeEngine) {
1080    match engine.gc_stats() {
1081        Ok(stats) => {
1082            crate::metrics::record_engine_gc(
1083                stats.collections,
1084                stats.pause_millis,
1085                stats.heap_used_bytes,
1086                stats.heap_max_bytes,
1087            );
1088            for (name, per) in &stats.by_collector {
1089                crate::metrics::record_engine_gc_collector(name, per.collections, per.pause_millis);
1090            }
1091        }
1092        Err(err) => debug!(%err, "engine gc stats unavailable"),
1093    }
1094}
1095
1096/// The subprocess backend already reports this by parsing its own GC log.
1097#[cfg(feature = "java-forge")]
1098fn report_engine_gc(_engine: &ForgeEngine) {}
1099
1100#[cfg(feature = "java-forge")]
1101fn drive_game_via_handle(
1102    handle: &JavaEngineHandle,
1103    session_id: &str,
1104    max_prompts: usize,
1105) -> Result<(), String> {
1106    let mut bots: HashMap<usize, SimpleAi> = HashMap::new();
1107    let mut last_prompt: Option<String> = None;
1108    let mut acted = 0usize;
1109    let mut seen_prompt = false;
1110    let max_iterations = max_prompts.saturating_mul(200).max(2_000);
1111
1112    for _ in 0..max_iterations {
1113        if let Some(prompt_json) = handle.get_prompt(session_id, 0)? {
1114            seen_prompt = true;
1115            if last_prompt.as_deref() == Some(prompt_json.as_str()) {
1116                if handle.is_game_over(session_id)? {
1117                    return Ok(());
1118                }
1119                std::thread::sleep(Duration::from_millis(20));
1120                continue;
1121            }
1122            let prompt: AgentPrompt = serde_json::from_str(&prompt_json)
1123                .map_err(|error| format!("failed to parse concurrent prompt: {error}"))?;
1124            let player = player_index(&prompt.deciding_player_id);
1125            if let Some(action) = bots.entry(player).or_default().decide(prompt) {
1126                let action_json = serde_json::to_string(&action).map_err(|err| err.to_string())?;
1127                handle.submit_action(session_id, &action_json)?;
1128                acted += 1;
1129                if acted >= max_prompts {
1130                    return Err(format!(
1131                        "concurrent game {session_id} did not finish within {max_prompts} decisions"
1132                    ));
1133                }
1134            }
1135            last_prompt = Some(prompt_json);
1136            continue;
1137        }
1138        if seen_prompt && handle.is_game_over(session_id)? {
1139            return Ok(());
1140        }
1141        std::thread::sleep(Duration::from_millis(20));
1142    }
1143    Err(format!(
1144        "concurrent game {session_id} exceeded its iteration cap"
1145    ))
1146}
1147
1148/// G1 keeps pauses roughly bounded as the live set grows; Serial does not.
1149const DEFAULT_JAVA_COLLECTOR: &str = "G1";
1150const DEFAULT_JAVA_GC_LOG: &str = "stderr";
1151const DEFAULT_JAVA_HEAP_MB: u64 = 1024;
1152const DEFAULT_JAVA_ACTIVE_PROCESSORS: u64 = 2;
1153
1154#[derive(Debug, Clone)]
1155pub struct JavaRuntimeConfig {
1156    pub assets_dir: PathBuf,
1157    pub harness_jar: PathBuf,
1158    pub java_home: Option<PathBuf>,
1159    pub extra_classpath: Vec<PathBuf>,
1160    pub heap_mb: Option<u64>,
1161    pub active_processor_count: Option<u64>,
1162    pub gc_log: Option<String>,
1163    pub collector: Option<String>,
1164    pub extra_jvm_args: Vec<String>,
1165}
1166
1167impl JavaRuntimeConfig {
1168    pub fn from_env() -> Self {
1169        let root = workspace_root();
1170        Self {
1171            assets_dir: env_path("SELF_HOSTED_NODE_FORGE_ASSETS_DIR")
1172                .or_else(|| env_path("MANA_BREW_FORGE_ASSETS_DIR"))
1173                .unwrap_or_else(|| root.join("forge/forge-gui")),
1174            harness_jar: env_path("SELF_HOSTED_NODE_FORGE_HARNESS_JAR")
1175                .or_else(|| env_path("MANA_BREW_FORGE_HARNESS_JAR"))
1176                .unwrap_or_else(|| {
1177                    root.join("forge-harness/target/forge-harness-jar-with-dependencies.jar")
1178                }),
1179            java_home: env_path("SELF_HOSTED_NODE_JAVA_HOME")
1180                .or_else(|| env_path("MANA_BREW_JAVA_HOME"))
1181                .or_else(|| env_path("JAVA_HOME")),
1182            extra_classpath: env_classpath("SELF_HOSTED_NODE_FORGE_EXTRA_CLASSPATH")
1183                .into_iter()
1184                .chain(env_classpath("MANA_BREW_FORGE_EXTRA_CLASSPATH"))
1185                .collect(),
1186            heap_mb: env_sizing("SELF_HOSTED_NODE_JAVA_HEAP_MB", DEFAULT_JAVA_HEAP_MB),
1187            active_processor_count: env_sizing(
1188                "SELF_HOSTED_NODE_JAVA_ACTIVE_PROCESSORS",
1189                DEFAULT_JAVA_ACTIVE_PROCESSORS,
1190            ),
1191            collector: env::var("SELF_HOSTED_NODE_JAVA_COLLECTOR")
1192                .ok()
1193                .map(|value| value.trim().to_string())
1194                .filter(|value| !value.is_empty())
1195                .or_else(|| Some(DEFAULT_JAVA_COLLECTOR.to_string())),
1196            // Defaults on: the fleet ships no logs, so this log exists to be
1197            // turned into metrics by the stderr pump (see metrics.rs).
1198            gc_log: env::var("SELF_HOSTED_NODE_JAVA_GC_LOG")
1199                .ok()
1200                .map(|value| value.trim().to_string())
1201                .filter(|value| !value.is_empty())
1202                .or_else(|| Some(DEFAULT_JAVA_GC_LOG.to_string())),
1203            extra_jvm_args: env::var("SELF_HOSTED_NODE_JAVA_OPTS")
1204                .unwrap_or_default()
1205                .split_whitespace()
1206                .map(str::to_string)
1207                .collect(),
1208        }
1209    }
1210
1211    #[cfg(feature = "java-forge")]
1212    fn jvm_args(&self) -> Vec<String> {
1213        let mut args = vec![
1214            "-Dfile.encoding=UTF-8".to_string(),
1215            "-Dsun.stdout.encoding=UTF-8".to_string(),
1216            "-Dsun.stderr.encoding=UTF-8".to_string(),
1217            "-Djava.awt.headless=true".to_string(),
1218            // Forge's endstep concurrency patches assume a heap exhaustion kills the process
1219            // instead of thrashing; a supervised node is restarted, a thrashing one is not.
1220            "-XX:+ExitOnOutOfMemoryError".to_string(),
1221            // Forge calls System.gc() on match boundaries (Match.java,
1222            // HostedMatch.java) and so does the harness. Each one is a full
1223            // collection, and a full collection costs about 1.3ms per MB of
1224            // live set, so on a large board they are seconds of stop-the-world
1225            // for no benefit the collector would not have reached anyway.
1226            "-XX:+DisableExplicitGC".to_string(),
1227        ];
1228        // Pause time is linear in the live set and Serial collects on one
1229        // thread, so a bigger heap under Serial means longer freezes, not
1230        // shorter ones. Name the collector rather than inheriting whatever the
1231        // host's environment happens to set.
1232        if let Some(collector) = &self.collector {
1233            args.push(format!("-XX:+Use{collector}GC"));
1234        }
1235        // A fleet runs many of these side by side, and every JVM sizes its heap and its GC
1236        // thread count from the whole machine unless told otherwise.
1237        if let Some(heap_mb) = self.heap_mb {
1238            args.push(format!("-Xmx{heap_mb}m"));
1239        }
1240        if let Some(processors) = self.active_processor_count {
1241            args.push(format!("-XX:ActiveProcessorCount={processors}"));
1242        }
1243        // In a container the only log that reaches Loki is the node's own, and the
1244        // subprocess owns stdout for the protocol, so "stderr" routes GC through the
1245        // stderr pump instead of a file nothing ships.
1246        match self.gc_log.as_deref() {
1247            None => {}
1248            Some("stderr") => args.push("-Xlog:gc*:stderr:time,uptime,level,tags".to_string()),
1249            Some(dir) => args.push(format!(
1250                "-Xlog:gc*:file={}:time,uptime,level,tags:filecount=5,filesize=20M",
1251                Path::new(dir).join("engine-gc-%p.log").display()
1252            )),
1253        }
1254        args.extend(self.extra_jvm_args.iter().cloned());
1255        args
1256    }
1257
1258    pub fn validate(&self) -> Result<(), String> {
1259        require_dir(&self.assets_dir, "Forge assets directory")?;
1260        require_file(&self.harness_jar, "Forge harness jar")?;
1261        if let Some(java_home) = &self.java_home {
1262            require_dir(java_home, "Java home")?;
1263        }
1264        for entry in &self.extra_classpath {
1265            if !entry.exists() {
1266                return Err(format!(
1267                    "Classpath entry does not exist: {}",
1268                    entry.display()
1269                ));
1270            }
1271        }
1272        Ok(())
1273    }
1274
1275    pub fn classpath_entries(&self) -> Vec<PathBuf> {
1276        let mut entries = Vec::with_capacity(1 + self.extra_classpath.len());
1277        entries.push(self.harness_jar.clone());
1278        entries.extend(self.extra_classpath.iter().cloned());
1279        entries
1280    }
1281}
1282
1283#[cfg(forge_backend)]
1284#[allow(clippy::too_many_arguments)]
1285pub fn run_hosted_engine_game(
1286    game_id: String,
1287    player_names: Vec<String>,
1288    decks: Vec<Deck>,
1289    commander_names: Vec<Option<String>>,
1290    commander_variant: bool,
1291    game_variant: String,
1292    local_player_index: Option<usize>,
1293    ai_player_indices: Vec<usize>,
1294    starting_life: i32,
1295    remote_prompt_tx: std_mpsc::Sender<(usize, AgentMessage)>,
1296    remote_response_rxs: Vec<(usize, std_mpsc::Receiver<ClientToServerMessage>)>,
1297    game_over_tx: std_mpsc::Sender<HostedGameOver>,
1298    cancel: Arc<AtomicBool>,
1299) -> Result<(), String> {
1300    run_hosted_engine_game_inner(
1301        game_id,
1302        player_names,
1303        decks,
1304        commander_names,
1305        commander_variant,
1306        game_variant,
1307        local_player_index,
1308        ai_player_indices,
1309        starting_life,
1310        remote_prompt_tx,
1311        remote_response_rxs,
1312        game_over_tx,
1313        cancel,
1314    )
1315}
1316
1317#[cfg(not(forge_backend))]
1318#[allow(clippy::too_many_arguments)]
1319pub fn run_hosted_engine_game(
1320    _game_id: String,
1321    _player_names: Vec<String>,
1322    _decks: Vec<Deck>,
1323    _commander_names: Vec<Option<String>>,
1324    _commander_variant: bool,
1325    _game_variant: String,
1326    _local_player_index: Option<usize>,
1327    _ai_player_indices: Vec<usize>,
1328    _starting_life: i32,
1329    _remote_prompt_tx: std_mpsc::Sender<(usize, AgentMessage)>,
1330    _remote_response_rxs: Vec<(usize, std_mpsc::Receiver<ClientToServerMessage>)>,
1331    _game_over_tx: std_mpsc::Sender<HostedGameOver>,
1332    _cancel: Arc<AtomicBool>,
1333) -> Result<(), String> {
1334    Err(unsupported_message().to_string())
1335}
1336
1337#[cfg(forge_backend)]
1338#[allow(clippy::too_many_arguments)]
1339fn run_hosted_engine_game_inner(
1340    game_id: String,
1341    player_names: Vec<String>,
1342    decks: Vec<Deck>,
1343    commander_names: Vec<Option<String>>,
1344    commander_variant: bool,
1345    game_variant: String,
1346    local_player_index: Option<usize>,
1347    ai_player_indices: Vec<usize>,
1348    starting_life: i32,
1349    remote_prompt_tx: std_mpsc::Sender<(usize, AgentMessage)>,
1350    remote_response_rxs: Vec<(usize, std_mpsc::Receiver<ClientToServerMessage>)>,
1351    game_over_tx: std_mpsc::Sender<HostedGameOver>,
1352    cancel: Arc<AtomicBool>,
1353) -> Result<(), String> {
1354    let engine = obtain_engine()?;
1355
1356    let mut players = Vec::with_capacity(player_names.len());
1357    for (index, name) in player_names.iter().enumerate() {
1358        let identities = deck_card_identities(&decks[index]);
1359        let seat_commander_names = if commander_variant {
1360            commander_names_for_java(&decks[index], commander_names[index].as_deref())
1361        } else {
1362            Vec::new()
1363        };
1364        players.push(PlayerConfig::new(
1365            name.clone(),
1366            &identities,
1367            seat_commander_names,
1368        ));
1369    }
1370    for &idx in &ai_player_indices {
1371        if let Some(player) = players.get_mut(idx) {
1372            player.ai = true;
1373        }
1374    }
1375    let request = StartGameRequest::new(
1376        game_id.clone(),
1377        game_variant,
1378        starting_life,
1379        rand::random(),
1380        players,
1381    );
1382    let session_id = engine.start_game(&request.to_json().map_err(|err| err.to_string())?)?;
1383    info!(game_id, session_id, "hosted java-forge session started");
1384
1385    struct SessionGuard {
1386        engine: ForgeEngine,
1387        session_id: String,
1388        armed: std::cell::Cell<bool>,
1389    }
1390    impl Drop for SessionGuard {
1391        fn drop(&mut self) {
1392            if self.armed.get() {
1393                if let Err(error) = self.engine.abort_game(&self.session_id) {
1394                    warn!(session_id = %self.session_id, %error, "failed to abort java session; context may leak");
1395                }
1396            }
1397        }
1398    }
1399    let guard = SessionGuard {
1400        engine: engine.clone(),
1401        session_id: session_id.clone(),
1402        armed: std::cell::Cell::new(true),
1403    };
1404
1405    let mut remote_response_rxs: HashMap<usize, std_mpsc::Receiver<ClientToServerMessage>> =
1406        remote_response_rxs.into_iter().collect();
1407    let mut last_prompt: Option<AgentPrompt> = None;
1408    let mut pending_roll_acks: usize = 0;
1409    let mut decision_received: Option<Instant> = None;
1410    let mut decision_submitted: Option<Instant> = None;
1411
1412    loop {
1413        if cancel.load(std::sync::atomic::Ordering::Relaxed) {
1414            info!(
1415                session_id,
1416                "hosted java-forge session cancelled; player left the game"
1417            );
1418            return Ok(());
1419        }
1420        for (player_index, rx) in &mut remote_response_rxs {
1421            loop {
1422                match rx.try_recv() {
1423                    Ok(ClientToServerMessage::Response {
1424                        action: PromptOutput::DiceRolled(DiceRolledOutput::DiceRolledAcknowledged),
1425                        ..
1426                    }) => {
1427                        if pending_roll_acks > 0 {
1428                            pending_roll_acks -= 1;
1429                            if pending_roll_acks == 0 {
1430                                let ack = serde_json::to_string(&PromptOutput::DiceRolled(
1431                                    DiceRolledOutput::DiceRolledAcknowledged,
1432                                ))
1433                                .map_err(|err| format!("failed to serialize roll ack: {err}"))?;
1434                                engine.submit_action(&session_id, &ack)?;
1435                            }
1436                        }
1437                    }
1438                    Ok(ClientToServerMessage::Response { prompt_id, action }) => {
1439                        // prompt_id 0 is transport-synthesized: the
1440                        // absent-player default, exempt from validation.
1441                        if prompt_id != 0 {
1442                            let Some(prompt) =
1443                                last_prompt.as_ref().filter(|p| p.prompt_id == prompt_id)
1444                            else {
1445                                reject_response(
1446                                    &remote_prompt_tx,
1447                                    *player_index,
1448                                    last_prompt.as_ref().filter(|p| {
1449                                        self::player_index(&p.deciding_player_id) == *player_index
1450                                    }),
1451                                    ProtocolErrorCode::StalePrompt,
1452                                    format!("response for prompt {prompt_id} is not open"),
1453                                );
1454                                continue;
1455                            };
1456                            if self::player_index(&prompt.deciding_player_id) != *player_index {
1457                                reject_response(
1458                                    &remote_prompt_tx,
1459                                    *player_index,
1460                                    None,
1461                                    ProtocolErrorCode::WrongPlayer,
1462                                    format!(
1463                                        "prompt {prompt_id} is for {}",
1464                                        prompt.deciding_player_id
1465                                    ),
1466                                );
1467                                continue;
1468                            }
1469                            match prompt.input.validate_response(&action) {
1470                                Ok(()) => {}
1471                                Err(ResponseViolation::WrongPromptType) => {
1472                                    reject_response(
1473                                        &remote_prompt_tx,
1474                                        *player_index,
1475                                        Some(prompt),
1476                                        ProtocolErrorCode::WrongPromptType,
1477                                        "response output does not match the prompt type"
1478                                            .to_string(),
1479                                    );
1480                                    continue;
1481                                }
1482                                Err(ResponseViolation::UnknownActionId(id)) => {
1483                                    reject_response(
1484                                        &remote_prompt_tx,
1485                                        *player_index,
1486                                        Some(prompt),
1487                                        ProtocolErrorCode::UnknownActionId,
1488                                        format!(
1489                                            "action id {id:?} was not advertised by the prompt"
1490                                        ),
1491                                    );
1492                                    continue;
1493                                }
1494                                Err(ResponseViolation::CancelNotAllowed) => {
1495                                    reject_response(
1496                                        &remote_prompt_tx,
1497                                        *player_index,
1498                                        Some(prompt),
1499                                        ProtocolErrorCode::CancelNotAllowed,
1500                                        "this prompt is not cancellable".to_string(),
1501                                    );
1502                                    continue;
1503                                }
1504                            }
1505                        }
1506                        decision_received = Some(Instant::now());
1507                        let action_json = serde_json::to_string(&action).map_err(|err| {
1508                            format!(
1509                                "failed to serialize prompt output for player {player_index}: {err}"
1510                            )
1511                        })?;
1512                        debug!(player_index, %action_json, "submitting remote response to java");
1513                        engine.submit_action(&session_id, &action_json)?;
1514                        decision_submitted = Some(Instant::now());
1515                    }
1516                    Ok(ClientToServerMessage::Directive {
1517                        directive: DirectiveInput::Concede,
1518                    }) => {
1519                        // A directive can arrive while another player's prompt
1520                        // is open (this loop drains every seat), so it names
1521                        // its seat.
1522                        let directive_json = directive_concede_json(*player_index);
1523                        debug!(player_index, %directive_json, "submitting concede directive to java");
1524                        engine.submit_action(&session_id, &directive_json)?;
1525                    }
1526                    Err(TryRecvError::Empty) => break,
1527                    Err(TryRecvError::Disconnected) => {
1528                        debug!(player_index, "java-forge response channel disconnected");
1529                        break;
1530                    }
1531                }
1532            }
1533        }
1534
1535        if let Some(prompt_json) = engine.get_prompt(&session_id, 0)? {
1536            let prompt: AgentPrompt = serde_json::from_str(&prompt_json)
1537                .map_err(|err| format!("failed to parse java prompt: {err}"))?;
1538            if last_prompt.as_ref().map(|p| p.prompt_id) != Some(prompt.prompt_id) {
1539                if let Some(started) = decision_submitted.take() {
1540                    crate::metrics::record_forge_decision_stage("next_prompt", started.elapsed());
1541                }
1542                if let Some(started) = decision_received.take() {
1543                    crate::metrics::record_forge_decision_stage(
1544                        "decision_total",
1545                        started.elapsed(),
1546                    );
1547                    report_engine_gc(&engine);
1548                }
1549                last_prompt = Some(prompt.clone());
1550                let player = player_index(&prompt.deciding_player_id);
1551                debug!(player, "forwarding java prompt to remote");
1552                if matches!(prompt.input, PromptInput::DiceRolled(_)) {
1553                    let snapshots_started = Instant::now();
1554                    let prompt_msg = AgentMessage::Prompt(prompt);
1555                    for &agent_index in remote_response_rxs.keys() {
1556                        let state = AgentMessage::State(state_via_handle(
1557                            &engine,
1558                            &session_id,
1559                            Some(agent_index),
1560                        )?);
1561                        let _ = remote_prompt_tx.send((agent_index, state));
1562                        let _ = remote_prompt_tx.send((agent_index, prompt_msg.clone()));
1563                    }
1564                    send_observer_state(&engine, &session_id, &remote_prompt_tx);
1565                    pending_roll_acks = remote_response_rxs.len();
1566                    if pending_roll_acks == 0 {
1567                        let ack = serde_json::to_string(&PromptOutput::DiceRolled(
1568                            DiceRolledOutput::DiceRolledAcknowledged,
1569                        ))
1570                        .map_err(|err| format!("failed to serialize roll ack: {err}"))?;
1571                        engine.submit_action(&session_id, &ack)?;
1572                    }
1573                    crate::metrics::record_forge_decision_stage(
1574                        "snapshots",
1575                        snapshots_started.elapsed(),
1576                    );
1577                } else if Some(player) == local_player_index {
1578                    if let Some(output) = auto_action(&prompt) {
1579                        let action_json = serde_json::to_string(&output)
1580                            .map_err(|err| format!("failed to serialize auto action: {err}"))?;
1581                        engine.submit_action(&session_id, &action_json)?;
1582                    }
1583                } else {
1584                    let snapshots_started = Instant::now();
1585                    for &agent_index in remote_response_rxs.keys() {
1586                        let state = AgentMessage::State(state_via_handle(
1587                            &engine,
1588                            &session_id,
1589                            Some(agent_index),
1590                        )?);
1591                        if remote_prompt_tx.send((agent_index, state)).is_err() {
1592                            return Ok(());
1593                        }
1594                    }
1595                    send_observer_state(&engine, &session_id, &remote_prompt_tx);
1596                    let prompt_msg = AgentMessage::Prompt(prompt);
1597                    if remote_prompt_tx.send((player, prompt_msg)).is_err() {
1598                        return Ok(());
1599                    }
1600                    crate::metrics::record_forge_decision_stage(
1601                        "snapshots",
1602                        snapshots_started.elapsed(),
1603                    );
1604                }
1605            }
1606        }
1607
1608        if engine.is_game_over(&session_id)? {
1609            info!("hosted java-forge session reached game over");
1610            let mut final_messages = Vec::new();
1611            for &agent_index in remote_response_rxs.keys() {
1612                match state_via_handle(&engine, &session_id, Some(agent_index)) {
1613                    Ok(state_update) => {
1614                        final_messages.push((agent_index, AgentMessage::State(state_update)));
1615                    }
1616                    Err(error) => {
1617                        warn!(%error, agent_index, "game over: final snapshot unavailable; sending game-over prompt only");
1618                    }
1619                }
1620            }
1621            if let Ok(state_update) = state_via_handle(&engine, &session_id, None) {
1622                final_messages.push((
1623                    crate::host::OBSERVER_SEAT,
1624                    AgentMessage::State(state_update),
1625                ));
1626            }
1627            let game_over = AgentMessage::Prompt(game_over_prompt());
1628            for &agent_index in remote_response_rxs.keys() {
1629                final_messages.push((agent_index, game_over.clone()));
1630            }
1631            let _ = game_over_tx.send(HostedGameOver {
1632                game_id: game_id.clone(),
1633                messages: final_messages,
1634            });
1635            engine.end_game(&session_id)?;
1636            guard.armed.set(false);
1637            return Ok(());
1638        }
1639
1640        std::thread::sleep(Duration::from_millis(50));
1641    }
1642}
1643
1644#[cfg(feature = "java-forge")]
1645fn wait_for_prompt<B: JavaBridge>(
1646    session: &mut JavaForgeSession<B>,
1647    max_polls: usize,
1648) -> Result<Option<String>, String> {
1649    for _ in 0..max_polls {
1650        if let Some(prompt) = session.get_prompt(0)? {
1651            return Ok(Some(prompt));
1652        }
1653        std::thread::sleep(Duration::from_millis(50));
1654    }
1655    Ok(None)
1656}
1657
1658#[cfg(forge_backend)]
1659fn send_observer_state(
1660    engine: &ForgeEngine,
1661    session_id: &str,
1662    remote_prompt_tx: &std_mpsc::Sender<(usize, AgentMessage)>,
1663) {
1664    match state_via_handle(engine, session_id, None) {
1665        Ok(state_update) => {
1666            let _ = remote_prompt_tx.send((
1667                crate::host::OBSERVER_SEAT,
1668                AgentMessage::State(state_update),
1669            ));
1670        }
1671        Err(error) => warn!(%error, "observer snapshot unavailable"),
1672    }
1673}
1674
1675#[cfg(forge_backend)]
1676fn player_index(deciding_player_id: &str) -> usize {
1677    deciding_player_id
1678        .strip_prefix("player-")
1679        .and_then(|n| n.parse().ok())
1680        .unwrap_or(0)
1681}
1682
1683#[cfg(forge_backend)]
1684fn reject_response(
1685    remote_prompt_tx: &std_mpsc::Sender<(usize, AgentMessage)>,
1686    seat: usize,
1687    reopen_prompt: Option<&AgentPrompt>,
1688    code: ProtocolErrorCode,
1689    message: String,
1690) {
1691    let _ = remote_prompt_tx.send((
1692        seat,
1693        AgentMessage::Error(ProtocolError {
1694            code,
1695            message,
1696            prompt_id: reopen_prompt.map(|p| p.prompt_id),
1697        }),
1698    ));
1699    if let Some(prompt) = reopen_prompt {
1700        let _ = remote_prompt_tx.send((seat, AgentMessage::Prompt(prompt.clone())));
1701    }
1702}
1703
1704#[cfg(forge_backend)]
1705fn auto_action(prompt: &AgentPrompt) -> Option<PromptOutput> {
1706    match prompt.input {
1707        PromptInput::ChooseAction(_) => {
1708            Some(PromptOutput::ChooseAction(ChooseActionOutput::Pass {
1709                until: None,
1710                exhaust_stack: false,
1711            }))
1712        }
1713        _ => None,
1714    }
1715}
1716
1717#[cfg(forge_backend)]
1718fn game_over_prompt() -> AgentPrompt {
1719    AgentPrompt {
1720        prompt_id: u32::MAX,
1721        deciding_player_id: "player-0".to_string(),
1722        source_card: None,
1723        input: PromptInput::GameOver(GameOverInput {}),
1724    }
1725}
1726
1727#[cfg(forge_backend)]
1728fn state_via_handle(
1729    engine: &ForgeEngine,
1730    session_id: &str,
1731    viewer: Option<usize>,
1732) -> Result<StateUpdate, String> {
1733    let game_view: GameViewDto = serde_json::from_str(&engine.get_snapshot(session_id, viewer)?)
1734        .map_err(|err| format!("failed to parse java snapshot: {err}"))?;
1735    Ok(StateUpdate { game_view })
1736}
1737
1738#[cfg(forge_backend)]
1739fn deck_card_identities(deck: &Deck) -> Vec<DeckCardIdentity> {
1740    deck.cards
1741        .iter()
1742        .chain(deck.commanders.iter().flatten())
1743        .map(|card| card.identity.clone())
1744        .collect()
1745}
1746
1747#[cfg(forge_backend)]
1748fn commander_names_for_java(deck: &Deck, fallback: Option<&str>) -> Vec<String> {
1749    let names: Vec<String> = deck
1750        .commanders
1751        .iter()
1752        .flatten()
1753        .map(|card| java_card_name(&card.identity.name))
1754        .collect();
1755    if !names.is_empty() {
1756        return names;
1757    }
1758    fallback
1759        .filter(|name| !name.is_empty())
1760        .map(|name| vec![java_card_name(name)])
1761        .unwrap_or_default()
1762}
1763
1764#[cfg(forge_backend)]
1765fn smoke_deck(land_name: &str, spell_name: &str) -> Vec<DeckCardIdentity> {
1766    (0..24)
1767        .map(|_| DeckCardIdentity {
1768            name: land_name.to_string(),
1769            ..Default::default()
1770        })
1771        .chain((0..36).map(|_| DeckCardIdentity {
1772            name: spell_name.to_string(),
1773            ..Default::default()
1774        }))
1775        .collect()
1776}
1777
1778#[cfg(feature = "java-forge")]
1779fn scenario_deck(land_name: &str) -> Vec<DeckCardIdentity> {
1780    (0..60)
1781        .map(|_| DeckCardIdentity {
1782            name: land_name.to_string(),
1783            ..Default::default()
1784        })
1785        .collect()
1786}
1787
1788#[cfg(feature = "java-forge")]
1789enum JavaScenario {
1790    KeepAndPlayLand {
1791        played_land: bool,
1792    },
1793    MulliganOncePlayLand {
1794        mulliganed: bool,
1795        kept_second_hand: bool,
1796        put_back_done: bool,
1797        played_land: bool,
1798    },
1799}
1800
1801#[cfg(feature = "java-forge")]
1802impl JavaScenario {
1803    fn from_name(name: &str) -> Result<Self, String> {
1804        match name {
1805            "keep-and-play-land" => Ok(Self::KeepAndPlayLand { played_land: false }),
1806            "mulligan-once-play-land" => Ok(Self::MulliganOncePlayLand {
1807                mulliganed: false,
1808                kept_second_hand: false,
1809                put_back_done: false,
1810                played_land: false,
1811            }),
1812            _ => Err(format!(
1813                "unknown java-forge scenario '{name}'. Supported scenarios: keep-and-play-land, mulligan-once-play-land"
1814            )),
1815        }
1816    }
1817
1818    fn name(&self) -> &'static str {
1819        match self {
1820            Self::KeepAndPlayLand { .. } => "keep-and-play-land",
1821            Self::MulliganOncePlayLand { .. } => "mulligan-once-play-land",
1822        }
1823    }
1824
1825    fn next_action(
1826        &mut self,
1827        prompt: &Value,
1828        game_view: &GameViewDto,
1829    ) -> Result<Option<PromptOutput>, String> {
1830        match self {
1831            Self::KeepAndPlayLand { played_land } => {
1832                if *played_land && battlefield_contains(game_view, "Swamp") {
1833                    return Ok(None);
1834                }
1835                match prompt_type(prompt) {
1836                    Some("mulligan") => Ok(Some(PromptOutput::Mulligan(
1837                        MulliganOutput::MulliganDecision { keep: true },
1838                    ))),
1839                    Some("chooseAction") => {
1840                        if let Some(action) = play_first_card_action(prompt, "Swamp")? {
1841                            *played_land = true;
1842                            Ok(Some(action))
1843                        } else {
1844                            Ok(Some(PromptOutput::ChooseAction(ChooseActionOutput::Pass {
1845                                until: None,
1846                                exhaust_stack: false,
1847                            })))
1848                        }
1849                    }
1850                    other => Err(format!(
1851                        "scenario '{}' expected mulligan or chooseAction, got {:?}",
1852                        self.name(),
1853                        other
1854                    )),
1855                }
1856            }
1857            Self::MulliganOncePlayLand {
1858                mulliganed,
1859                kept_second_hand,
1860                put_back_done,
1861                played_land,
1862            } => {
1863                if *played_land && battlefield_contains(game_view, "Swamp") {
1864                    return Ok(None);
1865                }
1866                match prompt_type(prompt) {
1867                    Some("mulligan") if !*mulliganed => {
1868                        *mulliganed = true;
1869                        Ok(Some(PromptOutput::Mulligan(MulliganOutput::MulliganDecision { keep: false })))
1870                    }
1871                    Some("mulligan") if !*kept_second_hand => {
1872                        *kept_second_hand = true;
1873                        Ok(Some(PromptOutput::Mulligan(MulliganOutput::MulliganDecision { keep: true })))
1874                    }
1875                    Some("mulliganPutBack") if !*put_back_done => {
1876                        let count = prompt
1877                            .get("input")
1878                            .and_then(|input| input.get("count"))
1879                            .and_then(Value::as_u64)
1880                            .unwrap_or(1) as usize;
1881                        let card_ids = prompt_card_ids(prompt, "handCardIds", count)?;
1882                        *put_back_done = true;
1883                        Ok(Some(PromptOutput::MulliganPutBack(MulliganPutBackOutput::MulliganPutBackDecision { card_ids })))
1884                    }
1885                    Some("chooseAction") => {
1886                        if let Some(action) = play_first_card_action(prompt, "Swamp")? {
1887                            *played_land = true;
1888                            Ok(Some(action))
1889                        } else {
1890                            Ok(Some(PromptOutput::ChooseAction(ChooseActionOutput::Pass { until: None, exhaust_stack: false })))
1891                        }
1892                    }
1893                    other => Err(format!(
1894                        "scenario '{}' expected mulligan, mulliganPutBack, or chooseAction, got {:?}",
1895                        self.name(),
1896                        other
1897                    )),
1898                }
1899            }
1900        }
1901    }
1902}
1903
1904#[cfg(feature = "java-forge")]
1905fn run_scenario_loop<B: JavaBridge>(
1906    session: &mut JavaForgeSession<B>,
1907    mut scenario: JavaScenario,
1908    max_prompts: usize,
1909) -> Result<(), String> {
1910    let mut prompts_seen = 0usize;
1911    let mut last_prompt_json: Option<String> = None;
1912    while prompts_seen < max_prompts {
1913        let Some(prompt_json) = wait_for_prompt(session, 600)? else {
1914            return Err(format!(
1915                "timed out waiting for java-forge scenario '{}' prompt",
1916                scenario.name()
1917            ));
1918        };
1919        if last_prompt_json.as_deref() == Some(prompt_json.as_str()) {
1920            std::thread::sleep(Duration::from_millis(50));
1921            continue;
1922        }
1923        last_prompt_json = Some(prompt_json.clone());
1924        prompts_seen += 1;
1925
1926        let prompt: AgentPrompt = serde_json::from_str(&prompt_json)
1927            .map_err(|err| format!("failed to parse java scenario prompt: {err}"))?;
1928        let player = player_index(&prompt.deciding_player_id);
1929        if player != 0 {
1930            if let Some(output) = auto_action(&prompt) {
1931                session.submit_action(
1932                    &serde_json::to_string(&output).map_err(|err| err.to_string())?,
1933                )?;
1934            }
1935            continue;
1936        }
1937
1938        let game_view: GameViewDto = serde_json::from_str(&session.get_snapshot(Some(0))?)
1939            .map_err(|err| format!("failed to parse java scenario snapshot: {err}"))?;
1940        let normalized_prompt = serde_json::to_value(&prompt).map_err(|err| err.to_string())?;
1941        info!(
1942            scenario = scenario.name(),
1943            prompts_seen,
1944            prompt_type = prompt_type(&normalized_prompt).unwrap_or("<missing>"),
1945            "java-forge scenario prompt"
1946        );
1947        let Some(action) = scenario.next_action(&normalized_prompt, &game_view)? else {
1948            info!(
1949                scenario = scenario.name(),
1950                prompts_seen, "java-forge scenario assertions satisfied"
1951            );
1952            return Ok(());
1953        };
1954        submit_player_action(session, &action)?;
1955    }
1956    Err(format!(
1957        "java-forge scenario '{}' did not complete within {max_prompts} prompts",
1958        scenario.name()
1959    ))
1960}
1961
1962#[cfg(feature = "java-forge")]
1963pub fn run_concede_smoke() -> Result<(), String> {
1964    let config = JavaRuntimeConfig::from_env();
1965    let assets_dir = config.assets_dir.to_string_lossy().to_string();
1966    let bridge = SubprocessBridge::spawn(&config)?;
1967    let mut session = JavaForgeSession::new(bridge);
1968    session.initialize(&assets_dir)?;
1969    run_concede_game(&mut session, 3, 2, 12)?;
1970    run_concede_game(&mut session, 2, 1, 8)?;
1971    Ok(())
1972}
1973
1974#[cfg(not(feature = "java-forge"))]
1975pub fn run_concede_smoke() -> Result<(), String> {
1976    Err(unsupported_message().to_string())
1977}
1978
1979#[cfg(forge_backend)]
1980fn directive_concede_json(player: usize) -> String {
1981    format!(r#"{{"type":"directive","directive":{{"type":"concede"}},"player":{player}}}"#)
1982}
1983
1984#[cfg(feature = "java-forge")]
1985fn run_concede_game<B: JavaBridge>(
1986    session: &mut JavaForgeSession<B>,
1987    seats: usize,
1988    conceder: usize,
1989    concede_after: usize,
1990) -> Result<(), String> {
1991    const POST_CONCEDE_DECISIONS: usize = 12;
1992    const STALL_REPEATS: usize = 300;
1993
1994    let cards: Vec<manabrew_protocol::deck_dto::DeckCard> = (0..60)
1995        .map(|_| manabrew_protocol::deck_dto::DeckCard {
1996            identity: DeckCardIdentity {
1997                name: "Mountain".to_string(),
1998                set_code: "M20".to_string(),
1999                ..Default::default()
2000            },
2001            ..Default::default()
2002        })
2003        .collect();
2004    let deck = Deck {
2005        name: "concede-smoke".to_string(),
2006        cards,
2007        ..Default::default()
2008    };
2009    let identities = deck_card_identities(&deck);
2010    let players: Vec<PlayerConfig> = (0..seats)
2011        .map(|i| PlayerConfig::new(format!("Concede {}", i + 1), &identities, Vec::new()))
2012        .collect();
2013    let request = StartGameRequest::new(
2014        format!("concede-smoke-{seats}p"),
2015        String::new(),
2016        20,
2017        7,
2018        players,
2019    );
2020    let session_id = session.start_game(&request)?;
2021    info!(
2022        session_id,
2023        seats, conceder, concede_after, "concede smoke game started"
2024    );
2025
2026    let mut bots: HashMap<usize, SimpleAi> = HashMap::new();
2027    let mut last_prompt_json: Option<String> = None;
2028    let mut acted = 0usize;
2029    let mut acted_after_concede = 0usize;
2030    let mut conceded = false;
2031    let mut repeat_count = 0usize;
2032
2033    for _ in 0..40_000 {
2034        if session.is_game_over()? {
2035            if !conceded {
2036                return Err("concede smoke: game ended before the concede fired".to_string());
2037            }
2038            break;
2039        }
2040        let Some(prompt_json) = session.get_prompt(0)? else {
2041            std::thread::sleep(Duration::from_millis(10));
2042            continue;
2043        };
2044        if last_prompt_json.as_deref() == Some(prompt_json.as_str()) {
2045            repeat_count += 1;
2046            if repeat_count > STALL_REPEATS {
2047                return Err(format!(
2048                    "concede smoke stalled on the same prompt (seats={seats} acted={acted} conceded={conceded} after={acted_after_concede}): {prompt_json}"
2049                ));
2050            }
2051            std::thread::sleep(Duration::from_millis(10));
2052            continue;
2053        }
2054        repeat_count = 0;
2055
2056        let prompt: AgentPrompt = serde_json::from_str(&prompt_json)
2057            .map_err(|err| format!("concede smoke: bad prompt: {err}"))?;
2058        let player = player_index(&prompt.deciding_player_id);
2059
2060        if !conceded && acted >= concede_after {
2061            info!(conceder, acted, "concede smoke: injecting concede");
2062            session.submit_action(&directive_concede_json(conceder))?;
2063            conceded = true;
2064            last_prompt_json = None;
2065            continue;
2066        }
2067        if conceded && player == conceder {
2068            session.submit_action(&directive_concede_json(conceder))?;
2069            last_prompt_json = Some(prompt_json);
2070            continue;
2071        }
2072
2073        if let Some(action) = bots.entry(player).or_default().decide(prompt) {
2074            submit_player_action(session, &action)?;
2075            acted += 1;
2076            if conceded {
2077                acted_after_concede += 1;
2078            }
2079        }
2080        last_prompt_json = Some(prompt_json);
2081        if conceded && seats > 2 && acted_after_concede >= POST_CONCEDE_DECISIONS {
2082            break;
2083        }
2084    }
2085
2086    if !conceded {
2087        return Err("concede smoke: never reached the injection point".to_string());
2088    }
2089    if seats == 2 {
2090        if !session.is_game_over()? {
2091            return Err("concede smoke: 2p concession did not end the game".to_string());
2092        }
2093        info!(acted, "concede smoke: 2p concession ended the game");
2094    } else {
2095        if acted_after_concede < POST_CONCEDE_DECISIONS && !session.is_game_over()? {
2096            return Err(format!(
2097                "concede smoke: game did not progress after the concession (only {acted_after_concede} decisions)"
2098            ));
2099        }
2100        let snapshot = parse_snapshot(session)?;
2101        let status = snapshot
2102            .pointer(&format!("/players/{conceder}/status"))
2103            .and_then(Value::as_str)
2104            .unwrap_or("");
2105        if status != "conceded" {
2106            return Err(format!(
2107                "concede smoke: seat {conceder} has status '{status}', expected 'conceded'"
2108            ));
2109        }
2110        info!(
2111            acted,
2112            acted_after_concede, "concede smoke: game continued past the concession"
2113        );
2114    }
2115    session.end_game()?;
2116    Ok(())
2117}
2118
2119#[cfg(feature = "java-forge")]
2120fn run_self_play_loop<B: JavaBridge>(
2121    session: &mut JavaForgeSession<B>,
2122    max_prompts: usize,
2123) -> Result<(), String> {
2124    const STALL_REPEATS: usize = 100;
2125
2126    let mut bots: HashMap<usize, SimpleAi> = HashMap::new();
2127    let mut last_prompt_json: Option<String> = None;
2128    let mut acted = 0usize;
2129    let mut repeat_count = 0usize;
2130    let mut seen_prompt = false;
2131    let max_iterations = max_prompts.saturating_mul(200).max(2_000);
2132
2133    for _ in 0..max_iterations {
2134        if let Some(prompt_json) = session.get_prompt(0)? {
2135            seen_prompt = true;
2136            if last_prompt_json.as_deref() == Some(prompt_json.as_str()) {
2137                if session.is_game_over()? {
2138                    info!(acted, "java-forge self-play reached game over");
2139                    return Ok(());
2140                }
2141                repeat_count += 1;
2142                if repeat_count > STALL_REPEATS {
2143                    let raw_value: Value =
2144                        serde_json::from_str(&prompt_json).unwrap_or(Value::Null);
2145                    dump_stuck(
2146                        "java re-emitted the same prompt after the bot acted (stall)",
2147                        &raw_value,
2148                        None,
2149                        session,
2150                    );
2151                    return Err(
2152                        "self-play stalled: java re-emitted the same prompt after the bot's action"
2153                            .to_string(),
2154                    );
2155                }
2156                std::thread::sleep(Duration::from_millis(20));
2157                continue;
2158            }
2159            repeat_count = 0;
2160
2161            let prompt: AgentPrompt = serde_json::from_str(&prompt_json)
2162                .map_err(|err| format!("failed to parse java self-play prompt: {err}"))?;
2163            let player = player_index(&prompt.deciding_player_id);
2164            let raw_value: Value = serde_json::from_str(&prompt_json).unwrap_or(Value::Null);
2165            let normalized = raw_value.clone();
2166
2167            match bots.entry(player).or_default().decide(prompt) {
2168                Some(action) => {
2169                    if let Err(err) = submit_player_action(session, &action) {
2170                        dump_stuck(
2171                            "java rejected the bot action",
2172                            &raw_value,
2173                            Some(&normalized),
2174                            session,
2175                        );
2176                        return Err(format!(
2177                            "self-play: java rejected action for player {player}: {err}"
2178                        ));
2179                    }
2180                    acted += 1;
2181                    if acted >= max_prompts {
2182                        dump_stuck(
2183                            "did not reach game over within max prompts",
2184                            &raw_value,
2185                            Some(&normalized),
2186                            session,
2187                        );
2188                        return Err(format!(
2189                            "self-play did not reach game over within {max_prompts} decisions"
2190                        ));
2191                    }
2192                }
2193                None => debug!(
2194                    player,
2195                    prompt_type = prompt_type(&normalized).unwrap_or("<missing>"),
2196                    "self-play: no action for prompt (display-only)"
2197                ),
2198            }
2199            last_prompt_json = Some(prompt_json);
2200            continue;
2201        }
2202
2203        if seen_prompt && session.is_game_over()? {
2204            info!(acted, "java-forge self-play reached game over");
2205            return Ok(());
2206        }
2207        std::thread::sleep(Duration::from_millis(20));
2208    }
2209
2210    dump_stuck(
2211        "self-play exceeded its iteration cap without game over",
2212        &Value::Null,
2213        None,
2214        session,
2215    );
2216    Err("self-play exceeded its iteration cap without reaching game over".to_string())
2217}
2218
2219#[cfg(feature = "java-forge")]
2220fn parse_snapshot<B: JavaBridge>(session: &mut JavaForgeSession<B>) -> Result<Value, String> {
2221    let snapshot_json = session.get_snapshot(Some(0))?;
2222    serde_json::from_str(&snapshot_json)
2223        .map_err(|err| format!("failed to parse java self-play snapshot: {err}"))
2224}
2225
2226#[cfg(feature = "java-forge")]
2227fn dump_stuck<B: JavaBridge>(
2228    reason: &str,
2229    prompt: &Value,
2230    normalized: Option<&Value>,
2231    session: &mut JavaForgeSession<B>,
2232) {
2233    let snapshot = parse_snapshot(session).unwrap_or(Value::Null);
2234    let artifact = json!({
2235        "reason": reason,
2236        "rawPrompt": prompt,
2237        "normalizedPrompt": normalized,
2238        "snapshot": snapshot,
2239    });
2240    let ts = std::time::SystemTime::now()
2241        .duration_since(std::time::UNIX_EPOCH)
2242        .map(|d| d.as_millis())
2243        .unwrap_or(0);
2244    let path = workspace_root().join(format!("target/self-play-stuck-{ts}.json"));
2245    match serde_json::to_string_pretty(&artifact) {
2246        Ok(body) => {
2247            if let Err(error) = std::fs::write(&path, body) {
2248                warn!(%error, reason, "self-play stuck; failed to write artifact");
2249            } else {
2250                warn!(path = %path.display(), reason, "self-play stuck; wrote artifact");
2251            }
2252        }
2253        Err(error) => warn!(%error, reason, "self-play stuck; failed to serialize artifact"),
2254    }
2255}
2256
2257#[cfg(feature = "java-forge")]
2258fn submit_player_action<B: JavaBridge>(
2259    session: &mut JavaForgeSession<B>,
2260    action: &PromptOutput,
2261) -> Result<(), String> {
2262    let action_json = serde_json::to_string(action)
2263        .map_err(|err| format!("failed to serialize scenario action: {err}"))?;
2264    session.submit_action(&action_json)?;
2265    Ok(())
2266}
2267
2268#[cfg(feature = "java-forge")]
2269fn prompt_type(prompt: &Value) -> Option<&str> {
2270    prompt
2271        .get("input")
2272        .and_then(|input| input.get("type"))
2273        .and_then(Value::as_str)
2274}
2275
2276#[cfg(feature = "java-forge")]
2277fn play_first_card_action(prompt: &Value, card_name: &str) -> Result<Option<PromptOutput>, String> {
2278    let Some(action) = prompt
2279        .get("input")
2280        .and_then(|input| input.get("actions"))
2281        .and_then(Value::as_array)
2282        .and_then(|actions| {
2283            actions.iter().find(|action| {
2284                action
2285                    .get("modeLabel")
2286                    .and_then(Value::as_str)
2287                    .is_some_and(|label| label.contains(card_name))
2288            })
2289        })
2290    else {
2291        return Ok(None);
2292    };
2293    let action_id = action
2294        .get("id")
2295        .and_then(Value::as_str)
2296        .ok_or_else(|| format!("playable action for '{card_name}' is missing id"))?;
2297    Ok(Some(PromptOutput::ChooseAction(ChooseActionOutput::Act {
2298        action_id: action_id.to_string(),
2299    })))
2300}
2301
2302#[cfg(feature = "java-forge")]
2303fn prompt_card_ids(prompt: &Value, field: &str, count: usize) -> Result<Vec<String>, String> {
2304    let card_ids = prompt
2305        .get("input")
2306        .and_then(|input| input.get(field))
2307        .and_then(Value::as_array)
2308        .ok_or_else(|| format!("prompt is missing {field}"))?;
2309    if card_ids.len() < count {
2310        return Err(format!(
2311            "prompt field {field} has {} cards, need {count}",
2312            card_ids.len()
2313        ));
2314    }
2315    Ok(card_ids
2316        .iter()
2317        .take(count)
2318        .filter_map(Value::as_str)
2319        .map(str::to_string)
2320        .collect())
2321}
2322
2323#[cfg(feature = "java-forge")]
2324fn battlefield_contains(game_view: &GameViewDto, card_name: &str) -> bool {
2325    use manabrew_agent_interface::game_view_dto::{CardView, ZoneKind};
2326    game_view.zones.iter().any(|zone| {
2327        zone.zone == ZoneKind::Battlefield
2328            && zone.owner_id == "player-0"
2329            && zone.cards.iter().any(|card| match card {
2330                CardView::Visible(dto) => dto.identity.name == card_name,
2331                CardView::Hidden { .. } => false,
2332            })
2333    })
2334}
2335
2336pub trait JavaBridge {
2337    fn initialize(&mut self, assets_dir: &str) -> Result<(), String>;
2338    fn start_game_json(&mut self, request_json: &str) -> Result<String, String>;
2339    fn submit_action(&mut self, session_id: &str, action_json: &str) -> Result<String, String>;
2340    fn get_prompt(
2341        &mut self,
2342        session_id: &str,
2343        player_index: usize,
2344    ) -> Result<Option<String>, String>;
2345    fn get_snapshot(&mut self, session_id: &str, viewer: Option<usize>) -> Result<String, String>;
2346    fn is_game_over(&mut self, session_id: &str) -> Result<bool, String>;
2347    fn end_game(&mut self, session_id: &str) -> Result<(), String>;
2348    fn abort_game(&mut self, session_id: &str) -> Result<(), String>;
2349}
2350
2351pub struct JavaForgeSession<B> {
2352    bridge: B,
2353    session_id: Option<String>,
2354}
2355
2356impl<B: JavaBridge> JavaForgeSession<B> {
2357    pub fn new(bridge: B) -> Self {
2358        Self {
2359            bridge,
2360            session_id: None,
2361        }
2362    }
2363
2364    pub fn initialize(&mut self, assets_dir: &str) -> Result<(), String> {
2365        self.bridge.initialize(assets_dir)
2366    }
2367
2368    pub fn start_game(&mut self, request: &StartGameRequest) -> Result<String, String> {
2369        let request_json = request.to_json().map_err(|err| err.to_string())?;
2370        let response_json = self.bridge.start_game_json(&request_json)?;
2371        let response: StartGameResponse =
2372            serde_json::from_str(&response_json).map_err(|err| err.to_string())?;
2373        self.session_id = Some(response.session_id.clone());
2374        Ok(response.session_id)
2375    }
2376
2377    pub fn submit_action(&mut self, action_json: &str) -> Result<String, String> {
2378        let session_id = self.require_session_id()?.to_string();
2379        self.bridge.submit_action(&session_id, action_json)
2380    }
2381
2382    pub fn get_prompt(&mut self, player_index: usize) -> Result<Option<String>, String> {
2383        let session_id = self.require_session_id()?.to_string();
2384        self.bridge.get_prompt(&session_id, player_index)
2385    }
2386
2387    pub fn get_snapshot(&mut self, viewer: Option<usize>) -> Result<String, String> {
2388        let session_id = self.require_session_id()?.to_string();
2389        self.bridge.get_snapshot(&session_id, viewer)
2390    }
2391
2392    pub fn is_game_over(&mut self) -> Result<bool, String> {
2393        let session_id = self.require_session_id()?.to_string();
2394        self.bridge.is_game_over(&session_id)
2395    }
2396
2397    pub fn end_game(&mut self) -> Result<(), String> {
2398        let Some(session_id) = self.session_id.take() else {
2399            return Ok(());
2400        };
2401        self.bridge.end_game(&session_id)
2402    }
2403
2404    fn require_session_id(&self) -> Result<&str, String> {
2405        self.session_id
2406            .as_deref()
2407            .ok_or_else(|| "java-forge session has not started".to_string())
2408    }
2409}
2410
2411pub struct UnavailableJavaBridge;
2412
2413impl JavaBridge for UnavailableJavaBridge {
2414    fn initialize(&mut self, _assets_dir: &str) -> Result<(), String> {
2415        Err(unsupported_message().to_string())
2416    }
2417
2418    fn start_game_json(&mut self, _request_json: &str) -> Result<String, String> {
2419        Err(unsupported_message().to_string())
2420    }
2421
2422    fn submit_action(&mut self, _session_id: &str, _action_json: &str) -> Result<String, String> {
2423        Err(unsupported_message().to_string())
2424    }
2425
2426    fn get_prompt(
2427        &mut self,
2428        _session_id: &str,
2429        _player_index: usize,
2430    ) -> Result<Option<String>, String> {
2431        Err(unsupported_message().to_string())
2432    }
2433
2434    fn get_snapshot(
2435        &mut self,
2436        _session_id: &str,
2437        _viewer: Option<usize>,
2438    ) -> Result<String, String> {
2439        Err(unsupported_message().to_string())
2440    }
2441
2442    fn is_game_over(&mut self, _session_id: &str) -> Result<bool, String> {
2443        Err(unsupported_message().to_string())
2444    }
2445
2446    fn end_game(&mut self, _session_id: &str) -> Result<(), String> {
2447        Err(unsupported_message().to_string())
2448    }
2449
2450    fn abort_game(&mut self, _session_id: &str) -> Result<(), String> {
2451        Err(unsupported_message().to_string())
2452    }
2453}
2454
2455#[cfg(feature = "java-forge")]
2456#[derive(serde::Deserialize)]
2457struct SubprocessReply {
2458    ok: bool,
2459    #[serde(default)]
2460    result: String,
2461    #[serde(default)]
2462    error: Option<String>,
2463}
2464
2465#[cfg(feature = "java-forge")]
2466const CALL_TIMEOUT: Duration = Duration::from_secs(60);
2467#[cfg(feature = "java-forge")]
2468const SHUTDOWN_GRACE: Duration = Duration::from_secs(5);
2469
2470#[cfg(feature = "java-forge")]
2471pub struct SubprocessBridge {
2472    child: Child,
2473    stdin: BufWriter<ChildStdin>,
2474    stdout_rx: std_mpsc::Receiver<String>,
2475    stdout_handle: Option<std::thread::JoinHandle<()>>,
2476    stderr_handle: Option<std::thread::JoinHandle<()>>,
2477}
2478
2479#[cfg(feature = "java-forge")]
2480impl SubprocessBridge {
2481    fn spawn(config: &JavaRuntimeConfig) -> Result<Self, String> {
2482        config.validate()?;
2483
2484        let java_bin = resolve_java_bin(config);
2485        let jvm_args = config.jvm_args();
2486        info!(target: "self_hosted_node::java", args = %jvm_args.join(" "), "spawning java engine");
2487        let mut cmd = Command::new(&java_bin);
2488        // The JVM applies JAVA_TOOL_OPTIONS before anything on the command line,
2489        // so whatever is set on the host silently governs every flag we do not
2490        // name. Production ran -Xmx512m -XX:+UseSerialGC that way for months
2491        // with nothing in the repo to show it. Use SELF_HOSTED_NODE_JAVA_OPTS
2492        // instead: it is passed here, logged above, and lives in config.
2493        cmd.env_remove("JAVA_TOOL_OPTIONS");
2494        cmd.args(&jvm_args);
2495        cmd.arg("-jar").arg(&config.harness_jar);
2496        cmd.arg("--interactive-server");
2497        cmd.arg("--forge-home")
2498            .arg(format!("{}/", config.assets_dir.display()));
2499
2500        let mut child = cmd
2501            .stdin(Stdio::piped())
2502            .stdout(Stdio::piped())
2503            .stderr(Stdio::piped())
2504            .spawn()
2505            .map_err(|err| format!("failed to spawn java subprocess: {err}"))?;
2506
2507        let stdin = child
2508            .stdin
2509            .take()
2510            .ok_or_else(|| "java subprocess has no stdin".to_string())?;
2511        let stdout = child
2512            .stdout
2513            .take()
2514            .ok_or_else(|| "java subprocess has no stdout".to_string())?;
2515        let stderr = child.stderr.take();
2516
2517        // Bounded so a chatty Java side can't grow the queue without bound.
2518        // Cap is generous — protocol replies are one line per request.
2519        let (stdout_tx, stdout_rx) = std_mpsc::sync_channel::<String>(1024);
2520        let stdout_handle = std::thread::spawn(move || {
2521            let reader = BufReader::new(stdout);
2522            for line in reader.lines().map_while(Result::ok) {
2523                if stdout_tx.send(line).is_err() {
2524                    break;
2525                }
2526            }
2527        });
2528
2529        let stderr_handle = std::thread::spawn(move || {
2530            if let Some(stderr) = stderr {
2531                let reader = BufReader::new(stderr);
2532                for line in reader.lines().map_while(Result::ok) {
2533                    if let Some(pause) = parse_gc_pause(&line) {
2534                        crate::metrics::record_jvm_gc(
2535                            pause.kind,
2536                            Duration::from_secs_f64(pause.millis / 1000.0),
2537                            pause.heap_after_mb,
2538                        );
2539                    }
2540                    if line.contains("Exception") || line.contains("ERROR") {
2541                        warn!(target: "self_hosted_node::java", "[java] {line}");
2542                    } else if line.contains("][gc") {
2543                        // -Xlog:gc*:stderr is opt-in, and debug would drop it before Loki.
2544                        info!(target: "self_hosted_node::java", "[java] {line}");
2545                    } else {
2546                        debug!(target: "self_hosted_node::java", "[java] {line}");
2547                    }
2548                }
2549            }
2550        });
2551
2552        Ok(Self {
2553            child,
2554            stdin: BufWriter::new(stdin),
2555            stdout_rx,
2556            stdout_handle: Some(stdout_handle),
2557            stderr_handle: Some(stderr_handle),
2558        })
2559    }
2560
2561    fn call(&mut self, request_json: &str) -> Result<String, String> {
2562        // Drain anything still queued from a prior request — a previous call()
2563        // that timed out may have left its reply in the channel, and consuming
2564        // it now would shift every subsequent call off-by-one.
2565        loop {
2566            match self.stdout_rx.try_recv() {
2567                Ok(stale) => {
2568                    debug!(target: "self_hosted_node::java", line = %stale, "discarding stale stdout line");
2569                }
2570                Err(TryRecvError::Empty) => break,
2571                Err(TryRecvError::Disconnected) => {
2572                    return Err("java subprocess closed stdout (crashed?)".to_string());
2573                }
2574            }
2575        }
2576
2577        self.stdin
2578            .write_all(request_json.as_bytes())
2579            .map_err(|err| format!("failed to write subprocess stdin: {err}"))?;
2580        self.stdin
2581            .write_all(b"\n")
2582            .map_err(|err| format!("failed to write subprocess newline: {err}"))?;
2583        self.stdin
2584            .flush()
2585            .map_err(|err| format!("failed to flush subprocess stdin: {err}"))?;
2586
2587        let deadline = Instant::now() + CALL_TIMEOUT;
2588        loop {
2589            let remaining = deadline.saturating_duration_since(Instant::now());
2590            if remaining.is_zero() {
2591                return Err(format!(
2592                    "java subprocess timed out after {}s",
2593                    CALL_TIMEOUT.as_secs()
2594                ));
2595            }
2596            match self.stdout_rx.recv_timeout(remaining) {
2597                Ok(line) => {
2598                    let trimmed = line.trim();
2599                    if trimmed.is_empty() {
2600                        continue;
2601                    }
2602                    match serde_json::from_str::<SubprocessReply>(trimmed) {
2603                        Ok(reply) if reply.ok => return Ok(reply.result),
2604                        Ok(reply) => {
2605                            return Err(reply.error.unwrap_or_else(|| "unknown java error".into()));
2606                        }
2607                        Err(_) => {
2608                            debug!(target: "self_hosted_node::java", line = trimmed, "non-protocol stdout line");
2609                        }
2610                    }
2611                }
2612                Err(RecvTimeoutError::Timeout) => {
2613                    return Err(format!(
2614                        "java subprocess timed out after {}s",
2615                        CALL_TIMEOUT.as_secs()
2616                    ));
2617                }
2618                Err(RecvTimeoutError::Disconnected) => {
2619                    return Err("java subprocess closed stdout (crashed?)".to_string());
2620                }
2621            }
2622        }
2623    }
2624
2625    fn reset(&mut self) -> Result<(), String> {
2626        self.call("{\"command\":\"reset\"}").map(|_| ())
2627    }
2628
2629    fn is_alive(&mut self) -> bool {
2630        matches!(self.child.try_wait(), Ok(None))
2631    }
2632
2633    fn shutdown(mut self) {
2634        let _ = self.stdin.write_all(b"{\"command\":\"quit\"}\n");
2635        let _ = self.stdin.flush();
2636        let deadline = Instant::now() + SHUTDOWN_GRACE;
2637        loop {
2638            match self.child.try_wait() {
2639                Ok(Some(_)) => break,
2640                Ok(None) if Instant::now() >= deadline => {
2641                    let _ = self.child.kill();
2642                    let _ = self.child.wait();
2643                    break;
2644                }
2645                Ok(None) => std::thread::sleep(Duration::from_millis(100)),
2646                Err(_) => break,
2647            }
2648        }
2649        if let Some(handle) = self.stdout_handle.take() {
2650            let _ = handle.join();
2651        }
2652        if let Some(handle) = self.stderr_handle.take() {
2653            let _ = handle.join();
2654        }
2655    }
2656}
2657
2658#[cfg(feature = "java-forge")]
2659impl JavaBridge for SubprocessBridge {
2660    fn initialize(&mut self, _assets_dir: &str) -> Result<(), String> {
2661        Ok(())
2662    }
2663
2664    fn start_game_json(&mut self, request_json: &str) -> Result<String, String> {
2665        let body = json!({ "command": "startGame", "payload": request_json });
2666        self.call(&body.to_string())
2667    }
2668
2669    fn submit_action(&mut self, session_id: &str, action_json: &str) -> Result<String, String> {
2670        let body = json!({
2671            "command": "submitAction",
2672            "sessionId": session_id,
2673            "payload": action_json,
2674        });
2675        self.call(&body.to_string())
2676    }
2677
2678    fn get_prompt(
2679        &mut self,
2680        session_id: &str,
2681        player_index: usize,
2682    ) -> Result<Option<String>, String> {
2683        let body = json!({
2684            "command": "getPrompt",
2685            "sessionId": session_id,
2686            "playerIndex": player_index,
2687        });
2688        let prompt = self.call(&body.to_string())?;
2689        Ok((!prompt.is_empty()).then_some(prompt))
2690    }
2691
2692    fn get_snapshot(&mut self, session_id: &str, viewer: Option<usize>) -> Result<String, String> {
2693        let mut body = json!({ "command": "getSnapshot", "sessionId": session_id });
2694        if let Some(viewer) = viewer {
2695            body["viewer"] = json!(viewer);
2696        }
2697        self.call(&body.to_string())
2698    }
2699
2700    fn is_game_over(&mut self, session_id: &str) -> Result<bool, String> {
2701        let body = json!({ "command": "getGameOver", "sessionId": session_id });
2702        let value = self.call(&body.to_string())?;
2703        Ok(value.trim() == "true")
2704    }
2705
2706    fn end_game(&mut self, session_id: &str) -> Result<(), String> {
2707        let body = json!({ "command": "endGame", "sessionId": session_id });
2708        self.call(&body.to_string()).map(|_| ())
2709    }
2710
2711    fn abort_game(&mut self, session_id: &str) -> Result<(), String> {
2712        let body = json!({ "command": "abortGame", "sessionId": session_id });
2713        self.call(&body.to_string()).map(|_| ())
2714    }
2715}
2716
2717#[cfg(feature = "java-forge")]
2718impl Drop for SubprocessBridge {
2719    fn drop(&mut self) {
2720        let _ = self.child.kill();
2721        let _ = self.child.wait();
2722        if let Some(handle) = self.stdout_handle.take() {
2723            let _ = handle.join();
2724        }
2725        if let Some(handle) = self.stderr_handle.take() {
2726            let _ = handle.join();
2727        }
2728    }
2729}
2730
2731#[cfg(feature = "java-forge")]
2732fn resolve_java_bin(config: &JavaRuntimeConfig) -> String {
2733    if let Some(home) = &config.java_home {
2734        let bin = home.join("bin").join("java");
2735        if bin.is_file() {
2736            return bin.to_string_lossy().to_string();
2737        }
2738    }
2739    if let Ok(home) = env::var("JAVA_HOME") {
2740        let bin = PathBuf::from(home).join("bin").join("java");
2741        if bin.is_file() {
2742            return bin.to_string_lossy().to_string();
2743        }
2744    }
2745    "java".to_string()
2746}
2747
2748#[derive(Debug, Clone, Serialize)]
2749#[serde(rename_all = "camelCase")]
2750pub struct StartGameRequest {
2751    game_id: String,
2752    variant: String,
2753    starting_life: i32,
2754    seed: u64,
2755    players: Vec<PlayerConfig>,
2756}
2757
2758#[derive(Debug, Clone, Serialize)]
2759#[serde(rename_all = "camelCase")]
2760pub struct PlayerConfig {
2761    name: String,
2762    deck: Vec<CardIdentityForJava>,
2763    commander_names: Vec<String>,
2764    ai: bool,
2765}
2766
2767#[derive(Debug, Clone, Serialize)]
2768#[serde(rename_all = "camelCase")]
2769pub struct CardIdentityForJava {
2770    name: String,
2771    set_code: Option<String>,
2772    collector_number: Option<String>,
2773    foil: bool,
2774}
2775
2776#[derive(Debug, serde::Deserialize)]
2777#[serde(rename_all = "camelCase")]
2778struct StartGameResponse {
2779    session_id: String,
2780    #[allow(dead_code)]
2781    player_indexes: Vec<usize>,
2782}
2783
2784impl StartGameRequest {
2785    pub fn new(
2786        game_id: String,
2787        variant: String,
2788        starting_life: i32,
2789        seed: u64,
2790        players: Vec<PlayerConfig>,
2791    ) -> Self {
2792        Self {
2793            game_id,
2794            variant,
2795            starting_life,
2796            seed,
2797            players,
2798        }
2799    }
2800
2801    pub fn to_json(&self) -> Result<String, serde_json::Error> {
2802        serde_json::to_string(self)
2803    }
2804}
2805
2806impl PlayerConfig {
2807    pub fn new(name: String, deck: &[DeckCardIdentity], commander_names: Vec<String>) -> Self {
2808        Self {
2809            name,
2810            deck: deck.iter().map(CardIdentityForJava::from).collect(),
2811            commander_names,
2812            ai: false,
2813        }
2814    }
2815}
2816
2817impl From<&DeckCardIdentity> for CardIdentityForJava {
2818    fn from(identity: &DeckCardIdentity) -> Self {
2819        Self {
2820            name: java_card_name(&identity.name),
2821            set_code: (!identity.set_code.is_empty()).then(|| identity.set_code.clone()),
2822            collector_number: (!identity.card_number.is_empty())
2823                .then(|| identity.card_number.clone()),
2824            foil: identity.foil.unwrap_or(false),
2825        }
2826    }
2827}
2828
2829fn java_card_name(name: &str) -> String {
2830    name.split_once(" // ")
2831        .map(|(front, _)| front.to_string())
2832        .unwrap_or_else(|| name.to_string())
2833}
2834
2835fn env_path(key: &str) -> Option<PathBuf> {
2836    env::var_os(key)
2837        .filter(|value| !value.is_empty())
2838        .map(PathBuf::from)
2839}
2840
2841fn env_sizing(key: &str, default: u64) -> Option<u64> {
2842    let value = match env::var(key) {
2843        Ok(raw) if !raw.trim().is_empty() => match raw.trim().parse::<u64>() {
2844            Ok(parsed) => parsed,
2845            Err(_) => {
2846                tracing::warn!(target: "self_hosted_node::java", key, raw, "ignoring unparseable jvm sizing");
2847                default
2848            }
2849        },
2850        _ => default,
2851    };
2852    (value > 0).then_some(value)
2853}
2854
2855fn env_classpath(key: &str) -> Vec<PathBuf> {
2856    let Some(value) = env::var_os(key) else {
2857        return Vec::new();
2858    };
2859    env::split_paths(&value).collect()
2860}
2861
2862fn require_dir(path: &Path, label: &str) -> Result<(), String> {
2863    if path.is_dir() {
2864        Ok(())
2865    } else {
2866        Err(format!("{label} does not exist: {}", path.display()))
2867    }
2868}
2869
2870fn require_file(path: &Path, label: &str) -> Result<(), String> {
2871    if path.is_file() {
2872        Ok(())
2873    } else {
2874        Err(format!("{label} does not exist: {}", path.display()))
2875    }
2876}
2877
2878/// One `Pause` line from the JVM's unified GC log.
2879struct GcPause {
2880    kind: &'static str,
2881    millis: f64,
2882    heap_after_mb: Option<u64>,
2883}
2884
2885/// Parse a unified-log GC pause, e.g.
2886/// `[12.3s][info][gc] GC(7) Pause Full (Allocation Failure) 240M->171M(247M) 226.856ms`
2887fn parse_gc_pause(line: &str) -> Option<GcPause> {
2888    let kind = if line.contains("Pause Full") {
2889        "full"
2890    } else if line.contains("Pause Young") {
2891        "young"
2892    } else {
2893        return None;
2894    };
2895    let millis = line
2896        .rsplit(' ')
2897        .find_map(|token| token.strip_suffix("ms")?.parse::<f64>().ok())?;
2898    // "240M->171M(247M)": the figure after the arrow is what survived.
2899    let heap_after_mb = line.split_once("->").and_then(|(_, rest)| {
2900        let digits: String = rest.chars().take_while(char::is_ascii_digit).collect();
2901        digits.parse::<u64>().ok()
2902    });
2903    Some(GcPause {
2904        kind,
2905        millis,
2906        heap_after_mb,
2907    })
2908}
2909
2910#[cfg(test)]
2911mod gc_log_tests {
2912    use super::{parse_gc_pause, PlayerConfig, StartGameRequest};
2913    use manabrew_protocol::deck_dto::DeckCardIdentity;
2914
2915    #[test]
2916    fn reads_a_full_collection() {
2917        let pause = parse_gc_pause(
2918            "[117.6s][info][gc] GC(21) Pause Full (Allocation Failure) 240M->171M(247M) 226.856ms",
2919        )
2920        .expect("parsed");
2921        assert_eq!(pause.kind, "full");
2922        assert_eq!(pause.heap_after_mb, Some(171));
2923        assert!((pause.millis - 226.856).abs() < f64::EPSILON);
2924    }
2925
2926    #[test]
2927    fn reads_a_young_collection() {
2928        let pause = parse_gc_pause(
2929            "[3.1s][info][gc] GC(2) Pause Young (Allocation Failure) 216M->91M(494M) 17.254ms",
2930        )
2931        .expect("parsed");
2932        assert_eq!(pause.kind, "young");
2933        assert_eq!(pause.heap_after_mb, Some(91));
2934    }
2935
2936    /// The node asks for `-Xlog:gc*:stderr:time,uptime,level,tags`, which is a
2937    /// richer line than plain `-Xlog:gc`. Real production output.
2938    #[test]
2939    fn reads_the_format_the_node_actually_requests() {
2940        let pause = parse_gc_pause(
2941            "[2026-08-18T20:41:42.907+0000][120025.819s][info][gc] GC(1202) Pause Full (Allocation Failure) 494M->487M(494M) 545.344ms",
2942        )
2943        .expect("parsed");
2944        assert_eq!(pause.kind, "full");
2945        assert_eq!(pause.heap_after_mb, Some(487));
2946        assert!((pause.millis - 545.344).abs() < f64::EPSILON);
2947    }
2948
2949    /// `gc*` also emits region and metaspace lines that carry `->`; they must
2950    /// not be mistaken for pauses.
2951    #[test]
2952    fn ignores_the_extra_lines_gc_star_adds() {
2953        assert!(
2954            parse_gc_pause("[120025.819s][info][gc,heap] GC(1202) Eden regions: 10->0(12)")
2955                .is_none()
2956        );
2957        assert!(
2958            parse_gc_pause("[120025.819s][info][gc,metaspace] Metaspace: 48M->48M(1088M)")
2959                .is_none()
2960        );
2961    }
2962
2963    #[test]
2964    fn ignores_other_output() {
2965        assert!(parse_gc_pause("[java] LOGGER ERROR: something").is_none());
2966        assert!(parse_gc_pause("[1.0s][info][gc,init] CardTable entry size: 512").is_none());
2967    }
2968
2969    #[test]
2970    fn serializes_exact_printing_identity_for_java() {
2971        let card = DeckCardIdentity {
2972            name: "Forest".to_string(),
2973            set_code: "EOE".to_string(),
2974            card_number: "266".to_string(),
2975            foil: Some(true),
2976            ..Default::default()
2977        };
2978        let request = StartGameRequest::new(
2979            "printing-test".to_string(),
2980            "Constructed".to_string(),
2981            20,
2982            42,
2983            vec![PlayerConfig::new("Player".to_string(), &[card], Vec::new())],
2984        );
2985        let json = request.to_json().expect("serialized");
2986
2987        assert!(json.contains(
2988            r#""deck":[{"name":"Forest","setCode":"EOE","collectorNumber":"266","foil":true}]"#
2989        ));
2990    }
2991}