summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--src/rabbit_disk_queue.erl47
-rw-r--r--src/rabbit_mixed_queue.erl37
-rw-r--r--src/rabbit_queue_mode_manager.erl6
3 files changed, 62 insertions, 28 deletions
diff --git a/src/rabbit_disk_queue.erl b/src/rabbit_disk_queue.erl
index d8c6580f01..9c7d35eb4e 100644
--- a/src/rabbit_disk_queue.erl
+++ b/src/rabbit_disk_queue.erl
@@ -41,7 +41,8 @@
-export([publish/3, deliver/1, phantom_deliver/1, ack/2,
tx_publish/1, tx_commit/3, tx_cancel/1,
requeue/2, requeue_with_seqs/2, purge/1, delete_queue/1,
- delete_non_durable_queues/1, auto_ack_next_message/1
+ delete_non_durable_queues/1, auto_ack_next_message/1,
+ requeue_next_n/2
]).
-export([length/1, filesync/0, cache_info/0]).
@@ -264,6 +265,7 @@
-spec(requeue_with_seqs/2 ::
(queue_name(),
[{{msg_id(), seq_id()}, {seq_id_or_next(), bool()}}]) -> 'ok').
+-spec(requeue_next_n/2 :: (queue_name(), non_neg_integer()) -> 'ok').
-spec(purge/1 :: (queue_name()) -> non_neg_integer()).
-spec(delete_non_durable_queues/1 :: (set()) -> 'ok').
-spec(stop/0 :: () -> 'ok').
@@ -315,6 +317,9 @@ requeue(Q, MsgSeqIds) when is_list(MsgSeqIds) ->
requeue_with_seqs(Q, MsgSeqSeqIds) when is_list(MsgSeqSeqIds) ->
gen_server2:cast(?SERVER, {requeue_with_seqs, Q, MsgSeqSeqIds}).
+requeue_next_n(Q, N) when is_integer(N) ->
+ gen_server2:cast(?SERVER, {requeue_next_n, Q, N}).
+
purge(Q) ->
gen_server2:call(?SERVER, {purge, Q}, infinity).
@@ -509,6 +514,9 @@ handle_cast({requeue, Q, MsgSeqIds}, State) ->
handle_cast({requeue_with_seqs, Q, MsgSeqSeqIds}, State) ->
{ok, State1} = internal_requeue(Q, MsgSeqSeqIds, State),
noreply(State1);
+handle_cast({requeue_next_n, Q, N}, State) ->
+ {ok, State1} = internal_requeue_next_n(Q, N, State),
+ noreply(State1);
handle_cast({delete_queue, Q}, State) ->
{ok, State1} = internal_delete_queue(Q, State),
noreply(State1);
@@ -747,13 +755,19 @@ find_next_seq_id(CurrentSeq, NextSeqId)
when NextSeqId > CurrentSeq ->
NextSeqId.
+%% the queue is empty, and we've just written exactly where we
+%% expected, so read it back
determine_next_read_id(CurrentReadWrite, CurrentReadWrite, CurrentReadWrite) ->
CurrentReadWrite;
+%% we've just written in the next slot, so the next read pos is unaltered
determine_next_read_id(CurrentRead, _CurrentWrite, next) ->
CurrentRead;
+%% queue is empty, but we've written somewhere else - a gap has formed
+%% - so read back from where we wrote, after the gap
determine_next_read_id(CurrentReadWrite, CurrentReadWrite, NextWrite)
when NextWrite > CurrentReadWrite ->
NextWrite;
+%% queue is not empty, and we've created a gap, so the read pos is unaltered
determine_next_read_id(CurrentRead, CurrentWrite, NextWrite)
when NextWrite >= CurrentWrite ->
CurrentRead.
@@ -1200,6 +1214,37 @@ requeue_message({{{MsgId, SeqIdOrig}, {SeqIdTo, NewIsDelivered}},
decrement_cache(MsgId, State),
{NextSeqIdTo1, Q, State}.
+%% move the next N messages from the front of the queue to the back.
+internal_requeue_next_n(Q, N, State = #dqstate { sequences = Sequences }) ->
+ {ReadSeqId, WriteSeqId, Length} = sequence_lookup(Sequences, Q),
+ ReadSeqId1 = determine_next_read_id(ReadSeqId, WriteSeqId, next),
+ if N >= Length -> {ok, State};
+ true ->
+ {atomic, {ReadSeqIdN, WriteSeqIdN}} =
+ mnesia:transaction(
+ fun() ->
+ ok = mnesia:write_lock_table(rabbit_disk_queue),
+ requeue_next_messages(Q, N, ReadSeqId1, WriteSeqId)
+ end
+ ),
+ true = ets:insert(Sequences, {Q, ReadSeqIdN, WriteSeqIdN, Length}),
+ {ok, State}
+ end.
+
+requeue_next_messages(_Q, 0, ReadSeq, WriteSeq) ->
+ {ReadSeq, WriteSeq};
+requeue_next_messages(Q, N, ReadSeq, WriteSeq) ->
+ WriteSeq1 = adjust_last_msg_seq_id(Q, WriteSeq, next, write),
+ NextWriteSeq = find_next_seq_id(WriteSeq1, next),
+ [Obj = #dq_msg_loc { next_seq_id = NextSeqIdOrig }] =
+ mnesia:read(rabbit_disk_queue, {Q, ReadSeq}, write),
+ ok = mnesia:write(rabbit_disk_queue,
+ Obj #dq_msg_loc {queue_and_seq_id = {Q, WriteSeq1},
+ next_seq_id = NextWriteSeq
+ }, write),
+ ok = mnesia:delete(rabbit_disk_queue, {Q, ReadSeq}, write),
+ requeue_next_messages(Q, N - 1, NextSeqIdOrig, NextWriteSeq).
+
internal_purge(Q, State = #dqstate { sequences = Sequences }) ->
case ets:lookup(Sequences, Q) of
[] -> {ok, 0, State};
diff --git a/src/rabbit_mixed_queue.erl b/src/rabbit_mixed_queue.erl
index f415472719..4a2803a4ef 100644
--- a/src/rabbit_mixed_queue.erl
+++ b/src/rabbit_mixed_queue.erl
@@ -126,7 +126,7 @@ to_disk_only_mode(TxnMessages, State =
%% message on disk.
%% Note we also batch together messages on disk so that we minimise
%% the calls to requeue.
- ok = send_messages_to_disk(Q, MsgBuf, [], 0, []),
+ ok = send_messages_to_disk(Q, MsgBuf, 0, 0, []),
%% tx_publish txn messages. Some of these will have been already
%% published if they really are durable and persistent which is
%% why we can't just use our own tx_publish/2 function (would end
@@ -141,47 +141,36 @@ to_disk_only_mode(TxnMessages, State =
garbage_collect(),
{ok, State #mqstate { mode = disk, msg_buf = queue:new() }}.
-send_messages_to_disk(Q, Queue, Requeue, PublishCount, Commit) ->
+send_messages_to_disk(Q, Queue, RequeueCount, PublishCount, Commit) ->
case queue:out(Queue) of
{empty, Queue} ->
ok = flush_messages_to_disk_queue(Q, Commit),
- [] = flush_requeue_to_disk_queue(Q, Requeue, []),
+ [] = flush_requeue_to_disk_queue(Q, RequeueCount, []),
ok;
- {{value, {Msg = #basic_message { guid = MsgId }, IsDelivered, OnDisk}},
+ {{value, {Msg = #basic_message { guid = MsgId }, _IsDelivered, OnDisk}},
Queue1} ->
case OnDisk of
true ->
- ok = flush_messages_to_disk_queue (Q, Commit),
- {MsgId, IsDelivered, AckTag, _PersistRemaining} =
- rabbit_disk_queue:phantom_deliver(Q),
+ ok = flush_messages_to_disk_queue(Q, Commit),
send_messages_to_disk(
- Q, Queue1, [{AckTag, {next, IsDelivered}} | Requeue],
- 0, []);
+ Q, Queue1, 1 + RequeueCount, 0, []);
false ->
Commit1 =
- flush_requeue_to_disk_queue(Q, Requeue, Commit),
+ flush_requeue_to_disk_queue(Q, RequeueCount, Commit),
ok = rabbit_disk_queue:tx_publish(Msg),
case PublishCount == ?TO_DISK_MAX_FLUSH_SIZE of
true ->
ok = flush_messages_to_disk_queue(Q, Commit1),
- send_messages_to_disk(Q, Queue1, [], 1, [MsgId]);
+ send_messages_to_disk(Q, Queue1, 0, 1, [MsgId]);
false ->
send_messages_to_disk
- (Q, Queue1, [], PublishCount + 1,
+ (Q, Queue1, 0, PublishCount + 1,
[MsgId | Commit1])
end
end;
{{value, {disk, Count}}, Queue2} ->
ok = flush_messages_to_disk_queue(Q, Commit),
- {Requeue1, 0} =
- rabbit_misc:unfold(
- fun (0) -> false;
- (N) ->
- {_MsgId, IsDelivered, AckTag, _PersistRemaining}
- = rabbit_disk_queue:phantom_deliver(Q),
- {true, {AckTag, {next, IsDelivered}}, N - 1}
- end, Count),
- send_messages_to_disk(Q, Queue2, Requeue1 ++ Requeue, 0, [])
+ send_messages_to_disk(Q, Queue2, RequeueCount + Count, 0, [])
end.
flush_messages_to_disk_queue(Q, Commit) ->
@@ -189,10 +178,10 @@ flush_messages_to_disk_queue(Q, Commit) ->
true -> rabbit_disk_queue:tx_commit(Q, lists:reverse(Commit), [])
end.
-flush_requeue_to_disk_queue(Q, Requeue, Commit) ->
- if [] == Requeue -> Commit;
+flush_requeue_to_disk_queue(Q, RequeueCount, Commit) ->
+ if 0 == RequeueCount -> Commit;
true -> ok = rabbit_disk_queue:tx_commit(Q, lists:reverse(Commit), []),
- rabbit_disk_queue:requeue_with_seqs(Q, lists:reverse(Requeue)),
+ rabbit_disk_queue:requeue_next_n(Q, RequeueCount),
[]
end.
diff --git a/src/rabbit_queue_mode_manager.erl b/src/rabbit_queue_mode_manager.erl
index 99f6e408d9..f5cc32b445 100644
--- a/src/rabbit_queue_mode_manager.erl
+++ b/src/rabbit_queue_mode_manager.erl
@@ -201,20 +201,20 @@ handle_call({pin_to_disk, Pid}, _From,
disk_mode_pins = Pins }) ->
{Res, State1} =
case sets:is_element(Pid, Pins) of
- true -> {already_pinned, State};
+ true -> {ok, State};
false ->
case find_queue(Pid, Mixed) of
{mixed, {OAlloc, _OActivity}} ->
{Module, Function, Args} = dict:fetch(Pid, Callbacks),
ok = erlang:apply(Module, Function, Args ++ [disk]),
- {convert_to_disk_mode,
+ {ok,
State #state { mixed_queues = dict:erase(Pid, Mixed),
available_tokens = Avail + OAlloc,
disk_mode_pins =
sets:add_element(Pid, Pins)
}};
disk ->
- {already_disk,
+ {ok,
State #state { disk_mode_pins =
sets:add_element(Pid, Pins) }}
end