Packages
macula
5.0.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/client/macula_relay_discovery.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Relay Discovery — geographic-aware relay selection.
%%%
%%% Maintains a ranked list of relay identities sorted by distance
%%% from the node's own location. Updated via three mechanisms:
%%%
%%% 1. Bootstrap: HTTP GET /topology from seed relays on startup
%%% 2. Real-time: _mesh.relay.up/down events via pubsub subscription
%%% 3. Periodic: Reconciliation poll every 5 minutes
%%%
%%% Nodes connect to the nearest available relay. On failover, the
%%% next nearest is selected instantly from the cached ranked list.
%%%
%%% Usage:
%%% {ok, Pid} = macula_relay_discovery:start_link(Opts).
%%% {ok, Url} = macula_relay_discovery:nearest().
%%% {ok, Url} = macula_relay_discovery:nearest_except(FailedHostname).
%%% Relays = macula_relay_discovery:ranked_relays().
%%% @end
%%%-------------------------------------------------------------------
-module(macula_relay_discovery).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
start_link/1,
nearest/0,
nearest_except/1,
ranked_relays/0,
lookup/1,
relay_count/0,
mark_offline/1
]).
%% Exported for tests
-export([
haversine_km/4,
extract_hostname/1,
build_topology_url/1,
relay_status/1,
to_float/1,
update_rtt/2
]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-define(SERVER, ?MODULE).
-define(TABLE, macula_relay_discovery_cache).
-define(RECONCILE_MS, 300_000). % 5 minutes
-define(STALENESS_CHECK_MS, 60_000). % 1 minute
-define(STALE_THRESHOLD_MS, 90_000). % 90 seconds without ping = stale
-define(BOOTSTRAP_TIMEOUT_MS, 15_000).
-record(relay_entry, {
hostname :: binary(),
url :: binary(),
lat :: float(),
lng :: float(),
status :: online | offline,
distance_km :: float(),
rtt_ms :: non_neg_integer() | undefined,
last_seen :: integer() | undefined % erlang:monotonic_time(millisecond)
}).
%%====================================================================
%% API
%%====================================================================
start_link(Opts) ->
gen_server:start_link({local, ?SERVER}, ?MODULE, Opts, []).
-spec nearest() -> {ok, binary()} | {error, no_relays}.
nearest() ->
case ranked_online() of
[#relay_entry{url = Url} | _] -> {ok, Url};
[] -> {error, no_relays}
end.
-spec nearest_except(binary()) -> {ok, binary()} | {error, no_relays}.
nearest_except(FailedHostname) ->
case ranked_online_except(FailedHostname) of
[#relay_entry{url = Url} | _] -> {ok, Url};
[] -> {error, no_relays}
end.
-spec ranked_relays() -> [#{hostname := binary(), url := binary(), distance_km := float(), status := atom(), lat := float(), lng := float(), rtt_ms := non_neg_integer() | undefined}].
ranked_relays() ->
[#{hostname => E#relay_entry.hostname,
url => E#relay_entry.url,
lat => E#relay_entry.lat,
lng => E#relay_entry.lng,
distance_km => E#relay_entry.distance_km,
rtt_ms => E#relay_entry.rtt_ms,
status => E#relay_entry.status}
|| E <- ranked_all()].
-spec lookup(binary()) -> {ok, map()} | {error, not_found}.
lookup(Hostname) ->
case ets:lookup(?TABLE, Hostname) of
[E] ->
{ok, #{hostname => E#relay_entry.hostname,
url => E#relay_entry.url,
lat => E#relay_entry.lat,
lng => E#relay_entry.lng,
distance_km => E#relay_entry.distance_km,
rtt_ms => E#relay_entry.rtt_ms,
status => E#relay_entry.status}};
[] ->
{error, not_found}
end.
relay_count() ->
case ets:info(?TABLE, size) of
undefined -> 0;
N -> N
end.
%%====================================================================
%% gen_server callbacks
%%====================================================================
init(Opts) ->
ets:new(?TABLE, [named_table, set, {keypos, #relay_entry.hostname},
public, {read_concurrency, true}]),
Seeds = maps:get(seeds, Opts, []),
MyLat = maps:get(lat, Opts, 0.0),
MyLng = maps:get(lng, Opts, 0.0),
State = #{seeds => Seeds, lat => MyLat, lng => MyLng, subscribed => false},
self() ! bootstrap,
schedule_staleness_check(),
{ok, State}.
handle_call(_, _From, State) ->
{reply, ok, State}.
handle_cast(_, State) ->
{noreply, State}.
%% Bootstrap — fetch topology from seeds
handle_info(bootstrap, #{seeds := Seeds, lat := Lat, lng := Lng} = State) ->
bootstrap_from_seeds(Seeds, Lat, Lng),
schedule_reconcile(),
{noreply, State};
%% Periodic reconciliation
handle_info(reconcile, #{seeds := Seeds, lat := Lat, lng := Lng} = State) ->
bootstrap_from_seeds(Seeds, Lat, Lng),
schedule_reconcile(),
{noreply, State};
%% Real-time relay events from mesh subscription
handle_info({relay_up, Hostname, RelayInfo}, #{lat := Lat, lng := Lng} = State) ->
add_relay(Hostname, RelayInfo, Lat, Lng),
{noreply, State};
handle_info({relay_down, Hostname}, State) ->
mark_offline(Hostname),
{noreply, State};
%% RTT measurement from relay_client PING/PONG
handle_info({ping_rtt, Hostname, RttMs}, State) ->
update_rtt(Hostname, RttMs),
{noreply, State};
%% Staleness check — detect silent relays (net split / ungraceful failure)
handle_info(check_staleness, State) ->
check_stale_relays(),
schedule_staleness_check(),
{noreply, State};
handle_info(_, State) ->
{noreply, State}.
terminate(_, _State) ->
ok.
%%====================================================================
%% Bootstrap — fetch /topology from seed URLs
%%====================================================================
bootstrap_from_seeds(Seeds, Lat, Lng) ->
Urls = build_topology_urls(Seeds),
Results = fetch_topologies(Urls),
process_topology_results(Results, Lat, Lng).
build_topology_urls(Seeds) ->
[build_topology_url(S) || S <- Seeds].
build_topology_url(Seed) when is_binary(Seed) ->
Host = extract_hostname(Seed),
"https://" ++ binary_to_list(Host) ++ "/topology?n=90&s=-90&e=180&w=-180&z=10".
extract_hostname(Url) ->
Stripped = re:replace(Url, "^https?://", "", [{return, binary}]),
re:replace(Stripped, ":\\d+.*", "", [{return, binary}]).
fetch_topologies(Urls) ->
lists:filtermap(fun(Url) ->
case fetch_topology(Url) of
{ok, Relays} -> {true, Relays};
_ -> false
end
end, Urls).
fetch_topology(Url) ->
case httpc:request(get, {Url, []},
[{timeout, ?BOOTSTRAP_TIMEOUT_MS},
{ssl, topology_ssl_opts()}], []) of
{ok, {{_, 200, _}, _, Body}} ->
parse_topology_relays(Body);
Other ->
?LOG_WARNING("[relay_discovery] Topology fetch failed ~s: ~p", [Url, Other]),
{error, fetch_failed}
end.
%% The topology response decides which relays a bootstrapping node
%% will dial — an unverified fetch lets a MITM steer the node onto
%% attacker relays. Verify against the OS trust store; the
%% `match_fun' is required for the fleet's wildcard (*.macula.io)
%% certs, which the default ssl hostname check rejects.
topology_ssl_opts() ->
[{verify, verify_peer},
{cacerts, public_key:cacerts_get()},
{depth, 3},
{customize_hostname_check,
[{match_fun, public_key:pkix_verify_hostname_match_fun(https)}]}].
parse_topology_relays(Body) ->
case json:decode(list_to_binary(Body)) of
#{<<"relays">> := Relays} when is_list(Relays) ->
{ok, Relays};
_ ->
{error, bad_format}
end.
process_topology_results([], _Lat, _Lng) ->
?LOG_WARNING("[relay_discovery] No topology data from any seed"),
ok;
process_topology_results(RelayLists, Lat, Lng) ->
AllRelays = lists:usort(fun(A, B) ->
maps:get(<<"hostname">>, A, <<>>) =< maps:get(<<"hostname">>, B, <<>>)
end, lists:append(RelayLists)),
lists:foreach(fun(R) -> ingest_relay(R, Lat, Lng) end, AllRelays),
?LOG_INFO("[relay_discovery] Loaded ~b relay identities", [length(AllRelays)]).
ingest_relay(RelayMap, MyLat, MyLng) ->
Hostname = maps:get(<<"hostname">>, RelayMap, <<>>),
RelayLat = to_float(maps:get(<<"lat">>, RelayMap, 0)),
RelayLng = to_float(maps:get(<<"lng">>, RelayMap, 0)),
Status = relay_status(maps:get(<<"status">>, RelayMap, <<"unknown">>)),
Distance = haversine_km(MyLat, MyLng, RelayLat, RelayLng),
Url = <<"https://", Hostname/binary, ":4433">>,
%% Preserve RTT and last_seen from live measurements
{ExistingRtt, ExistingLastSeen} = case ets:lookup(?TABLE, Hostname) of
[#relay_entry{rtt_ms = R, last_seen = L}] -> {R, L};
[] -> {undefined, undefined}
end,
Entry = #relay_entry{
hostname = Hostname, url = Url,
lat = RelayLat, lng = RelayLng,
status = Status, distance_km = Distance,
rtt_ms = ExistingRtt, last_seen = ExistingLastSeen
},
ets:insert(?TABLE, Entry).
%%====================================================================
%% Real-time updates
%%====================================================================
add_relay(Hostname, #{lat := Lat, lng := Lng} = _Info, MyLat, MyLng) ->
Distance = haversine_km(MyLat, MyLng, Lat, Lng),
Url = <<"https://", Hostname/binary, ":4433">>,
Entry = #relay_entry{
hostname = Hostname, url = Url,
lat = Lat, lng = Lng,
status = online, distance_km = Distance
},
ets:insert(?TABLE, Entry),
?LOG_DEBUG("[relay_discovery] Relay up: ~s (~.0f km)", [Hostname, Distance]).
mark_offline(Hostname) ->
case ets:lookup(?TABLE, Hostname) of
[Entry] ->
ets:insert(?TABLE, Entry#relay_entry{status = offline}),
?LOG_DEBUG("[relay_discovery] Relay down: ~s", [Hostname]);
[] ->
ok
end.
%%====================================================================
%% Ranked relay access (reads from ETS, no gen_server call)
%%====================================================================
%% Rank by measured RTT when available, fall back to haversine distance.
ranked_all() ->
Entries = ets:tab2list(?TABLE),
lists:sort(fun rank_compare/2, Entries).
rank_compare(#relay_entry{rtt_ms = A}, #relay_entry{rtt_ms = B})
when is_integer(A), is_integer(B) ->
A =< B;
rank_compare(#relay_entry{rtt_ms = A}, _) when is_integer(A) ->
true; %% measured RTT always preferred over estimate
rank_compare(_, #relay_entry{rtt_ms = B}) when is_integer(B) ->
false;
rank_compare(A, B) ->
A#relay_entry.distance_km =< B#relay_entry.distance_km.
ranked_online() ->
[E || E <- ranked_all(), E#relay_entry.status =:= online].
ranked_online_except(FailedHostname) ->
[E || E <- ranked_online(), E#relay_entry.hostname =/= FailedHostname].
%%====================================================================
%% RTT tracking
%%====================================================================
update_rtt(Hostname, RttMs) ->
Now = erlang:monotonic_time(millisecond),
case ets:lookup(?TABLE, Hostname) of
[Entry] ->
ets:insert(?TABLE, Entry#relay_entry{
rtt_ms = RttMs,
last_seen = Now,
status = online
});
[] ->
ok
end.
check_stale_relays() ->
Now = erlang:monotonic_time(millisecond),
Entries = ets:tab2list(?TABLE),
lists:foreach(fun(E) -> check_staleness_entry(E, Now) end, Entries).
check_staleness_entry(#relay_entry{status = offline}, _Now) ->
ok;
check_staleness_entry(#relay_entry{last_seen = undefined}, _Now) ->
ok;
check_staleness_entry(#relay_entry{hostname = H, last_seen = LastSeen} = E, Now) when
Now - LastSeen > ?STALE_THRESHOLD_MS ->
?LOG_WARNING("[relay_discovery] Relay stale (no ping for ~bs): ~s",
[(Now - LastSeen) div 1000, H]),
ets:insert(?TABLE, E#relay_entry{status = offline});
check_staleness_entry(_, _) ->
ok.
%%====================================================================
%% Helpers
%%====================================================================
schedule_reconcile() ->
erlang:send_after(?RECONCILE_MS, self(), reconcile).
schedule_staleness_check() ->
erlang:send_after(?STALENESS_CHECK_MS, self(), check_staleness).
relay_status(<<"online">>) -> online;
relay_status(_) -> offline.
to_float(V) when is_float(V) -> V;
to_float(V) when is_integer(V) -> float(V);
to_float(_) -> 0.0.
haversine_km(Lat1, Lng1, Lat2, Lng2) ->
R = 6371.0,
DLat = deg_to_rad(Lat2 - Lat1),
DLng = deg_to_rad(Lng2 - Lng1),
A = math:sin(DLat / 2) * math:sin(DLat / 2) +
math:cos(deg_to_rad(Lat1)) * math:cos(deg_to_rad(Lat2)) *
math:sin(DLng / 2) * math:sin(DLng / 2),
C = 2 * math:atan2(math:sqrt(A), math:sqrt(1 - A)),
R * C.
deg_to_rad(D) -> D * math:pi() / 180.0.