]> git-server-git.apps.pok.os.sepia.ceph.com Git - ceph.git/commitdiff
rgw: refactor RGWDeleteMultiObj formatter_flush_cond member to local var in execute()
authorCory Snyder <csnyder@iland.com>
Tue, 8 Nov 2022 15:26:00 +0000 (15:26 +0000)
committerCory Snyder <csnyder@1111systems.com>
Wed, 22 Feb 2023 11:13:21 +0000 (06:13 -0500)
Make formatter_flush_cond a local variable in RGWDeleteMultiObj::execute().

Signed-off-by: Cory Snyder <csnyder@iland.com>
(cherry picked from commit 9c11b5fb340eb2120ed311b2a9bc13d5fdc5ab43)

Conflicts:
src/rgw/rgw_op.cc

src/rgw/rgw_op.cc
src/rgw/rgw_op.h
src/rgw/rgw_rest_s3.cc
src/rgw/rgw_rest_s3.h

index c5230d67ae6585caff5507525e435b77481ce569..d95f5eb94f06ce96047aaee86825f8eb00b302dd 100644 (file)
@@ -6784,9 +6784,11 @@ void RGWDeleteMultiObj::write_ops_log_entry(rgw_log_entry& entry) const {
   entry.delete_multi_obj_meta.objects = std::move(ops_log_entries);
 }
 
-void RGWDeleteMultiObj::wait_flush(optional_yield y, std::function<bool()> predicate)
+void RGWDeleteMultiObj::wait_flush(optional_yield y,
+                                   boost::asio::deadline_timer *formatter_flush_cond,
+                                  std::function<bool()> predicate)
 {
-  if (y) {
+  if (y && formatter_flush_cond) {
     auto yc = y.get_yield_context();
     while (!predicate()) {
       boost::system::error_code error;
@@ -6796,7 +6798,8 @@ void RGWDeleteMultiObj::wait_flush(optional_yield y, std::function<bool()> predi
   }
 }
 
-void RGWDeleteMultiObj::handle_individual_object(const rgw_obj_key *o, optional_yield y)
+void RGWDeleteMultiObj::handle_individual_object(const rgw_obj_key *o, optional_yield y,
+                                                 boost::asio::deadline_timer *formatter_flush_cond)
 {
   RGWObjectCtx *obj_ctx = static_cast<RGWObjectCtx *>(s->obj_ctx);
   std::string version_id;
@@ -6808,8 +6811,8 @@ void RGWDeleteMultiObj::handle_individual_object(const rgw_obj_key *o, optional_
                                             rgw::IAM::s3DeleteObjectVersion,
                                             ARN(obj->get_obj()));
     if (identity_policy_res == Effect::Deny) {
-      send_partial_response(*o, false, "", -EACCES);
-      continue;
+      send_partial_response(*o, false, "", -EACCES, formatter_flush_cond);
+      return;
     }
 
     rgw::IAM::Effect e = Effect::Pass;
@@ -6825,8 +6828,8 @@ void RGWDeleteMultiObj::handle_individual_object(const rgw_obj_key *o, optional_
          princ_type);
     }
     if (e == Effect::Deny) {
-      send_partial_response(*o, false, "", -EACCES);
-            continue;
+      send_partial_response(*o, false, "", -EACCES, formatter_flush_cond);
+      return;
     }
 
     if (!s->session_policies.empty()) {
@@ -6836,35 +6839,35 @@ void RGWDeleteMultiObj::handle_individual_object(const rgw_obj_key *o, optional_
                                             rgw::IAM::s3DeleteObjectVersion,
                                             ARN(obj->get_obj()));
       if (session_policy_res == Effect::Deny) {
-        send_partial_response(*o, false, "", -EACCES);
-              continue;
+        send_partial_response(*o, false, "", -EACCES, formatter_flush_cond);
+        return;
       }
       if (princ_type == rgw::IAM::PolicyPrincipal::Role) {
         //Intersection of session policy and identity policy plus intersection of session policy and bucket policy
         if ((session_policy_res != Effect::Allow || identity_policy_res != Effect::Allow) &&
             (session_policy_res != Effect::Allow || e != Effect::Allow)) {
-          send_partial_response(*o, false, "", -EACCES);
-                continue;
+          send_partial_response(*o, false, "", -EACCES, formatter_flush_cond);
+          return;
         }
       } else if (princ_type == rgw::IAM::PolicyPrincipal::Session) {
         //Intersection of session policy and identity policy plus bucket policy
         if ((session_policy_res != Effect::Allow || identity_policy_res != Effect::Allow) && e != Effect::Allow) {
-          send_partial_response(*o, false, "", -EACCES);
-                continue;
+          send_partial_response(*o, false, "", -EACCES, formatter_flush_cond);
+          return;
         }
       } else if (princ_type == rgw::IAM::PolicyPrincipal::Other) {// there was no match in the bucket policy
         if (session_policy_res != Effect::Allow || identity_policy_res != Effect::Allow) {
-          send_partial_response(*o, false, "", -EACCES);
-                continue;
+          send_partial_response(*o, false, "", -EACCES, formatter_flush_cond);
+          return;
         }
       }
-      send_partial_response(*o, false, "", -EACCES);
-            continue;
+      send_partial_response(*o, false, "", -EACCES, formatter_flush_cond);
+      return;
     }
 
     if ((identity_policy_res == Effect::Pass && e == Effect::Pass && !acl_allowed)) {
-            send_partial_response(*o, false, "", -EACCES);
-            continue;
+      send_partial_response(*o, false, "", -EACCES, formatter_flush_cond);
+      return;
     }
   }
 
@@ -6882,8 +6885,8 @@ void RGWDeleteMultiObj::handle_individual_object(const rgw_obj_key *o, optional_
         check_obj_lock = false;
       } else {
         // Something went wrong.
-        send_partial_response(*o, false, "", ret);
-        continue;
+        send_partial_response(*o, false, "", ret, formatter_flush_cond);
+        return;
       }
     } else {
       obj_size = astate->size;
@@ -6894,8 +6897,8 @@ void RGWDeleteMultiObj::handle_individual_object(const rgw_obj_key *o, optional_
       ceph_assert(astate);
       int object_lock_response = verify_object_lock(this, astate->attrset, bypass_perm, bypass_governance_mode);
       if (object_lock_response != 0) {
-        send_partial_response(*o, false, "", object_lock_response);
-        continue;
+        send_partial_response(*o, false, "", object_lock_response, formatter_flush_cond);
+        return;
       }
     }
   }
@@ -6909,8 +6912,8 @@ void RGWDeleteMultiObj::handle_individual_object(const rgw_obj_key *o, optional_
     = store->get_notification(obj.get(), s->src_object.get(), s, event_type);
   op_ret = res->publish_reserve(this);
   if (op_ret < 0) {
-    send_partial_response(*o, false, "", op_ret);
-    continue;
+    send_partial_response(*o, false, "", op_ret, formatter_flush_cond);
+    return;
   }
 
   obj->set_atomic(obj_ctx);
@@ -6926,7 +6929,7 @@ void RGWDeleteMultiObj::handle_individual_object(const rgw_obj_key *o, optional_
     op_ret = 0;
   }
 
-  send_partial_response(*o, obj->get_delete_marker(), del_op->result.version_id, op_ret);
+  send_partial_response(*o, obj->get_delete_marker(), del_op->result.version_id, op_ret, formatter_flush_cond);
 
   // send request to notification manager
   int ret = res->publish_commit(this, obj_size, ceph::real_clock::now(), etag, version_id);
@@ -6942,8 +6945,9 @@ void RGWDeleteMultiObj::execute(optional_yield y)
   vector<rgw_obj_key>::iterator iter;
   RGWMultiDelXMLParser parser;
   uint32_t aio_count = 0;
-  uint32_t max_aio = s->cct->_conf->rgw_multi_obj_del_max_aio;
+  const uint32_t max_aio = s->cct->_conf->rgw_multi_obj_del_max_aio;
   char* buf;
+  std::unique_ptr<boost::asio::deadline_timer> formatter_flush_cond;
   if (y) {
     formatter_flush_cond = std::make_unique<boost::asio::deadline_timer>(y.get_io_context());  
   }
@@ -7009,19 +7013,19 @@ void RGWDeleteMultiObj::execute(optional_yield y)
         ++iter) {
     rgw_obj_key* obj_key = &*iter;
     if (y && max_aio > 1) {
-      wait_flush(y, [&aio_count, max_aio] {
+      wait_flush(y, formatter_flush_cond.get(), [&aio_count, max_aio] {
         return aio_count < max_aio;
       });
       aio_count++;
-      spawn::spawn(y.get_yield_context(), [this, &y, &aio_count, obj_key] (yield_context yield) {
-        handle_individual_object(obj_key, optional_yield { y.get_io_context(), yield }); 
+      spawn::spawn(y.get_yield_context(), [this, &y, &aio_count, obj_key, &formatter_flush_cond] (yield_context yield) {
+        handle_individual_object(obj_key, optional_yield { y.get_io_context(), yield }, formatter_flush_cond.get()); 
         aio_count--;
       }); 
     } else {
-      handle_individual_object(obj_key, y);
+      handle_individual_object(obj_key, y, formatter_flush_cond.get());
     }
   }
-  wait_flush(y, [this, n=multi_delete->objects.size()] {
+  wait_flush(y, formatter_flush_cond.get(), [this, n=multi_delete->objects.size()] {
     return n == ops_log_entries.size();
   });
 
index be89e6365904f5672222626d41767c5c01cd0cac..970f400c8d737320f9b8d94d6c0078e242e3dae7 100644 (file)
@@ -2026,7 +2026,9 @@ class RGWDeleteMultiObj : public RGWOp {
    * Handles the deletion of an individual object and uses
    * set_partial_response to record the outcome. 
    */
-  void handle_individual_object(const rgw_obj_key *o, optional_yield y);
+  void handle_individual_object(const rgw_obj_key *o,
+                               optional_yield y,
+                                boost::asio::deadline_timer *formatter_flush_cond);
   
   /**
    * When the request is being executed in a coroutine, performs
@@ -2039,20 +2041,12 @@ class RGWDeleteMultiObj : public RGWOp {
    * and saved on the req_state vs. one that is passed on the stack.
    * This is a no-op in the case where we're not executing as a coroutine.
    */
-  void wait_flush(optional_yield y, std::function<bool()> predicate);
+  void wait_flush(optional_yield y,
+                  boost::asio::deadline_timer *formatter_flush_cond,
+                  std::function<bool()> predicate);
 
 protected:
   std::vector<delete_multi_obj_entry> ops_log_entries;
-
-  /**
-   * Acts as an async condition variable when the request is being
-   * executed on a coroutine. Formatter flushing must happen on the main
-   * request coroutine vs. spawned coroutines, so spawned coroutines use
-   * the cancellation of this timer to notify the main coroutine when
-   * data is ready to flush. 
-   */
-  std::unique_ptr<boost::asio::deadline_timer> formatter_flush_cond;
-  
   bufferlist data;
   rgw::sal::Bucket* bucket;
   bool quiet;
@@ -2077,7 +2071,8 @@ public:
   virtual void send_status() = 0;
   virtual void begin_response() = 0;
   virtual void send_partial_response(const rgw_obj_key& key, bool delete_marker,
-                                     const std::string& marker_version_id, int ret) = 0;
+                                     const std::string& marker_version_id, int ret,
+                                     boost::asio::deadline_timer *formatter_flush_cond) = 0;
   virtual void end_response() = 0;
   const char* name() const override { return "multi_object_delete"; }
   RGWOpType get_type() override { return RGW_OP_DELETE_MULTI_OBJ; }
index 49146db60301202e66f965e270dba8598c82c78f..e933e234ccf4cda064af659481579fdaaa945e99 100644 (file)
@@ -4116,7 +4116,8 @@ void RGWDeleteMultiObj_ObjStore_S3::begin_response()
 void RGWDeleteMultiObj_ObjStore_S3::send_partial_response(const rgw_obj_key& key,
                                                          bool delete_marker,
                                                          const string& marker_version_id,
-                                                          int ret)
+                                                          int ret,
+                                                          boost::asio::deadline_timer *formatter_flush_cond)
 {
   if (!key.empty()) {
     delete_multi_obj_entry ops_log_entry;
index 54ef16397a89fbf9ed8045a236f38876ae626ebc..fe6334d8c33def30a3276eff994c3f7258b34527 100644 (file)
@@ -517,7 +517,8 @@ public:
   void send_status() override;
   void begin_response() override;
   void send_partial_response(const rgw_obj_key& key, bool delete_marker,
-                             const std::string& marker_version_id, int ret) override;
+                             const std::string& marker_version_id, int ret,
+                             boost::asio::deadline_timer *formatter_flush_cond) override;
   void end_response() override;
 };