Current section
Files
Jump to
Current section
Files
src/dream@services@broadcaster.erl
-module(dream@services@broadcaster).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/dream/services/broadcaster.gleam").
-export([subscribe/1, publish/2, unsubscribe/2, channel_to_selector/1, start_broadcaster/0]).
-export_type([broadcaster/1, channel/1, broadcaster_message/1]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
-opaque broadcaster(AIAY) :: {broadcaster,
gleam@erlang@process:subject(broadcaster_message(AIAY))}.
-opaque channel(AIAZ) :: {channel, gleam@erlang@process:subject(AIAZ)}.
-type broadcaster_message(AIBA) :: {subscribe,
gleam@erlang@process:subject(AIBA)} |
{unsubscribe, gleam@erlang@process:subject(AIBA)} |
{publish, AIBA}.
-file("src/dream/services/broadcaster.gleam", 80).
-spec wrap_broadcaster_subject(
gleam@otp@actor:started(gleam@erlang@process:subject(broadcaster_message(AIBF)))
) -> broadcaster(AIBF).
wrap_broadcaster_subject(Started) ->
{broadcaster, erlang:element(3, Started)}.
-file("src/dream/services/broadcaster.gleam", 97).
?DOC(
" Subscribe to receive messages from the broadcaster.\n"
"\n"
" Returns a channel that will receive all messages published\n"
" to the broadcaster.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" let channel = broadcaster.subscribe(my_broadcaster)\n"
" let selector = broadcaster.channel_to_selector(channel)\n"
" ```\n"
).
-spec subscribe(broadcaster(AIBK)) -> channel(AIBK).
subscribe(Broadcaster) ->
{broadcaster, Subject} = Broadcaster,
Subscriber = gleam@erlang@process:new_subject(),
gleam@erlang@process:send(Subject, {subscribe, Subscriber}),
{channel, Subscriber}.
-file("src/dream/services/broadcaster.gleam", 114).
?DOC(
" Publish a message to all subscribers.\n"
"\n"
" The message will be sent to all channels that have subscribed\n"
" to this broadcaster.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" broadcaster.publish(my_broadcaster, UserJoined(\"Alice\"))\n"
" ```\n"
).
-spec publish(broadcaster(AIBN), AIBN) -> nil.
publish(Broadcaster, Message) ->
{broadcaster, Subject} = Broadcaster,
gleam@erlang@process:send(Subject, {publish, Message}).
-file("src/dream/services/broadcaster.gleam", 130).
?DOC(
" Unsubscribe a channel from the broadcaster.\n"
"\n"
" The channel will no longer receive published messages.\n"
" Note: This is typically not needed as channels are automatically\n"
" cleaned up when processes terminate.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" broadcaster.unsubscribe(my_broadcaster, channel)\n"
" ```\n"
).
-spec unsubscribe(broadcaster(AIBP), channel(AIBP)) -> nil.
unsubscribe(Broadcaster, Channel) ->
{broadcaster, Subject} = Broadcaster,
{channel, Subscriber} = Channel,
gleam@erlang@process:send(Subject, {unsubscribe, Subscriber}).
-file("src/dream/services/broadcaster.gleam", 157).
-spec identity(AIBV) -> AIBV.
identity(Value) ->
Value.
-file("src/dream/services/broadcaster.gleam", 151).
?DOC(
" Convert a channel to a selector for use in WebSocket message loops.\n"
"\n"
" This allows the channel to be used with `process.Selector` to receive\n"
" messages in the WebSocket handler's event loop.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" let channel = broadcaster.subscribe(my_broadcaster)\n"
" let selector = broadcaster.channel_to_selector(channel)\n"
" #(initial_state, Some(selector))\n"
" ```\n"
).
-spec channel_to_selector(channel(AIBS)) -> gleam@erlang@process:selector(AIBS).
channel_to_selector(Channel) ->
{channel, Subject} = Channel,
_pipe = gleam_erlang_ffi:new_selector(),
gleam@erlang@process:select_map(_pipe, Subject, fun identity/1).
-file("src/dream/services/broadcaster.gleam", 213).
-spec send_and_continue(
gleam@erlang@process:subject(AICV),
list(gleam@erlang@process:subject(AICV)),
AICV
) -> nil.
send_and_continue(Subscriber, Rest, Broadcast_message) ->
gleam@erlang@process:send(Subscriber, Broadcast_message),
send_to_all_subscribers(Rest, Broadcast_message).
-file("src/dream/services/broadcaster.gleam", 202).
-spec send_to_all_subscribers(list(gleam@erlang@process:subject(AICS)), AICS) -> nil.
send_to_all_subscribers(Subscribers, Broadcast_message) ->
case Subscribers of
[] ->
nil;
[Subscriber | Rest] ->
send_and_continue(Subscriber, Rest, Broadcast_message)
end.
-file("src/dream/services/broadcaster.gleam", 191).
-spec case_remove_subscriber_item(
gleam@erlang@process:subject(AICL),
list(gleam@erlang@process:subject(AICL)),
gleam@erlang@process:subject(AICL)
) -> list(gleam@erlang@process:subject(AICL)).
case_remove_subscriber_item(Subscriber, Rest, To_remove) ->
case Subscriber =:= To_remove of
true ->
remove_subscriber(Rest, To_remove);
false ->
[Subscriber | remove_subscriber(Rest, To_remove)]
end.
-file("src/dream/services/broadcaster.gleam", 180).
-spec remove_subscriber(
list(gleam@erlang@process:subject(AICF)),
gleam@erlang@process:subject(AICF)
) -> list(gleam@erlang@process:subject(AICF)).
remove_subscriber(Subscribers, To_remove) ->
case Subscribers of
[] ->
[];
[Subscriber | Rest] ->
case_remove_subscriber_item(Subscriber, Rest, To_remove)
end.
-file("src/dream/services/broadcaster.gleam", 161).
-spec handle_broadcaster_message(
list(gleam@erlang@process:subject(AIBW)),
broadcaster_message(AIBW)
) -> gleam@otp@actor:next(list(gleam@erlang@process:subject(AIBW)), broadcaster_message(AIBW)).
handle_broadcaster_message(Subscribers, Message) ->
case Message of
{subscribe, Subscriber} ->
gleam@otp@actor:continue([Subscriber | Subscribers]);
{unsubscribe, Subscriber@1} ->
Remaining_subscribers = remove_subscriber(Subscribers, Subscriber@1),
gleam@otp@actor:continue(Remaining_subscribers);
{publish, Broadcast_message} ->
send_to_all_subscribers(Subscribers, Broadcast_message),
gleam@otp@actor:continue(Subscribers)
end.
-file("src/dream/services/broadcaster.gleam", 73).
?DOC(
" Start a new broadcaster service.\n"
"\n"
" Returns an error if the broadcaster fails to start.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" let assert Ok(broadcaster) = broadcaster.start_broadcaster()\n"
" ```\n"
).
-spec start_broadcaster() -> {ok, broadcaster(any())} |
{error, gleam@otp@actor:start_error()}.
start_broadcaster() ->
_pipe = gleam@otp@actor:new([]),
_pipe@1 = gleam@otp@actor:on_message(
_pipe,
fun handle_broadcaster_message/2
),
_pipe@2 = gleam@otp@actor:start(_pipe@1),
gleam@result:map(_pipe@2, fun wrap_broadcaster_subject/1).