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