From 5e57b45e37e502c5827258721fc51f2c46726a35 Mon Sep 17 00:00:00 2001 From: ebrahim savari <106956085+EbiRider@users.noreply.github.com> Date: Fri, 31 Jul 2026 16:54:08 +0330 Subject: [PATCH 1/5] fix: source node going offline suspends the migration instead of transiting to cleanup target --- .../services/tasks_runner_lvol_migration.py | 22 ++++++++++++++++++- 1 file changed, 21 insertions(+), 1 deletion(-) diff --git a/simplyblock_core/services/tasks_runner_lvol_migration.py b/simplyblock_core/services/tasks_runner_lvol_migration.py index e9d8024eb..cc02b11fd 100644 --- a/simplyblock_core/services/tasks_runner_lvol_migration.py +++ b/simplyblock_core/services/tasks_runner_lvol_migration.py @@ -2774,13 +2774,13 @@ def _budget_suspend(task, migration, migration_id, error_msg): """Charge retry budget and suspend; redirect to cleanup_target when exhausted.""" migration.retry_count += 1 migration.error_message = error_msg - task.retry += 1 task.function_result = error_msg if migration.retry_count >= migration.max_retries: logger.error( f"Migration {migration_id} exceeded max retries " f"({migration.max_retries}); entering cleanup_target" ) + task.retry += 1 migration.phase = LVolMigration.PHASE_CLEANUP_TARGET migration.current_job_id = "" migration.write_to_db(db.kv_store) @@ -2859,6 +2859,26 @@ def task_runner(task): LVolMigration.PHASE_CLEANUP_TARGET, LVolMigration.PHASE_CLEANUP_SOURCE) if not _is_cleanup_phase: if src_node.status not in (StorageNode.STATUS_ONLINE, StorageNode.STATUS_SUSPENDED): + if migration.phase in (LVolMigration.PHASE_SNAP_COPY, + LVolMigration.PHASE_LVOL_MIGRATE): + # Source offline during data-transfer: migration cannot continue. + # Enter cleanup_target immediately — same fast-path as when the + # target goes offline (see below). Burning retries here is wrong + # because snap_copy / lvol_migrate require the source to be alive; + # waiting for it to come back doesn't help. + logger.warning( + f"Migration {migration_id}: source node offline " + f"(status={src_node.status}) during {migration.phase}; " + f"entering cleanup_target immediately") + migration.error_message = ( + f"source node offline (status={src_node.status}); migration failed") + migration.phase = LVolMigration.PHASE_CLEANUP_TARGET + migration.current_job_id = "" + migration.write_to_db(db.kv_store) + task.function_result = migration.error_message + task.write_to_db(db.kv_store) + migration_events.migration_phase_changed(migration) + return False return _budget_suspend( task, migration, migration_id, f"source node not online (status={src_node.status})") From 07cf60fe1ce5a6cf9bc4115ac0561abd2a3e6e02 Mon Sep 17 00:00:00 2001 From: ebrahim savari <106956085+EbiRider@users.noreply.github.com> Date: Fri, 31 Jul 2026 17:03:18 +0330 Subject: [PATCH 2/5] fix: retries jumping up by 2 per try on lvol migration --- .../services/tasks_runner_lvol_migration.py | 25 +++---------------- 1 file changed, 4 insertions(+), 21 deletions(-) diff --git a/simplyblock_core/services/tasks_runner_lvol_migration.py b/simplyblock_core/services/tasks_runner_lvol_migration.py index cc02b11fd..b1be225cd 100644 --- a/simplyblock_core/services/tasks_runner_lvol_migration.py +++ b/simplyblock_core/services/tasks_runner_lvol_migration.py @@ -2781,10 +2781,13 @@ def _budget_suspend(task, migration, migration_id, error_msg): f"({migration.max_retries}); entering cleanup_target" ) task.retry += 1 + task.status = JobSchedule.STATUS_SUSPENDED migration.phase = LVolMigration.PHASE_CLEANUP_TARGET migration.current_job_id = "" - migration.write_to_db(db.kv_store) + # Write task first so the runner can always re-enter and attempt cleanup, + # even if the migration write below fails. task.write_to_db(db.kv_store) + migration.write_to_db(db.kv_store) migration_events.migration_phase_changed(migration) return False return _suspend_task(task, migration, error_msg) @@ -2859,26 +2862,6 @@ def task_runner(task): LVolMigration.PHASE_CLEANUP_TARGET, LVolMigration.PHASE_CLEANUP_SOURCE) if not _is_cleanup_phase: if src_node.status not in (StorageNode.STATUS_ONLINE, StorageNode.STATUS_SUSPENDED): - if migration.phase in (LVolMigration.PHASE_SNAP_COPY, - LVolMigration.PHASE_LVOL_MIGRATE): - # Source offline during data-transfer: migration cannot continue. - # Enter cleanup_target immediately — same fast-path as when the - # target goes offline (see below). Burning retries here is wrong - # because snap_copy / lvol_migrate require the source to be alive; - # waiting for it to come back doesn't help. - logger.warning( - f"Migration {migration_id}: source node offline " - f"(status={src_node.status}) during {migration.phase}; " - f"entering cleanup_target immediately") - migration.error_message = ( - f"source node offline (status={src_node.status}); migration failed") - migration.phase = LVolMigration.PHASE_CLEANUP_TARGET - migration.current_job_id = "" - migration.write_to_db(db.kv_store) - task.function_result = migration.error_message - task.write_to_db(db.kv_store) - migration_events.migration_phase_changed(migration) - return False return _budget_suspend( task, migration, migration_id, f"source node not online (status={src_node.status})") From 924a87bb72685ebd9b0c267e0903f9629ed8651f Mon Sep 17 00:00:00 2001 From: ebrahim savari <106956085+EbiRider@users.noreply.github.com> Date: Fri, 31 Jul 2026 18:56:21 +0330 Subject: [PATCH 3/5] fix: reaching 10 retries kills the process and making it suspend forever --- .../services/tasks_runner_batch_migration.py | 72 +++++++++++-------- .../services/tasks_runner_lvol_migration.py | 2 +- 2 files changed, 43 insertions(+), 31 deletions(-) diff --git a/simplyblock_core/services/tasks_runner_batch_migration.py b/simplyblock_core/services/tasks_runner_batch_migration.py index 8f41c80e7..88bcca9ce 100644 --- a/simplyblock_core/services/tasks_runner_batch_migration.py +++ b/simplyblock_core/services/tasks_runner_batch_migration.py @@ -710,6 +710,27 @@ def _try_delete(rpc, node_id, label): _try_delete(tert_rpc, tert_node.get_id(), "tertiary") +# --------------------------------------------------------------------------- +# Retry-budget helper +# --------------------------------------------------------------------------- + +def _batch_budget_suspend(task, group, group_id, error_msg): + """Charge retry budget and suspend; redirect to cleanup_target when exhausted.""" + task.retry += 1 + task.function_result = error_msg + if task.max_retry >= 0 and task.retry >= task.max_retry: + logger.error( + f"Group {group_id[:8]}: max retry reached " + f"({task.retry}/{task.max_retry}); entering cleanup_target: {error_msg}" + ) + group.phase = LVolMigrationGroup.PHASE_CLEANUP_TARGET + group.error_message = error_msg + group.write_to_db(db.kv_store) + task.status = JobSchedule.STATUS_SUSPENDED + task.write_to_db(db.kv_store) + return False + + # --------------------------------------------------------------------------- # Main task runner # --------------------------------------------------------------------------- @@ -762,12 +783,19 @@ def task_runner(task): task.write_to_db(db.kv_store) return False + phase = group.phase + _is_cleanup_phase = phase in ( + LVolMigrationGroup.PHASE_CLEANUP_TARGET, + LVolMigrationGroup.PHASE_CLEANUP_SOURCE, + ) + cluster = db.get_cluster_by_id(group.cluster_id) if cluster.status not in (Cluster.STATUS_ACTIVE, Cluster.STATUS_DEGRADED): - task.function_result = f"cluster not active (status={cluster.status})" - task.status = JobSchedule.STATUS_SUSPENDED - task.write_to_db(db.kv_store) - return False + if not _is_cleanup_phase: + task.function_result = f"cluster not active (status={cluster.status})" + task.status = JobSchedule.STATUS_SUSPENDED + task.write_to_db(db.kv_store) + return False if task.status in (JobSchedule.STATUS_NEW, JobSchedule.STATUS_SUSPENDED): task.status = JobSchedule.STATUS_RUNNING @@ -786,8 +814,6 @@ def task_runner(task): task.write_to_db(db.kv_store) return False - phase = group.phase - # If source or target went offline during data-transfer phases, enter # CLEANUP_TARGET immediately — same fast-path as the single-lvol runner. _data_transfer_phases = (LVolMigrationGroup.PHASE_SNAP_COPY, @@ -811,29 +837,17 @@ def task_runner(task): logger.warning( f"Group {group_id[:8]}: source node unavailable " f"(status={fresh_src.status}) during {phase}; suspending") - task.retry += 1 - if task.max_retry >= 0 and task.retry >= task.max_retry: - task.function_result = ( - f"max retry reached ({task.retry}/{task.max_retry}): " - f"source node unavailable (status={fresh_src.status})") - task.status = JobSchedule.STATUS_DONE - task.write_to_db(db.kv_store) - return True - task.function_result = f"source node unavailable (status={fresh_src.status})" - task.status = JobSchedule.STATUS_SUSPENDED - task.write_to_db(db.kv_store) - return False + return _batch_budget_suspend( + task, group, group_id, + f"source node unavailable (status={fresh_src.status})") try: # ── PHASE_SNAP_COPY: wait for all workers, then reconstruct tree ───────── if phase == LVolMigrationGroup.PHASE_SNAP_COPY: done, err = _handle_snap_copy_barrier(group, member_migrations, tgt_node, tgt_rpc) if err: - task.function_result = err - task.status = JobSchedule.STATUS_SUSPENDED - task.write_to_db(db.kv_store) logger.error(f"Group {group_id[:8]}: snap_copy barrier error: {err}") - return False + return _batch_budget_suspend(task, group, group_id, err) if not done: task.write_to_db(db.kv_store) return False @@ -850,11 +864,8 @@ def task_runner(task): group, member_migrations, src_node, tgt_node, src_rpc, tgt_rpc) if err: - task.function_result = err - task.status = JobSchedule.STATUS_SUSPENDED - task.write_to_db(db.kv_store) logger.error(f"Group {group_id[:8]}: intermediate barrier error: {err}") - return False + return _batch_budget_suspend(task, group, group_id, err) if batch_ok is None: # Still waiting for workers. @@ -883,10 +894,11 @@ def task_runner(task): group.phase = LVolMigrationGroup.PHASE_CLEANUP_TARGET group.error_message = f"target node offline during {phase}: {exc}" group.write_to_db(db.kv_store) - task.function_result = str(exc) - task.status = JobSchedule.STATUS_SUSPENDED - task.write_to_db(db.kv_store) - return False + task.function_result = str(exc) + task.status = JobSchedule.STATUS_SUSPENDED + task.write_to_db(db.kv_store) + return False + return _batch_budget_suspend(task, group, group_id, f"RPC error in phase {phase}: {exc}") # ── PHASE_CLEANUP_SOURCE: wait for workers, then delete source subsystem ─── if phase == LVolMigrationGroup.PHASE_CLEANUP_SOURCE: diff --git a/simplyblock_core/services/tasks_runner_lvol_migration.py b/simplyblock_core/services/tasks_runner_lvol_migration.py index b1be225cd..599eace1a 100644 --- a/simplyblock_core/services/tasks_runner_lvol_migration.py +++ b/simplyblock_core/services/tasks_runner_lvol_migration.py @@ -3017,7 +3017,6 @@ def task_runner(task): # Real operation failure – increment retry counter. migration.retry_count += 1 migration.error_message = error - task.retry += 1 task.function_result = error if migration.retry_count >= migration.max_retries: @@ -3033,6 +3032,7 @@ def task_runner(task): f"Migration {migration_id} exceeded max retries " f"({migration.max_retries}); entering cleanup_target" ) + task.retry += 1 migration.phase = LVolMigration.PHASE_CLEANUP_TARGET migration.current_job_id = "" migration.write_to_db(db.kv_store) From bcf59288550a45ab232bbbb57af29404ca4e031e Mon Sep 17 00:00:00 2001 From: ebrahim savari <106956085+EbiRider@users.noreply.github.com> Date: Sat, 1 Aug 2026 06:51:50 +0330 Subject: [PATCH 4/5] TEMP: trying a workaround to avoid process from being killed by backup monitor --- simplyblock_core/constants.py | 2 +- .../services/tasks_runner_lvol_migration.py | 115 ++++++++---------- 2 files changed, 51 insertions(+), 66 deletions(-) diff --git a/simplyblock_core/constants.py b/simplyblock_core/constants.py index 05df22882..b0d5b025e 100644 --- a/simplyblock_core/constants.py +++ b/simplyblock_core/constants.py @@ -500,7 +500,7 @@ def get_config_var(name, default=None): MIG_JOB_SIZE = 64 # Live volume migration constants -LVOL_MIG_MAX_RETRIES = 5 # max retry attempts before aborting +LVOL_MIG_MAX_RETRIES = 10 # total retry budget; cleanup_target fires at max_retries // 2 LVOL_MIG_DEADLINE_SEC = 3600 # 1-hour deadline (0 = no deadline) LVOL_MIG_MAX_INTERMEDIATE_SNAPS = 3 # max recursive "shrink" snapshot rounds LVOL_MIG_INTERMEDIATE_SNAP_THRESHOLD_BYTES = 500 * 1024 * 1024 # 500 MiB — skip if delta is smaller diff --git a/simplyblock_core/services/tasks_runner_lvol_migration.py b/simplyblock_core/services/tasks_runner_lvol_migration.py index 599eace1a..0dd1293d6 100644 --- a/simplyblock_core/services/tasks_runner_lvol_migration.py +++ b/simplyblock_core/services/tasks_runner_lvol_migration.py @@ -106,13 +106,6 @@ logger = utils.get_logger(__name__) db = db_mod.DBController() -# Sentinel used as the ``error`` return value when a phase handler wants to -# suspend the task WITHOUT incrementing the retry counter. This is distinct -# from a real operation failure: it signals a *transient external condition* -# (e.g. secondary node in unexpected state) that the runner should wait for, -# not charge against the retry budget. -_WAIT = object() - # Busy-poll settings for intermediate ("shrink") snapshot transfers. # Intermediate snapshots represent a small dirty delta so they should complete # quickly; we spin at _INTERMEDIATE_POLL_INTERVAL_S rather than waiting for @@ -730,9 +723,8 @@ def _ensure_target_nvmf_state(migration, lvol, src_node, tgt_node, src_rpc, tgt_ subsystem created by this migration — if that subsystem is entirely gone, something outside migration's control needs to fix it, so this only waits. - Returns None if everything is fine (or was successfully repaired), or the - _WAIT sentinel if the caller should suspend without charging the - migration's retry budget. + Returns None if everything is fine (or was successfully repaired), or an + error string describing the transient failure. """ nqn = lvol.nqn try: @@ -740,7 +732,7 @@ def _ensure_target_nvmf_state(migration, lvol, src_node, tgt_node, src_rpc, tgt_ except Exception as e: logger.warning(f"_ensure_target_nvmf_state: could not build target paths: {e}") migration.error_message = f"target topology unavailable: {e}" - return _WAIT + return migration.error_message owned_node_ids = set(migration.target_subsystem_node_ids or []) short_bdev = _lvol_tgt_bdev_name(lvol.lvol_bdev) @@ -760,7 +752,7 @@ def _ensure_target_nvmf_state(migration, lvol, src_node, tgt_node, src_rpc, tgt_ f"_ensure_target_nvmf_state: subsystem_get failed on " f"{label} {node_id[:8]}: {e}") migration.error_message = f"could not query subsystem on {label} target node: {e}" - return _WAIT + return migration.error_message if not sub: if not owns_subsystem: @@ -769,7 +761,7 @@ def _ensure_target_nvmf_state(migration, lvol, src_node, tgt_node, src_rpc, tgt_ f"node recovery") logger.warning(f"_ensure_target_nvmf_state: {msg}") migration.error_message = msg - return _WAIT + return migration.error_message logger.warning( f"_ensure_target_nvmf_state: subsystem {nqn} missing on " @@ -800,7 +792,7 @@ def _ensure_target_nvmf_state(migration, lvol, src_node, tgt_node, src_rpc, tgt_ f"_ensure_target_nvmf_state: recreate failed on {label} " f"{node_id[:8]}: {e}") migration.error_message = f"failed to recreate target subsystem on {label} node: {e}" - return _WAIT + return migration.error_message continue # Subsystem present — verify our listener survived. @@ -811,7 +803,7 @@ def _ensure_target_nvmf_state(migration, lvol, src_node, tgt_node, src_rpc, tgt_ f"_ensure_target_nvmf_state: listeners_list failed on " f"{label} {node_id[:8]}: {e}") migration.error_message = f"could not query listeners on {label} target node: {e}" - return _WAIT + return migration.error_message listener_addrs = { (ls.get('address', {}).get('traddr'), str(ls.get('address', {}).get('trsvcid'))) @@ -831,7 +823,7 @@ def _ensure_target_nvmf_state(migration, lvol, src_node, tgt_node, src_rpc, tgt_ f"_ensure_target_nvmf_state: listener recreate failed on " f"{label} {node_id[:8]}: {e}") migration.error_message = f"failed to recreate target listener on {label} node: {e}" - return _WAIT + return migration.error_message # Namespace check — only on nodes whose namespace lifecycle we own; # overlap nodes legitimately still point at the SRC bdev pre-cutover @@ -854,7 +846,7 @@ def _ensure_target_nvmf_state(migration, lvol, src_node, tgt_node, src_rpc, tgt_ f"_ensure_target_nvmf_state: namespace re-add errored on " f"{label} {node_id[:8]}: {e}") migration.error_message = f"failed to re-add namespace on {label} node: {e}" - return _WAIT + return migration.error_message return None @@ -1318,7 +1310,7 @@ def _handle_snap_copy(migration, src_node, tgt_node, src_rpc, tgt_rpc): if sec_err: migration.error_message = sec_err migration.write_to_db(db.kv_store) - return False, True, _WAIT + return False, True, sec_err if tgt_sec: sec_rpc = _make_rpc(tgt_sec) elif snap.lvol.ha_type == "ha3": @@ -1326,14 +1318,14 @@ def _handle_snap_copy(migration, src_node, tgt_node, src_rpc, tgt_rpc): if sec_err: migration.error_message = sec_err migration.write_to_db(db.kv_store) - return False, True, _WAIT + return False, True, sec_err if tgt_sec: sec_rpc = _make_rpc(tgt_sec) tgt_ter, ter_err = _get_target_tertiary_node(tgt_node, src_node.get_id()) if ter_err: migration.error_message = ter_err migration.write_to_db(db.kv_store) - return False, True, _WAIT + return False, True, ter_err if tgt_ter: ter_rpc = _make_rpc(tgt_ter) break # one check is enough @@ -1477,11 +1469,6 @@ def _handle_snap_copy(migration, src_node, tgt_node, src_rpc, tgt_rpc): tgt_ter=tgt_ter, ter_rpc=ter_rpc) if not ok: migration.transfer_context = {} - if err is _WAIT: - migration.error_message = ( - f"Secondary node not ready during post-process of {snap_uuid}") - migration.write_to_db(db.kv_store) - return False, True, _WAIT migration.write_to_db(db.kv_store) return False, True, err @@ -1555,7 +1542,7 @@ def _handle_snap_copy(migration, src_node, tgt_node, src_rpc, tgt_rpc): if sec_err: migration.error_message = sec_err migration.write_to_db(db.kv_store) - return False, True, _WAIT + return False, True, sec_err if tgt_sec: sec_rpc = _make_rpc(tgt_sec) if snap.lvol.ha_type == "ha3": @@ -1563,7 +1550,7 @@ def _handle_snap_copy(migration, src_node, tgt_node, src_rpc, tgt_rpc): if ter_err: migration.error_message = ter_err migration.write_to_db(db.kv_store) - return False, True, _WAIT + return False, True, ter_err if tgt_ter: ter_rpc = _make_rpc(tgt_ter) @@ -1649,11 +1636,6 @@ def _handle_snap_copy(migration, src_node, tgt_node, src_rpc, tgt_rpc): tgt_sec=tgt_sec, sec_rpc=sec_rpc, tgt_ter=tgt_ter, ter_rpc=ter_rpc) if not ok: - if err is _WAIT: - migration.error_message = ( - f"Secondary node not ready after intermediate snap {snap_uuid}") - migration.write_to_db(db.kv_store) - return False, True, _WAIT return False, True, err migration.next_snap_index = len(plan) @@ -1865,7 +1847,7 @@ def _revert_src_replicas(reason): if sec_err: migration.error_message = sec_err migration.write_to_db(db.kv_store) - return False, True, _WAIT + return False, True, sec_err # --- Start the final migration --- @@ -2771,14 +2753,18 @@ def _handle_cleanup_target(migration, tgt_node, tgt_rpc, src_rpc=None): # --------------------------------------------------------------------------- def _budget_suspend(task, migration, migration_id, error_msg): - """Charge retry budget and suspend; redirect to cleanup_target when exhausted.""" + """Charge retry budget and suspend; redirect to cleanup_target when exhausted. + + Ceiling fires at max_retries // 2 so cleanup_target has the remaining + budget before the backup runner kills at task.max_retry == max_retries. + """ migration.retry_count += 1 migration.error_message = error_msg task.function_result = error_msg - if migration.retry_count >= migration.max_retries: + if migration.retry_count >= migration.max_retries // 2: logger.error( - f"Migration {migration_id} exceeded max retries " - f"({migration.max_retries}); entering cleanup_target" + f"Migration {migration_id} exceeded retry ceiling " + f"({migration.retry_count}/{migration.max_retries // 2}); entering cleanup_target" ) task.retry += 1 task.status = JobSchedule.STATUS_SUSPENDED @@ -2902,14 +2888,16 @@ def task_runner(task): if cluster.status not in (Cluster.STATUS_ACTIVE, Cluster.STATUS_DEGRADED): if not _is_cleanup_phase: return _suspend_task( - task, migration, f"cluster not active (status={cluster.status})") + task, migration, f"cluster not active (status={cluster.status})", + charge_retry=False) # Expansion-first ordering: defer while a cluster expansion is open — # even between the expand task's retries, when the cluster status is # momentarily ACTIVE (see tasks_controller.defer_task_for_expansion). if tasks_controller.get_active_cluster_expand_task(task.cluster_id): return _suspend_task( - task, migration, "cluster expansion in progress, deferring") + task, migration, "cluster expansion in progress, deferring", + charge_retry=False) # --- Transition NEW/SUSPENDED → RUNNING --- if task.status in (JobSchedule.STATUS_NEW, JobSchedule.STATUS_SUSPENDED): @@ -2942,9 +2930,10 @@ def task_runner(task): except KeyError: return _budget_suspend(task, migration, migration.get_id(), f"LVol {migration.lvol_id} not found") - nvmf_wait = _ensure_target_nvmf_state(migration, lvol, src_node, tgt_node, src_rpc, tgt_rpc) - if nvmf_wait is _WAIT: - return _suspend_task(task, migration, migration.error_message or "waiting for target node to recover") + nvmf_err = _ensure_target_nvmf_state(migration, lvol, src_node, tgt_node, src_rpc, tgt_rpc) + if nvmf_err: + return _budget_suspend(task, migration, migration.get_id(), + migration.error_message or nvmf_err) try: if migration.migration_group_id: @@ -3008,18 +2997,13 @@ def task_runner(task): error = str(exc) # --- Handle error / suspend --- - if error is _WAIT: - # Transient external condition (e.g. secondary node not ready). - # Suspend without charging against the retry budget. - return _suspend_task(task, migration, migration.error_message or "waiting") - if error: - # Real operation failure – increment retry counter. + # Operation failure – increment retry counter. migration.retry_count += 1 migration.error_message = error task.function_result = error - if migration.retry_count >= migration.max_retries: + if migration.retry_count >= migration.max_retries // 2: if phase not in (LVolMigration.PHASE_SNAP_COPY, LVolMigration.PHASE_LVOL_MIGRATE): # Already past the migration — never redirect to CLEANUP_TARGET. @@ -3029,8 +3013,8 @@ def task_runner(task): f"{phase}; suspending for operator review (not entering cleanup_target)") return _suspend_task(task, migration, error) logger.error( - f"Migration {migration_id} exceeded max retries " - f"({migration.max_retries}); entering cleanup_target" + f"Migration {migration_id} exceeded retry ceiling " + f"({migration.retry_count}/{migration.max_retries // 2}); entering cleanup_target" ) task.retry += 1 migration.phase = LVolMigration.PHASE_CLEANUP_TARGET @@ -3176,11 +3160,11 @@ def _handle_group_snap_copy(migration, src_node, tgt_node, src_rpc, tgt_rpc): if _g_sec_err: migration.error_message = _g_sec_err migration.write_to_db(db.kv_store) - return False, True, _WAIT + return False, True, _g_sec_err if _g_ter_err: migration.error_message = _g_ter_err migration.write_to_db(db.kv_store) - return False, True, _WAIT + return False, True, _g_ter_err _g_sec_rpc = _make_rpc(_g_tgt_sec) if _g_tgt_sec else None _g_ter_rpc = _make_rpc(_g_tgt_ter) if _g_tgt_ter else None @@ -3310,11 +3294,11 @@ def _handle_group_intermediate(migration, src_node, tgt_node, src_rpc, tgt_rpc): if _g_sec_err: migration.error_message = _g_sec_err migration.write_to_db(db.kv_store) - return False, True, _WAIT + return False, True, _g_sec_err if _g_ter_err: migration.error_message = _g_ter_err migration.write_to_db(db.kv_store) - return False, True, _WAIT + return False, True, _g_ter_err _g_sec_rpc = _make_rpc(_g_tgt_sec) if _g_tgt_sec else None _g_ter_rpc = _make_rpc(_g_tgt_ter) if _g_tgt_ter else None @@ -3414,9 +3398,9 @@ def _group_worker_phase_dispatch(task, migration, phase, src_node, tgt_node, src # Still transferring owned snaps. done, suspend, error = _handle_group_snap_copy( migration, src_node, tgt_node, src_rpc, tgt_rpc) - if error and error is not _WAIT: + if error: return _suspend_task(task, migration, error) - if error is _WAIT or suspend: + if suspend: return _suspend_task(task, migration, migration.error_message or "waiting") if done: # Signal snap_copy_done to group. @@ -3457,9 +3441,9 @@ def _group_worker_phase_dispatch(task, migration, phase, src_node, tgt_node, src if migration_id not in group.intermediates_done: done, suspend, error = _handle_group_intermediate( migration, src_node, tgt_node, src_rpc, tgt_rpc) - if error and error is not _WAIT: + if error: return _suspend_task(task, migration, error) - if error is _WAIT or suspend: + if suspend: return _suspend_task(task, migration, migration.error_message or "waiting") if done: group = db.get_migration_group_by_id(group_id) @@ -3506,9 +3490,9 @@ def _group_worker_phase_dispatch(task, migration, phase, src_node, tgt_node, src except RPCException as exc: return _suspend_task(task, migration, str(exc)) - if error and error is not _WAIT: + if error: return _suspend_task(task, migration, error) - if error is _WAIT or suspend: + if suspend: return _suspend_task(task, migration, migration.error_message or "waiting") if done: # Signal cleanup_source_done to group. @@ -3539,9 +3523,9 @@ def _group_worker_phase_dispatch(task, migration, phase, src_node, tgt_node, src except RPCException as exc: return _suspend_task(task, migration, str(exc)) - if error and error is not _WAIT: + if error: return _suspend_task(task, migration, error) - if error is _WAIT or suspend: + if suspend: return _suspend_task(task, migration, migration.error_message or "waiting") if done: migration.status = (LVolMigration.STATUS_CANCELLED if migration.canceled @@ -3570,10 +3554,11 @@ def _make_rpc(node): return node.rpc_client(timeout=5, retry=2) -def _suspend_task(task, migration, reason): +def _suspend_task(task, migration, reason, charge_retry=True): task.status = JobSchedule.STATUS_SUSPENDED task.function_result = reason - task.retry += 1 + if charge_retry: + task.retry += 1 task.write_to_db(db.kv_store) migration.status = LVolMigration.STATUS_SUSPENDED migration.error_message = reason From 5504903e77dd76046c51c8cd3f2345278e69a892 Mon Sep 17 00:00:00 2001 From: ebrahim savari <106956085+EbiRider@users.noreply.github.com> Date: Sun, 2 Aug 2026 22:32:30 +0330 Subject: [PATCH 5/5] setting max retires -1 for add_lvol_mig_task so the monitor don't kill this process before it's completion --- simplyblock_core/constants.py | 2 +- .../controllers/tasks_controller.py | 9 +++++++-- .../services/tasks_runner_lvol_migration.py | 18 +++++++----------- 3 files changed, 15 insertions(+), 14 deletions(-) diff --git a/simplyblock_core/constants.py b/simplyblock_core/constants.py index b0d5b025e..2c15eb849 100644 --- a/simplyblock_core/constants.py +++ b/simplyblock_core/constants.py @@ -500,7 +500,7 @@ def get_config_var(name, default=None): MIG_JOB_SIZE = 64 # Live volume migration constants -LVOL_MIG_MAX_RETRIES = 10 # total retry budget; cleanup_target fires at max_retries // 2 +LVOL_MIG_MAX_RETRIES = 5 # max retries before entering cleanup_target LVOL_MIG_DEADLINE_SEC = 3600 # 1-hour deadline (0 = no deadline) LVOL_MIG_MAX_INTERMEDIATE_SNAPS = 3 # max recursive "shrink" snapshot rounds LVOL_MIG_INTERMEDIATE_SNAP_THRESHOLD_BYTES = 500 * 1024 * 1024 # 500 MiB — skip if delta is smaller diff --git a/simplyblock_core/controllers/tasks_controller.py b/simplyblock_core/controllers/tasks_controller.py index 8dd2aedd5..a89f68123 100644 --- a/simplyblock_core/controllers/tasks_controller.py +++ b/simplyblock_core/controllers/tasks_controller.py @@ -872,13 +872,18 @@ def get_active_lvol_mig_task_on_node(cluster_id, node_id): def add_lvol_mig_task(migration): - """Create the JobSchedule task that drives a live volume migration.""" + """Create the JobSchedule task that drives a live volume migration. + + max_retry=-1 disables the backup runner's retry-count kill switch. + The migration runner has its own internal ceiling via migration.retry_count; + the backup runner's time-based timeout is the only external backstop. + """ return _add_task( JobSchedule.FN_LVOL_MIG, migration.cluster_id, migration.source_node_id, "", - max_retry=migration.max_retries, + max_retry=-1, function_params={ "migration_id": migration.uuid, "lvol_id": migration.lvol_id, diff --git a/simplyblock_core/services/tasks_runner_lvol_migration.py b/simplyblock_core/services/tasks_runner_lvol_migration.py index 0dd1293d6..a2b9846a9 100644 --- a/simplyblock_core/services/tasks_runner_lvol_migration.py +++ b/simplyblock_core/services/tasks_runner_lvol_migration.py @@ -2753,18 +2753,14 @@ def _handle_cleanup_target(migration, tgt_node, tgt_rpc, src_rpc=None): # --------------------------------------------------------------------------- def _budget_suspend(task, migration, migration_id, error_msg): - """Charge retry budget and suspend; redirect to cleanup_target when exhausted. - - Ceiling fires at max_retries // 2 so cleanup_target has the remaining - budget before the backup runner kills at task.max_retry == max_retries. - """ + """Charge retry budget and suspend; redirect to cleanup_target when exhausted.""" migration.retry_count += 1 migration.error_message = error_msg task.function_result = error_msg - if migration.retry_count >= migration.max_retries // 2: + if migration.retry_count >= migration.max_retries: logger.error( - f"Migration {migration_id} exceeded retry ceiling " - f"({migration.retry_count}/{migration.max_retries // 2}); entering cleanup_target" + f"Migration {migration_id} exceeded max retries " + f"({migration.max_retries}); entering cleanup_target" ) task.retry += 1 task.status = JobSchedule.STATUS_SUSPENDED @@ -3003,7 +2999,7 @@ def task_runner(task): migration.error_message = error task.function_result = error - if migration.retry_count >= migration.max_retries // 2: + if migration.retry_count >= migration.max_retries: if phase not in (LVolMigration.PHASE_SNAP_COPY, LVolMigration.PHASE_LVOL_MIGRATE): # Already past the migration — never redirect to CLEANUP_TARGET. @@ -3013,8 +3009,8 @@ def task_runner(task): f"{phase}; suspending for operator review (not entering cleanup_target)") return _suspend_task(task, migration, error) logger.error( - f"Migration {migration_id} exceeded retry ceiling " - f"({migration.retry_count}/{migration.max_retries // 2}); entering cleanup_target" + f"Migration {migration_id} exceeded max retries " + f"({migration.max_retries}); entering cleanup_target" ) task.retry += 1 migration.phase = LVolMigration.PHASE_CLEANUP_TARGET