diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 6540426..c86d7ca 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -42,6 +42,11 @@ class Ingestor: self.heartbeat_ttl = int(getenv("HEARTBEAT_TTL", "20")) self.pod_id = getenv("HOSTNAME", "localhost") self.poll_interval = int(getenv("POLL_INTERVAL", "5")) + mongo_url = getenv("MONGO_URL", 'localhost:27017') + mongo_username = getenv("MONGO_USERNAME", 'sientia') + mongo_password = getenv("MONGO_PASSWORD", 'sientia') + self.mongo_database = getenv("MONGO_DATABASE", 'sientia') + self.mongo_connection_string = f"mongodb://{mongo_username}:{mongo_password}@{mongo_url}" self.kafka_servers = kafka_servers.split(",") self.logger = None @@ -138,6 +143,8 @@ class Ingestor: self.heartbeat_ttl, self.pod_id, self.poll_interval, + self.mongo_connection_string, + self.mongo_database, self.logger, self.notification_handler, self.redis_username, diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index 2f62710..6fe99c3 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -14,11 +14,12 @@ class IngestorManager(): def __init__(self, kafka_servers: str, redis_host: str, redis_port: int, lease_ttl: int, heartbeat_ttl: int, pod_id: str, - poll_interval: int, logger: Logger, notification_handler: NotificationHandler, + poll_interval: int, mongo_connection_string: str, mongo_database: str, + logger: Logger, notification_handler: NotificationHandler, redis_username: str = None, redis_password: str = None): self.data_manager = DataManager( - kafka_servers, logger, notification_handler) + kafka_servers, mongo_connection_string, mongo_database, logger, notification_handler) self.opc_managers = {} self.resource_manager = ResourceManager( redis_host, redis_port, lease_ttl, heartbeat_ttl, pod_id, redis_username, redis_password @@ -58,7 +59,8 @@ class IngestorManager(): manager = OpcManager( server_config['name'], server_config['url'], data_manager, logger, server_config['server_uri'], - self.notification_handler, self.pod_id, server_config.get('cert_path'), + self.notification_handler, self.pod_id, server_config.get( + 'cert_path'), server_config.get('private_key_path'), server_config.get('server_cert_path') ) @@ -163,7 +165,8 @@ class IngestorManager(): self.opc_managers[server].disconnect() self.opc_managers.pop(server, None) - metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers)) + metrics.OPC_MANAGERS_ACTIVE.labels( + pod_id=self.pod_id).set(len(self.opc_managers)) def check_opc_servers_integrity(self): """ @@ -182,7 +185,8 @@ class IngestorManager(): for slot, _config in self.managed_tags.items(): self.managed_tags[slot].pop(server, None) - metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers)) + metrics.OPC_MANAGERS_ACTIVE.labels( + pod_id=self.pod_id).set(len(self.opc_managers)) def declare_active(self): """ @@ -266,7 +270,8 @@ class IngestorManager(): if len(acquired) >= max_slots: self.managed_tags.update(acquired) - metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags)) + metrics.SLOTS_MANAGED.labels( + pod_id=self.pod_id).set(len(self.managed_tags)) return acquired self.logger.warning( @@ -274,7 +279,8 @@ class IngestorManager(): f"Only {acquired} slots were leased." ) self.managed_tags.update(acquired) - metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags)) + metrics.SLOTS_MANAGED.labels( + pod_id=self.pod_id).set(len(self.managed_tags)) return acquired def unsubscribe_slot(self, slot: str): @@ -400,7 +406,8 @@ class IngestorManager(): tags_to_sub ) except Exception as e: - metrics.OPC_SUBSCRIPTION_ERRORS.labels(pod_id=self.pod_id, server=server, slot=slot).inc() + metrics.OPC_SUBSCRIPTION_ERRORS.labels( + pod_id=self.pod_id, server=server, slot=slot).inc() trace = traceback.format_exc() self.notification_handler.build_and_send_notification( notification_id=f'OPC_SUBSCRIPTION_ERROR_{slot}:{server}', diff --git a/values.yaml b/values.yaml index d92d51e..e597849 100644 --- a/values.yaml +++ b/values.yaml @@ -11,7 +11,7 @@ image: # This sets the pull policy for images. pullPolicy: Always # Overrides the image tag whose default is the chart appVersion. - tag: "0.1.1" + tag: "0.2.2" # This is for the secrets for pulling an image from a private repository more information can be found here: https://kubernetes.io/docs/tasks/configure-pod-container/pull-image-private-registry/ imagePullSecrets: @@ -155,6 +155,16 @@ env: - name: HTTP_SERVER_PORT value: "4840" + + - name: MONGODB_USERNAME + value: "root" + - name: MONGODB_PASSWORD + value: "wKZDbMNU1c" + - name: MONGODB_URL + value: "my-release-mongodb.mongodb.svc.cluster.local:27017" + - name: MONGODB_DATABASE + value: "sientia" + ssh: enabled: true @@ -163,7 +173,7 @@ ssh: knownHostsPath: /mnt/known_hosts # kubectl create secret docker-registry docker-hub-secret --namespace sientia --docker-server=http://aignosi.azurecr.io --docker-username=aignosi --docker-password=5I5zpQ6sRaHqX1hD3dr+2mo647yO3FRc359/wu6gsP+ACRDRz5mp -# helm upgrade --install sientia-opc-ingestor sientia/sientia-module -n sientia --create-namespace -f ./values.yaml --version 0.1.0-uat +# helm upgrade --install sientia-opc-ingestor sientia/sientia-module -n sientia --create-namespace -f ./values.yaml --version 0.4.0-uat # kubectl create secret generic git-ssh-key-sientia-opc-ingestor \ # --namespace sientia \