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 {}",
#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"
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;
// 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;
#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"
}
}
+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);
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();
};
/**
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)