diff --git a/apps/flows/wb/backend.yaml b/apps/flows/wb/backend.yaml index 268ff2c..90f6889 100644 --- a/apps/flows/wb/backend.yaml +++ b/apps/flows/wb/backend.yaml @@ -99,6 +99,21 @@ spec: path: _default: authentication.py + - name: kafka-configmap + mountPath: + _default: /opt/src/flow/kafka.py + subPath: + _default: kafka.py + readOnly: + _default: true + configMap: + name: + _default: kafka-configmap + items: + - key: kafka.py + path: + _default: kafka.py + envs: - name: KAFKA_HOST value: diff --git a/apps/flows/wb/celery.yaml b/apps/flows/wb/celery.yaml index 695dcf8..faea438 100644 --- a/apps/flows/wb/celery.yaml +++ b/apps/flows/wb/celery.yaml @@ -70,6 +70,23 @@ spec: name: _default: dockerhub + volumes: + _default: + - name: kafka-configmap + mountPath: + _default: /opt/src/flow/kafka.py + subPath: + _default: kafka.py + readOnly: + _default: true + configMap: + name: + _default: kafka-configmap + items: + - key: kafka.py + path: + _default: kafka.py + envs: - name: KAFKA_HOST value: diff --git a/apps/flows/wb/kafka-configmap.yaml b/apps/flows/wb/kafka-configmap.yaml new file mode 100644 index 0000000..7cc221a --- /dev/null +++ b/apps/flows/wb/kafka-configmap.yaml @@ -0,0 +1,67 @@ +--- +apiVersion: v1 +kind: ConfigMap +metadata: + name: kafka-configmap + namespace: flows +data: + kafka.py: | + import ssl + from functools import partial + + import anyio + from faststream.kafka import KafkaBroker + from faststream.security import BaseSecurity + + from flow.config import logger, settings + from flow.models.events import BaseEvent + + + def _build_broker() -> KafkaBroker: + ssl_context = ssl.create_default_context() + + ssl_context.check_hostname = False + ssl_context.verify_mode = ssl.CERT_NONE + + if settings.kafka.ssl_cert: + ssl_context.load_verify_locations(cadata=settings.kafka.ssl_cert) + + security = BaseSecurity( + ssl_context=ssl_context, + use_ssl=True, + ) + + return KafkaBroker( + bootstrap_servers=[settings.kafka.host], + security=security, + ) + + + broker = _build_broker() + + + async def _publish(event: BaseEvent, key: str) -> None: + await broker.publish( + message=event, + topic=event.metadata.event_type, + key=key.encode(), + ) + + + def publish_event(event: BaseEvent, key: str) -> None: + if not settings.enable_events: + return + + try: + anyio.from_thread.run(partial(_publish, event, key)) + except Exception: + logger.exception( + "Failed to publish kafka event %s", + event.metadata.event_type, + ) + + + __all__ = [ + "broker", + "publish_event", + ] diff --git a/apps/flows/wb/kustomization.yaml b/apps/flows/wb/kustomization.yaml index ad62e58..bb93c0e 100644 --- a/apps/flows/wb/kustomization.yaml +++ b/apps/flows/wb/kustomization.yaml @@ -4,6 +4,7 @@ kind: Kustomization namespace: flows resources: - authentication-configmap.yaml + - kafka-configmap.yaml - backend.yaml - celery.yaml - frontend.yaml