add logging and retry to connect to the rabbitmq server

This commit is contained in:
2026-09-14 18:34:15 +03:00
parent a93c6d5fca
commit e124d897eb
10 changed files with 241 additions and 90 deletions
+8 -2
View File
@@ -2,8 +2,14 @@
import logging
from .logger import LoggerDB
from .logger import LoggerDaemon, LoggerDB
sql_logger = logging.getLogger("sqlalchemy.engine")
sql_logger.setLevel(logging.INFO)
sql_logger.addHandler(LoggerDB())
sql_logger.addHandler(LoggerDB())
sql_logger.addHandler(logging.StreamHandler())
daemon_logger=logging.getLogger("daemon")
daemon_logger.setLevel(logging.INFO)
daemon_logger.addHandler(LoggerDaemon())
daemon_logger.addHandler(logging.StreamHandler())
+63
View File
@@ -0,0 +1,63 @@
import json
from time import perf_counter
from typing import cast
from uuid import uuid4
from fastapi import Request
from fastapi.responses import JSONResponse
from starlette.concurrency import iterate_in_threadpool
from starlette.middleware.base import BaseHTTPMiddleware
from starlette.responses import Response, StreamingResponse
from src.logging.logger import log_queue, request_id_ctx
class ProcessingTimeMiddleware(BaseHTTPMiddleware, ):
async def dispatch(self, request: Request, call_next)->Response:
start_time = perf_counter()
response = await call_next(request)
process_time = perf_counter() - start_time
response.headers["X-Process-Time"] = str(process_time)
return response
class LoggingMiddleware(BaseHTTPMiddleware):
async def build_log_line(self, request_id, method, path, status_code, detail, client_ip) -> str:
return f"[{request_id}] [{method}] [{path}] [{status_code}] [{detail}] [{client_ip}]"
async def dispatch(self, request: Request, call_next) -> Response:
request_id=str(uuid4())
request_id_ctx.set(request_id)
client_ip = request.headers.get('x-forwarded-for', '').split(',')[0].strip() or (request.client.host if request.client else 'unknown')
method = request.method
path=request.url.path
try:
response = await call_next(request)
except Exception as exc: # noqa: BLE001
line=await self.build_log_line(request_id=request_id, method=method, path=path, status_code=500, detail=repr(exc), client_ip=client_ip)
log_queue.put_nowait(("endpoints", line))
return JSONResponse(
status_code=500,
content={"detail": "Internal Server Error", "request_id": request_id}
)
streaming_response = cast(StreamingResponse, response)
chunks = []
async for chunk in streaming_response.body_iterator:
chunks.append(chunk.encode() if isinstance(chunk, str) else bytes(chunk))
body_bytes = b"".join(chunks)
streaming_response.body_iterator = iterate_in_threadpool(iter([body_bytes]))
try:
parsed = json.loads(body_bytes)
body = parsed.get("detail", None) if not isinstance(parsed, bool) else None
except (json.JSONDecodeError, TypeError):
body = None
line=await self.build_log_line(request_id=request_id, method=method, path=path, status_code=response.status_code, detail=body, client_ip=client_ip)
log_queue.put_nowait(("endpoints", line))
return response
+12 -62
View File
@@ -1,19 +1,12 @@
import asyncio
import json
import logging
from contextvars import ContextVar
from time import gmtime, perf_counter, strftime
from typing import cast
from uuid import uuid4
from time import gmtime, strftime
import aiofiles
from fastapi import Request
from fastapi.responses import JSONResponse
from starlette.concurrency import iterate_in_threadpool
from starlette.middleware.base import BaseHTTPMiddleware
from starlette.responses import Response, StreamingResponse
request_id_ctx: ContextVar[str] = ContextVar("request_id", default="-")
message_id_ctx: ContextVar[str] = ContextVar("message_id", default="-")
log_queue=asyncio.Queue()
@@ -34,60 +27,17 @@ class LogWriter:
await f.write(f"[{current_time}] {msg}\n")
class ProcessingTimeMiddleware(BaseHTTPMiddleware, ):
async def dispatch(self, request: Request, call_next)->Response:
start_time = perf_counter()
response = await call_next(request)
process_time = perf_counter() - start_time
response.headers["X-Process-Time"] = str(process_time)
return response
class LoggingMiddleware(BaseHTTPMiddleware):
async def build_log_line(self, request_id, method, path, status_code, detail, client_ip) -> str:
return f"[{request_id}] [{method}] [{path}] [{status_code}] [{detail}] [{client_ip}]"
async def dispatch(self, request: Request, call_next) -> Response:
request_id=str(uuid4())
request_id_ctx.set(request_id)
client_ip = request.headers.get('x-forwarded-for', '').split(',')[0].strip() or (request.client.host if request.client else 'unknown')
method = request.method
path=request.url.path
try:
response = await call_next(request)
except Exception as exc: # noqa: BLE001
line=await self.build_log_line(request_id=request_id, method=method, path=path, status_code=500, detail=repr(exc), client_ip=client_ip)
log_queue.put_nowait(("endpoints", line))
return JSONResponse(
status_code=500,
content={"detail": "Internal Server Error", "request_id": request_id}
)
streaming_response = cast(StreamingResponse, response)
chunks = []
async for chunk in streaming_response.body_iterator:
chunks.append(chunk.encode() if isinstance(chunk, str) else bytes(chunk))
body_bytes = b"".join(chunks)
streaming_response.body_iterator = iterate_in_threadpool(iter([body_bytes]))
try:
parsed = json.loads(body_bytes)
body = parsed.get("detail", None) if not isinstance(parsed, bool) else None
except (json.JSONDecodeError, TypeError):
body = None
line=await self.build_log_line(request_id=request_id, method=method, path=path, status_code=response.status_code, detail=body, client_ip=client_ip)
log_queue.put_nowait(("endpoints", line))
return response
class LoggerDB(logging.Handler, LogWriter):
class LoggerDB(logging.Handler):
def emit(self, record: logging.LogRecord) -> None:
msg = self.format(record)
rid=request_id_ctx.get()
log_queue.put_nowait(("sql",f"[{rid}] {msg}"))
log_queue.put_nowait(("sql",f"[{rid}] {msg}"))
class LoggerDaemon(logging.Handler):
def emit(self, record: logging.LogRecord)->None:
msg= self.format(record)
mid=message_id_ctx.get()
log_queue.put_nowait(("daemon", f"[{mid}], {msg}"))
+9 -6
View File
@@ -1,6 +1,9 @@
import json
import smtplib
from uuid import uuid4
from src.logging import daemon_logger
from src.logging.logger import message_id_ctx
from src.messaging.rabbitmq_client import rabbitmq_client
from src.service.email.email_welcome import DaemonEmailSender
@@ -19,10 +22,10 @@ class WelcomeEmailConsumer:
async def process_message(self, message) -> None:
message_id_ctx.set(str(uuid4()))
async with message.process(ignore_processed=True):
data = json.loads(message.body)
print(f"Обрабатываю: {data}, метка: {message.routing_key}")
daemon_logger.info(f"Обрабатываю: {data}, метка: {message.routing_key}")
try:
await self.daemon.send_email(data.get("email"))
except (
@@ -31,10 +34,10 @@ class WelcomeEmailConsumer:
TimeoutError,
ConnectionRefusedError,
) as exc:
print(f"transient error, retrying: {exc!r}")
daemon_logger.exception(f"transient error, retrying: {exc!r}")
await message.nack(requeue=True)
except Exception as exc: # noqa: BLE001
print(f"permanent error, sending to DLQ: {exc!r}")
daemon_logger.exception(f"permanent error, sending to DLQ: {exc!r}")
await message.nack(requeue=False)
@@ -65,10 +68,10 @@ class ResetEmailConsumer:
async def process_message(self, message) -> None:
message_id_ctx.set(str(uuid4()))
async with message.process():
data = json.loads(message.body)
print(f"Обрабатываю: {data}, метка: {message.routing_key}")
daemon_logger.info(f"Обрабатываю: {data}, метка: {message.routing_key}")
async def start_consuming(self)->None:
+18 -7
View File
@@ -1,3 +1,4 @@
import asyncio
import aio_pika
from aio_pika.abc import AbstractChannel, AbstractRobustConnection
@@ -10,13 +11,23 @@ class RabbitMQClient:
self.channel: AbstractChannel | None = None
async def connect(self) -> None:
if self.connection is None or self.connection.is_closed:
self.connection = await aio_pika.connect_robust(
host=env_settings.RABBITMQ_HOST,
port=env_settings.RABBITMQ_PORT,
login=env_settings.RABBITMQ_LOGIN,
password=env_settings.RABBITMQ_PASSWORD,
)
for attempt in range(1,6):
try:
if self.connection is None or self.connection.is_closed:
self.connection = await aio_pika.connect_robust(
host=env_settings.RABBITMQ_HOST,
port=env_settings.RABBITMQ_PORT,
login=env_settings.RABBITMQ_LOGIN,
password=env_settings.RABBITMQ_PASSWORD,
)
break
except Exception as exc:
if attempt==5:
raise
print(f"RabbitMQ not ready yet (attempt {attempt}/5): {exc!r}, retrying...")
await asyncio.sleep(2**attempt)
if self.connection is None:
raise RuntimeError("Failed to esablish RabbitMQ connection. Check the server!")
if self.channel is None or self.channel.is_closed:
self.channel = await self.connection.channel()
-3
View File
@@ -74,10 +74,7 @@ async def logout(response:Response,
response.delete_cookie("refresh_token")
return await auth.logout(refresh_token, access_token)
from src.messaging.producers.producers import email_producer
@router.get("")
async def protected(current_user:UserOut=Depends(require_permissions()))->dict:
await email_producer.send_welcome_email(current_user.email)
return {"protected router": "Hello, this is a protected router"}
@@ -1,5 +1,6 @@
from fastapi import APIRouter, Depends
from src.messaging.producers.producers import email_producer
from src.models.pydantic_models.model import UserCreate, UserOut, UserUpdate
from src.service.users_crud.users_crud import CrudService, crud_service
from src.web.protected_routes.auth_routes import require_permissions
@@ -11,7 +12,8 @@ async def get_current_user_by_email(email:str, crud:CrudService=Depends(crud_ser
return await crud.get_user_by_email(email)
@router.post("/create_user")
async def create_user(data:UserCreate, crud:CrudService=Depends(crud_service), current_user=Depends(require_permissions("admin")))->UserOut:
async def create_user(data:UserCreate, crud:CrudService=Depends(crud_service), current_user=Depends(require_permissions("admin")))->UserOut:
await email_producer.send_welcome_email(current_user.email)
return await crud.create_user(data)
@router.post("/delete_user_soft")