Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209//! 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<Self, Self::Err> { 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<LoadedModule> { // 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<Term>, language: Language, timeout_ms: u64,) -> RpcResult<Term> { 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<Term>, language: Language, timeout_ms: u64,) -> RpcResult<Term> { 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<Term>) -> Vec<Vec<Term>> { 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<Term> { rpc::extract_ok_value(term) .cloned() .ok_or_else(|| unexpected(mode, what, term))}
async fn scan( connection_manager: &ConnectionManager, source: &str, timeout_ms: Option<u64>,) -> RpcResult<Vec<Term>> { 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<u64>,) -> RpcResult<Term> { 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<u64>,) -> RpcResult<LoadedModule> { 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<LoadedModule> { 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<token: **redacted**>", ), )) .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::<Language>().unwrap(), Language::Erlang); assert_eq!("ex".parse::<Language>().unwrap(), Language::Elixir); assert!("gleam".parse::<Language>().is_err()); }}