Packages
macula
0.12.3
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_chatter.erl
%%%-------------------------------------------------------------------
%%% @doc Macula Chatter - P2P Chat Demo for NAT Traversal Testing
%%%
%%% A simple chat application that demonstrates peer-to-peer messaging
%%% across NAT boundaries using Macula's pub/sub and RPC capabilities.
%%%
%%% Each chatter node:
%%% - Registers a "chat.receive" RPC handler
%%% - Subscribes to "chat.room.global" topic
%%% - Periodically broadcasts messages to all peers
%%% - Logs all received messages
%%%
%%% @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
-record(state, {
node_id :: binary(),
messages_sent = 0 :: non_neg_integer(),
messages_received = 0 :: non_neg_integer(),
peers_seen = #{} :: #{binary() => non_neg_integer()},
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),
io:format("[Chatter ~s] Starting up...~n", [NodeId]),
%% Schedule setup after gen_server is fully started
self() ! setup,
{ok, #state{
node_id = NodeId,
interval = Interval
}}.
handle_call(get_stats, _From, State) ->
Stats = #{
node_id => State#state.node_id,
messages_sent => State#state.messages_sent,
messages_received => State#state.messages_received,
peers_seen => State#state.peers_seen,
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) ->
%% Generate a random message
MsgNum = State#state.messages_sent + 1,
Message = iolist_to_binary([
<<"Hello from ">>, State#state.node_id,
<<" (#">>, integer_to_binary(MsgNum), <<")">>
]),
NewState = do_broadcast(Message, State),
%% Schedule next broadcast
TimerRef = erlang:send_after(State#state.interval, self(), broadcast_tick),
{noreply, NewState#state{timer_ref = TimerRef}};
handle_info({chat_message, FromNode, Message}, State) ->
io:format("[Chatter ~s] Received from ~s: ~s~n",
[State#state.node_id, FromNode, Message]),
%% Update stats
NewPeersSeen = maps:update_with(
FromNode,
fun(Count) -> Count + 1 end,
1,
State#state.peers_seen
),
NewState = State#state{
messages_received = State#state.messages_received + 1,
peers_seen = NewPeersSeen
},
{noreply, NewState};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, State) ->
io:format("[Chatter ~s] Shutting down. Stats: sent=~p, received=~p, peers=~p~n",
[State#state.node_id,
State#state.messages_sent,
State#state.messages_received,
maps:keys(State#state.peers_seen)]),
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, Message} when FromNode =/= NodeId ->
Self ! {chat_message, FromNode, Message};
{ok, _FromNode, _Message} ->
%% Ignore our own messages
ok;
{error, _Reason} ->
ok
end
end,
%% Find the bootstrap connection handlers
case get_peer_handlers() of
{ok, #{pubsub := PubSubPid}} ->
try
macula_pubsub_handler:subscribe(PubSubPid, ?CHAT_TOPIC, Handler),
ok
catch
_:Reason ->
{error, Reason}
end;
{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,
%% Find the bootstrap connection to register with
case get_peer_handlers() of
{ok, #{rpc := RpcPid}} ->
try
macula_rpc_handler:register_local_procedure(RpcPid, ?CHAT_RPC, Handler),
ok
catch
_:Reason ->
{error, Reason}
end;
{error, Reason} ->
{error, Reason}
end.
%% @private Broadcast a message to all peers
do_broadcast(Message, State) ->
NodeId = State#state.node_id,
Payload = encode_chat_message(NodeId, Message),
case get_peer_handlers() of
{ok, #{pubsub := PubSubPid}} ->
try
macula_pubsub_handler:publish(PubSubPid, ?CHAT_TOPIC, Payload, #{}),
io:format("[Chatter ~s] Broadcast: ~s~n", [NodeId, Message]),
State#state{messages_sent = State#state.messages_sent + 1}
catch
_:Reason ->
io:format("[Chatter ~s] Broadcast FAILED: ~p~n", [NodeId, Reason]),
State
end;
{error, Reason} ->
io:format("[Chatter ~s] No peer connection: ~p~n", [NodeId, 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}} ->
try
macula_rpc_handler:call(RpcPid, ?CHAT_RPC, Args)
catch
_:Reason ->
{error, Reason}
end;
{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, Message) ->
msgpack:pack(#{
<<"type">> => <<"chat">>,
<<"from">> => FromNode,
<<"message">> => Message,
<<"timestamp">> => erlang:system_time(millisecond)
}).
%% @private Decode a received chat message
decode_chat_message(Payload) ->
case msgpack:unpack(Payload) of
{ok, #{<<"type">> := <<"chat">>, <<"from">> := From, <<"message">> := Msg}} ->
{ok, From, Msg};
{ok, _Other} ->
{error, invalid_format};
{error, Reason} ->
{error, Reason}
end.