Skip to content
Open
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
2 changes: 1 addition & 1 deletion simplyblock_core/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 = 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
Expand Down
9 changes: 7 additions & 2 deletions simplyblock_core/controllers/tasks_controller.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
72 changes: 42 additions & 30 deletions simplyblock_core/services/tasks_runner_batch_migration.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
# ---------------------------------------------------------------------------
Expand Down Expand Up @@ -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
Expand All @@ -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,
Expand All @@ -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
Expand All @@ -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.
Expand Down Expand Up @@ -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:
Expand Down
Loading
Loading