Current section
Files
Jump to
Current section
Files
src/nessie_cluster.erl
-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(
<<<<<<<<Prefix/binary, "["/utf8>>/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) -> <<<<Basename/binary, "@"/utf8>>/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)).