%%%------------------------------------------------------------------- %%% @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.