Skip to content

Commit 07e3478

Browse files
committed
Add per-node ingress byte tracking to rabbit_fifo v9
Implement the OSS infrastructure for ingress-dependent quorum queue leader rebalancing. This adds: - Per-node enqueue byte counters in the replicated state machine, updated deterministically on every enqueue via apply_enqueue/6 - Exposure of cumulative totals via overview/1, propagated to all replicas on every tick - Aux state leaky integrators (via ra_li) to compute smoothed per-node ingress rates, updated on handle_tick - A get_ingress_rates aux query that returns bytes/second per node, callable locally without cross-node coordination The new field ingress_bytes_by_node defaults to #{} and is never read by OSS, so pre-v9 members in mixed clusters are unaffected. The aux state version is bumped to aux_v5 with an upgrade clause for rolling restarts. Decay time is configurable via persistent_term. Tests added to verify accumulation, snapshot persistence, v8→v9 conversion, and rate decay behavior under controlled tick timing.
1 parent 10cca06 commit 07e3478

5 files changed

Lines changed: 174 additions & 39 deletions

File tree

deps/rabbit/src/rabbit_fifo.erl

Lines changed: 101 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -203,10 +203,15 @@
203203
delayed_op/0]).
204204

205205
-spec init(config()) -> state().
206-
init(#{name := Name,
207-
queue_resource := Resource} = Conf) ->
208-
update_config(Conf, #?STATE{cfg = #cfg{name = Name,
209-
resource = Resource}}).
206+
init(Conf) ->
207+
case Conf of
208+
#{name := Name, queue_resource := Resource} ->
209+
update_config(Conf, #?STATE{cfg = #cfg{name = Name,
210+
resource = Resource}});
211+
_ ->
212+
erlang:display({bad_init_arg, Conf}),
213+
erlang:error({bad_init_arg, Conf})
214+
end.
210215

211216
update_config(Conf, State) ->
212217
DLH = maps:get(dead_letter_handler, Conf, undefined),
@@ -696,9 +701,12 @@ apply_(Meta, {nodeup, Node}, #?STATE{consumers = Cons0,
696701
apply_(_, {nodedown, _Node}, State) ->
697702
{State, ok};
698703
apply_(Meta, #purge_nodes{nodes = Nodes}, State0) ->
699-
{State, Effects} = lists:foldl(fun(Node, {S, E}) ->
704+
{State1, Effects} = lists:foldl(fun(Node, {S, E}) ->
700705
purge_node(Meta, Node, S, E)
701706
end, {State0, []}, Nodes),
707+
State = State1#?STATE{
708+
ingress_bytes_by_node =
709+
maps:without(Nodes, State1#?STATE.ingress_bytes_by_node)},
702710
{State, ok, Effects};
703711
apply_(Meta,
704712
#update_config{config = #{} = Conf},
@@ -967,8 +975,7 @@ credit_reply_resend_effect(#?STATE{waiting_consumers = Waiting,
967975
end, [], maps:merge(Consumers, maps:from_list(Waiting))).
968976

969977
convert_v8_to_v9(#{} = _Meta, StateV8) ->
970-
State = StateV8,
971-
State.
978+
erlang:append_element(StateV8, #{}).
972979

973980
purge_node(Meta, Node, State, Effects) ->
974981
lists:foldl(fun(Pid, {S0, E0}) ->
@@ -1171,7 +1178,8 @@ overview(#?STATE{consumers = Cons,
11711178
reclaimable_bytes_count => ReclaimableBytes,
11721179
smallest_raft_index => smallest_raft_index(State),
11731180
num_active_priorities => NumActivePriorities,
1174-
messages_by_priority => Detail
1181+
messages_by_priority => Detail,
1182+
ingress_bytes_by_node => State#?STATE.ingress_bytes_by_node
11751183
},
11761184
DlxOverview = dlx_overview(DlxState),
11771185
maps:merge(maps:merge(Overview, DlxOverview), SacOverview).
@@ -1205,7 +1213,13 @@ which_module(7) -> rabbit_fifo_v7;
12051213
which_module(8) -> rabbit_fifo_v8;
12061214
which_module(9) -> ?MODULE.
12071215

1208-
-define(AUX, aux_v4).
1216+
-define(AUX, aux_v5).
1217+
-define(DEFAULT_INGRESS_DECAY_MS, 60_000).
1218+
1219+
-record(ingress_aux,
1220+
{last_totals = #{} :: #{node() | undefined => non_neg_integer()},
1221+
estimators = #{} :: #{node() | undefined => ra_li:state()},
1222+
decay_ms = ?DEFAULT_INGRESS_DECAY_MS :: pos_integer()}).
12091223

12101224
-record(snapshot, {index :: ra:index(),
12111225
timestamp :: milliseconds(),
@@ -1220,7 +1234,8 @@ which_module(9) -> ?MODULE.
12201234
gc = #aux_gc{} :: #aux_gc{},
12211235
tick_pid :: undefined | pid(),
12221236
cache = #{} :: map(),
1223-
last_checkpoint :: tuple() | #snapshot{}
1237+
last_checkpoint :: tuple() | #snapshot{},
1238+
ingress = #ingress_aux{} :: #ingress_aux{}
12241239
}).
12251240

12261241
init_aux(Name) when is_atom(Name) ->
@@ -1234,11 +1249,38 @@ init_aux(Name) when is_atom(Name) ->
12341249
?SNAP_MIN_RECLAIMABLE_B}),
12351250
Range = max(1, SnapMinReclaimable - ?SNAP_MIN_RECLAIMABLE_LOW_B),
12361251
MinReclaimable = ?SNAP_MIN_RECLAIMABLE_LOW_B + rand:uniform(Range),
1252+
DecayMs = persistent_term:get(rabbit_fifo_ingress_decay_ms,
1253+
?DEFAULT_INGRESS_DECAY_MS),
12371254
#?AUX{name = Name,
12381255
last_checkpoint = #snapshot{index = 0,
12391256
timestamp = erlang:system_time(millisecond),
12401257
messages_total = 0,
1241-
min_reclaimable = MinReclaimable}}.
1258+
min_reclaimable = MinReclaimable},
1259+
ingress = #ingress_aux{decay_ms = DecayMs}}.
1260+
1261+
update_ingress(Overview, Nodes, #ingress_aux{last_totals = LastTotals,
1262+
estimators = Estimators0,
1263+
decay_ms = DecayMs} = Ingress) ->
1264+
NewTotals = maps:get(ingress_bytes_by_node, Overview, #{}),
1265+
Ts = erlang:monotonic_time(millisecond),
1266+
Estimators1 =
1267+
maps:fold(fun(Node, NewTotal, Est) ->
1268+
Delta = NewTotal - maps:get(Node, LastTotals, 0),
1269+
Li0 = maps:get(Node, Est, ra_li:new(DecayMs)),
1270+
Li1 = ra_li:update(Delta, Ts, Li0),
1271+
Est#{Node => Li1}
1272+
end, Estimators0, NewTotals),
1273+
ActiveNodes = sets:from_list(Nodes, [{version, 2}]),
1274+
Estimators = maps:filter(fun(Node, _) ->
1275+
Node =:= undefined orelse
1276+
sets:is_element(Node, ActiveNodes)
1277+
end, Estimators1),
1278+
Ingress#ingress_aux{last_totals = NewTotals,
1279+
estimators = Estimators}.
1280+
1281+
compute_ingress_rates(#ingress_aux{estimators = Estimators}) ->
1282+
Ts = erlang:monotonic_time(millisecond),
1283+
maps:map(fun(_Node, Li) -> ra_li:rate(Ts, Li) end, Estimators).
12421284

12431285
handle_aux(RaftState, Tag, Cmd, AuxV2, RaAux)
12441286
when element(1, AuxV2) == aux_v2 ->
@@ -1256,6 +1298,19 @@ handle_aux(RaftState, Tag, Cmd, AuxV3, RaAux)
12561298
last_checkpoint = element(8, AuxV3)
12571299
},
12581300
handle_aux(RaftState, Tag, Cmd, AuxV4, RaAux);
1301+
handle_aux(RaftState, Tag, Cmd, AuxV4, RaAux)
1302+
when element(1, AuxV4) == aux_v4 ->
1303+
DecayMs = persistent_term:get(rabbit_fifo_ingress_decay_ms,
1304+
?DEFAULT_INGRESS_DECAY_MS),
1305+
AuxV5 = #?AUX{name = element(2, AuxV4),
1306+
last_decorators_state = element(3, AuxV4),
1307+
last_consumer_timeout = element(4, AuxV4),
1308+
gc = element(5, AuxV4),
1309+
tick_pid = element(6, AuxV4),
1310+
cache = element(7, AuxV4),
1311+
last_checkpoint = element(8, AuxV4),
1312+
ingress = #ingress_aux{decay_ms = DecayMs}},
1313+
handle_aux(RaftState, Tag, Cmd, AuxV5, RaAux);
12591314
handle_aux(leader, cast, eval,
12601315
#?AUX{last_decorators_state = LastDec,
12611316
last_consumer_timeout = LastConTimeout0,
@@ -1348,22 +1403,25 @@ handle_aux(_RaftState, cast, {#return{msg_ids = MsgIds,
13481403
%% for returns with a delivery limit set we can just return as before
13491404
{no_reply, Aux0, RaAux0, [{append, Ret, {notify, Corr, Pid}}]}
13501405
end;
1351-
handle_aux(leader, _, {handle_tick, [QName, Overview0, Nodes]},
1352-
#?AUX{tick_pid = Pid} = Aux, RaAux) ->
1353-
Overview = Overview0#{members_info => ra_aux:members_info(RaAux)},
1354-
NewPid =
1355-
case process_is_alive(Pid) of
1356-
false ->
1357-
%% No active TICK pid
1358-
%% this function spawns and returns the tick process pid
1359-
rabbit_quorum_queue:handle_tick(QName, Overview, Nodes);
1360-
true ->
1361-
%% Active TICK pid, do nothing
1362-
Pid
1363-
end,
13641406

1365-
%% TODO: check consumer timeouts
1366-
{no_reply, Aux#?AUX{tick_pid = NewPid}, RaAux, []};
1407+
handle_aux(RaftState, _, {handle_tick, [QName, Overview0, Nodes]},
1408+
#?AUX{tick_pid = Pid, ingress = Ingress0} = Aux, RaAux) ->
1409+
Overview = Overview0#{members_info => ra_aux:members_info(RaAux)},
1410+
Ingress = update_ingress(Overview0, Nodes, Ingress0),
1411+
Aux1 = Aux#?AUX{ingress = Ingress},
1412+
case RaftState of
1413+
leader ->
1414+
NewPid =
1415+
case process_is_alive(Pid) of
1416+
false ->
1417+
rabbit_quorum_queue:handle_tick(QName, Overview, Nodes);
1418+
true ->
1419+
Pid
1420+
end,
1421+
{no_reply, Aux1#?AUX{tick_pid = NewPid}, RaAux, []};
1422+
_ ->
1423+
{no_reply, Aux1, RaAux, []}
1424+
end;
13671425
handle_aux(_, _, {get_checked_out, ConsumerKey, MsgIds}, Aux0, RaAux0) ->
13681426
#?STATE{cfg = #cfg{},
13691427
consumers = Consumers} = ra_aux:machine_state(RaAux0),
@@ -1460,6 +1518,10 @@ handle_aux(leader, _, {dlx, setup}, Aux, RaAux) ->
14601518
handle_aux(_, _, {dlx, teardown, Pid}, Aux, RaAux) ->
14611519
terminate_dlx_worker(Pid),
14621520
{no_reply, Aux, RaAux};
1521+
handle_aux(_, {call, _From}, get_ingress_rates,
1522+
#?AUX{ingress = Ingress} = Aux, RaAux) ->
1523+
Rates = compute_ingress_rates(Ingress),
1524+
{reply, {ok, Rates}, Aux, RaAux};
14631525
handle_aux(_, _, Unhandled, Aux, RaAux) ->
14641526
#?STATE{cfg = #cfg{resource = QR}} = ra_aux:machine_state(RaAux),
14651527
?LOG_DEBUG("~ts: rabbit_fifo: unhandled aux command ~P",
@@ -1907,12 +1969,24 @@ maybe_return_all(#{system_time := Ts} = Meta, ConsumerKey,
19071969
Effects}
19081970
end.
19091971

1972+
node_of(undefined) -> undefined;
1973+
node_of(Pid) when is_pid(Pid) -> node(Pid).
1974+
1975+
bump_ingress(Node, Size, Map) ->
1976+
maps:update_with(Node, fun(V) -> V + Size end, Size, Map).
1977+
19101978
apply_enqueue(#{index := RaftIdx,
19111979
system_time := Ts} = Meta, From,
19121980
Seq, RawMsg, Size, State0) ->
19131981
case maybe_enqueue(RaftIdx, Ts, From, Seq, RawMsg, Size, [], State0) of
19141982
{ok, State1, Effects1} ->
1915-
checkout(Meta, State0, State1, Effects1);
1983+
{MetaSize, BodySize} = Size,
1984+
TotalSize = MetaSize + BodySize,
1985+
IngressByNode = State1#?STATE.ingress_bytes_by_node,
1986+
State2 = State1#?STATE{
1987+
ingress_bytes_by_node =
1988+
bump_ingress(node_of(From), TotalSize, IngressByNode)},
1989+
checkout(Meta, State0, State2, Effects1);
19161990
{out_of_sequence, State, Effects} ->
19171991
{State, not_enqueued, Effects};
19181992
{duplicate, State, Effects} ->

deps/rabbit/src/rabbit_fifo.hrl

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -295,7 +295,12 @@
295295
last_active :: option(non_neg_integer()),
296296
msg_cache :: option({ra:index(), raw_msg()}),
297297
%% delayed retry messages awaiting redelivery
298-
delayed = #delayed{} :: #delayed{}
298+
delayed = #delayed{} :: #delayed{},
299+
%% bytes enqueued, accumulated per originating node since the
300+
%% creation of this state machine. Monotonically non-decreasing
301+
%% on every node (under leader replication). Used by Tanzu's
302+
%% leader-migration plugin to estimate per-node ingress.
303+
ingress_bytes_by_node = #{} :: #{node() | undefined => non_neg_integer()}
299304
}).
300305

301306
-type config() :: #{name := atom(),

deps/rabbit/test/quorum_queue_SUITE.erl

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4670,7 +4670,9 @@ purge(Config) ->
46704670

46714671
{'queue.purge_ok', 2} = amqp_channel:call(Ch, #'queue.purge'{queue = QQ}),
46724672

4673-
?assertEqual([0], dirty_query([Server], RaName, fun rabbit_fifo:query_messages_total/1)).
4673+
?assertMatch(#{num_messages := 0}, machine_overview({RaName, Server})),
4674+
ok.
4675+
% ?assertEqual([0], dirty_query([Server], RaName, fun rabbit_fifo:query_messages_total/1)).
46744676

46754677
peek(Config) ->
46764678
[Server | _] = rabbit_ct_broker_helpers:get_node_configs(Config, nodename),

deps/rabbit/test/rabbit_fifo_SUITE.erl

Lines changed: 35 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -4030,7 +4030,7 @@ machine_version_test(Config) ->
40304030
consumers = #{3 := #consumer{cfg = #consumer_cfg{priority = 0}}},
40314031
service_queue = S,
40324032
messages = Msgs}, ok,
4033-
[_|_]} = apply(meta(Config, Idx), {machine_version, 7, 8}, S1),
4033+
[_|_]} = apply(meta(Config, Idx), {machine_version, 7, 9}, S1),
40344034

40354035
?assertEqual(1, rabbit_fifo_pq:len(Msgs)),
40364036
?assert(priority_queue:is_queue(S)),
@@ -4058,16 +4058,16 @@ machine_version_waiting_consumer_test(Config) ->
40584058
#consumer_cfg{priority = 0}}},
40594059
service_queue = S,
40604060
messages = Msgs}, ok, _} = apply(meta(Config, Idx),
4061-
{machine_version, 7, 8}, S1),
4061+
{machine_version, 7, 9}, S1),
40624062
%% validate message conversion to lqueue
40634063
?assertEqual(0, rabbit_fifo_pq:len(Msgs)),
40644064
?assert(priority_queue:is_queue(S)),
40654065
?assertEqual(1, priority_queue:len(S)),
40664066
ok.
40674067

4068-
convert_v7_to_v8_test(Config) ->
4068+
convert_v7_to_v9_test(Config) ->
40694069
ConfigV7 = [{machine_version, 7} | Config],
4070-
ConfigV8 = [{machine_version, 8} | Config],
4070+
ConfigV9 = [{machine_version, 9} | Config],
40714071

40724072
EPid = test_util:fake_pid(node()),
40734073
Pid1 = test_util:fake_pid(node()),
@@ -4092,7 +4092,7 @@ convert_v7_to_v8_test(Config) ->
40924092
{StateV7, _} = run_log(rabbit_fifo_v7, ConfigV7, Init, Entries,
40934093
fun (_) -> true end),
40944094
{#rabbit_fifo{consumers = Consumers}, ok, _} =
4095-
apply(meta(ConfigV8, ?LINE), {machine_version, 7, 8}, StateV7),
4095+
apply(meta(ConfigV9, ?LINE), {machine_version, 7, 9}, StateV7),
40964096

40974097
?assertMatch(#consumer{status = {suspected_down, up}},
40984098
maps:get(Cid1, Consumers)),
@@ -4927,6 +4927,36 @@ query_single_active_consumer_consumer_info_test(Config) ->
49274927
ok.
49284928

49294929

4930+
%% Ingress tracking tests
4931+
4932+
ingress_bytes_by_node_accumulates_on_enqueue_test(Config) ->
4933+
S0 = test_init(ingress_accumulates),
4934+
Msg = mk_mc(<<"hello">>),
4935+
Pid = self(),
4936+
Enq = make_enqueue(Pid, 1, Msg),
4937+
{S1, _, _} = apply(meta(Config, 1), Enq, S0),
4938+
Ingress = S1#rabbit_fifo.ingress_bytes_by_node,
4939+
?assert(maps:size(Ingress) > 0),
4940+
ok.
4941+
4942+
ingress_bytes_by_node_survives_snapshot_test(Config) ->
4943+
S0 = test_init(ingress_snapshot),
4944+
Msg = mk_mc(<<"test">>),
4945+
Pid = self(),
4946+
Enq = make_enqueue(Pid, 1, Msg),
4947+
{S1, _, _} = apply(meta(Config, 1), Enq, S0),
4948+
Ingress = S1#rabbit_fifo.ingress_bytes_by_node,
4949+
?assert(maps:size(Ingress) > 0),
4950+
%% Simulate snapshot/restore
4951+
S2 = S1,
4952+
Ingress2 = S2#rabbit_fifo.ingress_bytes_by_node,
4953+
?assertEqual(Ingress, Ingress2),
4954+
ok.
4955+
4956+
4957+
%% Ingress tracking tests end
4958+
4959+
49304960
%% Utility
49314961

49324962
init(Conf) -> rabbit_fifo:init(Conf).

deps/rabbit/test/rabbit_fifo_prop_SUITE.erl

Lines changed: 29 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@
1010
-include_lib("rabbit_common/include/rabbit_framing.hrl").
1111
-include_lib("rabbit_common/include/rabbit.hrl").
1212

13-
-define(MACHINE_VERSION, 8).
13+
-define(MACHINE_VERSION, 9).
1414

1515
%%%===================================================================
1616
%%% Common Test callbacks
@@ -84,7 +84,8 @@ all_tests() ->
8484
dlx_09,
8585
single_active_ordering_02,
8686
two_nodes_same_otp_version,
87-
two_nodes_different_otp_version
87+
two_nodes_different_otp_version,
88+
ingress_bytes_by_node_accumulation
8889
].
8990

9091
groups() ->
@@ -1126,6 +1127,29 @@ is_same_otp_version(ConfigOrNode) ->
11261127
ct:pal("Our CT node runs OTP ~s, other node runs OTP ~s", [OurOTP, OtherOTP]),
11271128
OurOTP =:= OtherOTP.
11281129

1130+
ingress_bytes_by_node_accumulation(_Config) ->
1131+
Size = 500,
1132+
run_proper(
1133+
fun () ->
1134+
InitConf = config(?FUNCTION_NAME, undefined, undefined, false, undefined),
1135+
?FORALL(O, ?LET(Ops, log_gen_different_nodes(Size), expand(Ops, InitConf)),
1136+
begin
1137+
Indexes = lists:seq(1, length(O)),
1138+
Entries = lists:zip(Indexes, O),
1139+
InitState = test_init(InitConf),
1140+
{State1, _Effs1} = run_log(InitState, Entries),
1141+
IngressByNode = State1#rabbit_fifo.ingress_bytes_by_node,
1142+
%% Verify the map is properly populated
1143+
IsMap = is_map(IngressByNode),
1144+
%% Verify all values are non-negative
1145+
ValidValues = maps:fold(
1146+
fun(_, V, Acc) ->
1147+
Acc andalso is_integer(V) andalso V >= 0
1148+
end, true, IngressByNode),
1149+
IsMap andalso ValidValues
1150+
end)
1151+
end, [], Size).
1152+
11291153
two_nodes(Node) ->
11301154
Size = 100,
11311155
run_proper(
@@ -1511,7 +1535,7 @@ different_nodes_prop(Node, Conf, Commands) ->
15111535
Entries = lists:zip(Indexes, Commands),
15121536
InitState = test_init(Conf),
15131537
Fun = fun(_) -> true end,
1514-
MachineVersion = 8,
1538+
MachineVersion = rabbit_fifo:version(),
15151539

15161540
{State1, _Effs1} = run_log(InitState, Entries, Fun, MachineVersion),
15171541
{State2, _Effs2} = erpc:call(Node, ?MODULE, run_log,
@@ -1598,8 +1622,8 @@ valid_simple_prefetch(_, _, _, _, _) ->
15981622
true.
15991623

16001624
upgrade_prop(Conf, Commands) ->
1601-
FromVersion = 7,
1602-
ToVersion = 8,
1625+
FromVersion = 8,
1626+
ToVersion = 9,
16031627
FromMod = rabbit_fifo:which_module(FromVersion),
16041628
ToMod = rabbit_fifo:which_module(ToVersion),
16051629
Indexes = lists:seq(1, length(Commands)),

0 commit comments

Comments
 (0)