summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorMatthew Sackman <matthew@lshift.net>2009-10-28 15:25:14 +0000
committerMatthew Sackman <matthew@lshift.net>2009-10-28 15:25:14 +0000
commitbb19594cabc7bc4d6dcafc3afeaddf5ecddc3f76 (patch)
treebe49b7ced50a6ba058a06cd33be79b457da4f75d /src
parenteea836fb25d38fecff1b58896bd1a37b8d0a9081 (diff)
downloadrabbitmq-server-git-bb19594cabc7bc4d6dcafc3afeaddf5ecddc3f76.tar.gz
Some minor cosmetics in qi, but mainly extend the qi tests so to cover one other code path that is pretty easy to hit and deserves testing (auto deletion of full segment files).
Diffstat (limited to 'src')
-rw-r--r--src/rabbit_queue_index.erl20
-rw-r--r--src/rabbit_tests.erl20
2 files changed, 30 insertions, 10 deletions
diff --git a/src/rabbit_queue_index.erl b/src/rabbit_queue_index.erl
index 7c317b3046..e0634bee55 100644
--- a/src/rabbit_queue_index.erl
+++ b/src/rabbit_queue_index.erl
@@ -524,7 +524,6 @@ read_and_prune_segments(State = #qistate { dir = Dir }) ->
{TotalMsgCount, State1}.
scatter_journal(TotalMsgCount, State = #qistate { dir = Dir }) ->
- JournalPath = filename:join(Dir, ?ACK_JOURNAL_FILENAME),
{Hdl, State1 = #qistate { journal_ack_dict = JAckDict }} =
get_journal_handle(State),
%% ADict may well contain duplicates. However, this is ok, due to
@@ -533,9 +532,15 @@ scatter_journal(TotalMsgCount, State = #qistate { dir = Dir }) ->
State2 = close_handle(journal, State1),
{TotalMsgCount1, State3} =
dict:fold(fun replay_journal_acks_to_segment/3,
- {TotalMsgCount, State2}, ADict),
+ {TotalMsgCount,
+ %% supply empty dict so that when
+ %% replay_journal_acks_to_segment loads segments,
+ %% it gets all msgs, and ignores anything we've
+ %% found in the journal.
+ State2 #qistate { journal_ack_dict = dict:new() }}, ADict),
+ JournalPath = filename:join(Dir, ?ACK_JOURNAL_FILENAME),
ok = file:delete(JournalPath),
- {TotalMsgCount1, State3 #qistate { journal_ack_dict = dict:new() }}.
+ {TotalMsgCount1, State3}.
load_journal(Hdl, ADict) ->
case file_handle_cache:read(Hdl, ?SEQ_BYTES) of
@@ -547,18 +552,13 @@ load_journal(Hdl, ADict) ->
replay_journal_acks_to_segment(_, [], Acc) ->
Acc;
replay_journal_acks_to_segment(SegNum, Acks, {TotalMsgCount, State}) ->
- %% supply empty dict so that we get all msgs in SDict that have
- %% not been acked in the segment file itself
- {SDict, _AckCount, _HighRelSeq, State1} =
- load_segment(SegNum, State #qistate { journal_ack_dict = dict:new() }),
+ {SDict, _AckCount, _HighRelSeq, State1} = load_segment(SegNum, State),
ValidRelSeqIds = dict:fetch_keys(SDict),
ValidAcks = sets:to_list(sets:intersection(sets:from_list(ValidRelSeqIds),
sets:from_list(Acks))),
%% ValidAcks will not contain any duplicates at this point.
- State2 =
- State1 #qistate { journal_ack_dict = State #qistate.journal_ack_dict },
{TotalMsgCount - length(ValidAcks),
- append_acks_to_segment(SegNum, ValidAcks, State2)}.
+ append_acks_to_segment(SegNum, ValidAcks, State1)}.
drop_and_deliver(SegNum, SDict, CleanShutdown, State) ->
{AckMe, DeliverMe} =
diff --git a/src/rabbit_tests.erl b/src/rabbit_tests.erl
index 56dd3483f8..3bf6dd36f1 100644
--- a/src/rabbit_tests.erl
+++ b/src/rabbit_tests.erl
@@ -1084,4 +1084,24 @@ test_queue_index() ->
{0, Qi20} = rabbit_queue_index:init(test_queue()),
_Qi21 = rabbit_queue_index:terminate_and_erase(Qi20),
ok = stop_msg_store(),
+ ok = empty_test_queue(),
+ %% this next bit is just to hit the auto deletion of segment files
+ SeqIdsC = lists:seq(1,65536),
+ {0, Qi22} = rabbit_queue_index:init(test_queue()),
+ {Qi23, _SeqIdsMsgIdsC} = queue_index_publish(SeqIdsC, false, Qi22),
+ Qi24 = lists:foldl(
+ fun (SeqId, QiN) ->
+ rabbit_queue_index:write_delivered(SeqId, QiN)
+ end, Qi23, SeqIdsC),
+ Qi25 = rabbit_queue_index:write_acks(SeqIdsC, Qi24),
+ {_Oks, {false, Qi26}} =
+ rabbit_misc:unfold(
+ fun ({true, QiN}) ->
+ QiM = rabbit_queue_index:flush_journal(QiN),
+ {true, ok, {rabbit_queue_index:can_flush_journal(QiM), QiM}};
+ ({false, _QiN}) ->
+ false
+ end, {true, Qi25}),
+ _Qi27 = rabbit_queue_index:terminate_and_erase(Qi26),
+ ok = stop_msg_store(),
passed.