}
}
+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);
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()