summaryrefslogtreecommitdiff
path: root/src/orchestration/namespacehandler.py
diff options
context:
space:
mode:
Diffstat (limited to 'src/orchestration/namespacehandler.py')
-rw-r--r--src/orchestration/namespacehandler.py431
1 files changed, 431 insertions, 0 deletions
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
+