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