diff options
Diffstat (limited to 'src/orchestration/views.py')
| -rw-r--r-- | src/orchestration/views.py | 572 |
1 files changed, 568 insertions, 4 deletions
diff --git a/src/orchestration/views.py b/src/orchestration/views.py index 4c3910c..cc15b3b 100644 --- a/src/orchestration/views.py +++ b/src/orchestration/views.py @@ -1,13 +1,198 @@ +import sys + +import django +import django.db from django.shortcuts import render +from django.db.models import F from django.core import serializers from django.http import HttpResponse, JsonResponse +from django.conf import settings +from django.utils import timezone +from django.views.decorators.csrf import csrf_exempt + +import json +import threading +import traceback +import threading +import logging +import uuid +import multiprocessing +from multiprocessing import Process, Queue, SimpleQueue +from http import HTTPStatus +from pulsar.schema import * + from .models import ImageSync, RemoteRegionSetup +from .imagesynchandler import ImageSyncWorker +from .namespacehandler import NamespaceWorker +from .remoteregionhandler import RemoteRegionWorker +from .kubeconfighandler import KubeconfigGenerator +from .vmb_producer import VMBProducer +from .vmb_messages import ClusterStatus, ClusterStatusMessage +from .db_updater import DBUpdater +from .vmb_handler import VMBHandler +from .vendorhandler import Vendor +import datetime +from .utils import * + +import os +import signal + + +def cleaner_thread(q, pid_list): + while True: + record = q.get() + if record == "Done": + for pid in pid_list: + print("Terminating process " + str(pid)) + os.kill(pid, signal.SIGKILL) + +def logger_thread(q): + logger = logging.getLogger("orchestration") + level = logging.DEBUG + seen = [] + while True: + record = q.get() + if record is None: + break + else: + logger.handle(record) + +def create_logger_and_start_thread(): + logger = logging.getLogger("orchestration") + logger.propagate = False + return logger + +def setup_machinery(image_handling=False, namespace_handling=False, kubeconfig_handling=False): + kubeconfigQueue = None + namespaceQueue = None + requestQueue = None + doneQueue = None + pid_list = [] + try: + + doneQueue = Queue() + + django.db.close_old_connections() + print("Hello::") + loggerQueue = Queue() + qh = logging.handlers.QueueHandler(loggerQueue) + logger = logging.getLogger(__name__) + logger.propagate = False + logger.addHandler(qh) + + dbCoordinationQ = Queue() + dbQueue = Queue() + db_updater = DBUpdater(loggerQueue, dbQueue, dbCoordinationQ) + db_updater_p = Process(target=db_updater.run) + db_updater_p.start() + db_updater_pid = db_updater_p.pid + print("DB Updater PID:" + str(db_updater_pid)) + pid_list.append(db_updater_pid) + + vmbCoordinationQ = Queue() + vmbQueue = Queue() + vm_handler = VMBHandler(loggerQueue, vmbQueue, vmbCoordinationQ) + vm_handler_p = Process(target=vm_handler.run) + vm_handler_p.start() + vmb_handler_pid = vm_handler_p.pid + print("VMB Handler PID:" + str(vmb_handler_pid)) + pid_list.append(vmb_handler_pid) + + if image_handling: + requestQueue = Queue() + image_sync_handler = ImageSyncWorker(loggerQueue, requestQueue, dbQueue, vmbQueue, dbCoordinationQ, vmbCoordinationQ, doneQueue) + image_sync_handler_p = Process(target=image_sync_handler.run) + image_sync_handler_p.start() + image_sync_handler_pid = image_sync_handler_p.pid + print("Imagesync Handler PID:" + str(image_sync_handler_pid)) + pid_list.append(image_sync_handler_pid) + + if kubeconfig_handling: + kubeconfigQueue = Queue() + kubeconfig_handler = KubeconfigGenerator(loggerQueue, kubeconfigQueue, dbQueue, vmbQueue, dbCoordinationQ, vmbCoordinationQ, doneQueue) + kubeconfig_handler_p = Process(target=kubeconfig_handler.run) + kubeconfig_handler_p.start() + kubeconfig_handler_pid = kubeconfig_handler_p.pid + print("Kubeconfig Handler PID:" + str(kubeconfig_handler_pid)) + pid_list.append(kubeconfig_handler_pid) + + if namespace_handling: + namespaceQueue = Queue() + namespace_handler = NamespaceWorker(loggerQueue, namespaceQueue, kubeconfigQueue, dbQueue, vmbQueue, doneQueue) + namespace_handler_p = Process(target=namespace_handler.run) + namespace_handler_p.start() + namespace_handler_pid = namespace_handler_p.pid + print("Namespace Handler PID:" + str(namespace_handler_pid)) + pid_list.append(namespace_handler_pid) + + ct = threading.Thread(target=cleaner_thread, args=(doneQueue,pid_list,)) + ct.start() + + lp = threading.Thread(target=logger_thread, args=(loggerQueue,)) + lp.start() + return logger, requestQueue, vmbQueue, kubeconfigQueue, namespaceQueue + except KeyboardInterrupt: + os.kill(db_updater_p.pid, signal.SIGKILL) + os.kill(image_sync_handler_p.pid, signal.SIGKILL) + os.kill(namespace_handler_p.pid, signal.SIGKILL) + os.kill(kubeconfig_handler_p.pid, signal.SIGKILL) + os.exit() +# print(sys.argv) +# if len(sys.argv) > 1 and (sys.argv[1] == "migrate" or sys.argv[1] == "makemigrations" or sys.argv[1] == "collectstatic"): +# pass +# else: +# if ((len(sys.argv) > 1 and sys.argv[1] == "runserver") or not process_started): +# if len(sys.argv) > 1 and sys.argv[1] == "runserver": +# settings.RUNSERVER = True +# print("settings.RUNSERVER: " + str(settings.RUNSERVER)) +# print("utils.process_started: " + str(process_started)) +# set_process_started() +# try: +# django.db.close_old_connections() +# print("Hello::") +# loggerQueue = Queue() +# qh = logging.handlers.QueueHandler(loggerQueue) +# logger = logging.getLogger(__name__) +# logger.addHandler(qh) + +# dbQueue = Queue() +# db_updater = DBUpdater(loggerQueue, dbQueue) +# db_updater_p = Process(target=db_updater.run) +# db_updater_p.start() + +# vmbQueue = Queue() +# vm_handler = VMBHandler(loggerQueue, vmbQueue) +# vm_handler_p = Process(target=vm_handler.run) +# vm_handler_p.start() + +# requestQueue = Queue() +# image_sync_handler = ImageSyncWorker(loggerQueue, requestQueue, dbQueue, vmbQueue) +# image_sync_handler_p = Process(target=image_sync_handler.run) +# image_sync_handler_p.start() + +# kubeconfigQueue = Queue() +# kubeconfig_handler = KubeconfigGenerator(loggerQueue, kubeconfigQueue, dbQueue, vmbQueue) +# kubeconfig_handler_p = Process(target=kubeconfig_handler.run) +# kubeconfig_handler_p.start() + +# namespaceQueue = Queue() +# namespace_handler = NamespaceWorker(loggerQueue, namespaceQueue, kubeconfigQueue, dbQueue, vmbQueue) +# namespace_handler_p = Process(target=namespace_handler.run) +# namespace_handler_p.start() + +# lp = threading.Thread(target=logger_thread, args=(loggerQueue,)) +# lp.start() +# except KeyboardInterrupt: +# os.kill(db_updater_p.pid, signal.SIGKILL) +# os.kill(image_sync_handler_p.pid, signal.SIGKILL) +# os.kill(namespace_handler_p.pid, signal.SIGKILL) +# os.kill(kubeconfig_handler_p.pid, signal.SIGKILL) +# os.exit() def index(request): return HttpResponse("Hello world") - def count(request, count): data = { 'name': 'Vitor', @@ -17,8 +202,387 @@ def count(request, count): } return JsonResponse(data) - +@csrf_exempt def imagesync(request): - rec = ImageSync.objects.order_by('-remote_region_name') - return JsonResponse(serializers.serialize('json', rec), safe=False) + requestQueue = None + try: + logger.info("Inside imagesync") + except UnboundLocalError as error: + print(error) + logger, requestQueue, _, _, _ = setup_machinery(image_handling=True) + logger.info("Inside imagesync") + + if request.method == 'GET': + rec = ImageSync.objects.order_by('-remote_region_name') + resp = serializers.serialize('json', rec) + data = json.loads(resp) + fields = data[0]['fields'] + item = None # TODO: query database + fields['data'] = item + return JsonResponse(fields) + elif request.method == 'POST': + body_unicode = request.body.decode('utf-8') + body_data = json.loads(body_unicode) + transaction_id = str(uuid.uuid4()).replace("-", "") + imageList = body_data['images'] + regionList = body_data['remoteRegions'] + regionImagePairList = get_pairs(imageList, regionList) + for rI in regionImagePairList: + region = rI['region'] + image = rI['image'] + imageSyncReq = ImageSync(remote_region_name=region, + docker_image=image, + transaction_id=str(transaction_id), + upload_start_time=datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + ) + imageSyncReq.save() + body_data['transaction_id'] = transaction_id + body_data['command'] = 'upload' + requestQueue.put(body_data) + + data = { + 'transaction_id': transaction_id + } + return JsonResponse(data, status=HTTPStatus.ACCEPTED) + +def get_status_by_transaction_id(request, transaction_id): + """ + To-do: add error handling and logging + """ + if request.method == 'GET': + value_list = ImageSync.objects.filter(transaction_id=transaction_id).values_list('remote_region_name',flat=True).distinct() + remote_regions = [] + for value in value_list: + remote_region = {} + remote_region['remoteRegion'] = value + remote_region['images'] = list(ImageSync.objects.filter( + transaction_id=transaction_id, remote_region_name=value).values(image=F('docker_image'), status=F('upload_status'), message=F('upload_message'), + start_time=F('upload_start_time'), end_time=F('upload_end_time'))) + remote_regions.append(remote_region) + + return_json = {} + return_json['transaction_id'] = transaction_id + return_json['imagesStatus'] = remote_regions + return JsonResponse(return_json, status=HTTPStatus.OK) + +def get_status_by_remote_region(request, remote_region): + """ + /regions/{remoteRegionId}/status Get remote region image status for a given region id + """ + if request.method == 'GET': + value_list = ImageSync.objects.filter(remote_region_name=remote_region).values_list('docker_image',flat=True).distinct() + image_list = [] + for value in value_list: + q_list = list(ImageSync.objects.filter(remote_region_name=remote_region, docker_image=value).order_by('-upload_start_time')[:1].values(image=F('docker_image'), status=F('upload_status'), message=F('upload_message'), start_time=F('upload_start_time'), end_time=F('upload_end_time'))) + image_list += q_list + + return_json = {} + return_json['remoteRegion'] = remote_region + return_json['images'] = image_list + return JsonResponse(return_json, status=HTTPStatus.OK) + +def get_status_by_image_name_and_remote_region(request, image_name, remote_region): + """ + /regions/{remoteRegionId}/images/{imageName}/status Get remote region image status for a given image name + """ + logger = create_logger_and_start_thread() + if request.method == 'GET': + #image_status = list(ImageSync.objects.filter(remote_region_name=remote_region, docker_image=image_name).order_by('-upload_start_time')[:1].values(imageName=F('docker_image'), imageUploadStatus=F('upload_status'), imageUploadMessage=F('upload_message'),imageUploadTime=F('upload_start_time'), remoteRegionName=F('remote_region_name'))) + image_sync_handler = ImageSyncWorker('', '', '', '', '', '', '') + logger.info("Finding image tags on a sub-cloud..." + image_name + " " + remote_region) + image_tags_list = image_sync_handler.get_image_tags(remote_region, image_name, logger) + return_json = {} + return_json['remoteRegion'] = remote_region + return_json['image'] = image_name + return_json['tags'] = image_tags_list + return JsonResponse(return_json, safe=False, status=HTTPStatus.OK) + +@csrf_exempt +def delete_image_tag_remote_region(request, image_name, image_tag, remote_region): + """ + /regions/{remoteRegionId}/images/{imageName}/{imageTag} Get remote region image status for a given image name + """ + logger = create_logger_and_start_thread() + if request.method == 'DELETE': + image_sync_handler = ImageSyncWorker('', '', '', '', '', '', '') + logger.info("Deleting image tag on a sub-cloud..." + image_name + " " + image_tag + " " + remote_region) + output = image_sync_handler.delete_image_tag(remote_region, image_name, image_tag, logger) + logger.info("O/P:") + logger.info(output) + return_json = {} + return_json['remoteRegion'] = remote_region + return_json['image'] = image_name + return_json['tag'] = image_tag + return_json['output'] = output + logger.info(return_json) + return JsonResponse(return_json, safe=False, status=HTTPStatus.OK) + +@csrf_exempt +def setup_cluster(request, cluster_name): + namespaceQueue = None + vmbQueue = None + try: + logger.info("Inside setup_cluster") + except UnboundLocalError as error: + print(error) + logger, _, vmbQueue, kubeconfigQueue, namespaceQueue = setup_machinery(kubeconfig_handling=True, namespace_handling=True) + logger.info("Inside setup_cluster. Cluster name:" + cluster_name) + remote_region_worker = RemoteRegionWorker() + site_name, site_location = remote_region_worker.get_fuze_spm_site_details(cluster_name, logger) + + # WR CAAS version; CAAS availability; deploy_status + caasversion = '20.06' + if 'caasversion' in request.GET: + caasversion = request.GET['caasversion'] + + availability = 'ONLINE' + if 'availability' in request.GET: + availability = request.GET['availability'] + + deploystatus = 'COMPLETE' + if 'deploystatus' in request.GET: + deploystatus = request.GET['deploystatus'] + + created_at = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + if 'created_at' in request.GET: + created_at = request.GET['created_at'] + + updated_at = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + logger.info("CaaS Version:" + caasversion + " Availability:" + availability + " DeployStatus:" + deploystatus) + logger.info("Created at:" + str(created_at) + " Updated at:" + str(updated_at)) + + network_setup = 'true' + if 'network_setup' in request.GET: + network_setup = request.GET['network_setup'] + + message = {} + cluster_status = {} + cluster_status['name'] = cluster_name + cluster_status['description'] = site_name + cluster_status['location'] = site_location + cluster_status['software_version'] = caasversion + cluster_status['availability'] = availability + cluster_status['deploy_status'] = deploystatus + cluster_status['created_at'] = created_at + cluster_status['updated_at'] = updated_at + message['cluster_status'] = [cluster_status] + message['network_setup'] = network_setup + logger.info(message) + namespaceQueue.put(message) + logger.info("Done putting message on Namespace Queue") + + # Create vmb object + vmbProducer = VMBProducer() + + # Dump request obj to vmb object + cluster_status_message = ClusterStatusMessage() + cluster_status_message.reportName = 'vcp_fe_cluster_status' + cluster_status_message.reportGeneratedOn = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + cluster_status_message.reportDescription = site_name + cluster_status_list = [] + cluster_status_vmb = ClusterStatus() + cluster_status_vmb.name = cluster_name + cluster_status_vmb.description = site_name + cluster_status_vmb.location = site_location + cluster_status_vmb.software_version = caasversion + cluster_status_vmb.availability = availability + cluster_status_vmb.deploy_status = deploystatus + cluster_status_vmb.created_at = created_at + cluster_status_vmb.updated_at = updated_at + cluster_status_list.append(cluster_status_vmb) + cluster_status_message.reportDataRows = cluster_status_list + cluster_status_message.rowCount = 1 + + # Send to VMB + logger.info("Putting message on VMB Queue") + item = {} + item['message'] = 'ClusterStatus' + item['payload'] = cluster_status_message + vmbQueue.put(item) + logger.info("Done putting message on VMB Queue") + + return_json = {} + return_json['cluster_name'] = cluster_name + return JsonResponse(return_json, status=HTTPStatus.OK) + +@csrf_exempt +def caas_status(request): + ''' + /caas-status Send CaaS readiness message to VMB + + expected: request payload: + + {"cluster_status": [ + { + "name": "wsbomagj-d654321-001", + "description": "NE CONCORD 8_N", + "location": 654321, + "software_version": "20.06", + "availability": "ONLINE", + "deploy_status": "COMPLETE", + "created_at": "2020-01-07 04:16:15.743617", + "updated_at": "2020-01-07 04:16:15.743617" + } + ] + } + ''' + namespaceQueue = None + vmbQueue = None + try: + logger.info("Inside caas_status") + except UnboundLocalError as error: + print(error) + logger, _, vmbQueue, kubeconfigQueue, namespaceQueue = setup_machinery(kubeconfig_handling=True, namespace_handling=True) + if request.method == 'POST': + message_json = json.loads(request.body) + logger.info("Inside caas_status") + logger.info(message_json) + + # Trigger Namespace creation + message_json['network_setup'] = 'true' + namespaceQueue.put(message_json) + + remote_region_worker = RemoteRegionWorker() + + # 1. create vmb object + vmbProducer = VMBProducer() + + # 2. dump request obj to vmb object + cluster_status_message = ClusterStatusMessage() + cluster_status_message.reportName = 'vcp_fe_cluster_status' + cluster_status_message.reportGeneratedOn = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + try: + regions_list = message_json['cluster_status'] + cluster_status_message.rowCount = len(regions_list) + cluster_status_list = [] + for i, region in enumerate(regions_list): + cluster_status = ClusterStatus() + for k in region: + if k == 'name': + cluster_status.name = region[k] + site_name, site_location = remote_region_worker.get_fuze_spm_site_details(cluster_status.name, logger) + if k == 'software_version': + cluster_status.software_version = region[k] + if k == 'availability': + cluster_status.availability = region[k] + if k == 'deploy_status': + cluster_status.deploy_status = region[k] + if k == 'created_at': + cluster_status.created_at = region[k] + if k == 'updated_at': + cluster_status.updated_at = region[k] + + cluster_status.description = site_name + cluster_status_message.reportDescription = site_name + cluster_status.location = site_location + cluster_status_list.append(cluster_status) + cluster_status_message.reportDataRows = cluster_status_list + except KeyError: + return HttpResponse("Malformed data!", HTTPStatus.BAD_REQUEST) + + # 3. send to VMB + item = {} + item['message'] = 'ClusterStatus' + item['payload'] = cluster_status_message + vmbQueue.put(item) + + return HttpResponse(status=HTTPStatus.OK) + +def get_kubeconfig_by_remote_region(request, remote_region): + ''' + /regions/{remoteRegionId}/connection-details get remote region kubeconfig by remote region id + + ''' + namespaceQueue = None + vmbQueue = None + try: + logger.info("Inside get_kubeconfig_by_remote_region") + except UnboundLocalError as error: + print(error) + logger, _, vmbQueue, kubeconfigQueue, namespaceQueue = setup_machinery(kubeconfig_handling=True) + if request.method == 'GET': + logger.info("Inside get_kubeconfig_by_remote_region") + request_data = {} + transaction_id = str(uuid.uuid4()).replace("-", "") + + kubeconfigReq = RemoteRegionSetup(remote_region_name=remote_region, + transaction_id=str(transaction_id)) + kubeconfigReq.save() + + #team = 'orchestration' # default + team = '' #default + if 'team' in request.GET: + team = request.GET['team'] + #print("TEAM:" + team) + team = team + "-team" + print("Kubeconfig for TEAM:" + team) + + kubeconfig_approach = 'SA' + if 'kubeconfig_approach' in request.GET: + kubeconfig_approach = request.GET['kubeconfig_approach'] + + print("Kubeconfig approach:" + kubeconfig_approach) + logger.info("Kubeconfig for:" + team + " using approach:" + kubeconfig_approach) + + request_data['remote_region'] = remote_region + request_data['transaction_id'] = transaction_id + request_data['kubeconfig_for'] = team + request_data['kubeconfig_approach'] = kubeconfig_approach + kubeconfigQueue.put(request_data) + + data = { + 'transaction_id': transaction_id + } + return JsonResponse(data, status=HTTPStatus.ACCEPTED) + + +def get_kubeconfig_by_transaction_id(request, remote_region, transaction_id): + """ + /regions/{remoteRegionId}/connection-details/{transactionId} get remote region kubeconfig by transactionId + """ + if request.method == 'GET': + remote_region_setup_list = list(RemoteRegionSetup.objects.filter(transaction_id=transaction_id).values(kube_config=F('kubeconfig'), namespace=F('kubernetes_namespace'), service_account=F('serviceaccount'), kubeconfig_gen_status=F('status'), kubeconfig_gen_message=F('message'))) + + return_json = {} + return_json['transaction_id'] = transaction_id + return_json['remote_region'] = remote_region + return_json['remote_region_setup'] = remote_region_setup_list + return JsonResponse(return_json, status=HTTPStatus.OK) + +def get_setup_details_by_remote_region(request, remote_region): + ''' + /regions/{remoteRegionId}/setup-details get remote region info by remote region id + ''' + logger = create_logger_and_start_thread() + if request.method == 'GET': + logger.info("Finding network status for " + remote_region) + vendor_provisioner = Vendor(logger, '', '', '') + network_details, namespaces, ns_secrets, ns_service_accounts, online_status = vendor_provisioner.check_site(remote_region, logger) + + setup_details = {} + setup_details['network_details'] = network_details + setup_details['namespaces'] = namespaces + setup_details['ns_secrets'] = ns_secrets + setup_details['ns_service_accounts'] = ns_service_accounts + setup_details['online_status'] = online_status + + return_json = {} + return_json['remoteRegion'] = remote_region + return_json['setup_details'] = setup_details + return_json['caasVendorVersion'] = '20.06' + return_json['kubernetesVersion'] = '1.16' + return_json['helmVersion'] = '3.0' + return JsonResponse(return_json, status=HTTPStatus.OK) + +@csrf_exempt +def setup_network_for_remote_region(request, remote_region): + logger = create_logger_and_start_thread() + if request.method == 'PUT': + logger.info("Setting up network for " + remote_region) + vendor_provisioner = Vendor(logger, '', '', '') + vendor_provisioner.setup_network(remote_region, logger) + data = { + 'remote_region': remote_region, + } + return JsonResponse(data, status=HTTPStatus.ACCEPTED) |
