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