Skip to content

Commit 810e6db

Browse files
committed
QQ: annotate messages with the raft index they were written at
Add an "x-opt-index" message annotation to every message read from a quorum queue - via basic.get, normal delivery, or dead-lettering - containing the raft log index the message currently occupies. This lets a client identify exactly which point in the log a message came from, e.g. to correlate it with the queue's own operator-facing raft index.
1 parent 4409b16 commit 810e6db

3 files changed

Lines changed: 138 additions & 20 deletions

File tree

deps/rabbit/src/rabbit_fifo.erl

Lines changed: 32 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -1430,7 +1430,8 @@ handle_aux(_, _, {get_checked_out, ConsumerKey, MsgIds}, Aux0, RaAux0) ->
14301430
%% crashed and the message got removed
14311431
case ra_aux:log_fetch(Idx, S0) of
14321432
{{_Term, _Meta, Cmd}, S} ->
1433-
{S, [{MsgId, {Header, get_msg_from_cmd(Cmd)}} | Acc]};
1433+
RawMsg = annotate_index(Idx, get_msg_from_cmd(Cmd)),
1434+
{S, [{MsgId, {Header, RawMsg}} | Acc]};
14341435
{undefined, S} ->
14351436
{S, Acc}
14361437
end
@@ -2331,25 +2332,36 @@ get_header(Key, Header)
23312332
annotate_msg(Header, Msg0) ->
23322333
case mc:is(Msg0) of
23332334
true when is_map(Header) ->
2334-
Msg1 = maps:fold(fun (K, V, Acc) ->
2335-
mc:set_annotation(K, V, Acc)
2336-
end, Msg0, maps:get(anns, Header, #{})),
2337-
Msg = case Header of
2335+
Msg1 = maps:fold(fun mc:set_annotation/3, Msg0,
2336+
maps:get(anns, Header, #{})),
2337+
Msg2 = case Header of
23382338
#{acquired_count := AcqCount} ->
23392339
mc:set_annotation(acquired_count, AcqCount, Msg1);
23402340
_ ->
23412341
Msg1
23422342
end,
23432343
case Header of
23442344
#{delivery_count := DelCount} ->
2345-
mc:set_annotation(delivery_count, DelCount, Msg);
2345+
mc:set_annotation(delivery_count, DelCount, Msg2);
23462346
_ ->
2347-
Msg
2347+
Msg2
23482348
end;
23492349
_ ->
23502350
Msg0
23512351
end.
23522352

2353+
%% Stamps the raft index a message currently occupies onto the message
2354+
%% itself. Applied server-side (never over the wire) so that clients -
2355+
%% including ones running an older release that predates this annotation -
2356+
%% don't need to understand any new wire format to benefit from it.
2357+
annotate_index(Idx, Msg0) ->
2358+
case mc:is(Msg0) of
2359+
true ->
2360+
mc:set_annotation(<<"x-opt-index">>, Idx, Msg0);
2361+
false ->
2362+
Msg0
2363+
end.
2364+
23532365
return_one(#{system_time := Ts} = Meta, MsgId,
23542366
?C_MSG(Msg0), DeliveryFailed, Anns,
23552367
#?STATE{returns = Returns,
@@ -2798,19 +2810,18 @@ peek_next_msg(#?STATE{returns = Returns0,
27982810
delivery_effect(ConsumerKey, [{MsgId, ?MSG(Idx, Header)}],
27992811
#?STATE{msg_cache = {Idx, RawMsg}} = State) ->
28002812
{CTag, CPid} = consumer_id(ConsumerKey, State),
2801-
{send_msg, CPid, {delivery, CTag, [{MsgId, {Header, RawMsg}}]},
2813+
{send_msg, CPid, {delivery, CTag, [{MsgId, {Header, annotate_index(Idx, RawMsg)}}]},
28022814
?DELIVERY_SEND_MSG_OPTS};
28032815
delivery_effect(ConsumerKey, [{MsgId, Msg}],
28042816
#?STATE{msg_cache = {Idx, RawMsg}} = State)
28052817
when is_integer(Msg) andalso ?PACKED_IDX(Msg) == Idx ->
28062818
Header = get_msg_header(Msg),
28072819
{CTag, CPid} = consumer_id(ConsumerKey, State),
2808-
{send_msg, CPid, {delivery, CTag, [{MsgId, {Header, RawMsg}}]},
2820+
{send_msg, CPid, {delivery, CTag, [{MsgId, {Header, annotate_index(Idx, RawMsg)}}]},
28092821
?DELIVERY_SEND_MSG_OPTS};
28102822
delivery_effect(ConsumerKey, Msgs, #?STATE{} = State) ->
28112823
{CTag, CPid} = consumer_id(ConsumerKey, State),
28122824
{RaftIdxs, _Num} = lists:foldr(fun ({_, Msg}, {Acc, N}) ->
2813-
28142825
{[get_msg_idx(Msg) | Acc], N+1}
28152826
end, {[], 0}, Msgs),
28162827
{log_ext, RaftIdxs,
@@ -2837,9 +2848,10 @@ reply_log_effect(RaftIdx, MsgId, Header, Ready, From) ->
28372848
fun ([]) ->
28382849
[];
28392850
([Cmd]) ->
2851+
RawMsg = annotate_index(RaftIdx, get_msg_from_cmd(Cmd)),
28402852
[{reply, From,
28412853
{wrap_reply,
2842-
{dequeue, {MsgId, {Header, get_msg_from_cmd(Cmd)}}, Ready}}}]
2854+
{dequeue, {MsgId, {Header, RawMsg}}, Ready}}}]
28432855
end}.
28442856

28452857
checkout_one(#{system_time := Ts} = Meta, ExpiredMsg0, InitState0, Effects0) ->
@@ -3931,7 +3943,8 @@ exec_read(Flru0, ReadPlan, Msgs) ->
39313943
Idx = get_msg_idx(Msg),
39323944
Header = get_msg_header(Msg),
39333945
Cmd = maps:get(Idx, Entries),
3934-
{MsgId, {Header, get_msg_from_cmd(Cmd)}}
3946+
RawMsg = annotate_index(Idx, get_msg_from_cmd(Cmd)),
3947+
{MsgId, {Header, RawMsg}}
39353948
end, Msgs), Flru}
39363949
catch exit:{missing_key, _}
39373950
when Flru0 =/= undefined ->
@@ -4209,7 +4222,8 @@ discard_or_dead_letter(Msgs0, Reason, {at_most_once, {Mod, Fun, Args}}, State) -
42094222
Cmd = maps:get(Idx, Lookup),
42104223
%% ensure header delivery count
42114224
%% is copied to the message container
4212-
annotate_msg(Hdr, rabbit_fifo:get_msg_from_cmd(Cmd))
4225+
annotate_index(Idx,
4226+
annotate_msg(Hdr, rabbit_fifo:get_msg_from_cmd(Cmd)))
42134227
end || Msg <- Msgs0],
42144228
[{mod_call, Mod, Fun, Args ++ [Reason, Msgs]}]
42154229
end},
@@ -4274,14 +4288,15 @@ dlx_delivery_effects(_CPid, []) ->
42744288
[];
42754289
dlx_delivery_effects(CPid, Msgs0) ->
42764290
Msgs1 = lists:reverse(Msgs0),
4277-
{RaftIdxs, RsnIds} = lists:unzip(Msgs1),
4291+
{RaftIdxs, _RsnIds} = lists:unzip(Msgs1),
42784292
[{log, RaftIdxs,
42794293
fun(Log) ->
42804294
Msgs = lists:zipwith(
4281-
fun (Cmd, {Reason, H, MsgId}) ->
4295+
fun (Cmd, {Idx, {Reason, H, MsgId}}) ->
42824296
{MsgId, {Reason,
4283-
annotate_msg(H, rabbit_fifo:get_msg_from_cmd(Cmd))}}
4284-
end, Log, RsnIds),
4297+
annotate_index(Idx,
4298+
annotate_msg(H, rabbit_fifo:get_msg_from_cmd(Cmd)))}}
4299+
end, Log, Msgs1),
42854300
[{send_msg, CPid, {dlx_event, self(), {dlx_delivery, Msgs}}, [cast]}]
42864301
end}].
42874302

deps/rabbit/test/amqp_client_SUITE.erl

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -94,6 +94,7 @@ groups() ->
9494
sync_get_settled_classic_queue,
9595
sync_get_settled_quorum_queue,
9696
sync_get_settled_stream,
97+
x_opt_index_quorum_queue,
9798
timed_get_classic_queue,
9899
timed_get_quorum_queue,
99100
timed_get_stream,
@@ -2775,6 +2776,31 @@ sync_get_settled(QType, Config) ->
27752776
timed_get_classic_queue(Config) ->
27762777
timed_get(<<"classic">>, Config).
27772778

2779+
%% Every message carries the "x-opt-index" message annotation, the raft
2780+
%% index the message was written at.
2781+
x_opt_index_quorum_queue(Config) ->
2782+
QName = atom_to_binary(?FUNCTION_NAME),
2783+
{_Conn0, Ch} = rabbit_ct_client_helpers:open_connection_and_channel(Config),
2784+
#'queue.declare_ok'{} = amqp_channel:call(
2785+
Ch, #'queue.declare'{
2786+
queue = QName,
2787+
durable = true,
2788+
arguments = [{<<"x-queue-type">>, longstr, <<"quorum">>}]}),
2789+
OpnConf = connection_config(Config),
2790+
{ok, Connection} = amqp10_client:open_connection(OpnConf),
2791+
{ok, Session} = amqp10_client:begin_session_sync(Connection),
2792+
Address = rabbitmq_amqp_address:queue(QName),
2793+
{ok, Sender} = amqp10_client:attach_sender_link(Session, <<"test-sender">>, Address),
2794+
ok = wait_for_credit(Sender),
2795+
ok = amqp10_client:send_msg(Sender, amqp10_msg:new(<<"tag1">>, <<"m1">>, true)),
2796+
{ok, Receiver} = amqp10_client:attach_receiver_link(
2797+
Session, <<"test-receiver">>, Address, unsettled),
2798+
{ok, Msg} = amqp10_client:get_msg(Receiver),
2799+
?assertMatch(#{<<"x-opt-index">> := V} when is_integer(V),
2800+
amqp10_msg:message_annotations(Msg)),
2801+
ok = amqp10_client:settle_msg(Receiver, Msg, accepted),
2802+
ok = amqp10_client:close_connection(Connection).
2803+
27782804
timed_get_quorum_queue(Config) ->
27792805
timed_get(<<"quorum">>, Config).
27802806

deps/rabbit/test/quorum_queue_SUITE.erl

Lines changed: 80 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -91,6 +91,7 @@ groups() ->
9191
consume_in_minority,
9292
get_in_minority,
9393
reject_after_leader_transfer,
94+
x_opt_index_annotation,
9495
delete_members,
9596
rebalance,
9697
node_removal_is_not_quorum_critical,
@@ -2494,6 +2495,74 @@ reject_after_leader_transfer(Config) ->
24942495
requeue = true}),
24952496
ok.
24962497

2498+
%% Every message carries the "x-opt-index" annotation, the raft index the
2499+
%% message was written at. Check that it is present for basic.get, for
2500+
%% deliveries to a consumer on a node with a member, and for deliveries to a
2501+
%% consumer on a node without a member (i.e. read remotely via log_ext).
2502+
x_opt_index_annotation(Config) ->
2503+
check_quorum_queues_v9_compat(Config),
2504+
2505+
[Server0, Server1, Server2] =
2506+
rabbit_ct_broker_helpers:get_node_configs(Config, nodename),
2507+
QQ = ?config(queue_name, Config),
2508+
RaName = ra_name(QQ),
2509+
IndexHeader = <<"x-opt-index">>,
2510+
2511+
Ch0 = rabbit_ct_client_helpers:open_channel(Config, Server0),
2512+
?assertEqual({'queue.declare_ok', QQ, 0, 0},
2513+
declare(Ch0, QQ, [{<<"x-queue-type">>, longstr, <<"quorum">>}])),
2514+
?awaitMatch(3, count_online_nodes(Server0, <<"/">>, QQ), ?DEFAULT_AWAIT),
2515+
2516+
%% basic.get
2517+
publish(Ch0, QQ),
2518+
wait_for_messages_ready([Server0, Server1, Server2], RaName, 1),
2519+
{#'basic.get_ok'{delivery_tag = GetTag},
2520+
#amqp_msg{props = #'P_basic'{headers = GetHeaders}}} =
2521+
basic_get(Ch0, QQ, false, 10),
2522+
?assertMatch({IndexHeader, long, _},
2523+
rabbit_basic:header(IndexHeader, GetHeaders)),
2524+
ok = amqp_channel:call(Ch0, #'basic.ack'{delivery_tag = GetTag}),
2525+
2526+
%% delivery to a consumer on a node that has a member
2527+
publish(Ch0, QQ),
2528+
wait_for_messages_ready([Server0, Server1, Server2], RaName, 1),
2529+
Ch1 = rabbit_ct_client_helpers:open_channel(Config, Server1),
2530+
subscribe(Ch1, QQ, false),
2531+
receive
2532+
{#'basic.deliver'{delivery_tag = LocalTag},
2533+
#amqp_msg{props = #'P_basic'{headers = LocalHeaders}}} ->
2534+
?assertMatch({IndexHeader, long, _},
2535+
rabbit_basic:header(IndexHeader, LocalHeaders)),
2536+
amqp_channel:cast(Ch1, #'basic.ack'{delivery_tag = LocalTag})
2537+
after ?TIMEOUT ->
2538+
flush(10),
2539+
exit(basic_deliver_timeout)
2540+
end,
2541+
amqp_channel:close(Ch1),
2542+
2543+
%% remove Server2's member so it no longer holds a copy of the raft log
2544+
?assertEqual(ok,
2545+
rpc:call(Server0, rabbit_queue_type_ra, delete_member,
2546+
[<<"/">>, QQ, Server2])),
2547+
?awaitMatch(2, count_online_nodes(Server0, <<"/">>, QQ), ?DEFAULT_AWAIT),
2548+
2549+
%% delivery to a consumer on a node that does not have a member
2550+
publish(Ch0, QQ),
2551+
wait_for_messages_ready([Server0, Server1], RaName, 1),
2552+
Ch2 = rabbit_ct_client_helpers:open_channel(Config, Server2),
2553+
subscribe(Ch2, QQ, false),
2554+
receive
2555+
{#'basic.deliver'{delivery_tag = RemoteTag},
2556+
#amqp_msg{props = #'P_basic'{headers = RemoteHeaders}}} ->
2557+
?assertMatch({IndexHeader, long, _},
2558+
rabbit_basic:header(IndexHeader, RemoteHeaders)),
2559+
amqp_channel:cast(Ch2, #'basic.ack'{delivery_tag = RemoteTag})
2560+
after ?TIMEOUT ->
2561+
flush(10),
2562+
exit(basic_deliver_timeout)
2563+
end,
2564+
ok.
2565+
24972566
delete_members(Config) ->
24982567
[Server0, Server1, Server2] =
24992568
rabbit_ct_broker_helpers:get_node_configs(Config, nodename),
@@ -5436,8 +5505,10 @@ amqpl_headers(Config) ->
54365505
#amqp_msg{props = #'P_basic'{headers = Headers2Received}}
54375506
} = amqp_channel:call(Ch, #'basic.get'{queue = QQ}),
54385507

5439-
?assertEqual(Headers1Sent, Headers1Received),
5440-
?assertEqual(Headers2Sent, Headers2Received),
5508+
%% every message carries the raft index it was written at as the
5509+
%% "x-opt-index" header, even when the client didn't set any headers
5510+
?assertMatch([{<<"x-opt-index">>, long, _}], Headers1Received),
5511+
?assertMatch([{<<"x-opt-index">>, long, _}], Headers2Received),
54415512

54425513
ok = amqp_channel:cast(Ch, #'basic.ack'{delivery_tag = DeliveryTag,
54435514
multiple = true}).
@@ -6631,9 +6702,15 @@ basic_get(Ch, Q, NoAck, Attempt) ->
66316702
end.
66326703

66336704
check_quorum_queues_v8_compat(Config) ->
6705+
check_quorum_queues_vn_compat(8, Config).
6706+
6707+
check_quorum_queues_v9_compat(Config) ->
6708+
check_quorum_queues_vn_compat(9, Config).
6709+
6710+
check_quorum_queues_vn_compat(Version, Config) ->
66346711
Nodes = rabbit_ct_broker_helpers:get_node_configs(Config, nodename),
66356712
MacVer = lists:min([V || {ok, V} <- erpc:multicall(Nodes, rabbit_fifo, version, [])]),
6636-
case MacVer >= 8 of
6713+
case MacVer >= Version of
66376714
true ->
66386715
ok;
66396716
false ->

0 commit comments

Comments
 (0)