diff --git a/simplyblock_core/rpc_client.py b/simplyblock_core/rpc_client.py index 864dd46c39..325726ad3c 100644 --- a/simplyblock_core/rpc_client.py +++ b/simplyblock_core/rpc_client.py @@ -2034,14 +2034,23 @@ def bdev_lvol_s3_backup(self, s3_id, snapshot_names, cluster_batch=1): } return self._request("bdev_lvol_s3_backup", params) - # Backup/recovery/merge polling: use bdev_lvol_transfer_stat(lvol_name) - # which reads lvol->transfer_status on the data plane. Works for backup - # (pass snapshot bdev name) and recovery (pass target lvol name). - # Merge has lvol=NULL on data plane so transfer_stat cannot poll it. + # Backup/recovery polling: use bdev_lvol_transfer_stat(lvol_name) which + # reads lvol->transfer_status on the data plane. Works for backup (pass + # snapshot bdev name) and recovery (pass target lvol name). Merge has + # lvol=NULL on data plane, so it's polled separately via + # bdev_lvol_s3_merge_stat(s3_id, old_s3_id) below. - def bdev_lvol_s3_merge(self, s3_id, old_s3_id, cluster_batch, lvs_name=None): + def bdev_lvol_s3_merge(self, s3_id, old_s3_id, cluster_batch, lvs_name=None, allow_exist: bool = True): """Merge two backups: keep s3_id and merge old_s3_id into it. - This shortens the backup chain.""" + This shortens the backup chain. + + allow_exist: if True (default), an EEXIST response -- a matching + merge already queued/running on the data plane, e.g. from a prior + call whose RPC connection dropped before the response arrived -- is + treated the same as a fresh success (returns True) rather than + raised. A caller that needs to distinguish "started this call" from + "already running" can pass allow_exist=False. + """ params = { "s3_id": s3_id, "old_s3_id": old_s3_id, @@ -2049,7 +2058,21 @@ def bdev_lvol_s3_merge(self, s3_id, old_s3_id, cluster_batch, lvs_name=None): } if lvs_name: params["lvs_name"] = lvs_name - return self._request("bdev_lvol_s3_merge", params) + try: + return self._request3("bdev_lvol_s3_merge", **params) + except RPCRemoteError as e: + if allow_exist and e.code == -17: + logger.debug("Merge %s -> %s already in progress", old_s3_id, s3_id) + return True + raise + + def bdev_lvol_s3_merge_stat(self, s3_id, old_s3_id): + """Return merge status for the (s3_id, old_s3_id) pair. + + Result dict keys: + ``transfer_state``: "No process" | "In progress" | "Failed" | "Done" + """ + return self._request("bdev_lvol_s3_merge_stat", {"s3_id": s3_id, "old_s3_id": old_s3_id}) def bdev_lvol_s3_recovery(self, lvol_name, s3_ids, cluster_batch): """Restore a chain of S3 backups into a new lvol. diff --git a/simplyblock_core/services/tasks_runner_backup.py b/simplyblock_core/services/tasks_runner_backup.py index 8f0ed8af05..4e73bb405c 100644 --- a/simplyblock_core/services/tasks_runner_backup.py +++ b/simplyblock_core/services/tasks_runner_backup.py @@ -347,23 +347,62 @@ def _run_merge(task): task.function_params["merge_started"] = True task.write_to_db(db.kv_store) - # Give the data plane time to complete the merge before finalizing return - # The merge RPC is synchronous on the data plane — once it returned - # successfully, the S3 data has been merged. Finalize: update the - # chain links, remove the old backup, and mark the task done. - keep_backup.prev_backup_id = old_backup.prev_backup_id - keep_backup.status = Backup.STATUS_COMPLETED - keep_backup.write_to_db() + # The merge RPC only queues the merge on the data plane; the actual work + # runs asynchronously afterward. Poll bdev_lvol_s3_merge_stat (keyed by + # the same s3_id/old_s3_id pair, since a merge task has no lvol) instead + # of assuming success. + try: + stat = rpc_client.bdev_lvol_s3_merge_stat(keep_backup.s3_id, old_backup.s3_id) + except RPCException: + task.retry += 1 + task.status = JobSchedule.STATUS_SUSPENDED + task.write_to_db(db.kv_store) + return - old_backup.status = Backup.STATUS_MERGED - old_backup.write_to_db() + if not stat or not isinstance(stat, dict): + task.retry += 1 + task.status = JobSchedule.STATUS_SUSPENDED + task.write_to_db(db.kv_store) + return - task.function_result = "Merge completed" - task.status = JobSchedule.STATUS_DONE - task.write_to_db(db.kv_store) - logger.info(f"Merge completed: {old_backup_id} merged into {keep_backup_id}") + state = stat.get("transfer_state", "") + if state == "Done": + # Finalize: update the chain links, retire the old backup, and mark + # the task done. + keep_backup.prev_backup_id = old_backup.prev_backup_id + keep_backup.status = Backup.STATUS_COMPLETED + keep_backup.write_to_db() + + old_backup.status = Backup.STATUS_MERGED + old_backup.write_to_db() + + task.function_result = "Merge completed" + task.status = JobSchedule.STATUS_DONE + task.write_to_db(db.kv_store) + logger.info(f"Merge completed: {old_backup_id} merged into {keep_backup_id}") + elif state == "Failed": + # Terminal, not retried: XFER_STATE_FAILED can be reached after the + # data plane has already started deleting old_backup's S3 objects + # (its DELETE state runs near the end of the merge), so retrying the + # identical merge isn't known to be safe from here. + if old_backup.status == Backup.STATUS_MERGING: + old_backup.status = Backup.STATUS_COMPLETED + old_backup.write_to_db() + task.function_result = "Merge failed on data plane" + task.status = JobSchedule.STATUS_DONE + task.write_to_db(db.kv_store) + elif state == "No process": + # Never started, or its result was already swept — re-issue. + task.function_params["merge_started"] = False + task.retry += 1 + task.status = JobSchedule.STATUS_SUSPENDED + task.write_to_db(db.kv_store) + else: + # "In progress" — still running, retry later + task.status = JobSchedule.STATUS_SUSPENDED + task.write_to_db(db.kv_store) def _terminate_task(task, reason): diff --git a/tests/integration/test_backup.py b/tests/integration/test_backup.py index cb733e2307..fc5bdb3f76 100644 --- a/tests/integration/test_backup.py +++ b/tests/integration/test_backup.py @@ -1125,6 +1125,12 @@ def test_bdev_lvol_s3_merge_exists(self): from simplyblock_core.rpc_client import RPCClient self.assertTrue(hasattr(RPCClient, 'bdev_lvol_s3_merge')) + def test_bdev_lvol_s3_merge_stat_exists(self): + """bdev_lvol_s3_merge_stat is used to poll merge progress -- a merge + task has no lvol, so it can't use bdev_lvol_transfer_stat.""" + from simplyblock_core.rpc_client import RPCClient + self.assertTrue(hasattr(RPCClient, 'bdev_lvol_s3_merge_stat')) + def test_bdev_lvol_s3_recovery_exists(self): from simplyblock_core.rpc_client import RPCClient self.assertTrue(hasattr(RPCClient, 'bdev_lvol_s3_recovery')) diff --git a/tests/unit/rpc/test_client.py b/tests/unit/rpc/test_client.py index 1a146cfacb..34f9474183 100644 --- a/tests/unit/rpc/test_client.py +++ b/tests/unit/rpc/test_client.py @@ -89,5 +89,55 @@ def test_subsystem_get_other_rpc_error_propagates(self, mock_req): client.subsystem_get("nqn.b") +class TestBdevLvolS3Merge(unittest.TestCase): + + @patch.object(RPCClient, "_request3") + def test_merge_eexist_treated_as_success_by_default(self, mock_req): + # A prior call whose RPC connection dropped before the response + # arrived can leave a matching merge already queued on the data + # plane; retrying that call must not be treated as a hard failure. + mock_req.side_effect = RPCRemoteError("The same transfer task already exists.", code=-errno.EEXIST) + client = _make_client() + + self.assertTrue(client.bdev_lvol_s3_merge(1, 2, cluster_batch=16)) + + @patch.object(RPCClient, "_request3") + def test_merge_eexist_propagates_when_disallowed(self, mock_req): + mock_req.side_effect = RPCRemoteError("The same transfer task already exists.", code=-errno.EEXIST) + client = _make_client() + + with self.assertRaises(RPCRemoteError): + client.bdev_lvol_s3_merge(1, 2, cluster_batch=16, allow_exist=False) + + @patch.object(RPCClient, "_request3") + def test_merge_other_rpc_error_propagates_regardless_of_allow_exist(self, mock_req): + mock_req.side_effect = RPCRemoteError("Cannot find S3 transfer device.", code=-errno.EINVAL) + client = _make_client() + + with self.assertRaises(RPCRemoteError): + client.bdev_lvol_s3_merge(1, 2, cluster_batch=16) + + @patch.object(RPCClient, "_request3") + def test_merge_success_passes_through(self, mock_req): + mock_req.return_value = True + client = _make_client() + + self.assertTrue(client.bdev_lvol_s3_merge(1, 2, cluster_batch=16, lvs_name="lvs0")) + mock_req.assert_called_once_with("bdev_lvol_s3_merge", s3_id=1, old_s3_id=2, cluster_batch=16, lvs_name="lvs0") + + +class TestBdevLvolS3MergeStat(unittest.TestCase): + + @patch.object(RPCClient, "_request") + def test_merge_stat_calls_request_with_ids(self, mock_req): + mock_req.return_value = {"transfer_state": "In progress"} + client = _make_client() + + result = client.bdev_lvol_s3_merge_stat(1, 2) + + self.assertEqual(result["transfer_state"], "In progress") + mock_req.assert_called_once_with("bdev_lvol_s3_merge_stat", {"s3_id": 1, "old_s3_id": 2}) + + if __name__ == "__main__": unittest.main() diff --git a/tests/unit/tasks/test_retry_ceiling.py b/tests/unit/tasks/test_retry_ceiling.py index 0b4a429b70..8fdb8325d1 100644 --- a/tests/unit/tasks/test_retry_ceiling.py +++ b/tests/unit/tasks/test_retry_ceiling.py @@ -148,6 +148,127 @@ def test_no_process_poll_counts_as_a_retry(runner, monkeypatch): rpc.bdev_lvol_s3_backup.assert_not_called() # this poll only resets state +# -------------------------------------------------------------------------- +# Backup runner: _run_merge polls bdev_lvol_s3_merge_stat instead of +# assuming success. Merge has no lvol, so it can't use bdev_lvol_transfer_stat +# like backup/restore above -- these cover its own stat RPC's four states. +# -------------------------------------------------------------------------- + +def _merge_task(retry=0, max_retry=10, merge_started=True): + task = JobSchedule() + task.uuid = "task-1" + task.function_name = JobSchedule.FN_BACKUP_MERGE + task.function_params = { + "keep_backup_id": "keep-1", + "old_backup_id": "old-1", + "merge_started": merge_started, + } + task.retry = retry + task.max_retry = max_retry + task.canceled = False + task.date = int(__import__("time").time()) + return task + + +def _merge_backups(old_status=Backup.STATUS_MERGING): + keep_backup = Backup() + keep_backup.uuid = "keep-1" + keep_backup.s3_id = 10 + keep_backup.node_id = "node-1" + keep_backup.prev_backup_id = "old-1" + + old_backup = Backup() + old_backup.uuid = "old-1" + old_backup.s3_id = 20 + old_backup.prev_backup_id = "older-0" + old_backup.status = old_status + + return keep_backup, old_backup + + +def _merge_db(keep_backup, old_backup, snode, rpc): + snode.rpc_client.return_value = rpc + fake_db = MagicMock() + fake_db.get_backup_by_id.side_effect = lambda bid: { + "keep-1": keep_backup, "old-1": old_backup, + }[bid] + fake_db.get_storage_node_by_id.return_value = snode + return fake_db + + +def _online_snode(): + snode = MagicMock() + snode.status = StorageNode.STATUS_ONLINE + return snode + + +def test_merge_done_finalizes(runner, monkeypatch): + keep_backup, old_backup = _merge_backups() + rpc = MagicMock() + rpc.bdev_lvol_s3_merge_stat.return_value = {"transfer_state": "Done"} + monkeypatch.setattr(runner, "db", _merge_db(keep_backup, old_backup, _online_snode(), rpc)) + + task = _merge_task() + runner._run_merge(task) + + rpc.bdev_lvol_s3_merge_stat.assert_called_once_with(10, 20) + assert keep_backup.status == Backup.STATUS_COMPLETED + assert keep_backup.prev_backup_id == "older-0" + assert old_backup.status == Backup.STATUS_MERGED + assert task.status == JobSchedule.STATUS_DONE + + +def test_merge_failed_reverts_old_backup_and_terminates(runner, monkeypatch): + """A 'Failed' merge must not be silently recorded as a success (the bug + this whole runner rewrite exists to fix), and must not be blindly + retried -- XFER_STATE_FAILED can be reached after the data plane has + already started deleting old_backup's S3 objects.""" + keep_backup, old_backup = _merge_backups(old_status=Backup.STATUS_MERGING) + rpc = MagicMock() + rpc.bdev_lvol_s3_merge_stat.return_value = {"transfer_state": "Failed"} + monkeypatch.setattr(runner, "db", _merge_db(keep_backup, old_backup, _online_snode(), rpc)) + + task = _merge_task(retry=0) + runner._run_merge(task) + + assert old_backup.status == Backup.STATUS_COMPLETED, "must revert MERGING, not leave it stuck" + assert keep_backup.status != Backup.STATUS_COMPLETED, "keep_backup must not be finalized on failure" + assert task.status == JobSchedule.STATUS_DONE + assert task.retry == 0, "terminal failure is not a retry" + + +def test_merge_in_progress_stays_suspended(runner, monkeypatch): + keep_backup, old_backup = _merge_backups() + rpc = MagicMock() + rpc.bdev_lvol_s3_merge_stat.return_value = {"transfer_state": "In progress"} + monkeypatch.setattr(runner, "db", _merge_db(keep_backup, old_backup, _online_snode(), rpc)) + + task = _merge_task() + runner._run_merge(task) + + assert task.status == JobSchedule.STATUS_SUSPENDED + assert task.function_params["merge_started"] is True + assert old_backup.status == Backup.STATUS_MERGING, "unchanged while still running" + + +def test_merge_no_process_resets_and_counts_as_a_retry(runner, monkeypatch): + """Mirrors test_no_process_poll_counts_as_a_retry above: the re-issue + branch must advance task.retry, otherwise the ceiling never binds here + and a merge that keeps reporting 'No process' loops forever.""" + keep_backup, old_backup = _merge_backups() + rpc = MagicMock() + rpc.bdev_lvol_s3_merge_stat.return_value = {"transfer_state": "No process"} + monkeypatch.setattr(runner, "db", _merge_db(keep_backup, old_backup, _online_snode(), rpc)) + + task = _merge_task(retry=3) + runner._run_merge(task) + + assert task.retry == 4, "No-process re-issue must count toward the ceiling" + assert task.function_params["merge_started"] is False + assert task.status == JobSchedule.STATUS_SUSPENDED + rpc.bdev_lvol_s3_merge.assert_not_called() # this poll only resets state + + # -------------------------------------------------------------------------- # Drive-main harness: exercise each runner's real retry loop end-to-end. # --------------------------------------------------------------------------