-module(glixir@syn). -compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]). -define(FILEPATH, "src/glixir/syn.gleam"). -export([init_scopes/1, register/3, whereis/2, unregister/2, join/2, leave/2, publish/3, publish_json/4, lookup/2, members_detailed/2, members/2, is_member/3, member_count/2, register_worker/3, find_worker/2, broadcast_status/2, join_coordination/1, leave_coordination/1]). -export_type([syn_result/0, syn_ok/0, registry_error/0, pub_sub_error/0, whereis_result/0, syn_whereis_result/0, syn_register_result/0, syn_unregister_result/0, syn_lookup_result/0, syn_members_result/0, syn_is_member_result/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( "\n" " Distributed process coordination using Erlang's `syn` library.\n" " \n" " Provides type-safe wrappers for syn's registry and PubSub functionality,\n" " enabling distributed service discovery and event streaming across BEAM nodes.\n" "\n" " Inspired by the excellent glyn project's syn wrapper patterns.\n" " Check out glyn for more advanced OTP patterns: https://github.com/glyn-project/glyn\n" "\n" " This implementation is optimized for glixir's type-safe OTP wrappers.\n" "\n" ). -type syn_result() :: any(). -type syn_ok() :: any(). -type registry_error() :: {registration_failed, binary()} | {lookup_failed, binary()} | {unregistration_failed, binary()}. -type pub_sub_error() :: {join_failed, binary()} | {leave_failed, binary()} | {publish_failed, binary()}. -type whereis_result() :: not_found | {found_pid_only, gleam@erlang@process:pid_()} | {found_with_metadata, gleam@erlang@process:pid_(), gleam@dynamic:dynamic_()} | {whereis_error, binary()}. -type syn_whereis_result() :: syn_not_found | {syn_found_pid_only, gleam@erlang@process:pid_()} | {syn_found_with_metadata, gleam@erlang@process:pid_(), gleam@dynamic:dynamic_()} | {syn_error, gleam@dynamic:dynamic_()}. -type syn_register_result() :: syn_register_ok | {syn_register_error, gleam@dynamic:dynamic_()}. -type syn_unregister_result() :: syn_unregister_ok | {syn_unregister_error, gleam@dynamic:dynamic_()}. -type syn_lookup_result() :: syn_lookup_not_found | {syn_lookup_found, gleam@erlang@process:pid_(), gleam@dynamic:dynamic_()} | {syn_lookup_error, gleam@dynamic:dynamic_()}. -type syn_members_result() :: syn_members_empty | {syn_members_list, list({binary(), gleam@erlang@process:pid_(), gleam@dynamic:dynamic_()})} | {syn_members_error, gleam@dynamic:dynamic_()}. -type syn_is_member_result() :: syn_is_member_true | syn_is_member_false | {syn_is_member_error, gleam@dynamic:dynamic_()}. -file("src/glixir/syn.gleam", 152). -spec syn_result_to_gleam(syn_result()) -> {ok, nil} | {error, gleam@dynamic:dynamic_()}. syn_result_to_gleam(Result) -> case gleam_stdlib:identity(Result) =:= gleam_stdlib:identity( erlang:binary_to_atom(<<"ok"/utf8>>) ) of true -> {ok, nil}; false -> {_, Reason} = gleam_stdlib:identity(gleam_stdlib:identity(Result)), {error, Reason} end. -file("src/glixir/syn.gleam", 175). ?DOC( " Initialize scopes for distributed coordination.\n" " Call this once at application startup for each scope you plan to use.\n" " \n" " ## Example\n" " ```gleam\n" " syn.init_scopes([\"user_sessions\", \"game_lobbies\", \"worker_pool\"])\n" " ```\n" ). -spec init_scopes(list(binary())) -> nil. init_scopes(Scopes) -> Scope_atoms = gleam@list:map(Scopes, fun erlang:binary_to_atom/1), glixir@utils:debug_log_with_prefix( info, <<"syn"/utf8>>, <<"Initializing scopes: "/utf8, (gleam@string:inspect(Scopes))/binary>> ), syn:add_node_to_scopes(Scope_atoms), nil. -file("src/glixir/syn.gleam", 202). -spec register(binary(), binary(), any()) -> {ok, nil} | {error, registry_error()}. register(Scope, Name, Metadata) -> Current_pid = erlang:self(), glixir@utils:debug_log_with_prefix( debug, <<"syn"/utf8>>, <<<<<<"Registering process: "/utf8, Scope/binary>>/binary, "/"/utf8>>/binary, Name/binary>> ), case glixir_syn:register( Scope, Name, Current_pid, gleam_stdlib:identity(Metadata) ) of syn_register_ok -> glixir@utils:debug_log_with_prefix( info, <<"syn"/utf8>>, <<"Registered successfully: "/utf8, Name/binary>> ), {ok, nil}; {syn_register_error, Reason} -> Error_msg = <<"Registration failed: "/utf8, (gleam@string:inspect(Reason))/binary>>, glixir@utils:debug_log_with_prefix(error, <<"syn"/utf8>>, Error_msg), {error, {registration_failed, Error_msg}} end. -file("src/glixir/syn.gleam", 243). ?DOC( " Look up a registered process by scope and name.\n" " Returns the process PID and its metadata if found.\n" " \n" " ## Example \n" " ```gleam\n" " case syn.whereis(\"user_sessions\", user_id) {\n" " Ok(#(pid, #(username, last_active))) -> // Send message to user\n" " Error(_) -> // User not online\n" " }\n" " ```\n" ). -spec whereis(binary(), binary()) -> {ok, {gleam@erlang@process:pid_(), any()}} | {error, registry_error()}. whereis(Scope, Name) -> glixir@utils:debug_log_with_prefix( debug, <<"syn"/utf8>>, <<<<<<"Looking up process: "/utf8, Scope/binary>>/binary, "/"/utf8>>/binary, Name/binary>> ), case glixir_syn:whereis_name(Scope, Name) of syn_not_found -> Error_msg = <<<<<<"Process not found: "/utf8, Scope/binary>>/binary, "/"/utf8>>/binary, Name/binary>>, glixir@utils:debug_log_with_prefix(debug, <<"syn"/utf8>>, Error_msg), {error, {lookup_failed, Error_msg}}; {syn_found_pid_only, _} -> Error_msg@1 = <<<<<<"Process found but has no metadata: "/utf8, Scope/binary>>/binary, "/"/utf8>>/binary, Name/binary>>, glixir@utils:debug_log_with_prefix( warning, <<"syn"/utf8>>, Error_msg@1 ), {error, {lookup_failed, Error_msg@1}}; {syn_found_with_metadata, Pid, Metadata_dyn} -> Metadata = gleam_stdlib:identity(Metadata_dyn), glixir@utils:debug_log_with_prefix( debug, <<"syn"/utf8>>, <<<<<<"Found process with metadata: "/utf8, Scope/binary>>/binary, "/"/utf8>>/binary, Name/binary>> ), {ok, {Pid, Metadata}}; {syn_error, Reason} -> Error_msg@2 = <<"Syn error: "/utf8, (gleam@string:inspect(Reason))/binary>>, glixir@utils:debug_log_with_prefix( error, <<"syn"/utf8>>, Error_msg@2 ), {error, {lookup_failed, Error_msg@2}} end. -file("src/glixir/syn.gleam", 294). -spec unregister(binary(), binary()) -> {ok, nil} | {error, registry_error()}. unregister(Scope, Name) -> glixir@utils:debug_log_with_prefix( debug, <<"syn"/utf8>>, <<<<<<"Unregistering process: "/utf8, Scope/binary>>/binary, "/"/utf8>>/binary, Name/binary>> ), case glixir_syn:unregister(Scope, Name) of syn_unregister_ok -> glixir@utils:debug_log_with_prefix( info, <<"syn"/utf8>>, <<"Unregistered successfully: "/utf8, Name/binary>> ), {ok, nil}; {syn_unregister_error, Reason} -> Error_msg = <<"Unregistration failed: "/utf8, (gleam@string:inspect(Reason))/binary>>, glixir@utils:debug_log_with_prefix(error, <<"syn"/utf8>>, Error_msg), {error, {unregistration_failed, Error_msg}} end. -file("src/glixir/syn.gleam", 331). ?DOC( " Join a PubSub group to receive published messages.\n" " The current process will receive all messages published to this group.\n" " \n" " ## Example\n" " ```gleam\n" " syn.join(\"chat_rooms\", \"general\")\n" " // Process will now receive all messages published to \"general\" chat\n" " ```\n" ). -spec join(binary(), binary()) -> {ok, nil} | {error, pub_sub_error()}. join(Scope, Group) -> Scope_atom = erlang:binary_to_atom(Scope), Current_pid = erlang:self(), glixir@utils:debug_log_with_prefix( debug, <<"syn"/utf8>>, <<<<<<"Joining PubSub group: "/utf8, Scope/binary>>/binary, "/"/utf8>>/binary, Group/binary>> ), case syn:join(Scope_atom, Group, Current_pid) of Result -> case syn_result_to_gleam(Result) of {ok, nil} -> glixir@utils:debug_log_with_prefix( info, <<"syn"/utf8>>, <<"Joined group successfully: "/utf8, Group/binary>> ), {ok, nil}; {error, Reason} -> Error_msg = <<"Join failed: "/utf8, (gleam@string:inspect(Reason))/binary>>, glixir@utils:debug_log_with_prefix( error, <<"syn"/utf8>>, Error_msg ), {error, {join_failed, Error_msg}} end end. -file("src/glixir/syn.gleam", 368). ?DOC( " Leave a PubSub group to stop receiving messages.\n" " The current process will no longer receive messages published to this group.\n" " \n" " ## Example\n" " ```gleam\n" " syn.leave(\"chat_rooms\", \"general\")\n" " ```\n" ). -spec leave(binary(), binary()) -> {ok, nil} | {error, pub_sub_error()}. leave(Scope, Group) -> Scope_atom = erlang:binary_to_atom(Scope), Current_pid = erlang:self(), glixir@utils:debug_log_with_prefix( debug, <<"syn"/utf8>>, <<<<<<"Leaving PubSub group: "/utf8, Scope/binary>>/binary, "/"/utf8>>/binary, Group/binary>> ), case syn:leave(Scope_atom, Group, Current_pid) of Result -> case syn_result_to_gleam(Result) of {ok, nil} -> glixir@utils:debug_log_with_prefix( info, <<"syn"/utf8>>, <<"Left group successfully: "/utf8, Group/binary>> ), {ok, nil}; {error, Reason} -> Error_msg = <<"Leave failed: "/utf8, (gleam@string:inspect(Reason))/binary>>, glixir@utils:debug_log_with_prefix( error, <<"syn"/utf8>>, Error_msg ), {error, {leave_failed, Error_msg}} end end. -file("src/glixir/syn.gleam", 408). ?DOC( " Publish a string message to all members of a PubSub group.\n" " Returns the number of processes the message was delivered to.\n" " \n" " ## Example\n" " ```gleam\n" " case syn.publish(\"chat_rooms\", \"general\", \"Hello everyone!\") {\n" " Ok(count) -> // Message sent to `count` processes\n" " Error(_) -> // Publish failed\n" " }\n" " ```\n" ). -spec publish(binary(), binary(), binary()) -> {ok, integer()} | {error, pub_sub_error()}. publish(Scope, Group, Message) -> Scope_atom = erlang:binary_to_atom(Scope), glixir@utils:debug_log_with_prefix( debug, <<"syn"/utf8>>, <<<<<<"Publishing to group: "/utf8, Scope/binary>>/binary, "/"/utf8>>/binary, Group/binary>> ), case syn:publish(Scope_atom, Group, gleam_stdlib:identity(Message)) of {ok, Count} -> glixir@utils:debug_log_with_prefix( info, <<"syn"/utf8>>, <<<<"Published to "/utf8, (gleam@string:inspect(Count))/binary>>/binary, " processes"/utf8>> ), {ok, Count}; {error, Reason} -> Error_msg = <<"Publish failed: "/utf8, (gleam@string:inspect(Reason))/binary>>, glixir@utils:debug_log_with_prefix(error, <<"syn"/utf8>>, Error_msg), {error, {publish_failed, Error_msg}} end. -file("src/glixir/syn.gleam", 451). ?DOC( " Publish a typed message using JSON encoding.\n" " The message will be JSON-encoded before sending to subscribers.\n" " \n" " ## Example\n" " ```gleam\n" " let user_message = json.object([\n" " #(\"user\", json.string(\"alice\")),\n" " #(\"content\", json.string(\"Hello!\")),\n" " #(\"timestamp\", json.int(system_time))\n" " ])\n" " \n" " syn.publish_json(\"chat_rooms\", \"general\", user_message, fn(j) { j })\n" " ```\n" ). -spec publish_json(binary(), binary(), FQQ, fun((FQQ) -> gleam@json:json())) -> {ok, integer()} | {error, pub_sub_error()}. publish_json(Scope, Group, Message, Encoder) -> Json_string = begin _pipe = Encoder(Message), gleam@json:to_string(_pipe) end, publish(Scope, Group, Json_string). -file("src/glixir/syn.gleam", 474). -spec lookup(binary(), binary()) -> {ok, {gleam@erlang@process:pid_(), any()}} | {error, registry_error()}. lookup(Scope, Name) -> glixir@utils:debug_log_with_prefix( debug, <<"syn"/utf8>>, <<<<<<"Looking up (via lookup): "/utf8, Scope/binary>>/binary, "/"/utf8>>/binary, Name/binary>> ), case glixir_syn:lookup(Scope, Name) of syn_lookup_not_found -> {error, {lookup_failed, <<<<<<"Process not found: "/utf8, Scope/binary>>/binary, "/"/utf8>>/binary, Name/binary>>}}; {syn_lookup_found, Pid, Metadata_dyn} -> Metadata = gleam_stdlib:identity(Metadata_dyn), {ok, {Pid, Metadata}}; {syn_lookup_error, Reason} -> {error, {lookup_failed, <<"Lookup error: "/utf8, (gleam@string:inspect(Reason))/binary>>}} end. -file("src/glixir/syn.gleam", 501). -spec members_detailed(binary(), binary()) -> list({binary(), gleam@erlang@process:pid_(), any()}). members_detailed(Scope, Group) -> case glixir_syn:members(Scope, Group) of syn_members_empty -> []; {syn_members_list, Members} -> gleam@list:map( Members, fun(Member) -> {Name, Pid, Meta_dyn} = Member, {Name, Pid, gleam_stdlib:identity(Meta_dyn)} end ); {syn_members_error, _} -> [] end. -file("src/glixir/syn.gleam", 527). ?DOC( " Get list of member PIDs and metadata without names.\n" " Simpler version of members_detailed for when you don't need member names.\n" "\n" " ## Example\n" " ```gleam\n" " let followers = syn.members(\"JobDispatcherFollowers\", \"dispatchers\")\n" " // Returns: [#(pid1, metadata1), #(pid2, metadata2)]\n" " ```\n" ). -spec members(binary(), binary()) -> list({gleam@erlang@process:pid_(), any()}). members(Scope, Group) -> case glixir_syn:members(Scope, Group) of syn_members_empty -> []; {syn_members_list, Members} -> gleam@list:map( Members, fun(Member) -> {_, Pid, Meta_dyn} = Member, {Pid, gleam_stdlib:identity(Meta_dyn)} end ); {syn_members_error, _} -> [] end. -file("src/glixir/syn.gleam", 546). -spec is_member(binary(), binary(), gleam@erlang@process:pid_()) -> boolean(). is_member(Scope, Name, Pid) -> case glixir_syn:is_member(Scope, Name, Pid) of syn_is_member_true -> true; syn_is_member_false -> false; {syn_is_member_error, _} -> false end. -file("src/glixir/syn.gleam", 564). ?DOC( " Get count of processes subscribed to a group.\n" " More efficient than getting the full member list if you only need the count.\n" " \n" " ## Example\n" " ```gleam\n" " let user_count = syn.member_count(\"chat_rooms\", \"general\")\n" " if user_count > 100 {\n" " // Maybe split the room\n" " }\n" " ```\n" ). -spec member_count(binary(), binary()) -> integer(). member_count(Scope, Group) -> Scope_atom = erlang:binary_to_atom(Scope), syn:member_count(Scope_atom, Group). -file("src/glixir/syn.gleam", 580). ?DOC( " Register a worker process with load information.\n" " Common pattern for distributed worker pools.\n" " \n" " ## Example\n" " ```gleam\n" " syn.register_worker(\"image_processors\", worker_id, #(cpu_usage, queue_size))\n" " ```\n" ). -spec register_worker(binary(), binary(), any()) -> {ok, nil} | {error, registry_error()}. register_worker(Pool_name, Worker_id, Load_info) -> register( <<"worker_pools"/utf8>>, <<<>/binary, Worker_id/binary>>, Load_info ). -file("src/glixir/syn.gleam", 597). ?DOC( " Find a worker in a pool and get its load information.\n" " \n" " ## Example\n" " ```gleam\n" " case syn.find_worker(\"image_processors\", worker_id) {\n" " Ok(#(pid, #(cpu_usage, queue_size))) -> // Send work to worker\n" " Error(_) -> // Worker not available\n" " }\n" " ```\n" ). -spec find_worker(binary(), binary()) -> {ok, {gleam@erlang@process:pid_(), any()}} | {error, registry_error()}. find_worker(Pool_name, Worker_id) -> whereis( <<"worker_pools"/utf8>>, <<<>/binary, Worker_id/binary>> ). -file("src/glixir/syn.gleam", 617). ?DOC( " Broadcast a status update to all processes in a coordination group.\n" " Common pattern for distributed consensus or health monitoring.\n" " \n" " ## Example\n" " ```gleam\n" " let status = json.object([\n" " #(\"node\", json.string(node_name)),\n" " #(\"load\", json.float(cpu_load)),\n" " #(\"timestamp\", json.int(now))\n" " ])\n" " \n" " syn.broadcast_status(\"health_check\", status)\n" " ```\n" ). -spec broadcast_status(binary(), gleam@json:json()) -> {ok, integer()} | {error, pub_sub_error()}. broadcast_status(Group, Status_message) -> publish_json( <<"coordination"/utf8>>, Group, Status_message, fun(J) -> J end ). -file("src/glixir/syn.gleam", 632). ?DOC( " Join a coordination group for distributed algorithm participation.\n" " Common pattern for consensus, leader election, or load balancing.\n" " \n" " ## Example\n" " ```gleam\n" " syn.join_coordination(\"leader_election\")\n" " // Process will receive leader election messages\n" " ```\n" ). -spec join_coordination(binary()) -> {ok, nil} | {error, pub_sub_error()}. join_coordination(Group) -> join(<<"coordination"/utf8>>, Group). -file("src/glixir/syn.gleam", 642). ?DOC( " Leave a coordination group.\n" " \n" " ## Example\n" " ```gleam\n" " syn.leave_coordination(\"leader_election\") \n" " ```\n" ). -spec leave_coordination(binary()) -> {ok, nil} | {error, pub_sub_error()}. leave_coordination(Group) -> leave(<<"coordination"/utf8>>, Group).