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