jevsnes.git / packages / mcp / src / lib.rs
lib.rsannotatedlib.rssource763 lines · 32.4 KB · raw
1//! An MCP server inside the window, so that whoever is developing the app can
2//! ask the running console what it is doing instead of photographing it.
3//!
4//! The console belongs to the UI thread and stays there. A tool call becomes a
5//! [`Request`] on a channel, the UI thread answers it at a frame boundary, and
6//! the reply comes back on a oneshot. Nothing here touches the emulator.
7//!
8//! Who is developing the app is not who plays the game. The app's own loops
9//! play; a client watches. So the tools come in two sets: the ones that only
10//! look are always there, and the ones that change anything — the pad, the
11//! machine's state, the process — exist only while [`DevMode`] is on, which it
12//! is not until somebody turns it on. Turning it either way tells every client
13//! that its tool list has changed.
14//!
15//! How a client is told depends on which MCP it speaks, and the transport is
16//! streamable HTTP either way. Before 2026-07-28 a client holds a stream open
17//! on its session, so there are sessions, kept in a file ([`crate::sessions`])
18//! so that a client connected before the app restarted is still connected
19//! after it. From 2026-07-28 there are no sessions and a client asks to be told
20//! with `subscriptions/listen`, which it makes again itself after a restart.
21//!
22//! This crate is the whole native-only side of talking to the app from
23//! outside it: the server, `native mcp`'s stdio proxy ([`proxy`]) and the
24//! session/snapshot-name files ([`states`]). It is deliberately its OWN
25//! buck2 target, separate from `apps/native`'s tip crate — see this
26//! directory's CLAUDE.md for why that split exists and what it costs the
27//! hot-patch cycle to get wrong.
28
29pub mod proxy;
30mod sessions;
31pub mod states;
32
33use std::collections::VecDeque;
34use std::num::NonZeroU32;
35use std::sync::atomic::{AtomicBool, Ordering};
36use std::sync::{Arc, Mutex, MutexGuard, OnceLock, PoisonError, mpsc};
37use std::thread;
38use std::time::Instant;
39
40use base64::Engine;
41use console::SnesButton;
42use panels::pad::Confidence;
43use rmcp::handler::server::router::tool::ToolRouter;
44use rmcp::handler::server::wrapper::Parameters;
45use rmcp::handler::server::tool::ToolCallContext;
46use rmcp::model::{
47    CallToolRequestParams, CallToolResponse, CallToolResult, ContentBlock, Implementation,
48    InitializeRequestParams, InitializeResult, ListToolsResult, PaginatedRequestParams,
49    ServerCapabilities, ServerConfig, SubscriptionFilter, Tool,
50};
51use rmcp::service::{Peer, RequestContext, SubscriptionContext, SubscriptionSink};
52use rmcp::transport::streamable_http_server::session::local::LocalSessionManager;
53use rmcp::transport::streamable_http_server::{StreamableHttpServerConfig, StreamableHttpService};
54use rmcp::{ErrorData, RoleServer, ServerHandler, schemars, tool, tool_router};
55use tokio::sync::oneshot;
56
57use crate::sessions::FileSessionStore;
58
59/// Where the server listens. Loopback only: this drives the pad and reads
60/// memory, and is for the machine the window is on.
61pub const ADDRESS: &str = "127.0.0.1:7637";
62
63/// A request here is never answered, only dropped when the window goes.
64pub const ALIVE: &str = "/alive";
65
66/// The longest a single `press` may run, in frames: one minute of game time.
67/// A tool call that outlives its client's patience reports nothing at all.
68const LONGEST_PRESS: u32 = 3600;
69
70/// Work RAM as the SNES addresses it.
71const WRAM_BUS: std::ops::Range<u32> = 0x7E0000..0x800000;
72
73/// How many tool calls the window keeps to show.
74const CALLS_KEPT: usize = 200;
75
76/// Whether the server is there to be talked to.
77#[derive(Debug, Default)]
78pub enum Listening {
79    #[default]
80    Starting,
81    On,
82    /// Why not: the port was taken, most likely by another window.
83    Failed(String),
84}
85
86#[derive(Debug)]
87pub enum Outcome {
88    Running,
89    Done,
90    Failed(String),
91}
92
93/// One tool call, as it arrived.
94#[derive(Debug)]
95pub struct Call {
96    id: u64,
97    pub at: Instant,
98    pub tool: String,
99    /// The arguments as the JSON they came in as; empty when there were none.
100    pub arguments: String,
101    pub outcome: Outcome,
102}
103
104/// Whether the tools that change things exist. Off when the window opens,
105/// whatever it was when the last one closed.
106#[derive(Debug, Default)]
107pub struct DevMode(AtomicBool);
108
109impl DevMode {
110    pub fn is_on(&self) -> bool {
111        self.0.load(Ordering::Relaxed)
112    }
113}
114
115/// What the window shows about the server. The server thread writes it and
116/// the UI thread reads it.
117#[derive(Debug, Default)]
118pub struct Activity {
119    pub listening: Listening,
120    pub dev_mode: Arc<DevMode>,
121    /// The client that last introduced itself, by name and version. A client
122    /// that was already connected before this window started never does.
123    pub client: Option<String>,
124    pub last_heard: Option<Instant>,
125    /// Oldest first, the newest [`CALLS_KEPT`].
126    pub calls: VecDeque<Call>,
127    next_id: u64,
128}
129
130/// A panic elsewhere while the record was held leaves a record that is still
131/// worth showing, so poisoning is not an error here.
132pub fn record(activity: &Mutex<Activity>) -> MutexGuard<'_, Activity> {
133    activity.lock().unwrap_or_else(PoisonError::into_inner)
134}
135
136/// An RGBA picture, rows top to bottom.
137pub struct Picture {
138    pub width: u32,
139    pub height: u32,
140    pub rgba: Vec<u8>,
141}
142
143/// What the server asks of the UI thread.
144pub enum Request {
145    State(oneshot::Sender<alttp::State>),
146    /// The game's own picture, at the console's resolution.
147    Frame(oneshot::Sender<Picture>),
148    /// The whole window as drawn, panels included.
149    Window(oneshot::Sender<Picture>),
150    Wram {
151        /// Offset into work RAM; in range by construction.
152        at: usize,
153        len: usize,
154        reply: oneshot::Sender<Vec<u8>>,
155    },
156    /// The names of the snapshots there are.
157    States(oneshot::Sender<Vec<String>>),
158    /// What the zbanks bot says about itself (`zbanks::Report`), already
159    /// serialized to JSON by the app. Read-only.
160    Bot(oneshot::Sender<String>),
161    /// Jev's recent goal-choice history (`decisions::goal_choice::Chooser::history`),
162    /// already serialized to JSON by the app, newest first. Read-only.
163    JevHistory(oneshot::Sender<String>),
164    /// Hold `buttons` for `frames` emulated frames, release them, then reply
165    /// with the state the game is left in.
166    Press {
167        buttons: Vec<SnesButton>,
168        frames: NonZeroU32,
169        confidence: Option<Confidence>,
170        reply: oneshot::Sender<alttp::State>,
171    },
172    /// Keep the machine as it stands under a name.
173    SaveState { name: states::Name, reply: oneshot::Sender<Result<(), String>> },
174    /// Become the machine kept under a name, and reply with the game's state.
175    LoadState { name: states::Name, reply: oneshot::Sender<Result<alttp::State, String>> },
176    /// Turn the machine off and on again.
177    PowerCycle(oneshot::Sender<alttp::State>),
178    /// Keep where the game is, then replace this process with a fresh run of
179    /// the binary it was started from, which carries on from there.
180    Restart(oneshot::Sender<Result<(), String>>),
181    /// Apply a subsecond hot patch. Must run on the UI thread at a frame
182    /// boundary, like every other request here - `subsecond::apply_patch`
183    /// rewrites which code the running `App` value's patch points jump to,
184    /// and that jump table is read from the same thread that calls through
185    /// it (`apps/native/src/app.rs`'s `logic_impl`/`ui_impl`).
186    HotPatch { table: subsecond_types::JumpTable, reply: oneshot::Sender<Result<usize, String>> },
187}
188
189#[derive(Debug, Clone, Copy, serde::Deserialize, schemars::JsonSchema)]
190#[serde(rename_all = "lowercase")]
191enum Button {
192    Up,
193    Down,
194    Left,
195    Right,
196    A,
197    B,
198    X,
199    Y,
200    L,
201    R,
202    Start,
203    Select,
204}
205
206impl From<Button> for SnesButton {
207    fn from(button: Button) -> Self {
208        match button {
209            Button::Up => Self::Up,
210            Button::Down => Self::Down,
211            Button::Left => Self::Left,
212            Button::Right => Self::Right,
213            Button::A => Self::A,
214            Button::B => Self::B,
215            Button::X => Self::X,
216            Button::Y => Self::Y,
217            Button::L => Self::L,
218            Button::R => Self::R,
219            Button::Start => Self::Start,
220            Button::Select => Self::Select,
221        }
222    }
223}
224
225#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
226struct PressRequest {
227    #[schemars(description = "Buttons to hold together. Empty lets the game run with the pad idle.")]
228    buttons: Vec<Button>,
229    #[schemars(description = "Emulated frames to hold them for; the game runs at 60 a second.")]
230    frames: NonZeroU32,
231    #[schemars(description = "How sure you are of this press, as a probability from 0 to 1. The window shows it on the controller.")]
232    confidence: Option<f32>,
233}
234
235#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
236struct WramRequest {
237    #[schemars(description = "SNES bus address in hex, as the RAM map writes it: \"7EF36C\" or \"$7EF36C\".")]
238    address: String,
239    #[schemars(description = "How many bytes to read.")]
240    length: u16,
241}
242
243#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
244struct DevModeRequest {
245    #[schemars(description = "true adds the tools that change things; false takes them away.")]
246    on: bool,
247}
248
249#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
250struct StateName {
251    #[schemars(description = "Lower-case letters, digits, '-' and '_'. \"resume\" is the one the app keeps for itself and starts from.")]
252    name: String,
253}
254
255#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
256struct HotPatchRequest {
257    #[schemars(description = "Path to a JSON-serialized subsecond_types::JumpTable, as tools/hotpatch writes it (patch-N.json).")]
258    jump_table: String,
259}
260
261#[derive(Debug, serde::Serialize, schemars::JsonSchema)]
262struct PatchInfo {
263    pid: u32,
264    /// `subsecond::aslr_reference()` in THIS process: `main`'s live address,
265    /// which a patch build rebases its jump table against
266    /// (research/subsecond-patch-build.md §2/§7).
267    aslr_reference: String,
268}
269
270/// Somebody to tell that the tool list changed.
271enum Listener {
272    /// A session, which lasts until telling it fails.
273    Session(Peer<RoleServer>),
274    /// A `subscriptions/listen`, which lasts until its request is cancelled.
275    Subscription { id: u64, sink: SubscriptionSink },
276}
277
278impl Listener {
279    async fn tell(&self) -> bool {
280        match self {
281            Self::Session(peer) => peer.notify_tool_list_changed().await.is_ok(),
282            Self::Subscription { sink, .. } => sink.notify_tool_list_changed().await.is_ok(),
283        }
284    }
285}
286
287#[derive(Default)]
288struct Listeners {
289    all: Vec<Listener>,
290    next_subscription: u64,
291}
292
293type Peers = Arc<Mutex<Listeners>>;
294
295/// One per session.
296#[derive(Clone)]
297struct Server {
298    requests: mpsc::Sender<Request>,
299    ui: egui::Context,
300    activity: Arc<Mutex<Activity>>,
301    dev_mode: Arc<DevMode>,
302    peers: Peers,
303    /// Set once this session's peer is among `peers`.
304    known: Arc<OnceLock<()>>,
305    watching: Arc<ToolRouter<Self>>,
306    driving: Arc<ToolRouter<Self>>,
307}
308
309impl Server {
310    /// The router a tool is in, if it exists right now. A driving tool called
311    /// while dev mode is off is not refused, it is not there.
312    fn router_of(&self, tool: &str) -> Option<&ToolRouter<Self>> {
313        if self.watching.has_route(tool) {
314            Some(&self.watching)
315        } else if self.dev_mode.is_on() && self.driving.has_route(tool) {
316            Some(&self.driving)
317        } else {
318            None
319        }
320    }
321
322    /// Any request carries the peer to reach its session by, including the
323    /// first one from a session restored after a restart, which never
324    /// introduces itself again.
325    fn meet(&self, peer: &Peer<RoleServer>) {
326        self.known.get_or_init(|| {
327            let mut peers = self.peers.lock().unwrap_or_else(PoisonError::into_inner);
328            peers.all.push(Listener::Session(peer.clone()));
329        });
330    }
331
332    fn state_name(name: &str) -> Result<states::Name, ErrorData> {
333        states::Name::new(name).ok_or_else(|| {
334            ErrorData::invalid_params(
335                format!("{name:?} is not a state name: lower-case letters, digits, '-' and '_'"),
336                None,
337            )
338        })
339    }
340
341    /// Note that a client was heard from, record what it said, and have the
342    /// window redraw to show it.
343    fn heard<T>(&self, note: impl FnOnce(&mut Activity) -> T) -> T {
344        let noted = {
345            let mut activity = record(&self.activity);
346            activity.last_heard = Some(Instant::now());
347            note(&mut activity)
348        };
349        self.ui.request_repaint();
350        noted
351    }
352
353    /// Put a request to the UI thread and wait for its answer.
354    async fn ask<T>(
355        &self,
356        request: impl FnOnce(oneshot::Sender<T>) -> Request,
357    ) -> Result<T, ErrorData> {
358        fn gone<E>(_: E) -> ErrorData {
359            ErrorData::internal_error("the window has closed", None)
360        }
361        let (reply, answer) = oneshot::channel();
362        self.requests.send(request(reply)).map_err(gone)?;
363        // The UI thread sleeps between frames; this is what it is to wake for.
364        self.ui.request_repaint();
365        answer.await.map_err(gone)
366    }
367}
368
369fn png(picture: &Picture) -> Result<CallToolResult, ErrorData> {
370    let failed = |e: png::EncodingError| ErrorData::internal_error(e.to_string(), None);
371    let mut bytes = Vec::new();
372    let mut encoder = png::Encoder::new(&mut bytes, picture.width, picture.height);
373    encoder.set_color(png::ColorType::Rgba);
374    encoder.set_depth(png::BitDepth::Eight);
375    let mut writer = encoder.write_header().map_err(failed)?;
376    writer.write_image_data(&picture.rgba).map_err(failed)?;
377    writer.finish().map_err(failed)?;
378    let encoded = base64::engine::general_purpose::STANDARD.encode(bytes);
379    Ok(CallToolResult::success(vec![ContentBlock::image(encoded, "image/png")]))
380}
381
382#[tool_router(router = watching)]
383impl Server {
384    #[tool(description = "What the game is doing and where Link is, read from work RAM as of the last frame.")]
385    async fn state(&self) -> Result<CallToolResult, ErrorData> {
386        let state = self.ask(Request::State).await?;
387        Ok(CallToolResult::success(vec![ContentBlock::json(state)?]))
388    }
389
390    #[tool(description = "The game's picture as of the last frame, as a PNG at the console's own resolution.")]
391    async fn frame(&self) -> Result<CallToolResult, ErrorData> {
392        png(&self.ask(Request::Frame).await?)
393    }
394
395    #[tool(description = "The whole app window as it is drawn, panels included, as a PNG.")]
396    async fn window(&self) -> Result<CallToolResult, ErrorData> {
397        png(&self.ask(Request::Window).await?)
398    }
399
400    #[tool(description = "Read bytes of work RAM ($7E0000-$7FFFFF), returned as hex.")]
401    async fn read_wram(
402        &self,
403        Parameters(WramRequest { address, length }): Parameters<WramRequest>,
404    ) -> Result<CallToolResult, ErrorData> {
405        let bus = u32::from_str_radix(address.trim_start_matches('$'), 16).map_err(|e| {
406            ErrorData::invalid_params(format!("address {address:?} is not hex: {e}"), None)
407        })?;
408        // Checked in this order so the sum below cannot overflow.
409        let end = WRAM_BUS.contains(&bus).then(|| bus + u32::from(length));
410        if !end.is_some_and(|end| end <= WRAM_BUS.end) {
411            return Err(ErrorData::invalid_params(
412                format!("{length} bytes at ${bus:06X} is not inside work RAM, $7E0000..$800000"),
413                None,
414            ));
415        }
416        let at = (bus - WRAM_BUS.start) as usize;
417        let bytes = self
418            .ask(|reply| Request::Wram { at, len: usize::from(length), reply })
419            .await?;
420        let hex: Vec<String> = bytes.iter().map(|b| format!("{b:02X}")).collect();
421        Ok(CallToolResult::success(vec![ContentBlock::text(hex.join(" "))]))
422    }
423
424    #[tool(description = "The names of the machine snapshots kept beside the ROM.")]
425    async fn states(&self) -> Result<CallToolResult, ErrorData> {
426        let names = self.ask(Request::States).await?;
427        Ok(CallToolResult::success(vec![ContentBlock::json(names)?]))
428    }
429
430    #[tool(description = "What the bot playing the game (zbanks/alttp, C, linked in) says about itself, as JSON: its info line, its task list (current first), the first goals on its goal list with type, node, screen, attempts and last score, how many goals there are, Link where the bot places him (its own map coordinates), how many assert_bp traps have fired, and base() reads that mapped to nothing. Read-only.")]
431    async fn bot(&self) -> Result<CallToolResult, ErrorData> {
432        let json = self.ask(Request::Bot).await?;
433        Ok(CallToolResult::success(vec![ContentBlock::text(json)]))
434    }
435
436    #[tool(description = "Jev's recent goal-choice history, newest first, bounded to the same count the Bot panel's Jev history shows (older decisions stay in the run's own JSONL log on disk, not here): for each choice, the frame, every option in the words Jev was given (upstream's own pick first), Jev's probability per option when it asked, the one picked, and how the choice was made - asked (with the exact request sent, tokens, latency and cost), reused from a recent answer, one option after look-alikes collapsed, throttled, or upstream's own pick with why. Read-only.")]
437    async fn jev_history(&self) -> Result<CallToolResult, ErrorData> {
438        let json = self.ask(Request::JevHistory).await?;
439        Ok(CallToolResult::success(vec![ContentBlock::text(json)]))
440    }
441
442    #[tool(description = "Jev's spend guards, read from the shared ledger every process writes to (so every process's questions count, not only this window's): spend so far, plainly, with no cap or denominator (there is no self-imposed lifetime cap - the only hard stop is the vendor's own out-of-credit answer); the question bucket (questions in the last minute, its size, tokens available, and when empty how many seconds until the next token returns - it refills by itself as the minute slides); the hour breaker (dollars in the last hour, its limit, and when closed how many seconds until it reopens by itself); holds in flight; whether the vendor said the account is out of credit; the last rate-limit snapshot; the ledger's path; and one line saying what, if anything, is stopping questions. Read-only: nothing here needs clearing by hand.")]
443    async fn budget(&self) -> Result<CallToolResult, ErrorData> {
444        let opened = |e: String| ErrorData::internal_error(format!("the spend ledger: {e}"), None);
445        let ledger = jev_http::Ledger::open().map_err(opened)?;
446        let status = ledger
447            .status(jev_http::ledger::now(), jev_http::Guards::from_env())
448            .map_err(opened)?;
449        // Flat, so tools/watch/budget.sh reads top-level fields.
450        let mut budget = serde_json::to_value(&status)
451            .map_err(|e| ErrorData::internal_error(format!("the guards' status: {e}"), None))?;
452        if let Some(fields) = budget.as_object_mut() {
453            fields.insert("stopped".into(), serde_json::json!(status.stopped()));
454            fields.insert("stopped_for_good".into(), serde_json::json!(status.stopped_for_good()));
455            fields.insert("line".into(), serde_json::json!(status.line()));
456            fields.insert("last_rate_limit".into(), serde_json::json!(ledger.last_rate_limit()));
457            fields.insert("ledger_path".into(), serde_json::json!(ledger.path().display().to_string()));
458        }
459        Ok(CallToolResult::success(vec![ContentBlock::json(budget)?]))
460    }
461
462    #[tool(description = "This process's pid and subsecond::aslr_reference(), for tools/hotpatch/patch.sh to build a patch against - no log scraping needed.")]
463    async fn patch_info(&self) -> Result<CallToolResult, ErrorData> {
464        let info = PatchInfo {
465            pid: std::process::id(),
466            aslr_reference: format!("{:#x}", subsecond::aslr_reference()),
467        };
468        Ok(CallToolResult::success(vec![ContentBlock::json(info)?]))
469    }
470
471    #[tool(description = "Dev mode is off when the app starts: the app plays the game and these tools only watch. Turning it on adds the tools that change things (press, save_state, load_state, power_cycle, restart), for debugging the app while building it.")]
472    async fn dev_mode(
473        &self,
474        Parameters(DevModeRequest { on }): Parameters<DevModeRequest>,
475    ) -> Result<CallToolResult, ErrorData> {
476        if self.dev_mode.0.swap(on, Ordering::Relaxed) != on {
477            // Taken out so the lock is not held while sending. One that
478            // cannot be told has gone away, and is forgotten.
479            let lock = || self.peers.lock().unwrap_or_else(PoisonError::into_inner);
480            let listeners = std::mem::take(&mut lock().all);
481            let mut reached = Vec::with_capacity(listeners.len());
482            for listener in listeners {
483                if listener.tell().await {
484                    reached.push(listener);
485                }
486            }
487            lock().all.append(&mut reached);
488            self.ui.request_repaint();
489        }
490        let now = if on { "on" } else { "off" };
491        Ok(CallToolResult::success(vec![ContentBlock::text(format!("dev mode is {now}"))]))
492    }
493}
494
495#[tool_router(router = driving)]
496impl Server {
497    #[tool(description = "Hold buttons on controller one for a number of frames, release them, and return the state the game is left in. With no buttons it just lets the game run.")]
498    async fn press(
499        &self,
500        Parameters(PressRequest { buttons, frames, confidence }): Parameters<PressRequest>,
501    ) -> Result<CallToolResult, ErrorData> {
502        let confidence = confidence
503            .map(|probability| {
504                Confidence::new(probability).ok_or_else(|| {
505                    ErrorData::invalid_params(
506                        format!("confidence is {probability}; a probability is 0 to 1"),
507                        None,
508                    )
509                })
510            })
511            .transpose()?;
512        if frames.get() > LONGEST_PRESS {
513            return Err(ErrorData::invalid_params(
514                format!("frames is {frames}; the most one call may hold is {LONGEST_PRESS}"),
515                None,
516            ));
517        }
518        let buttons = buttons.into_iter().map(SnesButton::from).collect();
519        let state = self
520            .ask(|reply| Request::Press { buttons, frames, confidence, reply })
521            .await?;
522        Ok(CallToolResult::success(vec![ContentBlock::json(state)?]))
523    }
524
525    #[tool(description = "Keep the machine exactly as it stands, under a name, beside the ROM.")]
526    async fn save_state(
527        &self,
528        Parameters(StateName { name }): Parameters<StateName>,
529    ) -> Result<CallToolResult, ErrorData> {
530        let name = Self::state_name(&name)?;
531        self.ask(|reply| Request::SaveState { name, reply })
532            .await?
533            .map_err(|e| ErrorData::internal_error(e, None))?;
534        Ok(CallToolResult::success(vec![ContentBlock::text("kept")]))
535    }
536
537    #[tool(description = "Become the machine kept under a name, and return the state the game is then in.")]
538    async fn load_state(
539        &self,
540        Parameters(StateName { name }): Parameters<StateName>,
541    ) -> Result<CallToolResult, ErrorData> {
542        let name = Self::state_name(&name)?;
543        let state = self
544            .ask(|reply| Request::LoadState { name, reply })
545            .await?
546            .map_err(|e| ErrorData::invalid_params(e, None))?;
547        Ok(CallToolResult::success(vec![ContentBlock::json(state)?]))
548    }
549
550    #[tool(description = "Turn the SNES off and on again. The cartridge's battery save is kept; where the game was is not, unless save_state kept it first.")]
551    async fn power_cycle(&self) -> Result<CallToolResult, ErrorData> {
552        let state = self.ask(Request::PowerCycle).await?;
553        Ok(CallToolResult::success(vec![ContentBlock::json(state)?]))
554    }
555
556    #[tool(description = "Restart the app on the binary it was started from, to run a new build: the game carries on from the frame it was on, this session stays connected, and dev mode is off again.")]
557    async fn restart(&self) -> Result<CallToolResult, ErrorData> {
558        self.ask(Request::Restart).await?.map_err(|e| ErrorData::internal_error(e, None))?;
559        Ok(CallToolResult::success(vec![ContentBlock::text("restarting")]))
560    }
561
562    #[tool(description = "Apply a subsecond hot patch (tools/hotpatch's patch-N.json) to the running window: no restart, the game keeps running. Only rewrites function bodies behind a subsecond::call/HotFn boundary (apps/native/src/app.rs's logic_impl/ui_impl, packages/panels' patch points) - a struct-layout change needs restart instead.")]
563    async fn hot_patch(
564        &self,
565        Parameters(HotPatchRequest { jump_table }): Parameters<HotPatchRequest>,
566    ) -> Result<CallToolResult, ErrorData> {
567        let bytes = tokio::fs::read(&jump_table).await.map_err(|e| {
568            ErrorData::invalid_params(format!("reading {jump_table}: {e}"), None)
569        })?;
570        let table: subsecond_types::JumpTable = serde_json::from_slice(&bytes).map_err(|e| {
571            ErrorData::invalid_params(format!("parsing {jump_table} as a JumpTable: {e}"), None)
572        })?;
573        let mapped = self
574            .ask(|reply| Request::HotPatch { table, reply })
575            .await?
576            .map_err(|e| ErrorData::internal_error(e, None))?;
577        Ok(CallToolResult::success(vec![ContentBlock::text(format!(
578            "patched, {mapped} address(es) mapped"
579        ))]))
580    }
581}
582
583impl ServerHandler for Server {
584    fn get_info(&self) -> ServerConfig {
585        ServerConfig::new(
586            ServerCapabilities::builder().enable_tools().enable_tool_list_changed().build(),
587        )
588        .with_server_info(Implementation::new("jev", "0.0.0"))
589        .with_instructions(
590            "The running jev window: a SNES playing A Link to the Past, played by the zbanks/alttp C bot (the bot tool says what it is doing); \
591             these tools watch. The dev_mode tool adds the ones that change things. \
592             Addresses and their meanings are in research/alttp-ram-map.md.",
593        )
594    }
595
596    async fn list_tools(
597        &self,
598        _: Option<PaginatedRequestParams>,
599        context: RequestContext<RoleServer>,
600    ) -> Result<ListToolsResult, ErrorData> {
601        self.meet(&context.peer);
602        self.heard(|_| ());
603        let mut tools = self.watching.list_all();
604        if self.dev_mode.is_on() {
605            tools.extend(self.driving.list_all());
606        }
607        Ok(ListToolsResult::with_all_items(tools))
608    }
609
610    fn accepted_subscription_filter(&self, _: &SubscriptionFilter) -> Option<SubscriptionFilter> {
611        Some(SubscriptionFilter::builder().tools_list_changed().build())
612    }
613
614    async fn listen(&self, context: SubscriptionContext) -> Result<(), ErrorData> {
615        let lock = || self.peers.lock().unwrap_or_else(PoisonError::into_inner);
616        let id = {
617            let mut peers = lock();
618            let id = peers.next_subscription;
619            peers.next_subscription += 1;
620            peers.all.push(Listener::Subscription { id, sink: context.sink().clone() });
621            id
622        };
623        self.heard(|_| ());
624        context.cancelled().await;
625        lock().all.retain(
626            |listener| !matches!(listener, Listener::Subscription { id: gone, .. } if *gone == id),
627        );
628        Ok(())
629    }
630
631    fn get_tool(&self, name: &str) -> Option<Tool> {
632        self.router_of(name)?.get(name).cloned()
633    }
634
635    async fn initialize(
636        &self,
637        request: InitializeRequestParams,
638        context: RequestContext<RoleServer>,
639    ) -> Result<InitializeResult, ErrorData> {
640        context.peer.set_peer_info(request.clone());
641        self.meet(&context.peer);
642        let client = &request.client_info;
643        self.heard(|activity| {
644            activity.client = Some(format!("{} {}", client.name, client.version));
645        });
646        self.negotiate_initialize(&request)
647    }
648
649    /// Every tool call comes through here, so this is where they are logged.
650    async fn call_tool(
651        &self,
652        request: CallToolRequestParams,
653        context: RequestContext<RoleServer>,
654    ) -> Result<CallToolResponse, ErrorData> {
655        self.meet(&context.peer);
656        let tool = request.name.to_string();
657        let arguments = match &request.arguments {
658            Some(arguments) if !arguments.is_empty() => {
659                serde_json::Value::from(arguments.clone()).to_string()
660            }
661            _ => String::new(),
662        };
663        let id = self.heard(|activity| {
664            let id = activity.next_id;
665            activity.next_id += 1;
666            if activity.calls.len() == CALLS_KEPT {
667                activity.calls.pop_front();
668            }
669            activity.calls.push_back(Call {
670                id,
671                at: Instant::now(),
672                tool,
673                arguments,
674                outcome: Outcome::Running,
675            });
676            id
677        });
678
679        let response = match self.router_of(&request.name) {
680            Some(router) => router.call(ToolCallContext::new(self, request, context)).await,
681            None => Err(ErrorData::invalid_params(
682                format!("no tool named {:?}", request.name),
683                None,
684            )),
685        };
686
687        let outcome = match &response {
688            Ok(_) => Outcome::Done,
689            Err(e) => Outcome::Failed(e.message.to_string()),
690        };
691        self.heard(|activity| {
692            // Gone already if two hundred calls arrived while this one ran.
693            if let Some(call) = activity.calls.iter_mut().find(|call| call.id == id) {
694                call.outcome = outcome;
695            }
696        });
697        response
698    }
699}
700
701/// Start the server on its own thread and return the two ends the UI thread
702/// reads: what is asked of it, and what there is to show.
703///
704/// A port that is already taken is reported and the window carries on without
705/// a server: a second window is still a window.
706pub fn serve(ui: egui::Context) -> (mpsc::Receiver<Request>, Arc<Mutex<Activity>>) {
707    let (requests, inbox) = mpsc::channel();
708    let activity = Arc::new(Mutex::new(Activity::default()));
709    let dev_mode = Arc::clone(&record(&activity).dev_mode);
710    let peers = Peers::default();
711    let (watching, driving) = (Arc::new(Server::watching()), Arc::new(Server::driving()));
712
713    let spawned = thread::Builder::new().name("mcp".into()).spawn({
714        let activity = Arc::clone(&activity);
715        move || {
716            let runtime = tokio::runtime::Builder::new_current_thread()
717                .enable_all()
718                .build()
719                .expect("building the tokio runtime");
720            let served: std::io::Result<()> = runtime.block_on(async {
721                let mut config = StreamableHttpServerConfig::default();
722                config.session_store = Some(Arc::new(FileSessionStore::open()));
723                let service: StreamableHttpService<Server, LocalSessionManager> =
724                    StreamableHttpService::new(
725                        {
726                            let (ui, activity) = (ui.clone(), Arc::clone(&activity));
727                            move || {
728                                Ok(Server {
729                                    requests: requests.clone(),
730                                    ui: ui.clone(),
731                                    activity: Arc::clone(&activity),
732                                    dev_mode: Arc::clone(&dev_mode),
733                                    peers: Arc::clone(&peers),
734                                    known: Arc::default(),
735                                    watching: Arc::clone(&watching),
736                                    driving: Arc::clone(&driving),
737                                })
738                            }
739                        },
740                        Default::default(),
741                        config,
742                    );
743                // `/alive` never answers. Whoever asks is holding the line to
744                // know the moment this process is gone, which the MCP client
745                // library hides from them by quietly reconnecting.
746                let alive = axum::routing::get(std::future::pending::<()>);
747                let router = axum::Router::new().route(ALIVE, alive).nest_service("/mcp", service);
748                let listener = tokio::net::TcpListener::bind(ADDRESS).await?;
749                record(&activity).listening = Listening::On;
750                ui.request_repaint();
751                axum::serve(listener, router).await
752            });
753            if let Err(e) = served {
754                record(&activity).listening = Listening::Failed(e.to_string());
755                ui.request_repaint();
756            }
757        }
758    });
759    if let Err(e) = spawned {
760        record(&activity).listening = Listening::Failed(e.to_string());
761    }
762    (inbox, activity)
763}