Current section
Files
Jump to
Current section
Files
src/glyn@pubsub.erl
-module(glyn@pubsub).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-define(FILEPATH, "src/glyn/pubsub.gleam").
-export([new/2, subscribe/3, unsubscribe/1, publish/3, subscribers/2]).
-export_type([pub_sub/1, subscription/2]).
-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(
" Glyn PubSub - Type-Safe Distributed Event Streaming\n"
"\n"
" This module provides a type-safe wrapper around Erlang's `syn` PubSub system,\n"
" enabling distributed event streaming and one-to-many message broadcasting with\n"
" compile-time type safety.\n"
"\n"
" ## Actor Integration Pattern\n"
"\n"
" PubSub works seamlessly with Gleam's actor system using the selector pattern:\n"
"\n"
" ```gleam\n"
" import gleam/otp/actor\n"
"\n"
" pub type ChatActorMessage {\n"
" GetMessageCount(reply_with: Subject(Int))\n"
" ChatEvent(ChatMessage)\n"
" ChatActorShutdown\n"
" }\n"
"\n"
" fn start_chat_actor(\n"
" pubsub: pubsub.PubSub(ChatMessage),\n"
" group: String,\n"
" ) -> Result(actor.Started(Subject(ChatActorMessage)), actor.StartError) {\n"
" actor.new_with_initialiser(5000, fn(subject) {\n"
" let subscription = pubsub.subscribe(pubsub, group, process.self())\n"
"\n"
" let selector =\n"
" process.new_selector()\n"
" |> process.select(subject)\n"
" |> process.select_map(subscription.subject, ChatEvent)\n"
"\n"
" let initial_state = ChatActorState(message_count: 0, last_message: \"\")\n"
"\n"
" actor.initialised(initial_state)\n"
" |> actor.selecting(selector)\n"
" |> actor.returning(subject)\n"
" |> Ok\n"
" })\n"
" |> actor.on_message(handle_chat_message)\n"
" |> actor.start()\n"
" }\n"
" ```\n"
).
-opaque pub_sub(FFR) :: {pub_sub, gleam@erlang@atom:atom_(), integer()} |
{gleam_phantom, FFR}.
-type subscription(FFS, FFT) :: {subscription,
pub_sub(FFS),
FFT,
gleam@erlang@process:subject(FFS),
gleam@erlang@process:pid_()}.
-file("src/glyn/pubsub.gleam", 95).
?DOC(
" Create a new type-safe PubSub system with a message type for deterministic type identification\n"
" The message_type should be a MessageType that uniquely identifies the message type\n"
).
-spec new(binary(), glyn:message_type(FGF)) -> pub_sub(FGF).
new(Scope, Message_type) ->
Scope@1 = erlang:binary_to_atom(Scope),
syn:add_node_to_scopes([Scope@1]),
{pub_sub, Scope@1, erlang:phash2(erlang:element(2, Message_type))}.
-file("src/glyn/pubsub.gleam", 105).
?DOC(" Subscribe to a PubSub group and return a type-safe Subject\n").
-spec subscribe(pub_sub(FGI), FGK, gleam@erlang@process:pid_()) -> subscription(FGI, FGK).
subscribe(Pubsub, Group, Subscriber_pid) ->
Tagged_group = {Group, erlang:element(3, Pubsub)},
_assert_subject = erlang:binary_to_atom(<<"ok"/utf8>>),
_assert_subject@1 = syn:join(
erlang:element(2, Pubsub),
Tagged_group,
Subscriber_pid
),
case _assert_subject =:= _assert_subject@1 of
true -> nil;
false -> erlang:error(#{gleam_error => assert,
message => <<"Assertion failed."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"glyn/pubsub"/utf8>>,
function => <<"subscribe"/utf8>>,
line => 111,
kind => binary_operator,
operator => '==',
left => #{kind => expression,
value => _assert_subject,
start => 3339,
'end' => 3356
},
right => #{kind => expression,
value => _assert_subject@1,
start => 3364,
'end' => 3416
},
start => 3332,
'end' => 3416,
expression_start => 3339})
end,
Subject = gleam@erlang@process:unsafely_create_subject(
Subscriber_pid,
gleam_stdlib:identity(erlang:element(3, Pubsub))
),
{subscription, Pubsub, Group, Subject, Subscriber_pid}.
-file("src/glyn/pubsub.gleam", 127).
?DOC(" Unsubscribe from a PubSub group\n").
-spec unsubscribe(subscription(any(), any())) -> nil.
unsubscribe(Subscription) ->
Tagged_group = {erlang:element(3, Subscription),
erlang:element(3, erlang:element(2, Subscription))},
_assert_subject = erlang:binary_to_atom(<<"ok"/utf8>>),
_assert_subject@1 = syn:leave(
erlang:element(2, erlang:element(2, Subscription)),
Tagged_group,
erlang:element(5, Subscription)
),
case _assert_subject =:= _assert_subject@1 of
true -> nil;
false -> erlang:error(#{gleam_error => assert,
message => <<"Assertion failed."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"glyn/pubsub"/utf8>>,
function => <<"unsubscribe"/utf8>>,
line => 129,
kind => binary_operator,
operator => '==',
left => #{kind => expression,
value => _assert_subject,
start => 3858,
'end' => 3875
},
right => #{kind => expression,
value => _assert_subject@1,
start => 3883,
'end' => 3987
},
start => 3851,
'end' => 3987,
expression_start => 3858})
end,
nil.
-file("src/glyn/pubsub.gleam", 139).
?DOC(" Publish a type-safe message to all subscribers of a group\n").
-spec publish(pub_sub(FGR), any(), FGR) -> integer().
publish(Pubsub, Group, Message) ->
Tagged_message = {gleam_stdlib:identity(erlang:element(3, Pubsub)), Message},
Tagged_group = {Group, erlang:element(3, Pubsub)},
Subscriber_count@1 = case syn:publish(
erlang:element(2, Pubsub),
Tagged_group,
gleam_stdlib:identity(Tagged_message)
) of
{ok, Subscriber_count} -> Subscriber_count;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"glyn/pubsub"/utf8>>,
function => <<"publish"/utf8>>,
line => 148,
value => _assert_fail,
start => 4492,
'end' => 4597,
pattern_start => 4503,
pattern_end => 4523})
end,
Subscriber_count@1.
-file("src/glyn/pubsub.gleam", 154).
?DOC(" Get list of subscriber PIDs for a group (useful for debugging)\n").
-spec subscribers(pub_sub(any()), binary()) -> list(gleam@erlang@process:pid_()).
subscribers(Pubsub, Group) ->
Tagged_group = {Group, erlang:element(3, Pubsub)},
syn:members(erlang:element(2, Pubsub), Tagged_group).