| 1 | from typing import Optional |
| 2 | |
| 3 | import asyncgelf |
| 4 | |
| 5 | from app.connectors.utils import get_connector_info_from_db |
| 6 | from app.db.db_session import get_db_session |
| 7 | |
| 8 | |
| 9 | class GelfLogger: |
| 10 | def __init__(self, host: str, port: str, compress: Optional[bool] = False): |
| 11 | self.host = host |
| 12 | self.port = port |
| 13 | self.compress = compress |
| 14 | |
| 15 | async def tcp_handler(self, message): |
| 16 | if not isinstance(message, dict): |
| 17 | message = message.to_dict() |
| 18 | |
| 19 | handler = asyncgelf.GelfTcp( |
| 20 | host=self.host, |
| 21 | port=self.port, |
| 22 | compress=self.compress, |
| 23 | ) |
| 24 | |
| 25 | response = await handler.tcp_handler(message) |
| 26 | return response |
| 27 | |
| 28 | |
| 29 | async def create_gelf_logger(): |
| 30 | async with get_db_session() as session: |
| 31 | connector_info = await get_connector_info_from_db("Event Shipper", session) |
| 32 | return GelfLogger( |
| 33 | host=connector_info["connector_url"], |
| 34 | port=str(connector_info["connector_extra_data"]), |
| 35 | ) |