}
void OSDService::_maybe_split_pgid(OSDMapRef old_map,
- OSDMapRef new_map,
- spg_t pgid)
+ OSDMapRef new_map,
+ spg_t pgid,
+ set<spg_t> *new_children)
{
if (!old_map->have_pg_pool(pgid.pool())) {
return;
set<spg_t> children;
if (pgid.is_split(old_pgnum, new_pgnum, &children)) {
_start_split(pgid, children);
+ for (auto pgid : children) {
+ new_children->insert(pgid);
+ }
}
} else {
assert(pgid.ps() < static_cast<unsigned>(new_pgnum));
void OSDService::init_splits_between(spg_t pgid,
OSDMapRef frommap,
- OSDMapRef tomap)
+ OSDMapRef tomap,
+ set<spg_t> *new_children)
{
// First, check whether we can avoid this potentially expensive check
if (!frommap->have_pg_pool(pgid.pool())) {
&split_pgs)) {
start_split(*i, split_pgs);
even_newer_pgs.insert(split_pgs.begin(), split_pgs.end());
+ new_children->insert(split_pgs.begin(), split_pgs.end());
}
}
new_pgs.insert(even_newer_pgs.begin(), even_newer_pgs.end());
}
void OSDService::expand_pg_num(OSDMapRef old_map,
- OSDMapRef new_map)
+ OSDMapRef new_map,
+ set<spg_t> *new_children)
{
Mutex::Locker l(in_progress_split_lock);
for (auto pgid : in_progress_splits) {
- _maybe_split_pgid(old_map, new_map, pgid);
+ _maybe_split_pgid(old_map, new_map, pgid, new_children);
}
for (auto i : pending_splits) {
- _maybe_split_pgid(old_map, new_map, i.first);
+ _maybe_split_pgid(old_map, new_map, i.first, new_children);
}
}
RWLock::WLocker l(pg_map_lock);
for (auto& i : pg_map) {
pgs.insert(i.second);
+ i.second->put("PGMap");
}
pg_map.clear();
}
}
pg->ch.reset();
pg->unlock();
- pg->put("PGMap");
}
}
#ifdef PG_DEBUG_REFS
// make sure we register any splits that happened between when the pg
// was created and our latest map.
- service.init_splits_between(pgid, createmap, servicemap);
+ set<spg_t> new_children;
+ service.init_splits_between(pgid, createmap, servicemap, &new_children);
+ op_shardedwq.prime_splits(new_children);
}
return pg;
}
OSDMapRef createmap,
spg_t pgid)
{
- dout(10) << "_open_lock_pg " << pgid << dendl;
+ dout(10) << __func__ << " " << pgid << dendl;
pg_pool_t pi;
string name;
if (createmap->have_pg_pool(pgid.pool())) {
continue;
}
- service.init_splits_between(pg->pg_id, pg->get_osdmap(), osdmap);
+ set<spg_t> new_children;
+ service.init_splits_between(pg->pg_id, pg->get_osdmap(), osdmap, &new_children);
+ op_shardedwq.prime_splits(new_children);
pg->reg_next_scrub();
{
spg_t pgid = info->pgid;
- int up_primary, acting_primary;
- vector<int> up, acting;
- osdmap->pg_to_up_acting_osds(
- pgid.pgid, &up, &up_primary, &acting, &acting_primary);
-
if (maybe_wait_for_max_pg(osdmap, pgid, info->by_mon)) {
dout(10) << __func__ << " hit max pg, dropping" << dendl;
return nullptr;
PG::RecoveryCtx rctx = create_context();
- const pg_pool_t* pp = osdmap->get_pg_pool(pgid.pool());
+ OSDMapRef createmap = get_map(info->epoch);
+ int up_primary, acting_primary;
+ vector<int> up, acting;
+ createmap->pg_to_up_acting_osds(
+ pgid.pgid, &up, &up_primary, &acting, &acting_primary);
+
+ const pg_pool_t* pp = createmap->get_pg_pool(pgid.pool());
if (pp->has_flag(pg_pool_t::FLAG_EC_OVERWRITES) &&
store->get_type() != "bluestore") {
clog->warn() << "pg " << pgid
PG::_create(*rctx.transaction, pgid, pgid.get_split_bits(pp->get_pg_num()));
PG::_init(*rctx.transaction, pgid, pp);
- int role = osdmap->calc_pg_role(whoami, acting, acting.size());
+ int role = createmap->calc_pg_role(whoami, acting, acting.size());
if (!pp->is_replicated() && role != pgid.shard) {
role = -1;
}
- PGRef pg = _open_pg(get_map(info->epoch), osdmap, pgid);
+ PGRef pg = _open_pg(createmap, osdmap, pgid);
pg->lock(true);
pg->handle_initialize(&rctx);
pg->queue_null(e, e);
dispatch_context_transaction(rctx, pg);
- wake_pg_waiters(pg);
+ op_shardedwq.wake_pg_split_waiters(pg->get_pgid());
pg->unlock();
}
int num_pg_primary = 0, num_pg_replica = 0, num_pg_stray = 0;
// scan pg's
+ set<spg_t> new_children;
vector<spg_t> pgids;
{
RWLock::RLocker l(pg_map_lock);
if (pg->is_deleted()) {
continue;
}
- service.init_splits_between(it->first, service.get_osdmap(), osdmap);
+ service.init_splits_between(it->first, service.get_osdmap(), osdmap, &new_children);
// FIXME: this is lockless and racy, but we don't want to take pg lock
// here.
}
}
- service.expand_pg_num(service.get_osdmap(), osdmap);
+ service.expand_pg_num(service.get_osdmap(), osdmap, &new_children);
+ op_shardedwq.prime_splits(new_children);
service.pre_publish_map(osdmap);
service.await_reserved_maps();
service.maybe_inject_dispatch_delay();
- // remove any PGs which we no longer host from the session waiting_for_pg lists
- dout(20) << __func__ << " checking waiting_for_pg" << dendl;
+ // remove any PGs which we no longer host from the pg_slot wait lists
+ dout(20) << __func__ << " checking pg_slot waiters" << dendl;
op_shardedwq.prune_or_wake_pg_waiters(osdmap, whoami);
service.maybe_inject_dispatch_delay();
OSDMapRef nextmap,
PG::RecoveryCtx *rctx)
{
- // make sure to-be-split children are blocked in wq
- op_shardedwq.prime_splits(childpgids);
-
unsigned pg_num = nextmap->get_pg_num(
parent->pg_id.pool());
parent->update_snap_mapper_bits(
{
dout(20) << __func__ << " " << pgid
<< " to_process " << slot.to_process
- << " waiting_for_pg=" << (int)slot.waiting_for_pg << dendl;
+ << " waiting " << slot.waiting
+ << " waiting_nopg " << slot.waiting_peering << dendl;
for (auto& q : slot.to_process) {
*pushes_to_free += q.get_reserved_pushes();
}
+ for (auto& q : slot.waiting) {
+ *pushes_to_free += q.get_reserved_pushes();
+ }
+ for (auto& q : slot.waiting_peering) {
+ *pushes_to_free += q.get_reserved_pushes();
+ }
for (auto i = slot.to_process.rbegin();
i != slot.to_process.rend();
++i) {
sdata->_enqueue_front(std::move(*i), osd->op_prio_cutoff);
}
slot.to_process.clear();
- slot.waiting_for_pg = false;
+ for (auto i = slot.waiting.rbegin();
+ i != slot.waiting.rend();
+ ++i) {
+ sdata->_enqueue_front(std::move(*i), osd->op_prio_cutoff);
+ }
+ slot.waiting.clear();
+ for (auto i = slot.waiting_peering.rbegin();
+ i != slot.waiting_peering.rend();
+ ++i) {
+ sdata->_enqueue_front(std::move(*i), osd->op_prio_cutoff);
+ }
+ slot.waiting_peering.clear();
+ slot.pending_peering_epoch = 0;
+ slot.waiting_for_split = false;
++slot.requeue_seq;
}
-void OSD::ShardedOpWQ::wake_pg_waiters(spg_t pgid)
+void OSD::ShardedOpWQ::wake_pg_split_waiters(spg_t pgid)
{
uint32_t shard_index = pgid.hash_to_shard(shard_list.size());
auto sdata = shard_list[shard_index];
ShardData* sdata = shard_list[shard_index];
Mutex::Locker l(sdata->sdata_op_ordering_lock);
ShardData::pg_slot& slot = sdata->pg_slots[pgid];
- slot.waiting_for_pg = true;
+ slot.waiting_for_split = true;
}
}
auto p = sdata->pg_slots.begin();
while (p != sdata->pg_slots.end()) {
ShardData::pg_slot& slot = p->second;
- if (slot.pending_nopg_epoch &&
- slot.pending_nopg_epoch <= osdmap->get_epoch()) {
+ if (slot.waiting_for_split) {
dout(20) << __func__ << " " << p->first
- << " pending_nopg_epoch " << slot.pending_nopg_epoch
- << " < " << osdmap->get_epoch() << ", requeueing" << dendl;
- assert(slot.waiting_for_pg);
- assert(!slot.to_process.empty());
- for (auto& q : slot.to_process) {
- pushes_to_free += q.get_reserved_pushes();
- }
- for (auto i = slot.to_process.rbegin();
- i != slot.to_process.rend();
- ++i) {
- sdata->_enqueue_front(std::move(*i), osd->op_prio_cutoff);
+ << " waiting for split" << dendl;
+ ++p;
+ continue;
+ }
+ if (!slot.waiting_peering.empty()) {
+ assert(slot.pending_peering_epoch);
+ if (slot.pending_peering_epoch <= osdmap->get_epoch()) {
+ dout(20) << __func__ << " " << p->first
+ << " pending_peering_epoch " << slot.pending_peering_epoch
+ << " < " << osdmap->get_epoch() << ", requeueing" << dendl;
+ assert(!slot.waiting_peering.empty());
+ _wake_pg_slot(p->first, sdata, slot, &pushes_to_free);
+ queued = true;
}
- slot.to_process.clear();
- slot.waiting_for_pg = false;
- slot.pending_nopg_epoch = 0;
- ++slot.requeue_seq;
- queued = true;
++p;
continue;
}
- if (!slot.to_process.empty() && slot.num_running == 0) {
+ if (!slot.waiting.empty()) {
if (osdmap->is_up_acting_osd_shard(p->first, whoami)) {
dout(20) << __func__ << " " << p->first << " maps to us, keeping"
<< dendl;
++p;
continue;
}
- while (!slot.to_process.empty() &&
- slot.to_process.front().get_map_epoch() <= osdmap->get_epoch()) {
- auto& qi = slot.to_process.front();
+ while (!slot.waiting.empty() &&
+ slot.waiting.front().get_map_epoch() <= osdmap->get_epoch()) {
+ auto& qi = slot.waiting.front();
dout(20) << __func__ << " " << p->first
- << " item " << qi
+ << " waiting item " << qi
<< " epoch " << qi.get_map_epoch()
<< " <= " << osdmap->get_epoch()
<< ", stale, dropping" << dendl;
pushes_to_free += qi.get_reserved_pushes();
- slot.to_process.pop_front();
+ slot.waiting.pop_front();
+ }
+ if (slot.waiting.empty() &&
+ slot.num_running == 0 &&
+ !slot.pg) {
+ dout(20) << __func__ << " " << p->first << " empty, pruning" << dendl;
+ p = sdata->pg_slots.erase(p);
+ continue;
}
}
- if (slot.to_process.empty() &&
- slot.num_running == 0 &&
- !slot.pg) {
- dout(20) << __func__ << " " << p->first << " empty, pruning" << dendl;
- p = sdata->pg_slots.erase(p);
- } else {
- ++p;
- }
+ ++p;
}
if (queued) {
sdata->sdata_lock.Lock();
}
}
+void OSD::ShardedOpWQ::_add_slot_waiter(
+ spg_t pgid,
+ OSD::ShardedOpWQ::ShardData::pg_slot& slot,
+ OpQueueItem&& qi)
+{
+ if (qi.is_peering()) {
+ if (!slot.pending_peering_epoch ||
+ slot.pending_peering_epoch > qi.get_map_epoch()) {
+ slot.pending_peering_epoch = qi.get_map_epoch();
+ }
+ dout(20) << __func__ << " " << pgid
+ << " no pg, peering, item epoch is "
+ << qi.get_map_epoch()
+ << ", pending_peering_epoch now "
+ << slot.pending_peering_epoch
+ << ", will wait on " << qi << dendl;
+ slot.waiting_peering.push_back(std::move(qi));
+ } else {
+ dout(20) << __func__ << " " << pgid
+ << " no pg, item epoch is "
+ << qi.get_map_epoch()
+ << ", will wait on " << qi << dendl;
+ slot.waiting.push_back(std::move(qi));
+ }
+}
+
#undef dout_prefix
#define dout_prefix *_dout << "osd." << osd->whoami << " op_wq(" << shard_index << ") "
auto& slot = sdata->pg_slots[token];
dout(30) << __func__ << " " << token
<< " to_process " << slot.to_process
- << " waiting_for_pg=" << (int)slot.waiting_for_pg << dendl;
- bool can_wait = item.requires_pg() && !item.creates_pg();
- slot.to_process.push_back(std::move(item));
- // note the requeue seq now...
- requeue_seq = slot.requeue_seq;
- if (slot.waiting_for_pg && can_wait) {
- dout(20) << __func__ << slot.to_process.back()
- << " queued, already waiting_for_pg" << dendl;
+ << " waiting " << slot.waiting
+ << " waiting_peering " << slot.waiting_peering
+ << dendl;
+ if (slot.waiting_for_split
+ || (item.is_peering() && !slot.waiting_peering.empty())
+ || (!item.is_peering() && !slot.waiting.empty())) {
+ dout(20) << __func__ << " " << token << " already waiting, adding " << item
+ << dendl;
+ _add_slot_waiter(token, slot, std::move(item));
sdata->sdata_op_ordering_lock.Unlock();
return;
}
+ // note the requeue seq now...
+ requeue_seq = slot.requeue_seq;
pg = slot.pg;
+ slot.to_process.push_back(std::move(item));
dout(20) << __func__ << " " << slot.to_process.back()
<< " queued" << dendl;
++slot.num_running;
}
dout(30) << __func__ << " " << token
<< " to_process " << slot.to_process
- << " waiting_for_pg=" << (int)slot.waiting_for_pg << dendl;
+ << " waiting " << slot.waiting
+ << " waiting_peering " << slot.waiting_peering << dendl;
// make sure we're not already waiting for this pg
- if (slot.waiting_for_pg) {
+ /*if (!slot.waiting.empty()) {
dout(20) << __func__ << " " << token
- << " slot is waiting_for_pg" << dendl;
+ << " slot is waiting" << dendl;
if (pg) {
pg->unlock();
}
sdata->sdata_op_ordering_lock.Unlock();
return;
- }
+ }*/
ThreadPool::TPHandle tp_handle(osd->cct, hb, timeout_interval,
suicide_interval);
// should this pg shard exist on this osd in this (or a later) epoch?
OSDMapRef osdmap = sdata->waiting_for_pg_osdmap;
const PGCreateInfo *create_info = qi.creates_pg();
- if (qi.get_map_epoch() > osdmap->get_epoch()) {
- if (!!create_info || !qi.requires_pg()) {
- if (!slot.pending_nopg_epoch ||
- slot.pending_nopg_epoch > qi.get_map_epoch()) {
- slot.pending_nopg_epoch = qi.get_map_epoch();
- }
- dout(20) << __func__ << " " << token
- << " no pg, item epoch is "
- << qi.get_map_epoch() << " > " << osdmap->get_epoch()
- << ", will wait on " << qi
- << ", pending_nopg_epoch now "
- << slot.pending_nopg_epoch << dendl;
- } else {
- dout(20) << __func__ << " " << token
- << " no pg, item epoch is "
- << qi.get_map_epoch() << " > " << osdmap->get_epoch()
- << ", will wait on " << qi << dendl;
- }
- slot.to_process.push_front(std::move(qi));
- slot.waiting_for_pg = true;
- } else if (!qi.requires_pg()) {
- // for pg-less events, we run them under the ordering lock, since
- // we don't have the pg lock to keep them ordered.
- qi.run(osd, pg, tp_handle);
- sdata->sdata_op_ordering_lock.Unlock();
- return;
- } else if (osdmap->is_up_acting_osd_shard(token, osd->whoami)) {
- if (osd->service.splitting(token)) {
- dout(20) << __func__ << " " << token
- << " splitting, waiting on " << qi << dendl;
- slot.to_process.push_front(std::move(qi));
- slot.waiting_for_pg = true;
- } else if (create_info) {
- if (create_info->by_mon &&
- osdmap->get_pg_acting_primary(token.pgid) != osd->whoami) {
- dout(20) << __func__ << " " << token
- << " no pg, no longer primary, ignoring mon create on "
- << qi << dendl;
+ if (slot.waiting_for_split) {
+ dout(20) << __func__ << " " << token
+ << " splitting" << dendl;
+ _add_slot_waiter(token, slot, std::move(qi));
+ } else if (qi.get_map_epoch() > osdmap->get_epoch()) {
+ dout(20) << __func__ << " " << token
+ << " map " << qi.get_map_epoch() << " > "
+ << osdmap->get_epoch() << dendl;
+ _add_slot_waiter(token, slot, std::move(qi));
+ } else if (qi.is_peering()) {
+ if (!qi.peering_requires_pg()) {
+ // for pg-less events, we run them under the ordering lock, since
+ // we don't have the pg lock to keep them ordered.
+ qi.run(osd, pg, tp_handle);
+ } else if (osdmap->is_up_acting_osd_shard(token, osd->whoami)) {
+ if (create_info) {
+ if (create_info->by_mon &&
+ osdmap->get_pg_acting_primary(token.pgid) != osd->whoami) {
+ dout(20) << __func__ << " " << token
+ << " no pg, no longer primary, ignoring mon create on "
+ << qi << dendl;
+ } else {
+ dout(20) << __func__ << " " << token
+ << " no pg, should create on " << qi << dendl;
+ pg = osd->handle_pg_create_info(osdmap, create_info);
+ if (pg) {
+ // we created the pg! drop out and continue "normally"!
+ _wake_pg_slot(token, sdata, slot, &pushes_to_free);
+ break;
+ }
+ dout(20) << __func__ << " ignored create on " << qi << dendl;
+ }
} else {
dout(20) << __func__ << " " << token
- << " no pg, should create on " << qi << dendl;
- pg = osd->handle_pg_create_info(osdmap, create_info);
- if (pg) {
- // we created the pg! drop out and continue "normally"!
- _wake_pg_slot(token, sdata, slot, &pushes_to_free);
- break;
- }
- dout(20) << __func__ << " ignored create on " << qi << dendl;
+ << " no pg, peering, !create, discarding " << qi << dendl;
}
} else {
dout(20) << __func__ << " " << token
- << " no pg, should exist, will wait on " << qi << dendl;
- slot.to_process.push_front(std::move(qi));
- slot.waiting_for_pg = true;
+ << " no pg, peering, does't map here, discarding " << qi
+ << dendl;
}
+ } else if (osdmap->is_up_acting_osd_shard(token, osd->whoami)) {
+ dout(20) << __func__ << " " << token
+ << " no pg, should exist, will wait on " << qi << dendl;
+ _add_slot_waiter(token, slot, std::move(qi));
} else {
dout(20) << __func__ << " " << token
<< " no pg, shouldn't exist,"