Packages
macula
4.2.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
Current section
Files
src/resolve_address/macula_resolve_quorum.erl
%%%-------------------------------------------------------------------
%%% @doc Eclipse-mitigation wrapper for {@link macula_resolve_address}.
%%%
%%% Phase 4.4 — see PLAN_MACULA_NET_PHASE4_4_RESOLVE_QUORUM.md.
%%%
%%% Wraps a list of `find_fn' callbacks (each backed by a different
%%% DHT endpoint) into a single `find_fn' that fans out in parallel,
%%% requires M of N to return byte-equal records, and reports
%%% disagreement as `{error, eclipse_suspected}'.
%%%
%%% Signature-based comparison: two records "agree" iff their
%%% Ed25519 signatures are byte-equal — the strongest possible check
%%% short of full payload equality, and equivalent in practice
%%% because the signing operation is deterministic over the
%%% canonical CBOR encoding of the payload.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_resolve_quorum).
-export([wrap/2]).
-export_type([opts/0]).
-type find_fn() :: macula_resolve_address:find_fn().
-type opts() :: #{
threshold => pos_integer(),
timeout_ms => pos_integer(),
parallel => boolean()
}.
-define(DEFAULT_TIMEOUT_MS, 1000).
%% =============================================================================
%% Public API
%% =============================================================================
%% @doc Returns a `find_fn' that queries every Fn in `FindFns' (in
%% parallel by default) and applies a quorum check.
-spec wrap([find_fn()], opts()) -> find_fn().
wrap([], _Opts) ->
error(no_find_fns);
wrap([SingleFn], _Opts) when is_function(SingleFn, 1) ->
%% Degenerate single-endpoint case — passthrough with telemetry.
fun(Key) ->
Result = SingleFn(Key),
telemetry:execute([macula, net, resolve_quorum, decided],
#{count => 1},
#{outcome => <<"single_endpoint">>}),
Result
end;
wrap(FindFns, Opts) when is_list(FindFns) ->
N = length(FindFns),
Threshold = maps:get(threshold, Opts, default_threshold(N)),
Timeout = maps:get(timeout_ms, Opts, ?DEFAULT_TIMEOUT_MS),
Parallel = maps:get(parallel, Opts, true),
fun(Key) ->
Responses = collect(FindFns, Key, Parallel, Timeout),
decide(Responses, Threshold)
end.
default_threshold(N) ->
%% Strict majority: ceil((N + 1) / 2).
(N + 1) div 2 + ((N + 1) rem 2).
%% =============================================================================
%% Fanout
%% =============================================================================
collect(FindFns, Key, true, Timeout) ->
Self = self(),
Refs = [erlang:spawn_monitor(fun() ->
Self ! {self(), F(Key)}
end) || F <- FindFns],
gather(Refs, Timeout, []);
collect(FindFns, Key, false, _Timeout) ->
[F(Key) || F <- FindFns].
gather([], _Timeout, Acc) ->
Acc;
gather([{Pid, MonRef} | Rest], Timeout, Acc) ->
receive
{Pid, Result} ->
erlang:demonitor(MonRef, [flush]),
gather(Rest, Timeout, [Result | Acc]);
{'DOWN', MonRef, process, Pid, Reason} ->
gather(Rest, Timeout, [{error, {worker_died, Reason}} | Acc])
after Timeout ->
%% Per-worker timeout: kill the worker and tally as error.
exit(Pid, kill),
erlang:demonitor(MonRef, [flush]),
gather(Rest, Timeout, [{error, timeout} | Acc])
end.
%% =============================================================================
%% Decide
%% =============================================================================
decide(Responses, Threshold) ->
Records = [R || {ok, R} <- Responses],
case Records of
[] -> emit_and_return(<<"insufficient_responses">>);
_ -> tally_records(Records, Threshold)
end.
tally_records(Records, Threshold) ->
Tally = lists:foldl(
fun(Rec, Acc) ->
Sig = macula_record:signature(Rec),
maps:update_with(Sig, fun(V) -> V + 1 end, 1, Acc)
end, #{}, Records),
pick_winner(Tally, Records, Threshold).
pick_winner(Tally, Records, Threshold) ->
{WinnerSig, Count} =
lists:foldl(fun({Sig, C}, {_BestSig, BestC} = Best) ->
decide_better({Sig, C}, BestC, Best)
end,
{undefined, 0},
maps:to_list(Tally)),
pick_with_threshold(Count >= Threshold, WinnerSig, Records).
decide_better({Sig, C}, BestC, _Best) when C > BestC -> {Sig, C};
decide_better(_, _BestC, Best) -> Best.
pick_with_threshold(true, WinnerSig, Records) ->
%% Find any record with the winning signature — they're all equal.
[Winner | _] = [R || R <- Records,
macula_record:signature(R) =:= WinnerSig],
telemetry:execute([macula, net, resolve_quorum, decided],
#{count => 1},
#{outcome => <<"consistent">>}),
{ok, Winner};
pick_with_threshold(false, _WinnerSig, _Records) ->
emit_and_return(<<"disagreement">>).
emit_and_return(Outcome) ->
telemetry:execute([macula, net, resolve_quorum, decided],
#{count => 1},
#{outcome => Outcome}),
error_for(Outcome).
error_for(<<"insufficient_responses">>) -> {error, insufficient_responses};
error_for(<<"disagreement">>) -> {error, eclipse_suspected};
error_for(_) -> {error, unknown}.