]> git-server-git.apps.pok.os.sepia.ceph.com Git - ceph-ci.git/commitdiff
crimson/os/seastore/journal: refactor SegmentedJournal::replay() to
authorXuehan Xu <xuxuehan@qianxin.com>
Fri, 3 Jul 2026 14:05:54 +0000 (22:05 +0800)
committerXuehan Xu <xuxuehan@qianxin.com>
Fri, 10 Jul 2026 14:24:38 +0000 (22:24 +0800)
adapt to Journal::scan_valid_record_delta()

Signed-off-by: Xuehan Xu <xuxuehan@qianxin.com>
src/crimson/os/seastore/journal/segmented_journal.cc
src/crimson/os/seastore/journal/segmented_journal.h

index 664110333b5807db3f3491d6e35b334874d1d0bb..130054f0fd1a813626f5aeae1a838a48f2627f1f 100644 (file)
@@ -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<bool> {
+    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,
index 05c191b6c25256c792e4979e4cf6a0790421093f..2b7135b2711fdd0f257f8f4d6370f9533e88e620 100644 (file)
@@ -83,10 +83,7 @@ private:
   WritePipeline* write_pipeline = nullptr;
 
   /// return ordered vector of segments to replay
-  using replay_segments_t = std::vector<
-    std::pair<journal_seq_t, segment_header_t>>;
-  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<std::pair<segment_id_t, segment_header_t>> 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<journal_seq_t, segment_header_t>>;
+  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;
 };
 
 }