summaryrefslogtreecommitdiff
path: root/src/orchestration
diff options
context:
space:
mode:
authorckonstanski <kostcarl@isu.edu>2026-07-30 18:17:14 -0600
committerckonstanski <kostcarl@isu.edu>2026-07-30 18:17:14 -0600
commit640ff61422bc0ee3966941b93d9889ddbbd38d77 (patch)
treeabf1c08f5d38ee53ec8b29dc4f425722517505e6 /src/orchestration
parent22dae02a86c1fce71091bfa5289cf99abba1b217 (diff)
more filesHEADmaster
Diffstat (limited to 'src/orchestration')
-rw-r--r--src/orchestration/cluster-setup-files/application-ldap.yaml69
-rw-r--r--src/orchestration/cluster-setup-files/application-sa.yaml79
-rw-r--r--src/orchestration/cluster-setup-files/crd-role.yaml18
-rw-r--r--src/orchestration/cluster-setup-files/helm.yaml75
-rw-r--r--src/orchestration/cluster-setup-files/orchestration-ldap.yaml166
-rw-r--r--src/orchestration/cluster-setup-files/orchestration-sa.yaml181
-rw-r--r--src/orchestration/cluster-setup-files/rbd.yaml9
-rw-r--r--src/orchestration/cluster-setup-files/samsung-sa.yaml55
-rw-r--r--src/orchestration/cluster-setup-files/tiller-sa.yaml42
-rw-r--r--src/orchestration/cluster-setup-files/vdu_tiller_sa.yaml29
-rw-r--r--src/orchestration/data.json4
-rw-r--r--src/orchestration/db_updater.py83
-rw-r--r--src/orchestration/imagesynchandler.py622
-rw-r--r--src/orchestration/kubeconfighandler.py926
-rw-r--r--src/orchestration/migrations/0008_auto_20200514_1951.py19
-rw-r--r--src/orchestration/migrations/0008_imagesync_transaction_id.py18
-rw-r--r--src/orchestration/migrations/0009_auto_20200529_1929.py18
-rw-r--r--src/orchestration/migrations/0010_centralregiontoremoteregionmap.py21
-rw-r--r--src/orchestration/migrations/0011_auto_20200529_1944.py17
-rw-r--r--src/orchestration/migrations/0012_merge_20200702_0053.py14
-rw-r--r--src/orchestration/migrations/0013_imagesync_remoteregionsetup.py39
-rw-r--r--src/orchestration/migrations/0014_remoteregionsetup_transaction_id.py18
-rw-r--r--src/orchestration/migrations/0015_auto_20200707_2034.py23
-rw-r--r--src/orchestration/models.py11
-rw-r--r--src/orchestration/namespacehandler.py431
-rw-r--r--src/orchestration/remoteregionhandler.py375
-rw-r--r--src/orchestration/tests/cluster_status.json15
-rw-r--r--src/orchestration/tests/image_status.json22
-rw-r--r--src/orchestration/tests/kubeconfig.yaml17
-rw-r--r--src/orchestration/tests/namespace.json12
-rw-r--r--src/orchestration/urls.py12
-rw-r--r--src/orchestration/utils.py110
-rw-r--r--src/orchestration/vendorhandler.py308
-rw-r--r--src/orchestration/views.py572
-rw-r--r--src/orchestration/vmb_consumer.py88
-rw-r--r--src/orchestration/vmb_handler.py52
-rw-r--r--src/orchestration/vmb_messages.py174
-rw-r--r--src/orchestration/vmb_producer.py329
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