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