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