jevsnes.git / packages / mcp / src / proxy.rs
proxy.rsannotatedproxy.rssource209 lines · 8.6 KB · raw

native mcp: the same MCP server, over stdio, for a client to start.

A client that connects straight to the window's HTTP server has a connection that dies with the window, and how well it comes back is the client's business. Started this way instead, the client's connection is a pipe to this process, which outlives any number of windows: it keeps an HTTP client to whichever window is running, passes tool calls through, and says tools/list_changed whenever what is on the other side changes — the window turned dev mode on or off, went away, or came back.

"Went away" needs its own line to the window. The MCP client library reconnects by itself, and the window remembers sessions, so a window replaced by a new one (dev mode off again, so fewer tools) looks from here like nothing happened. A request to the window's /alive is never answered and ends when the process does, and that is how this knows.

With no window there are no tools. That is the truth, and it means a client may be started before the window is.

20use std::process::ExitCode;
21use std::sync::Arc;
22use std::time::Duration;
24use rmcp::model::{
25    CallToolRequestParams, CallToolResponse, ClientCapabilities, ClientInfo, Implementation,
26    ListToolsResult, PaginatedRequestParams, ServerCapabilities, ServerConfig,
27};
28use rmcp::service::{NotificationContext, Peer, RequestContext};
29use rmcp::transport::streamable_http_client::StreamableHttpClientTransportConfig;
30use rmcp::transport::{StreamableHttpClientTransport, stdio};
31use rmcp::{ClientHandler, ErrorData, RoleClient, RoleServer, ServerHandler, ServiceError, ServiceExt};
32use tokio::io::{AsyncReadExt, AsyncWriteExt};
33use tokio::net::TcpStream;
34use tokio::sync::{Notify, watch};
35
36use crate::{ADDRESS, ALIVE};

How long after finding no window before looking again.

39const LOOK_AGAIN: Duration = Duration::from_secs(1);

The window, when there is one.

42type Window = watch::Receiver<Option<Peer<RoleClient>>>;

What the client that started this process talks to.

45struct Proxy {
46    window: Window,

The client has said initialized; nothing may be sent to it before.

48    ready: Arc<Notify>,
49}
51impl ServerHandler for Proxy {
52    fn get_info(&self) -> ServerConfig {
53        ServerConfig::new(
54            ServerCapabilities::builder().enable_tools().enable_tool_list_changed().build(),
55        )
56        .with_server_info(Implementation::new("jev", "0.0.0"))
57        .with_instructions(
58            "The jev window: a SNES playing A Link to the Past. The app plays it; these tools \
59             watch. No tools at all means no window is running. The dev_mode tool adds the ones \
60             that change things. Addresses and their meanings are in research/alttp-ram-map.md.",
61        )
62    }
63
64    async fn on_initialized(&self, _: NotificationContext<RoleServer>) {
65        self.ready.notify_one();
66    }
67
68    async fn list_tools(
69        &self,
70        _: Option<PaginatedRequestParams>,
71        _: RequestContext<RoleServer>,
72    ) -> Result<ListToolsResult, ErrorData> {
73        let window = self.window.borrow().clone();
74        let tools = match window {
75            Some(window) => window.list_all_tools().await.unwrap_or_default(),
76            None => Vec::new(),
77        };
78        Ok(ListToolsResult::with_all_items(tools))
79    }
80
81    async fn call_tool(
82        &self,
83        request: CallToolRequestParams,
84        _: RequestContext<RoleServer>,
85    ) -> Result<CallToolResponse, ErrorData> {
86        let window = self.window.borrow().clone();
87        let Some(window) = window else {
88            return Err(ErrorData::internal_error("no jev window is running", None));
89        };
90        match window.call_tool(request).await {
91            Ok(result) => Ok(CallToolResponse::Complete(result)),
92            // The window's own answer, as it gave it.
93            Err(ServiceError::McpError(e)) => Err(e),
94            Err(e) => Err(ErrorData::internal_error(format!("reaching the jev window: {e}"), None)),
95        }
96    }
97}

What the window talks to.

100struct Watcher {

Who this is on behalf of, so the window's panel names the real client.

102    introduces: Implementation,
103    client: Peer<RoleServer>,
104}
106impl ClientHandler for Watcher {
107    fn get_info(&self) -> ClientInfo {
108        ClientInfo::new(ClientCapabilities::default(), self.introduces.clone())
109    }
110
111    async fn on_tool_list_changed(&self, _: NotificationContext<RoleClient>) {
112        let _ = self.client.notify_tool_list_changed().await;
113    }
114}

Hold a line to the window and return when the window is gone.

117async fn gone(mut line: TcpStream) {
118    let request = format!("GET {ALIVE} HTTP/1.1\r\nHost: {ADDRESS}\r\n\r\n");
119    if line.write_all(request.as_bytes()).await.is_err() {
120        return;
121    }
122    let mut ignored = [0; 256];
123    // Nothing is ever sent, so anything but more bytes is the end of it.
124    while matches!(line.read(&mut ignored).await, Ok(1..)) {}
125}

Find a window, stay with it until it is gone, and look again, forever.

128async fn follow(
129    found: watch::Sender<Option<Peer<RoleClient>>>,
130    client: Peer<RoleServer>,
131    ready: Arc<Notify>,
132) {
133    ready.notified().await;
134    let introduces = match client.peer_info() {
135        Some(info) => Implementation::new(
136            format!("{} via jev mcp", info.client_info.name),
137            info.client_info.version.clone(),
138        ),
139        None => Implementation::new("jev mcp", "0.0.0"),
140    };
141    loop {
142        // The line first: a window that is there for the line and gone for
143        // the client is noticed, the other way round is not.
144        let Ok(line) = TcpStream::connect(ADDRESS).await else {
145            tokio::time::sleep(LOOK_AGAIN).await;
146            continue;
147        };
148        let mut config = StreamableHttpClientTransportConfig::with_uri(format!("http://{ADDRESS}/mcp"));
149        // A window that has restarted may not know the session; starting a new
150        // one is all there is to do about that.
151        config.reinit_on_expired_session = true;
152        let watcher = Watcher { introduces: introduces.clone(), client: client.clone() };
153        if let Ok(window) = watcher.serve(StreamableHttpClientTransport::from_config(config)).await {
154            found.send_replace(Some(window.peer().clone()));
155            let _ = client.notify_tool_list_changed().await;
156            let cancel = window.cancellation_token();
157            tokio::select! {
158                _ = window.waiting() => {}
159                () = gone(line) => cancel.cancel(),
160            }
161            found.send_replace(None);
162            let _ = client.notify_tool_list_changed().await;
163        }
164        tokio::time::sleep(LOOK_AGAIN).await;
165    }
166}
168pub fn run() -> ExitCode {
169    // `StreamableHttpClientTransport::from_config`, in `follow` below, builds
170    // its own plain `reqwest::Client` with no preconfigured TLS (unlike
171    // `jev-http`, which supplies `rustls::ClientConfig` directly via
172    // `use_preconfigured_tls` and so never needs this). The workspace's
173    // reqwest is `rustls-no-provider` (`../../../third-party/rust/Cargo.toml`):
174    // it links rustls but installs no default crypto provider, and reqwest
175    // 0.13's `Client::builder().build()` panics rather than picking one on
176    // its own ("No rustls crypto provider is configured") - measured
177    // 2026-09-21, `native mcp` connecting to a live window: `follow`'s first
178    // reconnect attempt panicked inside its own spawned task, silently, and
179    // the proxy answered every tool call "no jev window is running" forever
180    // after, window running and reachable the whole time. Install the same
181    // provider `jev-http` uses (pure Rust, no ring/aws-lc) once, before
182    // `follow` ever runs.
183    let _ = rustls_rustcrypto::provider().install_default();
184    let runtime = tokio::runtime::Builder::new_current_thread()
185        .enable_all()
186        .build()
187        .expect("building the tokio runtime");
188    runtime.block_on(async {
189        let (found, window) = watch::channel(None);
190        let ready = Arc::new(Notify::new());
191        // Returns once the client has introduced itself. Nothing but MCP may
192        // go to stdout from here on; stderr is free.
193        let served = match (Proxy { window, ready: Arc::clone(&ready) }).serve(stdio()).await {
194            Ok(served) => served,
195            Err(e) => {
196                eprintln!("jev mcp: {e}");
197                return ExitCode::FAILURE;
198            }
199        };
200        tokio::spawn(follow(found, served.peer().clone(), ready));
201        match served.waiting().await {
202            Ok(_) => ExitCode::SUCCESS,
203            Err(e) => {
204                eprintln!("jev mcp: {e}");
205                ExitCode::FAILURE
206            }
207        }
208    })
209}