Packages

macula

0.48.4
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
macula src macula_relay_discovery.erl
Raw

src/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,
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()}].
ranked_relays() ->
[#{hostname => E#relay_entry.hostname,
url => E#relay_entry.url,
distance_km => E#relay_entry.distance_km,
status => E#relay_entry.status}
|| E <- ranked_all()].
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, [{verify, verify_none}]}], []) of
{ok, {{_, 200, _}, _, Body}} ->
parse_topology_relays(Body);
Other ->
?LOG_WARNING("[relay_discovery] Topology fetch failed ~s: ~p", [Url, Other]),
{error, fetch_failed}
end.
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.