diff --git a/apps/message-hub/wb/kafka-config-configmap.yaml b/apps/message-hub/wb/kafka-config-configmap.yaml new file mode 100644 index 0000000..c68692b --- /dev/null +++ b/apps/message-hub/wb/kafka-config-configmap.yaml @@ -0,0 +1,67 @@ +--- +apiVersion: v1 +kind: ConfigMap +metadata: + name: kafka-config + namespace: message-hub +data: + kafka.py: | + from ssl import SSLContext + + import ssl + + from aiokafka.helpers import create_ssl_context + from faststream.kafka import KafkaBroker + from faststream.security import BaseSecurity, SASLPlaintext, SASLScram512 + from pydantic_settings import BaseSettings, SettingsConfigDict + + + class KafkaSettings(BaseSettings): + HOST: str = 'localhost' + PORT: int = 9092 + USERNAME: str | None = None + PASSWORD: str | None = None + SECURITY_PROTOCOL: str = 'PLAINTEXT' + SASL_MECHANISM: str | None = None + SSL_CAFILE: str | None = None + + model_config = SettingsConfigDict(env_prefix='KAFKA_', env_file='.env', extra='ignore') + + @property + def bootstrap_servers(self) -> list[str]: + return [f'{self.HOST}:{str(self.PORT)}'] + + def get_ssl_context(self) -> SSLContext: + context = create_ssl_context(cafile=self.SSL_CAFILE) + context.check_hostname = False + context.verify_mode = ssl.CERT_NONE + return context + + def get_security(self) -> BaseSecurity | None: + use_ssl = self.SECURITY_PROTOCOL in ('SSL', 'SASL_SSL') + ssl_context = self.get_ssl_context() if use_ssl else None + if self.SASL_MECHANISM == 'PLAINTEXT': + return SASLPlaintext( + username=self.USERNAME or '', + password=self.PASSWORD or '', + ssl_context=ssl_context, + use_ssl=use_ssl, + ) + if self.SASL_MECHANISM == 'SCRAM-SHA-512': + return SASLScram512( + username=self.USERNAME or '', + password=self.PASSWORD or '', + ssl_context=ssl_context, + use_ssl=use_ssl, + ) + if use_ssl: + return BaseSecurity(ssl_context=ssl_context) + return None + + @property + def broker(self) -> KafkaBroker: + return KafkaBroker( + bootstrap_servers=self.bootstrap_servers, + security=self.get_security(), + logger=None, + ) diff --git a/apps/message-hub/wb/kustomization.yaml b/apps/message-hub/wb/kustomization.yaml index 2ac5608..360e6b0 100644 --- a/apps/message-hub/wb/kustomization.yaml +++ b/apps/message-hub/wb/kustomization.yaml @@ -4,4 +4,5 @@ kind: Kustomization namespace: message-hub resources: - kafka-cert-configmap.yaml + - kafka-config-configmap.yaml - message-hub.yaml diff --git a/apps/message-hub/wb/message-hub.yaml b/apps/message-hub/wb/message-hub.yaml index cb383da..13d5eb0 100644 --- a/apps/message-hub/wb/message-hub.yaml +++ b/apps/message-hub/wb/message-hub.yaml @@ -94,6 +94,21 @@ spec: name: _default: kafka-cert + - name: kafka-config-volume + mountPath: + _default: /opt/src/config/kafka.py + subPath: + _default: kafka.py + readOnly: + _default: true + configMap: + name: + _default: kafka-config + items: + - key: kafka.py + path: + _default: kafka.py + envs: - name: WORKER_TIMEOUT value: