Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 30 additions & 7 deletions simplyblock_core/rpc_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -2034,22 +2034,45 @@ 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,
"cluster_batch": cluster_batch,
}
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.
Expand Down
65 changes: 52 additions & 13 deletions simplyblock_core/services/tasks_runner_backup.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
6 changes: 6 additions & 0 deletions tests/integration/test_backup.py
Original file line number Diff line number Diff line change
Expand Up @@ -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'))
Expand Down
50 changes: 50 additions & 0 deletions tests/unit/rpc/test_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
121 changes: 121 additions & 0 deletions tests/unit/tasks/test_retry_ceiling.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
# --------------------------------------------------------------------------
Expand Down
Loading