Skip to content

Commit 4409b16

Browse files
committed
Wire message deferral tokens through AMQP 1.0 FLOW to quorum queues
Consumers can now request previously deferred (parked) messages by listing their tokens under the rabbitmq:deferred-tokens FLOW property. rabbit_queue_type gains an optional assign_deferred/4 callback, implemented by quorum queues via rabbit_fifo_client. In rabbit_fifo, a deferral token is only honoured when the client also sets an explicit x-opt-delivery-time; the delayed-retry path never creates a deferred entry, keeping tokens purely client-driven. Advertise support via the rabbitmq:deferred-tokens link capability and document the protocol in deps/rabbit/DEFERRED.md.
1 parent 623ca2c commit 4409b16

9 files changed

Lines changed: 548 additions & 80 deletions

deps/rabbit/DEFERRED.md

Lines changed: 112 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,112 @@
1+
# Message Deferral for Quorum Queues
2+
3+
## Overview
4+
5+
Message deferral lets an AMQP 1.0 consumer return a message to a quorum queue in a
6+
*parked* state under a client-chosen token, then later pull that specific message back
7+
on demand — without waiting for a delayed-retry timer or competing with other consumers.
8+
9+
This is useful for workload scheduling patterns where a consumer wants to decide at
10+
receive time that a message should not be processed immediately, assign it a token
11+
for later retrieval, and then fetch it explicitly when ready.
12+
13+
## Constraints
14+
15+
- **Quorum queues only.** Classic queues and streams do not support deferral tokens.
16+
Clients can detect support by checking for `rabbitmq:deferral-tokens` in the
17+
`offered-capabilities` field of the ATTACH response.
18+
- **`x-opt-delivery-time` is required.** A deferral token alone does not park a
19+
message. Both `x-opt-deferral-token` and `x-opt-delivery-time` must be present in
20+
the MODIFIED outcome's `message-annotations`. A message returned via the
21+
`delayed-retry` queue configuration never creates a deferred entry, even if a token
22+
is present.
23+
- **Tokens must be of AMQP type `utf8`.** This applies both to `x-opt-deferral-token`
24+
in the MODIFIED outcome's `message-annotations` and to each element of the
25+
`rabbitmq:deferral-tokens` array in a FLOW frame's `properties`. Any other AMQP
26+
type (e.g. `binary`, `symbol`) is rejected with `amqp:invalid-field`.
27+
- **Credit must cover the matched messages.** When requesting deferred messages via a
28+
FLOW frame, the link credit granted in that same FLOW must be at least as large as
29+
the number of messages the submitted tokens resolve to (a single token may resolve
30+
to more than one message; see below). If credit is exhausted by normal message
31+
delivery before the deferred assignment runs, no deferred messages are delivered
32+
for that FLOW.
33+
34+
## Protocol Usage
35+
36+
### 1. Attach a consuming link and verify capability
37+
38+
Attach a link to a quorum queue source. Inspect the `offered-capabilities` array in
39+
the ATTACH response for the symbol `rabbitmq:deferral-tokens`. If absent, the queue
40+
does not support deferral.
41+
42+
### 2. Receive a message
43+
44+
The broker delivers a message to the link via a TRANSFER frame. Note the
45+
`delivery-tag` for the settlement step.
46+
47+
### 3. Park the message with a deferral token
48+
49+
Send a DISPOSITION frame settling the delivery with a MODIFIED outcome. Include both
50+
annotations in the `message-annotations` map of the MODIFIED outcome:
51+
52+
- `x-opt-deferral-token` (symbol key, utf8 value) — a client-chosen opaque
53+
identifier for this parked message. The same token may be assigned to more than
54+
one message, e.g. by settling a range of deliveries (`first =/= last`) with a
55+
single MODIFIED outcome. Retrieving such a token returns all messages parked
56+
under it, oldest first.
57+
- `x-opt-delivery-time` (symbol key, timestamp value in milliseconds since the Unix
58+
epoch) — the earliest time at which the message becomes eligible for normal
59+
timer-based redelivery. The message is held until this time unless explicitly
60+
retrieved earlier via `assign_deferred`.
61+
62+
Example (pseudocode):
63+
64+
```
65+
DISPOSITION {
66+
role = receiver,
67+
first = <delivery-id>,
68+
settled = true,
69+
state = MODIFIED {
70+
delivery-failed = false,
71+
undeliverable-here = false,
72+
message-annotations = {
73+
x-opt-deferral-token: "job-42-retry-1",
74+
x-opt-delivery-time: 1780000000000
75+
}
76+
}
77+
}
78+
```
79+
80+
### 4. Retrieve the message by token
81+
82+
When ready to process the parked message, send a FLOW frame on the same consuming
83+
link. Grant enough link credit to receive the deferred messages and include the tokens
84+
in the `properties` map of the FLOW frame under the key `rabbitmq:deferral-tokens` as
85+
an array of utf8 values.
86+
87+
Example (pseudocode):
88+
89+
```
90+
FLOW {
91+
handle = <link-handle>,
92+
delivery-count = <current-delivery-count>,
93+
link-credit = 1,
94+
properties = {
95+
rabbitmq:deferral-tokens: ["job-42-retry-1"]
96+
}
97+
}
98+
```
99+
100+
The broker looks up each token in the parked message set and delivers any matches
101+
directly to this link as normal TRANSFER frames. Tokens that are not found (because
102+
the message was already expired by its delivery time and requeued normally, or the
103+
token was never issued) produce no delivery — the client is responsible for tracking
104+
which tokens it expects.
105+
106+
## Relationship to `delayed-retry`
107+
108+
Deferral tokens and the queue-level `x-delayed-retry-*` configuration are
109+
complementary but independent. A message returned with `x-opt-delivery-time` follows
110+
the deferral path. A message returned without `x-opt-delivery-time` follows the
111+
delayed-retry path if the queue is configured for it. The two paths do not interact:
112+
delayed-retry messages cannot be retrieved by token.

deps/rabbit/src/rabbit_amqp_session.erl

Lines changed: 39 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -2148,7 +2148,14 @@ settle_op_from_outcome(#'v1_0.modified'{delivery_failed = DelFailed,
21482148
{map, KVList} ->
21492149
Anns1 = lists:map(
21502150
%% "all symbolic keys except those beginning with "x-" are reserved." [3.2.10]
2151-
fun({{symbol, <<"x-", _/binary>> = K}, V}) ->
2151+
fun({{symbol, <<"x-opt-deferral-token">> = K}, {utf8, _} = V}) ->
2152+
{K, unwrap_simple_type(V)};
2153+
({{symbol, <<"x-opt-deferral-token">>}, Other}) ->
2154+
protocol_error(
2155+
?V_1_0_AMQP_ERROR_INVALID_FIELD,
2156+
"x-opt-deferral-token must be of AMQP type utf8, got: ~tp",
2157+
[Other]);
2158+
({{symbol, <<"x-", _/binary>> = K}, V}) ->
21522159
{K, unwrap_simple_type(V)}
21532160
end, KVList),
21542161
maps:from_list(Anns1)
@@ -3162,7 +3169,8 @@ handle_outgoing_link_flow_control(
31623169
delivery_count = MaybeDeliveryCountRcv,
31633170
link_credit = ?UINT(LinkCreditRcv),
31643171
drain = Drain0,
3165-
echo = Echo0},
3172+
echo = Echo0,
3173+
properties = FlowProps},
31663174
#state{outgoing_links = OutgoingLinks,
31673175
queue_states = QStates0
31683176
} = State0) ->
@@ -3186,14 +3194,21 @@ handle_outgoing_link_flow_control(
31863194
credit = CappedCredit,
31873195
drain = Drain},
31883196
at_least_one_credit_req_in_flight = true},
3189-
{ok, QStates, Actions} = rabbit_queue_type:credit(
3190-
QName, Ctag,
3191-
QFC#queue_flow_ctl.delivery_count,
3192-
CappedCredit, Drain, QStates0),
3197+
{ok, QStates1, Actions0} = rabbit_queue_type:credit(
3198+
QName, Ctag,
3199+
QFC#queue_flow_ctl.delivery_count,
3200+
CappedCredit, Drain, QStates0),
3201+
{ok, QStates, Actions1} = case parse_deferred_tokens(FlowProps) of
3202+
[] ->
3203+
{ok, QStates1, []};
3204+
Tokens ->
3205+
rabbit_queue_type:assign_deferred(
3206+
QName, Ctag, Tokens, QStates1)
3207+
end,
31933208
State = State0#state{
31943209
queue_states = QStates,
31953210
outgoing_links = OutgoingLinks#{HandleInt := Link}},
3196-
handle_queue_actions(Actions, State);
3211+
handle_queue_actions(Actions0 ++ Actions1, State);
31973212
true ->
31983213
%% A credit request is currently in-flight. Let's first process its reply
31993214
%% before sending the next request. This ensures our outgoing_pending
@@ -3203,6 +3218,8 @@ handle_outgoing_link_flow_control(
32033218
%% to reason about. Therefore, we stash the new request. If there is already a
32043219
%% stashed request, we replace it because the latest flow control state from the
32053220
%% client applies.
3221+
%% Deferral tokens are not stashed; they target the current state of the delayed
3222+
%% set and must not be replayed when the stashed credit request is processed.
32063223
Link = Link0#outgoing_link{
32073224
stashed_credit_req = #credit_req{
32083225
delivery_count = DeliveryCountRcv,
@@ -3212,6 +3229,21 @@ handle_outgoing_link_flow_control(
32123229
State0#state{outgoing_links = OutgoingLinks#{HandleInt := Link}}
32133230
end.
32143231

3232+
parse_deferred_tokens(undefined) ->
3233+
[];
3234+
parse_deferred_tokens({map, KVList}) ->
3235+
case lists:keyfind({symbol, <<"rabbitmq:deferral-tokens">>}, 1, KVList) of
3236+
{{symbol, <<"rabbitmq:deferral-tokens">>}, {array, utf8, Elems}} ->
3237+
[T || {utf8, T} <- Elems];
3238+
false ->
3239+
[];
3240+
{{symbol, <<"rabbitmq:deferral-tokens">>}, Other} ->
3241+
protocol_error(
3242+
?V_1_0_AMQP_ERROR_INVALID_FIELD,
3243+
"rabbitmq:deferral-tokens must be an array of AMQP type utf8, got: ~tp",
3244+
[Other])
3245+
end.
3246+
32153247
delivery_count_rcv(?UINT(DeliveryCount)) ->
32163248
DeliveryCount;
32173249
delivery_count_rcv(undefined) ->

deps/rabbit/src/rabbit_fifo.erl

Lines changed: 46 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -2403,24 +2403,20 @@ return_one(#{system_time := Ts} = Meta, MsgId,
24032403
end.
24042404

24052405
should_delay(DeliveryFailed, DelayedRetry, Ts, Header, Anns) ->
2406-
%% First check for explicit x-opt-delivery-time annotation.
2407-
%% This takes precedence over delayed_retry configuration.
2408-
DeferralToken = case Anns of
2409-
#{<<"x-opt-deferral-token">> := Token}
2410-
when is_binary(Token) ->
2411-
Token;
2412-
_ ->
2413-
undefined
2414-
end,
24152406
case Anns of
24162407
#{<<"x-opt-delivery-time">> := DeliveryTime}
24172408
when is_integer(DeliveryTime),
24182409
DeliveryTime > Ts ->
2410+
%% Deferral tokens are only honoured when the client explicitly
2411+
%% sets a delivery time; the delayed-retry path never creates a
2412+
%% deferred entry so that tokens remain a purely client-driven
2413+
%% mechanism.
2414+
DeferralToken = maps:get(<<"x-opt-deferral-token">>, Anns, undefined),
24192415
{true, DeliveryTime, DeferralToken};
24202416
_ ->
24212417
case should_delay0(DeliveryFailed, DelayedRetry, Ts, Header) of
24222418
{true, ReadyAt} ->
2423-
{true, ReadyAt, DeferralToken};
2419+
{true, ReadyAt, undefined};
24242420
false ->
24252421
false
24262422
end
@@ -2616,8 +2612,7 @@ take_next_delayed(Ts, #delayed{next = {ReadyAt, Idx, Msg},
26162612
{?TUPLE(NextReadyAt, NextIdx), V} = gb_trees:smallest(Tree),
26172613
{NextReadyAt, NextIdx, V}
26182614
end,
2619-
%% Remove any deferral token that maps to this key
2620-
Deferred = maps:filter(fun(_Token, K) -> K =/= Key end, Deferred0),
2615+
Deferred = remove_deferred_key(Key, Deferred0),
26212616
Delayed = #delayed{tree = Tree, next = Next, deferred = Deferred},
26222617
{Msg, Delayed};
26232618
take_next_delayed(_Ts, #delayed{}) ->
@@ -2660,12 +2655,22 @@ take_delayed_for_retry(N, Ts, #delayed{tree = Tree0,
26602655
?TUPLE(ReadyAt, Idx) = NextKey,
26612656
{ReadyAt, Idx, NextMsg}
26622657
end,
2663-
%% Remove any deferral token that maps to this key
2664-
Deferred = maps:filter(fun(_Token, K) -> K =/= Key end, Deferred0),
2658+
Deferred = remove_deferred_key(Key, Deferred0),
26652659
Delayed = #delayed{tree = Tree, next = Next, deferred = Deferred},
26662660
take_delayed_for_retry(N - 1, Ts, Delayed, [Msg | Acc])
26672661
end.
26682662

2663+
%% Drop a single tree key from every token's key list, dropping the token
2664+
%% entirely once its last key is removed.
2665+
remove_deferred_key(Key, Deferred0) ->
2666+
maps:filtermap(
2667+
fun(_Token, Keys) ->
2668+
case lists:delete(Key, Keys) of
2669+
[] -> false;
2670+
Remaining -> {true, Remaining}
2671+
end
2672+
end, Deferred0).
2673+
26692674
take_deferred(Tokens, Delayed) ->
26702675
take_deferred(Tokens, Delayed, [], []).
26712676

@@ -2675,20 +2680,30 @@ take_deferred([Token | Rest], #delayed{tree = Tree0,
26752680
deferred = Deferred0} = Delayed0,
26762681
MsgsAcc, NotFoundAcc) ->
26772682
case maps:take(Token, Deferred0) of
2678-
{Key, Deferred1} ->
2679-
case gb_trees:lookup(Key, Tree0) of
2680-
{value, Msg} ->
2681-
Tree = gb_trees:delete(Key, Tree0),
2682-
Next = update_delayed_next(Tree),
2683-
Delayed = Delayed0#delayed{tree = Tree,
2684-
next = Next,
2685-
deferred = Deferred1},
2686-
take_deferred(Rest, Delayed, [Msg | MsgsAcc], NotFoundAcc);
2687-
none ->
2688-
%% Key in deferred map but not in tree - inconsistent,
2689-
%% treat as not found and clean up
2690-
Delayed = Delayed0#delayed{deferred = Deferred1},
2691-
take_deferred(Rest, Delayed, MsgsAcc, [Token | NotFoundAcc])
2683+
{Keys, Deferred1} ->
2684+
%% Keys were prepended as they were parked, so reverse to
2685+
%% resolve them oldest first.
2686+
{Tree, MsgsAcc1, Found} =
2687+
lists:foldl(
2688+
fun(Key, {TreeAcc, Acc, FoundAcc}) ->
2689+
case gb_trees:lookup(Key, TreeAcc) of
2690+
{value, Msg} ->
2691+
{gb_trees:delete(Key, TreeAcc), [Msg | Acc], true};
2692+
none ->
2693+
%% Key in deferred map but not in tree -
2694+
%% inconsistent, skip and clean up
2695+
{TreeAcc, Acc, FoundAcc}
2696+
end
2697+
end, {Tree0, MsgsAcc, false}, lists:reverse(Keys)),
2698+
Next = update_delayed_next(Tree),
2699+
Delayed = Delayed0#delayed{tree = Tree,
2700+
next = Next,
2701+
deferred = Deferred1},
2702+
case Found of
2703+
true ->
2704+
take_deferred(Rest, Delayed, MsgsAcc1, NotFoundAcc);
2705+
false ->
2706+
take_deferred(Rest, Delayed, MsgsAcc1, [Token | NotFoundAcc])
26922707
end;
26932708
error ->
26942709
take_deferred(Rest, Delayed0, MsgsAcc, [Token | NotFoundAcc])
@@ -2765,7 +2780,9 @@ delayed_in(ReadyAt, Idx, Msg, DeferralToken, #delayed{tree = Tree0,
27652780
undefined ->
27662781
Deferred0;
27672782
_ ->
2768-
Deferred0#{DeferralToken => Key}
2783+
maps:update_with(DeferralToken,
2784+
fun(Keys) -> [Key | Keys] end,
2785+
[Key], Deferred0)
27692786
end,
27702787
#delayed{tree = Tree, next = Next, deferred = Deferred}.
27712788

deps/rabbit/src/rabbit_fifo.hrl

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -202,8 +202,11 @@
202202
{tree = gb_trees:empty() :: gb_trees:tree(delayed_key(), msg()),
203203
%% Cached smallest entry for O(1) readiness check in take_next_msg
204204
next = undefined :: option({milliseconds(), ra:index(), msg()}),
205-
%% Map from deferral token to tree key for direct message lookup
206-
deferred = #{} :: #{deferral_token() => delayed_key()}}).
205+
%% Map from deferral token to the tree keys of all messages parked
206+
%% under it. A single token may address more than one message,
207+
%% e.g. when a ranged DISPOSITION settles several deliveries with
208+
%% the same annotations.
209+
deferred = #{} :: #{deferral_token() => [delayed_key()]}}).
207210

208211
-record(enqueuer,
209212
{next_seqno = 1 :: msg_seqno(),

deps/rabbit/src/rabbit_fifo_client.erl

Lines changed: 4 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -527,26 +527,16 @@ retry_delayed(Server, Mode) when Mode =:= all orelse
527527
%% @param ConsumerTag the tag uniquely identifying the consumer.
528528
%% @param Tokens list of deferral tokens (from x-opt-deferral-token).
529529
%% @param State the {@module} state
530-
%% @returns `{ok, NumAssigned}' on success,
531-
%% `{partial, NumAssigned, NotFoundTokens}' if some tokens not found,
532-
%% `{error, Reason}' on failure.
530+
%% @returns `{state(), actions()}' - matched messages arrive as normal deliveries.
533531
-spec assign_deferred(rabbit_types:ctag(), [binary()], state()) ->
534-
{{ok, non_neg_integer()} |
535-
{partial, non_neg_integer(), [binary()]} |
536-
{error, term()},
537-
state()}.
532+
{state(), rabbit_queue_type:actions()}.
538533
assign_deferred(ConsumerTag, [_|_] = Tokens, State0) ->
539534
ConsumerKey = consumer_key(ConsumerTag, State0),
540535
ServerId = pick_server(State0),
541536
Cmd = rabbit_fifo:make_delayed({assign_deferred, ConsumerKey, Tokens}),
542-
case ra:process_command(ServerId, Cmd, ?COMMAND_TIMEOUT) of
543-
{ok, Reply, _} ->
544-
{Reply, State0};
545-
Err ->
546-
{Err, State0}
547-
end;
537+
{send_command(ServerId, undefined, Cmd, normal, State0), []};
548538
assign_deferred(_ConsumerTag, [], State) ->
549-
{{ok, 0}, State}.
539+
{State, []}.
550540

551541
-spec pending_size(state()) -> non_neg_integer().
552542
pending_size(#state{pending = Pend}) ->

0 commit comments

Comments
 (0)