Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685//! MCP server exposing an Erlang node's own compiler and evaluator.//!//! The server is bound to one node, named at startup, and//! [`crate::supervisor`] keeps the connection to it up. Four tools: the node's//! status, evaluate source, load a module, read the node's log. Everything//! else a caller might want is reachable through `eval`.
use crate::connection::{ConnectionManager, ConnectionState, DEFAULT_LOCAL_NODE};use crate::eval::{self, Language};use crate::formatter::{TermFormatter, get_formatter};use crate::logs;use crate::supervisor::Supervisor;use rmcp::handler::server::common::schema_for_output;use rmcp::handler::server::tool::ToolRouter;use rmcp::handler::server::wrapper::Parameters;use rmcp::model::{CallToolResult, ContentBlock, Implementation, ServerCapabilities, ServerInfo};use rmcp::{ErrorData as McpError, ServerHandler, tool, tool_handler, tool_router};use schemars::JsonSchema;use serde::{Deserialize, Serialize};use std::sync::Arc;
/// The output format mode for displaying Erlang terms.#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]pub enum FormatterMode { /// Standard Erlang syntax. #[default] Erlang, /// Elixir syntax. Elixir, /// Gleam syntax. Gleam, /// LFE (Lisp Flavoured Erlang) syntax. Lfe,}
impl std::fmt::Display for FormatterMode { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { FormatterMode::Erlang => write!(f, "erlang"), FormatterMode::Elixir => write!(f, "elixir"), FormatterMode::Gleam => write!(f, "gleam"), FormatterMode::Lfe => write!(f, "lfe"), } }}
impl std::str::FromStr for FormatterMode { type Err = String;
fn from_str(s: &str) -> Result<Self, Self::Err> { match s.to_lowercase().as_str() { "erlang" => Ok(FormatterMode::Erlang), "elixir" => Ok(FormatterMode::Elixir), "gleam" => Ok(FormatterMode::Gleam), "lfe" => Ok(FormatterMode::Lfe), _ => Err(format!( "invalid mode '{}', expected one of: erlang, elixir, gleam, lfe", s )), } }}
/// State shared by the MCP server.pub(crate) struct ServerState { /// The connection manager holding the link to the node. pub(crate) connection_manager: Arc<ConnectionManager>, /// Keeps that link up. pub(crate) supervisor: Arc<Supervisor>, /// The output format terms are rendered in, fixed at startup. pub(crate) mode: FormatterMode, /// The formatter for `mode`. pub(crate) formatter: Box<dyn TermFormatter>,}
impl ServerState { /// Creates a new server state bound to `node`. pub(crate) fn new(mode: FormatterMode, node: String, cookie: String) -> Self { let connection_manager = Arc::new( ConnectionManager::new(DEFAULT_LOCAL_NODE.to_string(), node).with_formatter_mode(mode), ); let supervisor = Supervisor::new(cookie, Arc::clone(&connection_manager));
Self { connection_manager, supervisor, mode, formatter: get_formatter(mode), } }}
impl std::fmt::Debug for ServerState { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("ServerState") .field("mode", &self.mode) .finish_non_exhaustive() }}
/// The MCP server for Erlang Distribution.#[derive(Clone)]pub struct BeamdevServer { /// Shared server state. state: Arc<ServerState>, /// Tool router for handling tool calls. /// The router is used by rmcp macros at runtime for dispatching tool calls. #[allow(dead_code)] tool_router: ToolRouter<Self>,}
impl std::fmt::Debug for BeamdevServer { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("BeamdevServer") .field("state", &self.state) .finish_non_exhaustive() }}
#[derive(Debug, Serialize, JsonSchema)]struct NodeInfo { name: String, status: String, #[serde(skip_serializing_if = "Option::is_none")] connected_at: Option<String>,}
#[derive(Debug, Serialize, JsonSchema)]struct ListNodesResponse { nodes: Vec<NodeInfo>,}
/// Request parameters for the eval tool.#[derive(Debug, Deserialize, JsonSchema)]pub struct EvalRequest { /// The source to evaluate. Erlang expressions must end with a period. pub code: String, /// The source language: "erlang" (default) or "elixir". #[serde(default)] pub language: Option<String>, /// How long the code may run on the node before it is killed, in /// milliseconds. Defaults to `eval::DEFAULT_TIMEOUT_MS`. #[serde(default)] pub timeout_ms: Option<u64>,}
#[derive(Debug, Serialize, JsonSchema)]struct EvalResponse { node: String, language: String, #[schemars(schema_with = "any_term_schema")] result: serde_json::Value,}
fn any_term_schema(_: &mut schemars::SchemaGenerator) -> schemars::Schema { schemars::json_schema!({ "type": ["object", "array", "string", "number", "boolean", "null"], "description": "The evaluated term, rendered in the server's formatter mode." })}
/// Request parameters for the load_module tool.#[derive(Debug, Deserialize, JsonSchema)]pub(crate) struct LoadModuleRequest { /// Complete module source, including the module and export attributes. pub source: String, /// The source language: "erlang" (default) or "elixir". #[serde(default)] pub language: Option<String>, /// How long the compile may run on the node, in milliseconds. Defaults to /// `eval::DEFAULT_TIMEOUT_MS`. #[serde(default)] pub timeout_ms: Option<u64>,}
#[derive(Debug, Serialize, JsonSchema)]struct LoadModuleResponse { node: String, module: String, beam_bytes: usize,}
/// Request parameters for the get_logs tool.#[derive(Debug, Deserialize, JsonSchema)]pub(crate) struct GetLogsRequest { /// How many of the most recent events to return. Defaults to /// `logs::DEFAULT_TAIL`. #[serde(default)] pub tail: Option<usize>, /// The lowest level to include, e.g. "warning". Defaults to /// `logs::DEFAULT_LEVEL`. #[serde(default)] pub level: Option<String>, /// An Erlang regular expression an event's text must match to be returned. #[serde(default)] pub grep: Option<String>,}
#[derive(Debug, Serialize, JsonSchema)]struct LogEntryResponse { time: String, level: String, message: String,}
#[derive(Debug, Serialize, JsonSchema)]struct GetLogsResponse { node: String, entries: Vec<LogEntryResponse>,}
fn tail_size(request: &GetLogsRequest) -> usize { request.tail.unwrap_or(logs::DEFAULT_TAIL)}
fn log_level(request: &GetLogsRequest) -> &str { request.level.as_deref().unwrap_or(logs::DEFAULT_LEVEL)}
fn parse_language(language: &Option<String>) -> Result<Language, McpError> { match language { None => Ok(Language::default()), Some(name) => name .parse() .map_err(|e: String| McpError::invalid_params(e, None)), }}
#[tool_router]impl BeamdevServer { /// Creates a new MCP server bound to `node`. pub fn new(mode: FormatterMode, node: String, cookie: String) -> Self { Self { state: Arc::new(ServerState::new(mode, node, cookie)), tool_router: Self::tool_router(), } }
/// Returns a reference to the server state. pub(crate) fn state(&self) -> &ServerState { &self.state }
/// The node this server is bound to. pub(crate) fn node(&self) -> &str { self.state().supervisor.node() }
/// Connects, and keeps the connection up for the life of the server. pub fn start(&self) { self.state().supervisor.spawn(); }
/// Returns the node, connected, or the reason it could not be reached. async fn ready(&self) -> Result<&str, CallToolResult> { match self.state().supervisor.ensure_connected().await { Ok(()) => Ok(self.node()), Err(e) => Err(tool_error(format!( "{} is not reachable: {e}. It may not be running yet; \ retrying the call will try the connection again.", self.node() ))), } }
/// Report the node this server is bound to. #[tool( name = "list_nodes", output_schema = schema_for_output::<ListNodesResponse>(), description = "Report the node this server is bound to and whether the connection to it is currently up.", annotations( title = "List connected nodes", read_only_hint = true, destructive_hint = false, idempotent_hint = true, open_world_hint = false ) )] pub(crate) async fn tool_list_nodes(&self) -> Result<CallToolResult, McpError> { let node = match self.state().connection_manager.status().await { Some(status) => NodeInfo { name: status.name, connected_at: (status.state == ConnectionState::Connected) .then(|| status.connected_at.map(|at| format_duration(at.elapsed()))) .flatten(), status: status.state.to_string(), }, None => NodeInfo { name: self.node().to_string(), status: ConnectionState::Disconnected.to_string(), connected_at: None, }, };
structured(&ListNodesResponse { nodes: vec![node] }) }
/// Evaluate source on the node. #[tool( name = "eval", output_schema = schema_for_output::<EvalResponse>(), description = "Evaluate Erlang or Elixir source on a connected node and return the result. Nothing is sandboxed: the code runs with the full authority of the node, including introspection, process manipulation, and side effects. Erlang expressions must end with a period.", annotations( title = "Evaluate code", read_only_hint = false, destructive_hint = true, idempotent_hint = false, open_world_hint = true ) )] pub async fn tool_eval( &self, Parameters(request): Parameters<EvalRequest>, ) -> Result<CallToolResult, McpError> { let language = parse_language(&request.language)?; let node = match self.ready().await { Ok(node) => node, Err(failure) => return Ok(failure), };
match eval::eval( &self.state().connection_manager, &request.code, language, request.timeout_ms, ) .await { Ok(term) => structured(&EvalResponse { node: node.to_string(), language: language.to_string(), result: term_to_json(&term, self.state().formatter.as_ref()), }), Err(e) => Ok(tool_error(e.to_string())), } }
/// Compile source and load the resulting module onto the node. #[tool( name = "load_module", output_schema = schema_for_output::<LoadModuleResponse>(), description = "Compile module source on a connected node and load the result into its code server, replacing any module of the same name. Takes complete Erlang or Elixir module source.", annotations( title = "Load module", read_only_hint = false, destructive_hint = true, idempotent_hint = true, open_world_hint = true ) )] pub(crate) async fn tool_load_module( &self, Parameters(request): Parameters<LoadModuleRequest>, ) -> Result<CallToolResult, McpError> { let language = parse_language(&request.language)?; let node = match self.ready().await { Ok(node) => node, Err(failure) => return Ok(failure), };
match eval::load_module( &self.state().connection_manager, &request.source, language, request.timeout_ms, ) .await { Ok(loaded) => structured(&LoadModuleResponse { node: node.to_string(), module: loaded.module, beam_bytes: loaded.beam_bytes, }), Err(e) => Ok(tool_error(e.to_string())), } }
/// Read the log events the node has buffered since we connected to it. #[tool( name = "get_logs", output_schema = schema_for_output::<GetLogsResponse>(), description = "Return the log events a connected node has produced, newest last. Only events logged since this server connected to the node are available, and only the most recent 1000 are kept. `level` filters to that level and above, and `grep` to events whose text matches an Erlang regular expression.", annotations( title = "Get node logs", read_only_hint = true, destructive_hint = false, idempotent_hint = false, open_world_hint = true ) )] pub(crate) async fn tool_get_logs( &self, Parameters(request): Parameters<GetLogsRequest>, ) -> Result<CallToolResult, McpError> { let node = match self.ready().await { Ok(node) => node, Err(failure) => return Ok(failure), };
match logs::tail( &self.state().connection_manager, tail_size(&request), log_level(&request), request.grep.as_deref(), ) .await { Ok(entries) => structured(&GetLogsResponse { node: node.to_string(), entries: entries .into_iter() .map(|entry| LogEntryResponse { time: entry.time, level: entry.level, message: entry.message, }) .collect(), }), Err(e) => Ok(tool_error(e.to_string())), } }}
/// Build a failed tool result: the tool ran, and the caller should see why it/// failed.fn tool_error(message: impl Into<String>) -> CallToolResult { CallToolResult::error(vec![ContentBlock::text(message.into())])}
/// Build a successful tool result whose structured content is `response`.fn structured<T: Serialize>(response: &T) -> Result<CallToolResult, McpError> { match serde_json::to_value(response) { Ok(value) => Ok(CallToolResult::structured(value)), Err(e) => Err(McpError::internal_error( format!("failed to serialise the tool response: {e}"), None, )), }}
fn format_duration(duration: std::time::Duration) -> String { let secs = duration.as_secs(); if secs < 60 { format!("{}s ago", secs) } else if secs < 3600 { format!("{}m ago", secs / 60) } else if secs < 86400 { format!("{}h ago", secs / 3600) } else { format!("{}d ago", secs / 86400) }}
fn format_pid(pid: &eetf::Pid) -> String { format!("<{}.{}.{}>", pid.node.name, pid.id, pid.serial)}
/// Convert an Erlang term into JSON, rendering anything without a JSON/// equivalent through the formatter.fn term_to_json(term: &eetf::Term, formatter: &dyn TermFormatter) -> serde_json::Value { use crate::rpc;
match term { eetf::Term::Atom(atom) => serde_json::Value::String(formatter.format_atom(&atom.name)),
eetf::Term::FixInteger(i) => serde_json::Value::Number(i.value.into()),
eetf::Term::BigInteger(i) => { use std::convert::TryInto; let result: Result<i64, _> = (&i.value).try_into(); if let Ok(val) = result { serde_json::Value::Number(serde_json::Number::from(val)) } else { serde_json::Value::String(format!("{}", i.value)) } }
eetf::Term::Float(f) => { if let Some(num) = serde_json::Number::from_f64(f.value) { serde_json::Value::Number(num) } else { serde_json::Value::String(f.value.to_string()) } }
eetf::Term::Binary(b) => { if let Ok(s) = String::from_utf8(b.bytes.clone()) { serde_json::Value::String(s) } else { serde_json::Value::String(formatter.format_binary(&b.bytes)) } }
// Erlang encodes small-integer lists as STRING_EXT, so they arrive as a // ByteList rather than a List. eetf::Term::ByteList(b) => { serde_json::Value::Array(b.bytes.iter().map(|byte| (*byte).into()).collect()) }
eetf::Term::List(l) => { if l.is_nil() { serde_json::Value::Array(vec![]) } else { let items: Vec<serde_json::Value> = l .elements .iter() .map(|e| term_to_json(e, formatter)) .collect(); serde_json::Value::Array(items) } }
eetf::Term::Tuple(t) => { // Special handling for {Key, Value} pairs if t.elements.len() == 2 && let Some(key) = rpc::extract_atom(&t.elements[0]) { let value = term_to_json(&t.elements[1], formatter); return serde_json::json!({ key: value }); }
// Otherwise, represent tuple as array let items: Vec<serde_json::Value> = t .elements .iter() .map(|e| term_to_json(e, formatter)) .collect(); serde_json::Value::Array(items) }
eetf::Term::Map(m) => { let mut obj = serde_json::Map::new(); for (k, v) in &m.map { let key_str = match k { eetf::Term::Atom(a) => a.name.to_string(), eetf::Term::Binary(b) => String::from_utf8_lossy(&b.bytes).to_string(), _ => formatter.format_term(k), }; obj.insert(key_str, term_to_json(v, formatter)); } serde_json::Value::Object(obj) }
eetf::Term::Pid(pid) => serde_json::Value::String(format_pid(pid)),
eetf::Term::Reference(r) => serde_json::Value::String(formatter.format_reference(r)),
_ => serde_json::Value::String(formatter.format_term(term)), }}
#[tool_handler]impl ServerHandler for BeamdevServer { fn get_info(&self) -> ServerInfo { ServerInfo::new(ServerCapabilities::builder().enable_tools().build()) .with_server_info(Implementation::new( env!("CARGO_PKG_NAME"), env!("CARGO_PKG_VERSION"), )) .with_instructions( "beamdev - evaluate source on an Erlang/BEAM node with \ `eval` or compile and load a module with `load_module`. Introspection, tracing, \ and debugging are all reachable by evaluating the node's own functions. Code \ runs unsandboxed, so use this against development nodes only.\n\n\ `get_logs` returns what the node has logged since this server connected to it, \ which is where a crash or a warning raised by something you evaluated will \ show up.\n\n\ The node is fixed at startup and the connection to it is kept up for you, so \ `eval` works without any setup and there is nothing to connect or configure. \ A node under development restarts often; the connection is re-established \ when it comes back, and a reconnect is logged, so a `get_logs` that starts \ with a reconnect notice means the state you were working with is gone. A tool \ that reports the node unreachable is worth retrying once the node is up \ again. `list_nodes` reports which node this is and whether it is currently \ connected.", ) }}
#[cfg(test)]#[allow(clippy::unwrap_used, clippy::expect_used)]mod tests { use super::*;
fn test_server() -> BeamdevServer { BeamdevServer::new( FormatterMode::Erlang, "fake@localhost".to_string(), "cookie".to_string(), ) }
fn tool_names() -> Vec<String> { let server = test_server(); let mut names: Vec<String> = server .tool_router .list_all() .into_iter() .map(|tool| tool.name.to_string()) .collect(); names.sort(); names }
#[test] fn the_server_exposes_four_tools() { assert_eq!( tool_names(), vec!["eval", "get_logs", "list_nodes", "load_module"] ); }
#[tokio::test] async fn a_node_that_was_never_reached_still_has_a_status() { let response = test_server().tool_list_nodes().await.unwrap(); let nodes = response.structured_content.unwrap();
assert_eq!( nodes, serde_json::json!({ "nodes": [{"name": "fake@localhost", "status": "disconnected"}] }) ); }
#[test] fn get_logs_defaults_its_tail_and_level() { let request: GetLogsRequest = serde_json::from_value(serde_json::json!({})).unwrap(); assert_eq!(tail_size(&request), logs::DEFAULT_TAIL); assert_eq!(log_level(&request), logs::DEFAULT_LEVEL); }
#[test] fn every_output_schema_property_is_an_object() { let server = test_server(); for tool in server.tool_router.list_all() { let schema = tool.output_schema.expect("tool has an output schema"); let properties = schema .get("properties") .and_then(serde_json::Value::as_object) .expect("output schema has properties"); for (name, subschema) in properties { assert!( subschema.is_object(), "{}.{} is not an object schema: {}", tool.name, name, subschema ); } } }
#[test] fn a_language_defaults_to_erlang_and_rejects_the_unknown() { assert_eq!(parse_language(&None).unwrap(), Language::Erlang); assert_eq!( parse_language(&Some("elixir".to_string())).unwrap(), Language::Elixir ); assert!(parse_language(&Some("gleam".to_string())).is_err()); }
#[test] fn a_byte_list_renders_as_numbers_not_a_string() { let formatter = get_formatter(FormatterMode::Erlang); let term = eetf::Term::ByteList(eetf::ByteList::from(vec![2, 4, 6])); assert_eq!( term_to_json(&term, formatter.as_ref()), serde_json::json!([2, 4, 6]) ); }
#[test] fn a_duration_reads_as_time_since() { assert_eq!(format_duration(std::time::Duration::from_secs(5)), "5s ago"); assert_eq!( format_duration(std::time::Duration::from_secs(3700)), "1h ago" ); }}