Packages
macula
0.20.0
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_bridge_system/macula_bridge_node.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Macula Bridge Node - Manages connection to parent mesh level.
%%%
%%% The Bridge Node is responsible for:
%%% - Connecting to parent mesh (street to neighborhood to city to etc.)
%%% - Escalating DHT queries when local DHT misses
%%% - Caching results from parent queries locally
%%% - Maintaining connection health to parent bridges
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_bridge_node).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
start_link/1,
escalate_query/2,
escalate_query/3,
store_to_parent/2,
is_connected/1,
get_stats/1,
get_parent_bridges/1,
set_parent_bridges/2
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2
]).
%% State
-record(state, {
mesh_level :: atom(), % cluster | street | neighborhood | city | etc.
parent_bridges :: [binary()], % List of parent bridge endpoints
connected_parent :: binary() | undefined, % Currently connected parent bridge
connection_pid :: pid() | undefined, % PID of QUIC connection to parent
escalation_timeout :: pos_integer(), % Timeout for parent queries (ms)
stats :: map() % Statistics
}).
-define(DEFAULT_ESCALATION_TIMEOUT, 5000).
-define(RECONNECT_INTERVAL, 5000).
-define(HEALTH_CHECK_INTERVAL, 30000).
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Start bridge node with registered name.
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Config) ->
gen_server:start_link({local, ?MODULE}, ?MODULE, Config, []).
%% @doc Escalate a DHT query to parent level.
%% Called when local DHT lookup fails.
-spec escalate_query(pid(), map()) -> {ok, term()} | {error, term()}.
escalate_query(Pid, Query) ->
escalate_query(Pid, Query, ?DEFAULT_ESCALATION_TIMEOUT).
-spec escalate_query(pid(), map(), pos_integer()) -> {ok, term()} | {error, term()}.
escalate_query(Pid, Query, Timeout) ->
gen_server:call(Pid, {escalate_query, Query}, Timeout + 1000).
%% @doc Store value to parent DHT (for advertisement propagation).
-spec store_to_parent(pid(), map()) -> ok | {error, term()}.
store_to_parent(Pid, StoreMsg) ->
gen_server:call(Pid, {store_to_parent, StoreMsg}, 10000).
%% @doc Check if connected to parent bridge.
-spec is_connected(pid()) -> boolean().
is_connected(Pid) ->
gen_server:call(Pid, is_connected).
%% @doc Get bridge node statistics.
-spec get_stats(pid()) -> {ok, map()}.
get_stats(Pid) ->
gen_server:call(Pid, get_stats).
%% @doc Get list of parent bridge endpoints.
-spec get_parent_bridges(pid()) -> [binary()].
get_parent_bridges(Pid) ->
gen_server:call(Pid, get_parent_bridges).
%% @doc Update parent bridge endpoints (for dynamic configuration).
-spec set_parent_bridges(pid(), [binary()]) -> ok.
set_parent_bridges(Pid, Bridges) ->
gen_server:call(Pid, {set_parent_bridges, Bridges}).
%%%===================================================================
%%% gen_server Callbacks
%%%===================================================================
init(Config) ->
MeshLevel = maps:get(mesh_level, Config, cluster),
ParentBridges = maps:get(parent_bridges, Config, []),
EscalationTimeout = maps:get(escalation_timeout, Config, ?DEFAULT_ESCALATION_TIMEOUT),
?LOG_INFO("[BridgeNode] Starting at level ~p with parent bridges: ~p",
[MeshLevel, ParentBridges]),
State = #state{
mesh_level = MeshLevel,
parent_bridges = ParentBridges,
connected_parent = undefined,
connection_pid = undefined,
escalation_timeout = EscalationTimeout,
stats = init_stats()
},
%% Schedule initial connection attempt if we have parent bridges
schedule_connect_if_needed(ParentBridges),
%% Schedule periodic health checks
erlang:send_after(?HEALTH_CHECK_INTERVAL, self(), health_check),
{ok, State}.
handle_call({escalate_query, Query}, _From, State) ->
{Reply, NewState} = do_escalate_query(Query, State),
{reply, Reply, NewState};
handle_call({store_to_parent, StoreMsg}, _From, State) ->
{Reply, NewState} = do_store_to_parent(StoreMsg, State),
{reply, Reply, NewState};
handle_call(is_connected, _From, #state{connected_parent = Parent} = State) ->
{reply, Parent =/= undefined, State};
handle_call(get_stats, _From, #state{stats = Stats, mesh_level = Level,
connected_parent = Parent} = State) ->
FullStats = Stats#{
mesh_level => Level,
connected_parent => Parent,
is_connected => Parent =/= undefined
},
{reply, {ok, FullStats}, State};
handle_call(get_parent_bridges, _From, #state{parent_bridges = Bridges} = State) ->
{reply, Bridges, State};
handle_call({set_parent_bridges, Bridges}, _From, State) ->
?LOG_INFO("[BridgeNode] Updating parent bridges to: ~p", [Bridges]),
NewState = State#state{parent_bridges = Bridges},
schedule_connect_if_needed(Bridges),
{reply, ok, NewState};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
handle_cast(_Request, State) ->
{noreply, State}.
handle_info(connect_to_parent, #state{parent_bridges = []} = State) ->
%% No parent bridges configured - nothing to connect to
{noreply, State};
handle_info(connect_to_parent, #state{parent_bridges = Bridges,
connected_parent = undefined} = State) ->
NewState = try_connect_to_parent(Bridges, State),
{noreply, NewState};
handle_info(connect_to_parent, State) ->
%% Already connected
{noreply, State};
handle_info(health_check, State) ->
NewState = perform_health_check(State),
erlang:send_after(?HEALTH_CHECK_INTERVAL, self(), health_check),
{noreply, NewState};
handle_info({'DOWN', _Ref, process, Pid, Reason},
#state{connection_pid = Pid} = State) ->
?LOG_WARNING("[BridgeNode] Parent connection lost: ~p", [Reason]),
NewState = State#state{
connected_parent = undefined,
connection_pid = undefined,
stats = increment_stat(disconnections, State#state.stats)
},
%% Schedule reconnection
erlang:send_after(?RECONNECT_INTERVAL, self(), connect_to_parent),
{noreply, NewState};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Initialize statistics map.
-spec init_stats() -> map().
init_stats() ->
#{
queries_escalated => 0,
queries_successful => 0,
queries_failed => 0,
stores_propagated => 0,
cache_hits => 0,
disconnections => 0,
started_at => erlang:system_time(second)
}.
%% @doc Increment a statistic counter.
-spec increment_stat(atom(), map()) -> map().
increment_stat(Key, Stats) ->
maps:update_with(Key, fun(V) -> V + 1 end, 1, Stats).
%% @doc Schedule connection attempt if bridges are configured.
-spec schedule_connect_if_needed([binary()]) -> ok.
schedule_connect_if_needed([]) -> ok;
schedule_connect_if_needed(_Bridges) ->
erlang:send_after(100, self(), connect_to_parent),
ok.
%% @doc Try to connect to one of the parent bridges.
-spec try_connect_to_parent([binary()], #state{}) -> #state{}.
try_connect_to_parent([], State) ->
?LOG_WARNING("[BridgeNode] Failed to connect to any parent bridge"),
erlang:send_after(?RECONNECT_INTERVAL, self(), connect_to_parent),
State;
try_connect_to_parent([Bridge | Rest], State) ->
?LOG_INFO("[BridgeNode] Attempting to connect to parent bridge: ~p", [Bridge]),
case connect_to_bridge(Bridge) of
{ok, ConnPid} ->
?LOG_INFO("[BridgeNode] Connected to parent bridge: ~p", [Bridge]),
erlang:monitor(process, ConnPid),
State#state{
connected_parent = Bridge,
connection_pid = ConnPid
};
{error, Reason} ->
?LOG_WARNING("[BridgeNode] Failed to connect to ~p: ~p", [Bridge, Reason]),
try_connect_to_parent(Rest, State)
end.
%% @doc Establish QUIC connection to parent bridge.
-spec connect_to_bridge(binary()) -> {ok, pid()} | {error, term()}.
connect_to_bridge(BridgeEndpoint) ->
%% Parse endpoint (format: "quic://host:port" or "host:port")
case parse_bridge_endpoint(BridgeEndpoint) of
{ok, Host, Port} ->
%% Use macula_connection to establish QUIC connection
case macula_connection:connect(Host, Port, #{}) of
{ok, ConnPid} -> {ok, ConnPid};
{error, _} = Error -> Error
end;
{error, _} = Error ->
Error
end.
%% @doc Parse bridge endpoint string.
-spec parse_bridge_endpoint(binary()) -> {ok, string(), pos_integer()} | {error, term()}.
parse_bridge_endpoint(Endpoint) when is_binary(Endpoint) ->
parse_bridge_endpoint(binary_to_list(Endpoint));
parse_bridge_endpoint("quic://" ++ Rest) ->
parse_host_port(Rest);
parse_bridge_endpoint(HostPort) ->
parse_host_port(HostPort).
parse_host_port(HostPort) ->
case string:split(HostPort, ":") of
[Host, PortStr] ->
case catch list_to_integer(PortStr) of
Port when is_integer(Port), Port > 0, Port < 65536 ->
{ok, Host, Port};
_ ->
{error, invalid_port}
end;
_ ->
{error, invalid_endpoint}
end.
%% @doc Perform DHT query escalation to parent.
-spec do_escalate_query(map(), #state{}) -> {{ok, term()} | {error, term()}, #state{}}.
do_escalate_query(_Query, #state{connected_parent = undefined} = State) ->
{{error, not_connected}, increment_stats(queries_failed, State)};
do_escalate_query(Query, #state{connection_pid = ConnPid,
escalation_timeout = Timeout} = State) ->
Stats1 = increment_stat(queries_escalated, State#state.stats),
%% Send query to parent bridge via QUIC connection
case send_dht_query(ConnPid, Query, Timeout) of
{ok, Result} ->
%% Cache result locally via bridge_cache
cache_result(Query, Result),
Stats2 = increment_stat(queries_successful, Stats1),
{{ok, Result}, State#state{stats = Stats2}};
{error, Reason} ->
?LOG_WARNING("[BridgeNode] Query escalation failed: ~p", [Reason]),
Stats2 = increment_stat(queries_failed, Stats1),
{{error, Reason}, State#state{stats = Stats2}}
end.
%% @doc Increment stats helper.
-spec increment_stats(atom(), #state{}) -> #state{}.
increment_stats(Key, #state{stats = Stats} = State) ->
State#state{stats = increment_stat(Key, Stats)}.
%% @doc Send DHT query via QUIC connection.
-spec send_dht_query(pid(), map(), pos_integer()) -> {ok, term()} | {error, term()}.
send_dht_query(ConnPid, Query, Timeout) ->
%% Encode query as DHT protocol message
QueryType = maps:get(type, Query, find_value),
Key = maps:get(key, Query, <<>>),
%% Use internal RPC to query parent's DHT
Procedure = <<"_dht.find_value">>,
Args = #{<<"key">> => Key, <<"query_type">> => QueryType},
case macula_connection:call(ConnPid, Procedure, Args, #{timeout => Timeout}) of
{ok, #{<<"value">> := Value}} -> {ok, Value};
{ok, #{<<"nodes">> := Nodes}} -> {nodes, Nodes};
{error, _} = Error -> Error
end.
%% @doc Cache query result locally.
-spec cache_result(map(), term()) -> ok.
cache_result(Query, Result) ->
case whereis(macula_bridge_cache) of
undefined -> ok;
CachePid ->
Key = maps:get(key, Query, <<>>),
macula_bridge_cache:put(CachePid, Key, Result)
end.
%% @doc Store value to parent DHT.
-spec do_store_to_parent(map(), #state{}) -> {ok | {error, term()}, #state{}}.
do_store_to_parent(_StoreMsg, #state{connected_parent = undefined} = State) ->
{{error, not_connected}, State};
do_store_to_parent(StoreMsg, #state{connection_pid = ConnPid} = State) ->
Key = maps:get(key, StoreMsg, <<>>),
Value = maps:get(value, StoreMsg, #{}),
%% Use internal RPC to store in parent's DHT
Procedure = <<"_dht.store">>,
Args = #{<<"key">> => Key, <<"value">> => Value},
case macula_connection:call(ConnPid, Procedure, Args, #{timeout => 10000}) of
{ok, _} ->
Stats = increment_stat(stores_propagated, State#state.stats),
{ok, State#state{stats = Stats}};
{error, _} = Error ->
{Error, State}
end.
%% @doc Perform periodic health check on parent connection.
-spec perform_health_check(#state{}) -> #state{}.
perform_health_check(#state{connected_parent = undefined} = State) ->
State;
perform_health_check(#state{connection_pid = ConnPid} = State) ->
case is_process_alive(ConnPid) of
true -> State;
false ->
?LOG_WARNING("[BridgeNode] Parent connection process dead, reconnecting"),
erlang:send_after(100, self(), connect_to_parent),
State#state{
connected_parent = undefined,
connection_pid = undefined
}
end.