This commit is contained in:
ivan 2026-08-29 12:53:22 +05:00
parent bee1882830
commit b8755c7399
3 changed files with 68 additions and 3 deletions

View File

@ -1,6 +1,4 @@
--- ---
# Скопировано из живого ConfigMap кластера ugok (namespace django) —
# монтируется frontend в /etc/nginx/nginx.conf.
apiVersion: v1 apiVersion: v1
kind: ConfigMap kind: ConfigMap
metadata: metadata:

View File

@ -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",
]

View File

@ -1,5 +1,4 @@
--- ---
# Не наследуем base (backend/celery vault-native) — см. комментарии в файлах.
apiVersion: kustomize.config.k8s.io/v1beta1 apiVersion: kustomize.config.k8s.io/v1beta1
kind: Kustomization kind: Kustomization
namespace: flows namespace: flows
@ -7,3 +6,4 @@ resources:
- backend.yaml - backend.yaml
- celery.yaml - celery.yaml
- frontend.yaml - frontend.yaml
- kafka-configmap.yaml