summaryrefslogtreecommitdiff
path: root/src/orchestration/db_updater.py
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/db_updater.py
parent22dae02a86c1fce71091bfa5289cf99abba1b217 (diff)
more filesHEADmaster
Diffstat (limited to 'src/orchestration/db_updater.py')
-rw-r--r--src/orchestration/db_updater.py83
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
+