Packages

macula

0.31.8
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_chatter.erl
Raw

src/macula_chatter.erl

%%%-------------------------------------------------------------------
%%% @doc Macula Chatter - P2P PubSub Demo for NAT Traversal Testing
%%%
%%% A pub/sub chat application that demonstrates broadcast messaging
%%% across NAT boundaries using Macula's pub/sub capabilities.
%%%
%%% Each chatter node:
%%% - Subscribes to "chat.room.global" topic
%%% - Periodically broadcasts numbered messages to all peers
%%% - Tracks delivery metrics per peer (by NAT type)
%%% - Reports delivery rates at shutdown
%%%
%%% PubSub Delivery Metrics:
%%% - Each broadcast includes a sequence number
%%% - Receivers track which sequence numbers they've seen per sender
%%% - Gaps in sequence numbers indicate missed messages
%%% - Delivery rate = received / expected (based on max seq seen)
%%%
%%% @end
%%%-------------------------------------------------------------------
-module(macula_chatter).
-behaviour(gen_server).
%% API
-export([start_link/0, start_link/1]).
-export([send_message/1, send_direct/2, get_stats/0]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-include_lib("kernel/include/logger.hrl").
-define(SERVER, ?MODULE).
-define(CHAT_TOPIC, <<"chat.room.global">>).
-define(CHAT_RPC, <<"chat.receive">>).
-define(DEFAULT_INTERVAL, 5000). % 5 seconds between messages
%% Per-peer tracking for delivery metrics
-record(peer_tracker, {
nat_type :: binary() | undefined,
max_seq = 0 :: non_neg_integer(), % Highest sequence number seen
received_count = 0 :: non_neg_integer(), % Total messages received
first_seen :: non_neg_integer() | undefined,
last_seen :: non_neg_integer() | undefined
}).
-record(state, {
node_id :: binary(),
nat_type :: binary() | undefined,
messages_sent = 0 :: non_neg_integer(),
messages_received = 0 :: non_neg_integer(),
peer_trackers = #{} :: #{binary() => #peer_tracker{}}, % peer_id => tracker
interval :: pos_integer(),
timer_ref :: reference() | undefined
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the chatter with default settings
-spec start_link() -> {ok, pid()} | {error, term()}.
start_link() ->
start_link(#{}).
%% @doc Start the chatter with options
%% Options:
%% - interval: milliseconds between broadcasts (default: 5000)
%% - node_id: custom node identifier (default: hostname)
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link({local, ?SERVER}, ?MODULE, Opts, []).
%% @doc Send a message to all peers via pubsub
-spec send_message(binary()) -> ok.
send_message(Message) ->
gen_server:cast(?SERVER, {send_message, Message}).
%% @doc Send a direct message to a specific peer via RPC
-spec send_direct(binary(), binary()) -> ok | {error, term()}.
send_direct(PeerId, Message) ->
gen_server:call(?SERVER, {send_direct, PeerId, Message}).
%% @doc Get statistics about messages sent/received
-spec get_stats() -> map().
get_stats() ->
gen_server:call(?SERVER, get_stats).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(Opts) ->
NodeId = get_node_id(Opts),
Interval = maps:get(interval, Opts, ?DEFAULT_INTERVAL),
NatType = get_nat_type(),
macula_console:info(NodeId, <<"PubSub chatter starting...">>),
%% Schedule setup after gen_server is fully started
self() ! setup,
{ok, #state{
node_id = NodeId,
nat_type = NatType,
interval = Interval
}}.
handle_call(get_stats, _From, State) ->
Stats = #{
node_id => State#state.node_id,
nat_type => State#state.nat_type,
messages_sent => State#state.messages_sent,
messages_received => State#state.messages_received,
peer_stats => format_peer_stats(State#state.peer_trackers),
delivery_by_nat => calc_delivery_by_nat(State#state.peer_trackers),
uptime_ms => erlang:system_time(millisecond)
},
{reply, Stats, State};
handle_call({send_direct, PeerId, Message}, _From, State) ->
Result = do_send_direct(PeerId, Message, State),
{reply, Result, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_call}, State}.
handle_cast({send_message, Message}, State) ->
NewState = do_broadcast(Message, State),
{noreply, NewState};
handle_cast(_Msg, State) ->
{noreply, State}.
handle_info(setup, State) ->
NewState = setup_chatter(State),
{noreply, NewState};
handle_info(broadcast_tick, State) ->
%% Sequence number for delivery tracking
SeqNum = State#state.messages_sent + 1,
NewState = do_broadcast(SeqNum, State),
%% Schedule next broadcast
TimerRef = erlang:send_after(State#state.interval, self(), broadcast_tick),
{noreply, NewState#state{timer_ref = TimerRef}};
handle_info({pubsub_message, FromNode, SeqNum, SenderNat}, State) ->
%% Update peer tracker with sequence number for delivery metrics
Now = erlang:system_time(millisecond),
Trackers = State#state.peer_trackers,
Tracker = maps:get(FromNode, Trackers, #peer_tracker{}),
UpdatedTracker = Tracker#peer_tracker{
nat_type = SenderNat,
max_seq = max(SeqNum, Tracker#peer_tracker.max_seq),
received_count = Tracker#peer_tracker.received_count + 1,
first_seen = case Tracker#peer_tracker.first_seen of
undefined -> Now;
T -> T
end,
last_seen = Now
},
%% Log with colored output
DeliveryRate = calc_peer_delivery_rate(UpdatedTracker),
macula_console:pubsub_recv(State#state.node_id, FromNode, SeqNum, SenderNat, DeliveryRate),
NewState = State#state{
messages_received = State#state.messages_received + 1,
peer_trackers = maps:put(FromNode, UpdatedTracker, Trackers)
},
{noreply, NewState};
%% Legacy handler for RPC-based chat messages (backwards compatibility)
handle_info({chat_message, FromNode, _Message}, State) ->
Trackers = State#state.peer_trackers,
Tracker = maps:get(FromNode, Trackers, #peer_tracker{}),
Now = erlang:system_time(millisecond),
UpdatedTracker = Tracker#peer_tracker{
received_count = Tracker#peer_tracker.received_count + 1,
first_seen = case Tracker#peer_tracker.first_seen of
undefined -> Now;
T -> T
end,
last_seen = Now
},
NewState = State#state{
messages_received = State#state.messages_received + 1,
peer_trackers = maps:put(FromNode, UpdatedTracker, Trackers)
},
{noreply, NewState};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, State) ->
NodeId = State#state.node_id,
macula_console:info(NodeId, <<"Shutting down, printing delivery stats...">>),
print_delivery_stats(State),
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @private Get node identifier
get_node_id(Opts) ->
case maps:get(node_id, Opts, undefined) of
undefined ->
case os:getenv("NODE_ID") of
false ->
{ok, Hostname} = inet:gethostname(),
list_to_binary(Hostname);
NodeId ->
list_to_binary(NodeId)
end;
NodeId when is_binary(NodeId) ->
NodeId;
NodeId when is_list(NodeId) ->
list_to_binary(NodeId)
end.
%% @private Setup subscriptions and RPC handlers
setup_chatter(State) ->
NodeId = State#state.node_id,
io:format("[Chatter ~s] Setting up pub/sub and RPC handlers...~n", [NodeId]),
%% Wait a bit for the mesh to stabilize
timer:sleep(2000),
%% Try to subscribe to chat topic
case setup_pubsub(NodeId) of
ok ->
io:format("[Chatter ~s] Subscribed to ~s~n", [NodeId, ?CHAT_TOPIC]);
{error, SubReason} ->
io:format("[Chatter ~s] Failed to subscribe: ~p~n", [NodeId, SubReason])
end,
%% Try to register RPC handler
case setup_rpc(NodeId) of
ok ->
io:format("[Chatter ~s] Registered RPC handler ~s~n", [NodeId, ?CHAT_RPC]);
{error, RpcReason} ->
io:format("[Chatter ~s] Failed to register RPC: ~p~n", [NodeId, RpcReason])
end,
%% Start broadcasting
io:format("[Chatter ~s] Starting broadcasts every ~pms~n", [NodeId, State#state.interval]),
TimerRef = erlang:send_after(State#state.interval, self(), broadcast_tick),
State#state{timer_ref = TimerRef}.
%% @private Setup pub/sub subscription
setup_pubsub(NodeId) ->
Self = self(),
%% Callback receives a single map: #{topic, matched_pattern, payload}
Handler = fun(#{payload := Payload}) ->
case decode_chat_message(Payload) of
{ok, FromNode, SeqNum, SenderNat} when FromNode =/= NodeId ->
Self ! {pubsub_message, FromNode, SeqNum, SenderNat};
{ok, _FromNode, _SeqNum, _SenderNat} ->
%% Ignore our own messages
ok;
{error, _Reason} ->
ok
end
end,
case get_peer_handlers() of
{ok, #{pubsub := PubSubPid}} ->
macula_pubsub_handler:subscribe(PubSubPid, ?CHAT_TOPIC, Handler),
ok;
{error, Reason} ->
{error, Reason}
end.
%% @private Setup RPC handler for direct messages
setup_rpc(NodeId) ->
Self = self(),
Handler = fun(Args) ->
FromNode = maps:get(<<"from">>, Args, <<"unknown">>),
Message = maps:get(<<"message">>, Args, <<"">>),
%% Only process if not from ourselves
case FromNode of
NodeId -> ok;
_ -> Self ! {chat_message, FromNode, Message}
end,
{ok, #{<<"status">> => <<"delivered">>, <<"to">> => NodeId}}
end,
case get_peer_handlers() of
{ok, #{rpc := RpcPid}} ->
macula_rpc_handler:register_local_procedure(RpcPid, ?CHAT_RPC, Handler),
ok;
{error, Reason} ->
{error, Reason}
end.
%% @private Broadcast a message to all peers
do_broadcast(SeqNum, State) ->
NodeId = State#state.node_id,
NatType = State#state.nat_type,
Payload = encode_chat_message(NodeId, SeqNum, NatType),
case get_peer_handlers() of
{ok, #{pubsub := PubSubPid}} ->
macula_pubsub_handler:publish(PubSubPid, ?CHAT_TOPIC, Payload, #{}),
macula_console:pubsub_send(NodeId, SeqNum, NatType),
State#state{messages_sent = SeqNum};
{error, Reason} ->
macula_console:warning(NodeId, iolist_to_binary([
<<"No peer connection: ">>, io_lib:format("~p", [Reason])
])),
State
end.
%% @private Send a direct message to a specific peer
do_send_direct(_PeerId, Message, State) ->
NodeId = State#state.node_id,
Args = #{
<<"from">> => NodeId,
<<"message">> => Message
},
case get_peer_handlers() of
{ok, #{rpc := RpcPid}} ->
macula_rpc_handler:call(RpcPid, ?CHAT_RPC, Args);
{error, Reason} ->
{error, Reason}
end.
%% @private Get the first available peer system and extract handler PIDs
%% Returns {ok, #{pubsub => Pid, rpc => Pid, connection => Pid}} or {error, Reason}
get_peer_handlers() ->
case whereis(macula_peers_sup) of
undefined ->
{error, no_peers_sup};
_Pid ->
case macula_peers_sup:list_peers() of
[] ->
{error, no_peers};
[PeerSystemPid | _] ->
%% Get child PIDs from the peer system supervisor
get_handlers_from_supervisor(PeerSystemPid)
end
end.
%% @private Extract handler PIDs from peer system supervisor
get_handlers_from_supervisor(SupPid) ->
case supervisor:which_children(SupPid) of
Children when is_list(Children) ->
PubSubPid = find_child_pid(Children, pubsub_handler),
RpcPid = find_child_pid(Children, rpc_handler),
ConnPid = find_child_pid(Children, connection_manager),
case {PubSubPid, RpcPid, ConnPid} of
{undefined, _, _} -> {error, no_pubsub_handler};
{_, undefined, _} -> {error, no_rpc_handler};
{_, _, undefined} -> {error, no_connection_manager};
_ -> {ok, #{pubsub => PubSubPid, rpc => RpcPid, connection => ConnPid}}
end;
_ ->
{error, no_children}
end.
%% @private Find a child PID by child ID
find_child_pid(Children, ChildId) ->
case lists:keyfind(ChildId, 1, Children) of
{ChildId, Pid, _Type, _Modules} when is_pid(Pid) -> Pid;
_ -> undefined
end.
%% @private Encode a chat message for transmission
encode_chat_message(FromNode, SeqNum, NatType) ->
msgpack:pack(#{
<<"type">> => <<"pubsub_chat">>,
<<"from">> => FromNode,
<<"seq">> => SeqNum,
<<"nat_type">> => NatType,
<<"timestamp">> => erlang:system_time(millisecond)
}).
%% @private Decode a received chat message
decode_chat_message(Payload) ->
case msgpack:unpack(Payload) of
{ok, #{<<"type">> := <<"pubsub_chat">>, <<"from">> := From,
<<"seq">> := Seq, <<"nat_type">> := Nat}} ->
{ok, From, Seq, Nat};
%% Backwards compatibility with old format
{ok, #{<<"type">> := <<"chat">>, <<"from">> := From}} ->
{ok, From, 0, undefined};
{ok, _Other} ->
{error, invalid_format};
{error, Reason} ->
{error, Reason}
end.
%%%===================================================================
%%% Delivery Metrics
%%%===================================================================
%% @private Get NAT type from environment or cache
get_nat_type() ->
case os:getenv("NAT_TYPE") of
false ->
case whereis(macula_nat_cache) of
undefined -> undefined;
_Pid ->
case macula_nat_cache:get_local() of
{ok, #{mapping_policy := M, filtering_policy := F}} ->
classify_nat_type(M, F);
_ -> undefined
end
end;
Type ->
list_to_binary(Type)
end.
%% @private Classify NAT type from policies
classify_nat_type(endpoint_independent, endpoint_independent) -> <<"full_cone">>;
classify_nat_type(endpoint_independent, address_dependent) -> <<"restricted">>;
classify_nat_type(endpoint_independent, address_and_port_dependent) -> <<"port_restricted">>;
classify_nat_type(_, _) -> <<"symmetric">>.
%% @private Calculate delivery rate for a single peer
calc_peer_delivery_rate(#peer_tracker{max_seq = 0}) -> 0.0;
calc_peer_delivery_rate(#peer_tracker{max_seq = MaxSeq, received_count = Received}) ->
(Received / MaxSeq) * 100.0.
%% @private Format peer stats for get_stats/0
format_peer_stats(Trackers) ->
maps:map(fun(_PeerId, Tracker) ->
#{
nat_type => Tracker#peer_tracker.nat_type,
received => Tracker#peer_tracker.received_count,
max_seq => Tracker#peer_tracker.max_seq,
delivery_rate => calc_peer_delivery_rate(Tracker)
}
end, Trackers).
%% @private Calculate delivery rate grouped by NAT type
calc_delivery_by_nat(Trackers) ->
%% Group trackers by NAT type
Grouped = maps:fold(fun(_PeerId, Tracker, Acc) ->
NatType = case Tracker#peer_tracker.nat_type of
undefined -> <<"unknown">>;
T -> T
end,
Existing = maps:get(NatType, Acc, {0, 0}),
{TotalReceived, TotalExpected} = Existing,
NewReceived = TotalReceived + Tracker#peer_tracker.received_count,
NewExpected = TotalExpected + Tracker#peer_tracker.max_seq,
maps:put(NatType, {NewReceived, NewExpected}, Acc)
end, #{}, Trackers),
%% Calculate rates
maps:map(fun(_NatType, {Received, Expected}) ->
case Expected of
0 -> #{received => 0, expected => 0, rate => 0.0};
_ -> #{received => Received, expected => Expected,
rate => (Received / Expected) * 100.0}
end
end, Grouped).
%% @private Print delivery stats on shutdown
print_delivery_stats(State) ->
NodeId = State#state.node_id,
NatType = State#state.nat_type,
Trackers = State#state.peer_trackers,
io:format("~n"),
io:format("╔══════════════════════════════════════════════════════════════╗~n"),
io:format("║ PUBSUB DELIVERY STATISTICS ║~n"),
io:format("╠══════════════════════════════════════════════════════════════╣~n"),
io:format("║ Node: ~-20s NAT Type: ~-15s ║~n",
[NodeId, format_nat(NatType)]),
io:format("║ Messages Sent: ~-10B Messages Received: ~-10B ║~n",
[State#state.messages_sent, State#state.messages_received]),
io:format("║ Unique Peers: ~-10B ║~n",
[maps:size(Trackers)]),
io:format("╠══════════════════════════════════════════════════════════════╣~n"),
io:format("║ DELIVERY BY NAT TYPE: ║~n"),
ByNat = calc_delivery_by_nat(Trackers),
maps:foreach(fun(Nat, #{received := R, expected := E, rate := Rate}) ->
io:format("║ ~-15s: ~5B/~-5B (~6.1f%%) ║~n",
[format_nat(Nat), R, E, Rate])
end, ByNat),
io:format("╚══════════════════════════════════════════════════════════════╝~n"),
io:format("~n").
%% @private Format NAT type for display
format_nat(undefined) -> "unknown";
format_nat(Nat) when is_binary(Nat) -> binary_to_list(Nat);
format_nat(Nat) -> Nat.