Packages

Denrei - a lightweight Erlang messaging system.

Current section

Files

Jump to
denrei src denrei_protocol.erl
Raw

src/denrei_protocol.erl

%%%-------------------------------------------------------------------
%%% @private
%%% @doc
%%% Denrei's protocol handler.
%%% @end
%%%-------------------------------------------------------------------
-module(denrei_protocol).
-behaviour(ranch_protocol).
%% Includes
-include("denrei.hrl").
%% ranch_protocol callbacks
-export([start_link/4, init/4]).
start_link(Ref, Socket, Transport, Opts) ->
Pid = spawn_link(?MODULE, init, [Ref, Socket, Transport, Opts]),
{ok, Pid}.
init(Ref, Socket, Transport, Opts) ->
ok = ranch:accept_ack(Ref),
ok = Transport:setopts(Socket, [binary, {packet, 4}]),
loop(Transport, Socket, proplists:get_value(tree, Opts)).
%%%===================================================================
%%% Internal functions
%%%===================================================================
loop(Transport, Socket, TreePid) ->
case Transport:recv(Socket, 0, 60000) of
{ok, Packet} ->
lager:debug([{transport, Transport}, {socket, Socket}], "~p received: ~p", [{Transport, Socket}, Packet]),
ok = handle_packet(Transport, Socket, TreePid, Packet),
loop(Transport, Socket, TreePid);
{error, closed} ->
lager:debug([{transport, Transport}, {socket, Socket}], "~p closed", [{Transport, Socket}]),
handle_close(Transport, Socket, TreePid);
{error, timeout} ->
lager:debug([{transport, Transport}, {socket, Socket}], "~p timed out", [{Transport, Socket}]),
ok = denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_PING>>, [denrei_utils:integer_to_binary(denrei_utils:timestamp(), 36)]),
loop(Transport, Socket, TreePid);
{error, Reason} ->
lager:warning([{transport, Transport}, {socket, Socket}], "~p error: ~p", [{Transport, Socket}, Reason]),
ok = Transport:close(Socket)
end.
handle_close(Transport, Socket, TreePid) ->
Subscriber = {Transport, Socket},
denrei_tree_server:delete(TreePid, Subscriber, all).
handle_packet(_Transport, _Socket, TreePid, Packet = <<?DENREI_PROTOCOL_PUBLISH, Fields/binary>>) ->
{Subject, _} = denrei_utils:get_next_field(Fields),
{ok, Subscribers} = denrei_tree_server:match(TreePid, denrei_utils:tokenize_subject(Subject)),
[denrei_utils:send(T, S, Packet) || {T, S} <- Subscribers],
ok;
handle_packet(Transport, Socket, TreePid, <<?DENREI_PROTOCOL_SUBSCRIBE, Fields/binary>>) ->
Subscriber = {Transport, Socket},
{Subject, <<>>} = denrei_utils:get_next_field(Fields),
lager:info([{transport, Transport}, {socket, Socket}], "~p subscribe to ~p", [Subscriber, Subject]),
denrei_tree_server:insert(TreePid, Subscriber, denrei_utils:tokenize_subject(Subject));
handle_packet(Transport, Socket, TreePid, <<?DENREI_PROTOCOL_UNSUBSCRIBE>>) ->
Subscriber = {Transport, Socket},
lager:info([{transport, Transport}, {socket, Socket}], "~p unsubscribe from all", [Subscriber]),
denrei_tree_server:delete(TreePid, Subscriber, all);
handle_packet(Transport, Socket, TreePid, <<?DENREI_PROTOCOL_UNSUBSCRIBE, Fields/binary>>) ->
Subscriber = {Transport, Socket},
{Subject, <<>>} = denrei_utils:get_next_field(Fields),
lager:info([{transport, Transport}, {socket, Socket}], "~p unsubscribe from ~p", [Subject]),
denrei_tree_server:delete(TreePid, Subscriber, denrei_utils:tokenize_subject(Subject));
handle_packet(Transport, Socket, TreePid, <<?DENREI_PROTOCOL_SUBSCRIBERS, Fields0/binary>>) ->
{Subject, Fields1} = denrei_utils:get_next_field(Fields0),
{CorrelationID, _} = denrei_utils:get_next_field(Fields1),
case denrei_tree_server:match(TreePid, denrei_utils:tokenize_subject(Subject)) of
{ok, []} ->
denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_RESPONSE>>, [CorrelationID, <<0:?DENREI_COUNT_BITS>>]);
{ok, Subscribers} ->
Count = length(Subscribers),
denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_RESPONSE>>, [CorrelationID, <<Count:?DENREI_COUNT_BITS>>])
end;
handle_packet(Transport, Socket, TreePid, <<?DENREI_PROTOCOL_REQUEST, Fields0/binary>>) ->
{Subject, Fields1} = denrei_utils:get_next_field(Fields0),
case denrei_tree_server:match(TreePid, denrei_utils:tokenize_subject(Subject)) of
{ok, []} ->
{CorrelationID, _} = denrei_utils:get_next_field(Fields1),
denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_NOSUBSCRIBER>>, [CorrelationID]);
{ok, Subscribers} ->
Index = random:uniform(length(Subscribers)),
{T, S} = lists:nth(Index, Subscribers),
EncodedCallerID = term_to_binary({Transport, Socket}),
denrei_utils:send(T, S, <<?DENREI_PROTOCOL_REQUEST>>, [EncodedCallerID], Fields0)
end;
handle_packet(_Transport, _Socket, _TreePid, <<?DENREI_PROTOCOL_RESPONSE, Fields0/binary>>) ->
{EncodedCallerID, Fields1} = denrei_utils:get_next_field(Fields0),
{T, S} = binary_to_term(EncodedCallerID),
denrei_utils:send(T, S, [<<?DENREI_PROTOCOL_RESPONSE>>, Fields1]);
handle_packet(_Transport, _Socket, _TreePid, <<?DENREI_PROTOCOL_NOCALLBACK, Fields0/binary>>) ->
{EncodedCallerID, Fields1} = denrei_utils:get_next_field(Fields0),
{T, S} = binary_to_term(EncodedCallerID),
denrei_utils:send(T, S, [<<?DENREI_PROTOCOL_NOCALLBACK>>, Fields1]);
handle_packet(_Transport, _Socket, _TreePid, <<?DENREI_PROTOCOL_TIMEOUT, Fields0/binary>>) ->
{EncodedCallerID, Fields1} = denrei_utils:get_next_field(Fields0),
{T, S} = binary_to_term(EncodedCallerID),
denrei_utils:send(T, S, [<<?DENREI_PROTOCOL_TIMEOUT>>, Fields1]);
handle_packet(Transport, Socket, _TreePid, <<?DENREI_PROTOCOL_PING, Fields0/binary>>) ->
{T1, Fields1} = denrei_utils:get_next_field(Fields0),
{T2, <<>>} = denrei_utils:get_next_field(Fields1),
ServerSendTime = denrei_utils:binary_to_integer(T1, 36),
ClientSendTime = denrei_utils:binary_to_integer(T2, 36),
ServerReceiveTime = denrei_utils:timestamp(),
lager:debug([{transport, Transport}, {socket, Socket}], "~p ping, total: ~p, server->client: ~p, client->server: ~p", [{Transport, Socket}, (ServerReceiveTime - ServerSendTime), (ClientSendTime - ServerSendTime), (ServerReceiveTime - ClientSendTime)]),
ok;
handle_packet(Transport, Socket, _TreePid, Packet) ->
lager:warning([{transport, Transport}, {socket, Socket}], "~p unknown packet: ~p", [{Transport, Socket}, Packet]).