Packages
diint_utilites_common_app
1.4.21
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.12
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.9
1.3.8
1.3.7
1.3.6
1.3.5
1.3.4
1.3.3
1.3.2
1.3.1
1.3.0
1.2.101
1.2.11
1.2.10
1.2.9
1.2.8
1.2.7
1.2.6
1.2.5
1.2.4
1.2.3
1.2.2
1.2.1
1.2.0
1.1.7
1.1.6
1.1.5
1.1.4
1.1.3
1.1.2
1.1.1
1.1.0
1.0.13
1.0.12
1.0.11
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
библиотеки для работы с ребитом и протоколы
Current section
Files
Jump to
Current section
Files
src/mq/rabbit_rpc2.erl
%%%-------------------------------------------------------------------
%%% @author cheese
%%% @copyright (C) 2016, <COMPANY>
%%% @doc
%%%
%%% @end
%%% Created : 18. Apr 2016 15:47
%%%-------------------------------------------------------------------
-module(rabbit_rpc2).
-author("cheese").
-behaviour(supervisor).
-include("include/types_rabbit.hrl").
-include_lib("kernel/include/logger.hrl").
%% API
-export([start_link/1, send_message_and_receive_response/5]).
%% Internal
-export([set_state_rmq_sender/3, check_rabbit_sender/3, create_new_sender/3, get_current_supervisor_name/1, on_rabbit_response/4]).
%% Supervisor callbacks
-export([init/1, get_kv_rabbit_sender_process_name/1, get_kv_message_storage_process_name/1]).
-define(CHILD(I, Type), {I, {I, start_link, []}, permanent, 5000, Type, [I]}).
-define(CHILD_SV(ProcName, ModuleName, Arguments), {ProcName, {ModuleName, start_link, Arguments}, permanent, 2000, supervisor, []}).
-define(CHILD_WRK(ProcName, ModuleName, Arguments), {ProcName, {ModuleName, start_link, Arguments}, permanent, 2000, worker, []}).
%%%===================================================================
%%% API functions
%%%===================================================================
start_link(RabbitConfig) ->
Process = get_current_supervisor_name(RabbitConfig),
supervisor:start_link({local, Process}, ?MODULE, [RabbitConfig]).
%%%===================================================================
%%% Supervisor callbacks
%%%===================================================================
init([RmqConnectionConfig]) ->
KVRmqProcName = get_kv_rabbit_sender_process_name(RmqConnectionConfig),
KVRmqSuperVisor = ?CHILD_WRK(KVRmqProcName, kv_storage, [25920000, KVRmqProcName]),
KVMsgProcName = get_kv_message_storage_process_name(RmqConnectionConfig),
KVMsgSuperVisor = ?CHILD_WRK(KVMsgProcName, kv_storage, [120000, KVMsgProcName]),
RPC2Reader = ?CHILD_WRK(rabbit_rpc2_reader, rabbit_rpc2_reader, [RmqConnectionConfig]),
RabbitStat = ?CHILD_WRK(rabbit_rpc2_stat, rabbit_rpc2_stat, []),
{ok, {{one_for_one, 5, 10}, [RabbitStat, KVRmqSuperVisor, KVMsgSuperVisor, RPC2Reader]}}.
%%%===================================================================
%%% Internal functions
%%%===================================================================
get_kv_rabbit_sender_process_name(#rabbit_connection_config{process_name = ProcessName}) ->
list_to_atom(atom_to_list(ProcessName) ++ "_rpc2_kv_rmq").
get_kv_message_storage_process_name(#rabbit_connection_config{process_name = ProcessName}) ->
list_to_atom(atom_to_list(ProcessName) ++ "_rpc2_kv_msg").
send_message_and_receive_response(RabbitConfig, ExchangeName, QueueName, Bytes, TimeoutMiliseconds) ->
case check_rabbit_sender(RabbitConfig, ExchangeName, QueueName) of
{ok, _State} ->
MessageId = bytes_extension:generate_uuid(),
?LOG_INFO("Create MessageId: ~s~n",[MessageId]),
KVMessage = get_kv_message_storage_process_name(RabbitConfig),
kv_storage:add(MessageId, self(), KVMessage),
RabbitResponseQueue = rabbit_rpc2_reader:get_queue_name(RabbitConfig),
?LOG_INFO("RabbitResponseQueue: ~s~n",[RabbitResponseQueue]),
rabbit_rpc2_stat:message_in(QueueName),
rabbit_rpc2_sender:send(RabbitConfig, Bytes, MessageId, QueueName, RabbitResponseQueue),
receive
Result ->
rabbit_rpc2_stat:message_out(QueueName),
{response, Result}
after
TimeoutMiliseconds + 500 ->
rabbit_rpc2_stat:message_timeout(QueueName),
{error, wait_response_timeout}
end;
{ok, started_new_reader} ->
send_message_and_receive_response(RabbitConfig, ExchangeName, QueueName, Bytes, TimeoutMiliseconds);
Other ->
?LOG_INFO("Error send message: ~w~n",[Other]),
Other
end.
on_rabbit_response(RabbitConfig, MessageId, Bytes, QueueName) ->
KVMessage = get_kv_message_storage_process_name(RabbitConfig),
rabbit_rpc2_stat:message_queue_response(QueueName),
case kv_storage:get(MessageId, true, KVMessage) of
{ok, Pid} ->
Pid ! Bytes,
{ok, sended};
{error, Reason} ->
?LOG_INFO("Unsolicited message ~s~n",[MessageId]),
{error, Reason}
end.
check_rabbit_sender(RabbitConfig, ExchangeName, QueueName) ->
KVProcess = get_kv_rabbit_sender_process_name(RabbitConfig),
RabbitSenderProcess = rabbit_rpc2_sender:get_process_name(RabbitConfig, QueueName),
case kv_storage:get(RabbitSenderProcess, false, KVProcess) of
{ok, connected} ->
{ok, connected};
{ok, initialization} ->
?LOG_INFO("Sender [~s] in state initialization~n", [RabbitSenderProcess]),
{error, rabbit_state_initialization};
{error, Reason} ->
?LOG_INFO("No sender [~s] found. Reason: ~w. Try to create~n", [RabbitSenderProcess, Reason]),
create_new_sender(RabbitConfig, ExchangeName, QueueName)
end.
create_new_sender(RabbitConfig, ExchangeName, QueueName) ->
RabbitSenderProcess = rabbit_rpc2_sender:get_process_name(RabbitConfig, QueueName),
NewSenderSpec = ?CHILD_WRK(RabbitSenderProcess, rabbit_rpc2_sender, [RabbitConfig, ExchangeName, QueueName]),
supervisor:start_child(get_current_supervisor_name(RabbitConfig), NewSenderSpec),
timer:sleep(1000),
{ok, started_new_reader}.
get_current_supervisor_name(RabbitConfig) ->
list_to_atom(atom_to_list(RabbitConfig#rabbit_connection_config.process_name) ++ "_rpc2_super").
-spec(set_state_rmq_sender(#rabbit_connection_config{}, term(), initialization | connected) -> ok | {error, term()}).
set_state_rmq_sender(RabbitConnectionConfig, RabbitSenderProcessName, State) ->
KVProcess = get_kv_rabbit_sender_process_name(RabbitConnectionConfig),
kv_storage:get(RabbitSenderProcessName, true, KVProcess),
kv_storage:add(RabbitSenderProcessName,State, KVProcess),
ok.