//! Evaluating source and loading modules on a remote node. //! //! Both work against a stock node with nothing pre-installed: the chain uses //! only stdlib modules, and every intermediate (tokens, forms, beam bytes) //! stays an EETF term as it passes back through us. use crate::connection::ConnectionManager; use crate::error::{RpcError, RpcResult}; use crate::rpc::{self, rpc_call}; use crate::server::FormatterMode; use eetf::Term; /// The source language a node evaluates. #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] pub enum Language { #[default] Erlang, Elixir, } impl std::fmt::Display for Language { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { Language::Erlang => write!(f, "erlang"), Language::Elixir => write!(f, "elixir"), } } } impl std::str::FromStr for Language { type Err = String; fn from_str(s: &str) -> Result { match s.to_lowercase().as_str() { "erlang" | "erl" => Ok(Language::Erlang), "elixir" | "ex" => Ok(Language::Elixir), other => Err(format!("unknown language: {}", other)), } } } /// The module we load onto each node to run evaluation under our control. pub(crate) const RUNNER_MODULE: &str = "beamdev_runner"; /// Source for [`RUNNER_MODULE`]. /// /// An RPC timeout bounds how long *we* wait, not how long the node works, and /// stdlib will not do the bounding for us: neither `rpc:call/5` nor /// `erpc:call/5` kills its worker when the timeout expires, so both leave the /// runaway process running. Hence a process we spawn, monitor, and kill. const RUNNER_SOURCE: &str = r#" -module(beamdev_runner). -export([run/4]). run(M, F, A, Opts) -> Timeout = maps:get(timeout, Opts), MaxHeap = maps:get(max_heap_size, Opts), Parent = self(), {Pid, Ref} = spawn_opt( fun() -> Result = try {ok, apply(M, F, A)} catch Class:Reason:Stack -> {exception, Class, Reason, Stack} end, Parent ! {self(), Result} end, [monitor, {max_heap_size, #{size => MaxHeap, kill => true, error_logger => false}}] ), receive {Pid, Result} -> erlang:demonitor(Ref, [flush]), Result; {'DOWN', Ref, process, Pid, Reason} -> {killed, Reason} after Timeout -> exit(Pid, kill), erlang:demonitor(Ref, [flush]), {killed, timeout} end. "#; /// How long a call may run on the node before it is killed. pub(crate) const DEFAULT_TIMEOUT_MS: u64 = 15_000; /// How much a call may allocate before it is killed, in words. pub(crate) const DEFAULT_MAX_HEAP_SIZE: u64 = 100_000_000; /// How much longer than the node's own budget we wait for the answer. const RPC_MARGIN_MS: u64 = 5_000; /// Compiles and loads [`RUNNER_MODULE`] onto a node. /// /// Called once per connection: the runner must be there before the first eval. pub(crate) async fn deploy_runner( connection_manager: &ConnectionManager, ) -> RpcResult { // Deliberately not via load_module: its Elixir branch evaluates, and // evaluating is what needs the runner. load_erlang_module(connection_manager, RUNNER_SOURCE, DEFAULT_TIMEOUT_MS).await } /// Runs `module:function(args)` on the node inside the runner's bounded process, /// loading the runner first if the node does not have it. /// /// Deploying lazily rather than on connect keeps a bare connection from writing /// to the node, and recovers on its own if the node restarts under us. async fn run_bounded( connection_manager: &ConnectionManager, module: &str, function: &str, args: Vec, language: Language, timeout_ms: u64, ) -> RpcResult { let generation = connection_manager.runner_generation(); let seen = *generation.lock().await; match call_runner( connection_manager, module, function, args.clone(), language, timeout_ms, ) .await { Err(RpcError::Undef { .. }) => { let mut current = generation.lock().await; if *current == seen { deploy_runner(connection_manager).await?; *current += 1; } drop(current); call_runner( connection_manager, module, function, args, language, timeout_ms, ) .await } other => other, } } async fn call_runner( connection_manager: &ConnectionManager, module: &str, function: &str, args: Vec, language: Language, timeout_ms: u64, ) -> RpcResult { let mode = connection_manager.formatter_mode(); let opts = rpc::map(vec![ ( rpc::atom("timeout"), Term::from(eetf::BigInteger::from(timeout_ms)), ), ( rpc::atom("max_heap_size"), Term::from(eetf::BigInteger::from(DEFAULT_MAX_HEAP_SIZE)), ), ]); let outcome = rpc_call( connection_manager, RUNNER_MODULE, "run", vec![ rpc::atom(module), rpc::atom(function), rpc::list(args), opts, ], Some(timeout_ms + RPC_MARGIN_MS), ) .await?; let tuple = rpc::extract_tuple(&outcome).ok_or_else(|| unexpected(mode, "runner failed", &outcome))?; match tuple { [head, value] if rpc::is_atom(head, "ok") => Ok(value.clone()), [head, reason] if rpc::is_atom(head, "killed") => Err(RpcError::Killed { reason: rpc::format_term_for_error(mode, reason), }), [head, class, reason, stack] if rpc::is_atom(head, "exception") && language == Language::Elixir && raised_undef_for(reason, stack, module, function) => { Err(RpcError::NoElixir { node: connection_manager.node().to_string(), module: module.to_string(), function: function.to_string(), }) } [head, class, reason, stack] if rpc::is_atom(head, "exception") => { Err(RpcError::RemoteException { class: rpc::format_term_for_error(mode, class), reason: describe_exception( connection_manager, language, class, reason, stack, timeout_ms, ) .await, }) } _ => Err(unexpected(mode, "runner failed", &outcome)), } } /// Recognises an `undef` raised by the call we made rather than by the code we /// asked it to run. /// /// A node without Elixir answers every `Elixir.*` call this way, and `undef` /// alone does not say which module was missing -- the top stack frame does. fn raised_undef_for(reason: &Term, stack: &Term, module: &str, function: &str) -> bool { if !rpc::is_atom(reason, "undef") { return false; } rpc::extract_list(stack) .and_then(<[Term]>::first) .and_then(rpc::extract_tuple) .is_some_and(|frame| match frame { [frame_module, frame_function, ..] => { rpc::is_atom(frame_module, module) && rpc::is_atom(frame_function, function) } _ => false, }) } /// Renders the reason an exception carries, for an error message. /// /// Elixir terms are rendered by the node itself: `Exception.format/3` goes /// through the `Inspect` protocol, so a struct's `@derive`d or hand-written /// implementation applies -- including Ecto's `redact: true`. Formatting the /// raw term here instead would dump whatever the exception was holding. async fn describe_exception( connection_manager: &ConnectionManager, language: Language, class: &Term, reason: &Term, stack: &Term, timeout_ms: u64, ) -> String { let raw = || rpc::format_term_for_error(connection_manager.formatter_mode(), reason); if language != Language::Elixir { return raw(); } let formatted = rpc_call( connection_manager, "Elixir.Exception", "format", vec![class.clone(), reason.clone(), stack.clone()], Some(timeout_ms), ) .await; match formatted.as_ref().map(rpc::extract_binary) { Ok(Some(bytes)) => String::from_utf8_lossy(bytes).into_owned(), _ => raw(), } } /// A module compiled and loaded onto a node. #[derive(Debug, Clone)] pub struct LoadedModule { pub module: String, pub beam_bytes: usize, } /// Encodes a string as an Erlang charlist, which is what `erl_scan` takes. fn charlist(s: &str) -> Term { rpc::list( s.chars() .map(|c| Term::from(eetf::FixInteger::from(c as i32))) .collect(), ) } fn is_dot(token: &Term) -> bool { rpc::extract_tuple(token) .and_then(|t| t.first()) .is_some_and(|head| rpc::is_atom(head, "dot")) } /// Splits a token stream into one chunk per form, on the trailing dot. fn split_forms(tokens: Vec) -> Vec> { let mut forms = Vec::new(); let mut current = Vec::new(); for token in tokens { let ends_form = is_dot(&token); current.push(token); if ends_form { forms.push(std::mem::take(&mut current)); } } if !current.is_empty() { forms.push(current); } forms } fn unexpected(mode: FormatterMode, what: &str, term: &Term) -> RpcError { RpcError::UnexpectedResponse { message: format!("{}: {}", what, rpc::format_term_for_error(mode, term)), } } /// Pulls `Value` out of `{ok, Value}`, reporting anything else as an error. fn expect_ok(mode: FormatterMode, what: &str, term: &Term) -> RpcResult { rpc::extract_ok_value(term) .cloned() .ok_or_else(|| unexpected(mode, what, term)) } async fn scan( connection_manager: &ConnectionManager, source: &str, timeout_ms: Option, ) -> RpcResult> { let mode = connection_manager.formatter_mode(); let scanned = rpc_call( connection_manager, "erl_scan", "string", vec![charlist(source)], timeout_ms, ) .await?; // erl_scan:string/1 returns {ok, Tokens, EndLocation}. let tuple = rpc::extract_tuple(&scanned).ok_or_else(|| unexpected(mode, "scan failed", &scanned))?; match tuple { [head, tokens, _] if rpc::is_atom(head, "ok") => rpc::extract_list(tokens) .map(<[Term]>::to_vec) .ok_or_else(|| unexpected(mode, "tokens are not a list", tokens)), _ => Err(unexpected(mode, "scan failed", &scanned)), } } /// Evaluates `source` on `node` and returns the resulting term. pub async fn eval( connection_manager: &ConnectionManager, source: &str, language: Language, timeout_ms: Option, ) -> RpcResult { let mode = connection_manager.formatter_mode(); let budget = timeout_ms.unwrap_or(DEFAULT_TIMEOUT_MS); let rpc_timeout = Some(budget + RPC_MARGIN_MS); match language { Language::Elixir => { let evaluated = run_bounded( connection_manager, "Elixir.Code", "eval_string", vec![rpc::binary_from_str(source)], language, budget, ) .await?; // Code.eval_string/1 returns {result, bindings}. match rpc::extract_tuple(&evaluated) { Some([result, _]) => Ok(result.clone()), _ => Err(unexpected(mode, "eval returned no result", &evaluated)), } } Language::Erlang => { let tokens = scan(connection_manager, source, rpc_timeout).await?; let parsed = rpc_call( connection_manager, "erl_parse", "parse_exprs", vec![rpc::list(tokens)], rpc_timeout, ) .await?; let exprs = expect_ok(mode, "parse failed", &parsed)?; // erl_eval:exprs/2 takes no function handler, so nothing is filtered. let evaluated = run_bounded( connection_manager, "erl_eval", "exprs", vec![exprs, rpc::list(vec![])], language, budget, ) .await?; // erl_eval:exprs/2 returns {value, Value, NewBindings}. match rpc::extract_tuple(&evaluated) { Some([head, value, _]) if rpc::is_atom(head, "value") => Ok(value.clone()), _ => Err(unexpected(mode, "eval returned no value", &evaluated)), } } } } /// Compiles `source` and loads the resulting module onto `node`. pub async fn load_module( connection_manager: &ConnectionManager, source: &str, language: Language, timeout_ms: Option, ) -> RpcResult { let mode = connection_manager.formatter_mode(); let budget = timeout_ms.unwrap_or(DEFAULT_TIMEOUT_MS); // Elixir defines modules as a side effect of evaluating `defmodule`, which // evaluates to {:module, Name, Beam, _}. if language == Language::Elixir { let result = eval(connection_manager, source, language, Some(budget)).await?; let defined = rpc::extract_tuple(&result) .filter(|t| t.len() == 4 && rpc::is_atom(&t[0], "module")) .ok_or_else(|| unexpected(mode, "no module defined", &result))?; let module = rpc::extract_atom(&defined[1]) .ok_or_else(|| unexpected(mode, "module name is not an atom", &defined[1]))?; let beam_bytes = rpc::extract_binary(&defined[2]) .ok_or_else(|| unexpected(mode, "beam is not a binary", &defined[2]))? .len(); return Ok(LoadedModule { module: module.to_string(), beam_bytes, }); } load_erlang_module(connection_manager, source, budget).await } /// Compiles Erlang `source` on the node and loads the module it defines. pub(crate) async fn load_erlang_module( connection_manager: &ConnectionManager, source: &str, budget: u64, ) -> RpcResult { let mode = connection_manager.formatter_mode(); let timeout_ms = Some(budget); let tokens = scan(connection_manager, source, timeout_ms).await?; let mut forms = Vec::new(); for form_tokens in split_forms(tokens) { let parsed = rpc_call( connection_manager, "erl_parse", "parse_form", vec![rpc::list(form_tokens)], timeout_ms, ) .await?; forms.push(expect_ok(mode, "parse failed", &parsed)?); } let compiled = rpc_call( connection_manager, "compile", "forms", vec![rpc::list(forms)], timeout_ms, ) .await?; // compile:forms/1 returns {ok, Module, Beam}. let tuple = rpc::extract_tuple(&compiled) .ok_or_else(|| unexpected(mode, "compile failed", &compiled))?; let (module, beam) = match tuple { [head, module, beam] if rpc::is_atom(head, "ok") => (module.clone(), beam.clone()), _ => return Err(unexpected(mode, "compile failed", &compiled)), }; let name = rpc::extract_atom(&module) .ok_or_else(|| unexpected(mode, "module name is not an atom", &module))? .to_string(); let beam_bytes = rpc::extract_binary(&beam) .ok_or_else(|| unexpected(mode, "beam is not a binary", &beam))? .len(); // The beam bytes go straight back out as the term they arrived as. let loaded = rpc_call( connection_manager, "code", "load_binary", vec![module, charlist(&format!("{}.erl", name)), beam], timeout_ms, ) .await?; match rpc::extract_tuple(&loaded) { Some([head, _]) if rpc::is_atom(head, "module") => Ok(LoadedModule { module: name, beam_bytes, }), _ => Err(unexpected(mode, "load failed", &loaded)), } } #[cfg(test)] #[allow(clippy::unwrap_used, clippy::expect_used)] mod tests { use super::*; use std::sync::Arc; fn token(kind: &str) -> Term { rpc::tuple(vec![rpc::atom(kind), Term::from(eetf::FixInteger::from(1))]) } #[test] fn charlist_encodes_each_character_as_an_integer() { let term = charlist("hi"); let elements = rpc::extract_list(&term).unwrap(); assert_eq!(elements.len(), 2); assert!(matches!(&elements[0], Term::FixInteger(i) if i.value == 'h' as i32)); } #[test] fn split_forms_breaks_on_each_dot() { let tokens = vec![ token("atom"), token("dot"), token("atom"), token("atom"), token("dot"), ]; let forms = split_forms(tokens); assert_eq!(forms.len(), 2); assert_eq!(forms[0].len(), 2); assert_eq!(forms[1].len(), 3); } #[test] fn split_forms_keeps_a_trailing_form_without_a_dot() { let forms = split_forms(vec![token("atom"), token("dot"), token("atom")]); assert_eq!(forms.len(), 2); } #[test] fn split_forms_of_nothing_is_nothing() { assert!(split_forms(vec![]).is_empty()); } use crate::fake_node; const NODE: &str = "fake@localhost"; fn int(value: i32) -> Term { Term::from(eetf::FixInteger::from(value)) } /// Answers the erlang eval chain, returning `value` from erl_eval:exprs. fn eval_chain(value: Term) -> impl Fn(&str, &str, &[Term]) -> Term + Send + 'static { move |module, function, _args| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), int(1), ]), ("erl_parse", "parse_exprs") => { rpc::tuple(vec![rpc::atom("ok"), rpc::list(vec![token("expr")])]) } (RUNNER_MODULE, "run") => rpc::tuple(vec![ rpc::atom("ok"), rpc::tuple(vec![rpc::atom("value"), value.clone(), rpc::list(vec![])]), ]), _ => rpc::tuple(vec![rpc::atom("badarg")]), } } async fn with_fake_node( handler: impl Fn(&str, &str, &[Term]) -> Term + Send + 'static, ) -> (ConnectionManager, fake_node::CallLog) { let manager = ConnectionManager::for_test(NODE); let peer = manager.connect_in_memory().await; let log = fake_node::spawn(peer, handler); (manager, log) } #[tokio::test] async fn eval_returns_the_value_erl_eval_produced() { let (manager, log) = with_fake_node(eval_chain(int(42))).await; let result = eval(&manager, "40 + 2.", Language::Erlang, Some(1_000)) .await .unwrap(); assert!(matches!(&result, Term::FixInteger(i) if i.value == 42)); assert_eq!( log.mfas(), vec![ ("erl_scan".to_string(), "string".to_string()), ("erl_parse".to_string(), "parse_exprs".to_string()), (RUNNER_MODULE.to_string(), "run".to_string()), ] ); } #[tokio::test] async fn eval_reports_a_scan_failure() { let (manager, _) = with_fake_node(|module, function, _| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("error"), rpc::atom("illegal_character"), rpc::tuple(vec![ rpc::atom("line"), Term::from(eetf::FixInteger::from(1)), ]), ]), _ => rpc::atom("unexpected"), }) .await; let err = eval(&manager, "\u{0}.", Language::Erlang, Some(1_000)) .await .unwrap_err(); assert!(err.to_string().contains("scan failed"), "{}", err); } #[tokio::test] async fn eval_reports_a_parse_failure() { let (manager, _) = with_fake_node(|module, function, _| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), Term::from(eetf::FixInteger::from(1)), ]), ("erl_parse", "parse_exprs") => { rpc::tuple(vec![rpc::atom("error"), rpc::atom("syntax_error")]) } _ => rpc::atom("unexpected"), }) .await; let err = eval(&manager, "not erlang.", Language::Erlang, Some(1_000)) .await .unwrap_err(); assert!(err.to_string().contains("parse failed"), "{}", err); } #[tokio::test] async fn eval_unwraps_the_elixir_result_from_its_bindings() { let (manager, log) = with_fake_node(|module, function, args| match (module, function) { (RUNNER_MODULE, "run") => { // The runner is asked to call Code.eval_string on our behalf. assert!(rpc::is_atom(&args[0], "Elixir.Code")); assert!(rpc::is_atom(&args[1], "eval_string")); rpc::tuple(vec![ rpc::atom("ok"), rpc::tuple(vec![int(7), rpc::list(vec![])]), ]) } _ => rpc::atom("unexpected"), }) .await; let result = eval(&manager, "3 + 4", Language::Elixir, Some(1_000)) .await .unwrap(); assert!(matches!(&result, Term::FixInteger(i) if i.value == 7)); assert_eq!( log.mfas(), vec![(RUNNER_MODULE.to_string(), "run".to_string())] ); } #[tokio::test] async fn load_module_compiles_each_form_then_loads_the_beam() { let (manager, log) = with_fake_node(move |module, function, _args| { match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), // Two forms: each a token then a dot. rpc::list(vec![ token("atom"), token("dot"), token("atom"), token("dot"), ]), int(1), ]), ("erl_parse", "parse_form") => rpc::tuple(vec![rpc::atom("ok"), token("form")]), ("compile", "forms") => rpc::tuple(vec![ rpc::atom("ok"), rpc::atom("my_mod"), rpc::binary(vec![0xBE, 0xEF]), ]), ("code", "load_binary") => { rpc::tuple(vec![rpc::atom("module"), rpc::atom("my_mod")]) } _ => rpc::atom("unexpected"), } }) .await; let loaded = load_module( &manager, "-module(my_mod).\nhi() -> ok.", Language::Erlang, Some(1_000), ) .await .unwrap(); assert_eq!(loaded.module, "my_mod"); assert_eq!(loaded.beam_bytes, 2); // One parse_form per form, not one parse_exprs for the lot. assert_eq!( log.mfas(), vec![ ("erl_scan".to_string(), "string".to_string()), ("erl_parse".to_string(), "parse_form".to_string()), ("erl_parse".to_string(), "parse_form".to_string()), ("compile".to_string(), "forms".to_string()), ("code".to_string(), "load_binary".to_string()), ] ); } #[tokio::test] async fn load_module_reports_a_compile_failure() { let (manager, _) = with_fake_node(|module, function, _| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), Term::from(eetf::FixInteger::from(1)), ]), ("erl_parse", "parse_form") => rpc::tuple(vec![rpc::atom("ok"), token("form")]), ("compile", "forms") => rpc::atom("error"), _ => rpc::atom("unexpected"), }) .await; let err = load_module(&manager, "-module(x).", Language::Erlang, Some(1_000)) .await .unwrap_err(); assert!(err.to_string().contains("compile failed"), "{}", err); } #[tokio::test] async fn eval_reports_code_the_node_killed() { let (manager, _) = with_fake_node(|module, function, _| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), int(1), ]), ("erl_parse", "parse_exprs") => { rpc::tuple(vec![rpc::atom("ok"), rpc::list(vec![token("expr")])]) } (RUNNER_MODULE, "run") => rpc::tuple(vec![rpc::atom("killed"), rpc::atom("timeout")]), _ => rpc::atom("unexpected"), }) .await; let err = eval( &manager, "timer:sleep(infinity).", Language::Erlang, Some(1_000), ) .await .unwrap_err(); assert!(err.to_string().contains("killed"), "{}", err); assert!(err.to_string().contains("timeout"), "{}", err); } #[tokio::test] async fn eval_reports_an_exception_the_code_raised() { let (manager, _) = with_fake_node(|module, function, _| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), int(1), ]), ("erl_parse", "parse_exprs") => { rpc::tuple(vec![rpc::atom("ok"), rpc::list(vec![token("expr")])]) } (RUNNER_MODULE, "run") => rpc::tuple(vec![ rpc::atom("exception"), rpc::atom("error"), rpc::atom("badarith"), rpc::list(vec![]), ]), _ => rpc::atom("unexpected"), }) .await; let err = eval(&manager, "1 / 0.", Language::Erlang, Some(1_000)) .await .unwrap_err(); assert!(err.to_string().contains("badarith"), "{}", err); } /// Answers an Elixir eval whose runner reports an exception carrying /// `reason`, and formats it with `formatted` if asked. fn elixir_exception( reason: Term, formatted: Term, ) -> impl Fn(&str, &str, &[Term]) -> Term + Send + 'static { move |module, function, _args| match (module, function) { (RUNNER_MODULE, "run") => rpc::tuple(vec![ rpc::atom("exception"), rpc::atom("error"), reason.clone(), rpc::list(vec![]), ]), ("Elixir.Exception", "format") => formatted.clone(), _ => rpc::atom("unexpected"), } } #[tokio::test] async fn an_elixir_exception_is_rendered_by_the_node() { let secret = rpc::map(vec![(rpc::atom("token"), rpc::binary_from_str("s3cr3t"))]); let (manager, log) = with_fake_node(elixir_exception( secret, rpc::binary_from_str( "** (KeyError) key :provider not found in: #Conn", ), )) .await; let err = eval(&manager, "boom", Language::Elixir, Some(1_000)) .await .unwrap_err(); assert!(err.to_string().contains("**redacted**"), "{}", err); assert!(!err.to_string().contains("s3cr3t"), "{}", err); assert_eq!( log.mfas(), vec![ (RUNNER_MODULE.to_string(), "run".to_string()), ("Elixir.Exception".to_string(), "format".to_string()), ] ); } #[tokio::test] async fn an_elixir_exception_falls_back_when_the_node_cannot_format_it() { let (manager, _) = with_fake_node(elixir_exception( rpc::atom("badarith"), rpc::tuple(vec![rpc::atom("badrpc"), rpc::atom("nodedown")]), )) .await; let err = eval(&manager, "1 / 0", Language::Elixir, Some(1_000)) .await .unwrap_err(); assert!(err.to_string().contains("badarith"), "{}", err); } /// A node without Elixir answers `undef` for the very function we called, /// which on its own says nothing about why. #[tokio::test] async fn an_undef_for_the_elixir_evaluator_names_the_missing_runtime() { let (manager, _) = with_fake_node(|module, function, _| match (module, function) { (RUNNER_MODULE, "run") => rpc::tuple(vec![ rpc::atom("exception"), rpc::atom("error"), rpc::atom("undef"), rpc::list(vec![rpc::tuple(vec![ rpc::atom("Elixir.Code"), rpc::atom("eval_string"), int(1), rpc::list(vec![]), ])]), ]), _ => rpc::atom("unexpected"), }) .await; let err = eval(&manager, "3 + 4", Language::Elixir, Some(1_000)) .await .unwrap_err(); assert!(err.to_string().contains("has no Elixir"), "{err}"); assert!(err.to_string().contains(NODE), "{err}"); } /// An `undef` raised by the evaluated code itself is the caller's problem, /// not a missing runtime. #[tokio::test] async fn an_undef_from_the_evaluated_code_is_left_alone() { let (manager, _) = with_fake_node(|module, function, _| match (module, function) { (RUNNER_MODULE, "run") => rpc::tuple(vec![ rpc::atom("exception"), rpc::atom("error"), rpc::atom("undef"), rpc::list(vec![rpc::tuple(vec![ rpc::atom("Elixir.NoSuchModule"), rpc::atom("nope"), int(0), rpc::list(vec![]), ])]), ]), _ => rpc::atom("unexpected"), }) .await; let err = eval( &manager, "NoSuchModule.nope()", Language::Elixir, Some(1_000), ) .await .unwrap_err(); assert!(!err.to_string().contains("has no Elixir"), "{err}"); } #[tokio::test] async fn an_erlang_exception_is_not_sent_to_the_node_to_format() { let (manager, log) = with_fake_node(|module, function, _| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), int(1), ]), ("erl_parse", "parse_exprs") => { rpc::tuple(vec![rpc::atom("ok"), rpc::list(vec![token("expr")])]) } (RUNNER_MODULE, "run") => rpc::tuple(vec![ rpc::atom("exception"), rpc::atom("error"), rpc::atom("badarith"), rpc::list(vec![]), ]), _ => rpc::atom("unexpected"), }) .await; eval(&manager, "1 / 0.", Language::Erlang, Some(1_000)) .await .unwrap_err(); assert!( !log.mfas() .iter() .any(|(module, _)| module == "Elixir.Exception"), "{:?}", log.mfas() ); } #[tokio::test] async fn an_error_term_is_rendered_in_the_configured_mode() { let (manager, _) = with_fake_node(|module, function, _| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), int(1), ]), ("erl_parse", "parse_exprs") => { rpc::tuple(vec![rpc::atom("ok"), rpc::list(vec![token("expr")])]) } (RUNNER_MODULE, "run") => rpc::tuple(vec![ rpc::atom("exception"), rpc::atom("error"), rpc::binary_from_str("boom"), rpc::list(vec![]), ]), _ => rpc::atom("unexpected"), }) .await; let manager = manager.with_formatter_mode(crate::server::FormatterMode::Elixir); let err = eval(&manager, "boom.", Language::Erlang, Some(1_000)) .await .unwrap_err(); assert!(err.to_string().contains("\"boom\""), "{}", err); } #[tokio::test] async fn the_runner_is_given_the_callers_time_budget() { let (manager, _) = with_fake_node(|module, function, args| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), int(1), ]), ("erl_parse", "parse_exprs") => { rpc::tuple(vec![rpc::atom("ok"), rpc::list(vec![token("expr")])]) } (RUNNER_MODULE, "run") => { let opts = rpc::extract_map(&args[3]).expect("opts map"); let timeout = opts.get(&rpc::atom("timeout")).expect("timeout"); assert!( matches!(timeout, Term::BigInteger(i) if i.value == 2_500u64.into()), "{:?}", timeout ); rpc::tuple(vec![ rpc::atom("ok"), rpc::tuple(vec![rpc::atom("value"), int(1), rpc::list(vec![])]), ]) } _ => rpc::atom("unexpected"), }) .await; eval(&manager, "1.", Language::Erlang, Some(2_500)) .await .unwrap(); } #[tokio::test] async fn the_runner_falls_back_to_one_default_budget() { let (manager, _) = with_fake_node(|module, function, args| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), int(1), ]), ("erl_parse", "parse_exprs") => { rpc::tuple(vec![rpc::atom("ok"), rpc::list(vec![token("expr")])]) } (RUNNER_MODULE, "run") => { let opts = rpc::extract_map(&args[3]).expect("opts map"); let timeout = opts.get(&rpc::atom("timeout")).expect("timeout"); assert!( matches!(timeout, Term::BigInteger(i) if i.value == DEFAULT_TIMEOUT_MS.into()), "{:?}", timeout ); rpc::tuple(vec![ rpc::atom("ok"), rpc::tuple(vec![rpc::atom("value"), int(1), rpc::list(vec![])]), ]) } _ => rpc::atom("unexpected"), }) .await; eval(&manager, "1.", Language::Erlang, None).await.unwrap(); } #[tokio::test] async fn a_missing_runner_is_deployed_then_the_call_retried() { use std::sync::atomic::{AtomicBool, Ordering}; let deployed = Arc::new(AtomicBool::new(false)); let seen = deployed.clone(); let (manager, log) = with_fake_node(move |module, function, _| { match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), int(1), ]), ("erl_parse", "parse_exprs") => { rpc::tuple(vec![rpc::atom("ok"), rpc::list(vec![token("expr")])]) } ("erl_parse", "parse_form") => rpc::tuple(vec![rpc::atom("ok"), token("form")]), ("compile", "forms") => rpc::tuple(vec![ rpc::atom("ok"), rpc::atom(RUNNER_MODULE), rpc::binary(vec![0xBE, 0xEF]), ]), ("code", "load_binary") => { seen.store(true, Ordering::SeqCst); rpc::tuple(vec![rpc::atom("module"), rpc::atom(RUNNER_MODULE)]) } (RUNNER_MODULE, "run") if !seen.load(Ordering::SeqCst) => { // The node has not got the runner yet. rpc::tuple(vec![ rpc::atom("badrpc"), rpc::tuple(vec![ rpc::atom("EXIT"), rpc::tuple(vec![rpc::atom("undef"), rpc::list(vec![])]), ]), ]) } (RUNNER_MODULE, "run") => rpc::tuple(vec![ rpc::atom("ok"), rpc::tuple(vec![rpc::atom("value"), int(9), rpc::list(vec![])]), ]), _ => rpc::atom("unexpected"), } }) .await; let result = eval(&manager, "9.", Language::Erlang, Some(1_000)) .await .unwrap(); assert!(matches!(&result, Term::FixInteger(i) if i.value == 9)); assert!(deployed.load(Ordering::SeqCst), "should have deployed"); let runs = log .mfas() .iter() .filter(|(m, f)| m == RUNNER_MODULE && f == "run") .count(); assert_eq!(runs, 2, "should call the runner once, deploy, then retry"); } /// Loading the runner again purges the version in front of it, killing the /// runner processes of calls already running -- so callers that all find it /// missing at once must put it back exactly once between them. #[tokio::test] async fn concurrent_calls_to_a_missing_runner_deploy_it_once() { use std::sync::atomic::{AtomicBool, Ordering}; let deployed = Arc::new(AtomicBool::new(false)); let seen = deployed.clone(); let (manager, log) = with_fake_node(move |module, function, _| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), int(1), ]), ("erl_parse", "parse_exprs") => { rpc::tuple(vec![rpc::atom("ok"), rpc::list(vec![token("expr")])]) } ("erl_parse", "parse_form") => rpc::tuple(vec![rpc::atom("ok"), token("form")]), ("compile", "forms") => rpc::tuple(vec![ rpc::atom("ok"), rpc::atom(RUNNER_MODULE), rpc::binary(vec![0xBE, 0xEF]), ]), ("code", "load_binary") => { seen.store(true, Ordering::SeqCst); rpc::tuple(vec![rpc::atom("module"), rpc::atom(RUNNER_MODULE)]) } (RUNNER_MODULE, "run") if !seen.load(Ordering::SeqCst) => rpc::tuple(vec![ rpc::atom("badrpc"), rpc::tuple(vec![ rpc::atom("EXIT"), rpc::tuple(vec![rpc::atom("undef"), rpc::list(vec![])]), ]), ]), (RUNNER_MODULE, "run") => rpc::tuple(vec![ rpc::atom("ok"), rpc::tuple(vec![rpc::atom("value"), int(9), rpc::list(vec![])]), ]), _ => rpc::atom("unexpected"), }) .await; let (first, second) = tokio::join!( eval(&manager, "9.", Language::Erlang, Some(1_000)), eval(&manager, "9.", Language::Erlang, Some(1_000)), ); assert!(matches!(&first.unwrap(), Term::FixInteger(i) if i.value == 9)); assert!(matches!(&second.unwrap(), Term::FixInteger(i) if i.value == 9)); let deploys = log .mfas() .iter() .filter(|(m, f)| m == "code" && f == "load_binary") .count(); assert_eq!(deploys, 1, "the runner should be deployed once"); } #[tokio::test] async fn deploy_runner_loads_the_runner_module() { let (manager, log) = with_fake_node(|module, function, _| match (module, function) { ("erl_scan", "string") => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![token("atom"), token("dot")]), int(1), ]), ("erl_parse", "parse_form") => rpc::tuple(vec![rpc::atom("ok"), token("form")]), ("compile", "forms") => rpc::tuple(vec![ rpc::atom("ok"), rpc::atom(RUNNER_MODULE), rpc::binary(vec![0xBE, 0xEF]), ]), ("code", "load_binary") => { rpc::tuple(vec![rpc::atom("module"), rpc::atom(RUNNER_MODULE)]) } _ => rpc::atom("unexpected"), }) .await; let loaded = deploy_runner(&manager).await.unwrap(); assert_eq!(loaded.module, RUNNER_MODULE); assert!( log.mfas() .contains(&("code".to_string(), "load_binary".to_string())) ); } #[test] fn language_parses_both_names_and_shorthands() { assert_eq!("erlang".parse::().unwrap(), Language::Erlang); assert_eq!("ex".parse::().unwrap(), Language::Elixir); assert!("gleam".parse::().is_err()); } }