]> git-server-git.apps.pok.os.sepia.ceph.com Git - ceph.git/commitdiff
osd: restructure pg waiting
authorSage Weil <sage@redhat.com>
Fri, 19 Jan 2018 19:23:01 +0000 (13:23 -0600)
committerSage Weil <sage@redhat.com>
Wed, 4 Apr 2018 13:26:50 +0000 (08:26 -0500)
Rethink the way we wait for PGs.  We need to order peering events relative to
each other; keep them in a separate queue in the pg_slot.

Signed-off-by: Sage Weil <sage@redhat.com>
src/osd/OSD.cc
src/osd/OSD.h
src/osd/OpQueueItem.h

index 1a0179c6194dc77c55b3f4e0c20fbdb50c1c6916..b7f0ddb61cea348e50911d53287f60debf05858d 100644 (file)
@@ -409,8 +409,9 @@ void OSDService::_cancel_pending_splits_for_parent(spg_t parent)
 }
 
 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;
@@ -421,6 +422,9 @@ void OSDService::_maybe_split_pgid(OSDMapRef old_map,
     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));
@@ -429,7 +433,8 @@ void OSDService::_maybe_split_pgid(OSDMapRef old_map,
 
 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())) {
@@ -467,6 +472,7 @@ void OSDService::init_splits_between(spg_t pgid,
                        &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());
@@ -478,14 +484,15 @@ void OSDService::init_splits_between(spg_t pgid,
 }
 
 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);
   }
 }
 
@@ -3487,6 +3494,7 @@ int OSD::shutdown()
       RWLock::WLocker l(pg_map_lock);
       for (auto& i : pg_map) {
        pgs.insert(i.second);
+       i.second->put("PGMap");
       }
       pg_map.clear();
     }
@@ -3511,7 +3519,6 @@ int OSD::shutdown()
       }
       pg->ch.reset();
       pg->unlock();
-      pg->put("PGMap");
     }
   }
 #ifdef PG_DEBUG_REFS
@@ -3813,7 +3820,9 @@ PGRef OSD::_open_pg(
 
     // 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;
 }
@@ -3832,7 +3841,7 @@ PG* OSD::_make_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())) {
@@ -4032,7 +4041,9 @@ void OSD::load_pgs()
       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();
 
@@ -4048,11 +4059,6 @@ PGRef OSD::handle_pg_create_info(OSDMapRef osdmap, const PGCreateInfo *info)
 {
   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;
@@ -4060,7 +4066,13 @@ PGRef OSD::handle_pg_create_info(OSDMapRef osdmap, const PGCreateInfo *info)
 
   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
@@ -4071,12 +4083,12 @@ PGRef OSD::handle_pg_create_info(OSDMapRef osdmap, const PGCreateInfo *info)
   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);
 
@@ -7882,7 +7894,7 @@ void OSD::_finish_splits(set<PGRef>& pgs)
     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();
   }
 
@@ -7959,6 +7971,7 @@ void OSD::consume_map()
   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);
@@ -7970,7 +7983,7 @@ void OSD::consume_map()
       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.
@@ -7993,7 +8006,8 @@ void OSD::consume_map()
     }
   }
 
-  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();
@@ -8005,8 +8019,8 @@ void OSD::consume_map()
 
   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();
@@ -8184,9 +8198,6 @@ void OSD::split_pgs(
   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(
@@ -9463,21 +9474,41 @@ void OSD::ShardedOpWQ::_wake_pg_slot(
 {
   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];
@@ -9509,7 +9540,7 @@ void OSD::ShardedOpWQ::prime_splits(const set<spg_t>& pgs)
     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;
   }
 }
 
@@ -9523,56 +9554,52 @@ void OSD::ShardedOpWQ::prune_or_wake_pg_waiters(OSDMapRef osdmap, int whoami)
     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();
@@ -9610,6 +9637,32 @@ void OSD::ShardedOpWQ::clear_pg_slots()
   }
 }
 
+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 << ") "
 
@@ -9654,18 +9707,22 @@ void OSD::ShardedOpWQ::_process(uint32_t thread_index, heartbeat_handle_d *hb)
     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;
@@ -9728,18 +9785,19 @@ void OSD::ShardedOpWQ::_process(uint32_t thread_index, heartbeat_handle_d *hb)
   }
   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);
@@ -9754,61 +9812,51 @@ void OSD::ShardedOpWQ::_process(uint32_t thread_index, heartbeat_handle_d *hb)
     // 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,"
index 0ca6e16df7fc9344e0f6636831fe017bd743428b..3ba6c59040f541843f32c12c5d2df551daf8ebba 100644 (file)
@@ -959,11 +959,14 @@ public:
   void _cancel_pending_splits_for_parent(spg_t parent);
   bool splitting(spg_t pgid);
   void expand_pg_num(OSDMapRef old_map,
-                    OSDMapRef new_map);
+                    OSDMapRef new_map,
+                    set<spg_t> *new_children);
   void _maybe_split_pgid(OSDMapRef old_map,
                         OSDMapRef new_map,
-                        spg_t pgid);
-  void init_splits_between(spg_t pgid, OSDMapRef frommap, OSDMapRef tomap);
+                        spg_t pgid,
+                        set<spg_t> *new_children);
+  void init_splits_between(spg_t pgid, OSDMapRef frommap, OSDMapRef tomap,
+                          set<spg_t> *new_chilren);
 
   // -- stats --
   Mutex stat_lock;
@@ -1580,16 +1583,18 @@ private:
        deque<OpQueueItem> to_process; ///< order items for this slot
        int num_running = 0;          ///< _process threads doing pg lookup/lock
 
-       /// true if pg does/did not exist. if so all new items go directly to
-       /// to_process.  cleared by prune_pg_waiters.
-       bool waiting_for_pg = false;
+       deque<OpQueueItem> waiting;         ///< waiting for pg (or map + pg)
+       deque<OpQueueItem> waiting_peering; ///< waiting for map (peering evt)
 
-       /// one or more queued items doesn't need a pg, only a map >= this
-       epoch_t pending_nopg_epoch = 0;
+       /// min required map across waiting_peering items
+       epoch_t pending_peering_epoch = 0;
 
        /// incremented by wake_pg_waiters; indicates racing _process threads
        /// should bail out (their op has been requeued)
        uint64_t requeue_seq = 0;
+
+       /// waiting for split child to materialize
+       bool waiting_for_split = false;
       };
 
       /// map of slots for each spg_t.  maintains ordering of items dequeued
@@ -1671,8 +1676,13 @@ private:
       }
     }
 
-    /// wake any pg waiters after a PG is created/instantiated
-    void wake_pg_waiters(spg_t pgid);
+    void _add_slot_waiter(
+      spg_t token,
+      ShardData::pg_slot& slot,
+      OpQueueItem&& qi);
+
+    /// wake any pg waiters after a PG is split
+    void wake_pg_split_waiters(spg_t pgid);
 
     void _wake_pg_slot(spg_t pgid, ShardData *sdata, ShardData::pg_slot& slot,
                       unsigned *pushes_to_free);
@@ -1920,11 +1930,6 @@ protected:
     int lastactingprimary
     ); ///< @return false if there was a map gap between from and now
 
-  // this must be called with pg->lock held on any pg addition to pg_map
-  void wake_pg_waiters(PGRef pg) {
-    assert(pg->is_locked());
-    op_shardedwq.wake_pg_waiters(pg->get_pgid());
-  }
   epoch_t last_pg_create_epoch;
 
   void handle_pg_create(OpRequestRef op);
index b55114ae32ed7a6bcd599e91af9619891af9dba6..35cfdc8cdcd52f5b12fc85bf5eb864305b3e6a7e 100644 (file)
  */
 
 
+/*
+
+  Ordering notes:
+
+  - everybody waits for split.
+
+  - client ops must remained ordered by client, regardless of map epoch
+  - client ops wait for pg to exist (or are discarded if we confirm the pg
+  no longer should).
+  - client ops must wait for the min epoch.
+    -> this happens under the PG itself, not as part of the queue.
+       currently in PrimaryLogPG::do_request()
+    -> the pg waiting queue is ordered by client, so other clients do not have to wait
+
+  - peering messages must wait for the required_map
+    - currently in do_peering_event(), PG::peering_waiters
+  - peering messages must remain ordered (globally or by peer?)
+  - some peering messages create the pg
+  - query does not need a pg.
+    - q: do any peering messages need to wait for the pg to exist?
+        pretty sure no!
+
+    ---
+
+    bool waiting_for_split -- everyone waits.
+
+    waiting -- client/mon ops
+    waiting_peering -- peering ops
+
+  */
+
 #pragma once
 
 #include <ostream>
@@ -65,8 +96,11 @@ public:
       return 0;
     }
 
-    virtual bool requires_pg() const {
-      return true;
+    virtual bool is_peering() const {
+      return false;
+    }
+    virtual bool peering_requires_pg() const {
+      ceph_abort();
     }
     virtual const PGCreateInfo *creates_pg() const {
       return nullptr;
@@ -148,13 +182,19 @@ public:
   epoch_t get_map_epoch() const { return map_epoch; }
   dmc::ReqParams get_qos_params() const { return qos_params; }
   void set_qos_params(dmc::ReqParams qparams) { qos_params =  qparams; }
-  bool requires_pg() const {
-    return qitem->requires_pg();
+
+  bool is_peering() const {
+    return qitem->is_peering();
   }
+
   const PGCreateInfo *creates_pg() const {
     return qitem->creates_pg();
   }
 
+  bool peering_requires_pg() const {
+    return qitem->peering_requires_pg();
+  }
+
   friend ostream& operator<<(ostream& out, const OpQueueItem& item) {
     return out << "OpQueueItem("
               << item.get_ordering_token() << " " << *item.qitem
@@ -225,7 +265,10 @@ public:
     return rhs << "PGPeeringEvent(" << evt->get_desc() << ")";
   }
   void run(OSD *osd, PGRef& pg, ThreadPool::TPHandle &handle) override final;
-  bool requires_pg() const override {
+  bool is_peering() const override {
+    return true;
+  }
+  bool peering_requires_pg() const override {
     return evt->requires_pg;
   }
   const PGCreateInfo *creates_pg() const override {