Packages
macula
0.34.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
Current section
Files
src/macula_routing_system/macula_routing_bucket.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% K-bucket for Kademlia routing table.
%%% Stores up to k nodes with LRU eviction policy.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_routing_bucket).
%% API
-export([
new/1,
add_node/2,
remove_node/2,
remove_by_endpoint/2,
remove_ghost_by_endpoint/3,
evict_stale/2,
get_nodes/1,
find_node/2,
find_closest/3,
has_node/2,
update_timestamp/2,
get_endpoint/1,
size/1,
capacity/1
]).
%% Types
-type node_info() :: #{
node_id := binary(),
address := {inet:ip_address(), inet:port_number()},
last_seen => integer() % Optional timestamp
}.
-type bucket() :: #{
capacity := pos_integer(),
nodes := [node_info()] % Ordered: head = oldest, tail = most recent
}.
-export_type([bucket/0, node_info/0]).
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Create a new bucket with capacity k.
-spec new(pos_integer()) -> bucket().
new(Capacity) ->
#{capacity => Capacity, nodes => []}.
%% @doc Add a node to the bucket.
%% If node exists (same node_id), move to tail (most recent).
%% If endpoint exists with different node_id, replace the stale entry.
%% If bucket full, replace oldest stale entry or return {error, bucket_full}.
-spec add_node(bucket(), node_info()) -> bucket() | {error, bucket_full}.
add_node(Bucket, NodeInfo) ->
NodeWithTimestamp = ensure_timestamp(NodeInfo),
NodeId = maps:get(node_id, NodeWithTimestamp),
Endpoint = get_endpoint(NodeWithTimestamp),
%% First, remove any entry with the same endpoint but different node_id (ghost cleanup)
CleanBucket = remove_ghost_by_endpoint(Bucket, Endpoint, NodeId),
do_add_node(CleanBucket, NodeId, NodeWithTimestamp).
%% Node already exists - move to tail
do_add_node(#{nodes := Nodes} = Bucket, NodeId, NodeInfo) when is_list(Nodes) ->
case lists:keymember(NodeId, 1, nodes_to_tuples(Nodes)) of
true ->
UpdatedNodes = remove_by_id(Nodes, NodeId) ++ [NodeInfo],
Bucket#{nodes => UpdatedNodes};
false ->
add_new_node(Bucket, NodeInfo)
end.
%% Bucket has space - add new node
add_new_node(#{capacity := Capacity, nodes := Nodes} = Bucket, NodeInfo)
when length(Nodes) < Capacity ->
Bucket#{nodes => Nodes ++ [NodeInfo]};
%% Bucket full - try replacing oldest stale entry (not seen in 5 min)
add_new_node(#{nodes := [Oldest | Rest]} = Bucket, NodeInfo) ->
StaleThreshold = erlang:system_time(millisecond) - 300000,
OldestSeen = maps:get(last_seen, Oldest, 0),
replace_if_stale(OldestSeen, StaleThreshold, Bucket, Rest, NodeInfo).
replace_if_stale(OldestSeen, StaleThreshold, Bucket, Rest, NodeInfo)
when OldestSeen < StaleThreshold ->
Bucket#{nodes => Rest ++ [NodeInfo]};
replace_if_stale(_OldestSeen, _StaleThreshold, _Bucket, _Rest, _NodeInfo) ->
{error, bucket_full}.
%% Add timestamp if not present
ensure_timestamp(#{last_seen := _} = NodeInfo) ->
NodeInfo;
ensure_timestamp(NodeInfo) ->
NodeInfo#{last_seen => erlang:system_time(millisecond)}.
%% @doc Remove a node from the bucket.
-spec remove_node(bucket(), binary()) -> bucket().
remove_node(#{nodes := Nodes} = Bucket, NodeId) ->
Bucket#{nodes => remove_by_id(Nodes, NodeId)}.
%% @doc Remove all nodes matching an endpoint (regardless of node_id).
-spec remove_by_endpoint(bucket(), binary() | undefined) -> bucket().
remove_by_endpoint(Bucket, undefined) ->
Bucket;
remove_by_endpoint(#{nodes := Nodes} = Bucket, Endpoint) ->
Bucket#{nodes => [N || N <- Nodes, get_endpoint(N) =/= Endpoint]}.
%% @doc Remove nodes not seen since StaleThreshold (millisecond timestamp).
-spec evict_stale(bucket(), integer()) -> bucket().
evict_stale(#{nodes := Nodes} = Bucket, StaleThreshold) ->
Bucket#{nodes => [N || N <- Nodes, maps:get(last_seen, N, 0) >= StaleThreshold]}.
%% @doc Get all nodes in the bucket (ordered: oldest first).
-spec get_nodes(bucket()) -> [node_info()].
get_nodes(#{nodes := Nodes}) ->
Nodes.
%% @doc Find a node by ID.
-spec find_node(bucket(), binary()) -> {ok, node_info()} | not_found.
find_node(#{nodes := Nodes}, NodeId) ->
find_by_id(Nodes, NodeId).
%% @doc Find n closest nodes to target (sorted by XOR distance).
-spec find_closest(bucket(), binary(), pos_integer()) -> [node_info()].
find_closest(#{nodes := Nodes}, Target, N) ->
WithDistance = [{distance_to(Target, Node), Node} || Node <- Nodes],
Sorted = lists:keysort(1, WithDistance),
[Node || {_Dist, Node} <- lists:sublist(Sorted, N)].
%% @doc Check if bucket contains node.
-spec has_node(bucket(), binary()) -> boolean().
has_node(#{nodes := Nodes}, NodeId) ->
lists:keymember(NodeId, 1, nodes_to_tuples(Nodes)).
%% @doc Update node's last_seen timestamp (moves to tail).
-spec update_timestamp(bucket(), binary()) -> bucket().
update_timestamp(#{nodes := Nodes} = Bucket, NodeId) ->
do_update_timestamp(Bucket, Nodes, NodeId).
do_update_timestamp(Bucket, Nodes, NodeId) ->
case find_by_id(Nodes, NodeId) of
{ok, Node} ->
UpdatedNode = Node#{last_seen => erlang:system_time(millisecond)},
Bucket#{nodes => remove_by_id(Nodes, NodeId) ++ [UpdatedNode]};
not_found ->
Bucket
end.
%% @doc Get number of nodes in bucket.
-spec size(bucket()) -> non_neg_integer().
size(#{nodes := Nodes}) ->
length(Nodes).
%% @doc Get bucket capacity.
-spec capacity(bucket()) -> pos_integer().
capacity(#{capacity := Capacity}) ->
Capacity.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @doc Convert nodes to tuples for efficient key-based lookups.
nodes_to_tuples(Nodes) ->
[{maps:get(node_id, N), N} || N <- Nodes].
%% @doc Find node by ID using list comprehension.
-spec find_by_id([node_info()], binary()) -> {ok, node_info()} | not_found.
find_by_id(Nodes, NodeId) ->
case [N || #{node_id := Id} = N <- Nodes, Id =:= NodeId] of
[Node | _] -> {ok, Node};
[] -> not_found
end.
%% @doc Remove node by ID using list comprehension.
-spec remove_by_id([node_info()], binary()) -> [node_info()].
remove_by_id(Nodes, NodeId) ->
[N || #{node_id := Id} = N <- Nodes, Id =/= NodeId].
%% @doc Calculate XOR distance between target and node.
%% Returns raw XOR binary (32 bytes) for Kademlia comparison.
-spec distance_to(binary(), node_info()) -> binary().
distance_to(Target, #{node_id := NodeId}) ->
macula_routing_nodeid:distance(Target, NodeId).
%% @doc Extract endpoint from node info (handles multiple key formats).
-spec get_endpoint(node_info()) -> binary() | undefined.
get_endpoint(#{endpoint := Endpoint}) when is_binary(Endpoint) -> Endpoint;
get_endpoint(#{<<"endpoint">> := Endpoint}) when is_binary(Endpoint) -> Endpoint;
get_endpoint(#{address := Address}) when is_binary(Address) -> Address;
get_endpoint(#{<<"address">> := Address}) when is_binary(Address) -> Address;
get_endpoint(_) -> undefined.
%% @doc Remove entries with matching endpoint but different node_id (ghost cleanup).
-spec remove_ghost_by_endpoint(bucket(), binary() | undefined, binary()) -> bucket().
remove_ghost_by_endpoint(Bucket, undefined, _NodeId) ->
Bucket;
remove_ghost_by_endpoint(#{nodes := Nodes} = Bucket, Endpoint, NodeId) ->
Bucket#{nodes => [N || N <- Nodes,
not (get_endpoint(N) =:= Endpoint andalso
maps:get(node_id, N) =/= NodeId)]}.