]> git-server-git.apps.pok.os.sepia.ceph.com Git - ceph.git/commitdiff
crimson/osd: Add functions to notify mon when PGs are ready to merge
authorAishwarya Mathuria <amathuri@redhat.com>
Wed, 7 Jan 2026 11:55:25 +0000 (11:55 +0000)
committerAishwarya Mathuria <amathuri@redhat.com>
Thu, 11 Jun 2026 04:49:14 +0000 (10:19 +0530)
When a PG is in the pending merge state it is >= pg_num_pending and <
pg_num. When this happens, IO is paused and once the PG peers we notify
the mon that we are idle and safe to merge.
Use Gated for merge notify callbacks.

Signed-off-by: Aishwarya Mathuria <amathuri@redhat.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
src/crimson/osd/osd.cc
src/crimson/osd/pg.h
src/crimson/osd/shard_services.cc
src/crimson/osd/shard_services.h

index 0f94d16bdfc8523517dbb7738ef6934fe4a0a70d..8e6d048be1d24d9459c224c7a394ec5148fec465 100644 (file)
@@ -1330,6 +1330,10 @@ seastar::future<> OSD::committed_osd_maps(
     old_map = osdmap;
   }
 
+  // Drop merge-ready bookkeeping for PGs removed by the map we just
+  // consumed, mirroring classic OSD::consume_map().
+  co_await get_shard_services().prune_sent_ready_to_merge();
+
   if (osdmap->is_up(whoami)) {
     const auto up_from = osdmap->get_up_from(whoami);
     INFO("osd.{}: map e {} marked me up: up_from {}, bind_epoch {}, state {}",
index e24427a6a0de762360d537d44d0e71703086627e..f3913c919df2c0a0707086d32f16e83dbd9c82ce 100644 (file)
@@ -26,6 +26,7 @@
 #include "osd/DynamicPerfStats.h"
 
 #include "crimson/common/interruptible_future.h"
+#include "crimson/common/gated.h"
 #include "crimson/common/log.h"
 #include "crimson/common/type_helpers.h"
 #include "crimson/os/futurized_collection.h"
@@ -563,12 +564,51 @@ public:
   std::pair<ghobject_t, bool>
   do_delete_work(ceph::os::Transaction &t, ghobject_t _next) final;
 
-  // merge/split not ready
-  void clear_ready_to_merge() final {}
-  void set_not_ready_to_merge_target(pg_t pgid, pg_t src) final {}
-  void set_not_ready_to_merge_source(pg_t pgid) final {}
-  void set_ready_to_merge_target(eversion_t lu, epoch_t les, epoch_t lec) final {}
-  void set_ready_to_merge_source(eversion_t lu) final {}
+  void clear_ready_to_merge() final {
+    LOG_PREFIX(PG::clear_ready_to_merge);
+    SUBDEBUGDPP(osd, "", *this);
+    merge_notify_gate.dispatch_in_background(
+      "clear_ready_to_merge", *this,
+      [this] {
+        return shard_services.clear_ready_to_merge(pgid.pgid);
+      });
+  }
+  void set_not_ready_to_merge_target(pg_t pgid, pg_t src) final {
+    LOG_PREFIX(PG::set_not_ready_to_merge_target);
+    SUBDEBUGDPP(osd, "", *this);
+    merge_notify_gate.dispatch_in_background(
+      "set_not_ready_to_merge_target", *this,
+      [this, pgid, src] {
+        return shard_services.set_not_ready_to_merge_target(pgid, src);
+      });
+  }
+  void set_not_ready_to_merge_source(pg_t pgid) final {
+    LOG_PREFIX(PG::set_not_ready_to_merge_source);
+    SUBDEBUGDPP(osd, "", *this);
+    merge_notify_gate.dispatch_in_background(
+      "set_not_ready_to_merge_source", *this,
+      [this, pgid] {
+        return shard_services.set_not_ready_to_merge_source(pgid);
+      });
+  }
+  void set_ready_to_merge_target(eversion_t lu, epoch_t les, epoch_t lec) final {
+    LOG_PREFIX(PG::set_ready_to_merge_target);
+    SUBDEBUGDPP(osd, "", *this);
+    merge_notify_gate.dispatch_in_background(
+      "set_ready_to_merge_target", *this,
+      [this, lu, les, lec] {
+        return shard_services.set_ready_to_merge_target(pgid.pgid, lu, les, lec);
+      });
+  }
+  void set_ready_to_merge_source(eversion_t lu) final {
+    LOG_PREFIX(PG::set_ready_to_merge_source);
+    SUBDEBUGDPP(osd, "", *this);
+    merge_notify_gate.dispatch_in_background(
+      "set_ready_to_merge_source", *this,
+      [this, lu] {
+        return shard_services.set_ready_to_merge_source(pgid.pgid, lu);
+      });
+  }
 
   void on_active_actmap() final;
   void on_active_advmap(const OSDMapRef &osdmap) final;
@@ -1148,6 +1188,11 @@ private:
   // continuations here.
   bool stopping = false;
 
+  // PeeringListener merge callbacks must remain void, but they trigger async
+  // mon notifies in ShardServices. Gate them here so failures are logged and
+  // PG::stop() waits for them to drain.
+  crimson::common::Gated merge_notify_gate;
+
   PGActivationBlocker wait_for_active_blocker;
   PglogBasedRecovery* pglog_based_recovery_op = nullptr;
 
index 3780cd626c643f66af49661d9f61888e991e1222..e4ff7d02a423a2122c65e15eba450de7b344aae4 100644 (file)
@@ -9,6 +9,7 @@
 #include "messages/MOSDMap.h"
 #include "messages/MOSDPGCreated.h"
 #include "messages/MOSDPGTemp.h"
+#include "messages/MOSDPGReadyToMerge.h"
 
 #include "osd/osd_perf_counters.h"
 #include "osd/PeeringState.h"
@@ -308,6 +309,150 @@ void OSDSingletonState::prune_pg_created()
   }
 }
 
+seastar::future<> OSDSingletonState::set_ready_to_merge_source(pg_t pgid,
+                                                  eversion_t version)
+{
+  LOG_PREFIX(OSDSingletonState::set_ready_to_merge_source);
+  DEBUG("{}", pgid);
+  ready_to_merge_source[pgid] = version;
+  ceph_assert(!not_ready_to_merge_source.contains(pgid));
+  return send_ready_to_merge();
+}
+
+seastar::future<> OSDSingletonState::set_ready_to_merge_target(pg_t pgid,
+                                           eversion_t version,
+                                           epoch_t last_epoch_started,
+                                           epoch_t last_epoch_clean)
+{
+  LOG_PREFIX(OSDSingletonState::set_ready_to_merge_target);
+  DEBUG("{}", pgid);
+  ready_to_merge_target.insert(std::make_pair(pgid,
+                                         std::make_tuple(version,
+                                                    last_epoch_started,
+                                                    last_epoch_clean)));
+  ceph_assert(!not_ready_to_merge_target.contains(pgid));
+  return send_ready_to_merge();
+}
+
+seastar::future<> OSDSingletonState::set_not_ready_to_merge_source(pg_t source)
+{
+  LOG_PREFIX(OSDSingletonState::set_not_ready_to_merge_source);
+  DEBUG("{}", source);
+  not_ready_to_merge_source.insert(source);
+  ceph_assert(!ready_to_merge_source.contains(source));
+  return send_ready_to_merge();
+}
+
+seastar::future<> OSDSingletonState::set_not_ready_to_merge_target(pg_t target, pg_t source)
+{
+  LOG_PREFIX(OSDSingletonState::set_not_ready_to_merge_target);
+  DEBUG("{} source {}", target, source);
+  not_ready_to_merge_target[target] = source;
+  ceph_assert(!ready_to_merge_target.contains(target));
+  return send_ready_to_merge();
+}
+
+seastar::future<> OSDSingletonState::send_ready_to_merge()
+{
+  LOG_PREFIX(OSDSingletonState::send_ready_to_merge);
+  DEBUG(" ready_to_merge_source: {} not_ready_to_merge_source: {} \
+          ready_to_merge_target: {} not_ready_to_merge_target: {} \
+          sent_ready_to_merge_source {}", ready_to_merge_source,
+          not_ready_to_merge_source, ready_to_merge_target, not_ready_to_merge_target,
+          sent_ready_to_merge_source);
+
+  struct ready_to_merge_send_t {
+    pg_t pgid;
+    eversion_t source_version;
+    eversion_t target_version;
+    epoch_t last_epoch_started = 0;
+    epoch_t last_epoch_clean = 0;
+    bool ready = false;
+  };
+
+  const epoch_t map_epoch = osdmap->get_epoch();
+  std::vector<ready_to_merge_send_t> pending;
+  pending.reserve(not_ready_to_merge_source.size() +
+                  not_ready_to_merge_target.size() +
+                  ready_to_merge_source.size());
+
+  for (auto src : not_ready_to_merge_source) {
+    if (!sent_ready_to_merge_source.contains(src)) {
+      sent_ready_to_merge_source.insert(src);
+      pending.push_back({src, {}, {}, 0, 0, false});
+    }
+  }
+  for (auto p : not_ready_to_merge_target) {
+    if (!sent_ready_to_merge_source.contains(p.second)) {
+      sent_ready_to_merge_source.insert(p.second);
+      pending.push_back({p.second, {}, {}, 0, 0, false});
+    }
+  }
+  for (auto& [src_pg, src_version] : ready_to_merge_source) {
+    if (not_ready_to_merge_source.contains(src_pg) ||
+        not_ready_to_merge_target.contains(src_pg.get_parent())) {
+      continue;
+    }
+    auto p = ready_to_merge_target.find(src_pg.get_parent());
+    if (p != ready_to_merge_target.end() &&
+        !sent_ready_to_merge_source.contains(src_pg)) {
+      sent_ready_to_merge_source.insert(src_pg);
+      pending.push_back({
+        src_pg,
+        src_version,
+        std::get<0>(p->second),
+        std::get<1>(p->second),
+        std::get<2>(p->second),
+        true});
+    }
+  }
+
+  if (pending.empty()) {
+    return seastar::now();
+  }
+  return seastar::parallel_for_each(
+    std::move(pending),
+    [this, map_epoch](const ready_to_merge_send_t& m) {
+      return monc.send_message(crimson::make_message<MOSDPGReadyToMerge>(
+        m.pgid,
+        m.source_version,
+        m.target_version,
+        m.last_epoch_started,
+        m.last_epoch_clean,
+        m.ready,
+        map_epoch));
+    });
+}
+
+void OSDSingletonState::clear_ready_to_merge(pg_t pgid)
+{
+  ready_to_merge_source.erase(pgid);
+  ready_to_merge_target.erase(pgid);
+  not_ready_to_merge_source.erase(pgid);
+  not_ready_to_merge_target.erase(pgid);
+  sent_ready_to_merge_source.erase(pgid);
+}
+
+void OSDSingletonState::clear_sent_ready_to_merge()
+{
+  sent_ready_to_merge_source.clear();
+}
+
+void OSDSingletonState::prune_sent_ready_to_merge()
+{
+  LOG_PREFIX(OSDSingletonState::prune_sent_ready_to_merge);
+  auto source = sent_ready_to_merge_source.begin();
+  while (source != sent_ready_to_merge_source.end()) {
+    if (!osdmap->pg_exists(*source)) {
+      DEBUG("{}", *source);
+      source = sent_ready_to_merge_source.erase(source);
+    } else {
+      DEBUG(" exist {}", *source);
+      ++source;
+    }
+  }
+}
+
 seastar::future<> OSDSingletonState::send_alive(const epoch_t want)
 {
   LOG_PREFIX(OSDSingletonState::send_alive);
index 19717ee55ef1d448eef2f5afe9be15cbdab5abfc..7ab1eef6ce15040450f29e368fe453882703f8fd 100644 (file)
@@ -371,6 +371,25 @@ private:
   seastar::future<> store_maps(ceph::os::Transaction& t,
                                epoch_t start, Ref<MOSDMap> m);
   void trim_maps(ceph::os::Transaction& t, OSDSuperblock& superblock);
+
+  // -- PG merging --
+  std::map<pg_t, eversion_t> ready_to_merge_source;
+  std::map<pg_t,std::tuple<eversion_t,epoch_t,epoch_t>> ready_to_merge_target;
+  std::set<pg_t> not_ready_to_merge_source;
+  std::map<pg_t,pg_t> not_ready_to_merge_target;
+  std::set<pg_t> sent_ready_to_merge_source;
+  seastar::future<> set_ready_to_merge_source(pg_t pgid,
+                                 eversion_t version);
+  seastar::future<> set_ready_to_merge_target(pg_t pgid,
+                                 eversion_t version,
+                                 epoch_t last_epoch_started,
+                                 epoch_t last_epoch_clean);
+  seastar::future<> set_not_ready_to_merge_source(pg_t source);
+  seastar::future<> set_not_ready_to_merge_target(pg_t target, pg_t source);
+  void clear_ready_to_merge(pg_t pgid);
+  seastar::future<> send_ready_to_merge();
+  void clear_sent_ready_to_merge();
+  void prune_sent_ready_to_merge();
 };
 
 /**
@@ -650,6 +669,15 @@ public:
     return local_state.ec_extent_cache_lru;
   }
 
+  FORWARD_TO_OSD_SINGLETON(set_ready_to_merge_source)
+  FORWARD_TO_OSD_SINGLETON(set_ready_to_merge_target)
+  FORWARD_TO_OSD_SINGLETON(set_not_ready_to_merge_source)
+  FORWARD_TO_OSD_SINGLETON(set_not_ready_to_merge_target)
+  FORWARD_TO_OSD_SINGLETON(clear_ready_to_merge)
+  FORWARD_TO_OSD_SINGLETON(send_ready_to_merge)
+  FORWARD_TO_OSD_SINGLETON(clear_sent_ready_to_merge)
+  FORWARD_TO_OSD_SINGLETON(prune_sent_ready_to_merge)
+
   FORWARD_TO_OSD_SINGLETON(get_pool_info)
   FORWARD(get_throttle, get_throttle, local_state.throttler)