Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515//! Capturing a node's own log output.//!//! A node logs into its own `logger`, and nothing of that reaches us: the//! distribution link carries calls and replies, not the node's log stream. So//! we install a handler on the node that buffers formatted events, and read the//! buffer back on demand.
use crate::connection::ConnectionManager;use crate::error::{RpcError, RpcResult};use crate::eval::{self, DEFAULT_TIMEOUT_MS};use crate::rpc::{self, rpc_call};use eetf::Term;
/// The module we load onto each node to buffer its log events.pub(crate) const LOGGER_MODULE: &str = "beamdev_logger";
/// Source for [`LOGGER_MODULE`].////// The buffer is an ETS table rather than a gen_server mailbox so that a/// logging process never blocks on us, and it is owned by a process of its own/// so that removing the handler is what frees it.////// No macros: this is compiled from parsed forms, which never see the/// preprocessor.////// Reports go through Elixir's `inspect` where the node has Elixir, for the/// same reason `eval` formats exceptions there: `logger_formatter` renders a/// struct with `~p`, dumping fields an `Inspect` implementation redacts.const LOGGER_SOURCE: &str = r#"-module(beamdev_logger).-export([install/0, remove/0, log/2, tail/1, own/1]).
install() -> ok = ensure_buffer(), case logger:add_handler(beamdev_logger, beamdev_logger, #{level => all}) of ok -> ok; {error, {already_exist, _}} -> ok; {error, Reason} -> {error, Reason} end.
remove() -> _ = logger:remove_handler(beamdev_logger), case whereis(beamdev_logger) of undefined -> ok; Pid -> Pid ! stop, ok end.
ensure_buffer() -> case whereis(beamdev_logger) of undefined -> Pid = spawn(beamdev_logger, own, [self()]), receive {Pid, ready} -> ok after 5000 -> {error, buffer_timeout} end; _ -> ok end.
own(Parent) -> try register(beamdev_logger, self()) of true -> beamdev_logger = ets:new(beamdev_logger, [ordered_set, public, named_table, {write_concurrency, true}]), Parent ! {self(), ready}, wait() catch error:badarg -> Parent ! {self(), ready} end.
wait() -> receive stop -> ok; _ -> wait() end.
log(Event, _Config) -> Level = maps:get(level, Event, info), Time = maps:get(time, maps:get(meta, Event, #{}), 0), Key = erlang:unique_integer([monotonic, positive]), try true = ets:insert(beamdev_logger, {Key, Level, Time, render(Event)}), trim() catch error:badarg -> ok end.
render(Event) -> case maps:get(msg, Event, undefined) of {report, Report} -> render_report(Report, Event); _ -> render_otp(Event) end.
render_report(Report, Event) -> try unicode:characters_to_binary('Elixir.Kernel':inspect(Report)) catch _:_ -> render_otp(Event) end.
render_otp(Event) -> try unicode:characters_to_binary( logger_formatter:format(Event, #{single_line => true, template => [msg]})) catch _:_ -> <<"[unformattable log event]">> end.
trim() -> case ets:info(beamdev_logger, size) of Size when is_integer(Size), Size > 1000 -> case ets:first(beamdev_logger) of '$end_of_table' -> ok; Key -> ets:delete(beamdev_logger, Key), trim() end; _ -> ok end.
tail(Opts) -> case ets:info(beamdev_logger, size) of undefined -> {error, <<"the log buffer is not installed on this node">>}; _ -> case compile_grep(maps:get(grep, Opts)) of {error, Reason} -> {error, Reason}; {ok, RE} -> {ok, collect(Opts, RE)} end end.
collect(Opts, RE) -> Floor = maps:get(level, Opts), Newest = lists:foldl( fun({_, Level, Time, Text}, Acc) -> case keep(Level, Floor, Text, RE) of true -> [entry(Level, Time, Text) | Acc]; false -> Acc end end, [], ets:tab2list(beamdev_logger) ), lists:reverse(lists:sublist(Newest, maps:get(limit, Opts))).
keep(Level, Floor, Text, RE) -> logger:compare_levels(Level, Floor) =/= lt andalso matches(Text, RE).
matches(_Text, undefined) -> true;matches(Text, RE) -> re:run(Text, RE, [{capture, none}]) =:= match.
compile_grep(undefined) -> {ok, undefined};compile_grep(Pattern) -> case re:compile(Pattern, [unicode]) of {ok, RE} -> {ok, RE}; {error, {Message, Position}} -> {error, unicode:characters_to_binary( io_lib:format("bad regex at ~p: ~ts", [Position, Message]))} end.
entry(Level, Time, Text) -> Stamp = calendar:system_time_to_rfc3339(Time, [{unit, microsecond}, {offset, "Z"}]), #{time => unicode:characters_to_binary(Stamp), level => Level, message => Text}."#;
/// The levels OTP's `logger` knows, coarsest first.const LEVELS: [&str; 8] = [ "emergency", "alert", "critical", "error", "warning", "notice", "info", "debug",];
/// How many events `get_logs` returns when the caller does not say.pub(crate) const DEFAULT_TAIL: usize = 100;
/// The level `get_logs` filters at when the caller does not say.pub(crate) const DEFAULT_LEVEL: &str = "debug";
/// One buffered log event, as the node rendered it.#[derive(Debug, Clone, PartialEq, Eq)]pub struct LogEntry { pub time: String, pub level: String, pub message: String,}
/// Loads [`LOGGER_MODULE`] onto a node and starts buffering its log events.////// Only events logged after this point are captured, so it runs on connect/// rather than on the first `get_logs`.pub async fn install(connection_manager: &ConnectionManager) -> RpcResult<()> { eval::load_erlang_module(connection_manager, LOGGER_SOURCE, DEFAULT_TIMEOUT_MS).await?; let outcome = call(connection_manager, "install", vec![]).await?; expect_ok(connection_manager, &outcome)}
/// Stops buffering and frees the buffer, leaving the node as we found it.pub async fn remove(connection_manager: &ConnectionManager) -> RpcResult<()> { let outcome = call(connection_manager, "remove", vec![]).await?; expect_ok(connection_manager, &outcome)}
/// Returns the last `limit` buffered events at or above `level`, keeping only/// those matching `grep`.////// Filtering happens on the node, so an unmatched event never crosses the wire.pub async fn tail( connection_manager: &ConnectionManager, limit: usize, level: &str, grep: Option<&str>,) -> RpcResult<Vec<LogEntry>> { if !LEVELS.contains(&level) { return Err(RpcError::UnexpectedResponse { message: format!( "unknown level '{}', expected one of: {}", level, LEVELS.join(", ") ), }); }
let opts = rpc::map(vec![ ( rpc::atom("limit"), Term::from(eetf::BigInteger::from(limit as u64)), ), (rpc::atom("level"), rpc::atom(level)), ( rpc::atom("grep"), match grep { Some(pattern) => rpc::binary_from_str(pattern), None => rpc::atom("undefined"), }, ), ]);
let outcome = match call(connection_manager, "tail", vec![opts.clone()]).await { Err(RpcError::Undef { .. }) => { install(connection_manager).await?; call(connection_manager, "tail", vec![opts]).await? } other => other?, };
let entries = expect_ok_value(connection_manager, &outcome)?; rpc::extract_list(&entries) .ok_or_else(|| unexpected(connection_manager, "logs are not a list", &entries))? .iter() .map(|term| decode_entry(connection_manager, term)) .collect()}
async fn call( connection_manager: &ConnectionManager, function: &str, args: Vec<Term>,) -> RpcResult<Term> { rpc_call( connection_manager, LOGGER_MODULE, function, args, Some(DEFAULT_TIMEOUT_MS), ) .await}
fn decode_entry(connection_manager: &ConnectionManager, term: &Term) -> RpcResult<LogEntry> { let map = rpc::extract_map(term) .ok_or_else(|| unexpected(connection_manager, "a log entry is not a map", term))?;
let field = |name: &str| -> RpcResult<String> { let value = map.get(&rpc::atom(name)).ok_or_else(|| { unexpected( connection_manager, &format!("no {name} in a log entry"), term, ) })?; match (rpc::extract_binary(value), rpc::extract_atom(value)) { (Some(bytes), _) => Ok(String::from_utf8_lossy(bytes).into_owned()), (_, Some(atom)) => Ok(atom.to_string()), _ => Err(unexpected( connection_manager, &format!("{name} is not text"), value, )), } };
Ok(LogEntry { time: field("time")?, level: field("level")?, message: field("message")?, })}
/// Reads `ok` or `{ok, Value}`, turning `{error, Reason}` into our own error.fn expect_ok_value(connection_manager: &ConnectionManager, outcome: &Term) -> RpcResult<Term> { match rpc::extract_tuple(outcome) { Some([head, value]) if rpc::is_atom(head, "ok") => Ok(value.clone()), Some([head, reason]) if rpc::is_atom(head, "error") => Err(RpcError::UnexpectedResponse { message: describe(connection_manager, reason), }), _ if rpc::is_atom(outcome, "ok") => Ok(rpc::list(vec![])), _ => Err(unexpected(connection_manager, "unexpected reply", outcome)), }}
fn expect_ok(connection_manager: &ConnectionManager, outcome: &Term) -> RpcResult<()> { expect_ok_value(connection_manager, outcome).map(|_| ())}
/// Renders a reason the node sent, preferring its own text where it gave any.fn describe(connection_manager: &ConnectionManager, reason: &Term) -> String { match rpc::extract_binary(reason) { Some(bytes) => String::from_utf8_lossy(bytes).into_owned(), None => rpc::format_term_for_error(connection_manager.formatter_mode(), reason), }}
fn unexpected(connection_manager: &ConnectionManager, what: &str, term: &Term) -> RpcError { RpcError::UnexpectedResponse { message: format!("{}: {}", what, describe(connection_manager, term)), }}
#[cfg(test)]#[allow(clippy::unwrap_used, clippy::expect_used)]mod tests { use super::*; use crate::fake_node; use crate::rpc; use eetf::Term;
const NODE: &str = "fake@localhost";
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) }
/// Answers the compile-and-load chain, then whatever `answer` returns for /// calls into the logger module itself. fn logger_node( answer: impl Fn(&str, &[Term]) -> Term + Send + 'static, ) -> 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![rpc::tuple(vec![ rpc::atom("dot"), Term::from(eetf::FixInteger::from(1)), ])]), Term::from(eetf::FixInteger::from(1)), ]), ("erl_parse", "parse_form") => rpc::tuple(vec![rpc::atom("ok"), rpc::atom("form")]), ("compile", "forms") => rpc::tuple(vec![ rpc::atom("ok"), rpc::atom(LOGGER_MODULE), rpc::binary(vec![0, 1, 2]), ]), ("code", "load_binary") => { rpc::tuple(vec![rpc::atom("module"), rpc::atom(LOGGER_MODULE)]) } (LOGGER_MODULE, name) => answer(name, args), _ => rpc::atom("unexpected"), } }
fn entry(level: &str, message: &str) -> Term { rpc::map(vec![ ( rpc::atom("time"), rpc::binary_from_str("2026-08-27T10:00:00Z"), ), (rpc::atom("level"), rpc::atom(level)), (rpc::atom("message"), rpc::binary_from_str(message)), ]) }
#[tokio::test] async fn install_loads_the_module_then_installs_the_handler() { let (manager, log) = with_fake_node(logger_node(|_, _| rpc::atom("ok"))).await;
install(&manager).await.unwrap();
let mfas = log.mfas(); assert!(mfas.contains(&("code".to_string(), "load_binary".to_string()))); assert_eq!( mfas.last(), Some(&(LOGGER_MODULE.to_string(), "install".to_string())) ); }
#[tokio::test] async fn tail_returns_the_entries_the_node_buffered() { let (manager, _) = with_fake_node(logger_node(|name, _| match name { "tail" => rpc::tuple(vec![ rpc::atom("ok"), rpc::list(vec![entry("error", "it broke")]), ]), _ => rpc::atom("ok"), })) .await;
let entries = tail(&manager, 10, "debug", None).await.unwrap();
assert_eq!( entries, vec![LogEntry { time: "2026-08-27T10:00:00Z".to_string(), level: "error".to_string(), message: "it broke".to_string(), }] ); }
#[tokio::test] async fn tail_passes_the_limit_level_and_grep_to_the_node() { let (manager, _) = with_fake_node(logger_node(|name, args| match name { "tail" => { let opts = rpc::extract_map(&args[0]).expect("tail takes a map"); assert_eq!( opts.get(&rpc::atom("grep")), Some(&rpc::binary_from_str("boom")) ); assert_eq!(opts.get(&rpc::atom("level")), Some(&rpc::atom("warning"))); assert_eq!( opts.get(&rpc::atom("limit")), Some(&Term::from(eetf::BigInteger::from(5u64))) ); rpc::tuple(vec![rpc::atom("ok"), rpc::list(vec![])]) } _ => rpc::atom("ok"), })) .await;
tail(&manager, 5, "warning", Some("boom")).await.unwrap(); }
#[tokio::test] async fn tail_reports_an_error_the_node_returns() { let (manager, _) = with_fake_node(logger_node(|name, _| match name { "tail" => rpc::tuple(vec![ rpc::atom("error"), rpc::binary_from_str("bad regex at 2: missing )"), ]), _ => rpc::atom("ok"), })) .await;
let err = tail(&manager, 10, "debug", None).await.unwrap_err(); assert!(err.to_string().contains("bad regex"), "{}", err); }
#[tokio::test] async fn tail_installs_the_logger_if_the_node_has_not_got_it() { let seen_tail = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); let counter = seen_tail.clone();
let (manager, log) = with_fake_node(logger_node(move |name, _| match name { "tail" => { if counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst) == 0 { rpc::tuple(vec![ rpc::atom("badrpc"), rpc::tuple(vec![ rpc::atom("EXIT"), rpc::tuple(vec![rpc::atom("undef"), rpc::list(vec![])]), ]), ]) } else { rpc::tuple(vec![rpc::atom("ok"), rpc::list(vec![])]) } } _ => rpc::atom("ok"), })) .await;
tail(&manager, 10, "debug", None).await.unwrap();
assert_eq!(seen_tail.load(std::sync::atomic::Ordering::SeqCst), 2); assert!( log.mfas() .contains(&(LOGGER_MODULE.to_string(), "install".to_string())) ); }
#[tokio::test] async fn remove_removes_the_handler_from_the_node() { let (manager, log) = with_fake_node(logger_node(|_, _| rpc::atom("ok"))).await;
remove(&manager).await.unwrap();
assert_eq!( log.mfas(), vec![(LOGGER_MODULE.to_string(), "remove".to_string())] ); }
#[tokio::test] async fn tail_rejects_a_level_the_node_would_not_understand() { let (manager, _) = with_fake_node(logger_node(|_, _| rpc::atom("ok"))).await;
let err = tail(&manager, 10, "chatty", None).await.unwrap_err(); assert!(err.to_string().contains("chatty"), "{}", err); }}