Packages

macula

4.7.0
7.1.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 resolve_address macula_resolve_quorum.erl
Raw

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}.