-module(nessie_cluster). -compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]). -export([with_name/2, with_query/2, with_logger/2, with_interval/2, with_resolver/2, discover_nodes/2, stop/2, has_ran/2, default_resolver/0, new/0, start_spec/2]). -export_type([resolver/0, dns_query/0, dns_cluster/0, node_connect_error/0, dns_cluster_state/0, message/0]). -type resolver() :: {resolver, fun((gleam@erlang@atom:atom_()) -> {ok, binary()} | {error, nil}), fun((gleam@erlang@atom:atom_()) -> {ok, gleam@erlang@node:node_()} | {error, gleam@erlang@node:connect_error()}), fun(() -> list(gleam@erlang@node:node_())), fun((binary()) -> list(binary()))}. -type dns_query() :: {dns_query, binary()} | ignore. -opaque dns_cluster() :: {dns_cluster, gleam@erlang@atom:atom_(), dns_query(), gleam@option:option(integer()), fun((binary(), binary()) -> nil), resolver()}. -type node_connect_error() :: {node_connect_error, gleam@erlang@atom:atom_(), gleam@erlang@node:connect_error()}. -type dns_cluster_state() :: {dns_cluster_state, boolean(), dns_cluster(), binary(), gleam@option:option(gleam@erlang@process:timer()), gleam@erlang@process:subject(message())}. -opaque message() :: {discover_nodes, gleam@option:option(gleam@erlang@process:subject({list(gleam@erlang@node:node_()), list(node_connect_error())})), boolean()} | {stop, gleam@erlang@process:subject(nil)} | {has_ran, gleam@erlang@process:subject(boolean())}. -spec with_name(dns_cluster(), gleam@erlang@atom:atom_()) -> dns_cluster(). with_name(Cluster, Name) -> erlang:setelement(2, Cluster, Name). -spec with_query(dns_cluster(), dns_query()) -> dns_cluster(). with_query(Cluster, Q) -> erlang:setelement(3, Cluster, Q). -spec with_logger(dns_cluster(), fun((binary(), binary()) -> nil)) -> dns_cluster(). with_logger(Cluster, Logger) -> erlang:setelement(5, Cluster, Logger). -spec with_interval(dns_cluster(), gleam@option:option(integer())) -> dns_cluster(). with_interval(Cluster, Interval) -> erlang:setelement(4, Cluster, Interval). -spec with_resolver(dns_cluster(), resolver()) -> dns_cluster(). with_resolver(Cluster, Resolver) -> erlang:setelement(6, Cluster, Resolver). -spec discover_nodes( gleam@erlang@process:subject(message()), gleam@option:option(integer()) ) -> {ok, {list(gleam@erlang@node:node_()), list(node_connect_error())}} | {error, gleam@erlang@process:call_error({list(gleam@erlang@node:node_()), list(node_connect_error())})}. discover_nodes(Subject, Timeout) -> case Timeout of {some, Timeout@1} -> gleam@erlang@process:try_call( Subject, fun(Client) -> {discover_nodes, {some, Client}, true} end, Timeout@1 ); none -> gleam@erlang@process:send(Subject, {discover_nodes, none, true}), {ok, {[], []}} end. -spec stop(gleam@erlang@process:subject(message()), integer()) -> {ok, nil} | {error, gleam@erlang@process:call_error(nil)}. stop(Subject, Timeout) -> gleam@erlang@process:try_call( Subject, fun(Field@0) -> {stop, Field@0} end, Timeout ). -spec has_ran(gleam@erlang@process:subject(message()), integer()) -> {ok, boolean()} | {error, gleam@erlang@process:call_error(boolean())}. has_ran(Subject, Timeout) -> gleam@erlang@process:try_call( Subject, fun(Field@0) -> {has_ran, Field@0} end, Timeout ). -spec default_resolver() -> resolver(). default_resolver() -> {resolver, fun(A) -> Split = begin _pipe = A, _pipe@1 = erlang:atom_to_binary(_pipe), gleam@string:split_once(_pipe@1, <<"@"/utf8>>) end, case Split of {ok, {Basename, _}} -> {ok, Basename}; _ -> {error, nil} end end, fun gleam_erlang_ffi:connect_node/1, fun() -> [erlang:node() | erlang:nodes()] end, fun(Q) -> Ipv4_addrs = begin _pipe@2 = Q, _pipe@3 = nessie:lookup_ipv4(_pipe@2, in, []), gleam@list:map(_pipe@3, fun(Field@0) -> {ipv4, Field@0} end) end, Ipv6_addrs = begin _pipe@4 = Q, _pipe@5 = nessie:lookup_ipv6(_pipe@4, in, []), gleam@list:map(_pipe@5, fun(Field@0) -> {ipv6, Field@0} end) end, {Ips, _} = begin _pipe@6 = [Ipv4_addrs, Ipv6_addrs], _pipe@7 = gleam@list:concat(_pipe@6), _pipe@8 = gleam@list:map(_pipe@7, fun nessie:ip_to_string/1), gleam@result:partition(_pipe@8) end, Ips end}. -spec default_logger(binary()) -> fun((binary(), binary()) -> nil). default_logger(Prefix) -> fun(Level, Msg) -> gleam@io:println( <<<<<<<>/binary, (gleam@string:uppercase(Level))/binary>>/binary, "] "/utf8>>/binary, Msg/binary>> ) end. -spec new() -> dns_cluster(). new() -> {dns_cluster, erlang:binary_to_atom(<<"nessie_cluster"/utf8>>), ignore, {some, 5000}, default_logger(<<"[nessie_cluster]"/utf8>>), default_resolver()}. -spec connect_error_to_string(gleam@erlang@node:connect_error()) -> binary(). connect_error_to_string(E) -> case E of failed_to_connect -> <<"failed to connect"/utf8>>; local_node_is_not_alive -> <<"local node is not alive"/utf8>> end. -spec do_discover_nodes( resolver(), fun((binary(), binary()) -> nil), binary(), binary() ) -> list(node_connect_error()). do_discover_nodes(Resolver, Logger, Basename, Query) -> Node_names = gleam@list:map( (erlang:element(4, Resolver))(), fun(N) -> erlang:atom_to_binary(gleam_erlang_ffi:identity(N)) end ), Peer_ips = (erlang:element(5, Resolver))(Query), {_, Errors} = begin _pipe = Peer_ips, _pipe@1 = gleam@list:map( _pipe, fun(Ip) -> <<<>/binary, Ip/binary>> end ), _pipe@2 = gleam@list:filter( _pipe@1, fun(Node_name) -> not gleam@list:contains(Node_names, Node_name) end ), _pipe@3 = gleam@list:map( _pipe@2, fun(Node_name@1) -> Atom_node_name = erlang:binary_to_atom(Node_name@1), case (erlang:element(3, Resolver))(Atom_node_name) of {ok, _} -> Logger( <<"info"/utf8>>, <<"Connected to node "/utf8, Node_name@1/binary>> ), {ok, Node_name@1}; {error, Err} -> Logger( <<"error"/utf8>>, <<<<<<"Failed to connect to node "/utf8, Node_name@1/binary>>/binary, ": "/utf8>>/binary, (connect_error_to_string(Err))/binary>> ), {error, {node_connect_error, Atom_node_name, Err}} end end ), gleam@result:partition(_pipe@3) end, Errors. -spec spec( dns_cluster(), gleam@option:option(gleam@erlang@process:subject(gleam@erlang@process:subject(message()))) ) -> gleam@otp@actor:spec(dns_cluster_state(), message()). spec(Cluster, Parent_subject) -> {spec, fun() -> Basename_result = begin _pipe = erlang:node(), _pipe@1 = gleam_erlang_ffi:identity(_pipe), (erlang:element(2, erlang:element(6, Cluster)))(_pipe@1) end, case Basename_result of {ok, Basename} -> _ = gleam_erlang_ffi:register_process( erlang:self(), erlang:element(2, Cluster) ), State = {dns_cluster_state, false, Cluster, Basename, none, gleam@erlang@process:new_subject()}, case {erlang:element(3, Cluster), erlang:element(4, Cluster)} of {_, none} -> nil; {ignore, _} -> nil; {{dns_query, _}, _} -> gleam@erlang@process:send( erlang:element(6, State), {discover_nodes, none, false} ) end, gleam@option:map( Parent_subject, fun(_capture) -> gleam@erlang@process:send( _capture, erlang:element(6, State) ) end ), Selector = gleam@erlang@process:selecting( gleam_erlang_ffi:new_selector(), erlang:element(6, State), fun gleam@function:identity/1 ), {ready, State, Selector}; {error, _} -> {failed, <<"Failed to get node basename"/utf8>>} end end, 10000, fun(Msg, State@1) -> case {Msg, erlang:element(3, erlang:element(3, State@1))} of {{stop, Client}, _} -> gleam@option:map( erlang:element(5, State@1), fun gleam@erlang@process:cancel_timer/1 ), _ = gleam_erlang_ffi:unregister_process( erlang:element(2, erlang:element(3, State@1)) ), gleam@erlang@process:send(Client, nil), (erlang:element(5, erlang:element(3, State@1)))( <<"warn"/utf8>>, <<"DNS cluster stopped."/utf8>> ), {stop, normal}; {{has_ran, Client@1}, _} -> gleam@erlang@process:send( Client@1, erlang:element(2, State@1) ), {continue, State@1, none}; {{discover_nodes, Maybe_client, Manual}, {dns_query, Query}} -> Cluster@1 = erlang:element(3, State@1), Errors = do_discover_nodes( erlang:element(6, Cluster@1), erlang:element(5, Cluster@1), erlang:element(4, State@1), Query ), State@2 = case {erlang:element(4, Cluster@1), Maybe_client, Manual} of {_, {some, Client@2}, _} -> Connected_nodes = (erlang:element( 4, erlang:element(6, Cluster@1) ))(), gleam@otp@actor:send( Client@2, {Connected_nodes, Errors} ), State@1; {_, _, true} -> State@1; {none, _, _} -> State@1; {{some, Interval_millis}, _, _} -> erlang:setelement( 5, State@1, {some, gleam@erlang@process:send_after( erlang:element(6, State@1), Interval_millis, {discover_nodes, none, false} )} ) end, State@3 = erlang:setelement(2, State@2, true), {continue, State@3, none}; {{discover_nodes, Maybe_client@1, _}, ignore} -> (erlang:element(5, erlang:element(3, State@1)))( <<"warn"/utf8>>, <<"DNS cluster is set to ignore, will not discover or connect to nodes."/utf8>> ), case Maybe_client@1 of {some, Client@3} -> Nodes = (erlang:element( 4, erlang:element(6, erlang:element(3, State@1)) ))(), gleam@erlang@process:send(Client@3, {Nodes, []}); none -> nil end, {continue, State@1, none} end end}. -spec start_spec( dns_cluster(), gleam@option:option(gleam@erlang@process:subject(gleam@erlang@process:subject(message()))) ) -> {ok, gleam@erlang@process:subject(message())} | {error, gleam@otp@actor:start_error()}. start_spec(Cluster, Parent_subject) -> gleam@otp@actor:start_spec(spec(Cluster, Parent_subject)).