Current section

Files

Jump to
distribute src distribute@cluster_monitor.erl
Raw

src/distribute@cluster_monitor.erl

-module(distribute@cluster_monitor).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/distribute/cluster_monitor.gleam").
-export([start_observed/1, start/0, subscribe/2, unsubscribe/2]).
-export_type([cluster_event/0, message/0, subscriber/0, state/0]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
?MODULEDOC(
" Cluster-event monitor.\n"
"\n"
" Wraps Erlang's `net_kernel:monitor_nodes/1` in a typed actor.\n"
" Subscribers register a `Subject(ClusterEvent)` and receive\n"
" `NodeUp(name)` / `NodeDown(name)` events for the lifetime of the\n"
" monitor. Subscribers are pruned proactively on owner death via\n"
" `process.monitor`, so the subscriber list is bounded under churn\n"
" regardless of cluster activity.\n"
).
-type cluster_event() :: {node_up, binary()} | {node_down, binary()}.
-opaque message() :: {subscribe, gleam@erlang@process:subject(cluster_event())} |
{unsubscribe, gleam@erlang@process:subject(cluster_event())} |
{subscriber_down, gleam@erlang@process:pid_()} |
{classified_node_event, cluster_event()} |
{unknown_message, gleam@dynamic:dynamic_()}.
-type subscriber() :: {subscriber,
gleam@erlang@process:subject(cluster_event()),
gleam@erlang@process:pid_(),
gleam@erlang@process:monitor()}.
-type state() :: {state, list(subscriber())}.
-file("src/distribute/cluster_monitor.gleam", 118).
-spec handle_message(state(), message(), fun((gleam@dynamic:dynamic_()) -> nil)) -> gleam@otp@actor:next(state(), message()).
handle_message(State, Msg, On_unknown_msg) ->
case Msg of
{subscribe, Sub} ->
case gleam@erlang@process:subject_owner(Sub) of
{ok, Pid} ->
case gleam@list:any(
erlang:element(2, State),
fun(S) -> erlang:element(2, S) =:= Sub end
) of
true ->
gleam@otp@actor:continue(State);
false ->
Mon = gleam@erlang@process:monitor(Pid),
Subscriber = {subscriber, Sub, Pid, Mon},
gleam@otp@actor:continue(
{state, [Subscriber | erlang:element(2, State)]}
)
end;
{error, nil} ->
gleam@otp@actor:continue(State)
end;
{unsubscribe, Sub@1} ->
Kept = gleam@list:filter(
erlang:element(2, State),
fun(S@1) -> case erlang:element(2, S@1) =:= Sub@1 of
true ->
gleam@erlang@process:demonitor_process(
erlang:element(4, S@1)
),
false;
false ->
true
end end
),
gleam@otp@actor:continue({state, Kept});
{subscriber_down, Pid@1} ->
Kept@1 = gleam@list:filter(
erlang:element(2, State),
fun(S@2) -> erlang:element(3, S@2) /= Pid@1 end
),
gleam@otp@actor:continue({state, Kept@1});
{classified_node_event, Event} ->
gleam@list:each(
erlang:element(2, State),
fun(S@3) ->
gleam@erlang@process:send(erlang:element(2, S@3), Event)
end
),
gleam@otp@actor:continue(State);
{unknown_message, Dyn} ->
On_unknown_msg(Dyn),
gleam@otp@actor:continue(State)
end.
-file("src/distribute/cluster_monitor.gleam", 191).
-spec classify_event(gleam@dynamic:dynamic_()) -> {ok, cluster_event()} |
{error, nil}.
classify_event(Dyn) ->
case cluster_ffi:decode_node_event(Dyn) of
{ok, {<<"nodeup"/utf8>>, Name}} ->
{ok, {node_up, Name}};
{ok, {<<"nodedown"/utf8>>, Name@1}} ->
{ok, {node_down, Name@1}};
_ ->
{error, nil}
end.
-file("src/distribute/cluster_monitor.gleam", 72).
?DOC(
" Like `start`, but fires `on_unknown_msg(dyn)` whenever the monitor\n"
" receives a mailbox term it cannot classify as a Subject message or\n"
" a recognised `nodeup`/`nodedown` event. Useful as a diagnostic hook\n"
" Silent drops in cluster discovery are a debugging nightmare.\n"
).
-spec start_observed(fun((gleam@dynamic:dynamic_()) -> nil)) -> {ok,
gleam@erlang@process:subject(message())} |
{error, gleam@otp@actor:start_error()}.
start_observed(On_unknown_msg) ->
_pipe@6 = gleam@otp@actor:new_with_initialiser(
5000,
fun(Self) ->
cluster_ffi:monitor_nodes(true),
Selector = begin
_pipe = gleam_erlang_ffi:new_selector(),
_pipe@1 = gleam@erlang@process:select(_pipe, Self),
_pipe@2 = gleam@erlang@process:select_monitors(
_pipe@1,
fun(Down) -> case Down of
{process_down, _, Pid, _} ->
{subscriber_down, Pid};
{port_down, _, _, _} ->
{unknown_message,
gleam_stdlib:identity(
<<"unexpected port down"/utf8>>
)}
end end
),
gleam@erlang@process:select_other(
_pipe@2,
fun(Dyn) -> case classify_event(Dyn) of
{ok, Event} ->
{classified_node_event, Event};
{error, nil} ->
{unknown_message, Dyn}
end end
)
end,
_pipe@3 = gleam@otp@actor:initialised({state, []}),
_pipe@4 = gleam@otp@actor:selecting(_pipe@3, Selector),
_pipe@5 = gleam@otp@actor:returning(_pipe@4, Self),
{ok, _pipe@5}
end
),
_pipe@7 = gleam@otp@actor:on_message(
_pipe@6,
fun(State, Msg) -> handle_message(State, Msg, On_unknown_msg) end
),
_pipe@8 = gleam@otp@actor:start(_pipe@7),
gleam@result:map(_pipe@8, fun(Started) -> erlang:element(3, Started) end).
-file("src/distribute/cluster_monitor.gleam", 64).
-spec start() -> {ok, gleam@erlang@process:subject(message())} |
{error, gleam@otp@actor:start_error()}.
start() ->
start_observed(fun(_) -> nil end).
-file("src/distribute/cluster_monitor.gleam", 204).
?DOC(
" Subscribe `listener` to receive `NodeUp`/`NodeDown` events from\n"
" `monitor`. Idempotent: subscribing the same `listener` twice\n"
" produces a single subscription (the handler dedups internally).\n"
"\n"
" See also: `unsubscribe/2`, `start/0`, `start_observed/1`.\n"
).
-spec subscribe(
gleam@erlang@process:subject(message()),
gleam@erlang@process:subject(cluster_event())
) -> nil.
subscribe(Monitor, Listener) ->
gleam@erlang@process:send(Monitor, {subscribe, Listener}).
-file("src/distribute/cluster_monitor.gleam", 211).
?DOC(
" Unsubscribe `listener` from `monitor`. The corresponding\n"
" `process.monitor` is demonitored so we no longer hear about the\n"
" owner's death.\n"
).
-spec unsubscribe(
gleam@erlang@process:subject(message()),
gleam@erlang@process:subject(cluster_event())
) -> nil.
unsubscribe(Monitor, Listener) ->
gleam@erlang@process:send(Monitor, {unsubscribe, Listener}).