Compare commits

...

3 Commits

Author SHA1 Message Date
ivan
9f7af7e39e ++ 2026-08-29 12:57:51 +05:00
ivan
dfcf9423f2 ++ 2026-08-29 12:55:04 +05:00
ivan
b8755c7399 ++ 2026-08-29 12:53:22 +05:00
5 changed files with 100 additions and 19 deletions

View File

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

View File

@ -65,6 +65,22 @@ spec:
enabled: false
readiness:
enabled: false
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
service:
enabled: true
@ -189,14 +205,6 @@ spec:
_default: /flows
secretEnvs:
- name: KAFKA_USERNAME
secretName:
_default: kafka-secret
secretKey: username
- name: KAFKA_PASSWORD
secretName:
_default: kafka-secret
secretKey: password
- name: ADMIN_PANEL_SECRET_KEY
secretName:
_default: admin-secret

View File

@ -69,6 +69,22 @@ spec:
enabled: false
readiness:
enabled: false
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
service:
enabled: false
@ -188,14 +204,6 @@ spec:
secretName:
_default: django-secret
secretKey: token
- name: KAFKA_USERNAME
secretName:
_default: kafka-secret
secretKey: username
- name: KAFKA_PASSWORD
secretName:
_default: kafka-secret
secretKey: password
- name: ADMIN_PANEL_SECRET_KEY
secretName:
_default: admin-secret

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
kind: Kustomization
namespace: flows
@ -7,3 +6,4 @@ resources:
- backend.yaml
- celery.yaml
- frontend.yaml
- kafka-configmap.yaml