Packages

Type-safe PubSub and Registry for Gleam actors with distributed clustering support

Current section

Files

Jump to
glyn src glyn@pubsub.erl
Raw

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/3, subscribe/2, unsubscribe/2, publish/3, subscribers/2, subscriber_count/2]).
-export_type([syn_result/0, syn_ok/0, pub_sub/1, pub_sub_error/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(
" Glyn PubSub - Selector-Based Type-Safe Event Streaming\n"
"\n"
" This module provides a selector-based wrapper around Erlang's `syn` PubSub system,\n"
" enabling distributed event streaming and one-to-many message broadcasting with\n"
" runtime type safety through dynamic decoding.\n"
"\n"
" ## Multi-Channel Actor Integration Pattern\n"
"\n"
" PubSub seamlessly composes with other message channels using selectors:\n"
"\n"
" ```gleam\n"
" import gleam/dynamic.{type Dynamic}\n"
" import gleam/dynamic/decode\n"
" import gleam/erlang/atom\n"
" import gleam/erlang/process.{type Subject}\n"
" import gleam/otp/actor\n"
" import glyn/pubsub\n"
" import glyn/registry\n"
"\n"
" // Define your event types\n"
" pub type ChatMessage {\n"
" UserJoined(username: String)\n"
" UserLeft(username: String)\n"
" Message(username: String, content: String)\n"
" }\n"
"\n"
" pub type MetricEvent {\n"
" CounterIncrement(name: String, value: Int)\n"
" GaugeUpdate(name: String, value: Float)\n"
" }\n"
"\n"
" pub type ActorMessage {\n"
" DirectCommand(String) // Direct commands\n"
" ChatEvent(ChatMessage) // Chat PubSub events\n"
" MetricEvent(MetricEvent) // Metrics PubSub events\n"
" }\n"
"\n"
" // Create decoders for your event types\n"
" fn expect_atom(expected: String) -> decode.Decoder(atom.Atom) {\n"
" use value <- decode.then(atom.decoder())\n"
" case atom.to_string(value) == expected {\n"
" True -> decode.success(value)\n"
" False -> decode.failure(value, \"Expected atom: \" <> expected)\n"
" }\n"
" }\n"
"\n"
" fn chat_message_decoder() -> decode.Decoder(ChatMessage) {\n"
" decode.one_of(\n"
" {\n"
" use _ <- decode.field(0, expect_atom(\"user_joined\"))\n"
" use username <- decode.field(1, decode.string)\n"
" decode.success(UserJoined(username))\n"
" },\n"
" or: [\n"
" {\n"
" use _ <- decode.field(0, expect_atom(\"message\"))\n"
" use username <- decode.field(1, decode.string)\n"
" use content <- decode.field(2, decode.string)\n"
" decode.success(Message(username, content))\n"
" },\n"
" // Add other variants as needed\n"
" ]\n"
" )\n"
" }\n"
"\n"
" fn metric_event_decoder() -> decode.Decoder(MetricEvent) {\n"
" decode.one_of(\n"
" {\n"
" use _ <- decode.field(0, expect_atom(\"counter_increment\"))\n"
" use name <- decode.field(1, decode.string)\n"
" use value <- decode.field(2, decode.int)\n"
" decode.success(CounterIncrement(name, value))\n"
" },\n"
" or: [\n"
" {\n"
" use _ <- decode.field(0, expect_atom(\"gauge_update\"))\n"
" use name <- decode.field(1, decode.string)\n"
" use value <- decode.field(2, decode.float)\n"
" decode.success(GaugeUpdate(name, value))\n"
" },\n"
" ]\n"
" )\n"
" }\n"
"\n"
" fn start_multi_channel_actor() {\n"
" actor.new_with_initialiser(5000, fn(_) {\n"
" let command_subject = process.new_subject()\n"
"\n"
" // Create base selector for direct commands\n"
" let base_selector =\n"
" process.new_selector()\n"
" |> process.select_map(command_subject, DirectCommand)\n"
"\n"
" // Add chat PubSub channel\n"
" let chat_pubsub = pubsub.new(\n"
" scope: \"chat_events\",\n"
" decoder: chat_message_decoder(),\n"
" error_default: UserJoined(\"unknown\")\n"
" )\n"
" let chat_selector = pubsub.subscribe(chat_pubsub, \"general\")\n"
" let with_chat = base_selector\n"
" |> process.merge_selector(\n"
" process.map_selector(chat_selector, ChatEvent)\n"
" )\n"
"\n"
" // Add metrics PubSub channel\n"
" let metrics_pubsub = pubsub.new(\n"
" scope: \"metrics_events\",\n"
" decoder: metric_event_decoder(),\n"
" error_default: CounterIncrement(\"unknown\", 0)\n"
" )\n"
" let metrics_selector = pubsub.subscribe(metrics_pubsub, \"system\")\n"
" let final_selector = with_chat\n"
" |> process.merge_selector(\n"
" process.map_selector(metrics_selector, MetricEvent)\n"
" )\n"
"\n"
" actor.initialised(initial_state)\n"
" |> actor.selecting(final_selector)\n"
" |> actor.returning(command_subject)\n"
" |> Ok\n"
" })\n"
" }\n"
"\n"
" // Publishing events to subscribers\n"
" let chat_pubsub = pubsub.new(\n"
" scope: \"chat_events\",\n"
" decoder: chat_message_decoder(),\n"
" error_default: UserJoined(\"unknown\")\n"
" )\n"
"\n"
" // Publish a chat message to all subscribers in \"general\" channel\n"
" let assert Ok(subscriber_count) = pubsub.publish(\n"
" chat_pubsub,\n"
" \"general\",\n"
" Message(\"alice\", \"Hello everyone!\")\n"
" )\n"
"\n"
" // Check how many subscribers received the message\n"
" let count = pubsub.subscriber_count(chat_pubsub, \"general\")\n"
" ```\n"
).
-type syn_result() :: any().
-type syn_ok() :: any().
-opaque pub_sub(FIP) :: {pub_sub,
gleam@erlang@atom:atom_(),
gleam@dynamic@decode:decoder(FIP),
FIP}.
-type pub_sub_error() :: {publish_failed, binary()}.
-file("src/glyn/pubsub.gleam", 198).
?DOC(" Create a new PubSub system for a given scope with dynamic decoding\n").
-spec new(binary(), gleam@dynamic@decode:decoder(FJD), FJD) -> pub_sub(FJD).
new(Scope, Decoder, Error_default) ->
Scope@1 = erlang:binary_to_atom(Scope),
syn:add_node_to_scopes([Scope@1]),
{pub_sub, Scope@1, Decoder, Error_default}.
-file("src/glyn/pubsub.gleam", 210).
?DOC(
" Subscribe to a PubSub group and compose into a selector\n"
" Creates an internal Subject(Dynamic) and uses select_map for type safety\n"
).
-spec subscribe(pub_sub(FJG), binary()) -> gleam@erlang@process:selector(FJG).
subscribe(Pubsub, Group) ->
Current_pid = erlang:self(),
Group_tag = gleam_stdlib:identity(Group),
case begin
_pipe = syn:join(erlang:element(2, Pubsub), Group, Current_pid),
syn_ffi:to_result(_pipe)
end of
{ok, nil} -> nil;
_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 => <<"subscribe"/utf8>>,
line => 218,
value => _assert_fail,
start => 6924,
'end' => 7002,
pattern_start => 6935,
pattern_end => 6942})
end,
Dynamic_subject = gleam@erlang@process:unsafely_create_subject(
Current_pid,
Group_tag
),
_pipe@1 = gleam_erlang_ffi:new_selector(),
gleam@erlang@process:select_map(
_pipe@1,
Dynamic_subject,
fun(Dynamic) ->
_pipe@2 = gleam@dynamic@decode:run(
Dynamic,
erlang:element(3, Pubsub)
),
gleam@result:unwrap(_pipe@2, erlang:element(4, Pubsub))
end
).
-file("src/glyn/pubsub.gleam", 228).
?DOC(" Unsubscribe from a PubSub group\n").
-spec unsubscribe(pub_sub(any()), binary()) -> nil.
unsubscribe(Pubsub, Group) ->
Current_pid = erlang:self(),
case begin
_pipe = syn:leave(erlang:element(2, Pubsub), Group, Current_pid),
syn_ffi:to_result(_pipe)
end of
{ok, nil} ->
nil;
{error, _} ->
nil
end.
-file("src/glyn/pubsub.gleam", 238).
?DOC(" Publish a type-safe message to all subscribers of a group\n").
-spec publish(pub_sub(FJL), binary(), FJL) -> {ok, integer()} |
{error, pub_sub_error()}.
publish(Pubsub, Group, Message) ->
Group_tag = gleam_stdlib:identity(Group),
Tagged_message = {Group_tag, Message},
case syn:publish(
erlang:element(2, Pubsub),
Group,
gleam_stdlib:identity(Tagged_message)
) of
{ok, Subscriber_count} ->
{ok, Subscriber_count};
{error, Reason} ->
{error,
{publish_failed,
<<"publish failed: "/utf8,
(gleam@string:inspect(Reason))/binary>>}}
end.
-file("src/glyn/pubsub.gleam", 256).
?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) ->
syn:members(erlang:element(2, Pubsub), Group).
-file("src/glyn/pubsub.gleam", 261).
?DOC(" Get the count of subscribers for a group\n").
-spec subscriber_count(pub_sub(any()), binary()) -> integer().
subscriber_count(Pubsub, Group) ->
syn:member_count(erlang:element(2, Pubsub), Group).