diff --git a/simplyblock_core/cluster_ops.py b/simplyblock_core/cluster_ops.py index 7e3b91376d..303f225c5e 100644 --- a/simplyblock_core/cluster_ops.py +++ b/simplyblock_core/cluster_ops.py @@ -2280,6 +2280,54 @@ def get_cluster(cl_id) -> dict: def update_cluster(cluster_id, mgmt_only=False, restart=False, spdk_image=None, mgmt_image=None, **kwargs) -> None: cluster = db_controller.get_cluster_by_id(cluster_id) # ensure exists + # stop jm replication on all leader nodes + for node in db_controller.get_storage_nodes_by_cluster_id(cluster_id): + if node.status in [StorageNode.STATUS_ONLINE, StorageNode.STATUS_DOWN, StorageNode.STATUS_SUSPENDED]: + ret = node.rpc_client().bdev_lvol_get_lvstores(node.lvstore) + if ret: + lvs_info = ret[0] + if "lvs leadership" in lvs_info and lvs_info['lvs leadership']: + ret, err = node.rpc_client().jc_suspend_compression(jm_vuid=node.jm_vuid, suspend=True) + if not ret: + logger.warning(f"Failed to stop JC compression on node: {node.get_id()}, JM:{node.jm_vuid}") + if node.lvstore_stack_secondary: + sec_node = db_controller.get_storage_node_by_id(node.lvstore_stack_secondary) + ret = node.rpc_client().bdev_lvol_get_lvstores(sec_node.lvstore) + if ret: + lvs_info = ret[0] + if "lvs leadership" in lvs_info and lvs_info['lvs leadership']: + ret, err = node.rpc_client().jc_suspend_compression(jm_vuid=sec_node.jm_vuid, suspend=True) + if not ret: + logger.warning(f"Failed to stop JC compression on node: {node.get_id()}, JM:{sec_node.jm_vuid}") + + jc_compression_is_active=False + # check for jm replication on all nodes () + for node in db_controller.get_storage_nodes_by_cluster_id(cluster_id): + if node.status in [StorageNode.STATUS_ONLINE, StorageNode.STATUS_DOWN, StorageNode.STATUS_SUSPENDED]: + try: + ret = node.rpc_client().bdev_lvol_get_lvstores(node.lvstore) + if ret: + lvs_info = ret[0] + if "lvs leadership" in lvs_info and lvs_info['lvs leadership']: + jc_compression_is_active = node.rpc_client().jc_compression_get_status(node.jm_vuid) + retries = 10 + while jc_compression_is_active: + if retries <= 0: + logger.warning("Timeout waiting for JC compression task to finish") + break + retries -= 1 + logger.info( + f"JC compression task found on node: {node.get_id()}, retrying in 60 seconds") + time.sleep(30) + jc_compression_is_active = node.rpc_client().jc_compression_get_status(node.jm_vuid) + except Exception as e: + logger.error(e) + return + + if jc_compression_is_active: + logger.error(f"JC compression task found.") + return + logger.info("Updating mgmt cluster") if cluster.mode == "docker": cluster_docker = utils.get_docker_client(cluster_id) diff --git a/simplyblock_core/services/tasks_runner_jc_comp.py b/simplyblock_core/services/tasks_runner_jc_comp.py index 764e1828cc..b31bda2d92 100644 --- a/simplyblock_core/services/tasks_runner_jc_comp.py +++ b/simplyblock_core/services/tasks_runner_jc_comp.py @@ -77,17 +77,35 @@ def main(): task.status = JobSchedule.STATUS_SUSPENDED task.write_to_db(db.kv_store) else: - logger.info("no task found on same node, resuming compression") node = db.get_storage_node_by_id(task.node_id) + storage_nodes_versions = set() + all_nodes_online = True for n in db.get_storage_nodes_by_cluster_id(node.cluster_id): - if n.status != StorageNode.STATUS_ONLINE: - msg = "Not all nodes are online, can not resume JC compression" - logger.info(msg) - task.function_result = msg - task.status = JobSchedule.STATUS_SUSPENDED - task.write_to_db(db.kv_store) + if n.status == StorageNode.STATUS_REMOVED: continue + if n.status != StorageNode.STATUS_ONLINE: + all_nodes_online = False + break + storage_nodes_versions.add(n.spdk_version) + + if not all_nodes_online: + msg = "Not all nodes are online, can not resume JC compression" + logger.info(msg) + task.function_result = msg + task.status = JobSchedule.STATUS_SUSPENDED + task.write_to_db(db.kv_store) + continue + + if len(storage_nodes_versions) > 1: + logger.error(f"Found multiple versions of storage nodes: {storage_nodes_versions}") + msg = "Not all nodes are updated yet, can not resume JC compression" + logger.info(msg) + task.function_result = msg + task.status = JobSchedule.STATUS_SUSPENDED + task.write_to_db(db.kv_store) + continue + logger.info("Resuming compression") rpc_client = node.rpc_client(timeout=5, retry=2) jm_vuid = node.jm_vuid if "jm_vuid" in task.function_params: diff --git a/simplyblock_core/storage_node_ops.py b/simplyblock_core/storage_node_ops.py index afc46480d4..f1fe5a7c0a 100644 --- a/simplyblock_core/storage_node_ops.py +++ b/simplyblock_core/storage_node_ops.py @@ -7576,12 +7576,10 @@ def _recreate_lvstore_on_non_leader_impl(snode: StorageNode, leader_node, primar logger.error("Error establishing hublvol: %s", e) # return False - # Resume JC compression for this LVS group on the restarting node - ret, err = snode.rpc_client().jc_suspend_compression(jm_vuid=primary_node.jm_vuid, suspend=False) - if not ret: - logger.info("Failed to resume JC compression adding task...") - tasks_controller.add_jc_comp_resume_task( - snode.cluster_id, snode.get_id(), jm_vuid=primary_node.jm_vuid) + # Add resume JC compression for this LVS group on the restarting node + logger.info("Adding JC compression resume task...") + tasks_controller.add_jc_comp_resume_task( + snode.cluster_id, snode.get_id(), jm_vuid=primary_node.jm_vuid) ### 2- create lvols nvmf subsystems (idempotent: skip existing) is_tertiary = (primary_node.tertiary_node_id == snode.get_id())