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