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
48 changes: 48 additions & 0 deletions simplyblock_core/cluster_ops.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
32 changes: 25 additions & 7 deletions simplyblock_core/services/tasks_runner_jc_comp.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
10 changes: 4 additions & 6 deletions simplyblock_core/storage_node_ops.py
Original file line number Diff line number Diff line change
Expand Up @@ -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())
Expand Down
Loading