diff options
| author | Matthew Sackman <matthew@lshift.net> | 2009-10-28 15:25:14 +0000 |
|---|---|---|
| committer | Matthew Sackman <matthew@lshift.net> | 2009-10-28 15:25:14 +0000 |
| commit | bb19594cabc7bc4d6dcafc3afeaddf5ecddc3f76 (patch) | |
| tree | be49b7ced50a6ba058a06cd33be79b457da4f75d /src | |
| parent | eea836fb25d38fecff1b58896bd1a37b8d0a9081 (diff) | |
| download | rabbitmq-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.erl | 20 | ||||
| -rw-r--r-- | src/rabbit_tests.erl | 20 |
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. |
