Packages

Denrei - a lightweight Erlang messaging system.

Current section

Files

Jump to
denrei src denrei_client.erl
Raw

src/denrei_client.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% Denrei's client module.
%%% @end
%%%-------------------------------------------------------------------
-module(denrei_client).
-behaviour(gen_server).
%% Includes
-include("denrei.hrl").
%% API
-export([start/0, start/1, start_link/0, start_link/1, stop/1]).
-export([publish/3, subscribe/3, unsubscribe/3, subscribers/4, request/5]).
%% gen_server callbacks
-export([init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2,
code_change/3]).
%% internal state
-record(state, {
transport,
socket,
tree,
messages
}).
%%%===================================================================
%%% API
%%%===================================================================
-spec(publish(Pid :: pid(), Subject :: iodata(), Message :: iodata()) -> ok).
publish(Pid, Subject, Message) ->
gen_server:cast(Pid, {publish, Subject, Message}).
-spec(subscribe(Pid :: pid(), Subject :: iodata(), Callback :: pid() | fun()) -> ok | {error, atom()}).
subscribe(Pid, Subject, Callback) ->
gen_server:call(Pid, {subscribe, Subject, Callback}).
-spec(unsubscribe(Pid :: pid(), Subject :: iodata() | any, Callback :: pid() | fun() | all) -> ok | {error, atom()}).
unsubscribe(Pid, Subject, Callback) ->
gen_server:call(Pid, {unsubscribe, Subject, Callback}).
-spec(subscribers(Pid :: pid(), Subject :: iodata(), Callback :: pid() | fun(), Timeout :: timeout()) -> {ok, reference()} | {error, atom()}).
subscribers(Pid, Subject, Callback, Timeout) ->
gen_server:call(Pid, {subscribers, Subject, Callback, Timeout}).
-spec(request(Pid :: pid(), Subject :: iodata(), Message :: iodata(), Callback :: pid() | fun(), Timeout :: timeout()) -> {ok, reference()} | {error, atom()}).
request(Pid, Subject, Message, Callback, Timeout) ->
gen_server:call(Pid, {request, Subject, Message, Callback, Timeout}).
%%--------------------------------------------------------------------
%% @doc
%% Starts the server
%%
%% @end
%%--------------------------------------------------------------------
-spec(start() ->
{ok, Pid :: pid()} | ignore | {error, Reason :: term()}).
start() ->
start([]).
-spec(start(Options :: term()) ->
{ok, Pid :: pid()} | ignore | {error, Reason :: term()}).
start(Options) ->
gen_server:start(?MODULE, Options, []).
-spec(start_link() ->
{ok, Pid :: pid()} | ignore | {error, Reason :: term()}).
start_link() ->
start_link([]).
-spec(start_link(Options :: term()) ->
{ok, Pid :: pid()} | ignore | {error, Reason :: term()}).
start_link(Options) ->
gen_server:start_link(?MODULE, Options, []).
-spec(stop(Pid :: pid()) ->
ok).
stop(Pid) ->
gen_server:cast(Pid, stop).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%%--------------------------------------------------------------------
%% @private
%% @doc
%% Initializes the server
%%
%% @spec init(Args) -> {ok, State} |
%% {ok, State, Timeout} |
%% ignore |
%% {stop, Reason}
%% @end
%%--------------------------------------------------------------------
-spec(init(Args :: term()) ->
{ok, State :: #state{}} | {ok, State :: #state{}, timeout() | hibernate} |
{stop, Reason :: term()} | ignore).
init(Args) ->
Transport = denrei_utils:get_opt("TRANSPORT", transport, Args, ?DENREI_TRANSPORT),
Address = denrei_utils:to_ip(denrei_utils:get_opt("IP", ip, Args, ?DENREI_IP)),
Port = denrei_utils:to_integer(denrei_utils:get_opt("PORT", port, Args, ?DENREI_PORT)),
{ok, Socket} = Transport:connect(Address, Port, []),
ok = Transport:setopts(Socket, [binary, {active, once}, {packet, 4}]),
lager:info([{transport, Transport}, {socket, Socket}], "~p connected to ~p:~p", [{Transport, Socket}, Address, Port]),
{ok, #state{transport = Transport, socket = Socket, tree = denrei_tree:new(), messages = Transport:messages()}}.
%%--------------------------------------------------------------------
%% @private
%% @doc
%% Handling call messages
%%
%% @end
%%--------------------------------------------------------------------
-spec(handle_call(Request :: term(), From :: {pid(), Tag :: term()},
State :: #state{}) ->
{reply, Reply :: term(), NewState :: #state{}} |
{reply, Reply :: term(), NewState :: #state{}, timeout() | hibernate} |
{noreply, NewState :: #state{}} |
{noreply, NewState :: #state{}, timeout() | hibernate} |
{stop, Reason :: term(), Reply :: term(), NewState :: #state{}} |
{stop, Reason :: term(), NewState :: #state{}}).
handle_call({subscribe, Subject, Callback}, _From, State = #state{transport = Transport, socket = Socket, tree = Tree}) ->
Tokens = denrei_utils:tokenize_subject(Subject),
Reply = case denrei_tree:fetch(Tokens, Tree) of
false ->
denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_SUBSCRIBE>>, [iolist_to_binary(Subject)]);
_ ->
ok
end,
NewTree = denrei_tree:insert(Callback, Tokens, Tree),
{reply, Reply, State#state{tree = NewTree}};
handle_call({unsubscribe, Subject, Callback}, _From, State = #state{transport = Transport, socket = Socket, tree = Tree}) ->
Tokens = denrei_utils:tokenize_subject(Subject),
NewTree = denrei_tree:delete(Callback, Tokens, Tree),
Reply = case denrei_tree:fetch(Tokens, NewTree) of
false when Subject == all ->
denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_UNSUBSCRIBE>>, []);
false ->
denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_UNSUBSCRIBE>>, [iolist_to_binary(Subject)]);
_ ->
ok
end,
{reply, Reply, State#state{tree = NewTree}};
handle_call({subscribers, Subject, Callback, infinity}, From, State) ->
handle_call({subscribers, Subject, Callback, 0}, From, State);
handle_call({subscribers, Subject, Callback, Timeout}, _From, State = #state{transport = Transport, socket = Socket}) ->
Reference = make_ref(),
CorrelationId = term_to_binary({Callback, Reference}),
Reply = case denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_SUBSCRIBERS>>, [iolist_to_binary(Subject), CorrelationId, <<Timeout:?DENREI_TIMEOUT_BITS>>]) of
ok ->
{ok, Reference};
Error ->
Error
end,
{reply, Reply, State};
handle_call({request, Subject, Message, Callback, infinity}, From, State) ->
handle_call({request, Subject, Message, Callback, 0}, From, State);
handle_call({request, Subject, Message, Callback, Timeout}, _From, State = #state{transport = Transport, socket = Socket}) ->
Reference = make_ref(),
CorrelationId = term_to_binary({Callback, Reference}),
Reply = case denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_REQUEST>>, [iolist_to_binary(Subject), CorrelationId, <<Timeout:?DENREI_TIMEOUT_BITS>>, iolist_to_binary(Message)]) of
ok ->
{ok, Reference};
Error ->
Error
end,
{reply, Reply, State};
handle_call(_Request, _From, State) ->
{reply, ok, State}.
%%--------------------------------------------------------------------
%% @private
%% @doc
%% Handling cast messages
%%
%% @end
%%--------------------------------------------------------------------
-spec(handle_cast(Request :: term(), State :: #state{}) ->
{noreply, NewState :: #state{}} |
{noreply, NewState :: #state{}, timeout() | hibernate} |
{stop, Reason :: term(), NewState :: #state{}}).
handle_cast({publish, Subject, Message}, State = #state{transport = Transport, socket = Socket}) ->
denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_PUBLISH>>, [iolist_to_binary(Subject), iolist_to_binary(Message)]),
{noreply, State};
handle_cast(stop, State) ->
{stop, normal, State};
handle_cast(_Request, State) ->
{noreply, State}.
%%--------------------------------------------------------------------
%% @private
%% @doc
%% Handling all non call/cast messages
%%
%% @spec handle_info(Info, State) -> {noreply, State} |
%% {noreply, State, Timeout} |
%% {stop, Reason, State}
%% @end
%%--------------------------------------------------------------------
-spec(handle_info(Info :: timeout() | term(), State :: #state{}) ->
{noreply, NewState :: #state{}} |
{noreply, NewState :: #state{}, timeout() | hibernate} |
{stop, Reason :: term(), NewState :: #state{}}).
handle_info({OK, Socket, Packet}, State = #state{transport = Transport, messages = {OK, _, _}}) ->
Transport:setopts(Socket, [{active, once}]),
lager:debug([{transport, Transport}, {socket, Socket}], "~p received: ~p", [{Transport, Socket}, Packet]),
{noreply, handle_packet(Transport, Socket, State, Packet)};
handle_info({Closed, Socket}, State = #state{transport = Transport, messages = {_, Closed, _}}) ->
lager:info([{transport, Transport}, {socket, Socket}], "~p closed", [{Transport, Socket}]),
{stop, Closed, State};
handle_info({Error, Socket, Reason}, State = #state{transport = Transport, messages = {_, _, Error}}) ->
Transport:setopts(Socket, [{active, once}]),
lager:warning([{transport, Transport}, {socket, Socket}], "~p error: ~p", [{Transport, Socket}, Reason]),
{stop, {Error, Reason, State}};
handle_info(_Info, State) ->
{noreply, State}.
%%--------------------------------------------------------------------
%% @private
%% @doc
%% This function is called by a gen_server when it is about to
%% terminate. It should be the opposite of Module:init/1 and do any
%% necessary cleaning up. When it returns, the gen_server terminates
%% with Reason. The return value is ignored.
%%
%% @spec terminate(Reason, State) -> void()
%% @end
%%--------------------------------------------------------------------
-spec(terminate(Reason :: (normal | shutdown | {shutdown, term()} | term()),
State :: #state{}) -> term()).
terminate(Reason, #state{transport = Transport, socket = Socket}) ->
lager:debug([{transport, Transport}, {socket, Socket}], "~p terminate: ~p", [{Transport, Socket}, Reason]),
Transport:close(Socket),
ok.
%%--------------------------------------------------------------------
%% @private
%% @doc
%% Convert process state when code is changed
%%
%% @spec code_change(OldVsn, State, Extra) -> {ok, NewState}
%% @end
%%--------------------------------------------------------------------
-spec(code_change(OldVsn :: term() | {down, term()}, State :: #state{},
Extra :: term()) ->
{ok, NewState :: #state{}} | {error, Reason :: term()}).
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
%%%===================================================================
%%% Internal functions
%%%===================================================================
handle_packet(_Transport, _Socket, State = #state{tree = Tree}, <<?DENREI_PROTOCOL_PUBLISH, Fields0/binary>>) ->
{Subject, Fields1} = denrei_utils:get_next_field(Fields0),
{Message, <<>>} = denrei_utils:get_next_field(Fields1),
[do_callback(Callback, {?DENREI_MESSAGE, Message}) || Callback <- denrei_tree:match(denrei_utils:tokenize_subject(Subject), Tree)],
State;
handle_packet(Transport, Socket, State = #state{tree = Tree}, <<?DENREI_PROTOCOL_REQUEST, Fields0/binary>>) ->
{CallerID, Fields1} = denrei_utils:get_next_field(Fields0),
{Subject, Fields2} = denrei_utils:get_next_field(Fields1),
{CorrelationID, Fields3} = denrei_utils:get_next_field(Fields2),
{<<Timeout:?DENREI_TIMEOUT_BITS>>, Fields4} = denrei_utils:get_next_field(Fields3),
{Message, <<>>} = denrei_utils:get_next_field(Fields4),
spawn(fun() ->
case denrei_tree:match(denrei_utils:tokenize_subject(Subject), Tree) of
[] ->
denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_NOCALLBACK>>, [CallerID, CorrelationID]);
Callbacks ->
Index = random:uniform(length(Callbacks)),
Callback = lists:nth(Index, Callbacks),
case do_request(Callback, Message, Timeout) of
{error, timeout} ->
denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_TIMEOUT>>, [CallerID, CorrelationID]);
Response ->
denrei_utils:send(Transport, Socket, <<?DENREI_PROTOCOL_RESPONSE>>, [CallerID, CorrelationID, Response])
end
end
end),
State;
handle_packet(_Transport, _Socket, State, <<?DENREI_PROTOCOL_RESPONSE, Fields0/binary>>) ->
{EncodedCorrelationID, Fields1} = denrei_utils:get_next_field(Fields0),
{Response, <<>>} = denrei_utils:get_next_field(Fields1),
{Callback, Reference} = binary_to_term(EncodedCorrelationID),
do_callback(Callback, {?DENREI_RESPONSE, Reference, Response}),
State;
handle_packet(_Transport, _Socket, State, <<?DENREI_PROTOCOL_NOSUBSCRIBER, Fields/binary>>) ->
{EncodedCorrelationID, <<>>} = denrei_utils:get_next_field(Fields),
{Callback, Reference} = binary_to_term(EncodedCorrelationID),
do_callback(Callback, {?DENREI_RESPONSE, Reference, {error, no_subscriber}}),
State;
handle_packet(_Transport, _Socket, State, <<?DENREI_PROTOCOL_NOCALLBACK, Fields/binary>>) ->
{EncodedCorrelationID, <<>>} = denrei_utils:get_next_field(Fields),
{Callback, Reference} = binary_to_term(EncodedCorrelationID),
do_callback(Callback, {?DENREI_RESPONSE, Reference, {error, no_callback}}),
State;
handle_packet(_Transport, _Socket, State, <<?DENREI_PROTOCOL_TIMEOUT, Fields/binary>>) ->
{EncodedCorrelationID, <<>>} = denrei_utils:get_next_field(Fields),
{Callback, Reference} = binary_to_term(EncodedCorrelationID),
do_callback(Callback, {?DENREI_RESPONSE, Reference, {error, timeout}}),
State;
handle_packet(Transport, Socket, State, Packet = <<?DENREI_PROTOCOL_PING, _/binary>>) ->
denrei_utils:send(Transport, Socket, Packet, [denrei_utils:integer_to_binary(denrei_utils:timestamp(), 36)]),
State;
handle_packet(Transport, Socket, State, Packet) ->
lager:warning([{transport, Transport}, {socket, Socket}], "~p unknown packet: ~p", [{Transport, Socket}, Packet]),
State.
do_callback(Callback, Message) when is_pid(Callback) ->
Callback ! Message;
do_callback(Callback, Message) ->
spawn(fun() -> Callback(Message) end).
do_request(Callback, Message, 0) ->
do_request(Callback, Message, infinity);
do_request(Callback, Message, Timeout) when is_pid(Callback) ->
Callback ! {?DENREI_REQUEST, self(), Message},
receive
Response -> Response
after
Timeout -> {error, timeout}
end;
do_request(Callback, Message, Timeout) ->
From = self(),
spawn(fun() ->
From ! Callback(Message)
end),
receive
Response -> Response
after
Timeout -> {error, timeout}
end.