summaryrefslogtreecommitdiff
path: root/deps/rabbitmq_management_agent/test
diff options
context:
space:
mode:
authordcorbacho <dparracorbacho@piotal.io>2020-11-18 14:27:41 +0000
committerdcorbacho <dparracorbacho@piotal.io>2020-11-18 14:27:41 +0000
commitf23a51261d9502ec39df0f8db47ba6b22aa7659f (patch)
tree53dcdf46e7dc2c14e81ee960bce8793879b488d3 /deps/rabbitmq_management_agent/test
parentafa2c2bf6c7e0e9b63f4fb53dc931c70388e1c82 (diff)
parent9f6d64ec4a4b1eeac24d7846c5c64fd96798d892 (diff)
downloadrabbitmq-server-git-stream-timestamp-offset.tar.gz
Merge remote-tracking branch 'origin/master' into stream-timestamp-offsetstream-timestamp-offset
Diffstat (limited to 'deps/rabbitmq_management_agent/test')
-rw-r--r--deps/rabbitmq_management_agent/test/exometer_slide_SUITE.erl631
-rw-r--r--deps/rabbitmq_management_agent/test/metrics_SUITE.erl157
-rw-r--r--deps/rabbitmq_management_agent/test/rabbit_mgmt_gc_SUITE.erl626
-rw-r--r--deps/rabbitmq_management_agent/test/rabbit_mgmt_slide_SUITE.erl174
4 files changed, 1588 insertions, 0 deletions
diff --git a/deps/rabbitmq_management_agent/test/exometer_slide_SUITE.erl b/deps/rabbitmq_management_agent/test/exometer_slide_SUITE.erl
new file mode 100644
index 0000000000..abdf24853d
--- /dev/null
+++ b/deps/rabbitmq_management_agent/test/exometer_slide_SUITE.erl
@@ -0,0 +1,631 @@
+%% This Source Code Form is subject to the terms of the Mozilla Public
+%% License, v. 2.0. If a copy of the MPL was not distributed with this
+%% file, You can obtain one at https://mozilla.org/MPL/2.0/.
+%%
+%% Copyright (c) 2016-2020 VMware, Inc. or its affiliates. All rights reserved.
+%%
+
+-module(exometer_slide_SUITE).
+
+-include_lib("proper/include/proper.hrl").
+-include_lib("eunit/include/eunit.hrl").
+
+-compile(export_all).
+
+all() ->
+ [
+ {group, parallel_tests}
+ ].
+
+groups() ->
+ [
+ {parallel_tests, [parallel], [
+ incremental_add_element_basics,
+ last_two_normalises_old_sample_timestamp,
+ incremental_last_two_returns_last_two_completed_samples,
+ incremental_sum,
+ incremental_sum_stale,
+ incremental_sum_stale2,
+ incremental_sum_with_drop,
+ incremental_sum_with_total,
+ foldl_realises_partial_sample,
+ foldl_and_to_list,
+ foldl_and_to_list_incremental,
+ optimize,
+ stale_to_list,
+ to_list_single_after_drop,
+ to_list_drop_and_roll,
+ to_list_with_drop,
+ to_list_simple,
+ foldl_with_drop,
+ sum_single,
+ to_normalized_list,
+ to_normalized_list_no_padding,
+ to_list_in_the_past,
+ sum_mgmt_352,
+ sum_mgmt_352_extra,
+ sum_mgmt_352_peak
+ ]}
+ ].
+
+%% -------------------------------------------------------------------
+%% Testsuite setup/teardown.
+%% -------------------------------------------------------------------
+init_per_suite(Config) ->
+ Config.
+
+end_per_suite(_Config) ->
+ ok.
+
+init_per_group(_, Config) ->
+ Config.
+
+end_per_group(_, _Config) ->
+ ok.
+
+init_per_testcase(_, Config) ->
+ Config.
+
+end_per_testcase(_, _Config) ->
+ ok.
+
+%% -------------------------------------------------------------------
+%% Generators.
+%% -------------------------------------------------------------------
+elements_gen() ->
+ ?LET(Length, oneof([1, 2, 3, 7, 8, 20]),
+ ?LET(Elements, list(vector(Length, int())),
+ [erlang:list_to_tuple(E) || E <- Elements])).
+
+%% -------------------------------------------------------------------
+%% Testcases.
+%% -------------------------------------------------------------------
+
+%% TODO: turn tests into properties
+
+incremental_add_element_basics(_Config) ->
+ Now = exometer_slide:timestamp(),
+ S0 = exometer_slide:new(Now, 10, [{incremental, true},
+ {interval, 100}]),
+
+ [] = exometer_slide:to_list(Now, S0),
+ % add element before next interval
+ S1 = exometer_slide:add_element(Now + 10, {1}, S0),
+
+ [] = exometer_slide:to_list(Now + 20, S1),
+
+ %% to_list is not empty, as we take the 'real' total if a full interval has passed
+ Now100 = Now + 100,
+ [{Now100, {1}}] = exometer_slide:to_list(Now100, S1),
+
+ Then = Now + 101,
+ % add element after interval
+ S2 = exometer_slide:add_element(Then, {1}, S1),
+
+ % contains single element with incremented value
+ [{Then, {2}}] = exometer_slide:to_list(Then, S2).
+
+last_two_normalises_old_sample_timestamp(_Config) ->
+ Now = 0,
+ S0 = exometer_slide:new(Now, 10, [{incremental, true},
+ {interval, 100}]),
+
+ S1 = exometer_slide:add_element(100, {1}, S0),
+ S2 = exometer_slide:add_element(500, {1}, S1),
+
+ [{500, {2}}, {400, {1}}] = exometer_slide:last_two(S2).
+
+
+incremental_last_two_returns_last_two_completed_samples(_Config) ->
+ Now = exometer_slide:timestamp(),
+ S0 = exometer_slide:new(Now, 10, [{incremental, true},
+ {interval, 100}]),
+
+ % add two full elements then a partial
+ Now100 = Now + 100,
+ Now200 = Now + 200,
+ S1 = exometer_slide:add_element(Now100, {1}, S0),
+ S2 = exometer_slide:add_element(Now200, {1}, S1),
+ S3 = exometer_slide:add_element(Now + 210, {1}, S2),
+
+ [{Now200, {2}}, {Now100, {1}}] = exometer_slide:last_two(S3).
+
+incremental_sum(_Config) ->
+ Now = exometer_slide:timestamp(),
+ S1 = lists:foldl(fun (Next, S) ->
+ exometer_slide:add_element(Now + Next, {1}, S)
+ end,
+ exometer_slide:new(Now, 1000, [{incremental, true}, {interval, 100}]),
+ lists:seq(100, 1000, 100)),
+ Now50 = Now - 50,
+ S2 = lists:foldl(fun (Next, S) ->
+ exometer_slide:add_element(Now50 + Next, {1}, S)
+ end,
+ exometer_slide:new(Now50, 1000, [{incremental, true}, {interval, 100}]),
+ lists:seq(100, 1000, 100)),
+ S3 = exometer_slide:sum([S1, S2]),
+
+ 10 = length(exometer_slide:to_list(Now + 1000, S1)),
+ 10 = length(exometer_slide:to_list(Now + 1000, S2)),
+ 10 = length(exometer_slide:to_list(Now + 1000, S3)).
+
+incremental_sum_stale(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 25, [{incremental, true}, {interval, 5}]),
+
+ S1 = lists:foldl(fun (Next, S) ->
+ exometer_slide:add_element(Now + Next, {1}, S)
+ end, Slide, [1, 8, 15, 21, 27]),
+
+ S2 = lists:foldl(fun (Next, S) ->
+ exometer_slide:add_element(Now + Next, {1}, S)
+ end, Slide, [2, 7, 14, 20, 25]),
+ S3 = exometer_slide:sum([S1, S2]),
+ [27,22,17,12,7] = lists:reverse([T || {T, _} <- exometer_slide:to_list(27, S3)]),
+ [10,8,6,4,2] = lists:reverse([V || {_, {V}} <- exometer_slide:to_list(27, S3)]).
+
+incremental_sum_stale2(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 25, [{incremental, true},
+ {max_n, 5},
+ {interval, 5}]),
+
+ S1 = lists:foldl(fun (Next, S) ->
+ exometer_slide:add_element(Now + Next, {1}, S)
+ end, Slide, [5]),
+
+ S2 = lists:foldl(fun (Next, S) ->
+ exometer_slide:add_element(Now + Next, {1}, S)
+ end, Slide, [500, 505, 510, 515, 520, 525, 527]),
+ S3 = exometer_slide:sum([S1, S2], {0}),
+ [500, 505, 510, 515, 520, 525] = [T || {T, _} <- exometer_slide:to_list(525, S3)],
+ [7,6,5,4,3,2] = lists:reverse([V || {_, {V}} <- exometer_slide:to_list(525, S3)]).
+
+incremental_sum_with_drop(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 25, [{incremental, true},
+ {max_n, 5},
+ {interval, 5}]),
+
+ S1 = lists:foldl(fun ({Next, Incr}, S) ->
+ exometer_slide:add_element(Now + Next, {Incr}, S)
+ end, Slide, [{1, 1}, {8, 0}, {15, 0}, {21, 1}, {27, 0}]),
+
+ S2 = lists:foldl(fun (Next, S) ->
+ exometer_slide:add_element(Now + Next, {1}, S)
+ end, Slide, [2, 7, 14, 20, 25]),
+ S3 = exometer_slide:sum([S1, S2]),
+ [27,22,17,12,7] = lists:reverse([T || {T, _} <- exometer_slide:to_list(27, S3)]),
+ [7,6,4,3,2] = lists:reverse([V || {_, {V}} <- exometer_slide:to_list(27, S3)]).
+
+incremental_sum_with_total(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 50, [{incremental, true}, {interval, 5}]),
+
+ S1 = lists:foldl(fun (Next, S) ->
+ exometer_slide:add_element(Now + Next, {1}, S)
+ end, Slide, [5, 10, 15, 20, 25]),
+
+ S2 = lists:foldl(fun (Next, S) ->
+ exometer_slide:add_element(Now + Next, {1}, S)
+ end, Slide, [7, 12, 17, 22, 23]),
+ S3 = exometer_slide:sum([S1, S2]),
+ {10} = exometer_slide:last(S3),
+ [25,20,15,10,5] = lists:reverse([T || {T, _} <- exometer_slide:to_list(26, S3)]),
+ [ 9, 7, 5, 3,1] = lists:reverse([V || {_, {V}} <- exometer_slide:to_list(26, S3)]).
+
+foldl_realises_partial_sample(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 25, [{incremental, true}, {interval, 5}]),
+ S = lists:foldl(fun (Next, S) ->
+ exometer_slide:add_element(Now + Next, {1}, S)
+ end, Slide, [5, 10, 15, 20, 23]),
+ Fun = fun(last, Acc) -> Acc;
+ ({TS, {X}}, Acc) -> [{TS, X} | Acc]
+ end,
+
+ [{25, 5}, {20, 4}, {15, 3}, {10, 2}, {5, 1}] =
+ exometer_slide:foldl(25, 5, Fun, [], S),
+ [{20, 4}, {15, 3}, {10, 2}, {5, 1}] =
+ exometer_slide:foldl(20, 5, Fun, [], S),
+ % do not realise sample unless Now is at least an interval beyond the last
+ % full sample
+ [{20, 4}, {15, 3}, {10, 2}, {5, 1}] =
+ exometer_slide:foldl(23, 5, Fun, [], S).
+
+optimize(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 25, [{interval, 5}, {max_n, 5}]),
+ S = lists:foldl(fun (Next, S) ->
+ exometer_slide:add_element(Now + Next, {Next}, S)
+ end, Slide, [5, 10, 15, 20, 25, 30, 35]),
+ OS = exometer_slide:optimize(S),
+ SRes = exometer_slide:to_list(35, S),
+ OSRes = exometer_slide:to_list(35, OS),
+ SRes = OSRes,
+ ?assert(S =/= OS).
+
+to_list_with_drop(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 25, [{interval, 5},
+ {incremental, true},
+ {max_n, 5}]),
+ S = exometer_slide:add_element(30, {1}, Slide),
+ S2 = exometer_slide:add_element(35, {1}, S),
+ S3 = exometer_slide:add_element(40, {0}, S2),
+ S4 = exometer_slide:add_element(45, {0}, S3),
+ [{30, {1}}, {35, {2}}, {40, {2}}, {45, {2}}] = exometer_slide:to_list(45, S4).
+
+to_list_simple(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 25, [{interval, 5},
+ {incremental, true},
+ {max_n, 5}]),
+ S = exometer_slide:add_element(30, {0}, Slide),
+ S2 = exometer_slide:add_element(35, {0}, S),
+ [{30, {0}}, {35, {0}}] = exometer_slide:to_list(38, S2).
+
+foldl_with_drop(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 25, [{interval, 5},
+ {incremental, true},
+ {max_n, 5}]),
+ S = exometer_slide:add_element(30, {1}, Slide),
+ S2 = exometer_slide:add_element(35, {1}, S),
+ S3 = exometer_slide:add_element(40, {0}, S2),
+ S4 = exometer_slide:add_element(45, {0}, S3),
+ Fun = fun(last, Acc) -> Acc;
+ ({TS, {X}}, Acc) -> [{TS, X} | Acc]
+ end,
+ [{45, 2}, {40, 2}, {35, 2}, {30, 1}] =
+ exometer_slide:foldl(45, 30, Fun, [], S4).
+
+foldl_and_to_list(_Config) ->
+ Now = 0,
+ Tests = [ % {input, expected, query range}
+ {[],
+ [],
+ {0, 10}},
+ {[{5, 1}],
+ [{5, {1}}],
+ {0, 5}},
+ {[{10, 1}],
+ [{10, {1}}],
+ {0, 10}},
+ {[{5, 1}, {10, 2}],
+ [{10, {2}}, {5, {1}}],
+ {0, 10}},
+ {[{5, 0}, {10, 0}], % drop 1
+ [{10, {0}}, {5, {0}}],
+ {0, 10}},
+ {[{5, 2}, {10, 1}, {15, 1}], % drop 2
+ [{15, {1}}, {10, {1}}, {5, {2}}],
+ {0, 15}},
+ {[{10, 0}, {15, 0}, {20, 0}], % drop
+ [{20, {0}}, {15, {0}}, {10, {0}}],
+ {0, 20}},
+ {[{5, 1}, {10, 5}, {15, 5}, {20, 0}], % drop middle
+ [{20, {0}}, {15, {5}}, {10, {5}}, {5, {1}}],
+ {0, 20}},
+ {[{5, 1}, {10, 5}, {15, 5}, {20, 1}], % drop middle filtered
+ [{20, {1}}, {15, {5}}, {10, {5}}],
+ {10, 20}},
+ {[{5, 1}, {10, 2}, {15, 3}, {20, 4}, {25, 4}, {30, 5}], % buffer roll over
+ [{30, {5}}, {25, {4}}, {20, {4}}, {15, {3}}, {10, {2}}, {5, {1}}],
+ {5, 30}}
+ ],
+ Slide = exometer_slide:new(Now, 25, [{interval, 5},
+ {max_n, 5}]),
+ Fun = fun(last, Acc) -> Acc;
+ (V, Acc) -> [V | Acc]
+ end,
+ [begin
+ S = lists:foldl(fun ({T, V}, Acc) ->
+ exometer_slide:add_element(T, {V}, Acc)
+ end, Slide, Inputs),
+ Expected = exometer_slide:foldl(To, From, Fun, [], S),
+ ExpRev = lists:reverse(Expected),
+ ExpRev = exometer_slide:to_list(To, From, S)
+ end || {Inputs, Expected, {From, To}} <- Tests].
+
+foldl_and_to_list_incremental(_Config) ->
+ Now = 0,
+ Tests = [ % {input, expected, query range}
+ {[],
+ [],
+ {0, 10}},
+ {[{5, 1}],
+ [{5, {1}}],
+ {0, 5}},
+ {[{10, 1}],
+ [{10, {1}}],
+ {0, 10}},
+ {[{5, 1}, {10, 1}],
+ [{10, {2}}, {5, {1}}],
+ {0, 10}},
+ {[{5, 0}, {10, 0}], % drop 1
+ [{10, {0}}, {5, {0}}],
+ {0, 10}},
+ {[{5, 1}, {10, 0}, {15, 0}], % drop 2
+ [{15, {1}}, {10, {1}}, {5, {1}}],
+ {0, 15}},
+ {[{10, 0}, {15, 0}, {20, 0}], % drop
+ [{20, {0}}, {15, {0}}, {10, {0}}],
+ {0, 20}},
+ {[{5, 1}, {10, 0}, {15, 0}, {20, 1}], % drop middle
+ [{20, {2}}, {15, {1}}, {10, {1}}, {5, {1}}],
+ {0, 20}},
+ {[{5, 1}, {10, 0}, {15, 0}, {20, 1}], % drop middle filtered
+ [{20, {2}}, {15, {1}}, {10, {1}}],
+ {10, 20}},
+ {[{5, 1}, {10, 1}, {15, 1}, {20, 1}, {25, 0}, {30, 1}], % buffer roll over
+ [{30, {5}}, {25, {4}}, {20, {4}}, {15, {3}}, {10, {2}}, {5, {1}}],
+ {5, 30}}
+ ],
+ Slide = exometer_slide:new(Now, 25, [{interval, 5},
+ {incremental, true},
+ {max_n, 5}]),
+ Fun = fun(last, Acc) -> Acc;
+ (V, Acc) -> [V | Acc]
+ end,
+ [begin
+ S = lists:foldl(fun ({T, V}, Acc) ->
+ exometer_slide:add_element(T, {V}, Acc)
+ end, Slide, Inputs),
+ Expected = exometer_slide:foldl(To, From, Fun, [], S),
+ ExpRev = lists:reverse(Expected),
+ ExpRev = exometer_slide:to_list(To, From, S)
+ end || {Inputs, Expected, {From, To}} <- Tests].
+
+stale_to_list(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 25, [{interval, 5}, {max_n, 5}]),
+ S = exometer_slide:add_element(50, {1}, Slide),
+ S2 = exometer_slide:add_element(55, {1}, S),
+ [] = exometer_slide:to_list(100, S2).
+
+to_list_single_after_drop(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 25, [{interval, 5},
+ {incremental, true},
+ {max_n, 5}]),
+ S = exometer_slide:add_element(5, {0}, Slide),
+ S2 = exometer_slide:add_element(10, {0}, S),
+ S3 = exometer_slide:add_element(15, {1}, S2),
+ Res = exometer_slide:to_list(17, S3),
+ [{5,{0}},{10,{0}},{15,{1}}] = Res.
+
+
+to_list_drop_and_roll(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 10, [{interval, 5},
+ {incremental, true},
+ {max_n, 5}]),
+ S = exometer_slide:add_element(5, {0}, Slide),
+ S2 = exometer_slide:add_element(10, {0}, S),
+ S3 = exometer_slide:add_element(15, {0}, S2),
+ [{10, {0}}, {15, {0}}] = exometer_slide:to_list(17, S3).
+
+
+sum_single(_Config) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 25, [{interval, 5},
+ {incremental, true},
+ {max_n, 5}]),
+ S = exometer_slide:add_element(Now + 5, {0}, Slide),
+ S2 = exometer_slide:add_element(Now + 10, {0}, S),
+ Summed = exometer_slide:sum([S2]),
+ [_,_] = exometer_slide:to_list(15, Summed).
+
+
+to_normalized_list(_Config) ->
+ Interval = 5,
+ Tests = [ % {input, expected, query range}
+ {[], % zero pad when slide has never seen any samples
+ [{10, {0}}, {5, {0}}, {0, {0}}],
+ {0, 10}},
+ {[{5, 1}], % zero pad before first known sample
+ [{5, {1}}, {0, {0}}],
+ {0, 5}},
+ {[{10, 1}, {15, 1}], % zero pad before last know sample
+ [{15, {2}}, {10, {1}}, {5, {0}}],
+ {5, 15}},
+ {[{5, 1}, {15, 1}], % insert missing sample using previous total
+ [{15, {2}}, {10, {1}}, {5, {1}}],
+ {5, 15}},
+ % {[{6, 1}, {11, 1}, {16, 1}], % align timestamps with query
+ % [{15, {3}}, {10, {2}}, {5, {1}}, {0, {0}}],
+ % {0, 15}},
+ {[{5, 1}, {10, 1}, {15, 1}, {20, 1}, {25, 1}, {30, 1}], % outside of max_n
+ [{30, {6}}, {25, {5}}, {20, {4}}, {15, {3}}, {10, {2}}], % we cannot possibly be expected deduce what 10 should be
+ {10, 30}},
+ {[{5, 1}, {20, 1}, {25, 1}], % as long as the past TS 5 sample still exists we should use to for padding
+ [{25, {3}}, {20, {2}}, {15, {1}}, {10, {1}}],
+ {10, 25}},
+ {[{5, 1}, {10, 1}], % pad based on total
+ [{35, {2}}, {30, {2}}],
+ {30, 35}},
+ {[{5, 1}], % make up future values to fill the window
+ [{10, {1}}, {5, {1}}],
+ {5, 10}},
+ {[{5, 1}, {7, 1}], % realise last sample
+ [{10, {2}}, {5, {1}}],
+ {5, 10}}
+ ],
+
+ Slide = exometer_slide:new(0, 20, [{interval, 5},
+ {incremental, true},
+ {max_n, 4}]),
+ [begin
+ S0 = lists:foldl(fun ({T, V}, Acc) ->
+ exometer_slide:add_element(T, {V}, Acc)
+ end, Slide, Inputs),
+ Expected = exometer_slide:to_normalized_list(To, From, Interval, S0, {0}),
+ S = exometer_slide:sum([exometer_slide:optimize(S0)], {0}), % also test it post sum
+ Expected = exometer_slide:to_normalized_list(To, From, Interval, S, {0})
+ end || {Inputs, Expected, {From, To}} <- Tests].
+
+to_normalized_list_no_padding(_Config) ->
+ Interval = 5,
+ Tests = [ % {input, expected, query range}
+ {[],
+ [],
+ {0, 10}},
+ {[{5, 1}],
+ [{5, {1}}],
+ {0, 5}},
+ {[{5, 1}, {15, 1}],
+ [{15, {2}}, {10, {1}}, {5, {1}}],
+ {5, 15}},
+ {[{10, 1}, {15, 1}],
+ [{15, {2}}, {10, {1}}],
+ {5, 15}},
+ {[{5, 1}, {20, 1}], % NB as 5 is outside of the query we can't pick the value up
+ [{20, {2}}],
+ {10, 20}}
+ ],
+
+ Slide = exometer_slide:new(0, 20, [{interval, 5},
+ {incremental, true},
+ {max_n, 4}]),
+ [begin
+ S = lists:foldl(fun ({T, V}, Acc) ->
+ exometer_slide:add_element(T, {V}, Acc)
+ end, Slide, Inputs),
+ Expected = exometer_slide:to_normalized_list(To, From, Interval, S, no_pad)
+ end || {Inputs, Expected, {From, To}} <- Tests].
+
+to_list_in_the_past(_Config) ->
+ Slide = exometer_slide:new(0, 20, [{interval, 5},
+ {incremental, true},
+ {max_n, 4}]),
+ % ensure firstTS is way in the past
+ S0 = exometer_slide:add_element(5, {1}, Slide),
+ S1 = exometer_slide:add_element(105, {0}, S0),
+ S = exometer_slide:add_element(110, {0}, S1), % create drop
+ % query into the past
+ % this could happen if a node with and incorrect clock joins the cluster
+ [] = exometer_slide:to_list(50, 10, S).
+
+sum_mgmt_352(_Config) ->
+ %% In bug mgmt#352 all the samples returned have the same vale
+ Slide = sum_mgmt_352_slide(),
+ Last = 1487689330000,
+ First = 1487689270000,
+ Incr = 5000,
+ Empty = {0},
+ Sum = exometer_slide:sum(Last, First, Incr, [Slide], Empty),
+ Values = sets:to_list(sets:from_list(
+ [V || {_, V} <- exometer_slide:buffer(Sum)])),
+ true = (length(Values) == 12),
+ ok.
+
+sum_mgmt_352_extra(_Config) ->
+ %% Testing previous case clause to the one that fixes mgmt_352
+ %% exometer_slide.erl#L463
+ %% In the buggy version, all but the last sample are the same
+ Slide = sum_mgmt_352_slide_extra(),
+ Last = 1487689330000,
+ First = 1487689260000,
+ Incr = 5000,
+ Empty = {0},
+ Sum = exometer_slide:sum(Last, First, Incr, [Slide], Empty),
+ Values = sets:to_list(sets:from_list(
+ [V || {_, V} <- exometer_slide:buffer(Sum)])),
+ true = (length(Values) == 13),
+ ok.
+
+sum_mgmt_352_peak(_Config) ->
+ %% When buf2 contains data, we were returning a too old sample that
+ %% created a massive rate peak at the beginning of the graph.
+ %% The sample used was captured during a debug session.
+ Slide = sum_mgmt_352_slide_peak(),
+ Last = 1487752040000,
+ First = 1487751980000,
+ Incr = 5000,
+ Empty = {0},
+ Sum = exometer_slide:sum(Last, First, Incr, [Slide], Empty),
+ [{LastV}, {BLastV} | _] =
+ lists:reverse([V || {_, V} <- exometer_slide:buffer(Sum)]),
+ Rate = (BLastV - LastV) div 5,
+ true = (Rate < 20000),
+ ok.
+
+%% -------------------------------------------------------------------
+%% Util
+%% -------------------------------------------------------------------
+
+ele(TS, V) -> {TS, {V}}.
+
+%% -------------------------------------------------------------------
+%% Data
+%% -------------------------------------------------------------------
+sum_mgmt_352_slide() ->
+ %% Provide slide as is, from a debug session triggering mgmt-352 bug
+ {slide,610000,45,122,true,5000,1487689328468,1487689106834,
+ [{1487689328468,{1574200}},
+ {1487689323467,{1538800}},
+ {1487689318466,{1500800}},
+ {1487689313465,{1459138}},
+ {1487689308463,{1419200}},
+ {1487689303462,{1379600}},
+ {1487689298461,{1340000}},
+ {1487689293460,{1303400}},
+ {1487689288460,{1265600}},
+ {1487689283458,{1231400}},
+ {1487689278457,{1215800}},
+ {1487689273456,{1215200}},
+ {1487689262487,drop},
+ {1487689257486,{1205600}}],
+ [],
+ {1591000}}.
+
+sum_mgmt_352_slide_extra() ->
+ {slide,610000,45,122,true,5000,1487689328468,1487689106834,
+ [{1487689328468,{1574200}},
+ {1487689323467,{1538800}},
+ {1487689318466,{1500800}},
+ {1487689313465,{1459138}},
+ {1487689308463,{1419200}},
+ {1487689303462,{1379600}},
+ {1487689298461,{1340000}},
+ {1487689293460,{1303400}},
+ {1487689288460,{1265600}},
+ {1487689283458,{1231400}},
+ {1487689278457,{1215800}},
+ {1487689273456,{1215200}},
+ {1487689272487,drop},
+ {1487689269486,{1205600}}],
+ [],
+ {1591000}}.
+
+sum_mgmt_352_slide_peak() ->
+ {slide,610000,96,122,true,5000,1487752038481,1487750936863,
+ [{1487752038481,{11994024}},
+ {1487752033480,{11923200}},
+ {1487752028476,{11855800}},
+ {1487752023474,{11765800}},
+ {1487752018473,{11702431}},
+ {1487752013472,{11636200}},
+ {1487752008360,{11579800}},
+ {1487752003355,{11494800}},
+ {1487751998188,{11441400}},
+ {1487751993184,{11381000}},
+ {1487751988180,{11320000}},
+ {1487751983178,{11263000}},
+ {1487751978177,{11187600}},
+ {1487751973172,{11123375}},
+ {1487751968167,{11071800}},
+ {1487751963166,{11006200}},
+ {1487751958162,{10939477}},
+ {1487751953161,{10882400}},
+ {1487751948140,{10819600}},
+ {1487751943138,{10751200}},
+ {1487751938134,{10744400}},
+ {1487751933129,drop},
+ {1487751927807,{10710200}},
+ {1487751922803,{10670000}}],
+ [{1487751553386,{6655800}},
+ {1487751548385,{6580365}},
+ {1487751543384,{6509358}}],
+ {11994024}}.
diff --git a/deps/rabbitmq_management_agent/test/metrics_SUITE.erl b/deps/rabbitmq_management_agent/test/metrics_SUITE.erl
new file mode 100644
index 0000000000..227a04b21c
--- /dev/null
+++ b/deps/rabbitmq_management_agent/test/metrics_SUITE.erl
@@ -0,0 +1,157 @@
+%% This Source Code Form is subject to the terms of the Mozilla Public
+%% License, v. 2.0. If a copy of the MPL was not distributed with this
+%% file, You can obtain one at https://mozilla.org/MPL/2.0/.
+%%
+%% Copyright (c) 2016-2020 VMware, Inc. or its affiliates. All rights reserved.
+%%
+-module(metrics_SUITE).
+-compile(export_all).
+
+-include_lib("common_test/include/ct.hrl").
+-include_lib("amqp_client/include/amqp_client.hrl").
+
+all() ->
+ [
+ {group, non_parallel_tests}
+ ].
+
+groups() ->
+ [
+ {non_parallel_tests, [], [
+ node,
+ storage_reset
+ ]}
+ ].
+
+%% -------------------------------------------------------------------
+%% Testsuite setup/teardown.
+%% -------------------------------------------------------------------
+
+merge_app_env(Config) ->
+ _Config1 = rabbit_ct_helpers:merge_app_env(Config,
+ {rabbit, [
+ {collect_statistics, fine},
+ {collect_statistics_interval, 500}
+ ]}).
+ %% rabbit_ct_helpers:merge_app_env(
+ %% Config1, {rabbitmq_management_agent, [{sample_retention_policies,
+ %% [{global, [{605, 500}]},
+ %% {basic, [{605, 500}]},
+ %% {detailed, [{10, 500}]}] }]}).
+
+init_per_suite(Config) ->
+ rabbit_ct_helpers:log_environment(),
+ Config1 = rabbit_ct_helpers:set_config(Config, [
+ {rmq_nodename_suffix, ?MODULE},
+ {rmq_nodes_count, 2}
+ ]),
+ rabbit_ct_helpers:run_setup_steps(Config1,
+ [ fun merge_app_env/1 ] ++
+ rabbit_ct_broker_helpers:setup_steps() ++
+ rabbit_ct_client_helpers:setup_steps()).
+
+end_per_suite(Config) ->
+ rabbit_ct_helpers:run_teardown_steps(Config,
+ rabbit_ct_client_helpers:teardown_steps() ++
+ rabbit_ct_broker_helpers:teardown_steps()).
+
+init_per_group(_, Config) ->
+ Config.
+
+end_per_group(_, Config) ->
+ Config.
+
+init_per_testcase(Testcase, Config) ->
+ rabbit_ct_helpers:testcase_started(Config, Testcase).
+
+end_per_testcase(Testcase, Config) ->
+ rabbit_ct_helpers:testcase_finished(Config, Testcase).
+
+
+%% -------------------------------------------------------------------
+%% Testcases.
+%% -------------------------------------------------------------------
+
+read_table_rpc(Config, Table) ->
+ rabbit_ct_broker_helpers:rpc(Config, 0, ?MODULE, read_table, [Table]).
+
+read_table(Table) ->
+ ets:tab2list(Table).
+
+force_stats() ->
+ rabbit_mgmt_external_stats ! emit_update.
+
+node(Config) ->
+ [A, B] = rabbit_ct_broker_helpers:get_node_configs(Config, nodename),
+ % force multipe stats refreshes
+ [ rabbit_ct_broker_helpers:rpc(Config, 0, ?MODULE, force_stats, [])
+ || _ <- lists:seq(0, 10)],
+ [_] = read_table_rpc(Config, node_persister_metrics),
+ [_] = read_table_rpc(Config, node_coarse_metrics),
+ [_] = read_table_rpc(Config, node_metrics),
+ true = wait_until(
+ fun() ->
+ Tab = read_table_rpc(Config, node_node_metrics),
+ lists:keymember({A, B}, 1, Tab)
+ end, 10).
+
+
+storage_reset(Config) ->
+ %% Ensures that core stats are reset, otherwise consume generates negative values
+ %% Doesn't really test if the reset does anything!
+ {Ch, Q} = publish_msg(Config),
+ wait_until(fun() ->
+ {1, 0, 1} == get_vhost_stats(Config)
+ end),
+ ok = rabbit_ct_broker_helpers:rpc(Config, 0, rabbit_mgmt_storage, reset, []),
+ wait_until(fun() ->
+ {1, 0, 1} == get_vhost_stats(Config)
+ end),
+ consume_msg(Ch, Q),
+ wait_until(fun() ->
+ {0, 0, 0} == get_vhost_stats(Config)
+ end),
+ rabbit_ct_client_helpers:close_channel(Ch).
+
+%% -------------------------------------------------------------------
+%% Helpers
+%% -------------------------------------------------------------------
+
+publish_msg(Config) ->
+ Ch = rabbit_ct_client_helpers:open_channel(Config, 0),
+ #'queue.declare_ok'{queue = Q} =
+ amqp_channel:call(Ch, #'queue.declare'{exclusive = true}),
+ amqp_channel:cast(Ch, #'basic.publish'{routing_key = Q},
+ #amqp_msg{props = #'P_basic'{delivery_mode = 2},
+ payload = Q}),
+ {Ch, Q}.
+
+consume_msg(Ch, Q) ->
+ amqp_channel:subscribe(Ch, #'basic.consume'{queue = Q,
+ no_ack = true}, self()),
+ receive #'basic.consume_ok'{} -> ok
+ end,
+ receive {#'basic.deliver'{}, #amqp_msg{payload = Q}} ->
+ ok
+ end.
+
+wait_until(Fun) ->
+ wait_until(Fun, 120).
+
+wait_until(_, 0) ->
+ false;
+wait_until(Fun, N) ->
+ case Fun() of
+ true ->
+ true;
+ false ->
+ timer:sleep(1000),
+ wait_until(Fun, N-1)
+ end.
+
+get_vhost_stats(Config) ->
+ Dict = rabbit_ct_broker_helpers:rpc(Config, 0, rabbit_mgmt_data, overview_data,
+ [none, all ,{no_range, no_range, no_range,no_range},
+ [<<"/">>]]),
+ {ok, {VhostMsgStats, _}} = maps:find(vhost_msg_stats, Dict),
+ exometer_slide:last(VhostMsgStats).
diff --git a/deps/rabbitmq_management_agent/test/rabbit_mgmt_gc_SUITE.erl b/deps/rabbitmq_management_agent/test/rabbit_mgmt_gc_SUITE.erl
new file mode 100644
index 0000000000..b5ee4e9ce2
--- /dev/null
+++ b/deps/rabbitmq_management_agent/test/rabbit_mgmt_gc_SUITE.erl
@@ -0,0 +1,626 @@
+%% This Source Code Form is subject to the terms of the Mozilla Public
+%% License, v. 2.0. If a copy of the MPL was not distributed with this
+%% file, You can obtain one at https://mozilla.org/MPL/2.0/.
+%%
+%% Copyright (c) 2007-2020 VMware, Inc. or its affiliates. All rights reserved.
+%%
+
+-module(rabbit_mgmt_gc_SUITE).
+
+-include_lib("common_test/include/ct.hrl").
+-include_lib("eunit/include/eunit.hrl").
+-include_lib("amqp_client/include/amqp_client.hrl").
+-include("rabbit_mgmt_metrics.hrl").
+
+-compile(export_all).
+
+all() ->
+ [
+ {group, non_parallel_tests}
+ ].
+
+groups() ->
+ [
+ {non_parallel_tests, [],
+ [ queue_stats,
+ quorum_queue_stats,
+ connection_stats,
+ channel_stats,
+ vhost_stats,
+ exchange_stats,
+ node_stats,
+ consumer_stats
+ ]
+ }
+ ].
+
+%% -------------------------------------------------------------------
+%% Testsuite setup/teardown.
+%% -------------------------------------------------------------------
+
+merge_app_env(Config) ->
+ Config1 = rabbit_ct_helpers:merge_app_env(
+ Config,
+ {rabbitmq_management_agent, [
+ {metrics_gc_interval, 6000000},
+ {rates_mode, detailed},
+ {sample_retention_policies,
+ %% List of {MaxAgeInSeconds, SampleEveryNSeconds}
+ [{global, [{605, 1}, {3660, 60}]},
+ {basic, [{605, 1}, {3600, 60}]},
+ {detailed, [{605, 1}]}] }]}),
+ rabbit_ct_helpers:merge_app_env(Config1,
+ {rabbit, [
+ {collect_statistics_interval, 100},
+ {collect_statistics, fine}
+ ]}).
+
+init_per_suite(Config) ->
+ rabbit_ct_helpers:log_environment(),
+ Config1 = rabbit_ct_helpers:set_config(Config, [
+ {rmq_nodename_suffix, ?MODULE},
+ {rmq_nodes_count, 2},
+ {rmq_nodes_clustered, true}
+ ]),
+ rabbit_ct_helpers:run_setup_steps(
+ Config1,
+ [ fun merge_app_env/1 ] ++ rabbit_ct_broker_helpers:setup_steps()).
+
+end_per_suite(Config) ->
+ rabbit_ct_helpers:run_teardown_steps(
+ Config,
+ rabbit_ct_broker_helpers:teardown_steps()).
+
+init_per_group(_, Config) ->
+ Config.
+
+end_per_group(_, Config) ->
+ Config.
+
+init_per_testcase(quorum_queue_stats = Testcase, Config) ->
+ case rabbit_ct_broker_helpers:enable_feature_flag(Config, quorum_queue) of
+ ok ->
+ rabbit_ct_helpers:testcase_started(Config, Testcase),
+ rabbit_ct_helpers:run_steps(
+ Config, rabbit_ct_client_helpers:setup_steps());
+ Skip ->
+ Skip
+ end;
+init_per_testcase(Testcase, Config) ->
+ rabbit_ct_helpers:testcase_started(Config, Testcase),
+ rabbit_ct_helpers:run_steps(Config,
+ rabbit_ct_client_helpers:setup_steps()).
+
+end_per_testcase(Testcase, Config) ->
+ rabbit_ct_helpers:testcase_finished(Config, Testcase),
+ rabbit_ct_helpers:run_teardown_steps(
+ Config,
+ rabbit_ct_client_helpers:teardown_steps()).
+
+%% -------------------------------------------------------------------
+%% Testcases.
+%% -------------------------------------------------------------------
+
+queue_stats(Config) ->
+ A = rabbit_ct_broker_helpers:get_node_config(Config, 0, nodename),
+ Ch = rabbit_ct_client_helpers:open_channel(Config, A),
+
+ amqp_channel:call(Ch, #'queue.declare'{queue = <<"queue_stats">>}),
+ amqp_channel:cast(Ch, #'basic.publish'{routing_key = <<"queue_stats">>},
+ #amqp_msg{payload = <<"hello">>}),
+ {#'basic.get_ok'{}, _} = amqp_channel:call(Ch, #'basic.get'{queue = <<"queue_stats">>,
+ no_ack = true}),
+ timer:sleep(1150),
+
+ Q = q(<<"myqueue">>),
+ X = x(<<"">>),
+
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [queue_stats, {Q, infos}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [queue_stats_publish, {{Q, 5}, slide}]),
+ 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, A, ets, insert,
+ [queue_stats_deliver_stats, {{Q, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [queue_exchange_stats_publish,
+ {{{Q, X}, 5}, slide}]),
+
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [queue_stats, Q]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [queue_stats_publish, {Q, 5}]),
+ [_] = 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, A, ets, lookup,
+ [queue_stats_deliver_stats, {Q, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [queue_exchange_stats_publish, {{Q, X}, 5}]),
+
+ %% Trigger gc. When the gen_server:call returns, the gc has already finished.
+ rabbit_ct_broker_helpers:rpc(Config, A, erlang, send, [rabbit_mgmt_gc, start_gc]),
+ rabbit_ct_broker_helpers:rpc(Config, A, gen_server, call, [rabbit_mgmt_gc, test]),
+
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [queue_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [queue_stats_publish]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [queue_msg_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [queue_process_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [queue_msg_rates]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [queue_stats_deliver_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [queue_exchange_stats_publish]),
+
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [queue_stats, Q]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [queue_stats_publish, {Q, 5}]),
+ [] = 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, A, ets, lookup,
+ [queue_stats_deliver_stats, {Q, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [queue_exchange_stats_publish, {{Q, X}, 5}]),
+
+ amqp_channel:call(Ch, #'queue.delete'{queue = <<"queue_stats">>}),
+ rabbit_ct_client_helpers:close_channel(Ch),
+
+ 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),
+
+ amqp_channel:call(Ch, #'queue.declare'{queue = <<"queue_stats">>}),
+ amqp_channel:cast(Ch, #'basic.publish'{routing_key = <<"queue_stats">>},
+ #amqp_msg{payload = <<"hello">>}),
+ timer:sleep(1150),
+
+ DeadPid = rabbit_ct_broker_helpers:rpc(Config, A, ?MODULE, dead_pid, []),
+
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [connection_stats_coarse_conn_stats,
+ {{DeadPid, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [connection_created_stats, {DeadPid, name, infos}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [connection_stats, {DeadPid, infos}]),
+
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [connection_stats_coarse_conn_stats, {DeadPid, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [connection_created_stats, DeadPid]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [connection_stats, DeadPid]),
+
+ %% Trigger gc. When the gen_server:call returns, the gc has already finished.
+ rabbit_ct_broker_helpers:rpc(Config, A, erlang, send, [rabbit_mgmt_gc, start_gc]),
+ rabbit_ct_broker_helpers:rpc(Config, A, gen_server, call, [rabbit_mgmt_gc, test]),
+
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [connection_stats_coarse_conn_stats, {DeadPid, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [connection_created_stats, DeadPid]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [connection_stats, DeadPid]),
+
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [connection_stats_coarse_conn_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [connection_created_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [connection_stats]),
+
+ amqp_channel:call(Ch, #'queue.delete'{queue = <<"queue_stats">>}),
+ rabbit_ct_client_helpers:close_channel(Ch),
+
+ ok.
+
+channel_stats(Config) ->
+ A = rabbit_ct_broker_helpers:get_node_config(Config, 0, nodename),
+ Ch = rabbit_ct_client_helpers:open_channel(Config, A),
+
+ amqp_channel:call(Ch, #'queue.declare'{queue = <<"queue_stats">>}),
+ amqp_channel:cast(Ch, #'basic.publish'{routing_key = <<"queue_stats">>},
+ #amqp_msg{payload = <<"hello">>}),
+ {#'basic.get_ok'{}, _} = amqp_channel:call(Ch, #'basic.get'{queue = <<"queue_stats">>,
+ no_ack=true}),
+ timer:sleep(1150),
+
+ DeadPid = rabbit_ct_broker_helpers:rpc(Config, A, ?MODULE, dead_pid, []),
+
+ X = x(<<"myexchange">>),
+
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [channel_created_stats, {DeadPid, name, infos}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [channel_stats, {DeadPid, infos}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [channel_stats_fine_stats, {{DeadPid, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [channel_exchange_stats_fine_stats,
+ {{{DeadPid, X}, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [channel_queue_stats_deliver_stats,
+ {{{DeadPid, X}, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [channel_stats_deliver_stats,
+ {{DeadPid, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [channel_process_stats,
+ {{DeadPid, 5}, slide}]),
+
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_created_stats, DeadPid]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_stats, DeadPid]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_stats_fine_stats, {DeadPid, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_exchange_stats_fine_stats,
+ {{DeadPid, X}, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_queue_stats_deliver_stats,
+ {{DeadPid, X}, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_stats_deliver_stats,
+ {DeadPid, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_process_stats,
+ {DeadPid, 5}]),
+
+ %% Trigger gc. When the gen_server:call returns, the gc has already finished.
+ rabbit_ct_broker_helpers:rpc(Config, A, erlang, send, [rabbit_mgmt_gc, start_gc]),
+ rabbit_ct_broker_helpers:rpc(Config, A, gen_server, call, [rabbit_mgmt_gc, test]),
+
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [channel_created_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [channel_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [channel_stats_fine_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [channel_exchange_stats_fine_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [channel_queue_stats_deliver_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [channel_stats_deliver_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [channel_process_stats]),
+
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_created_stats, DeadPid]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_stats, DeadPid]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_stats_fine_stats, {DeadPid, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_exchange_stats_fine_stats,
+ {{DeadPid, X}, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_queue_stats_deliver_stats,
+ {{DeadPid, X}, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_stats_deliver_stats,
+ {DeadPid, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [channel_process_stats,
+ {DeadPid, 5}]),
+
+ amqp_channel:call(Ch, #'queue.delete'{queue = <<"queue_stats">>}),
+ rabbit_ct_client_helpers:close_channel(Ch),
+
+ ok.
+
+vhost_stats(Config) ->
+ A = rabbit_ct_broker_helpers:get_node_config(Config, 0, nodename),
+ Ch = rabbit_ct_client_helpers:open_channel(Config, A),
+
+ amqp_channel:call(Ch, #'queue.declare'{queue = <<"queue_stats">>}),
+ amqp_channel:cast(Ch, #'basic.publish'{routing_key = <<"queue_stats">>},
+ #amqp_msg{payload = <<"hello">>}),
+ {#'basic.get_ok'{}, _} = amqp_channel:call(Ch, #'basic.get'{queue = <<"queue_stats">>,
+ no_ack=true}),
+ timer:sleep(1150),
+
+ VHost = <<"myvhost">>,
+
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [vhost_stats_coarse_conn_stats,
+ {{VHost, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [vhost_stats_fine_stats,
+ {{VHost, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [vhost_stats_deliver_stats,
+ {{VHost, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert, [vhost_msg_stats,
+ {{VHost, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert, [vhost_msg_rates,
+ {{VHost, 5}, slide}]),
+
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [vhost_stats_coarse_conn_stats, {VHost, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [vhost_stats_fine_stats, {VHost, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [vhost_stats_deliver_stats, {VHost, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [vhost_msg_stats, {VHost, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [vhost_msg_rates, {VHost, 5}]),
+
+ %% Trigger gc. When the gen_server:call returns, the gc has already finished.
+ rabbit_ct_broker_helpers:rpc(Config, A, erlang, send, [rabbit_mgmt_gc, start_gc]),
+ rabbit_ct_broker_helpers:rpc(Config, A, gen_server, call, [rabbit_mgmt_gc, test]),
+
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [vhost_stats_coarse_conn_stats, {VHost, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [vhost_stats_fine_stats, {VHost, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [vhost_stats_deliver_stats, {VHost, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [vhost_msg_stats, {VHost, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [vhost_msg_rates, {VHost, 5}]),
+
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [vhost_stats_coarse_conn_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [vhost_stats_fine_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [vhost_stats_deliver_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [vhost_msg_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [vhost_msg_rates]),
+
+ amqp_channel:call(Ch, #'queue.delete'{queue = <<"queue_stats">>}),
+ rabbit_ct_client_helpers:close_channel(Ch),
+
+ ok.
+
+exchange_stats(Config) ->
+ A = rabbit_ct_broker_helpers:get_node_config(Config, 0, nodename),
+ Ch = rabbit_ct_client_helpers:open_channel(Config, A),
+
+ amqp_channel:call(Ch, #'queue.declare'{queue = <<"queue_stats">>}),
+ amqp_channel:cast(Ch, #'basic.publish'{routing_key = <<"queue_stats">>},
+ #amqp_msg{payload = <<"hello">>}),
+ {#'basic.get_ok'{}, _} = amqp_channel:call(Ch, #'basic.get'{queue = <<"queue_stats">>,
+ no_ack=true}),
+ timer:sleep(1150),
+
+ Exchange = x(<<"myexchange">>),
+
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [exchange_stats_publish_out,
+ {{Exchange, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [exchange_stats_publish_in,
+ {{Exchange, 5}, slide}]),
+
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [exchange_stats_publish_out, {Exchange, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [exchange_stats_publish_in, {Exchange, 5}]),
+
+ %% Trigger gc. When the gen_server:call returns, the gc has already finished.
+ rabbit_ct_broker_helpers:rpc(Config, A, erlang, send, [rabbit_mgmt_gc, start_gc]),
+ rabbit_ct_broker_helpers:rpc(Config, A, gen_server, call, [rabbit_mgmt_gc, test]),
+
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [exchange_stats_publish_out, {Exchange, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [exchange_stats_publish_in, {Exchange, 5}]),
+
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [exchange_stats_publish_out]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [exchange_stats_publish_in]),
+
+ amqp_channel:call(Ch, #'queue.delete'{queue = <<"queue_stats">>}),
+ rabbit_ct_client_helpers:close_channel(Ch),
+
+ ok.
+
+node_stats(Config) ->
+ A = rabbit_ct_broker_helpers:get_node_config(Config, 0, nodename),
+
+ timer:sleep(150),
+
+ Node = 'mynode',
+
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert, [node_stats, {Node, infos}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert, [node_coarse_stats,
+ {{Node, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert, [node_persister_stats,
+ {{Node, 5}, slide}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert, [node_node_stats,
+ {{A, Node}, infos}]),
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert, [node_node_coarse_stats,
+ {{{A, Node}, 5}, slide}]),
+
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [node_stats, Node]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [node_coarse_stats, {Node, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [node_persister_stats, {Node, 5}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [node_node_stats, {A, Node}]),
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [node_node_coarse_stats, {{A, Node}, 5}]),
+
+ %% Trigger gc. When the gen_server:call returns, the gc has already finished.
+ rabbit_ct_broker_helpers:rpc(Config, A, erlang, send, [rabbit_mgmt_gc, start_gc]),
+ rabbit_ct_broker_helpers:rpc(Config, A, gen_server, call, [rabbit_mgmt_gc, test]),
+
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [node_stats, Node]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [node_coarse_stats, {Node, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [node_persister_stats, {Node, 5}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [node_node_stats, {A, Node}]),
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup,
+ [node_node_coarse_stats, {{A, Node}, 5}]),
+
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [node_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [node_coarse_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [node_persister_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [node_node_stats]),
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [node_node_coarse_stats]),
+
+ ok.
+
+consumer_stats(Config) ->
+ A = rabbit_ct_broker_helpers:get_node_config(Config, 0, nodename),
+ Ch = rabbit_ct_client_helpers:open_channel(Config, A),
+
+ amqp_channel:call(Ch, #'queue.declare'{queue = <<"queue_stats">>}),
+ amqp_channel:call(Ch, #'basic.consume'{queue = <<"queue_stats">>}),
+ timer:sleep(1150),
+
+ DeadPid = rabbit_ct_broker_helpers:rpc(Config, A, ?MODULE, dead_pid, []),
+ Q = q(<<"queue_stats">>),
+
+ Id = {Q, DeadPid, tag},
+ rabbit_ct_broker_helpers:rpc(Config, A, ets, insert,
+ [consumer_stats, {Id, infos}]),
+
+ [_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup, [consumer_stats, Id]),
+
+ %% Trigger gc. When the gen_server:call returns, the gc has already finished.
+ rabbit_ct_broker_helpers:rpc(Config, A, erlang, send, [rabbit_mgmt_gc, start_gc]),
+ rabbit_ct_broker_helpers:rpc(Config, A, gen_server, call, [rabbit_mgmt_gc, test]),
+
+ [_|_] = rabbit_ct_broker_helpers:rpc(Config, A, ets, tab2list,
+ [consumer_stats]),
+
+ [] = rabbit_ct_broker_helpers:rpc(Config, A, ets, lookup, [consumer_stats, Id]),
+
+ amqp_channel:call(Ch, #'queue.delete'{queue = <<"queue_stats">>}),
+ rabbit_ct_client_helpers:close_channel(Ch),
+
+ ok.
+
+dead_pid() ->
+ spawn(fun() -> ok end).
+
+q(Name) ->
+ #resource{ virtual_host = <<"/">>,
+ kind = queue,
+ name = Name }.
+
+x(Name) ->
+ #resource{ virtual_host = <<"/">>,
+ kind = exchange,
+ name = Name }.
diff --git a/deps/rabbitmq_management_agent/test/rabbit_mgmt_slide_SUITE.erl b/deps/rabbitmq_management_agent/test/rabbit_mgmt_slide_SUITE.erl
new file mode 100644
index 0000000000..0261606bd5
--- /dev/null
+++ b/deps/rabbitmq_management_agent/test/rabbit_mgmt_slide_SUITE.erl
@@ -0,0 +1,174 @@
+%%% This Source Code Form is subject to the terms of the Mozilla Public
+%% License, v. 2.0. If a copy of the MPL was not distributed with this
+%% file, You can obtain one at https://mozilla.org/MPL/2.0/.
+%%
+%% Copyright (c) 2016-2020 VMware, Inc. or its affiliates. All rights reserved.
+%%
+
+-module(rabbit_mgmt_slide_SUITE).
+
+-include_lib("proper/include/proper.hrl").
+
+-compile(export_all).
+
+all() ->
+ [
+ {group, parallel_tests}
+ ].
+
+groups() ->
+ [
+ {parallel_tests, [parallel], [
+ last_two_test,
+ last_two_incremental_test,
+ sum_test,
+ sum_incremental_test
+ ]}
+ ].
+
+%% -------------------------------------------------------------------
+%% Testsuite setup/teardown.
+%% -------------------------------------------------------------------
+init_per_suite(Config) ->
+ rabbit_ct_helpers:log_environment(),
+ Config.
+
+end_per_suite(_Config) ->
+ ok.
+
+init_per_group(_, Config) ->
+ Config.
+
+end_per_group(_, _Config) ->
+ ok.
+
+init_per_testcase(_, Config) ->
+ Config.
+
+end_per_testcase(_, _Config) ->
+ ok.
+
+%% -------------------------------------------------------------------
+%% Generators.
+%% -------------------------------------------------------------------
+elements_gen() ->
+ ?LET(Length, oneof([1, 2, 3, 7, 8, 20]),
+ ?LET(Elements, list(vector(Length, int())),
+ [erlang:list_to_tuple(E) || E <- Elements])).
+
+%% -------------------------------------------------------------------
+%% Testcases.
+%% -------------------------------------------------------------------
+last_two_test(_Config) ->
+ rabbit_ct_proper_helpers:run_proper(fun prop_last_two/0, [], 100).
+
+prop_last_two() ->
+ ?FORALL(
+ Elements, elements_gen(),
+ begin
+ Interval = 1,
+ Incremental = false,
+ {_LastTS, Slide} = new_slide(Interval, Incremental, Elements),
+ Expected = last_two(Elements),
+ ValuesOnly = [V || {_Timestamp, V} <- exometer_slide:last_two(Slide)],
+ ?WHENFAIL(io:format("Last two values obtained: ~p~nExpected: ~p~n"
+ "Slide: ~p~n", [ValuesOnly, Expected, Slide]),
+ Expected == ValuesOnly)
+ end).
+
+last_two_incremental_test(_Config) ->
+ rabbit_ct_proper_helpers:run_proper(fun prop_last_two_incremental/0, [], 100).
+
+prop_last_two_incremental() ->
+ ?FORALL(
+ Elements, non_empty(elements_gen()),
+ begin
+ Interval = 1,
+ Incremental = true,
+ {_LastTS, Slide} = new_slide(Interval, Incremental, Elements),
+ [{_Timestamp, Values} | _] = exometer_slide:last_two(Slide),
+ Expected = add_elements(Elements),
+ ?WHENFAIL(io:format("Expected a total of: ~p~nGot: ~p~n"
+ "Slide: ~p~n", [Expected, Values, Slide]),
+ Values == Expected)
+ end).
+
+sum_incremental_test(_Config) ->
+ rabbit_ct_proper_helpers:run_proper(fun prop_sum/1, [true], 100).
+
+sum_test(_Config) ->
+ rabbit_ct_proper_helpers:run_proper(fun prop_sum/1, [false], 100).
+
+prop_sum(Inc) ->
+ ?FORALL(
+ {Elements, Number}, {non_empty(elements_gen()), ?SUCHTHAT(I, int(), I > 0)},
+ begin
+ Interval = 1,
+ {LastTS, Slide} = new_slide(Interval, Inc, Elements),
+ %% Add the same so the timestamp matches. As the timestamps are handled
+ %% internally, we cannot guarantee on which interval they go otherwise
+ %% (unless we manually manipulate the slide content).
+ Sum = exometer_slide:sum([Slide || _ <- lists:seq(1, Number)]),
+ Values = [V || {_TS, V} <- exometer_slide:to_list(LastTS + 1, Sum)],
+ Expected = expected_sum(Slide, LastTS + 1, Number, Interval, Inc),
+ ?WHENFAIL(io:format("Expected: ~p~nGot: ~p~nSlide:~p~n",
+ [Expected, Values, Slide]),
+ Values == Expected)
+ end).
+
+expected_sum(Slide, Now, Number, _Int, false) ->
+ [sum_n_times(V, Number) || {_TS, V} <- exometer_slide:to_list(Now, Slide)];
+expected_sum(Slide, Now, Number, Int, true) ->
+ [{TSfirst, First} = F | Rest] = All = exometer_slide:to_list(Now, Slide),
+ {TSlast, _Last} = case Rest of
+ [] ->
+ F;
+ _ ->
+ lists:last(Rest)
+ end,
+ Seq = lists:seq(TSfirst, TSlast, Int),
+ {Expected, _} = lists:foldl(fun(TS0, {Acc, Previous}) ->
+ Actual = proplists:get_value(TS0, All, Previous),
+ {[sum_n_times(Actual, Number) | Acc], Actual}
+ end, {[], First}, Seq),
+ lists:reverse(Expected).
+%% -------------------------------------------------------------------
+%% Helpers
+%% -------------------------------------------------------------------
+new_slide(Interval, Incremental, Elements) ->
+ new_slide(Interval, Interval, Incremental, Elements).
+
+new_slide(PublishInterval, Interval, Incremental, Elements) ->
+ Now = 0,
+ Slide = exometer_slide:new(Now, 60 * 1000, [{interval, Interval},
+ {incremental, Incremental}]),
+ lists:foldl(
+ fun(E, {TS0, Acc}) ->
+ TS1 = TS0 + PublishInterval,
+ {TS1, exometer_slide:add_element(TS1, E, Acc)}
+ end, {Now, Slide}, Elements).
+
+last_two(Elements) when length(Elements) >= 2 ->
+ [F, S | _] = lists:reverse(Elements),
+ [F, S];
+last_two(Elements) ->
+ Elements.
+
+add_elements([H | T]) ->
+ add_elements(T, H).
+
+add_elements([], Acc) ->
+ Acc;
+add_elements([Tuple | T], Acc) ->
+ add_elements(T, sum(Tuple, Acc)).
+
+sum(T1, T2) ->
+ list_to_tuple(lists:zipwith(fun(A, B) -> A + B end, tuple_to_list(T1), tuple_to_list(T2))).
+
+sum_n_times(V, N) ->
+ sum_n_times(V, V, N - 1).
+
+sum_n_times(_V, Acc, 0) ->
+ Acc;
+sum_n_times(V, Acc, N) ->
+ sum_n_times(V, sum(V, Acc), N-1).