diff options
| author | ckonstanski <kostcarl@isu.edu> | 2026-07-30 18:17:14 -0600 |
|---|---|---|
| committer | ckonstanski <kostcarl@isu.edu> | 2026-07-30 18:17:14 -0600 |
| commit | 640ff61422bc0ee3966941b93d9889ddbbd38d77 (patch) | |
| tree | abf1c08f5d38ee53ec8b29dc4f425722517505e6 /src/orchestration | |
| parent | 22dae02a86c1fce71091bfa5289cf99abba1b217 (diff) | |
Diffstat (limited to 'src/orchestration')
38 files changed, 5067 insertions, 6 deletions
diff --git a/src/orchestration/cluster-setup-files/application-ldap.yaml b/src/orchestration/cluster-setup-files/application-ldap.yaml new file mode 100644 index 0000000..2155f3a --- /dev/null +++ b/src/orchestration/cluster-setup-files/application-ldap.yaml @@ -0,0 +1,69 @@ +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: application-svc-edge-eng + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: edit +subjects: + - kind: User + name: SVC-Edge-Eng +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: {{crd_cluster_role}}-for-app-team +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: {{crd_cluster_role}}-application +subjects: + - kind: User + name: SVC-Edge-Eng +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRole +metadata: + name: {{crd_cluster_role}}-application +rules: + - apiGroups: ["{{crd_api_group}}"] + resources: ["{{crd_api_resources}}"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: crd-rb-app +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: crd-view +subjects: + - kind: User + name: SVC-Edge-Eng +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: vdu-upgrade-role + namespace: {{ namespace }} +rules: + - apiGroups: ["", "apps"] + resources: ["configmaps", "services", "secrets", "persistentvolumeclaims", "deployments"] + verbs: ["get","list","watch","delete","update"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: vdu-upgrade-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: vdu-upgrade-role +subjects: + - kind: User + name: SVC-Edge-Eng +--- diff --git a/src/orchestration/cluster-setup-files/application-sa.yaml b/src/orchestration/cluster-setup-files/application-sa.yaml new file mode 100644 index 0000000..9d87a80 --- /dev/null +++ b/src/orchestration/cluster-setup-files/application-sa.yaml @@ -0,0 +1,79 @@ +apiVersion: v1 +kind: ServiceAccount +metadata: + name: {{application_sa}} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: {{application_sa}}-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: edit +subjects: + - kind: ServiceAccount + name: {{application_sa}} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: {{crd_cluster_role}}-for-app-team +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: {{ crd_cluster_role }}-application +subjects: + - kind: ServiceAccount + name: {{application_sa}} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRole +metadata: + name: {{ crd_cluster_role }}-application +rules: + - apiGroups: ["{{ crd_api_group }}"] + resources: ["{{ crd_api_resources }}"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: crd-rb-app +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: crd-view +subjects: + - kind: ServiceAccount + name: {{ application_sa }} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: vdu-upgrade-role + namespace: {{ namespace }} +rules: + - apiGroups: ["", "apps"] + resources: ["configmaps", "services", "secrets", "persistentvolumeclaims", "deployments"] + verbs: ["get","list","watch","delete","update"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: vdu-upgrade-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: vdu-upgrade-role +subjects: + - kind: ServiceAccount + name: {{ application_sa }} + namespace: default +--- diff --git a/src/orchestration/cluster-setup-files/crd-role.yaml b/src/orchestration/cluster-setup-files/crd-role.yaml new file mode 100644 index 0000000..b500e6d --- /dev/null +++ b/src/orchestration/cluster-setup-files/crd-role.yaml @@ -0,0 +1,18 @@ +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRole +metadata: + name: {{ crd_cluster_role }} +rules: + - apiGroups: ["{{ crd_api_group }}"] + resources: ["{{ crd_api_resources }}"] + verbs: ["{{ crd_api_verbs }}"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRole +metadata: + name: crd-view +rules: + - apiGroups: ["apiextensions.k8s.io"] + resources: ["crds","customresourcedefinitions"] + verbs: ["*"] +--- diff --git a/src/orchestration/cluster-setup-files/helm.yaml b/src/orchestration/cluster-setup-files/helm.yaml new file mode 100644 index 0000000..f0ec72c --- /dev/null +++ b/src/orchestration/cluster-setup-files/helm.yaml @@ -0,0 +1,75 @@ +--- +apiVersion: apps/v1 +kind: Deployment +metadata: + creationTimestamp: null + labels: + app: helm + name: tiller + name: tiller-deploy + namespace: {{ namespace }} +spec: + replicas: 1 + selector: {"matchLabels": {"app": "helm", "name": "tiller"}} + strategy: {} + template: + metadata: + creationTimestamp: null + labels: + app: helm + name: tiller + spec: + automountServiceAccountToken: true + containers: + - env: + - name: TILLER_NAMESPACE + value: {{ namespace }} + - name: TILLER_HISTORY_MAX + value: "0" + image: registry.local:9001/gcr.io/kubernetes-helm/tiller:v2.13.1 + imagePullPolicy: IfNotPresent + livenessProbe: + httpGet: + path: /liveness + port: 44135 + initialDelaySeconds: 1 + timeoutSeconds: 1 + name: tiller + ports: + - containerPort: 44134 + name: tiller + - containerPort: 44135 + name: http + readinessProbe: + httpGet: + path: /readiness + port: 44135 + initialDelaySeconds: 1 + timeoutSeconds: 1 + resources: {} + serviceAccountName: vdu-tiller +status: {} + +--- +apiVersion: v1 +kind: Service +metadata: + creationTimestamp: null + labels: + app: helm + name: tiller + name: tiller-deploy + namespace: {{ namespace }} +spec: + ports: + - name: tiller + port: 44134 + targetPort: tiller + selector: + app: helm + name: tiller + type: ClusterIP +status: + loadBalancer: {} + +... diff --git a/src/orchestration/cluster-setup-files/orchestration-ldap.yaml b/src/orchestration/cluster-setup-files/orchestration-ldap.yaml new file mode 100644 index 0000000..ee03ba3 --- /dev/null +++ b/src/orchestration/cluster-setup-files/orchestration-ldap.yaml @@ -0,0 +1,166 @@ +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: orchestration-user-crb +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: view +subjects: + - kind: User + name: SVC-FE-Atlas +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: orchestration-user-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: edit +subjects: + - kind: User + name: SVC-FE-Atlas +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: {{crd_cluster_role}}-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: {{ crd_cluster_role }} +subjects: + - kind: User + name: SVC-FE-Atlas +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: configmap-role + namespace: {{ namespace }} +rules: + - apiGroups: [""] + resources: ["configmaps"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: configmap-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: configmap-role +subjects: + - kind: User + name: SVC-FE-Atlas +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: service-role + namespace: {{ namespace }} +rules: + - apiGroups: [""] + resources: ["services"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: service-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: service-role +subjects: + - kind: User + name: SVC-FE-Atlas +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: deployment-role + namespace: {{ namespace }} +rules: + - apiGroups: ["apps"] + resources: ["deployments"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: deployment-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: deployment-role +subjects: + - kind: User + name: SVC-FE-Atlas +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: secret-role + namespace: {{ namespace }} +rules: + - apiGroups: [""] + resources: ["secrets"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: secret-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: secret-role +subjects: + - kind: User + name: SVC-FE-Atlas +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: pvc-role + namespace: {{ namespace }} +rules: + - apiGroups: [""] + resources: ["persistentvolumeclaims"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: pvc-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: pvc +subjects: + - kind: User + name: SVC-FE-Atlas +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: crd-rb-orch +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: crd-view +subjects: + - kind: User + name: SVC-FE-Atlas +--- + diff --git a/src/orchestration/cluster-setup-files/orchestration-sa.yaml b/src/orchestration/cluster-setup-files/orchestration-sa.yaml new file mode 100644 index 0000000..5531431 --- /dev/null +++ b/src/orchestration/cluster-setup-files/orchestration-sa.yaml @@ -0,0 +1,181 @@ +apiVersion: v1 +kind: ServiceAccount +metadata: + name: {{orchestration_sa}} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: {{orchestration_sa}}-crb +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: view +subjects: + - kind: ServiceAccount + name: {{orchestration_sa}} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: {{orchestration_sa}}-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: edit +subjects: + - kind: ServiceAccount + name: {{orchestration_sa}} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: {{crd_cluster_role}}-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: {{ crd_cluster_role }} +subjects: + - kind: ServiceAccount + name: {{orchestration_sa}} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: configmap-role + namespace: {{ namespace }} +rules: + - apiGroups: [""] + resources: ["configmaps"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: configmap-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: configmap-role +subjects: + - kind: ServiceAccount + name: {{ orchestration_sa }} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: service-role + namespace: {{ namespace }} +rules: + - apiGroups: [""] + resources: ["services"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: service-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: service-role +subjects: + - kind: ServiceAccount + name: {{ orchestration_sa }} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: deployment-role + namespace: {{ namespace }} +rules: + - apiGroups: ["apps"] + resources: ["deployments"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: deployment-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: deployment-role +subjects: + - kind: ServiceAccount + name: {{ orchestration_sa }} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: secret-role + namespace: {{ namespace }} +rules: + - apiGroups: [""] + resources: ["secrets"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: secret-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: secret-role +subjects: + - kind: ServiceAccount + name: {{ orchestration_sa }} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: pvc-role + namespace: {{ namespace }} +rules: + - apiGroups: [""] + resources: ["persistentvolumeclaims"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: pvc-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: pvc +subjects: + - kind: ServiceAccount + name: {{ orchestration_sa }} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: crd-rb-orch +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: crd-view +subjects: + - kind: ServiceAccount + name: {{ orchestration_sa }} + namespace: default +--- + diff --git a/src/orchestration/cluster-setup-files/rbd.yaml b/src/orchestration/cluster-setup-files/rbd.yaml new file mode 100644 index 0000000..45d497b --- /dev/null +++ b/src/orchestration/cluster-setup-files/rbd.yaml @@ -0,0 +1,9 @@ +classes: +- additionalNamespaces: [default, kube-public, {{ namespace }} ] + chunk_size: 64 + crush_rule_name: storage_tier_ruleset + name: general + pool_name: kube-rbdkube-system + replication: 1 + userId: ceph-pool-kube-rbd + userSecretName: ceph-pool-kube-rbd diff --git a/src/orchestration/cluster-setup-files/samsung-sa.yaml b/src/orchestration/cluster-setup-files/samsung-sa.yaml new file mode 100644 index 0000000..2525e8b --- /dev/null +++ b/src/orchestration/cluster-setup-files/samsung-sa.yaml @@ -0,0 +1,55 @@ +apiVersion: v1 +kind: ServiceAccount +metadata: + name: {{application_sa}} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: {{application_sa}}-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: edit +subjects: + - kind: ServiceAccount + name: {{application_sa}} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: {{crd_cluster_role}}-for-app-team +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: {{ crd_cluster_role }}-application +subjects: + - kind: ServiceAccount + name: {{application_sa}} + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRole +metadata: + name: {{ crd_cluster_role }}-application +rules: + - apiGroups: ["{{ crd_api_group }}"] + resources: ["{{ crd_api_resources }}"] + verbs: ["get", "list", "watch"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: crd-rb-app +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: crd-view +subjects: + - kind: ServiceAccount + name: {{ application_sa }} + namespace: default +--- diff --git a/src/orchestration/cluster-setup-files/tiller-sa.yaml b/src/orchestration/cluster-setup-files/tiller-sa.yaml new file mode 100644 index 0000000..62cc860 --- /dev/null +++ b/src/orchestration/cluster-setup-files/tiller-sa.yaml @@ -0,0 +1,42 @@ +apiVersion: v1 +kind: ServiceAccount +metadata: + name: vdu-tiller + namespace: {{ namespace }} +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: vdu-tiller-role + namespace: {{ namespace }} +rules: + - apiGroups: ["*"] + resources: ["*"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: vdu-tiller-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: vdu-tiller-role +subjects: + - kind: ServiceAccount + name: vdu-tiller + namespace: {{ namespace }} +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: vdu-tiller-cluster-crd-crb +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: {{crd_cluster_role}} +subjects: + - kind: ServiceAccount + name: vdu-tiller + namespace: {{ namespace }} diff --git a/src/orchestration/cluster-setup-files/vdu_tiller_sa.yaml b/src/orchestration/cluster-setup-files/vdu_tiller_sa.yaml new file mode 100644 index 0000000..3644856 --- /dev/null +++ b/src/orchestration/cluster-setup-files/vdu_tiller_sa.yaml @@ -0,0 +1,29 @@ +apiVersion: v1 +kind: ServiceAccount +metadata: + name: vdu-tiller + namespace: {{ namespace }} +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: vdu-tiller-role + namespace: {{ namespace }} +rules: + - apiGroups: ["*"] + resources: ["*"] + verbs: ["*"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: vdu-tiller-rb + namespace: {{ namespace }} +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: vdu-tiller-role +subjects: + - kind: ServiceAccount + name: vdu-tiller + namespace: {{ namespace }} diff --git a/src/orchestration/data.json b/src/orchestration/data.json new file mode 100644 index 0000000..927707c --- /dev/null +++ b/src/orchestration/data.json @@ -0,0 +1,4 @@ +{ + "images": ["i1", "i2", "i5"], + "remoteRegions": ["r1", "r2", "r3"] +}
\ No newline at end of file diff --git a/src/orchestration/db_updater.py b/src/orchestration/db_updater.py new file mode 100644 index 0000000..6837dd3 --- /dev/null +++ b/src/orchestration/db_updater.py @@ -0,0 +1,83 @@ +import django + +from django.conf import settings +from django.core.exceptions import AppRegistryNotReady +from django.db import transaction +import logging +from logging.handlers import QueueHandler + +try: + django.setup() + from .models import ImageSync, CentralToRemoteMap, RemoteRegionSetup +except django.core.exceptions.AppRegistryNotReady as exp: + pass + +class DBUpdater(): + + def __init__(self, loggerQueue, dbQueue, coordinateQueue): + self.loggerQueue = loggerQueue + self.dbQueue = dbQueue + self.coordinateQueue = coordinateQueue + + def run(self): + qh = QueueHandler(self.loggerQueue) + self.logger = logging.getLogger() + self.logger.addHandler(qh) + self.logger.setLevel(logging.DEBUG) + + self.logger.info("DBUpdater started...") + while True: + #self.logger.info("--------------------------") + item = self.dbQueue.get() + msg_type = item['type'] + self.logger.info("Message type:" + msg_type) + if msg_type == 'kubeconfig': + self._update_kubeconfig(item) + if msg_type == 'image': + self._update_image_status(item) + + def _update_image_status(self, item): + self.logger.info("Inside _update_image_status") + image = item['image'] + region = item['region'] + message = item['message'] + status = item['status'] + transaction_id = item['transaction_id'] + upload_end_time = item['upload_end_time'] + + #ImageSync.objects.filter(remote_region_name=region, + # transaction_id=transaction_id, + # docker_image=image).update(upload_message=message, + # upload_status=status, + # upload_end_time=upload_end_time) + + ImageSync.objects.filter( + transaction_id=transaction_id, + docker_image=image).update(upload_message=message, + upload_status=status, + upload_end_time=upload_end_time) + + done = item['done'] + if done: + self.coordinateQueue.put("Done") + return + + def _update_kubeconfig(self, item): + transaction_id = item['transaction_id'] + message = item['message'] + status = item['status'] + saName = item['saName'] + namespace = item['namespace'] + region = item['region'] + kubeconfig = item['kubeconfig'] + RemoteRegionSetup.objects.filter(transaction_id=transaction_id).update(message=message, + status=status, + serviceaccount=saName, + kubernetes_namespace=namespace, + remote_region_name=region, + kubeconfig=kubeconfig) + done = item['done'] + if done: + self.coordinateQueue.put("Done") + return + diff --git a/src/orchestration/imagesynchandler.py b/src/orchestration/imagesynchandler.py new file mode 100644 index 0000000..cddf1ac --- /dev/null +++ b/src/orchestration/imagesynchandler.py @@ -0,0 +1,622 @@ +import django +import threading + +from django.conf import settings +from django.core.exceptions import AppRegistryNotReady + +from multiprocessing import Process, Pool, Queue +import pexpect +import traceback +import logging +import multiprocessing +from logging.handlers import QueueHandler +import re +import sys +import time +import datetime +from .utils import * + +import os +import signal, psutil + +from .db_updater import DBUpdater +from .vmb_messages import ImageStatus, ImagesStatusMessage +from .remoteregionhandler import RemoteRegionWorker + +try: + django.setup() + from .models import ImageSync, CentralToRemoteMap, RemoteRegionSetup + from caas.models import Cluster +except django.core.exceptions.AppRegistryNotReady as exp: + pass + +class ImageSyncWorker(RemoteRegionWorker): + + def __init__(self, loggerQueue, requestQueue, dbQueue, vmbQueue, dbCoordinationQueue, vmbCoordinationQueue, doneQueue): + self.loggerQueue = loggerQueue + self.requestQueue = requestQueue + self.dbQueue = dbQueue + self.vmbQueue = vmbQueue + self.doneQueue = doneQueue + self.dbCoordinationQueue = dbCoordinationQueue + self.vmbCoordinationQueue = vmbCoordinationQueue + #self.dbQueue = Queue() + #self.db_updater = DBUpdater(loggerQueue, self.dbQueue) + #self.db_updater_p = Process(target=self.db_updater.run) + #self.db_updater_p.start() + + # https://stackoverflow.com/questions/3332043/obtaining-pid-of-child-process + # Current not used; Leaving it here if needed in the future + def _kill_child_processes(self, parent_pid, sig=signal.SIGTERM): + try: + parent = psutil.Process(parent_pid) + except psutil.NoSuchProcess: + return + children = parent.children(recursive=True) + for process in children: + #if not process.is_alive(): + process.send_signal(sig) + + def _kill_child_process(self, pid): + os.kill(pid, signal.SIGKILL) + + def run(self): + qh = QueueHandler(self.loggerQueue) + self.logger = logging.getLogger() + self.logger.addHandler(qh) + self.logger.setLevel(logging.DEBUG) + + self.logger.info("ImageSyncWorker 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): + command = request['command'] + if command == 'upload': + self._handle_request_image_upload(request) + if command == 'delete': + self._handle_request_image_delete(request) + if command == 'cancel': + self._handle_request_image_upload_cancel(request) + + def _handle_request_image_delete(self, request): + pass + + def _handle_request_image_upload_cancel(self, request): + pass + + def _setup_logger(self): + qh = QueueHandler(self.loggerQueue) + self.logger = logging.getLogger() + self.logger.addHandler(qh) + self.logger.setLevel(logging.INFO) + + def get_image_tags(self, remote_region, image_name, logger): + remote_region_oam_ip = self._get_remote_region_oam_ip(remote_region, "-1", logger) + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + env_string = self._read_remote_openrc(remote_region_oam_ip, logger) + basecmd = env_string + + cmds = [] + cmd = ssh_prefix + " " + basecmd + " system registry-image-tags " + image_name + logger.info("cmd:" + cmd) + cmds.append(cmd) + tags = [] + output_lines = self._run_command_get_all_lines(cmds, host_password, logger) + #logger.info("Returned output lines:") + #logger.info(output_lines) + print(output_lines) + if len(output_lines) > 0: + for line in output_lines.split("\n"): + #logger.info("Line:" + line) + print("Line:" + line) + if 'Connection to' not in line and 'Authorization failed' not in line: + parts = line.split(" ") + if len(parts) > 1: + if parts[1] != 'Image': + print("Line:" + line) + tags.append(parts[1]) + if 'Authorization failed' in line: + tags.append(line) + logger.info("Tags:" + str(tags)) + print("Tags:" + str(tags)) + return tags + + def delete_image_tag(self, remote_region, image_name, image_tag, logger): + remote_region_oam_ip = self._get_remote_region_oam_ip(remote_region, "-1", logger) + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + env_string = self._read_remote_openrc(remote_region_oam_ip, logger) + basecmd = env_string + + cmds = [] + cmd = ssh_prefix + " " + basecmd + " system registry-image-delete " + image_name + ":" + image_tag + logger.info("cmd:" + cmd) + cmds.append(cmd) + tags = [] + output_lines = self._run_command_get_all_lines(cmds, host_password, logger) + + cmds = [] + cmd = ssh_prefix + " " + basecmd + " system registry-garbage-collect " + logger.info("cmd:" + cmd) + cmds.append(cmd) + tags = [] + output_lines1 = self._run_command_get_all_lines(cmds, host_password, logger) + logger.info("Output lines1:") + logger.info(output_lines1) + + lines_to_return = [] + for line in output_lines.split('\n'): + if 'Connection to' not in line: + lines_to_return.append(line) + + for line in output_lines1.split('\n'): + if 'Connection to' not in line: + lines_to_return.append(line) + + logger.info("Output lines:") + logger.info(lines_to_return) + return lines_to_return + + def _handle_request_image_upload(self, request): + transaction_id = request['transaction_id'] + self.logger.info("Transction ID:" + str(transaction_id)) + imageList = request['images'] + remoteRegions = request['remoteRegions'] + centralRegions = [] + for region in remoteRegions: + centralRegionList = self._get_central_region_oam_ip(region, transaction_id) + for centralRegion in centralRegionList: + if centralRegion not in centralRegions: + centralRegions.append(centralRegion) + remoteregion_oam_ip = self._get_remote_region_oam_ip(region, transaction_id, self.logger) + self.logger.info("Remote region Name:" + region + " Remote region OAM IP:" + remoteregion_oam_ip) + self.logger.info(str(transaction_id) + " Downloading images to Central regions.." + ','.join(centralRegions)) + successful_image_list, failed_image_list = self._handle_central_regions(centralRegions, remoteRegions, imageList, transaction_id, region_type='central') + + if len(successful_image_list) > 0: + self.logger.info(str(transaction_id) + " Downloading images to Remote regions.." + ','.join(remoteRegions) + " " + ','.join(successful_image_list)) + self._handle_regions(remoteRegions, successful_image_list, transaction_id, region_type='remote') + else: + self.logger.info(str(transaction_id) + " Could not download images to Central region..So not progressing to Remote region") + self._wait_and_done(imageList, remoteRegions, transaction_id) + return + + def _handle_central_regions(self, regionList, remoteRegionList, imageList, transaction_id, region_type=''): + workers = [] + imageRegionList = get_pairs(imageList, regionList) + self.logger.info(str(transaction_id) + " Handling central region") + self.logger.info(str(transaction_id) + " ImageList:" + ','.join(imageList)) + self.logger.info(str(transaction_id) + " RegionList:" + ','.join(regionList)) + self.logger.info(str(transaction_id) + " ImageRegionList:" + str(imageRegionList)) + for region in regionList: + self.logger.info(str(transaction_id) + " Region:" + region) + successful_image_list, failed_image_list = self._handle_image_list(region, imageList, self.logger, region_type, transaction_id) + # We want to download from Artifactory in its own process so that the parent can terminate the process if needed; + # Without doing this in its own process, the parent will block foreever and we won't be able to process subsequent calls. + #worker = multiprocessing.Process(target=self._handle_image_list, + # args=(region, imageList, self.logger, region_type, transaction_id)) + #workers.append(worker) + #worker.start() + #self.logger.info("Central controller handling process pid:" + str(worker.pid)) + + #https://stackoverflow.com/questions/26063877/python-multiprocessing-module-join-processes-with-timeout + #TIMEOUT = 300 # 5 minutes + #start = time.time() + #while time.time() - start <= TIMEOUT: + # if not any(p.is_alive() for p in workers): + # # All the processes are done, break now. + # break + # time.sleep(1) # Just to avoid hogging the CPU + #else: + # # We only enter this if we didn't 'break' above. + # print("Central downloader timed out, terminating the process...") + # for w in workers: + # self.logger.info("Terminating process handling central region download - pid:" + str(w.pid)) + # self._kill_child_process(w.pid) + # return download_to_central_done + + #for w in workers: + # w.join() + + return successful_image_list, failed_image_list + + def _handle_regions(self, regionList, imageList, transaction_id, region_type=''): + self.logger.info(str(transaction_id) + " Handling remote regions") + workers = [] + imageRegionList = get_pairs(imageList, regionList) + self.logger.info(str(transaction_id) + " ImageList:" + ','.join(imageList)) + self.logger.info(str(transaction_id) + " RegionList:" + ','.join(regionList)) + self.logger.info(str(transaction_id) + " ImageRegionList:" + str(imageRegionList)) + for item in imageRegionList: + self.logger.info(str(transaction_id) + " Item:" + str(item)) + worker = multiprocessing.Process(target=self._worker_process, + args=(self.logger, self.loggerQueue, self._worker_configurer, item, region_type, transaction_id)) + workers.append(worker) + worker.start() + time.sleep(3) # Stagger the requests + + #https://stackoverflow.com/questions/26063877/python-multiprocessing-module-join-processes-with-timeout + TIMEOUT = 3600 # 1 hour + start = time.time() + while time.time() - start <= TIMEOUT: + if not any(p.is_alive() for p in workers): + # All the processes are done, break now. + break + time.sleep(1) # Just to avoid hogging the CPU + else: + # We only enter this if we didn't 'break' above. + print("timed out, killing all processes") + for w in workers: + self.logger.info(str(transaction_id) + " Terminating process handling remote region download - pid:" + str(w.pid)) + self._kill_child_process(w.pid) + return + + # No need to wait for processes to finish; otherwise we won't be able to pick-up any new incoming requests for an hour. + #for w in workers: + # w.join() + + def _wait_and_done(self, imageList, remoteRegions, transaction_id): + # Ensure that VMB messages have been sent + self.logger.info(str(transaction_id) + " About to be done..waiting for cleanup") + num_of_images = len(imageList) + num_of_regions = len(remoteRegions) + self.logger.info(str(transaction_id) + " Number of images.." + str(num_of_images)) + self.logger.info(str(transaction_id) + " Number of regions.." + str(num_of_regions)) + count = 0 + while count < num_of_images * num_of_regions: + vmb_coordination = self.vmbCoordinationQueue.get() + self.logger.info(str(transaction_id) + " VMB coordination message:" + vmb_coordination) + count = count + 1 + + # Ensure that all DB messages have been processed + db_coordination = self.dbCoordinationQueue.get() + self.logger.info(str(transaction_id) + " DB coordination message:" + db_coordination) + + self.logger.info(str(transaction_id) + " Done") + time.sleep(10) + self.doneQueue.put("Done") + + def _worker_process(self, lg, queue, configurer, item, region_type, transaction_id): + #configurer(queue) + name = multiprocessing.current_process().name + self._handle_image(item, lg, region_type, transaction_id) + + # Currently not used; left here if we want to configure the loggers to add any extra logging (formats, etc.) + def _worker_configurer(self, queue): + h = logging.handlers.QueueHandler(queue) + root = logging.getLogger() + root.addHandler(h) + root.setLevel(logging.INFO) + + def _handle_image(self, imageRegionPair, logger, region_type, transaction_id): + image = imageRegionPair["image"] + regionname = imageRegionPair["region"] + region = self._get_remote_region_oam_ip(regionname, transaction_id, logger) + name = multiprocessing.current_process().name + logger.info(str(transaction_id) + " Process:" + name + " " + image + " " + region) + + logger.info(str(transaction_id) + " Checking if image exist") + imageexists = self._check_image_exists(image, region) + logger.info(str(transaction_id) + " Image exists status:" + str(imageexists)) + if not imageexists: + message = 'Starting Image download from Central Controller' + status = 'STARTED' + self._update_status_async(image, region, message, status, transaction_id, logger) + success, err = self._download_image(image, region, logger, region_type, transaction_id) + if success: + really_there = False + dest_image_name_and_tag = self._get_dest_image_name(image) + parts = dest_image_name_and_tag.split(":") + dest_image_name = '' + if len(parts) > 0: + dest_image_name = parts[0] + logger.info("Destination image name:" + dest_image_name) + dest_image_tags = self.get_image_tags(regionname, dest_image_name, logger) + really_there = self._is_it_really_there(image, dest_image_tags, transaction_id, logger) + if really_there: + message = 'Image download to Remote Region complete.' + status = 'COMPLETED' + else: + message = 'Image download to Remote Region failed. Reason: Could not find image on remote region.' + status = 'FAILED' + else: + message = 'Image download to Remote Region failed. Reason:' + err + status = 'FAILED' + image_upload_completion_time = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + self._update_status_async(image, region, message, status, transaction_id, logger, upload_end_time=image_upload_completion_time, done=True) + + self._notify_vmb(region, regionname, image, status, message, image_upload_completion_time, transaction_id, logger) + + def _is_it_really_there(self, image, dest_image_tags, transaction_id, logger): + logger.info(str(transaction_id) + " Comparing image tags..") + logger.info(str(transaction_id) + " Image:" + image) + logger.info(str(transaction_id) + " Destination Image tags:" + str(dest_image_tags)) + present = False + for tag in dest_image_tags: + logger.info(str(transaction_id) + " Tag:" + tag) + logger.info(str(transaction_id) + " Image:" + image) + if tag in image: + present = True + break + return present + + def _notify_vmb(self, region, regionname, image, status, message, image_upload_completion_time, transaction_id, logger): + + # Notify VMB + image_status1 = ImageStatus( + cluster=regionname, + image=image, + action='UPLOAD', + status=status, + message=message, + created_at=str(image_upload_completion_time), + ) + + logger.info(str(transaction_id) + " Image status::" + str(image_status1)) + + site_name, site_location = self.get_fuze_spm_site_details(regionname, logger) + logger.info("Site name:" + site_name) + logger.info("Site location:" + site_location) + reportDescription = 'Image upload status : ' + site_name + + images_status_list = [] + images_status_message = ImagesStatusMessage(reportName='vcp_fe_imagestatus', + transactionId=transaction_id, + reportDescription=reportDescription, + reportGeneratedOn=str(datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S"))) + images_status_list.append(image_status1) + images_status_message.rowCount = 1 + images_status_message.reportDataRows = images_status_list + + logger.info(str(transaction_id) + " Image status message::" + str(images_status_message)) + + item = {} + item['message'] = 'ImageStatus' + item['payload'] = images_status_message + self.vmbQueue.put(item) + + def _get_central_region_oam_ip(self, remoteregion, transaction_id): + remoteclusterObj = Cluster.objects.filter(cluster_name=remoteregion) + parent_cluster_id = remoteclusterObj[0].parent_cluster_id + self.logger.info(str(transaction_id) + " Parent cluster id.." + str(parent_cluster_id)) + centralclusterObj = Cluster.objects.filter(id=parent_cluster_id) + self.logger.info(str(transaction_id) + " " + str(centralclusterObj)) + central_region_list = [] + for central_region in centralclusterObj: + central_region_list.append(central_region.oam_vip_address) + self.logger.info(str(transaction_id) + " Central Region List:" + ','.join(central_region_list)) + return central_region_list + + def _get_remote_region_oam_ip(self, regionname, transaction_id, logger): + remoteclusterObj = Cluster.objects.filter(cluster_name=regionname) + oam_ip = remoteclusterObj[0].oam_vip_address + logger.info(str(transaction_id) + " Remote region:" + regionname + " OAM IP:" + str(oam_ip)) + return oam_ip + + def _get_central_region_prev(self, remoteregion, transaction_id): + centralToRemoteMapObj = CentralToRemoteMap.objects.filter(remote_region_name=remoteregion) + self.logger.info(str(transaction_id) + " Inside _get_central_region...") + self.logger.info(str(transaction_id) + " " + str(centralToRemoteMapObj)) + central_region_list = [] + for central_region in centralToRemoteMapObj: + central_region_list.append(central_region.central_region_name) + #central_region_name = centralToRemoteMapObj[0].central_region_name + self.logger.info(str(transaction_id) + " Central Region List:" + ','.join(central_region_list)) + return central_region_list + + def _check_image_exists(self, image, region): + # TODO: Query database for a quick check; Query region for accurate check + return False + + def _handle_image_list(self, region, imageList, logger, region_type, transaction_id): + logger.info(str(transaction_id) + " Handling image list") + name = multiprocessing.current_process().name + logger.info(str(transaction_id) + " Process:" + name + " " + region) + + central_region_oam_ip = region + logger.info("1a") + central_region_name = self._get_central_region_name(central_region_oam_ip, logger) + logger.info("Central region name:" + central_region_name) + + host_username = settings.HOST_CREDS['username'] + host_password = settings.HOST_CREDS['password'] + art_username = settings.ARTIFACTORY_CREDS['username'] + art_password = settings.ARTIFACTORY_CREDS['password'] + dr_username = settings.CENTRAL_DR_CREDS['username'] + dr_password = settings.CENTRAL_DR_CREDS['password'] + + failed_image_list = [] + successful_image_list = [] + for image in imageList: + image_parts = image.split("/") + # Image: http://vnf-twb.vzwnet.com/docker-images/samsung_vdu_vdu_svr20aa5vvzwg06a3_r05_20.a.0-0101_6.0.0/vzw-adpf-rmp:svr20aa5vvzwg06a3_r05 + artifactory_host = image_parts[0] # ' vsp.vici.verizon.com:7443' + logger.info("Artifactory host:" + artifactory_host) + ssh_prefix = 'ssh -o StrictHostKeyChecking=no -t ' + host_username + '@' + region + login_to_artifactory = [] + login_to_artifactory.append(ssh_prefix + ' sudo docker login --username ' + art_username + ' --password ' + art_password + ' ' + artifactory_host) + artifactory_login_status, cmd_err = self._run_commands(login_to_artifactory, host_password, logger, transaction_id) + if not artifactory_login_status: + message = 'Could not download Image from Artifactory. Reason:' + cmd_err + status = 'FAILED' + failed_image_list.append(image) + image_upload_completion_time = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + self._update_status_async(image, region, message, status, transaction_id, logger, upload_end_time=image_upload_completion_time) + logger.info(str(transaction_id) + " Notifying VMB") + self._notify_vmb(central_region_name, central_region_name, image, status, message, image_upload_completion_time, transaction_id, logger) + else: + message = 'Starting Image download from Artifactory.' + status = 'STARTED' + self._update_status_async(image, region, message, status, transaction_id, logger) + + image_name = self._get_image_name(image) + download_from_artifactory = [] + download_from_artifactory.append(ssh_prefix + ' sudo docker pull ' + image) + logger.info(str(transaction_id) + " Downloading " + image) + image_download_status, cmd_err1 = self._run_commands(download_from_artifactory, host_password, logger, transaction_id, cmd_timeout=300) # 5 minutes + + if image_download_status: + message = 'Image download from Artifactory complete.' + status = 'COMPLETED' + successful_image_list.append(image) + else: + message = 'Could not download Image from Artifactory. Reason:' + cmd_err1 + status = 'FAILED' + failed_image_list.append(image) + logger.info(str(transaction_id) + " Notifying VMB") + self._notify_vmb(central_region_name, central_region_name, image, status, message, image_upload_completion_time, transaction_id, logger) + image_upload_completion_time = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + self._update_status_async(image, region, message, status, transaction_id, logger, upload_end_time=image_upload_completion_time) + + logger.info("**********") + logger.info("Successful Image list:" + str(len(successful_image_list))) + logger.info("----------") + logger.info("Failed Image list:" + str(len(failed_image_list))) + if len(successful_image_list) > 0: + login_to_dr_central = [] + login_to_dr_central.append(ssh_prefix + ' sudo docker login --username ' + dr_username + ' --password ' + dr_password + ' registry.local:9001') + self._run_commands(login_to_dr_central, host_password, logger, transaction_id) + + for image in successful_image_list: + message = 'Pushing Image to Central Controller Docker Registry.' + status = 'STARTED' + self._update_status_async(image, region, message, status, transaction_id, logger) + + image_name = self._get_image_name(image) + dest_image = 'registry.local:9001/' + image_name + push_to_local_dr_central = [] + push_to_local_dr_central.append(ssh_prefix + ' sudo docker tag ' + image + ' ' + dest_image) + push_to_local_dr_central.append(ssh_prefix + ' sudo docker push ' + dest_image) + logger.info(str(transaction_id) + " Pushing image to Central Docker Registry " + dest_image) + self._run_commands(push_to_local_dr_central, host_password, logger, transaction_id) + + message = 'Pushing Image to Central Controller Docker Registry complete.' + status = 'COMPLETED' + image_upload_completion_time = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + self._update_status_async(image, region, message, status, transaction_id, logger, upload_end_time=image_upload_completion_time) + + return successful_image_list, failed_image_list + + # TODO: Try two times + # https://gitlab.verizon.com/vcp/webscale/far-edge/caas/orchestration/-/blob/master/docker_image_sync_steps.txt + def _download_image(self, source_image, region, logger, region_type, transaction_id): + image_name = self._get_image_name(source_image) + central_image = 'registry.central:9001/' + image_name + dest_image_name = self._get_dest_image_name(image_name) + dest_image = 'registry.local:9001/' + dest_image_name + logger.info("Destination Image name:" + dest_image) + host_username = settings.HOST_CREDS['username'] + host_password = settings.HOST_CREDS['password'] + dr_username = settings.CENTRAL_DR_CREDS['username'] + dr_password = settings.CENTRAL_DR_CREDS['password'] + + ssh_prefix = 'ssh -o StrictHostKeyChecking=no -t ' + host_username + '@' + region + download_from_central_dr = [] + download_from_central_dr.append(ssh_prefix + ' sudo docker login --username ' + dr_username + ' --password ' + dr_password + ' registry.central:9001') + download_from_central_dr.append(ssh_prefix + ' sudo docker pull ' + central_image) + push_to_local_dr_remote = [] + push_to_local_dr_remote.append(ssh_prefix + ' sudo docker tag ' + central_image + ' ' + dest_image) + push_to_local_dr_remote.append(ssh_prefix + ' sudo docker login --username ' + dr_username + ' --password ' + dr_password + ' registry.local:9001') + push_to_local_dr_remote.append(ssh_prefix + ' sudo docker push ' + dest_image) + + image_download_complete = False + if region_type == 'remote': + logger.info(str(transaction_id) + " Downloading image from central region") + message = 'Downloading Image from Central Controller Docker Registry.' + status = 'STARTED' + self._update_status_async(central_image, region, message, status, transaction_id, logger) + image_download_complete, cmd_err = self._run_commands(download_from_central_dr, host_password, logger, transaction_id, cmd_timeout=3600) # 1 hour + if image_download_complete: + message = 'Downloading Image from Central Controller Docker Registry.' + status = 'COMPLETED' + image_upload_completion_time = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + self._update_status_async(central_image, region, message, status, transaction_id, logger, upload_end_time=image_upload_completion_time) + + message = 'Pushing Image to Remote region Docker Registry.' + status = 'STARTED' + self._update_status_async(central_image, region, message, status, transaction_id, logger) + logger.info(str(transaction_id) + " Push to local docker registry") + self._run_commands(push_to_local_dr_remote, host_password, logger, transaction_id) + message = 'Pushing Image to Remote region Docker Registry complete.' + status = 'COMPLETED' + image_upload_completion_time = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + self._update_status_async(central_image, region, message, status, transaction_id, logger, upload_end_time=image_upload_completion_time) + else: + message = 'Downloading Image from Central Controller Docker Registry Failed. Reason:' + cmd_err + status = 'FAILED' + image_upload_completion_time = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + self._update_status_async(central_image, region, message, status, transaction_id, logger, upload_end_time=image_upload_completion_time) + return image_download_complete, cmd_err + + def _run_commands(self, commands, host_password, logger, transaction_id, cmd_timeout=None): + #logger.info(str(transaction_id) + " " + str(commands)) + image_download_complete = False + image_download_error = '' + for command in commands: + logger.info(command) + child = pexpect.spawn(command) + #logger.info("-- Child PID:" + str(child.pid)) + child.timeout=cmd_timeout + try: + #child.expect(['password: '], timeout=cmd_timeout) + #child.expect('\r\n\r\n\w+')#, timeout=cmd_timeout) + #child.expect('\w+\r\n') + child.expect_exact('password: ', timeout=cmd_timeout) + child.sendline(host_password) + #child.expect(['Password: '], timeout=cmd_timeout) + #child.expect('\w+\r\n') + child.expect_exact('Password: ', timeout=cmd_timeout) + child.sendline(host_password) + #child.sendline("\r\n") + all_lines = child.read() + all_lines = all_lines.rstrip().lstrip() + all_lines = all_lines.decode('utf-8').replace('\r\n', '\n') + #logger.info(all_lines) + image_download_complete = True # Tentative + for line in all_lines.split("\n"): + logger.info(str(transaction_id) + " " + line) + if re.search('error', line, re.IGNORECASE) or 'tag does not exist' in line: #or 'Error response from daemon' in line: + image_download_complete = False + image_download_error = line + except: + logger.info(str(transaction_id) + " " + str(child)) + return image_download_complete, image_download_error + + def _get_image_name(self, image): + i = image.find("/") + image_name = image[i+1:] + return image_name + + def _get_dest_image_name(self, image): + i = image.rfind("/") + image_name = image[i+1:] + return image_name + + def _update_status_async(self, image, region, message, status, transaction_id, logger, upload_end_time=None, done=False): + item = {} + item['type'] = 'image' + item['image'] = image + item['region'] = region + item['message'] = message + item['status'] = status + item['transaction_id'] = transaction_id + item['upload_end_time'] = upload_end_time + item['done'] = done # done flag indicates when we are done using db_handler + self.dbQueue.put(item) + + def _update_status_sync(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) + return diff --git a/src/orchestration/kubeconfighandler.py b/src/orchestration/kubeconfighandler.py new file mode 100644 index 0000000..d7318b2 --- /dev/null +++ b/src/orchestration/kubeconfighandler.py @@ -0,0 +1,926 @@ +import django + +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 logging +import multiprocessing +from logging.handlers import QueueHandler +import sys +import time +import datetime +from .utils import * + +import os +import signal, psutil +import json + +from jinja2 import Environment, FileSystemLoader + +from .vmb_messages import Kubeconfig, KubeconfigMessage +from .remoteregionhandler import RemoteRegionWorker +from .vendorhandler import Vendor + +try: + django.setup() + from .models import RemoteRegionSetup + from caas.models import Cluster + from caas.models import Namespace +except django.core.exceptions.AppRegistryNotReady as exp: + pass + +class KubeconfigGenerator(RemoteRegionWorker): + + ORCHESTRATION_TEAM = 'orchestration-team' + APPLICATION_TEAM = 'application-team' + VENDOR_TEAM_SAMSUNG = 'samsung-team' + OPS_TEAM = 'ops-team' + + def __init__(self, loggerQueue, requestQueue, dbQueue, vmbQueue, dbCoordinationQueue, vmbCoordinationQueue, doneQueue): + self.loggerQueue = loggerQueue + self.requestQueue = requestQueue + self.dbQueue = dbQueue + self.vmbQueue = vmbQueue + self.dbCoordinationQueue = dbCoordinationQueue + self.vmbCoordinationQueue = vmbCoordinationQueue + 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("KubeconfigGenerator started...") + + while True: + try: + item = self.requestQueue.get(block=True) + if item: + self.logger.info(item) + self.get_kubeconfig(item) + self.logger.info("Done generating kubeconfig") + self.logger.info("Performing vendor setup...") + vendor_provisioner = Vendor(self.logger, self.dbCoordinationQueue, self.vmbCoordinationQueue, self.doneQueue) + vendor_provisioner.perform_vendor_setup(item) + except: + pass + time.sleep(1) + + def get_kubeconfig(self, request): + transaction_id = request['transaction_id'] + region = request['remote_region'] + kubeconfig_for = request['kubeconfig_for'] + kubeconfig_approach = request['kubeconfig_approach'] + namespace = self._get_namespace(region, self.logger) + if namespace != "": + self.logger.info(" Transaction:" + transaction_id) + self.logger.info(" Region:" + region) + self.logger.info(" Namespace:" + namespace) + kubeconfig = self._create_kubeconfig(kubeconfig_for, region, namespace, kubeconfig_approach, transaction_id) + return kubeconfig + else: + message = " Could not find Namespace for region " + region + self.logger.info(message) + status = "FAILED" + self._update_status_async(region, namespace, '', message, status, transaction_id, self.logger, kubeconfig='', done=True) + return json.dumps({}) + + def _get_serviceaccount_name(self, team, transaction_id): + if team == KubeconfigGenerator.ORCHESTRATION_TEAM: + return "orchestration-sa" + if team == KubeconfigGenerator.APPLICATION_TEAM: + return "application-sa-" + str(transaction_id) + if team == KubeconfigGenerator.OPS_TEAM: + return "ops-sa" + if team == KubeconfigGenerator.VENDOR_TEAM_SAMSUNG: + return "samsung-sa" + + def _get_user_name(self, team): + if team == KubeconfigGenerator.ORCHESTRATION_TEAM: + return "SVC-FE-Atlas" + if team == KubeconfigGenerator.APPLICATION_TEAM: + return "SVC-Edge-Eng" + + def _apply_application_rbac_policies_sa(self, region, remote_region_oam_ip, app_sa_namespace, namespace, saName, saFilePath, logger): + logger.info("Inside _apply_application_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, self.logger, block=True) + + logger.info("About to render Application 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) + + crd_cluster_role, crd_api_group, crd_api_resources, crd_api_verbs = self._get_crd_details() + app_config = {} + app_config["application_sa"] = saName + app_config["namespace"] = namespace + app_config["crd_cluster_role"] = crd_cluster_role + app_config["crd_api_group"] = crd_api_group + app_config["crd_api_resources"] = crd_api_resources + app_rbac_template = env.get_template('application-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-application-sa.yaml", "w") + fp.write(app_rbac_template.render(app_config)) + fp.close() + + logger.info("About to copy rendered-application-rbac.yaml") + cmds_scp = [] + cmds_scp.append('scp -o StrictHostKeyChecking=no ' + policy_file_location + '/rendered-application-sa.yaml ' + host_username + '@[' + remote_region_oam_ip + ']:~/' + folder + '/.') + self._run_commands_scp(cmds_scp, host_password, self.logger, block=True) + + logger.info("About to apply rendered-application-sa.yaml") + cmds = [] + cmds.append(ssh_prefix + ' kubectl apply --kubeconfig=/etc/kubernetes/admin.conf -f ./' + folder + '/rendered-application-sa.yaml') + self._run_commands(cmds, host_password, self.logger, block=True) + + def _apply_application_rbac_policies_ldap(self, region, remote_region_oam_ip, namespace, logger): + logger.info("Inside _apply_application_rbac_policies ldap") + + 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 Application RBAC files") + saFilePath = os.path.join(os.path.dirname(__file__), "./" + folder) + self.logger.info(" RBAC File Path:" + str(saFilePath)) + env = Environment(loader = FileSystemLoader(saFilePath), trim_blocks=True, lstrip_blocks=True) + logger.info(env) + + crd_cluster_role, crd_api_group, crd_api_resources, crd_api_verbs = self._get_crd_details() + orch_config = {} + orch_config["namespace"] = namespace + orch_config["crd_cluster_role"] = crd_cluster_role + orch_config["crd_api_group"] = crd_api_group + orch_config["crd_api_resources"] = crd_api_resources + orch_rbac_template = env.get_template('application-ldap.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-application-ldap.yaml", "w") + fp.write(orch_rbac_template.render(orch_config)) + fp.close() + + logger.info("About to copy rendered-application-ldap.yaml") + cmds_scp = [] + cmds_scp.append('scp -o StrictHostKeyChecking=no ' + policy_file_location + '/rendered-application-ldap.yaml ' + host_username + '@[' + remote_region_oam_ip + ']:~/' + folder + '/.') + self._run_commands_scp(cmds_scp, host_password, self.logger, block=True) + + logger.info("About to apply rendered-application-ldap.yaml") + cmds = [] + cmds.append(ssh_prefix + ' kubectl apply --kubeconfig=/etc/kubernetes/admin.conf -f ./' + folder + '/rendered-application-ldap.yaml') + self._run_commands(cmds, host_password, self.logger, block=True) + + def _apply_vendor_samsung_rbac_policies(self, region, remote_region_oam_ip, app_sa_namespace, namespace, saName, saFilePath, logger): + logger.info("Inside _apply_vendor_samsung_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, self.logger, block=True) + + logger.info("About to render Vendor Samsung 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) + + crd_cluster_role, crd_api_group, crd_api_resources, crd_api_verbs = self._get_crd_details() + app_config = {} + app_config["application_sa"] = saName + app_config["namespace"] = namespace + app_config["crd_cluster_role"] = crd_cluster_role + app_config["crd_api_group"] = crd_api_group + app_config["crd_api_resources"] = crd_api_resources + app_rbac_template = env.get_template('samsung-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-samsung-sa.yaml", "w") + fp.write(app_rbac_template.render(app_config)) + fp.close() + + logging.info("About to copy rendered-samsung-rbac.yaml") + cmds_scp = [] + cmds_scp.append('scp -o StrictHostKeyChecking=no ' + policy_file_location + '/rendered-samsung-sa.yaml ' + host_username + '@[' + remote_region_oam_ip + ']:~/' + folder + '/.') + self._run_commands_scp(cmds_scp, host_password, self.logger, block=True) + + logger.info("About to apply rendered-samsung-sa.yaml") + cmds = [] + cmds.append(ssh_prefix + ' kubectl apply --kubeconfig=/etc/kubernetes/admin.conf -f ./' + folder + '/rendered-samsung-sa.yaml') + self._run_commands(cmds, host_password, self.logger, block=True) + + def _apply_orchestration_rbac_policies_ldap(self, region, remote_region_oam_ip, namespace, logger): + logger.info("Inside _apply_orchestration_rbac_policies ldap") + + 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 Orchestration RBAC files") + saFilePath = os.path.join(os.path.dirname(__file__), "./" + folder) + self.logger.info(" RBAC File Path:" + str(saFilePath)) + env = Environment(loader = FileSystemLoader(saFilePath), trim_blocks=True, lstrip_blocks=True) + logger.info(env) + + crd_cluster_role, crd_api_group, crd_api_resources, crd_api_verbs = self._get_crd_details() + orch_config = {} + orch_config["namespace"] = namespace + orch_config["crd_cluster_role"] = crd_cluster_role + orch_rbac_template = env.get_template('orchestration-ldap.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-orchestration-ldap.yaml", "w") + fp.write(orch_rbac_template.render(orch_config)) + fp.close() + + logger.info("About to copy rendered-orchestration-ldap.yaml") + cmds_scp = [] + cmds_scp.append('scp -o StrictHostKeyChecking=no ' + policy_file_location + '/rendered-orchestration-ldap.yaml ' + host_username + '@[' + remote_region_oam_ip + ']:~/' + folder + '/.') + self._run_commands_scp(cmds_scp, host_password, self.logger, block=True) + + logger.info("About to apply rendered-orchestration-ldap.yaml") + cmds = [] + cmds.append(ssh_prefix + ' kubectl apply --kubeconfig=/etc/kubernetes/admin.conf -f ./' + folder + '/rendered-orchestration-ldap.yaml') + self._run_commands(cmds, host_password, self.logger, block=True) + + def _apply_orchestration_rbac_policies_sa(self, region, remote_region_oam_ip, orch_sa_namespace, namespace, saName, saFilePath, logger): + logger.info("Inside _apply_orchestration_rbac_policies sa") + + 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 Orchestration 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) + + crd_cluster_role, crd_api_group, crd_api_resources, crd_api_verbs = self._get_crd_details() + orch_config = {} + orch_config["orchestration_sa"] = saName + orch_config["namespace"] = namespace + orch_config["crd_cluster_role"] = crd_cluster_role + orch_rbac_template = env.get_template('orchestration-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-orchestration-sa.yaml", "w") + fp.write(orch_rbac_template.render(orch_config)) + fp.close() + + logger.info("About to copy rendered-orchestration-rbac.yaml") + cmds_scp = [] + cmds_scp.append('scp -o StrictHostKeyChecking=no ' + policy_file_location + '/rendered-orchestration-sa.yaml ' + host_username + '@[' + remote_region_oam_ip + ']:~/' + folder + '/.') + self._run_commands_scp(cmds_scp, host_password, self.logger, block=True) + + logger.info("About to apply rendered-orchestration-sa.yaml") + cmds = [] + cmds.append(ssh_prefix + ' kubectl apply --kubeconfig=/etc/kubernetes/admin.conf -f ./' + folder + '/rendered-orchestration-sa.yaml') + self._run_commands(cmds, host_password, self.logger, block=True) + + def _create_rbac_files_clusteradmin(self, region, namespace, saName, saFilePath, logger): + logger.info("Inside _create_rbac_files") + if not os.path.exists(saFilePath): + os.makedirs(saFilePath) + + sa_metadata = {} + sa_metadata["namespace"] = namespace + sa_metadata["name"] = saName + + subjects_list = [] + subjects = {} + subjects["kind"] = "ServiceAccount" + subjects["name"] = saName + subjects["namespace"] = namespace + subjects_list.append(subjects) + + 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 + "/sa-role.json", "w") + sa_role = {} + sa_role["apiVersion"] = "rbac.authorization.k8s.io/v1" + sa_role["kind"] = "Role" + sa_role["metadata"] = sa_metadata + sa_rules_list = [] + sa_rules = {} + sa_rules["apiGroups"] = ["*"] + sa_rules["resources"] = ["*"] + sa_rules["verbs"] = ["*"] + sa_rules_list.append(sa_rules) + sa_role["rules"] = sa_rules_list + sa_role_json = json.dumps(sa_role) + logger.info("sa_role_json:" + str(sa_role_json)) + fp.write(sa_role_json) + + fp = open(policy_file_location + "/sa-rolebinding.json", "w") + sa_role_binding = {} + sa_role_binding["apiVersion"] = "rbac.authorization.k8s.io/v1" + sa_role_binding["kind"] = "RoleBinding" + sa_role_binding["metadata"] = sa_metadata + sa_role_binding["subjects"] = subjects_list + role_ref = {} + role_ref["kind"] = "Role" + role_ref["name"] = saName + role_ref["apiGroup"] = "rbac.authorization.k8s.io" + sa_role_binding["roleRef"] = role_ref + sa_role_binding_json = json.dumps(sa_role_binding) + fp.write(sa_role_binding_json) + + fp = open(temp_file_location + "/sa-clusterrole.json", "w") + sa_clusterrole = {} + sa_clusterrole["apiVersion"] = "rbac.authorization.k8s.io/v1" + sa_clusterrole["kind"] = "ClusterRole" + sa_clusterrole["metadata"] = sa_metadata + sa_rules_list = [] + sa_rules = {} + sa_rules["apiGroups"] = [""] + sa_rules["resources"] = ["*"] + sa_rules["verbs"] = ["*"] + sa_rules_list.append(sa_rules) + sa_clusterrole["rules"] = sa_rules_list + sa_clusterrole_json = json.dumps(sa_clusterrole) + fp.write(sa_clusterrole_json) + + fp = open(temp_file_location + "/sa-clusterrolebinding.json", "w") + sa_clusterrole_binding = {} + sa_clusterrole_binding["apiVersion"] = "rbac.authorization.k8s.io/v1" + sa_clusterrole_binding["kind"] = "ClusterRoleBinding" + sa_clusterrole_binding["metadata"] = sa_metadata + sa_clusterrole_binding["subjects"] = subjects_list + clusterrole_ref = {} + clusterrole_ref["kind"] = "ClusterRole" + clusterrole_ref["name"] = saName + clusterrole_ref["apiGroup"] = "rbac.authorization.k8s.io" + sa_clusterrole_binding["roleRef"] = clusterrole_ref + sa_clusterrole_binding_json = json.dumps(sa_clusterrole_binding) + fp.write(sa_clusterrole_binding_json) + fp.close() + + def _wait_for_oidc_app(self, region, namespace): + self.logger.info("Inside checking _wait_for_oidc_app...") + + remote_region_oam_ip = self._get_remote_region_oam_ip(region, self.logger) + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + env_string = self._read_remote_openrc(remote_region_oam_ip, self.logger) + basecmd = env_string + + status_applied = False + while not status_applied: + cmds = [] + cmd = ssh_prefix + " " + basecmd + " system application-show oidc-auth-apps" + self.logger.info("cmd:" + cmd) + cmds.append(cmd) + output_lines = self._run_command_get_all_lines(cmds, host_password, self.logger) + self.logger.info("Returned output lines:") + self.logger.info(output_lines) + if len(output_lines) > 0: + for line in output_lines.split("\n"): + self.logger.info("Line:" + line) + if 'status' in line and ('applied' in line or 'apply-failed' in line): + self.logger.info("oidc-auth-apps applied: " + line) + if 'applied' in line: + status_applied = True + if 'apply-failed' in line: + cmds1 = [] + cmd1 = ssh_prefix + " " + basecmd + " system application-remove oidc-auth-apps" + self.logger.info("cmd:" + cmd1) + cmd2 = ssh_prefix + " " + basecmd + " system application-apply oidc-auth-apps" + self.logger.info("cmd:" + cmd2) + cmds1.append(cmd1) + cmds1.append(cmd2) + output_lines1 = self._run_command_get_all_lines(cmds1 , host_password, self.logger) + self.logger.info("Returned output lines:") + self.logger.info(output_lines1) + time.sleep(3) + + def _create_kubeconfig(self, kubeconfig_for, region, namespace, kubeconfig_approach, transaction_id): + self.logger.info("Kubeconfig for:" + kubeconfig_for) + kubeconfig_src = kubeconfig_approach + if kubeconfig_for == KubeconfigGenerator.ORCHESTRATION_TEAM: + if kubeconfig_src == 'SA': + self._create_kubeconfig_orchestration_sa(region, namespace, transaction_id) + if kubeconfig_src == 'LDAP': + #self._wait_for_oidc_app(region, namespace) + self._create_kubeconfig_orchestration_ldap(region, namespace, transaction_id) + elif kubeconfig_for == KubeconfigGenerator.APPLICATION_TEAM: + if kubeconfig_src == 'SA': + self._create_kubeconfig_application_sa(region, namespace, transaction_id) + if kubeconfig_src == 'LDAP': + self._create_kubeconfig_application_ldap(region, namespace, transaction_id) + self.vmbCoordinationQueue.put("Done") + elif kubeconfig_for == KubeconfigGenerator.OPS_TEAM: + self._create_kubeconfig_ops(region, namespace, transaction_id) + elif kubeconfig_for == KubeconfigGenerator.VENDOR_TEAM_SAMSUNG: + self._create_kubeconfig_samsung(region, namespace, transaction_id) + self.vmbCoordinationQueue.put("Done") + + def _create_kubeconfig_application_sa(self, region, namespace, transaction_id): + self.logger.info("Inside ..kubeconfig application") + saName = self._get_serviceaccount_name(KubeconfigGenerator.APPLICATION_TEAM, transaction_id) + self.logger.info("Service Account name:" + saName) + + message = 'Starting kubeconfig creation. ' + status = 'STARTING' + self.logger.info(message + ' ' + status) + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig='') + + remote_region_oam_ip = self._get_remote_region_oam_ip(region, self.logger) + + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + app_sa_namespace = 'default' # default ns is better as it already exists + + 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) + + saFilePath = os.path.join(os.path.dirname(__file__), "./" + folder) + + message = 'Applying RBAC to Service Account' + status = 'CREATING' + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig='') + self.logger.info(" Service Account File Path:" + str(saFilePath)) + + self._apply_application_rbac_policies_sa(region, remote_region_oam_ip, app_sa_namespace, namespace, saName, saFilePath, self.logger) + + cmds_token_name = [] + cmds_token_name.append(ssh_prefix + " kubectl describe serviceaccount --kubeconfig=/etc/kubernetes/admin.conf -n " + app_sa_namespace + " " + saName + "| grep Tokens ") + all_lines = self._run_commands(cmds_token_name, host_password, self.logger, block=True) + secretname = self._parse_token_name(all_lines) + + message = 'Parsing token' + status = 'CREATING' + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig='') + self.logger.info(" Secret name:" + secretname) + if secretname != None: + cmds1 = [] + cmds1.append(ssh_prefix + " kubectl describe secret --kubeconfig=/etc/kubernetes/admin.conf -n " + app_sa_namespace + " " + secretname + " | grep token:") + all_lines = self._run_commands(cmds1, host_password, self.logger, block=True) + token = self._parse_token(all_lines) + #self.logger.info("TOKEN TO USE:" + token) + token = token.strip() + self.logger.info(" Token:[" + str(token) + "]") + + # Generate kubeconfig + kubeconfig_value = self._generate_kubeconfig(region, namespace, saName, token, transaction_id, self.logger) + + # Update database + message = 'kubeconfig creation done. ' + status = 'COMPLETE' + self.logger.info(message + ' ' + status) + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig=kubeconfig_value, done=True) + + def _create_kubeconfig_application_ldap(self, region, namespace, transaction_id): + self.logger.info("Inside kubeconfig application ldap...") + message = 'Starting kubeconfig creation. ' + status = 'STARTING' + self.logger.info(message + ' ' + status) + app_user = self._get_user_name(KubeconfigGenerator.APPLICATION_TEAM) + self.logger.info("Application User:" + app_user) + self._update_status_async(region, namespace, app_user, message, status, transaction_id, self.logger, kubeconfig='') + + remote_region_oam_ip = self._get_remote_region_oam_ip(region, self.logger) + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + message = 'Applying RBAC to User Account' + status = 'CREATING' + self._update_status_async(region, namespace, app_user, message, status, transaction_id, self.logger, kubeconfig='') + + self.logger.info("Applying RBAC policies to the Application User...") + self._apply_application_rbac_policies_ldap(region, remote_region_oam_ip, namespace, self.logger) + + ldap_account = 'edge_eng' + token = self._get_ldap_token(host_password, ssh_prefix, app_user, ldap_account, remote_region_oam_ip) + + # Generate kubeconfig + kubeconfig_value = self._generate_kubeconfig(region, namespace, app_user, token, transaction_id, self.logger) + + # Update database + message = 'kubeconfig creation done. ' + status = 'COMPLETE' + self.logger.info(message + ' ' + status) + self._update_status_async(region, namespace, app_user, message, status, transaction_id, self.logger, kubeconfig=kubeconfig_value, done=True) + + #mv /home/sysadmin/.kube/config /home/sysadmin/.kube/config-for-orchestration + self.logger.info("Moving /home/sysadmin/.kube/config to /home/sysadmin/.kube/config-for-edge-eng...") + set_kubeconfig_cmd = [] + set_kubeconfig_cmd.append(ssh_prefix + ' mv /home/sysadmin/.kube/config /home/sysadmin/.kube/config-for-edge-eng') + self._run_commands(set_kubeconfig_cmd, host_password, self.logger, block=True) + + #export KUBECONFIG=/etc/kubernetes/admin.conf + self.logger.info("Resetting KUBECONFIG...") + set_kubeconfig_cmd = [] + set_kubeconfig_cmd.append(ssh_prefix + ' export KUBECONFIG=/etc/kubernetes/admin.conf ') + self._run_commands(set_kubeconfig_cmd, host_password, self.logger, block=True) + return kubeconfig_value + + def _create_kubeconfig_samsung(self, region, namespace, transaction_id): + self.logger.info("Inside ..kubeconfig samsung") + saName = self._get_serviceaccount_name(KubeconfigGenerator.VENDOR_TEAM_SAMSUNG, transaction_id) + self.logger.info("Service Account name:" + saName) + + message = 'Starting kubeconfig creation. ' + status = 'STARTING' + self.logger.info(message + ' ' + status) + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig='') + + remote_region_oam_ip = self._get_remote_region_oam_ip(region, self.logger) + + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + app_sa_namespace = 'default' # default ns is better as it already exists + + 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) + + saFilePath = os.path.join(os.path.dirname(__file__), "./" + folder) + + message = 'Applying RBAC to Service Account' + status = 'CREATING' + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig='') + self.logger.info(" Service Account File Path:" + str(saFilePath)) + + self._apply_vendor_samsung_rbac_policies(region, remote_region_oam_ip, app_sa_namespace, namespace, saName, saFilePath, self.logger) + + cmds_token_name = [] + cmds_token_name.append(ssh_prefix + " kubectl describe serviceaccount --kubeconfig=/etc/kubernetes/admin.conf -n " + app_sa_namespace + " " + saName + "| grep Tokens ") + all_lines = self._run_commands(cmds_token_name, host_password, self.logger, block=True) + secretname = self._parse_token_name(all_lines) + + message = 'Parsing token' + status = 'CREATING' + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig='') + self.logger.info(" Secret name:" + secretname) + if secretname != None: + cmds1 = [] + cmds1.append(ssh_prefix + " kubectl describe secret --kubeconfig=/etc/kubernetes/admin.conf -n " + app_sa_namespace + " " + secretname + " | grep token:") + all_lines = self._run_commands(cmds1, host_password, self.logger, block=True) + token = self._parse_token(all_lines) + #self.logger.info("TOKEN TO USE:" + token) + token = token.strip() + self.logger.info(" Token:[" + str(token) + "]") + + # Generate kubeconfig + kubeconfig_value = self._generate_kubeconfig(region, namespace, saName, token, transaction_id, self.logger) + + # Update database + message = 'kubeconfig creation done. ' + status = 'COMPLETE' + self.logger.info(message + ' ' + status) + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig=kubeconfig_value, done=True) + + + def _create_kubeconfig_ops(self, region, namespace, transaction_id): + pass + + def _create_kubeconfig_orchestration_ldap(self, region, namespace, transaction_id): + self.logger.info("Inside kubeconfig orchestration ldap...") + message = 'Starting kubeconfig creation. ' + status = 'STARTING' + self.logger.info(message + ' ' + status) + orch_user = self._get_user_name(KubeconfigGenerator.ORCHESTRATION_TEAM) + self.logger.info("Orchestration User:" + orch_user) + self._update_status_async(region, namespace, orch_user, message, status, transaction_id, self.logger, kubeconfig='') + + remote_region_oam_ip = self._get_remote_region_oam_ip(region, self.logger) + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + message = 'Applying RBAC to User Account' + status = 'CREATING' + self._update_status_async(region, namespace, orch_user, message, status, transaction_id, self.logger, kubeconfig='') + + self.logger.info("Applying RBAC policies to the Orchestration User...") + self._apply_orchestration_rbac_policies_ldap(region, remote_region_oam_ip, namespace, self.logger) + + ldap_account = 'orchestration' + token = self._get_ldap_token(host_password, ssh_prefix, orch_user, ldap_account, remote_region_oam_ip) + + # Generate kubeconfig + kubeconfig_value = self._generate_kubeconfig(region, namespace, orch_user, token, transaction_id, self.logger) + + # Update database + message = 'kubeconfig creation done. ' + status = 'COMPLETE' + self.logger.info(message + ' ' + status) + self._update_status_async(region, namespace, orch_user, message, status, transaction_id, self.logger, kubeconfig=kubeconfig_value, done=True) + + #mv /home/sysadmin/.kube/config /home/sysadmin/.kube/config-for-orchestration + self.logger.info("Moving /home/sysadmin/.kube/config to /home/sysadmin/.kube/config-for-orchestration...") + set_kubeconfig_cmd = [] + set_kubeconfig_cmd.append(ssh_prefix + ' mv /home/sysadmin/.kube/config /home/sysadmin/.kube/config-for-orchestration') + self._run_commands(set_kubeconfig_cmd, host_password, self.logger, block=True) + + #export KUBECONFIG=/etc/kubernetes/admin.conf + self.logger.info("Resetting KUBECONFIG...") + set_kubeconfig_cmd = [] + set_kubeconfig_cmd.append(ssh_prefix + ' export KUBECONFIG=/etc/kubernetes/admin.conf ') + self._run_commands(set_kubeconfig_cmd, host_password, self.logger, block=True) + + # Send notification on VMB + self._send_vmb_notification(region, namespace, transaction_id, kubeconfig_value) + return kubeconfig_value + + def _get_ldap_token(self, host_password, ssh_prefix, user, ldap_account, remote_region_oam_ip): + #cp /etc/kubernetes/admin.conf /home/sysadmin/.kube/config + self.logger.info("Starting token generation...") + cp_cmd = [] + cp_cmd.append(ssh_prefix + ' cp /etc/kubernetes/admin.conf /home/sysadmin/.kube/config') + self._run_commands(cp_cmd, host_password, self.logger, block=True) + + #export KUBECONFIG=/home/sysadmin/.kube/config + self.logger.info("Setting KUBECONFIG...") + set_kubeconfig_cmd = [] + set_kubeconfig_cmd.append(ssh_prefix + ' export KUBECONFIG=/home/sysadmin/.kube/config ') + self._run_commands(set_kubeconfig_cmd, host_password, self.logger, block=True) + + #kubectl config set-context --kubeconfig=/home/sysadmin/.kube/config SVC-Edge-Eng@kubernetes --cluster=kubernetes --user=SVC-Edge-Eng + self.logger.info("Performing set-context...") + set_context = ' kubectl config set-context --kubeconfig=/home/sysadmin/.kube/config ' + user + '@kubernetes --cluster=kubernetes --user=' + user + self.logger.info("Set context cmd:" + set_context) + set_context_cmd = [] + set_context_cmd.append(ssh_prefix + set_context) + self._run_commands(set_context_cmd, host_password, self.logger, block=True) + + #oidc-auth -c <OAM-IP> -u SVC-Edge-Eng -p 322C6v22acuhAGdyce22S3w282 + self.logger.info("Executing oidc-auth...") + orch_password = settings.LDAP_USER_CREDS[ldap_account] + oidc_auth = ' oidc-auth -c ' + remote_region_oam_ip + ' -u ' + user + ' -p ' + orch_password + self.logger.info("OIDC Auth..:" + oidc_auth) + oidc_auth_cmd = [] + oidc_auth_cmd.append(ssh_prefix + oidc_auth) + self._run_commands(oidc_auth_cmd, host_password, self.logger, block=True) + + #grep token /home/sysadmin/.kube/config + self.logger.info("Retrieving token...") + cmds1 = [] + cmds1.append(ssh_prefix + " grep token /home/sysadmin/.kube/config ") + all_lines = self._run_commands(cmds1, host_password, self.logger, block=True) + token = self._parse_token(all_lines) + #self.logger.info("TOKEN TO USE:" + token) + token = token.strip() + self.logger.info(" Token:[" + str(token) + "]") + return token + + def _create_kubeconfig_orchestration_sa(self, region, namespace, transaction_id): + # For orchestration SA, we use either the 'orchestration' namespace or the 'default' namespace + self.logger.info("Inside ..kubeconfig orchestration sa") + saName = self._get_serviceaccount_name(KubeconfigGenerator.ORCHESTRATION_TEAM, transaction_id) + self.logger.info("Service Account name:" + saName) + + message = 'Starting kubeconfig creation. ' + status = 'STARTING' + self.logger.info(message + ' ' + status) + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig='') + + remote_region_oam_ip = self._get_remote_region_oam_ip(region, self.logger) + orch_sa_namespace = 'default' # default ns is better as it already exists + + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + #cmds_sa = [] + #cmds_sa.append(ssh_prefix + ' kubectl create serviceaccount --kubeconfig=/etc/kubernetes/admin.conf ' + saName + ' -n ' + orch_sa_namespace) + #self._run_commands(cmds_sa, host_password, self.logger, block=True) + + 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) + + saFilePath = os.path.join(os.path.dirname(__file__), "./" + folder) + + message = 'Applying RBAC to Service Account' + status = 'CREATING' + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig='') + self.logger.info(" Service Account File Path:" + str(saFilePath)) + + self._apply_orchestration_rbac_policies_sa(region, remote_region_oam_ip, orch_sa_namespace, namespace, saName, saFilePath, self.logger) + + cmds_token_name = [] + secret_name_found = False + cmds_token_name.append(ssh_prefix + " kubectl describe serviceaccount --kubeconfig=/etc/kubernetes/admin.conf -n " + orch_sa_namespace + " " + saName + "| grep Tokens ") + while not secret_name_found: + all_lines = self._run_commands(cmds_token_name, host_password, self.logger, block=True) + # Check if secretname != <none>; if so, repeat the command + ## kubeconfighandler.py _create_kubeconfig_orchestration 351 Secret name:<none> + secretname = self._parse_token_name(all_lines) + if secretname != "<none>": + secret_name_found = True + else: + time.sleep(60) + + message = 'Parsing token' + status = 'CREATING' + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig='') + self.logger.info(" Secret name:" + secretname) + if secretname != None: + #self.logger.info("ABC") + cmds1 = [] + cmds1.append(ssh_prefix + " kubectl describe secret --kubeconfig=/etc/kubernetes/admin.conf -n " + orch_sa_namespace + " " + secretname + " | grep token:") + all_lines = self._run_commands(cmds1, host_password, self.logger, block=True) + token = self._parse_token(all_lines) + #self.logger.info("TOKEN TO USE:" + token) + token = token.strip() + self.logger.info(" Token:[" + str(token) + "]") + + # Generate kubeconfig + kubeconfig_value = self._generate_kubeconfig(region, namespace, saName, token, transaction_id, self.logger) + + # Update database + message = 'kubeconfig creation done. ' + status = 'COMPLETE' + self.logger.info(message + ' ' + status) + self._update_status_async(region, namespace, saName, message, status, transaction_id, self.logger, kubeconfig=kubeconfig_value, done=True) + + # Send notification on VMB + self._send_vmb_notification(region, namespace, transaction_id, kubeconfig_value) + return kubeconfig_value + + def _send_vmb_notification(self, region, namespace, transaction_id, kubeconfig_value): + created_at_time = datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S") + + site_name, site_location = self.get_fuze_spm_site_details(region, self.logger) + self.logger.info("Site name:" + site_name) + self.logger.info("Site location:" + site_location) + + # Send VMB Notification + kubeconfig = Kubeconfig( + cluster=region, + namespace=namespace, + location=site_location, + kubeconfig='Example kubeconfig', + created_at=str(created_at_time), + updated_at=str(created_at_time), + transactionId=str(transaction_id) + ) + + kubeconfig.kubeconfig = kubeconfig_value + + reportDescription = 'kubeconfig file: ' + site_name + kubeconfig_message = KubeconfigMessage( + reportName='vcp_fe_kubeconfig', + reportDescription=reportDescription, + reportGeneratedOn=str(created_at_time), + rowCount = 1, + reportDataRows=[kubeconfig], + ) + + item = {} + item['payload'] = kubeconfig_message + item['message'] = 'Kubeconfig' + self.logger.info("Sending to VMB") + self.logger.info(kubeconfig_message) + self.logger.info(kubeconfig) + self.vmbQueue.put(item) + + def _generate_kubeconfig(self, region, namespace, saName, token, transaction_id, logger): + logger.info(" Inside _generate_kubeconfig") + + user_list = [] + tokendata = {} + tokendata['token'] = token + userdata = {} + userdata['name'] = saName + userdata['user'] = tokendata + user_list.append(userdata) + + context_list = [] + contextdatawrapper = {} + contextdata = {} + contextdata['cluster'] = region + contextdata['user'] = saName + contextdata['namespace'] = namespace + contextdatawrapper['context'] = contextdata + contextdatawrapper['name'] = region + context_list.append(contextdatawrapper) + + cluster_list = [] + clusterdata = {} + clusterdata['name'] = region + cluster_detail = {} + cluster_detail['insecure-skip-tls-verify'] = True + remote_region_oam_ip = self._get_remote_region_oam_ip(region, logger) + cluster_detail['server'] = 'https://[' + remote_region_oam_ip + ']:6443' + clusterdata["cluster"] = cluster_detail + cluster_list.append(clusterdata) + + outer_dict = {} + outer_dict['apiVersion'] = 'v1' + outer_dict['kind'] = 'Config' + outer_dict['current-context'] = region + outer_dict['users'] = user_list + outer_dict['contexts'] = context_list + outer_dict['clusters'] = cluster_list + + kubeconfig_json = json.dumps(outer_dict) + logger.info(str(transaction_id) + " Kubeconfig:" + kubeconfig_json) + return kubeconfig_json + + def _parse_token_name(self, all_lines): + for line in all_lines.split("\n"): + if 'Tokens' in line: + parts = line.split(":") + tokenName = parts[1].rstrip().lstrip() + return tokenName + + def _parse_token(self, all_lines): + for line in all_lines.split("\n"): + if 'token:' in line: + parts = line.split(":") + token = parts[1].rstrip().lstrip() + return token + + 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 _run_commands(self, commands, host_password, logger, block=False): + #logger.info(commands) + all_lines1 = [] + for command in commands: + logger.info(" Executing.." + str(command)) + child = pexpect.spawn(command) + try: + if block: + child.timeout=None + child.expect(['password: '], timeout=None) + child.sendline(host_password) + all_lines1 = child.read() + all_lines1 = all_lines1.rstrip().lstrip() + all_lines1 = all_lines1.decode('utf-8').replace('\r\n', '\n') + logger.info(all_lines1) + return all_lines1 + else: + child.expect(['(yes/no)? ']) + child.sendline('yes') + child.expect(['password: ']) + child.sendline(host_password) + #child.interact() + #child.close() + child.expect(pexpect.EOF, timeout=5) + except: + pass + return all_lines1 + + def _update_status_async(self, region, namespace, saName, message, status, transaction_id, logger, kubeconfig='', done=False): + item = {} + item['type'] = 'kubeconfig' + item['region'] = region + item['namespace'] = namespace + item['saName'] = saName + item['message'] = message + item['status'] = status + item['transaction_id'] = transaction_id + item['kubeconfig'] = kubeconfig + item['done'] = done + self.dbQueue.put(item) + + def _update_status(self, region, namespace, saName, message, status, transaction_id, logger, kubeconfig=''): + RemoteRegionSetup.objects.filter(transaction_id=transaction_id).update(message=message, + status=status, + serviceaccount=saName, + kubernetes_namespace=namespace, + remote_region_name=region, + kubeconfig=kubeconfig) + return + diff --git a/src/orchestration/migrations/0008_auto_20200514_1951.py b/src/orchestration/migrations/0008_auto_20200514_1951.py new file mode 100644 index 0000000..d85773d --- /dev/null +++ b/src/orchestration/migrations/0008_auto_20200514_1951.py @@ -0,0 +1,19 @@ +# Generated by Django 3.0.6 on 2020-05-15 01:51 + +from django.db import migrations + + +class Migration(migrations.Migration): + + dependencies = [ + ('orchestration', '0007_auto_20200429_1837'), + ] + + operations = [ + migrations.DeleteModel( + name='ImageSync', + ), + migrations.DeleteModel( + name='RemoteRegionSetup', + ), + ] diff --git a/src/orchestration/migrations/0008_imagesync_transaction_id.py b/src/orchestration/migrations/0008_imagesync_transaction_id.py new file mode 100644 index 0000000..af58c73 --- /dev/null +++ b/src/orchestration/migrations/0008_imagesync_transaction_id.py @@ -0,0 +1,18 @@ +# Generated by Django 3.0.5 on 2020-05-29 16:01 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('orchestration', '0007_auto_20200429_1837'), + ] + + operations = [ + migrations.AddField( + model_name='imagesync', + name='transaction_id', + field=models.CharField(default='-1', max_length=20), + ), + ] diff --git a/src/orchestration/migrations/0009_auto_20200529_1929.py b/src/orchestration/migrations/0009_auto_20200529_1929.py new file mode 100644 index 0000000..c5d5aa1 --- /dev/null +++ b/src/orchestration/migrations/0009_auto_20200529_1929.py @@ -0,0 +1,18 @@ +# Generated by Django 3.0.5 on 2020-05-29 19:29 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('orchestration', '0008_imagesync_transaction_id'), + ] + + operations = [ + migrations.AlterField( + model_name='imagesync', + name='transaction_id', + field=models.CharField(default='-1', max_length=39), + ), + ] diff --git a/src/orchestration/migrations/0010_centralregiontoremoteregionmap.py b/src/orchestration/migrations/0010_centralregiontoremoteregionmap.py new file mode 100644 index 0000000..905f71c --- /dev/null +++ b/src/orchestration/migrations/0010_centralregiontoremoteregionmap.py @@ -0,0 +1,21 @@ +# Generated by Django 3.0.5 on 2020-05-29 19:43 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('orchestration', '0009_auto_20200529_1929'), + ] + + operations = [ + migrations.CreateModel( + name='CentralRegionToRemoteRegionMap', + fields=[ + ('id', models.AutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), + ('central_region_name', models.CharField(max_length=64)), + ('remote_region_name', models.CharField(max_length=64)), + ], + ), + ] diff --git a/src/orchestration/migrations/0011_auto_20200529_1944.py b/src/orchestration/migrations/0011_auto_20200529_1944.py new file mode 100644 index 0000000..0d44d36 --- /dev/null +++ b/src/orchestration/migrations/0011_auto_20200529_1944.py @@ -0,0 +1,17 @@ +# Generated by Django 3.0.5 on 2020-05-29 19:44 + +from django.db import migrations + + +class Migration(migrations.Migration): + + dependencies = [ + ('orchestration', '0010_centralregiontoremoteregionmap'), + ] + + operations = [ + migrations.RenameModel( + old_name='CentralRegionToRemoteRegionMap', + new_name='CentralToRemoteMap', + ), + ] diff --git a/src/orchestration/migrations/0012_merge_20200702_0053.py b/src/orchestration/migrations/0012_merge_20200702_0053.py new file mode 100644 index 0000000..f31df44 --- /dev/null +++ b/src/orchestration/migrations/0012_merge_20200702_0053.py @@ -0,0 +1,14 @@ +# Generated by Django 3.0.8 on 2020-07-02 00:53 + +from django.db import migrations + + +class Migration(migrations.Migration): + + dependencies = [ + ('orchestration', '0008_auto_20200514_1951'), + ('orchestration', '0011_auto_20200529_1944'), + ] + + operations = [ + ] diff --git a/src/orchestration/migrations/0013_imagesync_remoteregionsetup.py b/src/orchestration/migrations/0013_imagesync_remoteregionsetup.py new file mode 100644 index 0000000..987851e --- /dev/null +++ b/src/orchestration/migrations/0013_imagesync_remoteregionsetup.py @@ -0,0 +1,39 @@ +# Generated by Django 3.0.8 on 2020-07-02 03:34 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('orchestration', '0012_merge_20200702_0053'), + ] + + operations = [ + migrations.CreateModel( + name='ImageSync', + fields=[ + ('id', models.AutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), + ('transaction_id', models.CharField(default='-1', max_length=39)), + ('remote_region_name', models.CharField(max_length=64)), + ('docker_image', models.TextField(max_length=65535)), + ('upload_start_time', models.DateTimeField()), + ('upload_end_time', models.DateTimeField(null=True)), + ('upload_status', models.CharField(choices=[('STARTED', 'STARTED'), ('UPLOADING', 'UPLOADING'), ('COMPLETE', 'COMPLETE'), ('FAILED', 'FAILED')], max_length=9)), + ('upload_message', models.CharField(max_length=255)), + ], + ), + migrations.CreateModel( + name='RemoteRegionSetup', + fields=[ + ('id', models.AutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), + ('caas_vendor_version', models.CharField(max_length=20)), + ('kubernetes_version', models.CharField(max_length=20)), + ('central_region_name', models.CharField(max_length=64)), + ('remote_region_name', models.CharField(max_length=64)), + ('kubernetes_namespace', models.CharField(max_length=64)), + ('serviceaccount', models.CharField(max_length=64)), + ('kubeconfig', models.TextField(blank=True, max_length=16777215, null=True)), + ], + ), + ] diff --git a/src/orchestration/migrations/0014_remoteregionsetup_transaction_id.py b/src/orchestration/migrations/0014_remoteregionsetup_transaction_id.py new file mode 100644 index 0000000..2555269 --- /dev/null +++ b/src/orchestration/migrations/0014_remoteregionsetup_transaction_id.py @@ -0,0 +1,18 @@ +# Generated by Django 3.0.8 on 2020-07-07 20:11 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('orchestration', '0013_imagesync_remoteregionsetup'), + ] + + operations = [ + migrations.AddField( + model_name='remoteregionsetup', + name='transaction_id', + field=models.CharField(default='-1', max_length=39), + ), + ] diff --git a/src/orchestration/migrations/0015_auto_20200707_2034.py b/src/orchestration/migrations/0015_auto_20200707_2034.py new file mode 100644 index 0000000..7288f9a --- /dev/null +++ b/src/orchestration/migrations/0015_auto_20200707_2034.py @@ -0,0 +1,23 @@ +# Generated by Django 3.0.8 on 2020-07-07 20:34 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('orchestration', '0014_remoteregionsetup_transaction_id'), + ] + + operations = [ + migrations.AddField( + model_name='remoteregionsetup', + name='message', + field=models.CharField(max_length=255, null=True), + ), + migrations.AddField( + model_name='remoteregionsetup', + name='status', + field=models.CharField(choices=[('STARTED', 'STARTED'), ('CREATING', 'CREATING'), ('COMPLETE', 'COMPLETE'), ('FAILED', 'FAILED')], max_length=9, null=True), + ), + ] diff --git a/src/orchestration/models.py b/src/orchestration/models.py index 8e79d8d..fac240c 100644 --- a/src/orchestration/models.py +++ b/src/orchestration/models.py @@ -2,6 +2,7 @@ from django.db import models class ImageSync(models.Model): + transaction_id = models.CharField(max_length=39,default="-1") remote_region_name = models.CharField(max_length=64) docker_image = models.TextField(max_length=65535) upload_start_time = models.DateTimeField() @@ -10,10 +11,11 @@ class ImageSync(models.Model): upload_message = models.CharField(max_length=255) def __str__(self): - return "%s : %s" % (self.remote_region_name, self.docker_image) + return "%s : %s : %s" % (self.transaction_id, self.remote_region_name, self.docker_image) class RemoteRegionSetup(models.Model): + transaction_id = models.CharField(max_length=39,default="-1") caas_vendor_version = models.CharField(max_length=20) kubernetes_version = models.CharField(max_length=20) central_region_name = models.CharField(max_length=64) @@ -21,6 +23,13 @@ class RemoteRegionSetup(models.Model): kubernetes_namespace = models.CharField(max_length=64) serviceaccount = models.CharField(max_length=64) kubeconfig = models.TextField(max_length=16777215, null=True, blank=True) + status = models.CharField(max_length=9, null=True, choices=[('STARTED', 'STARTED'), ('CREATING', 'CREATING'), ('COMPLETE', 'COMPLETE'), ('FAILED', 'FAILED')]) + message = models.CharField(max_length=255, null=True) def __str__(self): return "%s : %s" % (self.remote_region_name, self.kubernetes_namespace) + + +class CentralToRemoteMap(models.Model): + central_region_name = models.CharField(max_length=64) + remote_region_name = models.CharField(max_length=64) 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 + diff --git a/src/orchestration/remoteregionhandler.py b/src/orchestration/remoteregionhandler.py new file mode 100644 index 0000000..6754672 --- /dev/null +++ b/src/orchestration/remoteregionhandler.py @@ -0,0 +1,375 @@ +import django +import os +import threading + +from django.conf import settings +from django.core.exceptions import AppRegistryNotReady +from django.utils import timezone + +from .utils import * + +try: + django.setup() + from .models import ImageSync + from caas.models import Cluster + from caas.models import Location + from caas.models import Namespace +except django.core.exceptions.AppRegistryNotReady as exp: + pass + +class RemoteRegionWorker: + + def __init__(self): + pass + + def _get_crd_details(self): + crd_cluster_role = "nad" + crd_api_group = "k8s.cni.cncf.io" + crd_api_resources = "network-attachment-definitions" + crd_api_verbs = "*" + return crd_cluster_role, crd_api_group, crd_api_resources, crd_api_verbs + + def _read_remote_openrc(self, remote_region_oam_ip, logger): + #cmd = "scp @[fd00:4888:2000:120c::290]:/etc/platform/openrc . + + logger.info("Inside _read_remote_openrc") + host_username = settings.HOST_CREDS['username'] + host_password = settings.HOST_CREDS['password'] + ssh_prefix = 'ssh -o StrictHostKeyChecking=no -t ' + host_username + '@' + remote_region_oam_ip + env_string = "" + cmds_scp = [] + folder = get_temp_file_location() + logger.info("Tmp file location:" + folder) + #tmpFilePath = os.path.join(os.path.dirname(__file__), "./" + folder) + tmpFileName = 'openrc_' + str(remote_region_oam_ip) + cmds_scp.append('scp -o StrictHostKeyChecking=no ' + host_username + '@[' + remote_region_oam_ip + ']:/etc/platform/openrc ' + folder + '/' + tmpFileName) + logger.info(cmds_scp) + run_commands_scp(cmds_scp, host_password, logger, block=True) + logger.info("About to create env string") + + fp = open(folder + '/' + tmpFileName) + lines = fp.readlines() + logger.info(lines) + password_line = '' + for line in lines: + line = line.lstrip().rstrip() + logger.info(line) + parts = line.split(' ') + if len(parts) == 2: + if parts[0] == 'export': + if parts[1]: + env_var = parts[1].lstrip().rstrip() + logger.info(env_var) + if 'OS_PASSWORD' not in env_var: + env_string = env_string + " " + env_var + if len(parts) >= 2 and 'OS_PASSWORD=' in parts[1]: + logger.info("Parsing password") + password_parts = line.split('PASSWORD=') + password_command_parts = password_parts[1].split(' ') + password_line = password_command_parts[1].rstrip().lstrip() + logger.info("Password line:" + password_line) + + logger.info("Deleting openrc file ") + + # Delete file + if os.path.exists(folder + '/' + tmpFileName): + os.remove(folder + '/' + tmpFileName) + + logger.info("Looking for password") + # Get password + password_cmd = [] + #cmd = 'TERM=linux /opt/platform/.keyring/20.06/.CREDENTIAL 2>/dev/null' + cmd_to_run = ssh_prefix + " " + password_line + logger.info(cmd_to_run) + successful, value = self._run_command_get_output([cmd_to_run], host_password, logger) + password_val = host_password + #if successful: + # if value != '': + # password_val = value + os_password = "OS_PASSWORD=" + password_val + env_string = env_string + " " + os_password + + logger.info("Env string:" + env_string) + return env_string + + def _run_command_get_all_lines(self, commands, host_password, logger, block=False, timeout=None): + logger.info("Inside _run_command_get_all_lines") + all_lines = [] + for command in commands: + logger.info(" Executing.." + str(command)) + child = pexpect.spawn(command) + child.timeout=timeout + try: + i = child.expect(['password: ','Connection refused\r\r\n'], timeout=timeout) + if i == 0: + child.sendline(host_password) + all_lines = child.read() + all_lines = all_lines.rstrip().lstrip() + all_lines = all_lines.decode('utf-8').replace('\r\n', '\n') + logger.info(all_lines) + if i == 1: + logger.info("Connection refused") + except: + logger.info(str(child)) + return all_lines + + def _run_command_get_output(self, commands, host_password, logger, block=False, timeout=None): + successful = False + value_to_return = '' + logger.info("Inside _run_command_get_output") + for command in commands: + logger.info(" Executing.." + str(command)) + child = pexpect.spawn(command) + child.timeout=timeout + try: + child.expect(['password: '], timeout=timeout) + child.sendline(host_password) + all_lines = child.read() + all_lines = all_lines.rstrip().lstrip() + all_lines = all_lines.decode('utf-8').replace('\r\n', '\n') + successful = True # Tentative + for line in all_lines.split("\n"): + logger.info(line) + if value_to_return == '': + value_to_return = line.strip() + if re.search('error', line, re.IGNORECASE): + successful = False + if re.search('unable', line, re.IGNORECASE): + successful = False + except: + logger.info(str(child)) + logger.info("Status:" + str(successful) + " value_to_return:" + value_to_return) + return successful, value_to_return + + def _get_remote_region_oam_ip(self, cluster_name, logger): + remoteclusterObj = Cluster.objects.filter(cluster_name=cluster_name) + oam_ip = remoteclusterObj[0].oam_vip_address + logger.info(" Remote region:" + cluster_name + " OAM IP:" + str(oam_ip)) + return oam_ip + + def _get_central_region_name(self, oam_vip_address, logger): + logger.info("1") + remoteclusterObj = Cluster.objects.filter(oam_vip_address=oam_vip_address) + logger.info("2") + logger.info(remoteclusterObj) + cluster_name = remoteclusterObj[0].cluster_name + logger.info(" Remote region:" + str(oam_vip_address) + " Cluster Name:" + str(cluster_name)) + return cluster_name + + def _get_namespace(self, region, logger): + # Lookup database and findout namespace given region + logger.info(" Inside _get_namespace") + remoteclusterObj = Cluster.objects.filter(cluster_name=region) + if len(remoteclusterObj) > 0: + logger.info(" RemoteClusterObj:" + str(remoteclusterObj)) + namespace_id = remoteclusterObj[0].namespace_id + logger.info(" Namespace id:" + str(namespace_id)) + namespaceObj = Namespace.objects.filter(id=namespace_id) + logger.info(" NamespaceObj:" + str(namespaceObj)) + namespace_name = namespaceObj[0].namespace_name + logger.info(" Remote region:" + region + " Namespace:" + namespace_name) + return namespace_name + else: + return "" + + def _get_central_region_oam_ip(self, remoteregion, transaction_id): + remoteclusterObj = Cluster.objects.filter(cluster_name=remoteregion) + parent_cluster_id = remoteclusterObj[0].parent_cluster_id + self.logger.info(str(transaction_id) + " Parent cluster id.." + str(parent_cluster_id)) + centralclusterObj = Cluster.objects.filter(id=parent_cluster_id) + self.logger.info(str(transaction_id) + " " + str(centralclusterObj)) + central_region_list = [] + for central_region in centralclusterObj: + central_region_list.append(central_region.oam_vip_address) + self.logger.info(str(transaction_id) + " Central Region List:" + ','.join(central_region_list)) + return central_region_list + + def get_fuze_spm_site_details(self, cluster_name, logger): + fuze_spm_site_name = '' + fuze_spm_site_id = '' + logger.info(" Inside _get_fuze_spm_site_name " + str(cluster_name)) + remoteclusterObj = Cluster.objects.filter(cluster_name=cluster_name) + if len(remoteclusterObj) > 0: + fuze_id = remoteclusterObj[0].location_id + locationObj = Location.objects.filter(id=fuze_id) + if len(locationObj) > 0: + fuze_spm_site_name = locationObj[0].fuze_spm_site_name + fuze_spm_site_id = locationObj[0].fuze_spm_site_id + logger.info(" Fuze site name:" + fuze_spm_site_name) + logger.info(" Fuze site id:" + fuze_spm_site_id) + return fuze_spm_site_name, fuze_spm_site_id + + def check_namespaces(self, remote_region_oam_ip, logger): + + 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) + + new_op_lines = [] + for line in output_lines.split("\n"): + if not 'Connection to' in line: + new_op_lines.append(line) + + return new_op_lines + + def check_namespace_secrets(self, region, remote_region_oam_ip, logger): + namespace = self._get_namespace(region, logger) + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + cmds = [] + cmds.append(ssh_prefix + " kubectl get secrets --kubeconfig=/etc/kubernetes/admin.conf -n " + namespace) + + output_lines = self._run_command_get_all_lines(cmds, host_password, logger, block=True) + #logger.info("Returned output lines:") + #logger.info(output_lines) + + new_op_lines = [] + for line in output_lines.split("\n"): + if not 'Connection to' in line: + new_op_lines.append(line) + + return new_op_lines + + def check_namespace_serviceaccounts(self, region, remote_region_oam_ip, logger): + namespace = self._get_namespace(region, logger) + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + cmds = [] + cmds.append(ssh_prefix + " kubectl get serviceaccounts --kubeconfig=/etc/kubernetes/admin.conf -n " + namespace) + + output_lines = self._run_command_get_all_lines(cmds, host_password, logger, block=True) + #logger.info("Returned output lines:") + #logger.info(output_lines) + + new_op_lines = [] + for line in output_lines.split("\n"): + if not 'Connection to' in line: + new_op_lines.append(line) + + return new_op_lines + + def check_online_status(self, region, remote_region_oam_ip, logger): + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + central_region_list = self._get_central_region_oam_ip(region, "-1") + region_online_status = [] + new_op_lines = [] + for central_region_ip in central_region_list: + cmds = [] + cmd = 'ssh -o StrictHostKeyChecking=no -t ' + host_username + '@' + central_region_ip + ' ' + '"source /etc/platform/openrc; dcmanager subcloud list | grep ' + region + '"' + cmds.append(cmd) + output_lines = self._run_command_get_all_lines(cmds, host_password, logger, block=True) + logger.info(output_lines) + + for line in output_lines.split("\n"): + if not 'Connection to' in line: + new_op_lines.append(line) + return new_op_lines + + def create_host_network(self, remote_region_oam_ip, logger): + logger.info("Inside _create_host_network") + + #host_username = settings.HOST_CREDS['username'] + #host_password = settings.HOST_CREDS['password'] + #ssh_prefix = 'ssh -o StrictHostKeyChecking=no -t ' + host_username + '@' + remote_region_oam_ip + + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + env_string = self._read_remote_openrc(remote_region_oam_ip, logger) + + #basecmd = " OS_ENDPOINT_TYPE=internalURL CINDER_ENDPOINT_TYPE=internalURL OS_USERNAME=admin" + #basecmd = basecmd + " OS_PASSWORD=`TERM=linux /opt/platform/.keyring/20.06/.CREDENTIAL 2>/dev/null`" + #basecmd = basecmd + " OS_AUTH_TYPE=password OS_AUTH_URL=http://[fd00:4888:2000:120b::220]:5000/v3" + #basecmd = basecmd + " OS_PROJECT_NAME=admin OS_USER_DOMAIN_NAME=Default OS_PROJECT_DOMAIN_NAME=Default" + #basecmd = basecmd + " OS_IDENTITY_API_VERSION=3 OS_REGION_NAME=subcloud2 OS_INTERFACE=internal" + + basecmd = env_string + + cmds = ["system host-lock controller-0", + "system host-if-modify -n f1u -c pci-sriov --num-vfs 8 controller-0 ens3f1 --vf-driver=vfio", + "system host-if-add -c pci-sriov controller-0 f1c vf f1u --num-vfs 4 --vf-driver=netdevice", + "system host-if-modify controller-0 f1u --imtu=1956", + "system host-if-modify controller-0 f1c --imtu=1956", + "system datanetwork-add f1u vlan --mtu=1956", + "system datanetwork-add f1c vlan --mtu=1956", + "system interface-datanetwork-assign controller-0 f1u f1u", + "system interface-datanetwork-assign controller-0 f1c f1c", + "system host-if-add -c pci-sriov controller-0 fh0m vf fh0 --num-vfs 4 --vf-driver=netdevice", + "system host-if-modify controller-0 fh0m --imtu=9000", + "system datanetwork-add fh0m flat", + "system interface-datanetwork-assign controller-0 fh0m fh0m", + "system host-if-modify -n fh1 -c pci-sriov --num-vfs 8 controller-0 enp181s0f0 --vf-driver=vfio", + "system host-if-modify controller-0 fh1 --imtu=9000", + "system datanetwork-add fh1 vlan --mtu=9000", + "system interface-datanetwork-assign controller-0 fh1 fh1"] + + # New steps proposed by Eddy - These do not seem to create fh0m so commenting out. + #cmds = ["system host-lock controller-0", + # "system host-if-modify -n f1c -c pci-sriov --num-vfs 8 controller-0 ens3f1 --vf-driver=netdevice", + # "system host-if-add -c pci-sriov controller-0 f1u vf f1c --num-vfs 4 --vf-driver=vfio", + # "system host-if-modify controller-0 f1u --imtu=1956", + # "system host-if-modify controller-0 f1c --imtu=1956", + # "system datanetwork-add f1u vlan --mtu=1956", + # "system datanetwork-add f1c vlan --mtu=1956", + # "system interface-datanetwork-assign controller-0 f1u f1u", + # "system interface-datanetwork-assign controller-0 f1c f1c", + # # this one is wrong as well but for some reason I think the fh0 is setup in Carlos deployment config. So we have to delete the fh0 and recreate both. + # "system host-if-modify controller0 fh0 -nc none" + # "system host-if-modify -n fh0m -c pci-sriov --num-vfs 8 controller-0 ens179s0f0 --vf-driver=netdevice", + # "system host-if-add -c pci-sriov controller-0 fh0 vf fh0m --num-vfs 4 --vf-driver=vfio", + # "system host-if-modify controller-0 fh0m --imtu=9000", + # "system datanetwork-add fh0m flat", + # "system interface-datanetwork-assign controller-0 fh0m fh0m", + # "system host-if-modify -n fh1 -c pci-sriov --num-vfs 8 controller-0 enp181s0f0 --vf-driver=vfio", + # "system host-if-modify controller-0 fh1 --imtu=9000", + # "system datanetwork-add fh1 vlan --mtu=9000", + # "system interface-datanetwork-assign controller-0 fh1 fh1"] + + for cmd in cmds: + cmd_to_run = ssh_prefix + " " + basecmd + " " + cmd + logger.info(cmd_to_run) + run_commands([cmd_to_run], host_password, logger, block=True) + +# host_network_configured = False +# while not host_network_configured: +# host_network_configured = self._verify_host_network(ssh_prefix, basecmd, host_password, logger) +# if not host_network_configured: +# cmds = ["system host-lock controller-0"] +# for cmd in cmds: +# cmd_to_run = ssh_prefix + " " + basecmd + " " + cmd +# logger.info(cmd_to_run) +# output_lines = self._run_command_get_all_lines([cmd_to_run], host_password, logger, block=True) +# logger.info("Returned output lines:") +# logger.info(output_lines) +# time.sleep(3) + + cmds = ["system host-unlock controller-0"] + while True: + unlock_wait = False + for cmd in cmds: + cmd_to_run = ssh_prefix + " " + basecmd + " " + cmd + logger.info(cmd_to_run) + output_lines = self._run_command_get_all_lines([cmd_to_run], host_password, logger, block=True) + 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 not unlock_wait: + if 'retry host-unlock' in line or 'Rejected' in line: + logger.info("Need to wait to call host-unlock ##### ") + unlock_wait = True + break + if unlock_wait: + time.sleep(60) + else: + break + +# for cmd in cmds: +# cmd_to_run = ssh_prefix + " " + basecmd + " " + cmd +# logger.info(cmd_to_run) +# output_lines = self._run_command_get_all_lines([cmd_to_run], host_password, logger, block=True) +# logger.info("Returned output lines:") +# logger.info(output_lines) + logger.info("Done setting up host network") diff --git a/src/orchestration/tests/cluster_status.json b/src/orchestration/tests/cluster_status.json new file mode 100644 index 0000000..db2821f --- /dev/null +++ b/src/orchestration/tests/cluster_status.json @@ -0,0 +1,15 @@ +{ + "reportName": "vcp_fe_cluster_status", + "reportDescription": "caas deployment status", + "reportGeneratedOn": "2020-06-28T09:06:32.1962088-04:00", + "reportDataRows": [{ + "name": "waeomagj-d654321-001", + "description": "NE CONCORD 8_NH", + "location": "654322", + "software_version": "19.12", + "availability": "online", + "deploy_status": "complete", + "created_at": "2020-06-17 04:16:15.743617", + "updated_at": "2020-06-17 06:03:10.854598" + }] +}
\ No newline at end of file diff --git a/src/orchestration/tests/image_status.json b/src/orchestration/tests/image_status.json new file mode 100644 index 0000000..adc1714 --- /dev/null +++ b/src/orchestration/tests/image_status.json @@ -0,0 +1,22 @@ +{ + "reportName": "vcp_fe_imagestatus", + "reportDescription": null, + "reportGeneratedOn": "2020-04-28T09:06:32.1962088-04:00", + "rowCount": 2, + "reportDataRows": [{ + "cluster": "wsbomagj-d654321-001", + "image": "i1", + "action": "UPLOAD", + "status": "SUCCESS", + "message":"Image Successfully Deleted", + "created_at": "2020-01-07 04:16:15.743617" + }, + { + "cluster": "wsbomagj-d654321-001", + "image": "i2", + "action": "UPLOAD", + "status": "FAILED", + "message":"Image Not Present in Artifactory", + "created_at": "2020-01-07 04:16:15.743617" + }] +}
\ No newline at end of file diff --git a/src/orchestration/tests/kubeconfig.yaml b/src/orchestration/tests/kubeconfig.yaml new file mode 100644 index 0000000..eba9a96 --- /dev/null +++ b/src/orchestration/tests/kubeconfig.yaml @@ -0,0 +1,17 @@ +apiVersion: v1 +kind: Config +users: +- name: ldap-user + user: + token: eyJhbGciOiJSUzI1NiIsImtpZCI6IiJ9.eyJpc3MiOiJrdWJlcm5ldGVzL3NlcnZpY2VhY2NvdW50Iiwia3ViZXJuZXRlcy5pby9zZXJ2aWNlYWNjb3VudC9uYW1lc3BhY2UiOiJub2tpYSIsImt1YmVybmV0ZXMuaW8vc2VydmljZWFjY291bnQvc2VjcmV0Lm5hbWUiOiJzYTEtdG9rZW4tbW5mMmoiLCJrdWJlcm5ldGVzLmlvL3NlcnZpY2VhY2NvdW50L3NlcnZpY2UtYWNjb3VudC5uYW1lIjoic2ExIiwia3ViZXJuZXRlcy5pby9zZXJ2aWNlYWNjb3VudC9zZXJ2aWNlLWFjY291bnQudWlkIjoiZTA5YTkwMDItMTg0Zi0xMWVhLWE5NzYtMDgwMDI3YTFjODc3Iiwic3ViIjoic3lzdGVtOnNlcnZpY2VhY2NvdW50Om5va2lhOnNhMSJ9.e9uD5kMmAT0HRboSTAbH5xlkETkltLclVQ2GedvoeUmH76WB6G5kGWQrhJjkjpMtPDKxWp6wTzZdEXXwGYGdB6aXbaxcAmau1qid5NGz725BtaoRbSVS2Uk6XrOSNfycFzqc8Z7GTX81VtKKSPYnjeMo47W6FqHw6qk0NEpLLxbpGfJHz8w2KZQiuvI-JRQXtA3PHW3tEaWq3ME3XnYgHNSRJmoiKA99bWjN-HKoOsBDMjhX7kw_VtycRYJ1gbqXmVsSl7BuAQjoplnLN_stRt7ZgpsV4aqmZueqJqHaglB91XOeqe_JcD5unLMb3B5VXpEXPp3V6tjLIyVD8F8XbA +clusters: +- cluster: + server: https://[fd00:4888:2000:120e::290]:8443 + name: wsbomagj-d654321-001 +contexts: +- context: + cluster: wsbomagj-d654321-001 + user: ldap-user + namespace: WSBOMAGJ-441352VZWcVDU-Y-SM-x-001 + name: ldap-user +current-context: ldab-user diff --git a/src/orchestration/tests/namespace.json b/src/orchestration/tests/namespace.json new file mode 100644 index 0000000..be11eb6 --- /dev/null +++ b/src/orchestration/tests/namespace.json @@ -0,0 +1,12 @@ +{ + "reportName": "vcp_fe_namespace", + "reportDescription": null, + "reportGeneratedOn": "2020-04-28T09:06:32.1962088-04:00", + "rowCount": 1, + "reportDataRows": [{ + "cluster": "wsbomagj-d654321-001", + "namespace": "WSBOMAGJ-441352VZWcVDU-Y-SM-x-001", + "location": "654321", + "created_at": "2020-01-07 04:16:15.743617" + }] +}
\ No newline at end of file diff --git a/src/orchestration/urls.py b/src/orchestration/urls.py index 6b7dbbc..78c6d3e 100644 --- a/src/orchestration/urls.py +++ b/src/orchestration/urls.py @@ -4,5 +4,15 @@ from . import views urlpatterns = [ path('', views.index, name='index'), path('count/<int:count>/', views.count, name='count'), - path('imagesync/', views.imagesync, name='imagesync') + path('image-upload', views.imagesync, name='images'), + path('image-upload/<str:transaction_id>', views.get_status_by_transaction_id, name='get_status_by_transaction_id'), + path('remote-regions/<str:remote_region>/status', views.get_status_by_remote_region, name='get_status_by_remote_region'), + path('remote-regions/<str:remote_region>/images/<str:image_name>/status', views.get_status_by_image_name_and_remote_region, name='get_status_by_image_name_and_remote_region'), + path('remote-regions/<str:remote_region>/images/<str:image_name>/tags/<str:image_tag>', views.delete_image_tag_remote_region, name='delete_image_tag_remote_region'), + path('caas-status', views.caas_status, name='caas-status'), + path('setup-cluster/<str:cluster_name>', views.setup_cluster, name='setup_cluster'), + path('remote-regions/<str:remote_region>/site-details', views.get_setup_details_by_remote_region, name='get_setup_details_by_remote_region'), + path('remote-regions/<str:remote_region>/setup-network', views.setup_network_for_remote_region, name='setup_network_for_remote_region'), + path('remote-regions/<str:remote_region>/connection-details', views.get_kubeconfig_by_remote_region, name='get_kubeconfig_by_remote_region'), + path('remote-regions/<str:remote_region>/connection-details/<str:transaction_id>', views.get_kubeconfig_by_transaction_id, name='get_kubeconfig_by_transaction_id'), ] diff --git a/src/orchestration/utils.py b/src/orchestration/utils.py new file mode 100644 index 0000000..24f6b9c --- /dev/null +++ b/src/orchestration/utils.py @@ -0,0 +1,110 @@ +import pexpect +import re +import django +from django.core.exceptions import AppRegistryNotReady + +from django.conf import settings + +try: + django.setup() + from caas.models import Cluster + from caas.models import Namespace +except django.core.exceptions.AppRegistryNotReady as exp: + pass + +process_started = False + +def set_process_started(): + global process_started + if not process_started: + process_started = True + return process_started + +def get_pairs(imageList, remoteRegions): + pairList = [] + for image in imageList: + for region in remoteRegions: + pair = {"image": image, "region": region} + pairList.append(pair) + return pairList + +def get_docker_reg_connection_details(): + dr_username = settings.CENTRAL_DR_CREDS['username'] + dr_password = settings.CENTRAL_DR_CREDS['password'] + return dr_username, dr_password + +def get_temp_file_location(): + file_loc = settings.ORCH_TEMP_FILE_LOCATION + return file_loc + +def get_host_connection_details(region_oam_ip): + host_username = settings.HOST_CREDS['username'] + host_password = settings.HOST_CREDS['password'] + ssh_prefix = 'ssh -o StrictHostKeyChecking=no -t ' + host_username + '@' + region_oam_ip + return host_username, host_password, ssh_prefix + +def run_commands(commands, host_password, logger, block=False, timeout=None): + successful = False + logger.info("Inside run_commands") + for command in commands: + logger.info(" Executing.." + str(command)) + child = pexpect.spawn(command) + child.timeout=timeout + try: + child.expect(['password: '], timeout=timeout) + child.sendline(host_password) + all_lines = child.read() + all_lines = all_lines.rstrip().lstrip() + all_lines = all_lines.decode('utf-8').replace('\r\n', '\n') + successful = True # Tentative + for line in all_lines.split("\n"): + logger.info(line) + if re.search('error', line, re.IGNORECASE): + successful = False + if re.search('unable', line, re.IGNORECASE): + successful = False + except: + logger.info(str(child)) + logger.info("Status:" + str(successful)) + return successful + +def run_commands_scp(commands, host_password, logger, block=False, timeout=settings.SCP_TIMEOUT): + successful = False + logger.info("Inside run_commands") + logger.info("settings.RUNSERVER:" + str(settings.RUNSERVER)) + for command in commands: + logger.info(" Executing.." + str(command)) + try: + child = pexpect.spawn(str(command)) + except: + logger.info(str(child)) + #child.timeout=timeout + child.timeout = None + try: + #child.expect(['(yes/no)? '], timeout=timeout) + #child.expect(['Are you sure you want to continue connecting (yes/no)? ']) + #child.expect(['.+'], timeout=timeout) + #child.expect([':b5:9c:21:14:6f:fa:45:20:c7:ba:49:3a:aa:31.\r\nAre you sure you want to continue connecting (yes/no)? '], timeout=timeout) + if not settings.RUNSERVER: + i = child.expect([b'', 'Are you .*(yes/no)? ', pexpect.EOF], timeout=None) + if i == 0: + child.sendline('') + if i == 1: + child.sendline('yes') + child.expect(['password: '], timeout=timeout) + child.sendline(host_password) + all_lines = child.read() + all_lines = all_lines.rstrip().lstrip() + all_lines = all_lines.decode('utf-8').replace('\r\n', '\n') + successful = True # Tentative + for line in all_lines.split("\n"): + logger.info(line) + if re.search('error', line, re.IGNORECASE): + successful = False + if re.search('unable', line, re.IGNORECASE): + successful = False + except: + logger.info(str(child)) + logger.info("Status:" + str(successful)) + return successful + diff --git a/src/orchestration/vendorhandler.py b/src/orchestration/vendorhandler.py new file mode 100644 index 0000000..0285441 --- /dev/null +++ b/src/orchestration/vendorhandler.py @@ -0,0 +1,308 @@ +from django.conf import settings + +import os +import pexpect +import re +import time +import json + +from .utils import * +from .remoteregionhandler import RemoteRegionWorker + +class Vendor(RemoteRegionWorker): + + def __init__(self, logger, dbCoordinationQueue, vmbCoordinationQueue, doneQueue): + self.logger = logger + self.doneQueue = doneQueue + self.dbCoordinationQueue = dbCoordinationQueue + self.vmbCoordinationQueue = vmbCoordinationQueue + self.logger.info(".. created Vendor handler") + + def perform_vendor_setup(self, data): + self.logger.info("Checking if any vendor setup needs to be done") + if 'namespace' in data: + namespace = data['namespace'] + region_oam_ip = data['region_oam_ip'] + network_setup = data['network_setup'] + + vendor = self._get_vendor(namespace, self.logger) + self.logger.info("Vendor:" + vendor) + self.logger.info("Network setup:" + network_setup) + self._perform_vendor_specific_actions(vendor, namespace, region_oam_ip, network_setup, self.logger) + self.logger.info("Done configuring vendor related things on the remote region") + else: + self.logger.info("Vendor setup not needed for this call.") + self._wait_and_done() + + def _wait_and_done(self): + if self.dbCoordinationQueue != '' and self.vmbCoordinationQueue != '' and self.doneQueue != '': + self.logger.info("About to be done..waiting for cleanup") + db_coordination = self.dbCoordinationQueue.get() + self.logger.info(" DB coordination message:" + db_coordination) + + vmb_coordination = self.vmbCoordinationQueue.get() + self.logger.info(" VMB coordination message:" + vmb_coordination) + self.doneQueue.put("Done") + self.logger.info("Done") + + def _get_vendor(self, namespace, logger): + logger.info(" Inside _get_vendor") + parts = namespace.split('-') + logger.info("parts:" + str(parts)) + vendor = "unknown" + if len(parts) >= 5: + if parts[3] == "ss" or parts[4] == "ss": + #vendor_shortform = parts[3] + #if vendor_shortform == "ss": + vendor = "Samsung" + return vendor + + def _perform_vendor_specific_actions(self, vendor, namespace, region_oam_ip, network_setup, logger): + logger.info(" Inside _perform_vendor_specific_actions...") + + if vendor == "Samsung": + logger.info(" Handling Samsung...") + samsung_provisioner = Samsung(logger) + samsung_provisioner.setup(namespace, region_oam_ip, network_setup) + + def setup_network(self, region, logger): + logger.info(" Setting up network...") + remote_region_oam_ip = self._get_remote_region_oam_ip(region, logger) + logger.info(" Received OAM IP: " + region + " " + remote_region_oam_ip) + self.create_host_network(remote_region_oam_ip, logger) + + def check_site(self, region, logger): + logger.info(" Checking network...") + remote_region_oam_ip = self._get_remote_region_oam_ip(region, logger) + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + cmds = [] + cmds.append(ssh_prefix + " kubectl describe nodes controller-0 --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) + new_op_lines = [] + if len(output_lines) > 0: + for line in output_lines.split("\n"): + if 'Connection to' not in line: + if 'Capacity:' in line or 'Allocatable:' in line or 'Allocated' in line or 'intel.com/pci_sriov_net' in line: + logger.info("LINE:" + line) + new_op_lines.append(line) + + namespaces = self.check_namespaces(remote_region_oam_ip, logger) + + #new_op_lines.append("----------") + #for namespace_line in namespaces.split("\n"): + # new_op_lines.append(namespace_line) + + namespace_secrets = self.check_namespace_secrets(region, remote_region_oam_ip, logger) + + namespace_service_accounts = self.check_namespace_serviceaccounts(region, remote_region_oam_ip, logger) + + online_status = self.check_online_status(region, remote_region_oam_ip, logger) + + return new_op_lines, namespaces, namespace_secrets, namespace_service_accounts, online_status + +class Samsung(RemoteRegionWorker): + + def __init__(self, logger): + self.logger = logger + self.logger.info(".. created Samsung handler") + + def setup(self, namespace, remote_region_oam_ip, network_setup): + self.logger.info(" Inside Samsung setup") + self._create_docker_reg_secret(namespace, remote_region_oam_ip, self.logger) + self._create_serviceaccount(namespace, remote_region_oam_ip, self.logger) + self._add_pac_crd_annotation(remote_region_oam_ip, self.logger) + self.logger.info(" Network setup:" + network_setup) + if network_setup == 'true': + self.create_host_network(remote_region_oam_ip, self.logger) + + def _add_pac_crd_annotation(self, remote_region_oam_ip, logger): + logger.info("Inside _add_pac_crd_annotation") + + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + cmd = " kubectl annotate --kubeconfig=/etc/kubernetes/admin.conf --overwrite crd network-attachment-definitions.k8s.cni.cncf.io " + cmd = cmd + " resource/annotation-relationship=\"on:Pod,key:k8s.v1.cni.cncf.io/networks,value:[{name:INSTANCE.metadata.name}]\"" + cmd = ssh_prefix + cmd + logger.info("Annotation cmd:" + cmd) + + cmds_annotate = [] + cmds_annotate.append(cmd) + run_commands(cmds_annotate, host_password, logger, block=True) + + def _create_host_network1(self, remote_region_oam_ip, logger): + logger.info("Inside _create_host_network") + + #host_username = settings.HOST_CREDS['username'] + #host_password = settings.HOST_CREDS['password'] + #ssh_prefix = 'ssh -o StrictHostKeyChecking=no -t ' + host_username + '@' + remote_region_oam_ip + + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + cmds = [] + cmds.append(ssh_prefix + " kubectl get nodes controller-0 --kubeconfig=/etc/kubernetes/admin.conf -o json") + + output_lines = self._run_command_get_all_lines(cmds, host_password, logger, block=True) + logger.info("Returned output lines:") + logger.info(output_lines) + new_op_lines = "" + for line in output_lines.split("\n"): + if not 'Connection to' in line: + new_op_lines = new_op_lines + line + "\n" + + allocatable_found = False + capacity_found = False + hugepg1G = "hugepages-1Gi" + hugepg2M = "hugepages-2Mi" + logger.info("New o/p lines:" + new_op_lines) + if len(new_op_lines) > 0: + json_op = json.loads(new_op_lines) + logger.info("JSON O/P:" + str(json_op)) + status = json_op["status"] + logger.info("^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^\n") + logger.info("Status:" + str(status)) + addresses = status["addresses"] + logger.info("###############################\n") + logger.info("Addresses:" + str(addresses)) + allocatable = status["allocatable"] + logger.info("********************************\n") + logger.info("Allocatable:" + str(allocatable)) + if hugepg1G in allocatable and hugepg2M in allocatable: + logger.info("Allocatable found..") + allocatable_found = True + capacity = status["capacity"] + if hugepg1G in capacity and hugepg2M in capacity: + logger.info("Capacity found..") + capacity_found = True + + result = allocatable_found and capacity_found + logger.info("Create host network result:" + str(result)) + return result + + def _verify_host_network(self, ssh_prefix, basecmd, host_password, logger): + cmds = ["system host-unlock controller-0"] + while True: + unlock_wait = False + for cmd in cmds: + cmd_to_run = ssh_prefix + " " + basecmd + " " + cmd + logger.info(cmd_to_run) + output_lines = self._run_command_get_all_lines([cmd_to_run], host_password, logger, block=True) + 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 not unlock_wait: + if 'retry host-unlock' in line or 'Rejected' in line: + logger.info("Need to wait to call host-unlock ##### ") + unlock_wait = True + break + if unlock_wait: + time.sleep(60) + else: + break + + cmds = ["system host-show controller-0"] + system_available = False + while True: + for cmd in cmds: + cmd_to_run = ssh_prefix + " " + basecmd + " " + cmd + logger.info(cmd_to_run) + output_lines = self._run_command_get_all_lines([cmd_to_run], host_password, logger, block=True) + 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 'available' in line: + logger.info("Found available #######") + system_available = True + break + if not system_available: + time.sleep(5) + else: + break + + cmds = [] + cmds.append(ssh_prefix + " kubectl get nodes controller-0 --kubeconfig=/etc/kubernetes/admin.conf -o json") + + output_lines = self._run_command_get_all_lines(cmds, host_password, logger, block=True) + logger.info("Returned output lines:") + logger.info(output_lines) + new_op_lines = "" + for line in output_lines.split("\n"): + if not 'Connection to' in line: + new_op_lines = new_op_lines + line + "\n" + + allocatable_found = False + capacity_found = False + + f1c = 'intel.com/pci_sriov_net_f1c' + f1u = 'intel.com/pci_sriov_net_f1u' + fh0 = 'intel.com/pci_sriov_net_fh0' + fh0m = 'intel.com/pci_sriov_net_fh0m' + fh1 = 'intel.com/pci_sriov_net_fh1' + logger.info("New o/p lines:" + new_op_lines) + if len(new_op_lines) > 0: + json_op = json.loads(new_op_lines) + logger.info("JSON O/P:" + str(json_op)) + status = json_op["status"] + logger.info("^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^\n") + logger.info("Status:" + str(status)) + addresses = status["addresses"] + logger.info("###############################\n") + logger.info("Addresses:" + str(addresses)) + allocatable = status["allocatable"] + logger.info("********************************\n") + logger.info("Allocatable:" + str(allocatable)) + if f1c in allocatable and f1u in allocatable and fh0 in allocatable and fh0m in allocatable and fh1 in allocatable: + logger.info("Allocatable found..") + allocatable_found = True + capacity = status["capacity"] + if f1c in capacity and f1u in capacity and fh0 in capacity and fh0m in capacity and fh1 in capacity: + logger.info("Capacity found..") + capacity_found = True + + result = allocatable_found and capacity_found + logger.info("Create host network result:" + str(result)) + return result + + def _create_serviceaccount(self, namespace, remote_region_oam_ip, logger): + logger.info("Inside _create_serviceaccount") + + #host_username = settings.HOST_CREDS['username'] + #host_password = settings.HOST_CREDS['password'] + #ssh_prefix = 'ssh -o StrictHostKeyChecking=no -t ' + host_username + '@' + remote_region_oam_ip + + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + + saName = 'vran-serviceaccount' + cmds_sa = [] + cmds_sa.append(ssh_prefix + ' kubectl create serviceaccount --kubeconfig=/etc/kubernetes/admin.conf ' + saName + ' -n ' + namespace) + run_commands(cmds_sa, host_password, logger, block=True) + + def _create_docker_reg_secret(self, namespace, remote_region_oam_ip, logger): + logger.info("Creating Docker registry secret in Namespace:" + namespace) + + #host_username = settings.HOST_CREDS['username'] + #host_password = settings.HOST_CREDS['password'] + #ssh_prefix = 'ssh -o StrictHostKeyChecking=no -t ' + host_username + '@' + remote_region_oam_ip + + host_username, host_password, ssh_prefix = get_host_connection_details(remote_region_oam_ip) + dr_username, dr_password = get_docker_reg_connection_details() + + #dr_username = settings.CENTRAL_DR_CREDS['username'] + #dr_password = settings.CENTRAL_DR_CREDS['password'] + + secret_name = 'admin-registry-secret' + cmd = ssh_prefix + " kubectl create secret --kubeconfig=/etc/kubernetes/admin.conf docker-registry " + secret_name + cmd = cmd + " --docker-server=registry.local:9001 --docker-username=" + dr_username + " --docker-password=" + dr_password + cmd = cmd + " -n " + namespace + + cmds = [] + cmds.append(cmd) + run_commands(cmds, host_password, logger, block=True) + 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) diff --git a/src/orchestration/vmb_consumer.py b/src/orchestration/vmb_consumer.py new file mode 100644 index 0000000..3d6ebe6 --- /dev/null +++ b/src/orchestration/vmb_consumer.py @@ -0,0 +1,88 @@ +import os +import pprint + +from django.conf import settings +from pulsar import Client, AuthenticationTLS +from .vmb_messages import * + +class VMBConsumer(object): + + os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'src.far_edge_ops_api.settings') + + TLS_CERT_FILE = os.path.join(settings.VMB["CERTS_PATH"], settings.VMB["TLS_CERT_FILE"]) + TLS_KEY_FILE = os.path.join(settings.VMB["CERTS_PATH"], settings.VMB["TLS_KEY_FILE"]) + TLS_TRUST_CERTS_FILE_PATH = os.path.join(settings.VMB["CERTS_PATH"], settings.VMB["TLS_TRUST_CERTS_FILE_PATH"]) + + # Each environment(Rocklin, Branchburg, AWS Non Prod/PLE/STG/PROD has its own service url + SERVICE_URL = settings.VMB['SERVICE_URL'] + #SERVICE_URL = "pulsar://localhost:6650/" + + auth = AuthenticationTLS(TLS_CERT_FILE,TLS_KEY_FILE) + + def get_topic_message(self, topic, topic_schema): + + client = Client(self.SERVICE_URL, tls_trust_certs_file_path=self.TLS_TRUST_CERTS_FILE_PATH, + tls_allow_insecure_connection=False, authentication=self.auth) + #client = Client(SERVICE_URL) + + topic = settings.VMB["VMB_TOPICS"][topic] + consumer = client.subscribe(topic=topic, subscription_name='sub-2', schema=JsonSchema(eval(topic_schema))) + pp = pprint.PrettyPrinter(indent=4) + try: + while True: + msg = consumer.receive() + ex = msg.value() + try: + print("Received message {}".format(ex)) + consumer.acknowledge(msg) + except: + # Message failed to be processed + consumer.negative_acknowledge(msg) + except KeyboardInterrupt: + pass + + client.close() + +import sys, getopt + +def main(argv): + try: + opts, args = getopt.getopt(argv,"ht:",["topic="]) + except getopt.GetoptError: + print('vmb_consumer.py -t <topic>') + sys.exit(2) + for opt, arg in opts: + if opt == '-h': + print('vmb_consumer.py -t [cluster_status|namespace|kubeconfig|image_status>') + sys.exit() + elif opt in ("-t", "--topic"): + topic = arg + print('topic is ' + topic) + + vmbConsumer = VMBConsumer() + + + if (topic == 'namespace'): + # 1. Namespace Creation + print('subscribe to namespace topic ...') + vmbConsumer.get_topic_message(topic='TOPIC_NAMESPACE_CREATION', topic_schema='NameSpaceMessage') + + elif(topic == 'cluster_status'): + # 2. Cluster Status + print('subscribe to cluster_status topic ...') + vmbConsumer.get_topic_message(topic='TOPIC_CLUSTER_STATUS', topic_schema='ClusterStatusMessage') + + elif(topic == 'kubeconfig'): + # 3. Kubeconfig Token + print('subscribe to kubeconfig topic ...') + vmbConsumer.get_topic_message(topic='TOPIC_KUBECONFIG_TOKEN', topic_schema='KubeconfigMessage') + + elif(topic == 'image_status'): + # 4. Images Status + print('subscribe to image_status topic ...') + vmbConsumer.get_topic_message(topic='TOPIC_IMAGE_STATUS', topic_schema='ImagesStatusMessage') + else: + print('unsupport topic: ' + topic) + +if __name__ == "__main__": + main(sys.argv[1:])
\ No newline at end of file diff --git a/src/orchestration/vmb_handler.py b/src/orchestration/vmb_handler.py new file mode 100644 index 0000000..fc003f2 --- /dev/null +++ b/src/orchestration/vmb_handler.py @@ -0,0 +1,52 @@ +import django + +from django.conf import settings +from django.core.exceptions import AppRegistryNotReady +from django.db import transaction +import logging +from logging.handlers import QueueHandler + +from .vmb_producer import VMBProducer + +try: + django.setup() + from .models import ImageSync, CentralToRemoteMap, RemoteRegionSetup +except django.core.exceptions.AppRegistryNotReady as exp: + pass + +class VMBHandler(): + + def __init__(self, loggerQueue, vmbQueue, coordinateQueue): + self.loggerQueue = loggerQueue + self.vmbQueue = vmbQueue + self.coordinateQueue = coordinateQueue + + def run(self): + qh = QueueHandler(self.loggerQueue) + self.logger = logging.getLogger() + self.logger.addHandler(qh) + self.logger.setLevel(logging.DEBUG) + self.vmbProducer = VMBProducer() + + self.logger.info("VMBHandler started...") + while True: + #self.logger.info("--------------------------") + item = self.vmbQueue.get() + self._send_message(item) + + def _send_message(self, item): + message = item['message'] + if message == 'ImageStatus': + images_status_message = item['payload'] + self.vmbProducer.images_status(images_status_message) + self.coordinateQueue.put("Done") + if message == 'Kubeconfig': + kubeconfig_message = item['payload'] + self.vmbProducer.kubeconfig(kubeconfig_message) + if message == 'Namespace': + namespace_message = item['payload'] + self.vmbProducer.namespace_creation(namespace_message) + if message == 'ClusterStatus': + cluster_status_message = item['payload'] + self.vmbProducer.cluster_status(cluster_status_message) + return diff --git a/src/orchestration/vmb_messages.py b/src/orchestration/vmb_messages.py new file mode 100644 index 0000000..515dc03 --- /dev/null +++ b/src/orchestration/vmb_messages.py @@ -0,0 +1,174 @@ +# -*- coding: UTF-8 -*- +from pulsar.schema import * + +''' +Topic - VCP Far Edge Namespace Provisioning +Payload +Format: JSON +Example Content: +{ + "reportName": “vcp_fe_namespace”, + "reportDescription": null, + "reportGeneratedOn": "2020-04-28T09:06:32.1962088-04:00", + "rowCount": 1, + "reportDataRows": [{ + "cluster”: "wsbomagj-d654321-001", + “namespace”: “WSBOMAGJ-441352VZWcVDU-Y-SM-x-001” + “location”: 654321 + “created_at”: 2020-01-07 04:16:15.743617 + }] +} +''' +class RemoteRegion(Record): + cluster = String() + namespace = String() + location = String() + created_at = String() + +class NameSpaceMessage(Record): + reportName = String() + reportDescription = String() + reportGeneratedOn = String() + rowCount = Integer() + reportDataRows = Array(RemoteRegion()) + +''' +Topic - VCP Far Edge Cluster Status +Payload +Format: JSON +Example Content: +{ + "reportName": “vcp_fe_cluster_status”, + "reportDescription": null, + "reportGeneratedOn": "2020-04-28T09:06:32.1962088-04:00", + "rowCount": 1, + "reportDataRows": [{ + “name”: "wsbomagj-d654321-001", + “description”: NE CONCORD 8_NH + “location”: 654321 + “software version”: 19.12 + “availability”: online + “deploy_status”: complete + “created_at”: 2020-01-07 04:16:15.743617 + “updated_at”: 2020-01-07 06:03:10.854598 + }] +} +''' +class ClusterStatus(Record): + name = String() + description = String() + software_version = String() + location = String() + availability = String() + deploy_status = String() + created_at = String() + updated_at = String() + +class ClusterStatusMessage(Record): + reportName = String() + reportDescription = String() + reportGeneratedOn = String() + rowCount = Integer() + reportDataRows = Array(ClusterStatus()) + +''' +Topic - VCP Far Edge Kubeconfig Token +Payload +Format: JSON +Example Content: +{ + "reportName": “vcp_fe_kubeconfig”, + "reportDescription": null, + "reportGeneratedOn": "2020-04-28T09:06:32.1962088-04:00", + "rowCount": 1, + "reportDataRows": [{ + “cluster”: "wsbomagj-d654321-001", + “namespace”: “WSBOMAGJ-441352VZWcVDU-Y-SM-x-001” + “location”: 654321 + “kubeconfig”: <Example kubeconfig> + “created_at”: 2020-01-07 04:16:15.743617 + }] +} +--- +Example kubeconfig: + +apiVersion: v1 +kind: Config +users: +- name: ldap-user + user: + token: eyJhbGciOiJSUzI1NiIsImtpZCI6IiJ9.eyJpc3MiOiJrdWJlcm5ldGVzL3NlcnZpY2VhY2NvdW50Iiwia3ViZXJuZXRlcy5pby9zZXJ2aWNlYWNjb3VudC9uYW1lc3BhY2UiOiJub2tpYSIsImt1YmVybmV0ZXMuaW8vc2VydmljZWFjY291bnQvc2VjcmV0Lm5hbWUiOiJzYTEtdG9rZW4tbW5mMmoiLCJrdWJlcm5ldGVzLmlvL3NlcnZpY2VhY2NvdW50L3NlcnZpY2UtYWNjb3VudC5uYW1lIjoic2ExIiwia3ViZXJuZXRlcy5pby9zZXJ2aWNlYWNjb3VudC9zZXJ2aWNlLWFjY291bnQudWlkIjoiZTA5YTkwMDItMTg0Zi0xMWVhLWE5NzYtMDgwMDI3YTFjODc3Iiwic3ViIjoic3lzdGVtOnNlcnZpY2VhY2NvdW50Om5va2lhOnNhMSJ9.e9uD5kMmAT0HRboSTAbH5xlkETkltLclVQ2GedvoeUmH76WB6G5kGWQrhJjkjpMtPDKxWp6wTzZdEXXwGYGdB6aXbaxcAmau1qid5NGz725BtaoRbSVS2Uk6XrOSNfycFzqc8Z7GTX81VtKKSPYnjeMo47W6FqHw6qk0NEpLLxbpGfJHz8w2KZQiuvI-JRQXtA3PHW3tEaWq3ME3XnYgHNSRJmoiKA99bWjN-HKoOsBDMjhX7kw_VtycRYJ1gbqXmVsSl7BuAQjoplnLN_stRt7ZgpsV4aqmZueqJqHaglB91XOeqe_JcD5unLMb3B5VXpEXPp3V6tjLIyVD8F8XbA +clusters: +- cluster: + server: https://<IPv6>:8443 + name: wsbomagj-d654321-001 +contexts: +- context: + cluster: wsbomagj-d654321-001 + user: ldap-user + namespace: WSBOMAGJ-441352VZWcVDU-Y-SM-x-001 + name: ldap-user +current-context: ldab-user +''' +class Kubeconfig(Record): + transactionId = String() + cluster = String() + namespace = String() + location = String() + kubeconfig = String() + created_at = String() + updated_at = String() + +class KubeconfigMessage(Record): + reportName = String() + reportDescription = String() + reportGeneratedOn = String() + rowCount = Integer() + reportDataRows = Array(Kubeconfig()) + +''' +Topic - VCP Far Edge Image Status +Payload +Format: JSON +Example Content: + +Image Upload Example +{ + "reportName": “vcp_fe_imagestatus”, + "reportDescription": null, + "reportGeneratedOn": "2020-04-28T09:06:32.1962088-04:00", + "rowCount": 2, + "reportDataRows": [{ + “cluster”: "wsbomagj-d654321-001", + “image”: “i1” + “action”: “UPLOAD” + “status”: “SUCCESS” + “message”:”Image Successfully Deleted” + “created_at”: 2020-01-07 04:16:15.743617 + } + { + “cluster”: "wsbomagj-d654321-001", + “image”: “i2” + “action”: “UPLOAD” + “status”: “FAILED” + “message”:”Image Not Present in Artifactory” + “created_at”: 2020-01-07 04:16:15.743617 + }] +} +''' +class ImageStatus(Record): + cluster = String() + image = String() + action = String() + status = String() + message = String() + created_at = String() + +class ImagesStatusMessage(Record): + reportName = String() + transactionId = String() + reportDescription = String() + reportGeneratedOn = String() + rowCount = Integer() + reportDataRows = Array(ImageStatus()) + diff --git a/src/orchestration/vmb_producer.py b/src/orchestration/vmb_producer.py new file mode 100644 index 0000000..709b193 --- /dev/null +++ b/src/orchestration/vmb_producer.py @@ -0,0 +1,329 @@ +import sys, getopt +import yaml +import os + +from django.conf import settings +from .vmb_messages import * +from pulsar import Client, AuthenticationTLS + +class VMBProducer(object): + + #os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'src.far_edge_ops_api.settings') + + TLS_CERT_FILE = os.path.join(settings.VMB["CERTS_PATH"], settings.VMB["TLS_CERT_FILE"]) + TLS_KEY_FILE = os.path.join(settings.VMB["CERTS_PATH"], settings.VMB["TLS_KEY_FILE"]) + TLS_TRUST_CERTS_FILE_PATH = os.path.join(settings.VMB["CERTS_PATH"], settings.VMB["TLS_TRUST_CERTS_FILE_PATH"]) + + #TLS_CERT_FILE = os.path.join(os.path.dirname(__file__), settings.VMB["TLS_CERT_FILE"]) + #TLS_KEY_FILE = os.path.join(os.path.dirname(__file__), settings.VMB["TLS_KEY_FILE"]) + #TLS_TRUST_CERTS_FILE_PATH = os.path.join(os.path.dirname(__file__), settings.VMB["TLS_TRUST_CERTS_FILE_PATH"]) + + # Each environment(Rocklin, Branchburg, AWS Non Prod/PLE/STG/PROD has its own service url + SERVICE_URL = settings.VMB['SERVICE_URL'] + #SERVICE_URL = "pulsar://localhost:6650/" + + auth = AuthenticationTLS(TLS_CERT_FILE,TLS_KEY_FILE) + + def vmb_client(self, topic, schema_name, message): + client = Client(self.SERVICE_URL, tls_trust_certs_file_path=self.TLS_TRUST_CERTS_FILE_PATH, tls_allow_insecure_connection=False, authentication=self.auth) + + topic = settings.VMB["VMB_TOPICS"][topic] + print("topic: " + topic) + producer = client.create_producer(topic=topic, schema=JsonSchema(eval(schema_name))) + consumer = client.subscribe(topic=topic, subscription_name='sub-fe', schema=JsonSchema(eval(schema_name))) + producer.send(message) + + msg = consumer.receive() + + try: + print("VMB-received message: %s" % (msg.value())) + consumer.acknowledge(msg) + except: + # Message failed to be processed + consumer.negative_acknowledge(msg) + + producer.close() + client.close() + + + def namespace_creation(self, message): + ''' + Payload + Format: JSON + Example Content: + { + "reportName": “vcp_fe_namespace”, + "reportDescription": null, + "reportGeneratedOn": "2020-04-28T09:06:32.1962088-04:00", + "rowCount": 1, + "reportDataRows": [{ + “cluster”: "wsbomagj-d654321-001", + “namespace”: “WSBOMAGJ-441352VZWcVDU-Y-SM-x-001” + “location”: 654321 + “created_at”: 2020-01-07 04:16:15.743617 + }] + } + ''' + + self.vmb_client("TOPIC_NAMESPACE_CREATION", "NameSpaceMessage", message) + + def cluster_status(self, message): + ''' + Format: JSON + Example Content: + + { + "reportName": “vcp_fe_cluster_status”, + "reportDescription": null, + "reportGeneratedOn": "2020-04-28T09:06:32.1962088-04:00", + "rowCount": 1, + "reportDataRows": [ + { + “name”: "wsbomagj-d654321-001", + “description”: NE CONCORD 8_NH + “location”: 654321 + “software version”: 19.12 + “availability”: online + “deploy_status”: complete + “created_at”: 2020-01-07 04:16:15.743617 + “updated_at”: 2020-01-07 06:03:10.854598 + } + ] + ''' + + self.vmb_client("TOPIC_CLUSTER_STATUS", "ClusterStatusMessage", message) + + def kubeconfig(self, message): + ''' + Topic - VCP Far Edge Kubeconfig Token + Payload + Format: JSON + Example Content: + { + "reportName": “vcp_fe_kubeconfig”, + "reportDescription": null, + "reportGeneratedOn": "2020-04-28T09:06:32.1962088-04:00", + "rowCount": 1, + "reportDataRows": [{ + “cluster”: "wsbomagj-d654321-001", + “namespace”: “WSBOMAGJ-441352VZWcVDU-Y-SM-x-001” + “location”: 654321 + “kubeconfig”: <Example kubeconfig> + “created_at”: 2020-01-07 04:16:15.743617 + }] + } + --- + Example kubeconfig: + + apiVersion: v1 + kind: Config + users: + - name: ldap-user + user: + token: eyJhbGciOiJSUzI1NiIsImtpZCI6IiJ9.eyJpc3MiOiJrdWJlcm5ldGVzL3NlcnZpY2VhY2NvdW50Iiwia3ViZXJuZXRlcy5pby9zZXJ2aWNlYWNjb3VudC9uYW1lc3BhY2UiOiJub2tpYSIsImt1YmVybmV0ZXMuaW8vc2VydmljZWFjY291bnQvc2VjcmV0Lm5hbWUiOiJzYTEtdG9rZW4tbW5mMmoiLCJrdWJlcm5ldGVzLmlvL3NlcnZpY2VhY2NvdW50L3NlcnZpY2UtYWNjb3VudC5uYW1lIjoic2ExIiwia3ViZXJuZXRlcy5pby9zZXJ2aWNlYWNjb3VudC9zZXJ2aWNlLWFjY291bnQudWlkIjoiZTA5YTkwMDItMTg0Zi0xMWVhLWE5NzYtMDgwMDI3YTFjODc3Iiwic3ViIjoic3lzdGVtOnNlcnZpY2VhY2NvdW50Om5va2lhOnNhMSJ9.e9uD5kMmAT0HRboSTAbH5xlkETkltLclVQ2GedvoeUmH76WB6G5kGWQrhJjkjpMtPDKxWp6wTzZdEXXwGYGdB6aXbaxcAmau1qid5NGz725BtaoRbSVS2Uk6XrOSNfycFzqc8Z7GTX81VtKKSPYnjeMo47W6FqHw6qk0NEpLLxbpGfJHz8w2KZQiuvI-JRQXtA3PHW3tEaWq3ME3XnYgHNSRJmoiKA99bWjN-HKoOsBDMjhX7kw_VtycRYJ1gbqXmVsSl7BuAQjoplnLN_stRt7ZgpsV4aqmZueqJqHaglB91XOeqe_JcD5unLMb3B5VXpEXPp3V6tjLIyVD8F8XbA + clusters: + - cluster: + server: https://<IPv6>:8443 + name: wsbomagj-d654321-001 + contexts: + - context: + cluster: wsbomagj-d654321-001 + user: ldap-user + namespace: WSBOMAGJ-441352VZWcVDU-Y-SM-x-001 + name: ldap-user + current-context: ldab-user + ''' + self.vmb_client("TOPIC_KUBECONFIG_TOKEN", "KubeconfigMessage", message) + + def images_status(self, message): + ''' + Topic - VCP Far Edge Image Status + Payload + Format: JSON + Example Content: + + Image Upload Example + { + "reportName": “vcp_fe_imagestatus”, + "reportDescription": null, + "reportGeneratedOn": "2020-04-28T09:06:32.1962088-04:00", + "rowCount": 2, + "reportDataRows": [{ + “cluster”: "wsbomagj-d654321-001", + “image”: “i1” + “action”: “UPLOAD” + “status”: “SUCCESS” + “message”:”Image Successfully Deleted” + “created_at”: 2020-01-07 04:16:15.743617 + } + { + “cluster”: "wsbomagj-d654321-001", + “image”: “i2” + “action”: “UPLOAD” + “status”: “FAILED” + “message”:”Image Not Present in Artifactory” + “created_at”: 2020-01-07 04:16:15.743617 + }] + } + ''' + + self.vmb_client("TOPIC_IMAGE_STATUS", "ImagesStatusMessage", message) + +def main(argv): + try: + opts, args = getopt.getopt(argv,"ht:i:",["topic=","ifile="]) + except getopt.GetoptError: + print('vmb_producer.py -t <topic> -i <inputfile>') + sys.exit(2) + for opt, arg in opts: + if opt == '-h': + print('vmb_producer.py -t [cluster_status|namespace|kubeconfig|image_status> -i <inputfile>') + sys.exit() + elif opt in ("-t", "--topic"): + topic = arg + print('topic is ' + topic) + elif opt in ("-i", "--ifile"): + inputfile = arg + print('Input file is ' + inputfile) + + + vmbProducer = VMBProducer() + + f = open (inputfile, "r") + message_json = json.loads(f.read()) + print(message_json) + + if (topic == 'namespace'): + # 1. Namespace Creation + print('return namespace to vmb ...') + remote_region = RemoteRegion(cluster='wsbomagj', + namespace='WSBOMAGcluster', + location='CILI654321', + created_at='2020-01-07 04:16:15.743617') + remote_region1 = RemoteRegion(cluster='wsbomagj1', + namespace='WSBOMAGcluster1', + location='123456', + created_at='2020-06-07 04:16:15.743617') + namespace_message = NameSpaceMessage(reportName='vcp_fe_namespace', + reportDescription='dRAN namespace', + reportGeneratedOn='2020-06-07 04:16:15.743617', + rowCount= 2, + reportDataRows=[remote_region, remote_region1]) + + vmbProducer.namespace_creation(namespace_message) + + elif(topic == 'cluster_status'): + # 2. Cluster Status + print('return cluster status to vmb ...') + + cluster_status_message = ClusterStatusMessage() + for k in message_json: + if k == 'reportName': + cluster_status_message.reportName = message_json[k] + if k == 'reportDescription': + cluster_status_message.reportDescription = message_json[k] + if k == 'reportGeneratedOn': + cluster_status_message.reportGeneratedOn = message_json[k] + if k == 'reportDataRows': + regions_list = message_json['reportDataRows'] + for region in regions_list: + cluster_status = ClusterStatus() + for k in region: + if k == 'name': + cluster_status.name = region[k] + if k == 'description': + cluster_status.description = region[k] + if k == 'location': + cluster_status.location = region[k] + 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_message.reportDataRows = [cluster_status] + cluster_status_message.rowCount = 1 + + + #cluster_status = ClusterStatus(name= "wsbomagj-d654321-001", + # description='NE CONCORD 8_NH', + # location='654321', + # software_version='19.12', + # availability='online', + # deploy_status='complete', + # created_at='2020-01-07 04:16:15.743617', + # updated_at='2020-01-07 06:03:10.854598') + + #cluster_status_message = ClusterStatusMessage(reportName='vcp_fe_cluster_status', + # reportDescription='cluster status', + # reportGeneratedOn='2020-04-28T09:06:32.1962088-04:00', + # reportDataRows=cluster_status + #) + + vmbProducer.cluster_status(cluster_status_message) + + elif(topic == 'kubeconfig'): + # 3. Kubeconfig Token + print('return kubeconfig to vmb ...') + kubeconfig = Kubeconfig( + cluster='wsbomagj-d654321-001', + namespace='WSBOMAGJ-441352VZWcVDU-Y-SM-x-001', + location='654321', + kubeconfig='Example kubeconfig', + created_at='2020-01-07 04:16:15.743617' + ) + + + with open('tests/kubeconfig.yaml') as f: + #kubeconfig_value = json.dumps(yaml.load(f,Loader=yaml.FullLoader), indent=2) + kubeconfig_value = json.dumps(yaml.load(f,Loader=yaml.FullLoader)) + print(kubeconfig_value) + kubeconfig.kubeconfig = kubeconfig_value + + kubeconfig_message = KubeconfigMessage( + reportName='vcp_fe_kubeconfig', + reportDescription='kubeconfig file', + reportGeneratedOn='2020-04-28T09:06:32.1962088-04:00', + rowCount = 1, + reportDataRows=[kubeconfig] + ) + + vmbProducer.kubeconfig(kubeconfig_message) + + elif(topic == 'image_status'): + # 4. Images Status + print('return image status to vmb ...') + image_status1 = ImageStatus( + cluster='wsbomagj-d654321-001', + image='i1', + action='UPLOAD', + status='SUCCESS', + message='Image Successfully uploaded', + created_at='2020-01-07 04:16:15.743617', + ) + + image_status2 = ImageStatus( + cluster='wsbomagj-d654321-002', + image='i1', + action='UPLOAD', + status='FAILED', + message='Image uploading failed', + created_at='2020-01-07 04:16:15.743617', + ) + images_status_message = ImagesStatusMessage(reportName='vcp_fe_imagestatus', + transactionId='13e11612a7494cd6a7f69ea46e3c0537', + reportDescription='cluster status', + reportGeneratedOn='2020-04-28T09:06:32.1962088-04:00', + rowCount = 2, + reportDataRows=[image_status1,image_status2]) + + vmbProducer.images_status(images_status_message) + else: + print('unsupport topic: ' + topic) + +if __name__ == "__main__": + main(sys.argv[1:])
\ No newline at end of file |
