// vim: ts=8 sw=2 sts=2 expandtab ft=cpp
#include <string.h>
+#include <algorithm>
+#include <chrono>
#include <iostream>
#include <map>
+#include <random>
+#include <thread>
+
+#include <boost/asio/steady_timer.hpp>
#include "common/XMLFormatter.h"
#include <common/errno.h>
using namespace std;
+int retry_on_busy(optional_yield y, const DoutPrefixProvider *dpp,
+ CephContext *cct, const char *op_name,
+ std::function<int()> op)
+{
+ const int max_attempts = cct->_conf.get_val<int64_t>("rgw_cloud_tier_retry_limit");
+ const int64_t initial_ms = cct->_conf.get_val<int64_t>("rgw_cloud_tier_retry_delay_ms");
+ const int64_t max_ms = cct->_conf.get_val<int64_t>("rgw_cloud_tier_retry_max_ms");
+ thread_local std::mt19937 rng{std::random_device{}()};
+
+ int ret = 0;
+ for (int i = 0; i < max_attempts; i++) {
+ ret = op();
+ if (ret != -EBUSY || i == max_attempts - 1) return ret;
+
+ int64_t base = std::min(initial_ms << std::min(i, 30), max_ms);
+ int delay_ms = base - std::uniform_int_distribution<int>(0, base / 10)(rng);
+
+ ldpp_dout(dpp, 1) << op_name << ": -EBUSY, attempt " << (i + 1) << "/"
+ << max_attempts << "; retrying after " << delay_ms
+ << "ms" << dendl;
+
+ if (y) {
+ auto& yc = y.get_yield_context();
+ boost::asio::steady_timer t(yc.get_executor());
+ t.expires_after(std::chrono::milliseconds(delay_ms));
+ boost::system::error_code ec;
+ t.async_wait(yc[ec]);
+ if (ec) return ret;
+ } else {
+ std::this_thread::sleep_for(std::chrono::milliseconds(delay_ms));
+ }
+ }
+ return ret;
+}
+
struct rgw_lc_multipart_part_info {
int part_num{0};
uint64_t ofs{0};
dest_bucket.name = tier_ctx.target_bucket_name;
target_obj_name = make_target_obj_name(tier_ctx);
- if (!in_progress) { // first time. Send RESTORE req.
-
- rgw_obj dest_obj(dest_bucket, rgw_obj_key(target_obj_name));
- ret = cloud_tier_restore(tier_ctx.dpp, tier_ctx.conn, dest_obj, days, glacier_params, tier_ctx.y);
-
- ldpp_dout(tier_ctx.dpp, 20) << __func__ << "Restoring object=" << target_obj_name << "returned ret = " << ret << dendl;
-
- if (ret < 0 ) {
- ldpp_dout(tier_ctx.dpp, -1) << __func__ << "ERROR: failed to restore object=" << dest_obj << "; ret = " << ret << dendl;
- return ret;
- }
- in_progress = true;
- }
-
- // now send HEAD request and verify if restore is complete on glacier/tape endpoint
- static constexpr int MAX_RETRIES = 2;
- uint32_t retries = 0;
- do {
- ret = rgw_cloud_tier_get_object(tier_ctx, true, headers, nullptr, etag,
- accounted_size, attrs, nullptr);
+ rgw_obj dest_obj(dest_bucket, rgw_obj_key(target_obj_name));
- if (ret < 0) {
- ldpp_dout(tier_ctx.dpp, 0) << __func__ << "ERROR: failed to fetch HEAD from cloud for obj=" << tier_ctx.obj << " , ret = " << ret << dendl;
- return ret;
+ ret = retry_on_busy(tier_ctx.y, tier_ctx.dpp, tier_ctx.cct, __func__, [&]() -> int {
+ if (!in_progress) { // first time. Send RESTORE req.
+ ret = cloud_tier_restore(tier_ctx.dpp, tier_ctx.conn, dest_obj, days, glacier_params, tier_ctx.y);
+ ldpp_dout(tier_ctx.dpp, 20) << __func__ << "Restoring object=" << target_obj_name << "returned ret = " << ret << dendl;
+ if (ret < 0 ) {
+ ldpp_dout(tier_ctx.dpp, -1) << __func__ << "ERROR: failed to restore object=" << dest_obj << "; ret = " << ret << dendl;
+ return ret;
+ }
+ in_progress = true;
}
- in_progress = is_restore_in_progress(tier_ctx.dpp, headers);
-
- } while(retries++ < MAX_RETRIES && in_progress);
+ // now send HEAD request and verify if restore is complete on glacier/tape endpoint
+ static constexpr int MAX_RETRIES = 2;
+ uint32_t retries = 0;
+ do {
+ ret = rgw_cloud_tier_get_object(tier_ctx, true, headers, nullptr, etag,
+ accounted_size, attrs, nullptr);
+ if (ret < 0) {
+ ldpp_dout(tier_ctx.dpp, 0) << __func__ << "ERROR: failed to fetch HEAD from cloud for obj=" << tier_ctx.obj << " , ret = " << ret << dendl;
+ return ret;
+ }
+ in_progress = is_restore_in_progress(tier_ctx.dpp, headers);
+ } while(retries++ < MAX_RETRIES && in_progress);
+ return 0;
+ });
+ if (ret < 0) return ret;
if (in_progress) {
ldpp_dout(tier_ctx.dpp, 20) << __func__ << "Restoring object=" << target_obj_name << " still in progress; returning " << dendl;
return 0;
- }
+ }
- // now do the actual GET
- ret = rgw_cloud_tier_get_object(tier_ctx, false, headers, pset_mtime, etag,
- accounted_size, attrs, cb);
+ ret = retry_on_busy(tier_ctx.y, tier_ctx.dpp, tier_ctx.cct, __func__, [&]() {
+ return rgw_cloud_tier_get_object(tier_ctx, false, headers, pset_mtime,
+ etag, accounted_size, attrs, cb);
+ });
ldpp_dout(tier_ctx.dpp, 20) << __func__ << "(): fetching object from cloud bucket:" << dest_bucket << ", object: " << target_obj_name << " returned ret:" << ret << dendl;
tier_ctx.obj->set_atomic(true);
- /* TODO: Define readf, writef as stack variables. For some reason,
- * when used as stack variables (esp., readf), the transition seems to
- * be taking lot of time eventually erroring out at times. */
- std::shared_ptr<RGWLCStreamRead> readf;
- readf.reset(new RGWLCStreamRead(tier_ctx.cct, tier_ctx.dpp,
- tier_ctx.obj, tier_ctx.o.meta.mtime, tier_ctx.y));
-
- std::shared_ptr<RGWLCCloudStreamPut> writef;
- writef.reset(new RGWLCCloudStreamPut(tier_ctx.dpp, obj_properties, tier_ctx.conn,
- dest_obj, tier_ctx.y));
-
- /* Prepare Read from source */
end = part_info.ofs + part_info.size - 1;
- readf->set_multipart(part_info.size, part_info.ofs, end);
-
- /* Prepare write */
- writef->set_multipart(upload_id, part_info.part_num, part_info.size);
-
- /* actual Read & Write */
- ret = cloud_tier_transfer_object(tier_ctx.dpp, readf.get(), writef.get());
- if (ret < 0) {
- return ret;
- }
+ std::shared_ptr<RGWLCCloudStreamPut> writef;
+ ret = retry_on_busy(tier_ctx.y, tier_ctx.dpp, tier_ctx.cct, __func__, [&]() {
+ auto rf = std::make_shared<RGWLCStreamRead>(tier_ctx.cct, tier_ctx.dpp,
+ tier_ctx.obj, tier_ctx.o.meta.mtime, tier_ctx.y);
+ auto wf = std::make_shared<RGWLCCloudStreamPut>(tier_ctx.dpp,
+ obj_properties, tier_ctx.conn, dest_obj, tier_ctx.y);
+ rf->set_multipart(part_info.size, part_info.ofs, end);
+ wf->set_multipart(upload_id, part_info.part_num, part_info.size);
+ int r = cloud_tier_transfer_object(tier_ctx.dpp, rf.get(), wf.get());
+ if (r == 0) writef = wf;
+ return r;
+ });
+ if (ret < 0) return ret;
if (!(writef->get_etag(petag))) {
ldpp_dout(tier_ctx.dpp, 0) << "ERROR: failed to get etag from PUT request" << dendl;
ret = dest_conn.send_resource(dpp, "POST", resource, params, nullptr,
out_bl, &bl, nullptr, y);
+ if (ret == -EBUSY) return ret;
if (ret < 0) {
ldpp_dout(dpp, 0) << __func__ << "ERROR: failed to send Restore request to cloud for obj=" << dest_obj << " , ret = " << ret << dendl;
} else {
cur_part_info,
&cur_part_info.etag);
+ if (ret == -EBUSY) return ret;
if (ret < 0) {
ldpp_dout(tier_ctx.dpp, 0) << "ERROR: failed to send multipart part of obj=" << tier_ctx.obj << ", sync via multipart upload, upload_id=" << status.upload_id << " part number " << cur_part << " (error: " << cpp_strerror(-ret) << ")" << dendl;
cloud_tier_abort_multipart_upload(tier_ctx, dest_obj, status_obj, status.upload_id);
}
ret = cloud_tier_complete_multipart(tier_ctx.dpp, tier_ctx.conn, dest_obj, status.upload_id, parts, tier_ctx.y);
+ if (ret == -EBUSY) return ret;
if (ret < 0) {
ldpp_dout(tier_ctx.dpp, 0) << "ERROR: failed to complete multipart upload of obj=" << tier_ctx.obj << " (error: " << cpp_strerror(-ret) << ")" << dendl;
cloud_tier_abort_multipart_upload(tier_ctx, dest_obj, status_obj, status.upload_id);
ret = tier_ctx.conn.send_resource(tier_ctx.dpp, "PUT", resource, nullptr, nullptr,
out_bl, &bl, nullptr, tier_ctx.y);
- if (ret < 0 ) {
+ if (ret == -EBUSY) return ret;
+ if (ret < 0) {
ldpp_dout(tier_ctx.dpp, 0) << "create target bucket : " << tier_ctx.target_bucket_name << " returned ret:" << ret << dendl;
}
if (out_bl.length() > 0) {
return 0;
}
-int rgw_cloud_tier_transfer_object(RGWLCCloudTierCtx& tier_ctx, std::set<std::string>& cloud_targets) {
+static int do_cloud_tier_transfer_object(RGWLCCloudTierCtx& tier_ctx, std::set<std::string>& cloud_targets) {
int ret = 0;
// check if target bucket is in local cache
bool already_tiered = false;
ret = cloud_tier_check_object(tier_ctx, already_tiered);
+ if (ret == -EBUSY) return ret;
if (ret < 0) {
ldpp_dout(tier_ctx.dpp, 0) << "ERROR: failed to check object on the cloud endpoint ret=" << ret << dendl;
}
return ret;
}
+
+int rgw_cloud_tier_transfer_object(RGWLCCloudTierCtx& tier_ctx, std::set<std::string>& cloud_targets) {
+ return retry_on_busy(tier_ctx.y, tier_ctx.dpp, tier_ctx.cct, __func__,
+ [&]() { return do_cloud_tier_transfer_object(tier_ctx, cloud_targets); });
+}