From: Xuehan Xu Date: Fri, 3 Jul 2026 14:05:54 +0000 (+0800) Subject: crimson/os/seastore/journal: refactor SegmentedJournal::replay() to X-Git-Url: http://git-server-git.apps.pok.os.sepia.ceph.com/?a=commitdiff_plain;h=ae0c6a9a906c33e51ff76d7e787243de3cff6189;p=ceph-ci.git crimson/os/seastore/journal: refactor SegmentedJournal::replay() to adapt to Journal::scan_valid_record_delta() Signed-off-by: Xuehan Xu --- diff --git a/src/crimson/os/seastore/journal/segmented_journal.cc b/src/crimson/os/seastore/journal/segmented_journal.cc index 664110333b5..130054f0fd1 100644 --- a/src/crimson/os/seastore/journal/segmented_journal.cc +++ b/src/crimson/os/seastore/journal/segmented_journal.cc @@ -10,6 +10,7 @@ #include "segmented_journal.h" #include "crimson/common/config_proxy.h" +#include "crimson/common/coroutine.h" #include "crimson/os/seastore/logging.h" SET_SUBSYS(seastore_journal); @@ -108,8 +109,8 @@ SegmentedJournal::prep_replay_segments( return scan_last_segment(last_segment_id, last_header ).safe_then([this, FNAME, segments=std::move(segments)] { INFO("dirty_tail={}, alloc_tail={}", - trimmer.get_dirty_tail(), - trimmer.get_alloc_tail()); + get_dirty_tail(), + get_alloc_tail()); auto journal_tail = trimmer.get_journal_tail(); auto journal_tail_paddr = journal_tail.offset; ceph_assert(journal_tail != JOURNAL_SEQ_NULL); @@ -129,9 +130,9 @@ SegmentedJournal::prep_replay_segments( auto num_segments = segments.end() - from; INFO("{} segments to replay", num_segments); - auto ret = replay_segments_t(num_segments); + replay_segments.resize(num_segments); std::transform( - from, segments.end(), ret.begin(), + from, segments.end(), replay_segments.begin(), [this](const auto &p) { auto ret = journal_seq_t{ p.second.segment_seq, @@ -141,10 +142,9 @@ SegmentedJournal::prep_replay_segments( }; return std::make_pair(ret, p.second); }); - ret[0].first.offset = journal_tail_paddr; + replay_segments[0].first.offset = journal_tail_paddr; return prep_replay_segments_fut( - replay_ertr::ready_future_marker{}, - std::move(ret)); + replay_ertr::ready_future_marker{}); }); } @@ -229,15 +229,14 @@ SegmentedJournal::replay_ertr::future<> SegmentedJournal::replay_segment( journal_seq_t seq, segment_header_t header, - delta_handler_t &handler, - replay_stats_t &stats) + scan_delta_handler_t &handler) { LOG_PREFIX(Journal::replay_segment); INFO("starting at {} -- {}", seq, header); return seastar::do_with( scan_valid_records_cursor(seq), SegmentManagerGroup::found_record_handler_t( - [&handler, this, &stats]( + [&handler, this]( record_locator_t locator, const record_group_header_t& header, const bufferlist& mdbuf) @@ -259,16 +258,14 @@ SegmentedJournal::replay_segment( [write_result=locator.write_result, this, FNAME, - &handler, - &stats](auto& record_deltas_list) + &handler](auto& record_deltas_list) { return crimson::do_for_each( record_deltas_list, [write_result, this, FNAME, - &handler, - &stats](record_deltas_t& record_deltas) + &handler](record_deltas_t& record_deltas) { ++stats.num_records; auto locator = record_locator_t{ @@ -281,30 +278,11 @@ SegmentedJournal::replay_segment( return crimson::do_for_each( record_deltas.deltas, [locator, - this, - &handler, - &stats](auto &p) + &handler](auto &p) { auto& modify_time = p.first; auto& delta = p.second; - return handler( - locator, - delta, - trimmer.get_dirty_tail(), - trimmer.get_alloc_tail(), - modify_time - ).safe_then([&stats, delta_type=delta.type](auto ret) { - auto [is_applied, ext] = ret; - if (is_applied) { - // see Cache::replay_delta() - assert(delta_type != extent_types_t::JOURNAL_TAIL); - if (delta_type == extent_types_t::ALLOC_INFO) { - ++stats.num_alloc_deltas; - } else { - ++stats.num_dirty_deltas; - } - } - }); + return handler(locator, delta, modify_time).discard_result(); }); }); }); @@ -325,6 +303,17 @@ SegmentedJournal::replay_segment( ); } +SegmentedJournal::replay_ret +SegmentedJournal::scan_valid_record_delta( + scan_delta_handler_t &&delta_handler, + journal_seq_t tail) +{ + auto handler = std::move(delta_handler); + for (auto &[seq, header] : replay_segments) { + co_await replay_segment(seq, header, handler); + } +} + SegmentedJournal::replay_ret SegmentedJournal::replay( delta_handler_t &&delta_handler) { @@ -332,11 +321,32 @@ SegmentedJournal::replay_ret SegmentedJournal::replay( auto handler = std::move(delta_handler); auto segment_headers = co_await sm_group.find_journal_segment_headers(); INFO("got {} segments", segment_headers.size()); - replay_stats_t stats; - auto segments = co_await prep_replay_segments(std::move(segment_headers)); - for (auto &[seq, header] : segments) { - co_await replay_segment(seq, header, handler, stats); - } + co_await prep_replay_segments(std::move(segment_headers)); + auto d_handler = [&handler, this]( + const record_locator_t &locator, + const delta_info_t &delta, + sea_time_point modify_time) -> replay_ertr::future { + auto ret = co_await handler( + locator, + delta, + get_dirty_tail(), + get_alloc_tail(), + modify_time); + auto [is_applied, ext] = ret; + if (is_applied) { + // see Cache::replay_delta() + assert(delta.type != extent_types_t::JOURNAL_TAIL); + if (delta.type == extent_types_t::ALLOC_INFO) { + ++stats.num_alloc_deltas; + } else { + ++stats.num_dirty_deltas; + } + } + co_return true; + }; + journal_seq_t tail = get_dirty_tail() <= get_alloc_tail() ? + get_dirty_tail() : get_alloc_tail(); + co_return co_await scan_valid_record_delta(std::move(d_handler), tail); INFO("replay done, record_groups={}, records={}, " "alloc_deltas={}, dirty_deltas={}", stats.num_record_groups, diff --git a/src/crimson/os/seastore/journal/segmented_journal.h b/src/crimson/os/seastore/journal/segmented_journal.h index 05c191b6c25..2b7135b2711 100644 --- a/src/crimson/os/seastore/journal/segmented_journal.h +++ b/src/crimson/os/seastore/journal/segmented_journal.h @@ -83,10 +83,7 @@ private: WritePipeline* write_pipeline = nullptr; /// return ordered vector of segments to replay - using replay_segments_t = std::vector< - std::pair>; - using prep_replay_segments_fut = replay_ertr::future< - replay_segments_t>; + using prep_replay_segments_fut = replay_ertr::future<>; prep_replay_segments_fut prep_replay_segments( std::vector> segments); @@ -100,15 +97,18 @@ private: std::size_t num_records = 0; std::size_t num_alloc_deltas = 0; std::size_t num_dirty_deltas = 0; - }; + } stats; + + using replay_segments_t = std::vector< + std::pair>; + replay_segments_t replay_segments; /// replays records starting at start through end of segment replay_ertr::future<> replay_segment( journal_seq_t start, ///< [in] starting addr, seq segment_header_t header, ///< [in] segment header - delta_handler_t &delta_handler, ///< [in] processes deltas in order - replay_stats_t &stats ///< [out] replay stats + scan_delta_handler_t &delta_handler ///< [in] processes deltas in order ); journal_seq_t get_dirty_tail() const final { @@ -121,9 +121,7 @@ private: replay_ret scan_valid_record_delta( scan_delta_handler_t &&delta_handler, - journal_seq_t tail) final { - return replay_ertr::now(); - } + journal_seq_t tail) final; }; }