Packages

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

Current section

Files

Jump to
glyn src glyn@registry.erl
Raw

src/glyn@registry.erl

-module(glyn@registry).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-define(FILEPATH, "src/glyn/registry.gleam").
-export([new/2, register/4, unregister/1, lookup/2, send/3, call/4]).
-export_type([registry/2, registration/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 Registry - Type-Safe Distributed Process Registry\n"
"\n"
" This module provides a type-safe wrapper around Erlang's `syn` process registry,\n"
" enabling distributed service discovery and direct process communication with\n"
" compile-time type safety.\n"
"\n"
" ## Actor Integration Pattern\n"
"\n"
" The registry works seamlessly with Gleam's actor system and composes with PubSub:\n"
"\n"
" ```gleam\n"
" import gleam/otp/actor\n"
" import glyn/pubsub\n"
"\n"
" pub type ActorMessage {\n"
" CommandMessage(Command) // Direct commands via Registry\n"
" PubSubMessage(Event) // Events via PubSub\n"
" ActorShutdown\n"
" }\n"
"\n"
" fn start_integration_actor(\n"
" registry: registry.Registry(Command, String),\n"
" pubsub: pubsub.PubSub(Event),\n"
" ) -> Result(actor.Started(Subject(ActorMessage)), actor.StartError) {\n"
" actor.new_with_initialiser(5000, fn(subject) {\n"
" // Create command subject for Registry\n"
" let command_subject = process.new_subject()\n"
" let assert Ok(_registration) = registry.register(\n"
" registry, \"my_service\", command_subject, \"service_metadata\"\n"
" )\n"
"\n"
" // Subscribe to PubSub events\n"
" let event_subscription = pubsub.subscribe(pubsub, \"events\", process.self())\n"
"\n"
" // Compose both message sources\n"
" let selector =\n"
" process.new_selector()\n"
" |> process.select(subject)\n"
" |> process.select_map(command_subject, CommandMessage)\n"
" |> process.select_map(event_subscription.subject, PubSubMessage)\n"
"\n"
" actor.initialised(initial_state)\n"
" |> actor.selecting(selector)\n"
" |> actor.returning(subject)\n"
" |> Ok\n"
" })\n"
" |> actor.on_message(handle_message)\n"
" |> actor.start()\n"
" }\n"
" ```\n"
).
-opaque registry(FIQ, FIR) :: {registry, gleam@erlang@atom:atom_(), integer()} |
{gleam_phantom, FIQ, FIR}.
-type registration(FIS, FIT) :: {registration,
registry(FIS, FIT),
binary(),
gleam@erlang@process:subject(FIS),
FIT}.
-file("src/glyn/registry.gleam", 105).
?DOC(
" Create a new Registry system for a given scope\n"
" The message_type should be a MessageType that uniquely identifies the message type\n"
).
-spec new(binary(), glyn:message_type(FIZ)) -> registry(FIZ, any()).
new(Scope, Message_type) ->
Scope@1 = erlang:binary_to_atom(Scope),
syn:add_node_to_scopes([Scope@1]),
{registry, Scope@1, erlang:phash2(erlang:element(2, Message_type))}.
-file("src/glyn/registry.gleam", 116).
?DOC(
" Register a process with a name using a caller-supplied Subject\n"
" Note: Registration will replace any existing registration with the same name\n"
).
-spec register(
registry(FJE, FJF),
binary(),
gleam@erlang@process:subject(FJE),
FJF
) -> {ok, registration(FJE, FJF)} | {error, binary()}.
register(Registry, Name, Subject, Metadata) ->
case gleam@erlang@process:subject_owner(Subject) of
{ok, Pid} ->
Result = syn:register(
erlang:element(2, Registry),
Name,
Pid,
gleam_stdlib:identity(
{Subject, erlang:element(3, Registry), Metadata}
)
),
case Result =:= gleam_stdlib:identity(
erlang:binary_to_atom(<<"ok"/utf8>>)
) of
true ->
{ok, {registration, Registry, Name, Subject, Metadata}};
false ->
{error,
<<"Registration failed: "/utf8,
(gleam@string:inspect(Result))/binary>>}
end;
{error, _} ->
{error, <<"Invalid subject: process may have terminated"/utf8>>}
end.
-file("src/glyn/registry.gleam", 152).
?DOC(" Unregister a process\n").
-spec unregister(registration(any(), any())) -> {ok, nil} | {error, binary()}.
unregister(Registration) ->
Result = syn:unregister(
erlang:element(2, erlang:element(2, Registration)),
erlang:element(3, Registration)
),
case Result =:= gleam_stdlib:identity(erlang:binary_to_atom(<<"ok"/utf8>>)) of
true ->
{ok, nil};
false ->
{error,
<<"Unregistration failed: "/utf8,
(gleam@string:inspect(Result))/binary>>}
end.
-file("src/glyn/registry.gleam", 163).
?DOC(" Look up a registered process and return a type-safe Subject with metadata\n").
-spec lookup(registry(FJT, FJU), binary()) -> {ok,
{gleam@erlang@process:subject(FJT), FJU}} |
{error, binary()}.
lookup(Registry, Name) ->
Result = syn:lookup(erlang:element(2, Registry), Name),
case Result =:= gleam_stdlib:identity(
erlang:binary_to_atom(<<"undefined"/utf8>>)
) of
true ->
{error, <<"Process not found: "/utf8, Name/binary>>};
false ->
{_, Stored_data} = gleam_stdlib:identity(Result),
{Subject, Tag, Metadata} = gleam_stdlib:identity(Stored_data),
case Tag =:= erlang:element(3, Registry) of
true ->
{ok, {Subject, Metadata}};
false ->
{error,
<<"Process registered under incompatible type: "/utf8,
Name/binary>>}
end
end.
-file("src/glyn/registry.gleam", 184).
?DOC(" Send a type-safe message to a registered process\n").
-spec send(registry(FKA, any()), binary(), FKA) -> {ok, nil} | {error, binary()}.
send(Registry, Name, Message) ->
case lookup(Registry, Name) of
{ok, {Subject, _}} ->
gleam@erlang@process:send(Subject, Message),
{ok, nil};
{error, Reason} ->
{error, Reason}
end.
-file("src/glyn/registry.gleam", 199).
?DOC(" Call a registered process and wait for a reply, similar to actor.call\n").
-spec call(
registry(FKG, any()),
binary(),
integer(),
fun((gleam@erlang@process:subject(FKK)) -> FKG)
) -> {ok, FKK} | {error, binary()}.
call(Registry, Name, Timeout, Message_fn) ->
case lookup(Registry, Name) of
{ok, {Subject, _}} ->
Reply = gleam@otp@actor:call(Subject, Timeout, Message_fn),
{ok, Reply};
{error, Reason} ->
{error, Reason}
end.