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