From b8755c7399c5db17301a2f0158006bc1c0928026 Mon Sep 17 00:00:00 2001 From: ivan Date: Sat, 29 Aug 2026 12:53:22 +0500 Subject: [PATCH] ++ --- apps/django/ugok/nginx-configmap.yaml | 2 - apps/flows/ugok/kafka-configmap.yaml | 67 +++++++++++++++++++++++++++ apps/flows/ugok/kustomization.yaml | 2 +- 3 files changed, 68 insertions(+), 3 deletions(-) create mode 100644 apps/flows/ugok/kafka-configmap.yaml diff --git a/apps/django/ugok/nginx-configmap.yaml b/apps/django/ugok/nginx-configmap.yaml index 6208eb1..cadd040 100644 --- a/apps/django/ugok/nginx-configmap.yaml +++ b/apps/django/ugok/nginx-configmap.yaml @@ -1,6 +1,4 @@ --- -# Скопировано из живого ConfigMap кластера ugok (namespace django) — -# монтируется frontend в /etc/nginx/nginx.conf. apiVersion: v1 kind: ConfigMap metadata: diff --git a/apps/flows/ugok/kafka-configmap.yaml b/apps/flows/ugok/kafka-configmap.yaml new file mode 100644 index 0000000..7cc221a --- /dev/null +++ b/apps/flows/ugok/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/ugok/kustomization.yaml b/apps/flows/ugok/kustomization.yaml index d27e688..e0eaf35 100644 --- a/apps/flows/ugok/kustomization.yaml +++ b/apps/flows/ugok/kustomization.yaml @@ -1,5 +1,4 @@ --- -# Не наследуем base (backend/celery vault-native) — см. комментарии в файлах. apiVersion: kustomize.config.k8s.io/v1beta1 kind: Kustomization namespace: flows @@ -7,3 +6,4 @@ resources: - backend.yaml - celery.yaml - frontend.yaml + - kafka-configmap.yaml