use concinnity_core::components::Camera3D;
use concinnity_core::ecs::World;
use concinnity_engine::shutdown::ShutdownToken;
use concinnity_host::thread::asset_id;
use std::io::BufReader;
use std::net::{TcpListener, TcpStream};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use crate::debug::hot_reload;
use crate::debug::queue::RuntimeQueue;
use crate::debug::state::{
AssetEntry, BudgetMemory, BudgetSnapshot, BudgetThreads, CameraSnapshot, DebugState,
PressureSnapshot,
};
use crate::debug::verbs::camera::{self, CameraMotion};
use crate::frame_hook::FrameHook;
use crate::mcp::AppServer;
const CONNECTION_TIMEOUT: Duration = Duration::from_secs(10);
const SNAPSHOT_INTERVAL: u64 = 30;
pub(crate) struct DebugServer {
shared: Arc<Mutex<DebugState>>,
queue: RuntimeQueue,
frame: u64,
reload: hot_reload::HotReloadDriver,
camera_motion: Option<CameraMotion>,
}
impl DebugServer {
pub(crate) fn start(port: u16) -> std::io::Result<Self> {
let listener = TcpListener::bind(("127.0.0.1", port))?;
let state = DebugState::default();
let queue = state.queue.clone();
let shared = Arc::new(Mutex::new(state));
let shared_for_thread = Arc::clone(&shared);
std::thread::Builder::new()
.name("debug-server".to_string())
.spawn(move || serve(listener, shared_for_thread))?;
tracing::info!("debug server listening on http://127.0.0.1:{port}/mcp");
Ok(Self {
shared,
queue,
frame: 0,
reload: hot_reload::HotReloadDriver::new(),
camera_motion: None,
})
}
pub(crate) fn with_notifier(mut self, notifier: crate::editor::notify::Notifier) -> Self {
self.reload = self.reload.with_notifier(notifier);
self
}
pub(crate) fn with_world_path(mut self, handle: hot_reload::WorldPathHandle) -> Self {
self.reload = self.reload.with_world_path(handle);
self
}
pub(crate) fn with_reload_reports(mut self, reports: hot_reload::ReloadReports) -> Self {
self.reload = self.reload.with_reload_reports(reports);
self
}
}
impl DebugServer {
fn drive_jobs(&mut self, world: &mut World) {
let jobs = self.queue.take();
let handoff = concinnity_engine::live_edit::render_handoff(world);
match handoff.backend {
Some(backend) => {
for job in jobs.backend {
job(&mut *backend, handoff.texture_slots);
}
}
None => self.queue.requeue_backend(jobs.backend),
}
for job in jobs.world {
job(world, &mut self.camera_motion);
}
if let Some(motion) = self.camera_motion.take() {
self.camera_motion = camera::advance_camera_motion(motion, world);
}
}
}
impl FrameHook for DebugServer {
fn tick(&mut self, world: &mut World) {
self.frame += 1;
self.drive_jobs(world);
self.reload.drive(world);
let mut state = match self.shared.lock() {
Ok(s) => s,
Err(poisoned) => poisoned.into_inner(),
};
state.frame = self.frame;
state.streaming = concinnity_engine::ecs::streaming_stats(world).unwrap_or_default();
state.scratch = world.scratch_stats();
state.streaming_pressure =
concinnity_engine::ecs::streaming_pressure(world).map(|p| PressureSnapshot {
rss_bytes: p.rss_bytes,
budget_bytes: p.budget_bytes,
under_pressure: p.under_pressure,
});
if let (Some(threads), Some(memory)) = (
concinnity_engine::ecs::thread_budget(world),
concinnity_engine::ecs::memory_budget(world),
) {
state.budget = Some(BudgetSnapshot {
threads: BudgetThreads {
total_cores: threads.total_cores,
job_threads: concinnity_host::thread::jobs::pool().thread_count(),
},
memory: BudgetMemory {
total_ram_mib: memory.total_ram_bytes.map(|b| b / (1024 * 1024)),
budget_mib: memory.budget_mib(),
overridden: memory.overridden,
rss_mib: concinnity_engine::app::sysmem::process_resident_bytes()
.map(|b| b / (1024 * 1024)),
},
});
}
if state.shader_reload.is_none()
&& let Some(flag) = concinnity_engine::live_edit::shader_reload_flag(world)
{
state.shader_reload = Some(flag);
}
if state.reload.is_none() {
state.reload = self.reload.signals();
}
let profile = world.profile();
state.profile_systems = profile
.system_timings()
.iter()
.map(|&(name, micros)| (name.to_string(), micros))
.collect();
state.profile_allocs = profile
.system_allocs()
.iter()
.map(|&(name, allocs)| (name.to_string(), allocs))
.collect();
state.profile_frame_allocs = profile.frame_allocs();
state.profile_render = profile.render;
state.camera = world.query::<Camera3D>().next().map(|c| CameraSnapshot {
position: c.position,
yaw: c.yaw,
pitch: c.pitch,
fov_y_degrees: c.fov_y_degrees,
near: c.near,
view_distance: c.view_distance,
});
if self.frame % SNAPSHOT_INTERVAL == 1 {
state.system_count = world.system_count();
state.component_count = world.component_count();
state.systems = world
.systems()
.iter()
.map(|s| s.name().to_string())
.collect();
state.assets = world
.component_census()
.into_iter()
.map(|(discriminant, count)| AssetEntry {
discriminant,
count,
})
.collect();
if state.names.is_empty() {
state.names = std::sync::Arc::new(asset_id::name_table());
}
}
}
fn attach_shutdown(&mut self, shutdown: ShutdownToken) {
let mut state = match self.shared.lock() {
Ok(s) => s,
Err(poisoned) => poisoned.into_inner(),
};
state.shutdown_token = Some(shutdown);
}
}
fn serve(listener: TcpListener, shared: Arc<Mutex<DebugState>>) {
let server = Arc::new(AppServer::new(shared));
for stream in listener.incoming() {
let stream = match stream {
Ok(s) => s,
Err(e) => {
tracing::debug!("debug server accept error: {e}");
continue;
}
};
let server = Arc::clone(&server);
std::thread::spawn(move || {
if let Err(e) = handle_conn(stream, &server) {
tracing::debug!("debug client closed: {e}");
}
});
}
}
fn handle_conn(stream: TcpStream, server: &AppServer) -> std::io::Result<()> {
stream.set_read_timeout(Some(CONNECTION_TIMEOUT))?;
stream.set_write_timeout(Some(CONNECTION_TIMEOUT))?;
let mut input = BufReader::new(stream.try_clone()?);
let mut output = stream;
server.serve(&mut input, &mut output)
}