From 640ff61422bc0ee3966941b93d9889ddbbd38d77 Mon Sep 17 00:00:00 2001 From: ckonstanski Date: Thu, 30 Jul 2026 18:17:14 -0600 Subject: more files --- src/orchestration/namespacehandler.py | 431 ++++++++++++++++++++++++++++++++++ 1 file changed, 431 insertions(+) create mode 100644 src/orchestration/namespacehandler.py (limited to 'src/orchestration/namespacehandler.py') diff --git a/src/orchestration/namespacehandler.py b/src/orchestration/namespacehandler.py new file mode 100644 index 0000000..59597ac --- /dev/null +++ b/src/orchestration/namespacehandler.py @@ -0,0 +1,431 @@ +import django +import threading + +from django.conf import settings +from django.core.exceptions import AppRegistryNotReady +from django.utils import timezone + +from multiprocessing import Process, Pool, Queue +import pexpect +import traceback +import logging +import multiprocessing +from logging.handlers import QueueHandler +import sys +import time +import datetime +from .utils import * + +import os +import signal, psutil + +from jinja2 import Environment, FileSystemLoader + +from .vmb_messages import RemoteRegion, NameSpaceMessage +from .remoteregionhandler import RemoteRegionWorker +from .vendorhandler import Samsung + +try: + django.setup() + from .models import ImageSync + from caas.models import Cluster + from caas.models import Namespace +except django.core.exceptions.AppRegistryNotReady as exp: + pass + +class NamespaceWorker(RemoteRegionWorker): + + def __init__(self, loggerQueue, requestQueue, kubeconfigQueue, dbQueue, vmbQueue, doneQueue): + self.loggerQueue = loggerQueue + self.requestQueue = requestQueue + self.kubeconfigQueue = kubeconfigQueue + self.dbQueue = dbQueue + self.vmbQueue = vmbQueue + self.doneQueue = doneQueue + + def run(self): + qh = QueueHandler(self.loggerQueue) + self.logger = logging.getLogger() + self.logger.addHandler(qh) + self.logger.setLevel(logging.DEBUG) + + self.logger.info("NamespaceWorker started...") + + while True: + try: + item = self.requestQueue.get(block=True) + if item: + self.logger.info(item) + self._handle_request(item) + except: + pass + time.sleep(1) + + def _handle_request(self, request): + self.logger.info("Inside handle request") + + remoteRegions = [] + self.logger.info(request) + regions_list = request['cluster_status'] + network_setup = request['network_setup'] + self.logger.info(regions_list) + for region in regions_list: + self.logger.info("region " + region['name']) + remote_region_oam_ip = self._get_remote_region_oam_ip(region['name'], self.logger) + self.logger.info(" _handle_request: " + remote_region_oam_ip) + remoteRegions.append(remote_region_oam_ip) + self._create_namespace(region['name'], remote_region_oam_ip, network_setup, self.logger) + + def _create_namespace(self, region, region_oam_ip, network_setup, logger): + logger.info("Inside _create_namespace " + region) + + host_username, host_password, ssh_prefix = get_host_connection_details(region_oam_ip) + + namespace = self._get_namespace(region, logger) + logger.info("_create_namespace " + namespace + " region: " + region + " oam_ip:" + region_oam_ip) + + namespace_creation_command = [] + namespace_creation_command.append(ssh_prefix + ' kubectl create namespace --kubeconfig=/etc/kubernetes/admin.conf ' + namespace) + + namespace_exist = self._check_namespaces(region_oam_ip, namespace, logger) + if not namespace_exist: + message = 'Starting namespace creation. ' + status = 'STARTING' + logger.info(message + ' ' + status) + #self._update_status(central_image, region, message, status, transaction_id, logger) + namespace_created = self._run_commands(namespace_creation_command, host_password, logger, block=True) + + if namespace_created: + self._create_crd_rbac(namespace, region, region_oam_ip, logger) + self._create_rbd_provisioner(namespace, region, region_oam_ip, logger) + + message = 'Namespace creation done. ' + status = 'COMPLETE' + logger.info(message + ' ' + status) + #self._update_status1(central_image, region, message, status, transaction_id, logger) + else: + logger.info("Namespace creation failed.") + + if namespace_exist or namespace_created: + created_at_time = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + + site_name, site_location = self.get_fuze_spm_site_details(region, logger) + logger.info("Site name:" + site_name) + logger.info("Site location:" + site_location) + + # Send VMB Notification + remote_region = RemoteRegion(cluster=region, + namespace=namespace, + location=site_location, + created_at=str(created_at_time), + updated_at=str(created_at_time)) + + namespace_message = NameSpaceMessage(reportName='vcp_fe_namespace', + reportDescription=site_name, + reportGeneratedOn=str(created_at_time), + rowCount=1, + reportDataRows=[remote_region]) + + item = {} + item['payload'] = namespace_message + item['message'] = 'Namespace' + self.vmbQueue.put(item) + + # Trigger Kubeconfig generation + kubeconfig_request = {} + kubeconfig_request['remote_region'] = region + kubeconfig_request['transaction_id'] = "-1" + kubeconfig_request['kubeconfig_for'] = "orchestration-team" + kubeconfig_request['kubeconfig_approach'] = settings.KUBECONFIG_SRC + kubeconfig_request['namespace'] = namespace + kubeconfig_request['region_oam_ip'] = region_oam_ip + kubeconfig_request['network_setup'] = network_setup + self.kubeconfigQueue.put(kubeconfig_request) + + def _create_rbd_provisioner(self, namespace, region, remote_region_oam_ip, logger): + logger.info("Creating RBD provisioner in Namespace:" + namespace) + + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + sa_files_dir = [] + folder = 'cluster-setup-files' + sa_files_dir.append(ssh_prefix + ' mkdir -p /home/sysadmin/' + folder) + self._run_commands(sa_files_dir, host_password, self.logger, block=True) + + logger.info("About to render rbd files") + saFilePath = os.path.join(os.path.dirname(__file__), "./" + folder) + env = Environment(loader = FileSystemLoader(saFilePath), trim_blocks=True, lstrip_blocks=True) + logger.info(env) + config_data = {} + config_data["namespace"] = namespace + + rbd_template = env.get_template('rbd.yaml') + logger.info(rbd_template) + temp_file_location = get_temp_file_location() + + policy_file_location = temp_file_location + "/" + region + logger.info("Policy file location:" + policy_file_location) + if not os.path.exists(policy_file_location): + os.makedirs(policy_file_location) + fp = open(policy_file_location + "/rendered-rbd.yaml", "w") + + fp.write(rbd_template.render(config_data)) + fp.close() + + cmds_scp = [] + cmds_scp.append('scp -o StrictHostKeyChecking=no ' + policy_file_location + '/rendered-rbd.yaml ' + host_username + '@[' + remote_region_oam_ip + ']:~/' + folder + '/.') + self._run_commands_scp(cmds_scp, host_password, self.logger, block=True) + + cmds = [] + + env_string = self._read_remote_openrc(remote_region_oam_ip, logger) + basecmd = env_string + + logger.info("basecmd:" + basecmd) + cmd = ssh_prefix + " " + basecmd + " system helm-override-update --values ./" + folder + "/rendered-rbd.yaml platform-integ-apps rbd-provisioner kube-system" + logger.info("cmd:" + cmd) + cmds.append(cmd) + self._run_commands(cmds, host_password, self.logger, block=True) + + cmds = [] + cmd = ssh_prefix + " " + basecmd + " system application-apply platform-integ-apps" + logger.info("cmd:" + cmd) + cmds.append(cmd) + self._run_commands(cmds, host_password, self.logger, block=True) + + # Wait for the platform-integ-apps to status to become 'applied' + cmds = [] + cmd = ssh_prefix + " " + basecmd + " system application-show platform-integ-apps" + logger.info("cmd:" + cmd) + cmds.append(cmd) + status_applied = False + while not status_applied: + output_lines = self._run_command_get_all_lines(cmds, host_password, logger) + logger.info("Returned output lines:") + logger.info(output_lines) + if len(output_lines) > 0: + for line in output_lines.split("\n"): + logger.info("Line:" + line) + if 'status' in line and 'applied' in line: + logger.info("platform-integ-apps application applied.") + status_applied = True + time.sleep(3) + + # Check if the rbd secret has been created or not; If not create it + cmds = [] + cmd = ssh_prefix + " kubectl get secrets --kubeconfig=/etc/kubernetes/admin.conf -n " + namespace + logger.info("cmd:" + cmd) + cmds.append(cmd) + found = False + count = 0 + # Will wait for 30 seconds to ensure that the ceph rbd secret is created. + while not found and count < 10: + output_lines = self._run_command_get_all_lines(cmds, host_password, logger) + logger.info("Returned output lines:") + logger.info(output_lines) + if len(output_lines) > 0: + for line in output_lines.split("\n"): + logger.info("Line:" + line) + if 'ceph-pool-kube-rbd' in line: + logger.info("ceph-pool-kube-rbd secret created in the namespace:" + namespace) + found = True + count = count + 1 + time.sleep(3) + + if not found: + logger.info("ceph-pool-rbd-secret not found in the namespace..creating one") + cmds = [] + rbd_secret_cmd = " kubectl --kubeconfig=/etc/kubernetes/admin.conf get secret ceph-pool-kube-rbd -n default -o yaml " + #rbd_secret_cmd = rbd_secret_cmd + " | sed 's/namespace: default/namespace: " + namespace + "/' | " + #rbd_secret_cmd = rbd_secret_cmd + " sed 's/namespaces\/default/namespaces\/" + namespace + "/'" + #rbd_secret_cmd = rbd_secret_cmd + " kubectl --kubeconfig=/etc/kubernetes/admin.conf apply -n " + namespace + " -f - " + cmd = ssh_prefix + rbd_secret_cmd + logger.info("cmd:" + cmd) + cmds.append(cmd) + all_lines = self._run_command_get_all_lines(cmds, host_password, logger) + temp_file_location = get_temp_file_location() + policy_file_location = temp_file_location + "/" + region + + fp = open(policy_file_location + "/rendered-rbd-secret.yaml", "w") + for line in all_lines.split("\n"): + if 'Connection to' not in line: + if 'default' not in line: + fp.write(line + "\n") + elif 'namespace: default' in line: + fp.write(' namespace: ' + namespace + "\n") + else: + fp.write(' selfLink: /api/v1/namespaces/' + namespace + '/secrets/ceph-pool-kube-rbd' + "\n") + fp.close() + + logger.info("About to copy rendered-rbd-secret.yaml") + cmds_scp = [] + cmds_scp.append('scp -o StrictHostKeyChecking=no ' + policy_file_location + '/rendered-rbd-secret.yaml ' + host_username + '@[' + remote_region_oam_ip + ']:~/' + folder + '/.') + self._run_commands_scp(cmds_scp, host_password, logger, block=True) + + logger.info("About to apply rendered-rbd-secret.yaml") + cmds = [] + cmds.append(ssh_prefix + ' kubectl apply --kubeconfig=/etc/kubernetes/admin.conf -f ./' + folder + '/rendered-rbd-secret.yaml') + self._run_commands(cmds, host_password, logger, block=True) + + + def _create_crd_rbac(self, namespace, region, remote_region_oam_ip, logger): + logger.info("Applyin CRD RBAC policies") + + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + sa_files_dir = [] + folder = 'cluster-setup-files' + sa_files_dir.append(ssh_prefix + ' mkdir -p /home/sysadmin/' + folder) + self._run_commands(sa_files_dir, host_password, logger, block=True) + + logger.info("About to render CRD RBAC files") + saFilePath = os.path.join(os.path.dirname(__file__), "./" + folder) + env = Environment(loader = FileSystemLoader(saFilePath), trim_blocks=True, lstrip_blocks=True) + logger.info(env) + + logger.info("Getting crd details") + crd_cluster_role, crd_api_group, crd_api_resources, crd_api_verbs = self._get_crd_details() + logger.info(crd_cluster_role + ' ' + crd_api_group + ' ' + crd_api_resources + ' ' + crd_api_verbs) + crd_config = {} + crd_config["crd_cluster_role"] = crd_cluster_role + crd_config["crd_api_group"] = crd_api_group + crd_config["crd_api_resources"] = crd_api_resources + crd_config["crd_api_verbs"] = crd_api_verbs + logger.info(crd_config) + crd_role_template = env.get_template('crd-role.yaml') + logger.info(crd_role_template) + temp_file_location = get_temp_file_location() + + policy_file_location = temp_file_location + "/" + region + logger.info("Policy file location:" + policy_file_location) + if not os.path.exists(policy_file_location): + os.makedirs(policy_file_location) + + fp = open(policy_file_location + "/rendered-crd-role.yaml", "w") + logger.info("ABC") + try: + rendered_crd_template = crd_role_template.render(crd_config) + logger.info(rendered_crd_template) + except Exception as e: + logger.info(e) + fp.write(rendered_crd_template) + logger.info("DEF") + fp.close() + + logger.info("About to copy rendered-crd-role.yaml") + cmds_scp = [] + cmds_scp.append('scp -o StrictHostKeyChecking=no ' + policy_file_location + '/rendered-crd-role.yaml ' + host_username + '@[' + remote_region_oam_ip + ']:~/' + folder + '/.') + self._run_commands_scp(cmds_scp, host_password, logger, block=True) + + logger.info("About to apply rendered-crd-role.yaml") + cmds = [] + cmds.append(ssh_prefix + ' kubectl apply --kubeconfig=/etc/kubernetes/admin.conf -f ./' + folder + '/rendered-crd-role.yaml') + self._run_commands(cmds, host_password, logger, block=True) + + def _deploy_helm(self, namespace, region, remote_region_oam_ip, logger): + logger.info("Instantiating Tiller in Namespace:" + namespace) + + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + sa_files_dir = [] + folder = 'cluster-setup-files' + sa_files_dir.append(ssh_prefix + ' mkdir -p /home/sysadmin/' + folder) + self._run_commands(sa_files_dir, host_password, logger, block=True) + + logger.info("About to render Tiller files") + saFilePath = os.path.join(os.path.dirname(__file__), "./" + folder) + logger.info("SA File Path:" + saFilePath) + env = Environment(loader = FileSystemLoader(saFilePath), trim_blocks=True, lstrip_blocks=True) + logger.info(env) + + # Create Tiller Service Account and RBAC + logger.info("Getting crd details") + crd_cluster_role, crd_api_group, crd_api_resources, crd_api_verbs = self._get_crd_details() + logger.info(crd_cluster_role + ' ' + crd_api_group + ' ' + crd_api_resources + ' ' + crd_api_verbs) + tiller_config = {} + tiller_config["namespace"] = namespace + tiller_config["crd_cluster_role"] = crd_cluster_role + vdu_tiller_sa_template = env.get_template('tiller-sa.yaml') + temp_file_location = get_temp_file_location() + + policy_file_location = temp_file_location + "/" + region + logger.info("Policy file location:" + policy_file_location) + if not os.path.exists(policy_file_location): + os.makedirs(policy_file_location) + + fp = open(policy_file_location + "/rendered-tiller-sa.yaml", "w") + fp.write(vdu_tiller_sa_template.render(tiller_config)) + fp.close() + + cmds_scp = [] + cmds_scp.append('scp -o StrictHostKeyChecking=no ' + policy_file_location + '/rendered-tiller-sa.yaml ' + host_username + '@[' + remote_region_oam_ip + ']:~/' + folder + '/.') + self._run_commands_scp(cmds_scp, host_password, logger, block=True) + + cmds = [] + cmds.append(ssh_prefix + ' kubectl apply --kubeconfig=/etc/kubernetes/admin.conf -f ./' + folder + '/rendered-tiller-sa.yaml') + self._run_commands(cmds, host_password, logger, block=True) + + # Instantiate Tiller Pod + logger.info("About to render Helm files") + config_data = {} + config_data["namespace"] = namespace + + helm_template = env.get_template('helm.yaml') + logger.info(helm_template) + fp = open(policy_file_location + "/rendered-helm.yaml", "w") + fp.write(helm_template.render(config_data)) + fp.close() + logger.info("ABC") + + cmds_scp = [] + cmds_scp.append('scp -o StrictHostKeyChecking=no ' + policy_file_location + '/rendered-helm.yaml ' + host_username + '@[' + remote_region_oam_ip + ']:~/' + folder + '/.') + self._run_commands_scp(cmds_scp, host_password, logger, block=True) + + cmds = [] + cmds.append(ssh_prefix + ' kubectl apply --kubeconfig=/etc/kubernetes/admin.conf -f ./' + folder + '/rendered-helm.yaml') + self._run_commands(cmds, host_password, logger, block=True) + + def _run_commands(self, commands, host_password, logger, block=False, timeout=None): + logger.info("Inside _run_commands..timeout:" + str(timeout)) + status = run_commands(commands, host_password, logger, block=False, timeout=timeout) + return status + + def _run_commands_scp(self, commands, host_password, logger, block=False): + logger.info("Inside _run_commands_scp..") + status = run_commands_scp(commands, host_password, logger, block=False) + return status + + def _update_status(self, image, region, message, status, transaction_id, logger, upload_end_time=None): + item = {} + item['image'] = image + item['region'] = region + item['message'] = message + item['status'] = status + item['transaction_id'] = transaction_id + item['upload_end_time'] = upload_end_time + self.dbQueue.put(item) + + def _update_status1(self, image, region, message, status, transaction_id, logger): + ImageSync.objects.filter( + transaction_id=transaction_id, + docker_image=image).update(upload_message=message, + upload_status=status) + + def _check_namespaces(self, remote_region_oam_ip, namespace, logger): + + namespace_exist = False + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + cmds = [] + cmds.append(ssh_prefix + " kubectl get namespaces --kubeconfig=/etc/kubernetes/admin.conf ") + + output_lines = self._run_command_get_all_lines(cmds, host_password, logger, block=True) + logger.info("Returned output lines:") + logger.info(output_lines) + + for line in output_lines.split("\n"): + if namespace in line: + print(namespace + ' already exists') + namespace_exist = True + + return namespace_exist + -- cgit v1.3