Current section
Files
Jump to
Current section
Files
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/3, register/3, unregister/2, whereis/2, send/3, call/4]).
-export_type([syn_ok/0, syn_result/0, registry/2, registry_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 Registry - Selector-Based Type-Safe Process Registry\n"
"\n"
" This module provides a selector-based wrapper around Erlang's `syn` process registry,\n"
" enabling distributed service discovery and direct process communication with\n"
" runtime type safety through dynamic decoding.\n"
"\n"
" ## Multi-Channel Actor Integration Pattern\n"
"\n"
" The registry 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/registry\n"
" import glyn/pubsub\n"
"\n"
" // Define your message types\n"
" pub type ServiceMessage {\n"
" ProcessOrder(id: String, reply_with: Subject(Bool))\n"
" GetStatus(reply_with: Subject(String))\n"
" Shutdown\n"
" }\n"
"\n"
" pub type SystemEvent {\n"
" ServiceStarted(name: String)\n"
" ServiceStopped(name: String)\n"
" }\n"
"\n"
" pub type ActorMessage {\n"
" DirectCommand(String) // Direct commands\n"
" RegistryMessage(ServiceMessage) // Registry messages (decoded)\n"
" PubSubEvent(SystemEvent) // PubSub events\n"
" }\n"
"\n"
" // Create decoders for your message 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 service_message_decoder() -> decode.Decoder(ServiceMessage) {\n"
" decode.one_of(\n"
" {\n"
" use _ <- decode.field(0, expect_atom(\"shutdown\"))\n"
" decode.success(Shutdown)\n"
" },\n"
" or: [\n"
" {\n"
" use _ <- decode.field(0, expect_atom(\"get_status\"))\n"
" use reply_with <- decode.field(1, subject_decoder())\n"
" decode.success(GetStatus(reply_with))\n"
" },\n"
" // Add other variants as needed\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 registry channel\n"
" let user_registry = registry.new(\n"
" scope: \"user_services\",\n"
" decoder: service_message_decoder(),\n"
" error_default: Shutdown\n"
" )\n"
" let assert Ok(registry_selector) = registry.register(\n"
" user_registry,\n"
" \"order_processor\",\n"
" \"v1.0\"\n"
" )\n"
" let with_registry = base_selector\n"
" |> process.merge_selector(\n"
" process.map_selector(registry_selector, RegistryMessage)\n"
" )\n"
"\n"
" // Add pubsub channel for system events\n"
" let system_pubsub = pubsub.new(\n"
" scope: \"system_events\",\n"
" decoder: system_event_decoder(),\n"
" error_default: ServiceStarted(\"unknown\")\n"
" )\n"
" let pubsub_selector = pubsub.subscribe(system_pubsub, \"services\")\n"
" let final_selector = with_registry\n"
" |> process.merge_selector(\n"
" process.map_selector(pubsub_selector, PubSubEvent)\n"
" )\n"
"\n"
" actor.initialised(initial_state)\n"
" |> actor.selecting(final_selector)\n"
" |> actor.returning(command_subject)\n"
" |> Ok\n"
" })\n"
" }\n"
"\n"
" // Send messages to registered services\n"
" let user_registry = registry.new(\n"
" scope: \"user_services\",\n"
" decoder: service_message_decoder(),\n"
" error_default: Shutdown\n"
" )\n"
"\n"
" // Send a message\n"
" let assert Ok(_) = registry.send(\n"
" user_registry,\n"
" \"order_processor\",\n"
" ProcessOrder(\"order-123\", reply_subject)\n"
" )\n"
"\n"
" // Make a call and wait for reply\n"
" let assert Ok(status) = registry.call(\n"
" user_registry,\n"
" \"order_processor\",\n"
" waiting: 5000,\n"
" sending: GetStatus(_)\n"
" )\n"
" ```\n"
).
-type syn_ok() :: any().
-type syn_result() :: any().
-opaque registry(FLV, FLW) :: {registry,
gleam@erlang@atom:atom_(),
gleam@dynamic@decode:decoder(FLV),
FLV} |
{gleam_phantom, FLW}.
-type registry_error() :: timeout |
{process_not_found, binary()} |
{registration_failed, binary()} |
{unregistration_failed, binary()}.
-file("src/glyn/registry.gleam", 189).
?DOC(" Create a new Registry system for a given scope with dynamic decoding\n").
-spec new(binary(), gleam@dynamic@decode:decoder(FMD), FMD) -> registry(FMD, any()).
new(Scope, Decoder, Error_default) ->
Scope@1 = erlang:binary_to_atom(Scope),
syn:add_node_to_scopes([Scope@1]),
{registry, Scope@1, Decoder, Error_default}.
-file("src/glyn/registry.gleam", 201).
?DOC(
" Register a process with a name and return a selector for receiving messages\n"
" Creates an internal Subject(Dynamic) and uses select_map for type safety\n"
).
-spec register(registry(FMI, FMJ), binary(), FMJ) -> {ok,
gleam@erlang@process:selector(FMI)} |
{error, registry_error()}.
register(Registry, Actor_name, Metadata) ->
Current_pid = erlang:self(),
Dynamic_subject = gleam@erlang@process:new_subject(),
Result = syn:register(
erlang:element(2, Registry),
Actor_name,
Current_pid,
gleam_stdlib:identity({Dynamic_subject, Metadata})
),
case syn_ffi:to_result(Result) of
{ok, nil} ->
Selector = begin
_pipe = gleam_erlang_ffi:new_selector(),
gleam@erlang@process:select_map(
_pipe,
Dynamic_subject,
fun(Dynamic_msg) ->
_pipe@1 = gleam@dynamic@decode:run(
Dynamic_msg,
erlang:element(3, Registry)
),
gleam@result:unwrap(
_pipe@1,
erlang:element(4, Registry)
)
end
)
end,
{ok, Selector};
{error, E} ->
{error,
{registration_failed,
<<"syn registration failed: "/utf8,
(gleam@string:inspect(E))/binary>>}}
end.
-file("src/glyn/registry.gleam", 236).
?DOC(" Unregister a process by name\n").
-spec unregister(registry(any(), any()), binary()) -> {ok, nil} |
{error, registry_error()}.
unregister(Registry, Actor_name) ->
Result = syn:unregister(erlang:element(2, Registry), Actor_name),
case syn_ffi:to_result(Result) of
{ok, nil} ->
{ok, nil};
{error, E} ->
{error,
{unregistration_failed,
<<"syn unregistration failed: "/utf8,
(gleam@string:inspect(E))/binary>>}}
end.
-file("src/glyn/registry.gleam", 251).
?DOC(" Look up a registered process and return PID with metadata\n").
-spec whereis(registry(any(), FMW), binary()) -> {ok,
{gleam@erlang@process:pid_(), FMW}} |
{error, registry_error()}.
whereis(Registry, Actor_name) ->
Result = syn:lookup(erlang:element(2, Registry), Actor_name),
case Result =:= gleam_stdlib:identity(
erlang:binary_to_atom(<<"undefined"/utf8>>)
) of
true ->
{error, {process_not_found, Actor_name}};
false ->
{Pid, Stored_data} = gleam_stdlib:identity(Result),
{_, Metadata} = gleam_stdlib:identity(Stored_data),
{ok, {Pid, Metadata}}
end.
-file("src/glyn/registry.gleam", 267).
?DOC(" Send a message to a registered process using the stored dynamic subject\n").
-spec send(registry(FNB, any()), binary(), FNB) -> {ok, nil} |
{error, registry_error()}.
send(Registry, Actor_name, Message) ->
Result = syn:lookup(erlang:element(2, Registry), Actor_name),
case Result =:= gleam_stdlib:identity(
erlang:binary_to_atom(<<"undefined"/utf8>>)
) of
true ->
{error, {process_not_found, Actor_name}};
false ->
{_, Stored_data} = gleam_stdlib:identity(Result),
{Dynamic_subject, _} = gleam_stdlib:identity(Stored_data),
gleam@erlang@process:send(
Dynamic_subject,
gleam_stdlib:identity(Message)
),
{ok, nil}
end.
-file("src/glyn/registry.gleam", 286).
?DOC(" Call a registered process and wait for a reply, similar to actor.call\n").
-spec call(
registry(FNH, any()),
binary(),
integer(),
fun((gleam@erlang@process:subject(FNL)) -> FNH)
) -> {ok, FNL} | {error, registry_error()}.
call(Registry, Actor_name, Timeout, Message_fn) ->
Reply_subject = gleam@erlang@process:new_subject(),
Message = Message_fn(Reply_subject),
case send(Registry, Actor_name, Message) of
{ok, _} ->
case gleam@erlang@process:'receive'(Reply_subject, Timeout) of
{ok, Reply} ->
{ok, Reply};
{error, nil} ->
{error, timeout}
end;
{error, Error} ->
{error, Error}
end.