diff options
| -rw-r--r-- | include/rabbit.hrl | 2 | ||||
| -rw-r--r-- | src/rabbit_disk_queue.erl | 251 | ||||
| -rw-r--r-- | src/rabbit_mixed_queue.erl | 18 | ||||
| -rw-r--r-- | src/rabbit_tests.erl | 2 |
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]), |
