diff --git a/bindings/python/elliptics_python.cpp b/bindings/python/elliptics_python.cpp index fde86f666..790142c9e 100644 --- a/bindings/python/elliptics_python.cpp +++ b/bindings/python/elliptics_python.cpp @@ -44,11 +44,13 @@ enum elliptics_iterator_types { }; enum elliptics_iterator_flags { - iflag_default = 0, - iflag_data = DNET_IFLAGS_DATA, - iflag_key_range = DNET_IFLAGS_KEY_RANGE, - iflag_ts_range = DNET_IFLAGS_TS_RANGE, - iflag_no_meta = DNET_IFLAGS_NO_META, + iflag_default = 0, + iflag_data = DNET_IFLAGS_DATA, + iflag_key_range = DNET_IFLAGS_KEY_RANGE, + iflag_ts_range = DNET_IFLAGS_TS_RANGE, + iflag_no_meta = DNET_IFLAGS_NO_META, + iflags_move = DNET_IFLAGS_MOVE, + iflags_overwrite = DNET_IFLAGS_OVERWRITE }; enum elliptics_cflags { @@ -377,13 +379,20 @@ BOOST_PYTHON_MODULE(core) "default\n There no filtering should be while iteration. All keys will be presented\n" "data\n Iteration results should also includes objects datas\n" "key_range\n elliptics.Id ranges should be used for filtering keys on the node while iteration\n" - "ts_range\n Time range should be used for filtering keys on the node while iteration" - "no_meta\n Iteration results will have empty key's metadata (user_flags and timestamp)") + "ts_range\n Time range should be used for filtering keys on the node while iteration\n" + "no_meta\n Iteration results will have empty key's metadata (user_flags and timestamp)\n" + "move\n Server-send iterator should move data not copy. This will force iterator/server-send logic\n" + " to queue REMOVE command locally if remote write has succeeeded.\n" + "overwrite\n Overwrite data. If this flag is NOT set, we only write data if remote timestamp is less\n" + " than in data being written. When NOT set, data will still be transferred over the network,\n" + " even if remote timestamp doesn't allow us to overwrite data.") .value("default", iflag_default) .value("data", iflag_data) .value("key_range", iflag_key_range) .value("ts_range", iflag_ts_range) .value("no_meta", iflag_no_meta) + .value("move", iflags_move) + .value("overwrite", iflags_overwrite) ; bp::enum_("iterator_types", diff --git a/bindings/python/elliptics_session.cpp b/bindings/python/elliptics_session.cpp index 5aff374b8..60a62c443 100644 --- a/bindings/python/elliptics_session.cpp +++ b/bindings/python/elliptics_session.cpp @@ -502,6 +502,17 @@ class elliptics_session: public session, public bp::wrapper { return create_result(std::move(session::start_iterator(transform(id).id(), std_ranges, type, flags, time_begin.m_time, time_end.m_time))); } + python_iterator_result start_copy_iterator(const bp::api::object &id, const bp::api::object &ranges, + const bp::api::object &dst_groups, + uint64_t flags, + const elliptics_time& time_begin = elliptics_time(0, 0), + const elliptics_time& time_end = elliptics_time(-1, -1)) { + auto std_ranges = convert_to_vector(ranges); + auto std_dst_groups = convert_to_vector(dst_groups); + + return create_result(std::move(session::start_copy_iterator(transform(id).id(), std_ranges, flags, time_begin.m_time, time_end.m_time, std_dst_groups))); + } + python_iterator_result pause_iterator(const bp::api::object &id, const uint64_t &iterator_id) { return create_result(std::move(session::pause_iterator(transform(id).id(), iterator_id))); } @@ -514,6 +525,18 @@ class elliptics_session: public session, public bp::wrapper { return create_result(std::move(session::cancel_iterator(transform(id).id(), iterator_id))); } + python_iterator_result server_send(const bp::api::object &keys, uint64_t iflags, const bp::api::object &groups) { + auto std_groups = convert_to_vector(groups); + std::vector std_keys; + std_keys.reserve(bp::len(keys)); + + for (bp::stl_input_iterator it(keys), end; it != end; ++it) { + std_keys.push_back(transform(*it).raw_id()); + } + + return create_result(std::move(session::server_send(std_keys, iflags, std_groups))); + } + python_exec_result exec(const bp::api::object &id_or_context, const std::string &event, const bp::api::object &data, const int src_key) { dnet_id* raw_id = NULL; dnet_id conv_id; @@ -1521,6 +1544,41 @@ void init_elliptics_session() { " result.response.timestamp.tnsec,\n" " result.response_data))\n") + .def("start_copy_iterator", &elliptics_session::start_copy_iterator, + bp::args("id", "ranges", "dst_groups", "flags", "time_begin", "time_end"), + "start_copy_iterator(id, ranges, dst_groups, flags, time_begin, time_end)\n" + " Start copy iterator on the Elliptics node specified by @id. Return elliptics.AsyncResult.\n" + " -- id - elliptics.Id of the node where iteration should be executed\n" + " -- ranges - list of elliptics.IteratorRange by which keys on the node should be filtered\n" + " -- dst_groups - list of remote groups where data will be copied/moved\n" + " -- flags - bits set of elliptics.iterator_flags\n" + " -- time_begin - start of time range by which keys on the node should be filtered\n" + " -- time_end - end of time range by which keys on the node should be filtered\n\n" + " flags = elliptics.iterator_flags.key_range\n" + " id = session.routes.get_address_id(Address.from_host_port('host.com:1025'))\n" + " range = elliptics.IteratorRange()\n" + " range.key_begin = elliptics.Id([0] * 64, 1)\n" + " range.key_end = elliptics.Id([255] * 64, 1)\n" + " dst_groups = [2,3]\n" + " iterator = session.start_copy_iterator(id,\n" + " [range],\n" + " type,\n" + " flags,\n" + " elliptics.Time(0,0),\n" + " elliptics.Time(0,0),\n" + " dst_groups)\n\n" + " for result in iterator:\n" + " if result.status != 0:\n" + " raise AssertionError('Wrong status: {0}'.format(result.status))\n\n" + " iterator_id = result.id\n" + " print ('node: {0}, key: {1}, flags: {2}, ts: {3}/{4}, data: {5}'\n" + " .format(node,\n" + " result.response.key,\n" + " result.response.user_flags,\n" + " result.response.timestamp.tsec,\n" + " result.response.timestamp.tnsec,\n" + " result.response_data))\n") + .def("pause_iterator", &elliptics_session::pause_iterator, bp::args("id", "iterator_id"), "pause_iterator(id, iterator_id)\n" @@ -1561,6 +1619,19 @@ void init_elliptics_session() { " iterator = session.cancel_iterator(id, iterator_id)\n" " iterator.wait()\n") +// Server send operations + + .def("server_send", &elliptics_session::server_send, + bp::args("keys", "iflags", "groups"), + "server_send(keys, iflags, groups)\n" + " Similar to iterator, but instead of running over all keys on remote backend,\n" + " remote server nodes will read all specified keys (which live on local backends)\n" + " and send them to remote nodes.\n" + " Returns elliptics.AsyncResult.\n" + " -- keys - iterable object which provides set of elliptics keys (elliptics.Id)\n" + " -- iflags - bits set of elliptics.iterator_flags\n" + " -- groups - iterable object which specifies groups to which data should be send\n") + // Index operations .def("set_indexes", &elliptics_session::set_indexes, diff --git a/example/eblob_backend.c b/example/eblob_backend.c index 62b143133..5bc253666 100644 --- a/example/eblob_backend.c +++ b/example/eblob_backend.c @@ -965,6 +965,7 @@ static int blob_send(struct eblob_backend_config *cfg, void *state, struct dnet_ struct dnet_server_send_ctl *ctl; int *groups; int i, err; + int backend_id; struct dnet_ext_list elist; static const size_t ehdr_size = sizeof(struct dnet_ext_list_hdr); @@ -981,9 +982,6 @@ static int blob_send(struct eblob_backend_config *cfg, void *state, struct dnet_ ids = (struct dnet_raw_id *)(req + 1); groups = (int *)(ids + req->id_num); - memset(&re, 0, sizeof(struct dnet_iterator_response)); - re.total_keys = req->id_num; - /* * Set NEED_ACK bit to signal server-send controller that we want * to send final ACK when controller will be destroyed, which in turn @@ -994,7 +992,8 @@ static int blob_send(struct eblob_backend_config *cfg, void *state, struct dnet_ */ cmd->flags |= DNET_FLAGS_NEED_ACK; - ctl = dnet_server_send_alloc(state, cmd, req->iflags, groups, req->group_num); + backend_id = cfg->data.stat_id; + ctl = dnet_server_send_alloc(state, cmd, req->iflags, groups, req->group_num, backend_id); if (!ctl) { err = -ENOMEM; goto err_out_exit; @@ -1018,6 +1017,14 @@ static int blob_send(struct eblob_backend_config *cfg, void *state, struct dnet_ for (i = 0; i < req->id_num; ++i) { + memset(&re, 0, sizeof(struct dnet_iterator_response)); + // set iterator response id to differentiate various commands + // client can use cmd->backend_id from reply though + re.id = cmd->backend_id; + re.key = ids[i]; + re.iterated_keys = i; + re.total_keys = req->id_num; + memcpy(key.id, ids[i].id, EBLOB_ID_SIZE); err = blob_lookup(b, &key, &wc); @@ -1027,14 +1034,8 @@ static int blob_send(struct eblob_backend_config *cfg, void *state, struct dnet_ goto err_out_send_fail_reply; } - re.key = ids[i]; re.flags = wc.flags; // these flags correspond to DNET_RECORD_FLAGS_* - re.status = 0; - re.iterated_keys = i; re.size = wc.total_data_size; - // set iterator response id to differentiate various commands - // client can use cmd->backend_id from reply though - re.id = cmd->backend_id; data_offset = wc.data_offset; record_offset = 0; diff --git a/include/elliptics/interface.h b/include/elliptics/interface.h index b48de0471..4c02131e3 100644 --- a/include/elliptics/interface.h +++ b/include/elliptics/interface.h @@ -943,7 +943,7 @@ int dnet_get_vm_stat(dnet_logger *l, struct dnet_vm_stat *st); struct dnet_server_send_ctl; struct dnet_server_send_ctl *dnet_server_send_alloc(void *state, struct dnet_cmd *cmd, uint64_t iflags, - int *groups, int group_num); + int *groups, int group_num, int backend_id); struct dnet_server_send_ctl *dnet_server_send_get(struct dnet_server_send_ctl *ctl); int dnet_server_send_put(struct dnet_server_send_ctl *ctl); int dnet_server_send_write(struct dnet_server_send_ctl *send, diff --git a/include/elliptics/packet.h b/include/elliptics/packet.h index 880fed2ef..dadba8dc6 100644 --- a/include/elliptics/packet.h +++ b/include/elliptics/packet.h @@ -1025,7 +1025,7 @@ enum { #define DNET_IFLAGS_MOVE (1<<4) /* * Overwrite data. If this flag is NOT set, we only write data if remote timestamp is less - * that that in data being written. When NOT set, data will still be transferred over the network, + * than in data being written. When NOT set, data will still be transferred over the network, * even if remote timestamp doesn't allow us to overwrite data. */ #define DNET_IFLAGS_OVERWRITE (1<<5) diff --git a/library/dnet.c b/library/dnet.c index 0144a308f..169c1f977 100644 --- a/library/dnet.c +++ b/library/dnet.c @@ -476,9 +476,9 @@ static int dnet_iterator_server_send_complete(struct dnet_addr *addr, struct dne lc = r->header; dnet_setup_id(&lc->id, send->cmd.id.group_id, cmd->id.id); lc->cmd = DNET_CMD_DEL; - lc->backend_id = -1; + lc->backend_id = send->backend_id; lc->trace_id = cmd->trace_id; - lc->flags = DNET_FLAGS_NOLOCK; + lc->flags = DNET_FLAGS_NOLOCK | DNET_FLAGS_DIRECT | DNET_FLAGS_DIRECT_BACKEND; if (send->cmd.flags & DNET_FLAGS_TRACE_BIT) lc->flags |= DNET_FLAGS_TRACE_BIT; lc->size = sizeof(struct dnet_io_attr); @@ -530,7 +530,7 @@ static int dnet_iterator_server_send_complete(struct dnet_addr *addr, struct dne } struct dnet_server_send_ctl *dnet_server_send_alloc(void *state, struct dnet_cmd *cmd, uint64_t iflags, - int *groups, int group_num) + int *groups, int group_num, int backend_id) { int err; struct dnet_net_state *st = state; @@ -549,6 +549,7 @@ struct dnet_server_send_ctl *dnet_server_send_alloc(void *state, struct dnet_cmd ctl->state = state; ctl->cmd = *cmd; ctl->iflags = iflags; + ctl->backend_id = backend_id; ctl->groups = (int *)(ctl + 1); memcpy(ctl->groups, groups, sizeof(int) * group_num); ctl->group_num = group_num; @@ -1112,7 +1113,7 @@ static int dnet_iterator_start(struct dnet_backend_io *backend, struct dnet_net_ * dnet_cmd, thus it will store command structure without NEED_ACK bit. */ cmd->flags &= ~DNET_FLAGS_NEED_ACK; - sspriv = dnet_server_send_alloc(st, cmd, ireq->flags, dst_groups, ireq->group_num); + sspriv = dnet_server_send_alloc(st, cmd, ireq->flags, dst_groups, ireq->group_num, backend->backend_id); cmd->flags |= DNET_FLAGS_NEED_ACK; if (!sspriv) { diff --git a/library/elliptics.h b/library/elliptics.h index 9876a8c6c..524c93eae 100644 --- a/library/elliptics.h +++ b/library/elliptics.h @@ -1049,6 +1049,7 @@ struct dnet_server_send_ctl { uint64_t iflags; /* Iterator flags */ + int backend_id; /* Source backend_id */ int *groups; /* Groups to send WRITE commands */ int group_num; diff --git a/recovery/elliptics_recovery/iterator.py b/recovery/elliptics_recovery/iterator.py index 7207b0466..5cc2d8ecb 100644 --- a/recovery/elliptics_recovery/iterator.py +++ b/recovery/elliptics_recovery/iterator.py @@ -239,7 +239,7 @@ def __init__(self, node, group, separately=False, trace_id=0): self.session.trace_id = trace_id self.separately = separately - def get_key_range_id(self, key): + def _get_key_range_id(self, key): if not self.separately: return 0 @@ -260,7 +260,6 @@ def get_key_range_id(self, key): def start(self, eid=IdRange.ID_MIN, - itype=elliptics.iterator_types.network, flags=elliptics.iterator_flags.key_range | elliptics.iterator_flags.ts_range, key_ranges=(IdRange(IdRange.ID_MIN, IdRange.ID_MAX),), timestamp_range=(Time.time_min().to_etime(), Time.time_max().to_etime()), @@ -270,7 +269,6 @@ def start(self, group_id=0, leave_file=False, batch_size=1024): - assert itype == elliptics.iterator_types.network, "Only network iterator is supported for now" assert flags & elliptics.iterator_flags.data == 0, "Only metadata iterator is supported for now" assert len(key_ranges) > 0, "There should be at least one iteration range." self.ranges = key_ranges @@ -300,12 +298,11 @@ def start(self, leave_file=leave_file) ranges = [IdRange.elliptics_range(start, stop) for start, stop in key_ranges] - records = self.session.start_iterator(eid, - ranges, - itype, - flags, - timestamp_range[0], - timestamp_range[1]) + records = self._start_iterator(eid, + ranges, + flags, + timestamp_range) + iterated_keys = 0 total_keys = 0 @@ -322,9 +319,8 @@ def start(self, if iterated_keys % batch_size == 0: yield (iterated_keys, total_keys, start, end) - if record.response.status != 0: - continue - results[self.get_key_range_id(record.response.key)].append(record) + + self._on_key_response(results, record) end = time.time() elapsed_time = records.elapsed_time() @@ -339,23 +335,34 @@ def start(self, .format(address, backend_id, repr(e), traceback.format_exc())) yield None - @classmethod - def iterate_with_stats(cls, node, eid, timestamp_range, + def _start_iterator(self, eid, ranges, flags, timestamp_range): + return self.session.start_iterator(eid, + ranges, + elliptics.iterator_types.network, + flags, + timestamp_range[0], + timestamp_range[1]) + + def _on_key_response(self, results, record): + if record.response.status == 0: + self._save_record(results, record) + + def _save_record(self, results, record): + results[self._get_key_range_id(record.response.key)].append(record) + + def iterate_with_stats(self, eid, timestamp_range, key_ranges, tmp_dir, address, group_id, backend_id, batch_size, - stats, flags, leave_file=False, - separately=False, trace_id=0): - iterator = cls(node, group_id, separately, trace_id=trace_id) - result = iterator.start(eid=eid, - timestamp_range=timestamp_range, - flags=flags, - key_ranges=key_ranges, - tmp_dir=tmp_dir, - address=address, - backend_id=backend_id, - group_id=group_id, - batch_size=batch_size, - leave_file=leave_file, - ) + stats, flags, leave_file=False): + result = self.start(eid=eid, + flags=flags, + key_ranges=key_ranges, + timestamp_range=timestamp_range, + tmp_dir=tmp_dir, + address=address, + backend_id=backend_id, + group_id=group_id, + leave_file=leave_file, + batch_size=batch_size,) result_len = 0 for it in result: if it is None: @@ -379,6 +386,24 @@ def iterate_with_stats(cls, node, eid, timestamp_range, return result, result_len +class MergeRecoveryIterator(Iterator): + ''' + This class is used in merge recovery for backend iteratation on ranges which are not belong to + it using copy iterator. Every iterated key is moved to the backend, where it should exists. + If moving of some key was failed, then it saves the key to the results container. + ''' + def __init__(self, *args, **kwargs): + super(MergeRecoveryIterator, self).__init__(*args, **kwargs) + + def _start_iterator(self, eid, ranges, flags, timestamp_range): + flags |= elliptics.iterator_flags.move + return self.session.start_copy_iterator(eid, ranges, [eid.group_id], flags, timestamp_range[0], timestamp_range[1]) + + def _on_key_response(self, results, record): + if record.response.status != 0: + self._save_record(results, record) + + class MergeData(object): """ Assist class for IteratorResult.__merge__ diff --git a/recovery/elliptics_recovery/types/dc.py b/recovery/elliptics_recovery/types/dc.py index 26cb8277a..fcc7d7e86 100644 --- a/recovery/elliptics_recovery/types/dc.py +++ b/recovery/elliptics_recovery/types/dc.py @@ -57,8 +57,8 @@ def iterate_node(arg): flags |= elliptics.iterator_flags.ts_range log.debug("Running iterator on node: {0}/{1}".format(address, backend_id)) - results, results_len = Iterator.iterate_with_stats( - node=node, + iterator = Iterator(node, node_id.group_id, separately=True, trace_id=ctx.trace_id) + results, results_len = iterator.iterate_with_stats( eid=node_id, timestamp_range=timestamp_range, key_ranges=ranges, @@ -69,9 +69,7 @@ def iterate_node(arg): batch_size=ctx.batch_size, stats=stats, flags=flags, - leave_file=True, - separately=True, - trace_id=ctx.trace_id) + leave_file=True) if results is None or results_len == 0: return None diff --git a/recovery/elliptics_recovery/types/merge.py b/recovery/elliptics_recovery/types/merge.py index 6a47db9a5..59f035d4e 100755 --- a/recovery/elliptics_recovery/types/merge.py +++ b/recovery/elliptics_recovery/types/merge.py @@ -30,11 +30,13 @@ from itertools import groupby import traceback import threading +import errno +from bisect import bisect from ..etime import Time from ..utils.misc import elliptics_create_node, RecoverStat, LookupDirect, RemoveDirect, WindowedRecovery from ..route import RouteList -from ..iterator import Iterator +from ..iterator import MergeRecoveryIterator from ..range import IdRange import elliptics @@ -328,14 +330,10 @@ def iterate_node(ctx, node, address, backend_id, ranges, eid, stats): try: log.debug("Running iterator on node: {0}/{1}".format(address, backend_id)) timestamp_range = ctx.timestamp.to_etime(), Time.time_max().to_etime() - flags = elliptics.iterator_flags.key_range - if ctx.no_meta: - flags |= elliptics.iterator_flags.no_meta - else: - flags |= elliptics.iterator_flags.ts_range + flags = elliptics.iterator_flags.key_range | elliptics.iterator_flags.ts_range key_ranges = [IdRange(r[0], r[1]) for r in ranges] - result, result_len = Iterator.iterate_with_stats(node=node, - eid=eid, + iterator = MergeRecoveryIterator(node, eid.group_id, trace_id=ctx.trace_id) + result, result_len = iterator.iterate_with_stats(eid=eid, timestamp_range=timestamp_range, key_ranges=key_ranges, tmp_dir=ctx.tmp_dir, @@ -345,8 +343,7 @@ def iterate_node(ctx, node, address, backend_id, ranges, eid, stats): batch_size=ctx.batch_size, stats=stats, flags=flags, - leave_file=False, - trace_id=ctx.trace_id) + leave_file=False,) if result is None: return None log.info("Iterator {0}/{1} obtained: {2} record(s)" @@ -686,6 +683,130 @@ def succeeded(self): return self.result +class ServerSendRecovery(object): + ''' + Special recovery class that tries to recover keys from backends that + should not contain this keys to proper backend via server-send operation. + ''' + def __init__(self, ctx, node, group): + self.routes = self._prepare_routes(ctx, group) + self.session = elliptics.Session(node) + self.session.exceptions_policy = elliptics.exceptions_policy.no_exceptions + self.session.set_filter(elliptics.filters.all) + self.session.timeout = 60 + self.session.groups = [group] + self.session.trace_id = ctx.trace_id + self.ctx = ctx + + def _prepare_routes(self, ctx, group): + ''' + Returns list of triplets (address, backend, [ranges]), + where ranges are sorted by their left boundary. + ''' + def sort_ranges(ranges): + ranges = sorted(ranges, key=lambda r: r[0]) + return reduce(lambda x, y: x + y, ranges, tuple()) + + group_routes = ctx.routes.filter_by_groups([group]) + routes = [] + if ctx.one_node: + if ctx.backend_id is not None: + ranges = group_routes.get_address_backend_ranges(ctx.address, ctx.backend_id) + routes = [(ctx.address, ctx.backend_id, sort_ranges(ranges))] + else: + for backend_id in group_routes.get_address_backends(ctx.address): + ranges = group_routes.get_address_backend_ranges(ctx.address, backend_id) + routes.append((ctx.address, backend_id, sort_ranges(ranges))) + else: + for addr, backend_id in group_routes.addresses_with_backends(): + ranges = group_routes.get_address_backend_ranges(addr, backend_id) + routes.append((addr, backend_id, sort_ranges(ranges))) + + log.info("Server-send recovery: group: {0}, num addresses: {1}".format(group, len(routes))) + return routes + + def recover(self, keys): + ''' + Tries to recover keys from every backend via server-send. Then it + removes keys with older timestamp or invalid checksum. + Returns list of keys that was not recovered via server-send. + ''' + log.info("Server-send bucket: num keys: {0}".format(len(keys))) + + def contain(key, ranges): + index = bisect(ranges, key) + return index % 2 == 1 + + responses = dict([(str(k), []) for k in keys]) # key -> [list of responses] + for addr, backend_id, backend_ranges in self.routes: + key_candidates = [k for k in keys if not contain(k, backend_ranges)] + if key_candidates: + self._server_send(key_candidates, addr, backend_id, responses) + + self._remove_bad_keys(responses) + return self._get_unrecovered_keys(responses) + + def _server_send(self, keys, addr, backend_id, responses): + ''' + Calls server-send with a given list of keys to the specific backend. + ''' + log.debug("Server-send: address: {0}, backend: {1}, num keys: {2}".format(addr, backend_id, len(keys))) + + self.session.set_direct_id(addr, backend_id) + iterator = self.session.server_send(keys, elliptics.iterator_flags.move, list(self.session.groups)) + for result in iterator: + status = result.response.status + key = result.response.key + r = (key, status, addr, backend_id) + log.debug("Server-send result: key: {0}, status: {1}".format(str(key), status)) + responses[str(key)].append(r) + + def _remove_bad_keys(self, responses): + ''' + Removes invalid keys with older timestamp or invalid checksum. + ''' + bad_keys = [] + for val in responses.itervalues(): + bad_keys.extend([r for r in val if self._check_bad_key(r)]) + + results = [] + for k in bad_keys: + key, _, addr, backend_id = k + self.session.set_direct_id(addr, backend_id) + result = self.session.remove(key) + results.append(result) + + for i, r in enumerate(results): + status = r.get()[0].status + log.info("Removing key: {0}, status: ".format(bad_keys[i], status)) + + def _check_bad_key(self, response): + status = response[1] + return status == -errno.EBADFD or status == -errno.EILSEQ + + def _get_unrecovered_keys(self, responses): + ''' + Returns keys that was not recovered via server-send. + ''' + keys = [] + for val in responses.iteritems(): + key_responses = val[1] + if not key_responses or self._check_unrecovered_key(key_responses): + key = elliptics.Id(val[0]) + keys.append(key) + return keys + + def _check_unrecovered_key(self, responses): + ''' + Returns True, if a valid key exists on the backend, but the key could not be recovered by any reason. + ''' + for r in responses: + status = r[1] + if status < 0 and status != -errno.ENOENT and not self._check_bad_key(r): + return True + return False + + def dump_process_group((ctx, group)): log.debug("Processing group: {0}".format(group)) stats = ctx.stats['group_{0}'.format(group)] @@ -702,12 +823,15 @@ def dump_process_group((ctx, group)): remotes=ctx.remotes) ret = True with open(ctx.dump_file, 'r') as dump: + ss_rec = ServerSendRecovery(ctx, node, group) # splits ids from dump file in batchs and recovers it for batch_id, batch in groupby(enumerate(dump), key=lambda x: x[0] / ctx.batch_size): recovers = [] rs = RecoverStat() - for _, val in batch: - rec = DumpRecover(routes=ctx.routes, node=node, id=elliptics.Id(val), group=group, ctx=ctx) + keys = [elliptics.Id(val) for _, val in batch] + keys = ss_rec.recover(keys) + for k in keys: + rec = DumpRecover(routes=ctx.routes, node=node, id=k, group=group, ctx=ctx) recovers.append(rec) rec.run() for r in recovers: diff --git a/tests/pytests/test_recovery.py b/tests/pytests/test_recovery.py index b381050d2..3f5dc90b3 100644 --- a/tests/pytests/test_recovery.py +++ b/tests/pytests/test_recovery.py @@ -17,6 +17,7 @@ import os import sys +import errno sys.path.insert(0, "") # for running from cmake import pytest from conftest import make_session @@ -102,6 +103,27 @@ def write_data(scope, session, keys, datas): r.wait() +def check_keys_absence(scope, session, keys): + ''' + Checks that merge recovery removes moved @keys from the source backend. + ''' + session = session.clone() + session.exceptions_policy = elliptics.core.exceptions_policy.no_exceptions + session.set_filter(elliptics.filters.all) + session.set_direct_id(scope.test_address, scope.test_backend) + + routes = session.routes.filter_by_group(scope.test_group) + results = [] + for k in keys: + addr, _, backend = routes.get_id_routes(session.transform(k))[0] + if addr != scope.test_address or backend != scope.test_backend: + results.append(session.lookup(k)) + + assert len(results) > 0 + for r in results: + assert r.get()[0].status == -errno.ENOENT + + def check_data(scope, session, keys, datas, timestamp): ''' Reads @keys from the session. Reads all keys async at once and waits/checks results at the end. @@ -311,6 +333,7 @@ def test_merge_two_backends(self, server, simple_node): session.groups = (scope.test_group,) check_data(scope, session, self.keys, self.datas, self.timestamp) + check_keys_absence(scope, session, self.keys) def test_enable_another_one_backend(self, server, simple_node): ''' @@ -353,6 +376,7 @@ def test_merge_from_dump_3_backends(self, server, simple_node): session.groups = (scope.test_group,) check_data(scope, session, self.keys, self.datas, self.timestamp) + check_keys_absence(scope, session, self.keys) def test_enable_all_group_backends(self, server, simple_node): ''' @@ -384,6 +408,7 @@ def test_merge_one_group(self, server, simple_node): session.groups = (scope.test_group,) check_data(scope, session, self.keys, self.datas, self.timestamp) + check_keys_absence(scope, session, self.keys) def test_enable_second_group_one_backend(self, server, simple_node): '''