Merge pull request #66885 from amathuria/wip-amat-crimson-merge-support

Crimson PG merging support
This commit is contained in:
Aishwarya Mathuria 2026-06-12 15:18:37 +05:30 committed by GitHub
commit cb95853a1d
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
20 changed files with 1136 additions and 40 deletions

View File

@ -0,0 +1,106 @@
===========================================
PG Merge Synchronization in Crimson
===========================================
This document describes the infrastructure used to coordinate the merging
of Placement Groups (PGs) in Crimson.
Background
----------
Crimson utilizes a shared-nothing memory model where each PG is owned by
a specific CPU core. Merging requires cross-core coordination while
maintaining memory safety and epoch consistency.
Core Concepts
-------------
.. _migration_safety:
Migration Safety (Memory Ownership)
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
PGs are managed by their ``birth_shard``. To prevent cross-shard
deallocation crashes:
1. **Extraction**: ``ShardServices::extract_pg()`` removes the source PG
from its shard map so it stops receiving messages.
2. **Transportation**: The PG crosses cores in a ``seastar::foreign_ptr``.
3. **Storage**: The target shard stores each source in
``crimson::local_shared_foreign_ptr<Ref<PG>>`` inside the target PG's
rendezvous map. When the target drops the last reference after
``PG::merge_from()``, destruction is routed back to the source
``birth_shard``.
.. _synchronization:
Target and Source PG Synchronization
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
Merging involves a race between source PGs checking in and the target PG
reaching the merge point in ``PGAdvanceMap``.
Synchronization state lives on the **target PG**, not in a per-shard
registry:
* **``merge_rendezvous_t``**: A map of arrived sources plus a
``seastar::semaphore`` counting first-time registrations.
* **``PG::add_merge_source()``**: Called on the target's shard when
``ShardServices::register_merge_source()`` delivers a source PG.
Duplicate source pgids are ignored (idempotent under replay).
* **``PG::collect_merge_sources(n)``**: Waits on the rendezvous semaphore
for ``n`` arrivals (one signal per first insert), asserts the map has
``n`` entries, then returns sources and clears the rendezvous. If
``reset_merge_rendezvous()`` breaks the wait, returns an empty map.
* **``PG::reset_merge_rendezvous()``**: Clears in-flight handoffs and
resets the semaphore. Called on ``PG::stop()`` or after a Seastore
cross-shard abort so a failed try cannot leave stale sources for a
later epoch.
**Cross-shard handoff**: ``register_merge_source()`` resolves the target
shard, extracts the source locally, and uses ``invoke_on`` to call
``add_merge_source()`` on the target PG.
.. _map_consistency:
Map Advancement and Pipeline Stalling
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
To ensure the merge occurs between PGs at the same logical time:
* ``PGAdvanceMap::check_for_merges()`` detects ``pg_num`` shrink between
map epochs and returns a ``merge_result_t`` (source, target, or none).
* ``merge_pg()`` performs merge-specific work only (Seastore eligibility,
rendezvous collection, ``PG::merge_from()`` on the target).
* A single ``PGAdvanceMap`` operation may batch multiple map epochs
(``from`` through ``to``). Only a **merge source** stops the epoch
loop early: ``finish_merge_source()`` commits ``rctx``, stops the PG,
and calls ``register_merge_source()``. A **merge target** runs
``merge_from()`` at the merge epoch, then continues advancing through
the remaining epochs up to ``to`` (matching classic OSD behavior).
Normal completion at the end of ``start()`` activates the final map
and commits ``rctx`` (also covering a completed merge target).
* The target does not call ``merge_from()`` until sources frozen at the
merge epoch have arrived on its rendezvous.
.. _seastore_cross_shard:
Seastore Cross-Shard Merges
~~~~~~~~~~~~~~~~~~~~~~~~~~~
Seastore does not support merging collections across reactor shards.
Before a source hands off or a target collects sources, ``merge_pg()``
calls ``ShardServices::seastore_merge_shards_ok()`` with the full sibling
set. If the target and any source map to different shards, the OSD aborts
the attempt: clears ready-to-merge state, sends
``MOSDPGReadyToMerge{ ready=false }`` for the source so the monitor does
not commit the unsafe ``pg_num`` decrement, resets the target rendezvous,
and sends ``MOSDPGStopMerge`` to the monitor (once per pool per OSD).
The monitor clamps ``pg_num_target`` to ``pg_num`` and clears
``FLAG_CRIMSON_ALLOW_PG_MERGE`` for that pool until operators re-enable
merging (``crimson_allow_pg_merge``).
Cleanup
-------
After ``complete_rctx()``, the target releases its ``local_shared_foreign_ptr``
references; source PG objects are destroyed on their ``birth_shard``.
When the OSD finishes consuming a new OSDMap in ``OSD::committed_osd_maps()``,
``OSDSingletonState::prune_sent_ready_to_merge()`` drops ``sent_ready_to_merge_source``
entries (and stale ``pools_merge_stopped_reported`` pool ids) for PGs that no
longer exist in the committed map, mirroring classic ``OSD::consume_map()``.

View File

@ -2,6 +2,8 @@ overrides:
ceph:
fs: xfs
conf:
global:
crimson_allow_pg_merge: true
osd:
# crimson's osd objectstore option
osd objectstore: bluestore

View File

@ -0,0 +1,30 @@
roles:
- - mon.a
- mon.b
- mon.c
- mgr.x
- osd.0
- osd.1
- osd.2
- client.0
openstack:
- volumes: # attached to each instance
count: 3
size: 10 # GB
tasks:
- install:
- ceph:
pre-mgr-commands:
- sudo ceph config set mgr mgr_pool false --force
log-ignorelist:
- but it is still running
- overall HEALTH_
- \(PG_
conf:
osd:
osd min pg log entries: 5
crimson cpu num: 2
- workunit:
clients:
client.0:
- rados/test_pg_merging.sh

View File

@ -0,0 +1,97 @@
#!/bin/bash
. $(dirname $0)/../../standalone/ceph-helpers.sh
set -x
function wait_for_merge_completion() {
local pool=$1
local target=$2
local stall_timeout=60
local actual=0
local last_actual=0
local last_progress_time
wait_for_clean || return 1
last_actual=$(ceph pg ls-by-pool $pool --format=json 2>/dev/null | jq -r ".pg_stats | length")
last_progress_time=$(date +%s)
echo "Waiting for pool $pool to reach $target PGs..."
while true; do
actual=$last_actual
echo "Current PG count: $actual (Target: $target)"
if [ "$actual" -eq "$target" ]; then
break
fi
sleep 2
wait_for_clean || return 1
actual=$(ceph pg ls-by-pool $pool --format=json 2>/dev/null | jq -r ".pg_stats | length")
if [ "$actual" -lt "$last_actual" ]; then
last_progress_time=$(date +%s)
fi
last_actual=$actual
if [ $(( $(date +%s) - last_progress_time )) -ge "$stall_timeout" ]; then
echo "PG merge made no progress for ${stall_timeout}s in pool $pool (still at $actual PGs)"
return 1
fi
done
}
function test_pg_merging() {
local pool="merge_pool"
ceph osd pool delete $pool $pool --yes-i-really-really-mean-it || true
create_pool $pool 16
wait_for_clean || return 1
ceph osd pool set $pool nopgchange 0
ceph osd pool set $pool crimson_allow_pg_merge true
# Induce merges
ceph osd pool set $pool pg_num 4
ceph osd pool set $pool pg_num_min 4
wait_for_merge_completion $pool 4 || return 1
ceph osd pool set $pool pgp_num 4
wait_for_clean || return 1
}
function test_pg_merging_with_radosbench() {
local pool="merge_bench"
ceph osd pool delete $pool $pool --yes-i-really-really-mean-it || true
create_pool $pool 16
wait_for_clean || return 1
# Start radosbench writes
timeout 120 rados bench -p $pool 60 write -b 4096 --no-cleanup &
BENCH_PID=$!
sleep 5
# merge
ceph osd pool set $pool nopgchange 0
ceph osd pool set $pool crimson_allow_pg_merge true
ceph osd pool set $pool pg_num 4
ceph osd pool set $pool pg_num_min 4
# Use the loop to ensure we don't stop until all 12 PGs are merged away
wait_for_merge_completion $pool 4 || return 1
ceph osd pool set $pool pgp_num 4
wait_for_clean || return 1
# Ensure the bench finished successfully
wait $BENCH_PID
# Final check
actual=$(ceph pg ls-by-pool $pool --format=json 2>/dev/null | jq -r ".pg_stats | length")
test "$actual" -eq 4 || return 1
}
test_pg_merging || { echo "test_pg_merging failed"; exit 1; }
test_pg_merging_with_radosbench || { echo "test_pg_merging_with_radosbench failed"; exit 1; }
echo "OK"

View File

@ -1316,6 +1316,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 {}",

View File

@ -82,20 +82,39 @@ seastar::future<> PGAdvanceMap::start()
pg->handle_initialize(rctx);
pg->handle_activate_map(rctx);
}
for (epoch_t next_epoch = from + 1; next_epoch <= to; next_epoch++) {
DEBUG("{}: start: getting map {}",
*this, next_epoch);
ceph_assert(std::cmp_less_equal(from, to));
for (epoch_t current_epoch = from; current_epoch < to; ++current_epoch) {
epoch_t next_epoch = current_epoch + 1;
DEBUG("{}: start: getting map {}", *this, next_epoch);
cached_map_t next_map = co_await shard_services.get_map(next_epoch);
DEBUG("{}: advancing map to {}",
*this, next_map->get_epoch());
// Use the current OSDMap epoch to check for splits consecutively.
epoch_t current_epoch = pg->get_osdmap_epoch();
merge_result_t merge_result =
co_await check_for_merges(current_epoch, next_map, rctx);
if (merge_result.role == merge_role_t::Source) {
// Merge source: hand the PG off to the target and finish. The PG
// stops existing locally, so there is nothing left to advance here.
DEBUG("{}: merge source, handing off after e{}", *this, current_epoch);
co_await finish_merge_source(merge_result.parent, rctx);
co_await handle.complete();
exit_handle.cancel();
co_return;
}
// A successful merge target (merge already applied by check_for_merges)
// and the no-merge case both advance this epoch and keep iterating, so
// the PG advances all the way to `to`, matching the classic OSD instead
// of stalling at an intermediate epoch.
DEBUG("{}: advancing map to {}", *this, next_map->get_epoch());
// Baseline epoch for split detection: PG map epoch before advancing to next_map.
epoch_t split_epoch = pg->get_osdmap_epoch();
pg->handle_advance_map(next_map, rctx);
DEBUG("{}: checking for splits between {} and {}",
*this, current_epoch, next_map->get_epoch());
co_await check_for_splits(current_epoch, std::move(next_map));
*this, split_epoch, next_map->get_epoch());
co_await check_for_splits(split_epoch, std::move(next_map));
}
// Normal completion, which also covers a completed merge target: the PG was
// fully advanced to `to` in the loop above, so just activate the final map
// and commit.
pg->handle_activate_map(rctx);
DEBUG("{}: map activated", *this);
if (do_init) {
@ -140,7 +159,6 @@ seastar::future<> PGAdvanceMap::check_for_splits(
}
}
seastar::future<> PGAdvanceMap::split_pg(
std::set<spg_t> split_children,
cached_map_t next_map)
@ -236,5 +254,89 @@ void PGAdvanceMap::split_stats(std::set<Ref<PG>> children_pgs,
pg->finish_split_stats(*stat_iter, rctx.transaction);
}
seastar::future<PGAdvanceMap::merge_result_t> PGAdvanceMap::check_for_merges(
epoch_t old_epoch,
cached_map_t next_map,
PeeringCtx &rctx)
{
LOG_PREFIX(PGAdvanceMap::check_for_merges);
using cached_map_t = OSDMapService::cached_map_t;
cached_map_t old_map = co_await shard_services.get_map(old_epoch);
if (!old_map->have_pg_pool(pg->get_pgid().pool())) {
DEBUG("{} pool doesn't exist in epoch {}", pg->get_pgid(),
old_epoch);
co_return merge_result_t{};
}
auto old_pg_num = old_map->get_pg_num(pg->get_pgid().pool());
if (!next_map->have_pg_pool(pg->get_pgid().pool())) {
DEBUG("{} pool doesn't exist in epoch {}", pg->get_pgid(),
next_map->get_epoch());
co_return merge_result_t{};
}
auto new_pg_num = next_map->get_pg_num(pg->get_pgid().pool());
DEBUG("{} pg_num change in e{} {} -> {}", pg->get_pgid(), next_map->get_epoch(),
old_pg_num, new_pg_num);
if (new_pg_num && new_pg_num < old_pg_num) {
co_return co_await merge_pg(next_map, new_pg_num, old_pg_num, rctx);
}
co_return merge_result_t{};
}
seastar::future<PGAdvanceMap::merge_result_t> PGAdvanceMap::merge_pg(
cached_map_t next_map,
unsigned new_pg_num,
unsigned old_pg_num,
PeeringCtx &rctx)
{
LOG_PREFIX(PGAdvanceMap::merge_pg);
DEBUG("{}: start", *this);
spg_t parent;
std::set<spg_t> merge_sources;
if (pg->pgid.is_merge_source(old_pg_num,
new_pg_num,
&parent)) {
parent.is_split(new_pg_num, old_pg_num, &merge_sources);
if (!co_await shard_services.seastore_merge_shards_ok(
parent, merge_sources)) {
co_return merge_result_t{};
}
co_return merge_result_t{merge_role_t::Source, parent};
} else if (pg->pgid.is_merge_target(old_pg_num,
new_pg_num)) {
DEBUG("Target PG {} identified. Waiting for sources...", pg->get_pgid());
pg->pgid.is_split(new_pg_num, old_pg_num, &merge_sources);
if (!co_await shard_services.seastore_merge_shards_ok(
pg->get_pgid(), merge_sources)) {
co_return merge_result_t{};
}
// Block until all source PGs (potentially from other shards) arrive
// on this PG's rendezvous
auto sources = co_await pg->collect_merge_sources(merge_sources.size());
if (sources.empty()) {
co_return merge_result_t{};
}
unsigned split_bits = pg->get_pgid().get_split_bits(new_pg_num);
const auto& merge_meta =
next_map->get_pg_pool(pg->get_pgid().pool())->last_pg_merge_meta;
pg->merge_from(sources, rctx, split_bits, merge_meta);
co_return merge_result_t{merge_role_t::Target};
} else {
co_return merge_result_t{};
}
}
seastar::future<> PGAdvanceMap::finish_merge_source(
spg_t parent,
PeeringCtx &rctx)
{
LOG_PREFIX(PGAdvanceMap::finish_merge_source);
DEBUG("{}: committing rctx before handoff to target {}", *this, parent);
co_await pg->complete_rctx(std::move(rctx));
co_await pg->stop();
co_await shard_services.register_merge_source(parent, pg->get_pgid());
}
}

View File

@ -52,8 +52,23 @@ public:
PipelineHandle &get_handle() { return handle; }
using cached_map_t = OSDMapService::cached_map_t;
enum class merge_role_t {
None,
Source,
Target,
};
struct merge_result_t {
merge_role_t role = merge_role_t::None;
spg_t parent;
};
seastar::future<> check_for_splits(epoch_t old_epoch,
cached_map_t next_map);
seastar::future<merge_result_t> check_for_merges(epoch_t old_epoch,
cached_map_t next_map,
PeeringCtx &rctx);
seastar::future<> split_pg(std::set<spg_t> split_children,
cached_map_t next_map);
void split_stats(std::set<Ref<PG>> child_pgs,
@ -69,10 +84,17 @@ public:
}
private:
seastar::future<merge_result_t> merge_pg(cached_map_t next_map,
unsigned new_pg_num,
unsigned old_pg_num,
PeeringCtx &rctx);
PGPeeringPipeline &peering_pp(PG &pg);
seastar::future<Ref<PG>> handle_split_pg_creation(
spg_t child_pgid,
cached_map_t next_map);
seastar::future<> finish_merge_source(
spg_t parent,
PeeringCtx &rctx);
};
}

View File

@ -1683,20 +1683,28 @@ seastar::future<> PG::stop()
{
logger().info("PG {} {}", pgid, __func__);
stopping = true;
context_registry_on_change();
clear_primary_state();
if (is_primary()) {
clear_ready_to_merge();
}
reset_merge_rendezvous();
cancel_local_background_io_reservation();
cancel_remote_recovery_reservation();
check_readable_timer.cancel();
renew_lease_timer.cancel();
backend->on_actingset_changed(false);
return osdmap_gate.stop().then([this] {
return wait_for_active_blocker.stop();
}).then([this] {
return recovery_handler->stop();
}).then([this] {
return recovery_backend->stop();
}).then([this] {
return backend->stop();
});
co_await merge_notify_gate.close();
co_await osdmap_gate.stop();
co_await wait_for_active_blocker.stop();
client_request_orderer.clear_and_cancel(*this);
co_await recovery_handler->stop();
co_await recovery_backend->stop();
co_await backend->stop();
}
void PG::on_change(ceph::os::Transaction &t) {
@ -2031,4 +2039,89 @@ void PG::PGLogEntryHandler::partial_write(pg_info_t *info,
__func__, info->partial_writes_last_complete);
}
void PG::add_merge_source(
spg_t source,
core_id_t birth_shard,
seastar::foreign_ptr<Ref<PG>> source_pg)
{
LOG_PREFIX(PG::add_merge_source);
auto wrapped =
crimson::make_local_shared_foreign<Ref<PG>>(std::move(source_pg));
auto [_, inserted] = merge_rendezvous.sources.emplace(
source, std::make_pair(birth_shard, std::move(wrapped)));
if (inserted) {
DEBUG("target {} source {} arrived from shard {} ({} total)",
get_pgid(), source, birth_shard,
merge_rendezvous.sources.size());
merge_rendezvous.arrivals.signal(1);
} else {
DEBUG("target {} source {} already registered, ignoring",
get_pgid(), source);
}
}
seastar::future<PG::merge_source_map_t>
PG::collect_merge_sources(std::size_t n)
{
LOG_PREFIX(PG::collect_merge_sources);
DEBUG("target {} waiting for {} sources ({} already arrived)",
get_pgid(), n, merge_rendezvous.sources.size());
try {
co_await merge_rendezvous.arrivals.wait(n);
} catch (const seastar::broken_semaphore&) {
DEBUG("target {} merge rendezvous broken, aborting wait", get_pgid());
co_return merge_source_map_t{};
}
ceph_assert(merge_rendezvous.sources.size() == n);
auto sources = std::move(merge_rendezvous.sources);
merge_rendezvous.sources.clear();
co_return sources;
}
void PG::reset_merge_rendezvous()
{
// Unblock any waiter in collect_merge_sources() with broken(); then
// replace the semaphore so the next merge attempt starts at zero signals.
merge_rendezvous.arrivals.broken();
merge_rendezvous = merge_rendezvous_t{};
}
void PG::merge_from(
merge_source_map_t& sources,
PeeringCtx &rctx,
unsigned split_bits,
const pg_merge_meta_t& last_pg_merge_meta)
{
LOG_PREFIX(PG::merge_from);
DEBUG("target {}", get_pgid());
std::map<spg_t, PeeringState*> source_states;
for (auto& [pgid, entry] : sources) {
auto& src_pg = entry.second;
source_states.emplace(pgid, &src_pg->peering_state);
}
// Updates the target's PeeringState and stats (Synchronous)
peering_state.merge_from(source_states, rctx, split_bits, last_pg_merge_meta);
// We iterate through each source and move its objects into the target collection
for (auto& [pgid, entry] : sources) {
auto& src_pg = entry.second;
auto src_coll = src_pg->get_collection_ref()->get_cid();
auto dst_coll = coll_ref->get_cid();
DEBUG("merging source {}", pgid);
// Remove source-specific metadata objects that are no longer needed
// now that the collections are being collapsed.
rctx.transaction.remove(src_coll, src_pg->get_pgid().make_snapmapper_oid());
rctx.transaction.remove(src_coll, src_pg->pgmeta_oid);
// Move source collection into target collection.
rctx.transaction.merge_collection(src_coll, dst_coll, split_bits);
}
// Adjust the collection and snap_mapper to reflect the
// new, smaller PG count (reducing bitmask).
rctx.transaction.collection_set_bits(coll_ref->get_cid(), split_bits);
snap_mapper.update_bits(split_bits);
}
}

View File

@ -8,6 +8,8 @@
#include <boost/smart_ptr/intrusive_ref_counter.hpp>
#include <seastar/core/future.hh>
#include <seastar/core/shared_future.hh>
#include <seastar/core/semaphore.hh>
#include <seastar/core/sharded.hh>
#include "common/dout.h"
#include "common/ostream_temp.h"
@ -26,6 +28,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 +566,83 @@ 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 {}
// Per-PG rendezvous used to collect source PGs converging on this PG
// during a merge. Producers (source-side coroutines on other shards)
// push entries via add_merge_source(); the target-side coroutine waits
// via collect_merge_sources(n).
using merge_source_entry_t =
std::pair<core_id_t, crimson::local_shared_foreign_ptr<Ref<PG>>>;
using merge_source_map_t = std::map<spg_t, merge_source_entry_t>;
// Producer side: called on the target's shard with the foreign PG
// already wrapped. The first registration for a given source signals
// the semaphore; duplicate registrations (e.g. from replay) are no-ops.
void add_merge_source(
spg_t source,
core_id_t birth_shard,
seastar::foreign_ptr<Ref<PG>> source_pg);
// Consumer side: wait until `n` distinct sources have arrived, then
// return them and clear the rendezvous state. Returns an empty map if
// reset_merge_rendezvous() breaks the wait (e.g. PG stop or merge cancel).
seastar::future<merge_source_map_t> collect_merge_sources(std::size_t n);
// Drop in-flight handoffs and reset the semaphore. Call on PG stop or
// after Seastore cross-shard cancel so a failed try cannot leave stale
// sources for the next epoch.
void reset_merge_rendezvous();
void merge_from(
merge_source_map_t& sources,
PeeringCtx &rctx,
unsigned split_bits,
const pg_merge_meta_t& last_pg_merge_meta);
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 +1222,20 @@ private:
// continuations here.
bool stopping = false;
// Rendezvous state owned by the target PG of a pending merge.
// sources is keyed by source pgid so re-registration is naturally
// idempotent; arrivals is signaled exactly once per first insert.
struct merge_rendezvous_t {
merge_source_map_t sources;
seastar::semaphore arrivals{0};
};
merge_rendezvous_t merge_rendezvous;
// 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;

View File

@ -9,6 +9,8 @@
#include "messages/MOSDMap.h"
#include "messages/MOSDPGCreated.h"
#include "messages/MOSDPGTemp.h"
#include "messages/MOSDPGReadyToMerge.h"
#include "messages/MOSDPGStopMerge.h"
#include "osd/osd_perf_counters.h"
#include "osd/PeeringState.h"
@ -105,6 +107,125 @@ seastar::future<> PerShardState::broadcast_map_to_pgs(
});
}
seastar::future<Ref<PG>> ShardServices::extract_pg(spg_t pgid) {
auto pg = local_state.pg_map.get_pg(pgid);
ceph_assert(pg);
co_await remove_pg(pgid);
co_return pg;
}
seastar::future<> ShardServices::register_merge_source(
spg_t target,
spg_t source)
{
LOG_PREFIX(ShardServices::register_merge_source);
core_id_t birth_shard = seastar::this_shard_id();
core_id_t target_core = co_await get_pg_mapping(target);
// Remove the source pg from pg_to_shard_mapping so that
// it no longer receives any messages
auto pg_to_move = co_await extract_pg(source);
// Wrap the PG in a foreign_ptr to cross cores safely. If the target
// happens to be on the same shard, invoke_on simply runs the lambda
// inline.
auto foreign_pg = seastar::make_foreign(std::move(pg_to_move));
DEBUG("Target {} on shard {} (from shard {}); handing off source {}",
target, target_core, birth_shard, source);
co_await container().invoke_on(
target_core,
[target, source, birth_shard, foreign_pg = std::move(foreign_pg)]
(ShardServices& target_svc) mutable {
auto target_pg = target_svc.local_state.pg_map.get_pg(target);
// PGAdvanceMap on the target shard is expected to have instantiated
// the target PG by this point in the merge protocol.
ceph_assert(target_pg);
target_pg->add_merge_source(
source, birth_shard, std::move(foreign_pg));
});
}
bool ShardServices::is_seastore_objectstore()
{
return crimson::common::local_conf().get_val<std::string>(
"osd_objectstore") == "seastore";
}
seastar::future<> ShardServices::reset_target_merge_rendezvous(spg_t target)
{
const core_id_t target_core = co_await get_pg_mapping(target);
co_await container().invoke_on(
target_core,
[target](ShardServices& target_svc) {
if (auto target_pg = target_svc.local_state.pg_map.get_pg(target)) {
target_pg->reset_merge_rendezvous();
}
});
}
seastar::future<> ShardServices::send_stop_pool_merge(pg_t source_pgid)
{
co_await with_singleton(
[pool = source_pgid.pool(), source_pgid](OSDSingletonState& singleton) {
return singleton.send_stop_pool_pg_merge(pool, source_pgid);
});
}
seastar::future<> ShardServices::abort_seastore_cross_shard_merge(
spg_t target,
pg_t source_pgid,
const std::set<spg_t>* all_sources,
bool reset_target_rendezvous)
{
if (all_sources) {
for (const auto& src : *all_sources) {
co_await clear_ready_to_merge(src.pgid);
}
} else {
co_await clear_ready_to_merge(source_pgid);
}
co_await clear_ready_to_merge(target.pgid);
// Ensure the monitor backs off before it can commit a pg_num decrement
// for this source. "Stop merge" is permanent policy; "not ready" prevents
// an already-in-flight decrement decision from being committed.
co_await set_not_ready_to_merge_source(source_pgid);
co_await send_stop_pool_merge(source_pgid);
if (reset_target_rendezvous) {
co_await reset_target_merge_rendezvous(target);
}
}
seastar::future<bool> ShardServices::seastore_merge_shards_ok(
spg_t target,
const std::set<spg_t>& merge_sources)
{
LOG_PREFIX(ShardServices::seastore_merge_shards_ok);
if (!is_seastore_objectstore()) {
co_return true;
}
const core_id_t target_core = co_await get_pg_mapping(target);
std::optional<pg_t> cross_shard_source;
for (const auto& src : merge_sources) {
if (co_await get_pg_mapping(src) != target_core) {
cross_shard_source = src.pgid;
break;
}
}
if (!cross_shard_source) {
co_return true;
}
DEBUG("seastore: target {} has a cross-shard merge source {}; stop pool merge",
target, *cross_shard_source);
co_await abort_seastore_cross_shard_merge(
target, *cross_shard_source, &merge_sources, true);
co_return false;
}
Ref<PG> PerShardState::get_pg(spg_t pgid)
{
assert_core();
@ -307,6 +428,190 @@ 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();
}
seastar::future<> OSDSingletonState::send_stop_pool_pg_merge(int64_t pool, pg_t pgid)
{
LOG_PREFIX(OSDSingletonState::send_stop_pool_pg_merge);
const pg_pool_t *pi = osdmap->get_pg_pool(pool);
if (!pi || !pi->is_crimson()) {
co_return;
}
if (!pi->has_flag(pg_pool_t::FLAG_CRIMSON_ALLOW_PG_MERGE)) {
pools_merge_stopped_reported.insert(pool);
co_return;
}
// Claim the pool before co_await: otherwise concurrent callers can all pass
// the check and each send MOSDPGStopMerge while yielded on send_message().
if (!pools_merge_stopped_reported.insert(pool).second) {
co_return;
}
DEBUG("seastore: asking monitor to stop PG merge for pool {} (pg {})",
pool, pgid);
co_await monc.send_message(crimson::make_message<MOSDPGStopMerge>(
pool, pgid, MOSDPGStopMerge::REASON_CROSS_SHARD, osdmap->get_epoch()));
}
void OSDSingletonState::prune_pools_merge_stopped_reported()
{
// send_stop_pool_pg_merge() inserts a pool after one MOSDPGStopMerge so we
// do not spam the mon on every cross-shard source/target pair. Keep that
// entry while the map still has merge disabled; forget it if the pool
// disappears or crimson_allow_pg_merge is set again.
auto pool = pools_merge_stopped_reported.begin();
while (pool != pools_merge_stopped_reported.end()) {
const pg_pool_t *pi = osdmap->get_pg_pool(*pool);
if (!pi || pi->has_flag(pg_pool_t::FLAG_CRIMSON_ALLOW_PG_MERGE)) {
pool = pools_merge_stopped_reported.erase(pool);
} else {
++pool;
}
}
}
void OSDSingletonState::prune_sent_ready_to_merge()
{
LOG_PREFIX(OSDSingletonState::prune_sent_ready_to_merge);
prune_pools_merge_stopped_reported();
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);

View File

@ -370,12 +370,35 @@ 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;
std::set<int64_t> pools_merge_stopped_reported;
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();
seastar::future<> send_stop_pool_pg_merge(int64_t pool, pg_t pgid);
void clear_sent_ready_to_merge();
void prune_pools_merge_stopped_reported();
void prune_sent_ready_to_merge();
};
/**
* Represents services available to each PG
*/
class ShardServices : public OSDMapService {
class ShardServices : public OSDMapService,
public seastar::peering_sharded_service<ShardServices> {
friend class PGShardManager;
friend class OSD;
using cached_map_t = OSDMapService::cached_map_t;
@ -525,6 +548,13 @@ public:
return pg_to_shard_mapping.get_or_create_pg_mapping(pgid, core, store_index);
}
seastar::future<core_id_t> get_pg_mapping(spg_t pgid) {
return pg_to_shard_mapping.get_or_create_pg_mapping(pgid).then(
[](std::pair<core_id_t, store_index_t> mapping) {
return mapping.first;
});
}
auto remove_pg(spg_t pgid) {
local_state.pg_map.remove_pg(pgid);
return pg_to_shard_mapping.remove_pg_mapping(pgid);
@ -649,6 +679,30 @@ public:
return local_state.ec_extent_cache_lru;
}
seastar::future<Ref<PG>> extract_pg(spg_t pgid);
// Hand the source PG off to the target PG's rendezvous, hopping shards
// if needed. The consumer side (target PG) waits via
// PG::collect_merge_sources(); cleanup happens when the target drops the
// local_shared_foreign_ptr it received.
seastar::future<> register_merge_source(spg_t target, spg_t source);
static bool is_seastore_objectstore();
// True if every merge source is co-located on the target PG's shard
// (Seastore only). Uses the full sibling set so one source cannot commit
// while another would abort.
seastar::future<bool> seastore_merge_shards_ok(
spg_t target, const std::set<spg_t>& merge_sources);
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)
@ -810,6 +864,15 @@ public:
invoke_context_on_core(seastar::this_shard_id(), on_reserved));
}
private:
seastar::future<> reset_target_merge_rendezvous(spg_t target);
seastar::future<> send_stop_pool_merge(pg_t source_pgid);
seastar::future<> abort_seastore_cross_shard_merge(
spg_t target,
pg_t source_pgid,
const std::set<spg_t>* all_sources,
bool reset_target_rendezvous);
#undef FORWARD_CONST
#undef FORWARD
#undef FORWARD_TO_OSD_SINGLETON

View File

@ -0,0 +1,53 @@
// -*- mode:C++; tab-width:8; c-basic-offset:2; indent-tabs-mode:nil -*-
// vim: ts=8 sw=2 sts=2 expandtab
#pragma once
#include "osd/osd_types.h"
#include "messages/PaxosServiceMessage.h"
/// OSD -> mon: permanently stop PG merge shrink for a Crimson pool.
class MOSDPGStopMerge : public PaxosServiceMessage {
public:
static constexpr uint8_t REASON_CROSS_SHARD = 1;
int64_t pool = -1;
pg_t pgid;
uint8_t reason = REASON_CROSS_SHARD;
MOSDPGStopMerge()
: PaxosServiceMessage{MSG_OSD_PG_STOP_MERGE, 0}
{}
MOSDPGStopMerge(int64_t pool, pg_t pgid, uint8_t reason, epoch_t epoch)
: PaxosServiceMessage{MSG_OSD_PG_STOP_MERGE, epoch},
pool(pool),
pgid(pgid),
reason(reason)
{}
void encode_payload(uint64_t features) override {
using ceph::encode;
paxos_encode();
encode(pool, payload);
encode(pgid, payload);
encode(reason, payload);
}
void decode_payload() override {
using ceph::decode;
auto p = payload.cbegin();
paxos_decode(p);
decode(pool, p);
decode(pgid, p);
decode(reason, p);
}
std::string_view get_type_name() const override { return "osd_pg_stop_merge"; }
void print(std::ostream &out) const {
out << get_type_name()
<< "(pool " << pool
<< " pg " << pgid
<< " reason " << (unsigned)reason
<< " v" << version << ")";
}
private:
template<class T, typename... Args>
friend boost::intrusive_ptr<T> ceph::make_message(Args&&... args);
};

View File

@ -1184,6 +1184,7 @@ COMMAND("osd pool get "
"|compression_min_blob_size"
"|compression_mode"
"|compression_required_ratio"
"|crimson_allow_pg_merge"
"|crush_rule"
"|csum_max_block"
"|csum_min_block"
@ -1250,6 +1251,7 @@ COMMAND("osd pool set "
"|compression_min_blob_size"
"|compression_mode"
"|compression_required_ratio"
"|crimson_allow_pg_merge"
"|crush_rule"
"|csum_max_block"
"|csum_min_block"

View File

@ -4782,6 +4782,7 @@ void Monitor::dispatch_op(MonOpRequestRef op)
case MSG_REMOVE_SNAPS:
case MSG_MON_GET_PURGED_SNAPS:
case MSG_OSD_PG_READY_TO_MERGE:
case MSG_OSD_PG_STOP_MERGE:
paxos_service[PAXOS_OSDMAP]->dispatch(op);
return;

View File

@ -54,6 +54,7 @@
#include "messages/MOSDPGCreated.h"
#include "messages/MOSDPGTemp.h"
#include "messages/MOSDPGReadyToMerge.h"
#include "messages/MOSDPGStopMerge.h"
#include "messages/MMonCommand.h"
#include "messages/MRemoveSnaps.h"
#include "messages/MRoute.h"
@ -2771,6 +2772,8 @@ bool OSDMonitor::preprocess_query(MonOpRequestRef op)
return preprocess_pg_created(op);
case MSG_OSD_PG_READY_TO_MERGE:
return preprocess_pg_ready_to_merge(op);
case MSG_OSD_PG_STOP_MERGE:
return preprocess_pg_stop_merge(op);
case MSG_OSD_PGTEMP:
return preprocess_pgtemp(op);
case MSG_OSD_BEACON:
@ -2817,6 +2820,8 @@ bool OSDMonitor::prepare_update(MonOpRequestRef op)
return prepare_pgtemp(op);
case MSG_OSD_PG_READY_TO_MERGE:
return prepare_pg_ready_to_merge(op);
case MSG_OSD_PG_STOP_MERGE:
return prepare_pg_stop_merge(op);
case MSG_OSD_BEACON:
return prepare_beacon(op);
@ -4078,7 +4083,19 @@ bool OSDMonitor::prepare_pg_ready_to_merge(MonOpRequestRef op)
return false; /* nothing to propose, yet */
}
if (m->ready) {
bool allow_merge = true;
if (m->ready && p.has_flag(pg_pool_t::FLAG_CRIMSON)) {
if (!p.has_flag(pg_pool_t::FLAG_CRIMSON_ALLOW_PG_MERGE)) {
allow_merge = false;
mon.clog->warn() << "blocking crimson pg merge for " << m->pgid
<< " (pool '" << osdmap.get_pool_name(m->pgid.pool())
<< "') because pool flag 'crimson_allow_pg_merge' is not set";
dout(1) << __func__ << " blocking crimson pg merge for " << m->pgid
<< " because pool flag crimson_allow_pg_merge is not set" << dendl;
}
}
if (m->ready && allow_merge) {
p.dec_pg_num(m->pgid,
pending_inc.epoch,
m->source_version,
@ -4088,6 +4105,16 @@ bool OSDMonitor::prepare_pg_ready_to_merge(MonOpRequestRef op)
p.last_change = pending_inc.epoch;
} else {
// back off the merge attempt!
if (!m->ready && !p.has_flag(pg_pool_t::FLAG_CRIMSON)) {
mon.clog->warn() << "osd." << m->get_orig_source().num()
<< " reported pg " << m->pgid
<< " not ready to merge; backing off pg_num decrease"
<< " for pool '"
<< osdmap.get_pool_name(m->pgid.pool()) << "'";
dout(1) << __func__ << " osd." << m->get_orig_source().num()
<< " pg " << m->pgid << " not ready to merge, backing off"
<< dendl;
}
p.set_pg_num_pending(p.get_pg_num());
}
@ -4117,6 +4144,79 @@ bool OSDMonitor::prepare_pg_ready_to_merge(MonOpRequestRef op)
return true;
}
bool OSDMonitor::preprocess_pg_stop_merge(MonOpRequestRef op)
{
op->mark_osdmon_event(__func__);
auto m = op->get_req<MOSDPGStopMerge>();
dout(10) << __func__ << " " << *m << dendl;
auto session = op->get_session();
if (!session) {
dout(10) << __func__ << ": no monitor session!" << dendl;
goto ignore;
}
if (!session->is_capable("osd", MON_CAP_X)) {
derr << __func__ << " received from entity "
<< "with insufficient privileges " << session->caps << dendl;
goto ignore;
}
if (!osdmap.get_pg_pool(m->pool)) {
derr << __func__ << " pool " << m->pool << " dne" << dendl;
goto ignore;
}
return false;
ignore:
mon.no_reply(op);
return true;
}
bool OSDMonitor::prepare_pg_stop_merge(MonOpRequestRef op)
{
op->mark_osdmon_event(__func__);
auto m = op->get_req<MOSDPGStopMerge>();
dout(10) << __func__ << " " << *m << dendl;
pg_pool_t p;
if (pending_inc.new_pools.count(m->pool))
p = pending_inc.new_pools[m->pool];
else
p = *osdmap.get_pg_pool(m->pool);
if (!p.is_crimson()) {
dout(10) << __func__ << " pool " << m->pool << " is not crimson, ignoring"
<< dendl;
wait_for_finished_proposal(op, new C_ReplyMap(this, op, m->version));
return false;
}
const char *reason_str = "unknown";
if (m->reason == MOSDPGStopMerge::REASON_CROSS_SHARD) {
reason_str = "cross-shard PG merge not supported on Seastore";
}
mon.clog->warn() << "osd." << m->get_orig_source().num()
<< " stopped PG merge for pool '"
<< osdmap.get_pool_name(m->pool)
<< "' (" << reason_str << ", source pg " << m->pgid
<< "); no further pg_num decrease will be attempted";
// Cancel any in-flight shrink; keep current pg_num as-is.
p.set_pg_num_pending(p.get_pg_num());
if (p.get_pg_num_target() < p.get_pg_num()) {
p.set_pg_num_target(p.get_pg_num());
}
if (p.get_pgp_num_target() < p.get_pgp_num()) {
p.set_pgp_num_target(p.get_pgp_num());
}
p.unset_flag(pg_pool_t::FLAG_CRIMSON_ALLOW_PG_MERGE);
p.last_pg_merge_meta = pg_merge_meta_t{};
p.last_change = pending_inc.epoch;
p.last_force_op_resend_prenautilus = pending_inc.epoch;
pending_inc.new_pools[m->pool] = p;
wait_for_finished_proposal(op, new C_ReplyMap(this, op, m->version));
return true;
}
// -------------
// pg_temp changes
@ -8742,11 +8842,14 @@ int OSDMonitor::prepare_command_pool_set(const cmdmap_t& cmdmap,
return -EPERM;
}
// check for Crimson pools
// pg merging is not yet supported in Crimson
// pg merging is only supported when explicitly enabled per-pool (crimson_allow_pg_merge)
if (p.has_flag(pg_pool_t::FLAG_CRIMSON)) {
if (n < (int)p.get_pg_num()) {
ss << "crimson-osd does not support decreasing pg_num_actual (shrinking)";
return -ENOTSUP;
if (!p.has_flag(pg_pool_t::FLAG_CRIMSON_ALLOW_PG_MERGE)) {
ss << "crimson-osd does not support decreasing pg_num_actual (shrinking) "
<< "unless the pool flag crimson_allow_pg_merge is set";
return -ENOTSUP;
}
}
if (n > (int)p.get_pg_num() && !g_conf().get_val<bool>("crimson_allow_pg_split")) {
ss << "crimson_allow_pg_split is false; pg_num_actual increase denied";
@ -8805,11 +8908,14 @@ int OSDMonitor::prepare_command_pool_set(const cmdmap_t& cmdmap,
return -EPERM;
}
// check for Crimson pools
// pg merging is not yet supported in Crimson
// pg merging is only supported when explicitly enabled per-pool (crimson_allow_pg_merge)
if (p.has_flag(pg_pool_t::FLAG_CRIMSON)) {
if (n < (int)p.get_pg_num_target()) {
ss << "crimson-osd does not support decreasing pg_num";
return -ENOTSUP;
if (!p.has_flag(pg_pool_t::FLAG_CRIMSON_ALLOW_PG_MERGE)) {
ss << "crimson-osd does not support decreasing pg_num "
<< "unless the pool flag crimson_allow_pg_merge is set";
return -ENOTSUP;
}
}
if (n > (int)p.get_pg_num_target() && !g_conf().get_val<bool>("crimson_allow_pg_split")) {
ss << "crimson_allow_pg_split is false; pg_num increase denied for crimson pool";
@ -8886,11 +8992,14 @@ int OSDMonitor::prepare_command_pool_set(const cmdmap_t& cmdmap,
return -EPERM;
}
// check for Crimson pools
// pg merging is not yet supported in Crimson
// pg merging is only supported when explicitly enabled per-pool (crimson_allow_pg_merge)
if (p.has_flag(pg_pool_t::FLAG_CRIMSON)) {
if (n < (int)p.get_pgp_num()) {
ss << "crimson-osd does not support decreasing pgp_num_actual";
return -ENOTSUP;
if (!p.has_flag(pg_pool_t::FLAG_CRIMSON_ALLOW_PG_MERGE)) {
ss << "crimson-osd does not support decreasing pgp_num_actual "
<< "unless the pool flag crimson_allow_pg_merge is set";
return -ENOTSUP;
}
}
if (n > (int)p.get_pgp_num() && !g_conf().get_val<bool>("crimson_allow_pg_split")) {
ss << "crimson_allow_pg_split is false; pgp_num_actual increase denied";
@ -8921,11 +9030,14 @@ int OSDMonitor::prepare_command_pool_set(const cmdmap_t& cmdmap,
return -EPERM;
}
// check for Crimson pools
// pg merging is not yet supported in Crimson
// pg merging is only supported when explicitly enabled per-pool (crimson_allow_pg_merge)
if (p.has_flag(pg_pool_t::FLAG_CRIMSON)) {
if (n < (int)p.get_pgp_num_target()) {
ss << "crimson-osd does not support decreasing pgp_num";
return -ENOTSUP;
if (!p.has_flag(pg_pool_t::FLAG_CRIMSON_ALLOW_PG_MERGE)) {
ss << "crimson-osd does not support decreasing pgp_num "
<< "unless the pool flag crimson_allow_pg_merge is set";
return -ENOTSUP;
}
}
if (n > (int)p.get_pgp_num_target() && !g_conf().get_val<bool>("crimson_allow_pg_split")) {
ss << "crimson_allow_pg_split is false; pgp_num increase denied";
@ -8978,7 +9090,8 @@ int OSDMonitor::prepare_command_pool_set(const cmdmap_t& cmdmap,
p.crush_rule = id;
} else if (var == "nodelete" || var == "nopgchange" ||
var == "nosizechange" || var == "write_fadvise_dontneed" ||
var == "noscrub" || var == "nodeep-scrub" || var == "bulk") {
var == "noscrub" || var == "nodeep-scrub" || var == "bulk" ||
var == "crimson_allow_pg_merge") {
uint64_t flag = pg_pool_t::get_flag_by_name(var);
// make sure we only compare against 'n' if we didn't receive a string
if (val == "true" || (interr.empty() && n == 1)) {

View File

@ -470,6 +470,9 @@ private:
bool preprocess_pg_ready_to_merge(MonOpRequestRef op);
bool prepare_pg_ready_to_merge(MonOpRequestRef op);
bool preprocess_pg_stop_merge(MonOpRequestRef op);
bool prepare_pg_stop_merge(MonOpRequestRef op);
int _check_remove_pool(int64_t pool_id, const pg_pool_t &pool, std::ostream *ss);
bool _check_become_tier(
int64_t tier_pool_id, const pg_pool_t *tier_pool,

View File

@ -101,6 +101,7 @@
#include "messages/MOSDPGRecoveryDelete.h"
#include "messages/MOSDPGRecoveryDeleteReply.h"
#include "messages/MOSDPGReadyToMerge.h"
#include "messages/MOSDPGStopMerge.h"
#include "messages/MRemoveSnaps.h"
@ -647,6 +648,9 @@ Message *decode_message(CephContext *cct,
case MSG_OSD_PG_READY_TO_MERGE:
m = make_message<MOSDPGReadyToMerge>();
break;
case MSG_OSD_PG_STOP_MERGE:
m = make_message<MOSDPGStopMerge>();
break;
case MSG_OSD_EC_WRITE:
m = make_message<MOSDECSubOpWrite>();
break;

View File

@ -152,6 +152,7 @@
#define MSG_OSD_SCRUB2 121
#define MSG_OSD_PG_READY_TO_MERGE 122
#define MSG_OSD_PG_STOP_MERGE 124
#define MSG_OSD_PG_LEASE 133
#define MSG_OSD_PG_LEASE_ACK 134

View File

@ -157,6 +157,7 @@ class MOSDPGPush;
class MOSDPGPushReply;
class MOSDPGQuery;
class MOSDPGReadyToMerge;
class MOSDPGStopMerge;
class MOSDPGRecoveryDelete;
class MOSDPGRecoveryDeleteReply;
class MOSDPGRemove;

View File

@ -1330,6 +1330,9 @@ struct pg_pool_t {
FLAG_EC_OPTIMIZATIONS = 1<<19, // enable optimizations, once enabled, cannot be disabled
FLAG_CLIENT_SPLIT_READS = 1<<20, // Optimized EC is permitted to do direct reads.
FLAG_OMAP = 1<<21, // Pool is permitted to perform OMAP operations
// Allow decreasing pg_num/pgp_num (PG merge) for crimson pools.
// Note: requires that the pool is currently all bluestore.
FLAG_CRIMSON_ALLOW_PG_MERGE = 1<<22,
};
static const char *get_flag_name(uint64_t f) {
@ -1356,6 +1359,7 @@ struct pg_pool_t {
case FLAG_EC_OPTIMIZATIONS: return "ec_optimizations";
case FLAG_CLIENT_SPLIT_READS: return "split_reads";
case FLAG_OMAP: return "supports_omap";
case FLAG_CRIMSON_ALLOW_PG_MERGE: return "crimson_allow_pg_merge";
default: return "???";
}
}
@ -1412,6 +1416,8 @@ struct pg_pool_t {
return FLAG_BULK;
if (name == "crimson")
return FLAG_CRIMSON;
if (name == "crimson_allow_pg_merge")
return FLAG_CRIMSON_ALLOW_PG_MERGE;
if (name == "ec_optimizations")
return FLAG_EC_OPTIMIZATIONS;
if (name == "split_reads")