Packages
macula
0.11.3
7.0.0
6.0.0
5.2.2
5.2.1
5.2.0
5.1.0
5.0.0
4.8.0
4.7.1
4.7.0
4.6.0
4.5.0
4.4.10
4.4.9
4.4.8
4.4.7
4.4.6
4.4.5
4.4.4
4.4.3
4.4.2
4.4.1
4.4.0
4.3.1
4.3.0
4.2.9
4.2.8
4.2.7
4.2.6
4.2.5
4.2.4
4.2.3
4.2.2
4.2.1
4.2.0
4.1.1
4.1.0
4.0.0
3.16.0
3.15.3
3.15.2
3.15.1
3.14.0
3.13.0
3.12.1
3.12.0
3.11.1
3.11.0
3.10.3
3.10.2
3.10.1
3.9.0
3.8.0
3.7.0
3.5.0
3.4.0
3.3.0
3.2.0
3.1.0
3.0.0
2.1.1
2.1.0
2.0.0
1.5.2
1.5.1
1.4.30
1.4.29
1.4.28
1.4.27
1.4.26
1.4.25
1.4.24
1.4.23
1.4.22
1.4.21
1.4.20
1.4.19
1.4.18
1.4.17
1.4.16
1.4.15
1.4.14
1.4.13
1.4.11
1.4.10
1.4.9
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.1
1.3.0
1.2.0
1.1.0
1.0.10
1.0.9
1.0.8
1.0.7
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.48.6
0.48.5
0.48.4
0.48.3
0.48.2
0.48.1
0.48.0
0.47.1
0.47.0
0.46.3
0.46.1
0.46.0
0.45.3
0.45.2
0.45.1
0.45.0
0.44.2
0.44.1
0.44.0
0.43.3
0.43.2
0.43.1
0.43.0
0.42.9
0.42.8
0.42.7
0.42.6
0.42.5
0.42.4
0.42.3
0.42.2
0.42.1
0.42.0
0.41.1
0.41.0
0.40.1
0.40.0
0.39.9
0.39.8
0.39.7
0.39.6
0.39.5
0.39.4
0.39.3
0.39.2
0.39.1
0.39.0
0.38.8
0.38.7
0.38.6
0.38.5
0.38.4
0.38.3
0.38.2
0.38.1
0.38.0
0.37.7
0.37.6
0.37.5
0.37.4
0.37.3
0.37.2
0.37.1
0.37.0
0.36.6
0.36.5
0.36.4
0.36.3
0.36.2
0.36.1
0.36.0
0.35.4
0.35.3
0.35.2
0.35.1
0.35.0
0.34.1
0.34.0
0.33.1
0.33.0
0.32.5
0.32.4
0.32.3
0.32.2
0.32.1
0.32.0
0.31.9
0.31.8
0.31.7
0.31.6
0.31.5
0.31.4
0.31.3
0.31.2
0.31.1
0.31.0
0.30.10
0.30.9
0.30.8
0.30.7
0.30.6
0.30.5
0.30.4
0.30.3
0.30.2
0.30.1
0.30.0
0.29.0
0.28.3
0.28.2
0.28.1
0.28.0
0.27.1
0.27.0
0.26.1
0.26.0
0.25.6
0.25.5
0.25.4
0.25.3
0.25.2
0.25.1
0.25.0
0.24.6
0.24.5
0.24.4
0.24.3
0.24.2
0.24.1
0.24.0
0.23.3
0.23.2
0.23.1
0.23.0
0.22.12
0.22.11
0.22.10
0.22.9
0.22.8
0.22.7
0.22.6
0.22.5
0.22.4
0.22.3
0.22.2
0.22.1
0.22.0
0.21.7
0.21.6
0.21.5
0.21.4
0.21.2
0.21.1
0.21.0
0.20.25
0.20.24
0.20.23
0.20.22
0.20.21
0.20.20
0.20.19
0.20.18
0.20.17
0.20.16
0.20.15
0.20.14
0.20.13
0.20.12
0.20.11
0.20.10
0.20.9
0.20.8
0.20.7
0.20.6
0.20.5
0.20.3
0.20.2
0.20.1
0.20.0
0.19.2
0.19.1
0.19.0
0.18.1
0.18.0
0.17.4
0.17.3
0.17.2
0.17.1
0.17.0
0.16.6
0.16.5
0.16.4
0.16.3
0.16.2
0.16.1
0.16.0
0.15.1
0.15.0
0.14.3
0.14.2
0.14.1
0.14.0
0.12.6
0.12.5
0.12.3
0.11.3
0.10.2
0.10.1
0.10.0
0.9.2
0.9.1
0.9.0
0.8.25
0.8.24
0.8.23
0.8.22
0.8.21
0.8.20
0.8.19
0.8.18
0.8.17
0.8.16
0.8.15
0.8.14
0.8.13
0.8.12
0.8.11
0.8.10
0.8.9
0.8.8
0.8.7
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.30
0.7.29
0.7.28
0.7.27
0.7.26
0.7.25
0.7.24
0.7.23
0.7.22
0.7.21
0.7.20
0.7.19
0.7.18
0.7.17
0.7.16
0.7.15
0.7.14
0.7.13
0.7.12
0.7.11
0.7.10
0.7.9
0.7.8
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.7
0.6.6
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.0
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.4
0.3.3
0.3.2
0.3.1
Macula HTTP/3 Mesh SDK — connect, subscribe, publish, call, advertise
Current section
Files
Jump to
Current section
Files
src/macula_rpc_system/macula_rpc_server.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% RPC server managing registrations and calls.
%%% GenServer that integrates registry, cache, discovery, router, and executor.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_rpc_server).
-behaviour(gen_server).
%% API
-export([
start_link/2,
stop/1,
register/4,
unregister/3,
call/4,
list_registrations/1
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2
]).
%% Types
-type config() :: #{
routing_strategy => macula_rpc_router:strategy(),
cache_enabled => boolean(),
dht_lookup_fun => macula_rpc_discovery:dht_lookup_fun(),
send_fun => macula_rpc_executor:send_fun()
}.
-type state() :: #{
local_node_id := binary(),
registry := macula_rpc_registry:registry(),
cache := macula_rpc_cache:cache(),
router_state := macula_rpc_router:router_state(),
config := config()
}.
-export_type([config/0]).
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Start RPC server.
-spec start_link(binary(), config()) -> {ok, pid()} | {error, term()}.
start_link(LocalNodeId, Config) ->
gen_server:start_link(?MODULE, {LocalNodeId, Config}, []).
%% @doc Stop RPC server.
-spec stop(pid()) -> ok.
stop(Pid) ->
gen_server:stop(Pid).
%% @doc Register procedure.
-spec register(pid(), binary(), macula_rpc_registry:handler_fn(), map()) -> ok.
register(Pid, Uri, Handler, Metadata) ->
gen_server:call(Pid, {register, Uri, Handler, Metadata}).
%% @doc Unregister procedure.
-spec unregister(pid(), binary(), macula_rpc_registry:handler_fn()) -> ok.
unregister(Pid, Uri, Handler) ->
gen_server:call(Pid, {unregister, Uri, Handler}).
%% @doc Synchronous call to procedure.
-spec call(pid(), binary(), map(), pos_integer()) -> {ok, term()} | {error, term()}.
call(Pid, Uri, Args, Timeout) ->
gen_server:call(Pid, {call, Uri, Args, Timeout}, Timeout + 1000).
%% @doc List local registrations.
-spec list_registrations(pid()) -> [macula_rpc_registry:registration()].
list_registrations(Pid) ->
gen_server:call(Pid, list_registrations).
%%%===================================================================
%%% gen_server Callbacks
%%%===================================================================
%% @doc Initialize server state.
init({LocalNodeId, Config}) ->
%% Create initial state
Registry = macula_rpc_registry:new(),
Cache = macula_rpc_cache:new(1000),
Strategy = maps:get(routing_strategy, Config, local_first),
RouterState = macula_rpc_router:new_state(Strategy),
State = #{
local_node_id => LocalNodeId,
registry => Registry,
cache => Cache,
router_state => RouterState,
config => Config
},
{ok, State}.
%% @doc Handle synchronous calls.
handle_call({register, Uri, Handler, Metadata}, _From, State) ->
#{registry := Registry} = State,
%% Validate URI
case macula_rpc_names:validate(Uri) of
ok ->
%% Add to registry
NewRegistry = macula_rpc_registry:register(Registry, Uri, Handler, Metadata),
NewState = State#{registry => NewRegistry},
{reply, ok, NewState};
{error, _Reason} = Error ->
{reply, Error, State}
end;
handle_call({unregister, Uri, Handler}, _From, State) ->
#{registry := Registry} = State,
%% Remove from registry
NewRegistry = macula_rpc_registry:unregister(Registry, Uri, Handler),
NewState = State#{registry => NewRegistry},
{reply, ok, NewState};
handle_call({call, Uri, Args, Timeout}, _From, State) ->
#{
local_node_id := LocalNodeId,
registry := Registry,
cache := Cache,
router_state := RouterState,
config := Config
} = State,
%% Validate URI
case macula_rpc_names:validate(Uri) of
ok ->
%% Execute call
{Result, NewState} = execute_call(
Uri, Args, Timeout, LocalNodeId, Registry, Cache, RouterState, Config
),
{reply, Result, NewState};
{error, _Reason} = Error ->
{reply, Error, State}
end;
handle_call(list_registrations, _From, State) ->
#{registry := Registry} = State,
Registrations = macula_rpc_registry:list_registrations(Registry),
{reply, Registrations, State}.
%% @doc Handle asynchronous casts (none implemented).
handle_cast(_Msg, State) ->
{noreply, State}.
%% @doc Handle info messages (none expected).
handle_info(_Info, State) ->
{noreply, State}.
%% @doc Cleanup on termination.
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Execute RPC call (local or remote).
-spec execute_call(
binary(), map(), pos_integer(), binary(),
macula_rpc_registry:registry(),
macula_rpc_cache:cache(),
macula_rpc_router:router_state(),
config()
) -> {{ok, term()} | {error, term()}, state()}.
execute_call(Uri, Args, Timeout, LocalNodeId, Registry, Cache, RouterState, Config) ->
%% Check cache first (if enabled)
CacheEnabled = maps:get(cache_enabled, Config, false),
case try_cache_get(Uri, Args, Cache, CacheEnabled) of
{ok, Result, NewCache} ->
%% Cache hit
NewState = #{
local_node_id => LocalNodeId,
registry => Registry,
cache => NewCache,
router_state => RouterState,
config => Config
},
{{ok, Result}, NewState};
not_found ->
%% Cache miss, execute call
execute_call_no_cache(
Uri, Args, Timeout, LocalNodeId, Registry, Cache, RouterState, Config
)
end.
%% @doc Try to get result from cache.
-spec try_cache_get(binary(), map(), macula_rpc_cache:cache(), boolean()) ->
{ok, term(), macula_rpc_cache:cache()} | not_found.
try_cache_get(_Uri, _Args, _Cache, false) ->
%% Cache disabled
not_found;
try_cache_get(Uri, Args, Cache, true) ->
%% Try cache
macula_rpc_cache:get(Cache, Uri, Args).
%% @doc Execute call without cache.
-spec execute_call_no_cache(
binary(), map(), pos_integer(), binary(),
macula_rpc_registry:registry(),
macula_rpc_cache:cache(),
macula_rpc_router:router_state(),
config()
) -> {{ok, term()} | {error, term()}, state()}.
execute_call_no_cache(Uri, Args, Timeout, LocalNodeId, Registry, Cache, RouterState, Config) ->
%% Find local handlers
LocalHandlers = macula_rpc_registry:find(Registry, Uri),
%% Find remote providers
DhtLookupFun = maps:get(dht_lookup_fun, Config, fun default_dht_lookup/1),
RemoteProviders = unwrap_providers(DhtLookupFun(Uri)),
%% Route call
Strategy = maps:get(routing_strategy, Config, local_first),
{ProviderResult, NewRouterState} = route_call(
Strategy, LocalHandlers, RemoteProviders, RouterState
),
%% Execute based on routing result
case ProviderResult of
{local, Registration} ->
%% Execute local handler
Handler = maps:get(handler, Registration),
Result = macula_rpc_executor:execute_local(Handler, Args, Timeout),
%% Cache result if enabled and successful
NewCache = maybe_cache_result(Uri, Args, Result, Registration, Cache, Config),
NewState = #{
local_node_id => LocalNodeId,
registry => Registry,
cache => NewCache,
router_state => NewRouterState,
config => Config
},
{Result, NewState};
{remote, Provider} ->
%% Execute remote call
SendFun = maps:get(send_fun, Config, fun default_send/4),
Result = macula_rpc_executor:execute_remote(Uri, Args, Provider, SendFun, Timeout),
%% Don't cache remote results (for now)
NewState = #{
local_node_id => LocalNodeId,
registry => Registry,
cache => Cache,
router_state => NewRouterState,
config => Config
},
{Result, NewState};
{error, no_provider} ->
%% No provider available
NewState = #{
local_node_id => LocalNodeId,
registry => Registry,
cache => Cache,
router_state => NewRouterState,
config => Config
},
{{error, no_provider}, NewState}
end.
%% @doc Route call to provider.
-spec route_call(
macula_rpc_router:strategy(),
[macula_rpc_registry:registration()],
[macula_rpc_discovery:provider_info()],
macula_rpc_router:router_state()
) -> {
{local, macula_rpc_registry:registration()} |
{remote, macula_rpc_discovery:provider_info()} |
{error, no_provider},
macula_rpc_router:router_state()
}.
route_call(round_robin, LocalHandlers, RemoteProviders, RouterState) ->
%% Round robin requires stateful API
macula_rpc_router:select_provider_stateful(RouterState, LocalHandlers, RemoteProviders);
route_call(Strategy, LocalHandlers, RemoteProviders, RouterState) ->
%% Other strategies are stateless
Result = macula_rpc_router:select_provider(Strategy, LocalHandlers, RemoteProviders),
{Result, RouterState}.
%% @doc Cache result if caching enabled and result is successful.
-spec maybe_cache_result(
binary(), map(),
{ok, term()} | {error, term()},
macula_rpc_registry:registration(),
macula_rpc_cache:cache(),
config()
) -> macula_rpc_cache:cache().
maybe_cache_result(Uri, Args, {ok, Result}, Registration, Cache, Config) ->
CacheEnabled = maps:get(cache_enabled, Config, false),
Metadata = maps:get(metadata, Registration),
CacheTTL = maps:get(cache_ttl, Metadata, 0),
cache_if_enabled(CacheEnabled andalso CacheTTL > 0, Cache, Uri, Args, Result, CacheTTL);
maybe_cache_result(_Uri, _Args, {error, _Reason}, _Registration, Cache, _Config) ->
%% Don't cache errors
Cache.
%% @doc Cache result if caching is enabled.
cache_if_enabled(true, Cache, Uri, Args, Result, CacheTTL) ->
%% Cache result with TTL
macula_rpc_cache:put(Cache, Uri, Args, Result, CacheTTL);
cache_if_enabled(false, Cache, _Uri, _Args, _Result, _CacheTTL) ->
%% Don't cache
Cache.
%% @doc Unwrap provider list from ok/error result.
unwrap_providers({ok, Providers}) -> Providers;
unwrap_providers({error, _}) -> [].
%% @doc Default DHT lookup (returns empty list - for testing).
-spec default_dht_lookup(binary()) -> {ok, [macula_rpc_discovery:provider_info()]}.
default_dht_lookup(_Uri) ->
{ok, []}.
%% @doc Default send function (returns error - for testing).
-spec default_send(binary(), map(), macula_rpc_discovery:address(), pos_integer()) ->
{ok, term()} | {error, term()}.
default_send(_Uri, _Args, _Address, _Timeout) ->
{error, no_transport}.