From 92e7939f87fedc435bf3c4351915602f912155af Mon Sep 17 00:00:00 2001
From: thijsheijden <hi@thijsheijden.nl>
Date: Fri, 21 May 2021 15:09:42 +0200
Subject: [PATCH] Added queryID as optional header key for backwards
 compatibility

Also added the Minio environment variables to the Kubernetes deployment.
---
 deployments/deployment.yml                 |  6 ++++++
 internal/usecases/consume/handlemessage.go | 16 +++++++++-------
 2 files changed, 15 insertions(+), 7 deletions(-)

diff --git a/deployments/deployment.yml b/deployments/deployment.yml
index cf1c8bb..a2f5bf4 100644
--- a/deployments/deployment.yml
+++ b/deployments/deployment.yml
@@ -38,6 +38,12 @@ spec:
           value: redis.redis.svc.cluster.local:6379
         - name: LOG_MESSAGES
           value: "true"
+        - name: MINIO_ADDRESS
+          value: minio:9000
+        - name: MINIO_ACCESSKEYID
+          value: root
+        - name: MINIO_ACCESSKEY
+          value: DikkeDraak
         resources:
             requests:
               memory: "100Mi"
diff --git a/internal/usecases/consume/handlemessage.go b/internal/usecases/consume/handlemessage.go
index 3a3bc11..339e060 100644
--- a/internal/usecases/consume/handlemessage.go
+++ b/internal/usecases/consume/handlemessage.go
@@ -27,9 +27,10 @@ func (s *Service) HandleMessage(msg *broker.Message) {
 		return
 	}
 
-	queryID, ok := msg.Headers["queryID"].(string)
+	var queryID string
+	queryID, ok = msg.Headers["queryID"].(string)
 	if !ok {
-		return
+		queryID = ""
 	}
 
 	s.sendStatus("Received", &sessionID)
@@ -117,10 +118,9 @@ func (s *Service) HandleMessage(msg *broker.Message) {
 		return
 	}
 
-	s.sendStatus("Done", &sessionID)
+	s.sendStatus("Completed", &sessionID)
 
 	// Add type indicator to result from database
-	// TODO: Change key 'values' to key 'value'
 	var res interface{}
 	json.Unmarshal(*result, &res)
 	resultMsg := entity.MessageStruct{
@@ -134,9 +134,11 @@ func (s *Service) HandleMessage(msg *broker.Message) {
 	}
 
 	// Caching message in object store
-	cacheQueryResultContext, cancelCacheQueryResult := context.WithTimeout(context.Background(), time.Second*20)
-	defer cancelCacheQueryResult()
-	s.objectStore.Put(cacheQueryResultContext, "cached-queries", fmt.Sprintf("%s-%s", sessionID, queryID), bytes.NewReader(resultMsgBytes))
+	if queryID != "" {
+		cacheQueryResultContext, cancelCacheQueryResult := context.WithTimeout(context.Background(), time.Second*20)
+		defer cancelCacheQueryResult()
+		s.objectStore.Put(cacheQueryResultContext, "cached-queries", fmt.Sprintf("%s-%s", sessionID, queryID), bytes.NewReader(resultMsgBytes))
+	}
 
 	s.producer.PublishMessage(&resultMsgBytes, &sessionID)
 }
-- 
GitLab