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_stat.erl
%%%-------------------------------------------------------------------
%%% @author cheese
%%% @copyright (C) 2016, <COMPANY>
%%% @doc
%%%
%%% @end
%%% Created : 22. Apr 2016 14:15
%%%-------------------------------------------------------------------
-module(rabbit_rpc2_stat).
-author("cheese").
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
%% API
-export([start_link/0, message_in/1, message_out/1, message_timeout/1, get_state/0, message_queue_response/1]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3, add_counter/4,
add_value/2, convert_record_list/2, update_speed/2, fold_response_queue/0, fold_queue/2]).
-define(SERVER, ?MODULE).
-define(SpeedCalcTimeout, 300).
-record(queue_info, {
name,
count_in,
count_out,
count_timeout,
count_in_prev,
speed_in,
last_date
}
).
-record(state, {
queues_request :: [#queue_info{}],
queues_response :: [#queue_info{}]
}).
%%%===================================================================
%%% API
%%%===================================================================
start_link() ->
gen_server:start_link({local, ?SERVER}, ?MODULE, [], []).
message_in(QueueName) ->
gen_server:cast(?SERVER, {add_value, in, QueueName}).
message_out(QueueName) ->
gen_server:cast(?SERVER, {add_value, out, QueueName}).
fold_response_queue() ->
gen_server:cast(?SERVER, fold_response_queue).
message_timeout(QueueName) ->
gen_server:cast(?SERVER, {add_value, timeout, QueueName}).
message_queue_response(QueueName) ->
gen_server:cast(?SERVER, {add_value_resp, QueueName}).
get_state() ->
gen_server:call(?SERVER, get_state).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init([]) ->
?LOG_INFO("Start stat_rmq~n",[]),
erlang:send_after(?SpeedCalcTimeout*1000, self(), calc_speed),
{ok, #state{queues_request = [], queues_response = []}}.
handle_call(get_state, _From, State) ->
{reply, {ok, convert_record_list(State#state.queues_request, []), convert_record_list(State#state.queues_response, [])}, State};
handle_call(_Request, _From, State) ->
{reply, ok, State}.
handle_cast(fold_response_queue, #state{queues_response = Arr} = State) ->
?LOG_INFO("fold_response_queue~n",[]),
Queue = #queue_info{count_in = 0, count_in_prev = 0, count_out = 0, count_timeout = 0, name = <<"rmq_responses_other">>, last_date = calendar:local_time(), speed_in = 0},
NewStat = State#state{queues_response = fold_queue(Arr, Queue)},
{noreply, NewStat};
handle_cast({add_value_resp, QueueName}, #state{queues_response = Arr} = State) ->
?LOG_INFO("add_value Type: ~w, QueueName: ~s~n",[in, QueueName]),
NewStat = State#state{queues_response = add_counter(in, QueueName, Arr, [])},
{noreply, NewStat};
handle_cast({add_value, Type, QueueName}, #state{queues_request = Arr} = State) ->
?LOG_INFO("add_value Type: ~w, QueueName: ~s~n",[Type, QueueName]),
NewStat = State#state{queues_request = add_counter(Type, QueueName, Arr, [])},
{noreply, NewStat};
handle_cast(_Request, State) ->
{noreply, State}.
handle_info(calc_speed, State = #state{queues_response = ResponseQueues, queues_request = RequestQueues}) ->
erlang:send_after(?SpeedCalcTimeout*1000, self(), calc_speed),
{noreply, State#state{queues_request = update_speed(RequestQueues, []), queues_response = update_speed(ResponseQueues, [])}};
handle_info(Info, State) ->
?LOG_INFO("Stat unknown message: ~w. State: ~w~n",[Info, State]),
{noreply, State}.
terminate(_Reason, _State) ->
ok.
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
%%%===================================================================
%%% Internal functions
%%%===================================================================
add_counter(Type, QueueName, [], Res) ->
QueueInfo = case Type of
in ->
#queue_info{count_in = 1, name = QueueName, count_timeout = 0, count_out = 0, count_in_prev = 0, speed_in = 0, last_date = calendar:local_time()};
out ->
#queue_info{count_out = 1, name = QueueName, count_timeout = 0, count_in = 0, count_in_prev = 0, speed_in = 0};
timeout ->
#queue_info{count_timeout = 1, name = QueueName, count_in = 0, count_out = 0, count_in_prev = 0, speed_in = 0}
end,
[QueueInfo | Res];
add_counter(Type, QueueName, [#queue_info{name = QueueName} = QueueInfo | Tail], Res) ->
Res ++ [add_value(Type, QueueInfo) | Tail];
add_counter(Type, QueueName, [Head | Tail], Res) ->
add_counter(Type, QueueName, Tail, [Head | Res]).
add_value(Type, QueueInfo = #queue_info{count_in = CountIn, count_out = CountOut, count_timeout = CountTimeout}) ->
case Type of
in ->
QueueInfo#queue_info{count_in = CountIn + 1, last_date = calendar:local_time()};
out ->
QueueInfo#queue_info{count_out = CountOut + 1};
timeout ->
QueueInfo#queue_info{count_timeout = CountTimeout + 1}
end.
convert_record_list([], Res) ->
Res;
convert_record_list([#queue_info{count_in = CountIn, count_out = CountOut, count_timeout = CountTimeout, name = QueueName, speed_in = Speed, last_date = LastDate} | Tail], Res) ->
convert_record_list(Tail, [{QueueName, CountIn, CountOut, CountTimeout, Speed, LastDate} | Res]).
update_speed([], Res) ->
Res;
update_speed([#queue_info{count_in = CountIn, count_in_prev = CountInPrev} = QueueInfo | Tail], Res) ->
Speed = (CountIn-CountInPrev)/?SpeedCalcTimeout,
update_speed(Tail, [QueueInfo#queue_info{speed_in = Speed, count_in_prev = CountIn} | Res]).
fold_queue(undefined, QueueResult) ->
?LOG_INFO("fold_queue undefined~n",[]),
[QueueResult];
fold_queue([], QueueResult) ->
?LOG_INFO("fold_queue []~n",[]),
[QueueResult];
fold_queue([HeadQ | TailQ], #queue_info{count_in = CountIn, count_out = CountOut, count_timeout = CountTimeout} = QueueResult) ->
#queue_info{count_in = HCountIn, count_out = HCountOut, count_timeout = HCountTimeout} = HeadQ,
NewQueueResult = QueueResult#queue_info{count_in = CountIn+HCountIn, count_out = CountOut+HCountOut, count_timeout = CountTimeout+HCountTimeout},
?LOG_INFO("==============================~n",[QueueResult]),
?LOG_INFO("fold_queue QueueResult: ~w~n",[QueueResult]),
?LOG_INFO("fold_queue NewQueueResult: ~w~n",[NewQueueResult]),
fold_queue(TailQ, NewQueueResult).