concinnity-dev 0.19.119

The Concinnity dev tooling library: world authoring, the in-engine editor, the debug server, docs and packaging
//! The localhost debug listener: `DebugServer` (the `FrameHook` the run loop
//! ticks), the accept / per-connection threads, and the per-frame drive of the
//! queued verb jobs plus the owned hot-reload driver. Each connection is handed
//! to `crate::mcp::AppServer`, which parses the HTTP request and answers the
//! MCP message it carried from the verb table in `super::super::catalog`.

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;

// Bound each connection's reads so a client that opens a socket and stalls
// mid-request releases its thread instead of holding it for the session.
const CONNECTION_TIMEOUT: Duration = Duration::from_secs(10);

// How often `tick` rebuilds the asset/system snapshot (in frames). The frame
// counter still advances every tick; only the heavier lists are throttled.
const SNAPSHOT_INTERVAL: u64 = 30;

// A running debug server. Implements `FrameHook`, so the run loop owns it as
// `Box<dyn FrameHook>` and ticks it each frame.
pub(crate) struct DebugServer {
    shared: Arc<Mutex<DebugState>>,
    // The snapshot's job queue, run every tick.
    queue: RuntimeQueue,
    frame: u64,
    // The asset / shader / world.jsonl reload drive. The server owns the
    // session's one driver so the `reload-assets` command can reach its
    // pending flag; the drive itself is shared with the plain `cn editor`
    // path (see `crate::debug::hot_reload::HotReloadDriver`).
    reload: hot_reload::HotReloadDriver,
    // Active camera-move motion installed by a `camera-move` call, advanced
    // once per frame by `drive_jobs` until exhausted or cleared by a
    // `camera-stop`. `None` when no motion is in progress. Main-thread only.
    camera_motion: Option<CameraMotion>,
}

impl DebugServer {
    // Bind the localhost MCP endpoint on `port` and spawn its accept thread.
    // Binds `127.0.0.1` only: the debug surface is never exposed off-box.
    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,
        })
    }

    // Report reload results through an editor session's toast queue as well as
    // the log.
    pub(crate) fn with_notifier(mut self, notifier: crate::editor::notify::Notifier) -> Self {
        self.reload = self.reload.with_notifier(notifier);
        self
    }

    // Watch the world.jsonl the handle names.
    pub(crate) fn with_world_path(mut self, handle: hot_reload::WorldPathHandle) -> Self {
        self.reload = self.reload.with_world_path(handle);
        self
    }

    // Publish each subject's latest reload outcome to `reports`.
    pub(crate) fn with_reload_reports(mut self, reports: hot_reload::ReloadReports) -> Self {
        self.reload = self.reload.with_reload_reports(reports);
        self
    }
}

impl DebugServer {
    // Run the jobs the verbs queued once per frame: backend jobs against the
    // parked backend, then world jobs, then the per-frame camera-move advance.
    // The asset / shader / world.jsonl reload passes live on `self.reload`,
    // driven separately by `tick`.
    fn drive_jobs(&mut self, world: &mut World) {
        let jobs = self.queue.take();
        // Backend jobs wait in the queue until a backend is parked; world jobs
        // never wait on one.
        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),
        }

        // tick() runs before the world step, so the Camera3DSystem step this
        // frame sees a new pose, and a freshly installed camera-move also steps
        // this same frame just below.
        for job in jobs.world {
            job(world, &mut self.camera_motion);
        }

        // Advance an in-progress camera-move one step, every frame before the
        // world step, so the renderer sees sustained motion across temporal
        // passes.
        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;

        // The queued verb jobs, then the asset / shader / world.jsonl reload
        // passes. Both run only from this hook, so a `cn run` (no debug hook)
        // never touches them.
        self.drive_jobs(world);
        self.reload.drive(world);

        let mut state = match self.shared.lock() {
            Ok(s) => s,
            // A panicked client thread should never take down the engine.
            Err(poisoned) => poisoned.into_inner(),
        };
        state.frame = self.frame;

        // Streaming counts change every frame in the early load-in, so refresh
        // them every tick -- `streaming_stats` is just a few small count loops
        // over the parked StreamingState (StreamingSystem owns the pools).
        state.streaming = concinnity_engine::ecs::streaming_stats(world).unwrap_or_default();
        state.scratch = world.scratch_stats();
        // Live RAM back-off pressure, refreshed alongside the streaming counts.
        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,
            });

        // Process thread + memory budgets (fixed at start) plus the live RSS
        // (one cheap syscall per tick, dev-only), for the `budget` query.
        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)),
                },
            });
        }
        // Opportunistically pick up the shader-reload flag the backend exposes
        // (Some only under `cn debug` on hot-reload backends); once captured,
        // the `reload-shaders` verb can fire the flag. The backend sits in
        // the world's parked slot between ticks.
        if state.shader_reload.is_none()
            && let Some(flag) = concinnity_engine::live_edit::shader_reload_flag(world)
        {
            state.shader_reload = Some(flag);
        }
        // The reload driver keeps one set of signals across re-arms, so the
        // `reload-assets` verb captures them once the driver is armed.
        if state.reload.is_none() {
            state.reload = self.reload.signals();
        }

        // The profiler snapshot is small (one entry per system + a handful of
        // render counters), so refresh it every tick like the streaming stats.
        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;

        // Active-camera pose for `camera-get`. One component read, so refresh
        // every tick like the streaming / profiler snapshots above.
        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();
            // The AssetId -> name table is the build interner snapshot; it is
            // stable once the world is built, so capture it just once.
            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);
    }
}

// Accept loop. Each client is handled on its own thread so a slow or stuck
// client never blocks another, and both threads are detached: the listener
// blocks in `accept` for the life of the process, which exits out from under
// it. Errors are logged and dropped: a debug client disconnecting is routine,
// not a fault.
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}");
            }
        });
    }
}

// One request, one response, then the connection closes with the stream.
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)
}