-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 => <>, 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).