]> git-server-git.apps.pok.os.sepia.ceph.com Git - ceph.git/commitdiff
crimson/osd: coroutinize OSD::start()
authorKefu Chai <k.chai@proxmox.com>
Mon, 8 Jun 2026 08:24:40 +0000 (16:24 +0800)
committerKefu Chai <k.chai@proxmox.com>
Tue, 9 Jun 2026 07:30:29 +0000 (15:30 +0800)
OSD::start() is a long, deeply nested .then() continuation chain. let's
rewrite it as a coroutine to make it readable. the chain was already
sequential and the one concurrent step keeps its when_all_succeed(), so
the rewrite preserves both ordering and concurrency.

start() runs once at boot, off the i/o path, so the small overhead of
co_await over a hand-rolled continuation chain is a fine price for the
readability.

Signed-off-by: Kefu Chai <k.chai@proxmox.com>
src/crimson/osd/osd.cc
src/crimson/osd/osd.h

index 0f94d16bdfc8523517dbb7738ef6934fe4a0a70d..1404a85bfc7bb19cdb276cf90e9d17501700f618 100644 (file)
@@ -461,6 +461,27 @@ namespace {
   }
 }
 
+seastar::future<> OSD::report_osd_stats()
+{
+  LOG_PREFIX(OSD::report_osd_stats);
+  co_await shard_services.invoke_on_all(
+    [this](auto &local_service) {
+      auto stats = local_service.report_stats();
+      shard_stats[seastar::this_shard_id()] = stats;
+    });
+  std::ostringstream oss;
+  double agg_ru = 0;
+  int cnt = 0;
+  for (const auto &stats : shard_stats) {
+    agg_ru += stats.reactor_utilization;
+    ++cnt;
+    oss << int(stats.reactor_utilization);
+    oss << ",";
+  }
+  INFO("reactor_utilizations: {}({})",
+       int(agg_ru/cnt), oss.str());
+}
+
 seastar::future<> OSD::start()
 {
   LOG_PREFIX(OSD::start);
@@ -473,175 +494,142 @@ seastar::future<> OSD::start()
   startup_time = ceph::mono_clock::now();
   ceph_assert(seastar::this_shard_id() == PRIMARY_CORE);
   DEBUG("starting store");
-  return store.start().then([this] (auto store_shard_nums) {
-    return pg_to_shard_mappings.start(0, seastar::smp::count, store_shard_nums
-    ).then([this] {
-      return osd_singleton_state.start_single(
+  uint32_t store_shards_num = co_await store.start();
+  co_await pg_to_shard_mappings.start(0, seastar::smp::count, store_shards_num);
+  co_await osd_singleton_state.start_single(
         whoami, std::ref(*cluster_msgr), std::ref(*public_msgr),
         std::ref(*monc), std::ref(*mgrc));
-    }).then([this] {
-      return osd_states.start();
-    }).then([this, store_shard_nums] {
-      ceph::mono_time startup_time = ceph::mono_clock::now();
-      return shard_services.start(
+  co_await osd_states.start();
+  ceph::mono_time startup_time = ceph::mono_clock::now();
+  co_await shard_services.start(
         std::ref(osd_singleton_state),
         std::ref(pg_to_shard_mappings),
-        store_shard_nums,
+        store_shards_num,
         whoami,
         startup_time,
         osd_singleton_state.local().perf,
         osd_singleton_state.local().recoverystate_perf,
         std::ref(store),
         std::ref(osd_states));
-    });
-  }).then([this, FNAME] {
-    heartbeat.reset(new Heartbeat{
-       whoami, get_shard_services(),
-       *monc, *hb_front_msgr, *hb_back_msgr});
-    DEBUG("mounting store");
-    return store.mount().handle_error(
+  heartbeat = std::make_unique<Heartbeat>(
+    whoami, get_shard_services(),
+    *monc, *hb_front_msgr, *hb_back_msgr);
+  DEBUG("mounting store");
+  co_await store.mount().handle_error(
       crimson::stateful_ec::assert_failure(fmt::format(
         "{} error mounting object store in {}",
         FNAME, local_conf().get_val<std::string>("osd_data")).c_str())
     );
-  }).then([this, FNAME] {
-    auto stats_seconds = local_conf().get_val<int64_t>("crimson_osd_stat_interval");
-    if (stats_seconds > 0) {
-      shard_stats.resize(seastar::smp::count);
-      stats_timer.set_callback([this, FNAME] {
-        gate.dispatch_in_background("stats_osd", *this, [this, FNAME] {
-          return shard_services.invoke_on_all(
-            [this](auto &local_service) {
-            auto stats = local_service.report_stats();
-            shard_stats[seastar::this_shard_id()] = stats;
-          }).then([this, FNAME] {
-            std::ostringstream oss;
-            double agg_ru = 0;
-            int cnt = 0;
-            for (const auto &stats : shard_stats) {
-              agg_ru += stats.reactor_utilization;
-              ++cnt;
-              oss << int(stats.reactor_utilization);
-              oss << ",";
-            }
-            INFO("reactor_utilizations: {}({})",
-                 int(agg_ru/cnt), oss.str());
-          });
-        });
-        gate.dispatch_in_background("stats_store", *this, [this] {
-          return store.report_stats();
-        });
+  auto stats_seconds = local_conf().get_val<int64_t>("crimson_osd_stat_interval");
+  if (stats_seconds > 0) {
+    shard_stats.resize(seastar::smp::count);
+    stats_timer.set_callback([this] {
+      gate.dispatch_in_background("stats_osd", *this, [this] {
+        return report_osd_stats();
       });
-      stats_timer.arm_periodic(std::chrono::seconds(stats_seconds));
-    }
+      gate.dispatch_in_background("stats_store", *this, [this] {
+        return store.report_stats();
+      });
+    });
+    stats_timer.arm_periodic(std::chrono::seconds(stats_seconds));
+  }
+
+  DEBUG("open metadata collection");
+  co_await open_meta_coll();
 
-    DEBUG("open metadata collection");
-    return open_meta_coll();
-  }).then([this, FNAME] {
-    DEBUG("loading superblock"); 
-    return pg_shard_manager.get_meta_coll().load_superblock(
+  DEBUG("loading superblock");
+  superblock = co_await pg_shard_manager.get_meta_coll().load_superblock(
     ).handle_error(
       crimson::ct_error::assert_all("open_meta_coll error")
     );
-  }).then([this](OSDSuperblock&& sb) {
-    superblock = std::move(sb);
-    if (!superblock.cluster_osdmap_trim_lower_bound) {
-      superblock.cluster_osdmap_trim_lower_bound = superblock.get_oldest_map();
-    }
-    return pg_shard_manager.set_superblock(superblock);
-  }).then([this] {
-    return pg_shard_manager.get_local_map(superblock.current_epoch);
-  }).then([this](OSDMapService::local_cached_map_t&& map) {
-    osdmap = make_local_shared_foreign(OSDMapService::local_cached_map_t(map));
-    return pg_shard_manager.update_map(std::move(map));
-  }).then([this] {
-    return shard_services.invoke_on_all([this](auto &local_service) {
-      local_service.local_state.osdmap_gate.got_map(osdmap->get_epoch());
-    });
-  }).then([this, FNAME] {
-    bind_epoch = osdmap->get_epoch();
-    DEBUG("loading PGs");
-    return pg_shard_manager.load_pgs(store);
-  }).then([this, FNAME] {
-    uint64_t osd_required =
-      CEPH_FEATURE_UID |
-      CEPH_FEATURE_PGID64 |
-      CEPH_FEATURE_OSDENC;
-    using crimson::net::SocketPolicy;
-
-    public_msgr->set_default_policy(SocketPolicy::stateless_server(0));
-    public_msgr->set_policy(entity_name_t::TYPE_MON,
-                            SocketPolicy::lossy_client(osd_required));
-    public_msgr->set_policy(entity_name_t::TYPE_MGR,
-                            SocketPolicy::lossy_client(osd_required));
-    public_msgr->set_policy(entity_name_t::TYPE_OSD,
-                            SocketPolicy::stateless_server(0));
-
-    cluster_msgr->set_default_policy(SocketPolicy::stateless_server(0));
-    cluster_msgr->set_policy(entity_name_t::TYPE_MON,
-                             SocketPolicy::lossy_client(0));
-    cluster_msgr->set_policy(entity_name_t::TYPE_OSD,
-                             SocketPolicy::lossless_peer(osd_required));
-    cluster_msgr->set_policy(entity_name_t::TYPE_CLIENT,
-                             SocketPolicy::stateless_server(0));
-
-    crimson::net::dispatchers_t dispatchers{this, monc.get(), mgrc.get()};
-    return seastar::when_all_succeed(
-      cluster_msgr->bind(pick_addresses(CEPH_PICK_ADDRESS_CLUSTER))
-        .safe_then([this, dispatchers]() mutable {
-         return cluster_msgr->start(dispatchers);
-        }, crimson::net::Messenger::bind_ertr::assert_all_func(
-            [FNAME] (const std::error_code& e) {
-          ERROR("cluster messenger bind(): {}", e);
-        })),
-      public_msgr->bind(pick_addresses(CEPH_PICK_ADDRESS_PUBLIC))
-        .safe_then([this, dispatchers]() mutable {
-         return public_msgr->start(dispatchers);
-        }, crimson::net::Messenger::bind_ertr::assert_all_func(
-            [FNAME] (const std::error_code& e) {
-          ERROR("public messenger bind(): {}", e);
-        })));
-  }).then_unpack([this, FNAME] {
+  if (!superblock.cluster_osdmap_trim_lower_bound) {
+    superblock.cluster_osdmap_trim_lower_bound = superblock.get_oldest_map();
+  }
+  co_await pg_shard_manager.set_superblock(superblock);
+  auto local_map = co_await pg_shard_manager.get_local_map(superblock.current_epoch);
+  osdmap = make_local_shared_foreign(OSDMapService::local_cached_map_t(local_map));
+  co_await pg_shard_manager.update_map(std::move(local_map));
+  co_await shard_services.invoke_on_all([this](auto &local_service) {
+    local_service.local_state.osdmap_gate.got_map(osdmap->get_epoch());
+  });
+
+  bind_epoch = osdmap->get_epoch();
+  DEBUG("loading PGs");
+  co_await pg_shard_manager.load_pgs(store);
+
+  uint64_t osd_required =
+    CEPH_FEATURE_UID |
+    CEPH_FEATURE_PGID64 |
+    CEPH_FEATURE_OSDENC;
+  using crimson::net::SocketPolicy;
+
+  public_msgr->set_default_policy(SocketPolicy::stateless_server(0));
+  public_msgr->set_policy(entity_name_t::TYPE_MON,
+                          SocketPolicy::lossy_client(osd_required));
+  public_msgr->set_policy(entity_name_t::TYPE_MGR,
+                          SocketPolicy::lossy_client(osd_required));
+  public_msgr->set_policy(entity_name_t::TYPE_OSD,
+                          SocketPolicy::stateless_server(0));
+
+  cluster_msgr->set_default_policy(SocketPolicy::stateless_server(0));
+  cluster_msgr->set_policy(entity_name_t::TYPE_MON,
+                           SocketPolicy::lossy_client(0));
+  cluster_msgr->set_policy(entity_name_t::TYPE_OSD,
+                           SocketPolicy::lossless_peer(osd_required));
+  cluster_msgr->set_policy(entity_name_t::TYPE_CLIENT,
+                           SocketPolicy::stateless_server(0));
+
+  crimson::net::dispatchers_t dispatchers{this, monc.get(), mgrc.get()};
+  co_await seastar::when_all_succeed(
+    cluster_msgr->bind(pick_addresses(CEPH_PICK_ADDRESS_CLUSTER))
+    .safe_then([this, dispatchers]() mutable {
+      return cluster_msgr->start(dispatchers);
+    }, crimson::net::Messenger::bind_ertr::assert_all_func(
+      [FNAME] (const std::error_code& e) {
+        ERROR("cluster messenger bind(): {}", e);
+      })),
+    public_msgr->bind(pick_addresses(CEPH_PICK_ADDRESS_PUBLIC))
+    .safe_then([this, dispatchers]() mutable {
+      return public_msgr->start(dispatchers);
+    }, crimson::net::Messenger::bind_ertr::assert_all_func(
+      [FNAME] (const std::error_code& e) {
+        ERROR("public messenger bind(): {}", e);
+      }))
+  ).then_unpack([this, FNAME] {
     DEBUG("starting mon and mgr clients");
     return seastar::when_all_succeed(monc->start(),
                                      mgrc->start());
   }).then_unpack([this, FNAME] {    
     DEBUG("adding to crush");
     return _add_me_to_crush();
-  }).then([this] {
-    return _add_device_class();
-   }).then([this] {
-    if (is_rotational.has_value()) {
-      return shard_services.invoke_on_all([this](auto &local_service) {
-        local_service.local_state.initialize_scheduler(local_service.get_cct(), *is_rotational);
-      });
-    } else {
-      throw std::runtime_error("No device class is set");
-    }
-   }).then([this] {
-    monc->sub_want("osd_pg_creates", last_pg_create_epoch, 0);
-    monc->sub_want("mgrmap", 0, 0);
-    monc->sub_want("osdmap", 0, 0);
-    return monc->renew_subs();
-  }).then([FNAME, this] {
-    if (auto [addrs, changed] =
-        replace_unknown_addrs(cluster_msgr->get_myaddrs(),
-                              public_msgr->get_myaddrs()); changed) {
-      DEBUG("replacing unkwnown addrs of cluster messenger");
-      cluster_msgr->set_myaddrs(addrs);
-    }
-    return heartbeat->start(pick_addresses(CEPH_PICK_ADDRESS_PUBLIC),
-                            pick_addresses(CEPH_PICK_ADDRESS_CLUSTER));
-  }).then([this] {
-    // create the admin-socket server, and the objects that register
-    // to handle incoming commands
-    return start_asok_admin();
-  }).then([this] {
-    return log_client.set_fsid(monc->get_fsid());
-  }).then([this, FNAME] {
-    DEBUG("starting boot");
-    return start_boot();
   });
+  co_await _add_device_class();
+  if (is_rotational.has_value()) {
+    co_await shard_services.invoke_on_all([this](auto &local_service) {
+      local_service.local_state.initialize_scheduler(local_service.get_cct(), *is_rotational);
+    });
+  } else {
+    throw std::runtime_error("No device class is set");
+  }
+  monc->sub_want("osd_pg_creates", last_pg_create_epoch, 0);
+  monc->sub_want("mgrmap", 0, 0);
+  monc->sub_want("osdmap", 0, 0);
+  co_await monc->renew_subs();
+
+  if (auto [addrs, changed] =
+      replace_unknown_addrs(cluster_msgr->get_myaddrs(),
+                            public_msgr->get_myaddrs()); changed) {
+    DEBUG("replacing unkwnown addrs of cluster messenger");
+    cluster_msgr->set_myaddrs(addrs);
+  }
+  co_await heartbeat->start(pick_addresses(CEPH_PICK_ADDRESS_PUBLIC),
+                            pick_addresses(CEPH_PICK_ADDRESS_CLUSTER));
+  // create the admin-socket server, and the objects that register
+  // to handle incoming commands
+  co_await start_asok_admin();
+  co_await log_client.set_fsid(monc->get_fsid());
+  DEBUG("starting boot");
+  co_await start_boot();
 }
 
 seastar::future<> OSD::start_boot()
index 703ccd3ae5ea5b74d6722dd104c5e0723ad1afb8..01b31295bcc57ef8ea6f3a6f23f40b59c29aeb73 100644 (file)
@@ -134,6 +134,10 @@ class OSD final : public crimson::net::Dispatcher,
 
   seastar::timer<seastar::lowres_clock> stats_timer;
   std::vector<ShardServices::shard_stats_t> shard_stats;
+  // collect per-shard stats and log reactor utilization; a member coroutine
+  // so its captures (this) live in the frame, since it runs as a detached
+  // gated task a capturing lambda coroutine would dangle (see report_osd_stats).
+  seastar::future<> report_osd_stats();
 
   std::vector<std::string> get_tracked_keys() const noexcept final;
   void handle_conf_change(const ConfigProxy& conf,