diff options
| author | Diana Corbacho <diana@rabbitmq.com> | 2019-04-08 20:50:29 +0100 |
|---|---|---|
| committer | Diana Corbacho <diana@rabbitmq.com> | 2019-04-08 20:50:29 +0100 |
| commit | ff266f0c8c0807d5c8530fb497307339e9044569 (patch) | |
| tree | 1a097ad99e5241b93b6ddf9ae0248f82a6bc905a /deps/rabbitmq_management_agent | |
| parent | 57133f66397d1558e545499351ba954253316d7c (diff) | |
| download | rabbitmq-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.erl | 10 | ||||
| -rw-r--r-- | deps/rabbitmq_management_agent/test/rabbit_mgmt_gc_SUITE.erl | 83 |
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), |
