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', 'location': 'Finland', 'is_active': True, 'count': count } return JsonResponse(data) @csrf_exempt def imagesync(request): 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)