jevsnes.git / packages / mcp / src / proxy.rs
proxy.rsannotatedproxy.rssource209 lines · 8.6 KB · raw
1//! `native mcp`: the same MCP server, over stdio, for a client to start.
2//!
3//! A client that connects straight to the window's HTTP server has a connection
4//! that dies with the window, and how well it comes back is the client's
5//! business. Started this way instead, the client's connection is a pipe to
6//! this process, which outlives any number of windows: it keeps an HTTP client
7//! to whichever window is running, passes tool calls through, and says
8//! `tools/list_changed` whenever what is on the other side changes — the window
9//! turned dev mode on or off, went away, or came back.
10//!
11//! "Went away" needs its own line to the window. The MCP client library
12//! reconnects by itself, and the window remembers sessions, so a window replaced
13//! by a new one (dev mode off again, so fewer tools) looks from here like
14//! nothing happened. A request to the window's `/alive` is never answered and
15//! ends when the process does, and that is how this knows.
16//!
17//! With no window there are no tools. That is the truth, and it means a client
18//! may be started before the window is.
19
20use std::process::ExitCode;
21use std::sync::Arc;
22use std::time::Duration;
23
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};
37
38/// How long after finding no window before looking again.
39const LOOK_AGAIN: Duration = Duration::from_secs(1);
40
41/// The window, when there is one.
42type Window = watch::Receiver<Option<Peer<RoleClient>>>;
43
44/// What the client that started this process talks to.
45struct Proxy {
46    window: Window,
47    /// The client has said `initialized`; nothing may be sent to it before.
48    ready: Arc<Notify>,
49}
50
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}
98
99/// What the window talks to.
100struct Watcher {
101    /// Who this is on behalf of, so the window's panel names the real client.
102    introduces: Implementation,
103    client: Peer<RoleServer>,
104}
105
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}
115
116/// 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}
126
127/// 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}
167
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}