/
laritsiuriumov
/
aef-manager-python-services
Обзор
Документация
Войти
/
laritsiuriumov
/
aef-manager-python-services
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
develop
services/gr-processor/src/logging_utils.py
175 строк
5 KB
Цюрюмов Лари Валерьевич [B]
Add JSON file logging for guardrails services
25 май 2026, 15:12
25 май 2026, 15:12
7d5c50f
Код
Авторство
О чём код?
import json import logging import sys from collections.abc import Mapping from datetime import datetime, timezone from os import PathLike from typing import Any ECS_VERSION = "1.2.0" SERVICE_HANDLER_MARKER = "_guardrails_service_logging" OPTIONAL_RECORD_FIELDS = { "app_agent_id": "app.agent_id", "app_distributive": "app.distributive", "app_hub_id": "app.hub_id", "kafka_from_topic": "kafka.from.topic", "kafka_from_key": "kafka.from.key", "kafka_from_partition": "kafka.from.partition", "kafka_to_topic": "kafka.to.topic", "kafka_to_key": "kafka.to.key", } class JsonFileFormatter(logging.Formatter): def __init__(self, service_name: str) -> None: super().__init__() self.service_name = service_name def format(self, record: logging.LogRecord) -> str: payload: dict[str, Any] = { "@timestamp": datetime.fromtimestamp( record.created, tz=timezone.utc, ) .isoformat(timespec="milliseconds") .replace("+00:00", "Z"), "ecs.version": ECS_VERSION, "log.level": record.levelname.lower(), "message": record.getMessage(), "process.thread.name": record.threadName, "log.logger": record.name, "service.name": self.service_name, } for record_attr, log_field in OPTIONAL_RECORD_FIELDS.items(): value = getattr(record, record_attr, None) if value is not None: payload[log_field] = value if record.exc_info: payload["error.stack_trace"] = self.formatException(record.exc_info) return json.dumps(payload, ensure_ascii=False, default=str) class PlainTextFormatter(logging.Formatter): def __init__(self) -> None: super().__init__("%(asctime)s - %(name)s - %(levelname)s - %(message)s") def configure_logging( service_name: str, log_path: str | PathLike[str], level: int | str = logging.INFO, ) -> None: root_logger = logging.getLogger() for handler in root_logger.handlers[:]: if getattr(handler, SERVICE_HANDLER_MARKER, False): root_logger.removeHandler(handler) handler.close() root_logger.setLevel(_coerce_log_level(level)) file_handler = logging.FileHandler(log_path, encoding="utf-8") file_handler.setFormatter(JsonFileFormatter(service_name)) setattr(file_handler, SERVICE_HANDLER_MARKER, True) stream_handler = logging.StreamHandler(sys.stdout) stream_handler.setFormatter(PlainTextFormatter()) setattr(stream_handler, SERVICE_HANDLER_MARKER, True) root_logger.addHandler(file_handler) root_logger.addHandler(stream_handler) def build_app_log_extra( *, agent_id: object | None = None, distributive: object | None = None, hub_id: object | None = None, ) -> dict[str, object]: extra: dict[str, object] = {} if agent_id is not None: extra["app_agent_id"] = agent_id if distributive is not None: extra["app_distributive"] = distributive if hub_id is not None: extra["app_hub_id"] = hub_id return extra def build_payload_log_extra(payload: Mapping[str, Any] | None) -> dict[str, object]: if payload is None: return {} agent_id = _first_present(payload, ("x-agent-id", "agent-id", "agent_id")) nested_key = payload.get("key") if agent_id is None and isinstance(nested_key, Mapping): agent_id = _first_present(nested_key, ("agent_id", "x-agent-id", "agent-id")) return build_app_log_extra( agent_id=agent_id, distributive=_first_present(payload, ("distributive",)), hub_id=_first_present(payload, ("ai-hub-id", "hub-id", "hub_id")), ) def build_kafka_consume_log_extra(record: object) -> dict[str, object]: extra: dict[str, object] = {} topic = getattr(record, "topic", None) partition = getattr(record, "partition", None) key = _decode_kafka_key(getattr(record, "key", None)) if topic is not None: extra["kafka_from_topic"] = topic if key is not None: extra["kafka_from_key"] = key if partition is not None: extra["kafka_from_partition"] = partition return extra def build_kafka_produce_log_extra( topic: str, *, key: object | None = None, ) -> dict[str, object]: extra: dict[str, object] = {"kafka_to_topic": topic} decoded_key = _decode_kafka_key(key) if decoded_key is not None: extra["kafka_to_key"] = decoded_key return extra def merge_log_extra(*extras: Mapping[str, object] | None) -> dict[str, object]: merged: dict[str, object] = {} for extra in extras: if extra: merged.update({key: value for key, value in extra.items() if value is not None}) return merged def _coerce_log_level(level: int | str) -> int: if isinstance(level, int): return level resolved_level = logging.getLevelName(level.upper()) return resolved_level if isinstance(resolved_level, int) else logging.INFO def _first_present(payload: Mapping[str, Any], keys: tuple[str, ...]) -> object | None: for key in keys: value = payload.get(key) if value is not None: return value return None def _decode_kafka_key(key: object | None) -> str | None: if key is None: return None if isinstance(key, bytes): return key.decode("utf-8", errors="replace") return str(key)