-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}).