summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--include/rabbit.hrl2
-rw-r--r--src/rabbit_disk_queue.erl251
-rw-r--r--src/rabbit_mixed_queue.erl18
-rw-r--r--src/rabbit_tests.erl2
4 files changed, 83 insertions, 190 deletions
diff --git a/include/rabbit.hrl b/include/rabbit.hrl
index b8425bafae..0ba31cb5e9 100644
--- a/include/rabbit.hrl
+++ b/include/rabbit.hrl
@@ -65,7 +65,7 @@
-record(basic_message, {exchange_name, routing_key, content,
guid, is_persistent}).
--record(dq_msg_loc, {queue_and_seq_id, is_delivered, msg_id, next_seq_id}).
+-record(dq_msg_loc, {queue_and_seq_id, is_delivered, msg_id}).
-record(delivery, {mandatory, immediate, txn, sender, message}).
diff --git a/src/rabbit_disk_queue.erl b/src/rabbit_disk_queue.erl
index 9c7d35eb4e..96889dbdf4 100644
--- a/src/rabbit_disk_queue.erl
+++ b/src/rabbit_disk_queue.erl
@@ -40,7 +40,7 @@
-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,
+ requeue/2, purge/1, delete_queue/1,
delete_non_durable_queues/1, auto_ack_next_message/1,
requeue_next_n/2
]).
@@ -106,7 +106,7 @@
%% FileSummary: this is an ets table which contains:
%% {File, ValidTotalSize, ContiguousTop, Left, Right}
%% Sequences: this is an ets table which contains:
-%% {Q, ReadSeqId, WriteSeqId, QueueLength}
+%% {Q, ReadSeqId, WriteSeqId}
%% rabbit_disk_queue: this is an mnesia table which contains:
%% #dq_msg_loc { queue_and_seq_id = {Q, SeqId},
%% is_delivered = IsDelivered,
@@ -245,7 +245,6 @@
-ifdef(use_specs).
-type(seq_id() :: non_neg_integer()).
--type(seq_id_or_next() :: ( seq_id() | 'next' )).
-spec(start_link/0 :: () ->
({'ok', pid()} | 'ignore' | {'error', any()})).
@@ -261,10 +260,7 @@
-spec(tx_commit/3 :: (queue_name(), [msg_id()], [{msg_id(), seq_id()}]) ->
'ok').
-spec(tx_cancel/1 :: ([msg_id()]) -> 'ok').
--spec(requeue/2 :: (queue_name(), [{msg_id(), seq_id()}]) -> 'ok').
--spec(requeue_with_seqs/2 ::
- (queue_name(),
- [{{msg_id(), seq_id()}, {seq_id_or_next(), bool()}}]) -> 'ok').
+-spec(requeue/2 :: (queue_name(), [{{msg_id(), seq_id()}, 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').
@@ -314,9 +310,6 @@ tx_cancel(MsgIds) when is_list(MsgIds) ->
requeue(Q, MsgSeqIds) when is_list(MsgSeqIds) ->
gen_server2:cast(?SERVER, {requeue, Q, 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}).
@@ -454,9 +447,8 @@ handle_call({phantom_deliver, Q}, _From, State) ->
{ok, Result, State1} = internal_deliver(Q, false, false, State),
reply(Result, State1);
handle_call({tx_commit, Q, PubMsgIds, AckSeqIds}, From, State) ->
- PubMsgSeqIds = zip_with_tail(PubMsgIds, {duplicate, next}),
{Reply, State1} =
- internal_tx_commit(Q, PubMsgSeqIds, AckSeqIds, From, State),
+ internal_tx_commit(Q, PubMsgIds, AckSeqIds, From, State),
case Reply of
true -> reply(ok, State1);
false -> noreply(State1)
@@ -483,8 +475,8 @@ handle_call(to_disk_only_mode, _From, State) ->
handle_call(to_ram_disk_mode, _From, State) ->
reply(ok, to_ram_disk_mode(State));
handle_call({length, Q}, _From, State = #dqstate { sequences = Sequences }) ->
- {_ReadSeqId, _WriteSeqId, Length} = sequence_lookup(Sequences, Q),
- reply(Length, State);
+ {ReadSeqId, WriteSeqId} = sequence_lookup(Sequences, Q),
+ reply(WriteSeqId - ReadSeqId, State);
handle_call({delete_non_durable_queues, DurableQueues}, _From, State) ->
{ok, State1} = internal_delete_non_durable_queues(DurableQueues, State),
reply(ok, State1);
@@ -492,8 +484,7 @@ handle_call(cache_info, _From, State = #dqstate { message_cache = Cache }) ->
reply(ets:info(Cache), State).
handle_cast({publish, Q, Message, IsDelivered}, State) ->
- {ok, _MsgSeqId, State1} =
- internal_publish(Q, Message, next, IsDelivered, State),
+ {ok, _MsgSeqId, State1} = internal_publish(Q, Message, IsDelivered, State),
noreply(State1);
handle_cast({ack, Q, MsgSeqIds}, State) ->
{ok, State1} = internal_ack(Q, MsgSeqIds, State),
@@ -508,11 +499,7 @@ handle_cast({tx_cancel, MsgIds}, State) ->
{ok, State1} = internal_tx_cancel(MsgIds, State),
noreply(State1);
handle_cast({requeue, Q, MsgSeqIds}, State) ->
- MsgSeqSeqIds = zip_with_tail(MsgSeqIds, {duplicate, {next, true}}),
- {ok, State1} = internal_requeue(Q, MsgSeqSeqIds, State),
- noreply(State1);
-handle_cast({requeue_with_seqs, Q, MsgSeqSeqIds}, State) ->
- {ok, State1} = internal_requeue(Q, MsgSeqSeqIds, State),
+ {ok, State1} = internal_requeue(Q, MsgSeqIds, State),
noreply(State1);
handle_cast({requeue_next_n, Q, N}, State) ->
{ok, State1} = internal_requeue_next_n(Q, N, State),
@@ -695,13 +682,6 @@ form_filename(Name) ->
base_directory() ->
filename:join(mnesia:system_info(directory), "rabbit_disk_queue/").
-zip_with_tail(List1, List2) when length(List1) =:= length(List2) ->
- lists:zip(List1, List2);
-zip_with_tail(List = [_|Tail], {last, E}) ->
- zip_with_tail(List, Tail ++ [E]);
-zip_with_tail(List, {duplicate, E}) ->
- zip_with_tail(List, lists:duplicate(erlang:length(List), E)).
-
dets_ets_lookup(#dqstate { msg_location_dets = MsgLocationDets,
operation_mode = disk_only },
Key) ->
@@ -749,29 +729,6 @@ dets_ets_match_object(#dqstate { msg_location_ets = MsgLocationEts,
Obj) ->
ets:match_object(MsgLocationEts, Obj).
-find_next_seq_id(CurrentSeq, next) ->
- CurrentSeq + 1;
-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.
-
get_read_handle(File, Offset, State =
#dqstate { read_file_handles = {ReadHdls, ReadHdlsAge},
read_file_handles_limit = ReadFileHandlesLimit,
@@ -809,32 +766,12 @@ get_read_handle(File, Offset, State =
{FileHdl,
State1 #dqstate { read_file_handles = {ReadHdls2, ReadHdlsAge3} }}.
-adjust_last_msg_seq_id(_Q, ExpectedSeqId, next, _Mode) ->
- ExpectedSeqId;
-adjust_last_msg_seq_id(_Q, 0, SuppliedSeqId, _Mode) ->
- SuppliedSeqId;
-adjust_last_msg_seq_id(_Q, ExpectedSeqId, ExpectedSeqId, _Mode) ->
- ExpectedSeqId;
-adjust_last_msg_seq_id(Q, ExpectedSeqId, SuppliedSeqId, dirty)
- when SuppliedSeqId > ExpectedSeqId ->
- [Obj] = mnesia:dirty_read(rabbit_disk_queue, {Q, ExpectedSeqId - 1}),
- ok = mnesia:dirty_write(rabbit_disk_queue,
- Obj #dq_msg_loc { next_seq_id = SuppliedSeqId }),
- SuppliedSeqId;
-adjust_last_msg_seq_id(Q, ExpectedSeqId, SuppliedSeqId, Lock)
- when SuppliedSeqId > ExpectedSeqId ->
- [Obj] = mnesia:read(rabbit_disk_queue, {Q, ExpectedSeqId - 1}, Lock),
- ok = mnesia:write(rabbit_disk_queue,
- Obj #dq_msg_loc { next_seq_id = SuppliedSeqId },
- Lock),
- SuppliedSeqId.
-
sequence_lookup(Sequences, Q) ->
case ets:lookup(Sequences, Q) of
[] ->
- {0, 0, 0};
- [{Q, ReadSeqId, WriteSeqId, Length}] ->
- {ReadSeqId, WriteSeqId, Length}
+ {0, 0};
+ [{Q, ReadSeqId, WriteSeqId}] ->
+ {ReadSeqId, WriteSeqId}
end.
start_commit_timer(State = #dqstate { commit_timer_ref = undefined }) ->
@@ -910,14 +847,14 @@ internal_deliver(Q, ReadMsg, FakeDeliver,
State = #dqstate { sequences = Sequences }) ->
case ets:lookup(Sequences, Q) of
[] -> {ok, empty, State};
- [{Q, SeqId, SeqId, 0}] -> {ok, empty, State};
- [{Q, ReadSeqId, WriteSeqId, Length}] when Length > 0 ->
- Remaining = Length - 1,
- {ok, Result, NextReadSeqId, State1} =
+ [{Q, SeqId, SeqId}] -> {ok, empty, State};
+ [{Q, ReadSeqId, WriteSeqId}] when WriteSeqId >= ReadSeqId ->
+ Remaining = WriteSeqId - ReadSeqId - 1,
+ {ok, Result, State1} =
internal_read_message(
Q, ReadSeqId, FakeDeliver, ReadMsg, State),
true = ets:insert(Sequences,
- {Q, NextReadSeqId, WriteSeqId, Remaining}),
+ {Q, ReadSeqId+1, WriteSeqId}),
{ok,
case Result of
{MsgId, Delivered, {MsgId, ReadSeqId}} ->
@@ -930,8 +867,7 @@ internal_deliver(Q, ReadMsg, FakeDeliver,
internal_read_message(Q, ReadSeqId, FakeDeliver, ReadMsg, State) ->
[Obj =
- #dq_msg_loc {is_delivered = Delivered, msg_id = MsgId,
- next_seq_id = NextReadSeqId}] =
+ #dq_msg_loc {is_delivered = Delivered, msg_id = MsgId}] =
mnesia:dirty_read(rabbit_disk_queue, {Q, ReadSeqId}),
[{MsgId, RefCount, File, Offset, TotalSize}] =
dets_ets_lookup(State, MsgId),
@@ -959,13 +895,13 @@ internal_read_message(Q, ReadSeqId, FakeDeliver, ReadMsg, State) ->
_ -> insert_into_cache(Message, BodySize, State1)
end,
{ok, {Message, BodySize, Delivered, {MsgId, ReadSeqId}},
- NextReadSeqId, State1};
+ State1};
{Message, BodySize, _RefCount} ->
{ok, {Message, BodySize, Delivered, {MsgId, ReadSeqId}},
- NextReadSeqId, State}
+ State}
end;
false ->
- {ok, {MsgId, Delivered, {MsgId, ReadSeqId}}, NextReadSeqId, State}
+ {ok, {MsgId, Delivered, {MsgId, ReadSeqId}}, State}
end.
internal_auto_ack(Q, State) ->
@@ -1064,25 +1000,15 @@ internal_tx_publish(MsgId, Message,
{ok, State}
end.
-%% can call this with PubMsgSeqIds as zip(PubMsgIds, duplicate(N, next))
-internal_tx_commit(Q, PubMsgSeqIds, AckSeqIds, From,
+internal_tx_commit(Q, PubMsgIds, AckSeqIds, From,
State = #dqstate { sequences = Sequences,
current_file_name = CurFile,
current_dirty = IsDirty,
on_sync_froms = SyncFroms,
last_sync_offset = SyncOffset
}) ->
- {PubList, PubAcc, ReadSeqId, Length} =
- case PubMsgSeqIds of
- [] -> {[], undefined, undefined, undefined};
- [{_, FirstSeqIdTo}|_] ->
- {InitReadSeqId, InitWriteSeqId, InitLength} =
- sequence_lookup(Sequences, Q),
- InitReadSeqId1 = determine_next_read_id(
- InitReadSeqId, InitWriteSeqId, FirstSeqIdTo),
- { zip_with_tail(PubMsgSeqIds, {last, {next, next}}),
- InitWriteSeqId, InitReadSeqId1, InitLength}
- end,
+ {InitReadSeqId, InitWriteSeqId} = sequence_lookup(Sequences, Q),
+ WriteSeqId = InitWriteSeqId + erlang:length(PubMsgIds),
{atomic, {InCurFile, WriteSeqId, State1}} =
mnesia:transaction(
fun() ->
@@ -1094,34 +1020,27 @@ internal_tx_commit(Q, PubMsgSeqIds, AckSeqIds, From,
%% order which _could_not_ have happened.
{InCurFile1, WriteSeqId1} =
lists:foldl(
- fun ({{MsgId, SeqId}, {_NextMsgId, NextSeqId}},
- {InCurFileAcc, ExpectedSeqId}) ->
+ fun (MsgId, {InCurFileAcc, SeqId}) ->
[{MsgId, _RefCount, File, Offset,
_TotalSize}] = dets_ets_lookup(State, MsgId),
- SeqId1 = adjust_last_msg_seq_id(
- Q, ExpectedSeqId, SeqId, write),
- NextSeqId1 =
- find_next_seq_id(SeqId1, NextSeqId),
ok = mnesia:write(
rabbit_disk_queue,
#dq_msg_loc { queue_and_seq_id =
- {Q, SeqId1},
+ {Q, SeqId},
msg_id = MsgId,
- is_delivered = false,
- next_seq_id = NextSeqId1
+ is_delivered = false
},
write),
{InCurFileAcc orelse (File =:= CurFile andalso
Offset >= SyncOffset),
- NextSeqId1}
- end, {false, PubAcc}, PubList),
+ SeqId + 1}
+ end, {false, InitWriteSeqId}, PubMsgIds),
{ok, State2} = remove_messages(Q, AckSeqIds, txn, State),
{InCurFile1, WriteSeqId1, State2}
end),
- true = case PubList of
+ true = case PubMsgIds of
[] -> true;
- _ -> ets:insert(Sequences, {Q, ReadSeqId, WriteSeqId,
- Length + erlang:length(PubList)})
+ _ -> ets:insert(Sequences, {Q, InitReadSeqId, WriteSeqId})
end,
if IsDirty andalso InCurFile ->
{false, State1 #dqstate { on_sync_froms = [From | SyncFroms] }};
@@ -1129,34 +1048,28 @@ internal_tx_commit(Q, PubMsgSeqIds, AckSeqIds, From,
{true, State1}
end.
-%% SeqId can be 'next'
-internal_publish(Q, Message = #basic_message { guid = MsgId }, SeqId,
+internal_publish(Q, Message = #basic_message { guid = MsgId },
IsDelivered, State) ->
{ok, State1 = #dqstate { sequences = Sequences }} =
internal_tx_publish(MsgId, Message, State),
- {ReadSeqId, WriteSeqId, Length} =
- sequence_lookup(Sequences, Q),
- ReadSeqId3 = determine_next_read_id(ReadSeqId, WriteSeqId, SeqId),
- WriteSeqId3 = adjust_last_msg_seq_id(Q, WriteSeqId, SeqId, dirty),
- WriteSeqId3Next = WriteSeqId3 + 1,
+ {ReadSeqId, WriteSeqId} = sequence_lookup(Sequences, Q),
ok = mnesia:dirty_write(rabbit_disk_queue,
- #dq_msg_loc { queue_and_seq_id = {Q, WriteSeqId3},
+ #dq_msg_loc { queue_and_seq_id = {Q, WriteSeqId},
msg_id = MsgId,
- next_seq_id = WriteSeqId3Next,
is_delivered = IsDelivered}),
- true = ets:insert(Sequences, {Q, ReadSeqId3, WriteSeqId3Next, Length + 1}),
- {ok, {MsgId, WriteSeqId3}, State1}.
+ true = ets:insert(Sequences, {Q, ReadSeqId, WriteSeqId + 1}),
+ {ok, {MsgId, WriteSeqId}, State1}.
internal_tx_cancel(MsgIds, State) ->
%% we don't need seq ids because we're not touching mnesia,
%% because seqids were never assigned
- MsgSeqIds = zip_with_tail(MsgIds, {duplicate, undefined}),
+ MsgSeqIds = lists:zip(MsgIds,
+ lists:duplicate(erlang:length(MsgIds), undefined)),
remove_messages(undefined, MsgSeqIds, false, State).
internal_requeue(_Q, [], State) ->
{ok, State};
-internal_requeue(Q, MsgSeqIds = [{_, {FirstSeqIdTo, _}}|_],
- State = #dqstate { sequences = Sequences }) ->
+internal_requeue(Q, MsgSeqIds, State = #dqstate { sequences = Sequences }) ->
%% We know that every seq_id in here is less than the ReadSeqId
%% you'll get if you look up this queue in Sequences (i.e. they've
%% already been delivered). We also know that the rows for these
@@ -1179,76 +1092,59 @@ internal_requeue(Q, MsgSeqIds = [{_, {FirstSeqIdTo, _}}|_],
%% MsgLocation and FileSummary stay put (which makes further sense
%% as they have no concept of sequence id anyway).
- {ReadSeqId, WriteSeqId, Length} = sequence_lookup(Sequences, Q),
- ReadSeqId1 = determine_next_read_id(ReadSeqId, WriteSeqId, FirstSeqIdTo),
- MsgSeqIdsZipped = zip_with_tail(MsgSeqIds, {last, {next, {next, true}}}),
+ {ReadSeqId, WriteSeqId} = sequence_lookup(Sequences, Q),
{atomic, {WriteSeqId1, Q, State}} =
mnesia:transaction(
fun() ->
ok = mnesia:write_lock_table(rabbit_disk_queue),
lists:foldl(fun requeue_message/2, {WriteSeqId, Q, State},
- MsgSeqIdsZipped)
+ MsgSeqIds)
end),
- true = ets:insert(Sequences, {Q, ReadSeqId1, WriteSeqId1,
- Length + erlang:length(MsgSeqIds)}),
+ true = ets:insert(Sequences, {Q, ReadSeqId, WriteSeqId1}),
{ok, State}.
-requeue_message({{{MsgId, SeqIdOrig}, {SeqIdTo, NewIsDelivered}},
- {_NextMsgSeqId, {NextSeqIdTo, _NextNewIsDelivered}}},
- {ExpectedSeqIdTo, Q, State}) ->
- SeqIdTo1 = adjust_last_msg_seq_id(Q, ExpectedSeqIdTo, SeqIdTo, write),
- NextSeqIdTo1 = find_next_seq_id(SeqIdTo1, NextSeqIdTo),
- [Obj = #dq_msg_loc { is_delivered = true, msg_id = MsgId,
- next_seq_id = NextSeqIdOrig }] =
- mnesia:read(rabbit_disk_queue, {Q, SeqIdOrig}, write),
- if SeqIdTo1 == SeqIdOrig andalso NextSeqIdTo1 == NextSeqIdOrig -> ok;
- true ->
- ok = mnesia:write(rabbit_disk_queue,
- Obj #dq_msg_loc {queue_and_seq_id = {Q, SeqIdTo1},
- next_seq_id = NextSeqIdTo1,
- is_delivered = NewIsDelivered
- },
- write),
- ok = mnesia:delete(rabbit_disk_queue, {Q, SeqIdOrig}, write)
- end,
+requeue_message({{MsgId, SeqId}, IsDelivered}, {WriteSeqId, Q, State}) ->
+ [Obj = #dq_msg_loc { is_delivered = true, msg_id = MsgId }] =
+ mnesia:read(rabbit_disk_queue, {Q, SeqId}, write),
+ ok = mnesia:write(rabbit_disk_queue,
+ Obj #dq_msg_loc {queue_and_seq_id = {Q, WriteSeqId},
+ is_delivered = IsDelivered
+ },
+ write),
+ ok = mnesia:delete(rabbit_disk_queue, {Q, SeqId}, write),
decrement_cache(MsgId, State),
- {NextSeqIdTo1, Q, State}.
+ {WriteSeqId + 1, 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};
+ {ReadSeqId, WriteSeqId} = sequence_lookup(Sequences, Q),
+ if N >= (WriteSeqId - ReadSeqId) -> {ok, State};
true ->
{atomic, {ReadSeqIdN, WriteSeqIdN}} =
mnesia:transaction(
fun() ->
ok = mnesia:write_lock_table(rabbit_disk_queue),
- requeue_next_messages(Q, N, ReadSeqId1, WriteSeqId)
+ requeue_next_messages(Q, N, ReadSeqId, WriteSeqId)
end
),
- true = ets:insert(Sequences, {Q, ReadSeqIdN, WriteSeqIdN, Length}),
+ true = ets:insert(Sequences, {Q, ReadSeqIdN, WriteSeqIdN}),
{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),
+ [Obj] = 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),
+ Obj #dq_msg_loc {queue_and_seq_id = {Q, WriteSeq}},
+ write),
ok = mnesia:delete(rabbit_disk_queue, {Q, ReadSeq}, write),
- requeue_next_messages(Q, N - 1, NextSeqIdOrig, NextWriteSeq).
+ requeue_next_messages(Q, N - 1, ReadSeq + 1, WriteSeq + 1).
internal_purge(Q, State = #dqstate { sequences = Sequences }) ->
case ets:lookup(Sequences, Q) of
[] -> {ok, 0, State};
- [{Q, ReadSeqId, WriteSeqId, _Length}] ->
+ [{Q, ReadSeqId, WriteSeqId}] ->
{atomic, {ok, State1}} =
mnesia:transaction(
fun() ->
@@ -1257,15 +1153,14 @@ internal_purge(Q, State = #dqstate { sequences = Sequences }) ->
rabbit_misc:unfold(
fun (SeqId) when SeqId == WriteSeqId -> false;
(SeqId) ->
- [#dq_msg_loc { msg_id = MsgId,
- next_seq_id = NextSeqId }
- ] = mnesia:read(rabbit_disk_queue,
+ [#dq_msg_loc { msg_id = MsgId }] =
+ mnesia:read(rabbit_disk_queue,
{Q, SeqId}, write),
- {true, {MsgId, SeqId}, NextSeqId}
+ {true, {MsgId, SeqId}, SeqId + 1}
end, ReadSeqId),
remove_messages(Q, MsgSeqIds, txn, State)
end),
- true = ets:insert(Sequences, {Q, WriteSeqId, WriteSeqId, 0}),
+ true = ets:insert(Sequences, {Q, WriteSeqId, WriteSeqId}),
{ok, WriteSeqId - ReadSeqId, State1}
end.
@@ -1282,8 +1177,7 @@ internal_delete_queue(Q, State) ->
rabbit_disk_queue,
#dq_msg_loc { queue_and_seq_id = {Q, '_'},
msg_id = '_',
- is_delivered = '_',
- next_seq_id = '_'
+ is_delivered = '_'
},
write),
MsgSeqIds =
@@ -1298,7 +1192,7 @@ internal_delete_queue(Q, State) ->
internal_delete_non_durable_queues(
DurableQueues, State = #dqstate { sequences = Sequences }) ->
ets:foldl(
- fun ({Q, _Read, _Write, _Length}, {ok, State1}) ->
+ fun ({Q, _Read, _Write}, {ok, State1}) ->
case sets:is_element(Q, DurableQueues) of
true -> {ok, State1};
false -> internal_delete_queue(Q, State1)
@@ -1667,10 +1561,9 @@ remove_gaps_in_sequences(#dqstate { sequences = Sequences }) ->
fun ({Q, ReadSeqId, WriteSeqId, _Length}) ->
Gap = shuffle_up(Q, ReadSeqId-1, WriteSeqId-1, 0),
ReadSeqId1 = ReadSeqId + Gap,
- Length = WriteSeqId - ReadSeqId1,
true =
ets:insert(Sequences,
- {Q, ReadSeqId1, WriteSeqId, Length})
+ {Q, ReadSeqId1, WriteSeqId})
end, ets:match_object(Sequences, '_'))
end).
@@ -1685,9 +1578,7 @@ shuffle_up(Q, BaseSeqId, SeqId, Gap) ->
0 -> ok;
_ -> mnesia:write(rabbit_disk_queue,
Obj #dq_msg_loc {
- queue_and_seq_id = {Q, SeqId + Gap },
- next_seq_id = SeqId + Gap + 1
- },
+ queue_and_seq_id = {Q, SeqId + Gap }},
write),
mnesia:delete(rabbit_disk_queue, {Q, SeqId}, write)
end,
@@ -1720,8 +1611,7 @@ load_messages(Left, [File|Files],
(rabbit_disk_queue,
#dq_msg_loc { msg_id = MsgId,
queue_and_seq_id = '_',
- is_delivered = '_',
- next_seq_id = '_'
+ is_delivered = '_'
},
msg_id)) of
0 -> {VMAcc, VTSAcc};
@@ -1761,8 +1651,7 @@ verify_messages_in_mnesia(MsgIds) ->
(rabbit_disk_queue,
#dq_msg_loc { msg_id = MsgId,
queue_and_seq_id = '_',
- is_delivered = '_',
- next_seq_id = '_'
+ is_delivered = '_'
},
msg_id))
end, MsgIds).
diff --git a/src/rabbit_mixed_queue.erl b/src/rabbit_mixed_queue.erl
index 4a2803a4ef..61487c9d56 100644
--- a/src/rabbit_mixed_queue.erl
+++ b/src/rabbit_mixed_queue.erl
@@ -180,9 +180,13 @@ flush_messages_to_disk_queue(Q, 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_next_n(Q, RequeueCount),
- []
+ true ->
+ ok = if [] == Commit -> ok;
+ true -> rabbit_disk_queue:tx_commit
+ (Q, lists:reverse(Commit), [])
+ end,
+ rabbit_disk_queue:requeue_next_n(Q, RequeueCount),
+ []
end.
to_mixed_mode(_TxnMessages, State = #mqstate { mode = mixed }) ->
@@ -226,7 +230,7 @@ purge_non_persistent_messages(State = #mqstate { mode = disk, queue = Q,
deliver_all_messages(Q, IsDurable, [], [], 0, 0),
ok = if Requeue == [] -> ok;
true ->
- rabbit_disk_queue:requeue_with_seqs(Q, lists:reverse(Requeue))
+ rabbit_disk_queue:requeue(Q, lists:reverse(Requeue))
end,
ok = if Acks == [] -> ok;
true -> rabbit_disk_queue:ack(Q, Acks)
@@ -241,7 +245,7 @@ deliver_all_messages(Q, IsDurable, Acks, Requeue, Length, QSize) ->
OnDisk = IsPersistent andalso IsDurable,
{Acks1, Requeue1, Length1, QSize1} =
if OnDisk -> { Acks,
- [{AckTag, {next, IsDelivered}} | Requeue],
+ [{AckTag, IsDelivered} | Requeue],
Length + 1, QSize + size_of_message(Msg) };
true -> { [AckTag | Acks], Requeue, Length, QSize }
end,
@@ -484,7 +488,7 @@ requeue(MessagesWithAckTags, State = #mqstate { mode = disk, queue = Q,
= lists:foldl(
fun ({#basic_message { is_persistent = IsPersistent }, AckTag}, RQ)
when IsDurable andalso IsPersistent ->
- [AckTag | RQ];
+ [{AckTag, true} | RQ];
({Msg, _AckTag}, RQ) ->
ok = case RQ == [] of
true -> ok;
@@ -508,7 +512,7 @@ requeue(MessagesWithAckTags, State = #mqstate { mode = mixed, queue = Q,
{Acc, MsgBuf2}) ->
OnDisk = IsDurable andalso IsPersistent,
Acc1 =
- if OnDisk -> [AckTag | Acc];
+ if OnDisk -> [{AckTag, true} | Acc];
true -> Acc
end,
{Acc1, queue:in({Msg, true, OnDisk}, MsgBuf2)}
diff --git a/src/rabbit_tests.erl b/src/rabbit_tests.erl
index f108285056..8ab8267725 100644
--- a/src/rabbit_tests.erl
+++ b/src/rabbit_tests.erl
@@ -924,7 +924,7 @@ rdq_test_redeliver() ->
%% now requeue every other message (starting at the _first_)
%% and ack the other ones
lists:foldl(fun (SeqId2, true) ->
- rabbit_disk_queue:requeue(q, [SeqId2]),
+ rabbit_disk_queue:requeue(q, [{SeqId2, true}]),
false;
(SeqId2, false) ->
rabbit_disk_queue:ack(q, [SeqId2]),