Current section
Files
Jump to
Current section
Files
src/esdee@discoverer.erl
-module(esdee@discoverer).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/esdee/discoverer.gleam").
-export([start_timeout/2, poll_timeout/2, build/1, named/2, stop/1, subscribe_to_service_types_mapping/3, subscribe_to_service_types/2, poll_service_types/1, subscribe_to_service_details_mapping/4, subscribe_to_service_details/3, poll_service_details/2, unsubscribe/2, start/1, supervised/2]).
-export_type([discoverer/0, builder/0, subscription/0, state/0, msg/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(" Provides an actor-based DNS-SD discovery mechanism.\n").
-opaque discoverer() :: {discoverer,
gleam@erlang@process:subject(msg()),
integer()}.
-opaque builder() :: {builder,
esdee:options(),
gleam@option:option(gleam@erlang@process:name(msg())),
integer(),
integer(),
gleam@option:option(fun((discoverer()) -> nil))}.
-opaque subscription() :: {service_type_subscription, fun((binary()) -> nil)} |
{service_details_subscription,
binary(),
fun((esdee:service_description()) -> nil)}.
-type state() :: {state,
esdee:options(),
esdee:sockets(),
esdee@internal@dispatcher:dispatcher()}.
-opaque msg() :: stop |
{subscribe_to_service_types, fun((binary()) -> nil), boolean()} |
{subscribe_to_service_details,
binary(),
fun((esdee:service_description()) -> nil),
boolean()} |
{poll_service_types,
gleam@erlang@process:subject({ok, nil} | {error, toss:error()})} |
{poll_service_details,
binary(),
gleam@erlang@process:subject({ok, nil} | {error, toss:error()})} |
{upd_update, esdee:udp_message()}.
-file("src/esdee/discoverer.gleam", 33).
?DOC(
" Sets a timeout for the actor to start.\n"
" Should be very fast, as no incoming data is waited for.\n"
).
-spec start_timeout(builder(), integer()) -> builder().
start_timeout(Builder, Timeout) ->
{builder,
erlang:element(2, Builder),
erlang:element(3, Builder),
Timeout,
erlang:element(5, Builder),
erlang:element(6, Builder)}.
-file("src/esdee/discoverer.gleam", 39).
?DOC(
" Sets a timeout for the actor to respond to poll requests.\n"
" Should be very fast, as no incoming data is waited for.\n"
).
-spec poll_timeout(builder(), integer()) -> builder().
poll_timeout(Builder, Timeout) ->
{builder,
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
Timeout,
erlang:element(6, Builder)}.
-file("src/esdee/discoverer.gleam", 44).
?DOC(" Starts building a discoverer from the base options.\n").
-spec build(esdee:options()) -> builder().
build(Options) ->
{builder, Options, none, 1000, 1000, none}.
-file("src/esdee/discoverer.gleam", 50).
?DOC(
" Configures the builder to use a named process with the actor,\n"
" and returns a discoverer instance that can be used across potential restarts.\n"
).
-spec named(builder(), gleam@erlang@process:name(msg())) -> {builder(),
discoverer()}.
named(Builder, Name) ->
Subject = gleam@erlang@process:named_subject(Name),
Discoverer = {discoverer, Subject, erlang:element(5, Builder)},
Builder@1 = {builder,
erlang:element(2, Builder),
{some, Name},
erlang:element(4, Builder),
erlang:element(5, Builder),
erlang:element(6, Builder)},
{Builder@1, Discoverer}.
-file("src/esdee/discoverer.gleam", 112).
?DOC(" Stops the service discovery actor.\n").
-spec stop(discoverer()) -> nil.
stop(Discoverer) ->
gleam@erlang@process:send(erlang:element(2, Discoverer), stop).
-file("src/esdee/discoverer.gleam", 139).
?DOC(
" Subscribes the given subject to all discovered service types,\n"
" using the given function to map to another message type.\n"
" Note that the same service type might be reported by multiple peers.\n"
" You will also need to call `poll_service_types` to discover services quickly.\n"
).
-spec subscribe_to_service_types_mapping(
discoverer(),
gleam@erlang@process:subject(JVD),
fun((binary()) -> JVD)
) -> subscription().
subscribe_to_service_types_mapping(Discoverer, Subject, Mapper) ->
Callback = fun(Msg) -> gleam@erlang@process:send(Subject, Mapper(Msg)) end,
gleam@erlang@process:send(
erlang:element(2, Discoverer),
{subscribe_to_service_types, Callback, true}
),
{service_type_subscription, Callback}.
-file("src/esdee/discoverer.gleam", 128).
?DOC(
" Subscribes the given subject to all discovered service types.\n"
" Note that the same service type might be reported by multiple peers.\n"
" You will also need to call `poll_service_types` to discover services quickly.\n"
).
-spec subscribe_to_service_types(
discoverer(),
gleam@erlang@process:subject(binary())
) -> subscription().
subscribe_to_service_types(Discoverer, Subject) ->
subscribe_to_service_types_mapping(
Discoverer,
Subject,
fun gleam@function:identity/1
).
-file("src/esdee/discoverer.gleam", 151).
?DOC(
" Sends a DNS-SD question querying all the available service types in the local network.\n"
" If there are errors with the socket(s), returns the first error.\n"
).
-spec poll_service_types(discoverer()) -> {ok, nil} | {error, toss:error()}.
poll_service_types(Discoverer) ->
gleam@erlang@process:call(
erlang:element(2, Discoverer),
erlang:element(3, Discoverer),
fun(Field@0) -> {poll_service_types, Field@0} end
).
-file("src/esdee/discoverer.gleam", 173).
?DOC(
" Subscribes the given subject to all discovered service details,\n"
" using the given function to map to another message type.\n"
" You will also need to call `poll_service_details` to discover services quickly.\n"
).
-spec subscribe_to_service_details_mapping(
discoverer(),
binary(),
gleam@erlang@process:subject(JVI),
fun((esdee:service_description()) -> JVI)
) -> subscription().
subscribe_to_service_details_mapping(Discoverer, Service_type, Subject, Mapper) ->
Callback = fun(Msg) -> gleam@erlang@process:send(Subject, Mapper(Msg)) end,
gleam@erlang@process:send(
erlang:element(2, Discoverer),
{subscribe_to_service_details, Service_type, Callback, true}
),
{service_details_subscription, Service_type, Callback}.
-file("src/esdee/discoverer.gleam", 157).
?DOC(
" Subscribes the given subject to all discovered service details.\n"
" You will also need to call `poll_service_details` to discover services quickly.\n"
).
-spec subscribe_to_service_details(
discoverer(),
binary(),
gleam@erlang@process:subject(esdee:service_description())
) -> subscription().
subscribe_to_service_details(Discoverer, Service_type, Subject) ->
subscribe_to_service_details_mapping(
Discoverer,
Service_type,
Subject,
fun gleam@function:identity/1
).
-file("src/esdee/discoverer.gleam", 189).
?DOC(
" Sends a DNS-SD question querying the given service type in the local network.\n"
" If there are errors with the socket(s), returns the first error.\n"
).
-spec poll_service_details(discoverer(), binary()) -> {ok, nil} |
{error, toss:error()}.
poll_service_details(Discoverer, Service_type) ->
gleam@erlang@process:call(
erlang:element(2, Discoverer),
erlang:element(3, Discoverer),
fun(_capture) -> {poll_service_details, Service_type, _capture} end
).
-file("src/esdee/discoverer.gleam", 200).
?DOC(" Terminates the given subscription.\n").
-spec unsubscribe(discoverer(), subscription()) -> nil.
unsubscribe(Discoverer, Subscription) ->
case Subscription of
{service_type_subscription, Callback} ->
gleam@erlang@process:send(
erlang:element(2, Discoverer),
{subscribe_to_service_types, Callback, false}
);
{service_details_subscription, Service_type, Callback@1} ->
gleam@erlang@process:send(
erlang:element(2, Discoverer),
{subscribe_to_service_details, Service_type, Callback@1, false}
)
end.
-file("src/esdee/discoverer.gleam", 287).
-spec poll(
state(),
binary(),
gleam@erlang@process:subject({ok, nil} | {error, toss:error()})
) -> gleam@otp@actor:next(state(), msg()).
poll(State, Service_type, Respond_to) ->
Result = esdee:broadcast_service_question(
erlang:element(3, State),
Service_type
),
gleam@erlang@process:send(Respond_to, Result),
gleam@otp@actor:continue(State).
-file("src/esdee/discoverer.gleam", 324).
-spec describe_toss_error(toss:error()) -> binary().
describe_toss_error(Error) ->
<<"UDP socket failed: "/utf8, (toss:describe_error(Error))/binary>>.
-file("src/esdee/discoverer.gleam", 319).
-spec receive_next_datagram_as_message(esdee:sockets()) -> {ok, nil} |
{error, binary()}.
receive_next_datagram_as_message(Sockets) ->
_pipe = esdee:receive_next_datagram_as_message(Sockets),
gleam@result:map_error(_pipe, fun describe_toss_error/1).
-file("src/esdee/discoverer.gleam", 297).
-spec handle_udp_update(
esdee@internal@dispatcher:dispatcher(),
esdee:sockets(),
esdee:udp_message()
) -> {ok, nil} | {error, binary()}.
handle_udp_update(Dispatcher, Sockets, Update) ->
gleam@result:'try'(case Update of
{dns_sd_message, Update@1} ->
esdee@internal@dispatcher:dispatch(Dispatcher, Update@1),
{ok, nil};
{other_udp_message, Message} ->
case Message of
{datagram, _, _, _, _} ->
{ok, nil};
{udp_error, _, Error} ->
{error, describe_toss_error(Error)}
end
end, fun(_) -> receive_next_datagram_as_message(Sockets) end).
-file("src/esdee/discoverer.gleam", 235).
-spec handle_message(state(), msg()) -> gleam@otp@actor:next(state(), msg()).
handle_message(State, Msg) ->
case Msg of
{poll_service_details, Service_type, Reply_to} ->
poll(State, Service_type, Reply_to);
{poll_service_types, Reply_to@1} ->
poll(State, <<"_services._dns-sd._udp.local"/utf8>>, Reply_to@1);
{subscribe_to_service_types, Callback, Subscribe} ->
Dispatcher = case Subscribe of
true ->
esdee@internal@dispatcher:subscribe_to_service_types(
erlang:element(4, State),
Callback
);
false ->
esdee@internal@dispatcher:unsubscribe_from_service_types(
erlang:element(4, State),
Callback
)
end,
gleam@otp@actor:continue(
{state,
erlang:element(2, State),
erlang:element(3, State),
Dispatcher}
);
{subscribe_to_service_details, Service_type@1, Callback@1, Subscribe@1} ->
Dispatcher@1 = case Subscribe@1 of
true ->
esdee@internal@dispatcher:subscribe_to_service_details(
erlang:element(4, State),
Service_type@1,
Callback@1
);
false ->
esdee@internal@dispatcher:unsubscribe_from_service_details(
erlang:element(4, State),
Service_type@1,
Callback@1
)
end,
gleam@otp@actor:continue(
{state,
erlang:element(2, State),
erlang:element(3, State),
Dispatcher@1}
);
{upd_update, Update} ->
case handle_udp_update(
erlang:element(4, State),
erlang:element(3, State),
Update
) of
{ok, _} ->
gleam@otp@actor:continue(State);
{error, E} ->
gleam@otp@actor:stop_abnormal(E)
end;
stop ->
esdee:close_sockets(erlang:element(3, State)),
gleam@otp@actor:stop()
end.
-file("src/esdee/discoverer.gleam", 61).
?DOC(" Starts a DNS-SD service discovery actor\n").
-spec start(builder()) -> {ok, gleam@otp@actor:started(discoverer())} |
{error, gleam@otp@actor:start_error()}.
start(Builder) ->
Actor_builder = begin
_pipe@5 = gleam@otp@actor:new_with_initialiser(
erlang:element(4, Builder),
fun(Self) ->
gleam@result:'try'(
begin
_pipe = esdee:set_up_sockets(erlang:element(2, Builder)),
gleam@result:map_error(
_pipe,
fun esdee:describe_setup_error/1
)
end,
fun(Sockets) ->
Selctor = begin
_pipe@1 = gleam_erlang_ffi:new_selector(),
_pipe@2 = gleam@erlang@process:select(_pipe@1, Self),
esdee:select_processed_udp_messages(
_pipe@2,
fun(Field@0) -> {upd_update, Field@0} end
)
end,
gleam@result:'try'(
receive_next_datagram_as_message(Sockets),
fun(_) ->
{ok,
begin
_pipe@3 = gleam@otp@actor:initialised(
{state,
erlang:element(2, Builder),
Sockets,
esdee@internal@dispatcher:new()}
),
_pipe@4 = gleam@otp@actor:selecting(
_pipe@3,
Selctor
),
gleam@otp@actor:returning(
_pipe@4,
{discoverer,
Self,
erlang:element(5, Builder)}
)
end}
end
)
end
)
end
),
gleam@otp@actor:on_message(_pipe@5, fun handle_message/2)
end,
Result = begin
_pipe@6 = case erlang:element(3, Builder) of
{some, Name} ->
gleam@otp@actor:named(Actor_builder, Name);
none ->
Actor_builder
end,
gleam@otp@actor:start(_pipe@6)
end,
case {Result, erlang:element(6, Builder)} of
{{ok, Started}, {some, On_start}} ->
On_start(erlang:element(3, Started));
{_, _} ->
nil
end,
Result.
-file("src/esdee/discoverer.gleam", 103).
?DOC(
" Returns a child specification for running the actor with supervision.\n"
" If specified, the `on_start` function will be called after the actor starts or restarts.\n"
).
-spec supervised(builder(), gleam@option:option(fun((discoverer()) -> nil))) -> gleam@otp@supervision:child_specification(discoverer()).
supervised(Builder, On_start) ->
Builder@1 = {builder,
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
erlang:element(5, Builder),
On_start},
gleam@otp@supervision:worker(fun() -> start(Builder@1) end).