summaryrefslogtreecommitdiff
path: root/deps/rabbitmq_management_agent
diff options
context:
space:
mode:
authorDiana Corbacho <diana@rabbitmq.com>2019-04-08 20:50:29 +0100
committerDiana Corbacho <diana@rabbitmq.com>2019-04-08 20:50:29 +0100
commitff266f0c8c0807d5c8530fb497307339e9044569 (patch)
tree1a097ad99e5241b93b6ddf9ae0248f82a6bc905a /deps/rabbitmq_management_agent
parent57133f66397d1558e545499351ba954253316d7c (diff)
downloadrabbitmq-server-git-ff266f0c8c0807d5c8530fb497307339e9044569.tar.gz
Clean up of non-local queue stats
Followers/slaves should not hold stats for any non-local queue. Ensure clean up happens if any has been left behind [#165153327]
Diffstat (limited to 'deps/rabbitmq_management_agent')
-rw-r--r--deps/rabbitmq_management_agent/src/rabbit_mgmt_gc.erl10
-rw-r--r--deps/rabbitmq_management_agent/test/rabbit_mgmt_gc_SUITE.erl83
2 files changed, 89 insertions, 4 deletions
diff --git a/deps/rabbitmq_management_agent/src/rabbit_mgmt_gc.erl b/deps/rabbitmq_management_agent/src/rabbit_mgmt_gc.erl
index c65f95866e..c6d5748238 100644
--- a/deps/rabbitmq_management_agent/src/rabbit_mgmt_gc.erl
+++ b/deps/rabbitmq_management_agent/src/rabbit_mgmt_gc.erl
@@ -85,11 +85,13 @@ gc_channels() ->
gc_queues() ->
Queues = rabbit_amqqueue:list_names(),
GbSet = gb_sets:from_list(Queues),
+ LocalQueues = rabbit_amqqueue:list_local_names(),
+ LocalGbSet = gb_sets:from_list(LocalQueues),
gc_entity(queue_stats_publish, GbSet),
- gc_entity(queue_stats, GbSet),
- gc_entity(queue_msg_stats, GbSet),
- gc_entity(queue_process_stats, GbSet),
- gc_entity(queue_msg_rates, GbSet),
+ gc_entity(queue_stats, LocalGbSet),
+ gc_entity(queue_msg_stats, LocalGbSet),
+ gc_entity(queue_process_stats, LocalGbSet),
+ gc_entity(queue_msg_rates, LocalGbSet),
gc_entity(queue_stats_deliver_stats, GbSet),
gc_process_and_entity(channel_queue_stats_deliver_stats_queue_index, GbSet),
gc_process_and_entity(consumer_stats_queue_index, GbSet),
diff --git a/deps/rabbitmq_management_agent/test/rabbit_mgmt_gc_SUITE.erl b/deps/rabbitmq_management_agent/test/rabbit_mgmt_gc_SUITE.erl
index d5507e5566..cbd385f0d2 100644
--- a/deps/rabbitmq_management_agent/test/rabbit_mgmt_gc_SUITE.erl
+++ b/deps/rabbitmq_management_agent/test/rabbit_mgmt_gc_SUITE.erl
@@ -32,6 +32,7 @@ groups() ->
[
{non_parallel_tests, [],
[ queue_stats,
+ quorum_queue_stats,
connection_stats,
channel_stats,
vhost_stats,
@@ -184,6 +185,88 @@ queue_stats(Config) ->
ok.
+quorum_queue_stats(Config) ->
+ A = rabbit_ct_broker_helpers:get_node_config(Config, 0, nodename),
+ B = rabbit_ct_broker_helpers:get_node_config(Config, 1, nodename),
+ Ch = rabbit_ct_client_helpers:open_channel(Config, A),
+
+ amqp_channel:call(Ch,
+ #'queue.declare'{queue = <<"quorum_queue_stats">>,
+ durable = true,
+ arguments = [{<<"x-queue-type">>, longstr, <<"quorum">>}]}),
+ timer:sleep(1150),
+
+ Q = q(<<"quorum_queue_stats">>),
+
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [queue_stats, {Q, infos}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [queue_msg_stats, {{Q, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [queue_process_stats, {{Q, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [queue_msg_rates, {{Q, 5}, slide}]),
+
+ rabbit_ct_broker_helpers:rpc(Config, B, ets, insert,
+ [queue_stats, {Q, infos}]),
+ rabbit_ct_broker_helpers:rpc(Config, B, ets, insert,
+ [queue_msg_stats, {{Q, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, B, ets, insert,
+ [queue_process_stats, {{Q, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, B, ets, insert,
+ [queue_msg_rates, {{Q, 5}, slide}]),
+
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [queue_stats, Q]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [queue_msg_stats, {Q, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [queue_process_stats, {Q, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [queue_msg_rates, {Q, 5}]),
+
+ [_] = rabbit_ct_broker_helpers:rpc(Config, B, ets, lookup,
+ [queue_stats, Q]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, B, ets, lookup,
+ [queue_msg_stats, {Q, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, B, ets, lookup,
+ [queue_process_stats, {Q, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, B, ets, lookup,
+ [queue_msg_rates, {Q, 5}]),
+
+ {ok, _, {_, Leader}} = rabbit_ct_broker_helpers:rpc(Config, B, ra, members,
+ [{'%2F_quorum_queue_stats', B}]),
+ [Follower] = [A, B] -- [Leader],
+
+ %% Trigger gc. When the gen_server:call returns, the gc has already finished.
+ rabbit_ct_broker_helpers:rpc(Config, Leader, erlang, send, [rabbit_mgmt_gc, start_gc]),
+ rabbit_ct_broker_helpers:rpc(Config, Leader, gen_server, call, [rabbit_mgmt_gc, test]),
+ rabbit_ct_broker_helpers:rpc(Config, Follower, erlang, send, [rabbit_mgmt_gc, start_gc]),
+ rabbit_ct_broker_helpers:rpc(Config, Follower, gen_server, call, [rabbit_mgmt_gc, test]),
+
+ [] = rabbit_ct_broker_helpers:rpc(Config, Follower, ets, lookup,
+ [queue_stats, Q]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, Follower, ets, lookup,
+ [queue_msg_stats, {Q, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, Follower, ets, lookup,
+ [queue_process_stats, {Q, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, Follower, ets, lookup,
+ [queue_msg_rates, {Q, 5}]),
+
+ [_] = rabbit_ct_broker_helpers:rpc(Config, Leader, ets, lookup,
+ [queue_stats, Q]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, Leader, ets, lookup,
+ [queue_msg_stats, {Q, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, Leader, ets, lookup,
+ [queue_process_stats, {Q, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, Leader, ets, lookup,
+ [queue_msg_rates, {Q, 5}]),
+
+ amqp_channel:call(Ch, #'queue.delete'{queue = <<"quorum_queue_stats">>}),
+ rabbit_ct_client_helpers:close_channel(Ch),
+
+ ok.
+
connection_stats(Config) ->
A = rabbit_ct_broker_helpers:get_node_config(Config, 0, nodename),
Ch = rabbit_ct_client_helpers:open_channel(Config, A),