Packages

macula

0.45.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
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
]).
%% 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(BOOTSTRAP_TIMEOUT_MS, 15_000).
-record(relay_entry, {
hostname :: binary(),
url :: binary(),
lat :: float(),
lng :: float(),
status :: online | offline,
distance_km :: float()
}).
%%====================================================================
%% 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,
{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};
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">>,
Entry = #relay_entry{
hostname = Hostname,
url = Url,
lat = RelayLat,
lng = RelayLng,
status = Status,
distance_km = Distance
},
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)
%%====================================================================
ranked_all() ->
Entries = ets:tab2list(?TABLE),
lists:sort(fun(A, B) -> A#relay_entry.distance_km =< B#relay_entry.distance_km end, Entries).
ranked_online() ->
[E || E <- ranked_all(), E#relay_entry.status =:= online].
ranked_online_except(FailedHostname) ->
[E || E <- ranked_online(), E#relay_entry.hostname =/= FailedHostname].
%%====================================================================
%% Helpers
%%====================================================================
schedule_reconcile() ->
erlang:send_after(?RECONCILE_MS, self(), reconcile).
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.