diff options
Diffstat (limited to 'src/orchestration/db_updater.py')
| -rw-r--r-- | src/orchestration/db_updater.py | 83 |
1 files changed, 83 insertions, 0 deletions
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 + |
