Add logging and add mongodb connection pool

This commit is contained in:
IgorVolochay
2026-08-25 14:53:34 +03:00
parent 48fd5ea19e
commit b1cca197d5
7 changed files with 202 additions and 26 deletions
+102
View File
@@ -0,0 +1,102 @@
"""
Centralized logging configuration using Loguru.
Outputs structured JSON to stdout (INFO/WARNING) and stderr (ERROR/CRITICAL).
Designed for Docker + Grafana Loki / Promtail.
Usage:
from logger import logger, setup_logging
setup_logging() # call once at application entry point
logger.info("message")
"""
import os
import sys
import logging
from dotenv import load_dotenv
from loguru import logger
def get_log_level() -> str:
"""Reads LOG_LEVEL or log_level from .env or environment, defaults to 'INFO'."""
load_dotenv()
level = os.getenv("LOG_LEVEL") or os.getenv("log_level") or "INFO"
return level.strip().upper()
# ── Stdout / stderr filters ───────────────────────────────────────────────────
def _stdout_filter(record: dict) -> bool:
"""Pass DEBUG / INFO / WARNING to stdout."""
return record["level"].no < logging.ERROR
def _stderr_filter(record: dict) -> bool:
"""Pass ERROR / CRITICAL to stderr."""
return record["level"].no >= logging.ERROR
# ── Stdlib → Loguru bridge ────────────────────────────────────────────────────
class InterceptHandler(logging.Handler):
"""Redirect all stdlib logging calls into Loguru."""
def emit(self, record: logging.LogRecord) -> None:
try:
level = logger.level(record.levelname).name
except ValueError:
level = record.levelno
frame, depth = sys._getframe(6), 6
while frame and frame.f_code.co_filename == logging.__file__:
frame = frame.f_back
depth += 1
logger.opt(depth=depth, exception=record.exc_info).log(
level, record.getMessage()
)
# ── Public setup function ─────────────────────────────────────────────────────
def setup_logging(level: str | None = None) -> None:
"""
Configure Loguru sinks and intercept all stdlib loggers.
Call once at the very start of the application entry point.
"""
if not level:
level = get_log_level()
else:
level = level.strip().upper()
logger.remove() # remove default sink
common: dict = {
"level": level,
"serialize": True, # JSON output
"backtrace": False,
"diagnose": False,
}
# stdout — DEBUG / INFO / WARNING
logger.add(sys.stdout, filter=_stdout_filter, **common)
# stderr — ERROR / CRITICAL
logger.add(sys.stderr, filter=_stderr_filter, **{**common, "level": "ERROR"})
# Redirect all stdlib loggers (uvicorn, motor, aiogram, aio_pika …)
logging.basicConfig(handlers=[InterceptHandler()], level=0, force=True)
# Suppress noisy third-party loggers — we handle HTTP access via middleware
_quiet = {
"uvicorn.access": logging.WARNING, # replaced by our middleware
"motor": logging.WARNING,
"aio_pika": logging.WARNING,
"aiormq": logging.WARNING,
}
for name, lvl in _quiet.items():
_lib_logger = logging.getLogger(name)
_lib_logger.handlers = [InterceptHandler()]
_lib_logger.setLevel(lvl)
_lib_logger.propagate = False
+7 -5
View File
@@ -1,6 +1,5 @@
import os
import secrets
import logging
import uvicorn
import asyncio
@@ -16,9 +15,10 @@ from schemas.base_schemas import Card
from mongo_worker import MongoWorker
from rabbit_worker import RabbitWorker
from tools.base_moderation import moderate_text
from logger import logger, setup_logging
from middleware import RequestLoggingMiddleware
logger = logging.getLogger(__name__)
setup_logging()
load_dotenv()
disable_docs = os.getenv("DISABLE_DOCS", "true").lower() == "true"
@@ -56,6 +56,7 @@ config = SecurityConfig(
guard_deco = SecurityDecorator(config)
app.add_middleware(SecurityMiddleware, config=config)
app.add_middleware(RequestLoggingMiddleware)
app.state.guard_decorator = guard_deco
mongo_worker = MongoWorker()
_rabbit_worker: Optional[RabbitWorker] = None
@@ -84,6 +85,7 @@ async def verify_moderation_secret(
@app.on_event("startup")
async def startup_event():
await mongo_worker.create_indexes()
logger.info("Application started on :5000")
@app.get("/check_user", status_code=200)
@@ -170,7 +172,7 @@ async def add_card(
try:
await get_rabbit_worker().send_to_moderation(card)
except Exception as exc:
logger.error("Failed to send card %s to moderation queue: %s", card.card_id, exc)
logger.error("Failed to send card {} to moderation queue: {}", card.card_id, exc)
return BaseResponse(result=card)
response.status_code = status.HTTP_400_BAD_REQUEST
@@ -277,7 +279,7 @@ async def get_comments(
async def main():
config = uvicorn.Config("main:app", host="0.0.0.0", port=5000, log_level="info")
config = uvicorn.Config("main:app", host="0.0.0.0", port=5000, log_level="warning")
server = uvicorn.Server(config)
await server.serve()
+62
View File
@@ -0,0 +1,62 @@
"""
HTTP request/response logging middleware for FastAPI.
Normal requests (2xx/3xx):
INFO — method, path, query_params, status_code, duration_ms
Error responses (4xx):
WARNING — all above + request_body (truncated to 1000 chars)
Server errors (5xx):
ERROR — all above + request_body (truncated to 1000 chars)
"""
import time
from starlette.middleware.base import BaseHTTPMiddleware
from starlette.requests import Request
from starlette.responses import Response
from logger import logger
_BODY_METHODS = frozenset({"POST", "PUT", "PATCH"})
_BODY_MAX_LEN = 1000
class RequestLoggingMiddleware(BaseHTTPMiddleware):
async def dispatch(self, request: Request, call_next) -> Response:
start = time.perf_counter()
# Read body only for methods that carry a payload
body: str | None = None
if request.method in _BODY_METHODS:
raw = await request.body()
body = raw.decode(errors="replace")[:_BODY_MAX_LEN]
response = await call_next(request)
duration_ms = round((time.perf_counter() - start) * 1000, 1)
status = response.status_code
base_fields = {
"method": request.method,
"path": request.url.path,
"query": str(request.query_params) or None,
"status": status,
"duration_ms": duration_ms,
}
if status >= 500:
logger.bind(**base_fields, request_body=body).error(
"{method} {path}{status} ({duration_ms}ms)", **base_fields
)
elif status >= 400:
logger.bind(**base_fields, request_body=body).warning(
"{method} {path}{status} ({duration_ms}ms)", **base_fields
)
else:
logger.bind(**base_fields).info(
"{method} {path}{status} ({duration_ms}ms)", **base_fields
)
return response
+16 -5
View File
@@ -1,5 +1,4 @@
import os
import logging
import motor.motor_asyncio
@@ -10,9 +9,7 @@ from pymongo import ReturnDocument
from schemas.base_schemas import User, Visited, Card, Comment
from schemas.api_schemas import BaseResponse
logger = logging.getLogger(__name__)
from logger import logger
class MongoWorker:
@@ -25,7 +22,12 @@ class MongoWorker:
password=os.getenv('MONGO_PASS'),
serverSelectionTimeoutMS=5000,
connectTimeoutMS=5000,
maxPoolSize=50,
minPoolSize=5,
maxIdleTimeMS=60000,
waitQueueTimeoutMS=5000
)
logger.info("MongoDB connection established.")
self.db = self.client["data"]
self.users_data = self.db["users"]
self.visited_data = self.db["visited"]
@@ -58,6 +60,7 @@ class MongoWorker:
registration_date=datetime.now().isoformat(),
)
await self.users_data.insert_one(new_user.model_dump())
logger.debug("User added: user_id={}, username={}", user_id, username)
return new_user
async def get_user(self, user_id: int) -> User:
@@ -130,6 +133,7 @@ class MongoWorker:
creation_date=datetime.now().isoformat(),
)
await self.game_data.insert_one(new_card.model_dump())
logger.debug("Card created by API: card_id={}, author_id={}", new_card.card_id, author_id)
return new_card
async def add_card_by_base_model(self, new_card: Card) -> Optional[Card]:
@@ -138,7 +142,7 @@ class MongoWorker:
await self.game_data.insert_one(new_card.model_dump())
return new_card
except Exception as exc:
logger.error("Failed to insert card: %s", exc)
logger.error("Failed to insert card: {}", exc)
raise
async def accept_card(self, card_id: int) -> BaseResponse:
@@ -152,14 +156,18 @@ class MongoWorker:
return_document=ReturnDocument.AFTER,
)
if not result:
logger.debug("Attempted to accept non-existent card: card_id={}", card_id)
return BaseResponse(result="Card doesn't exist", error=True)
logger.debug("Card accepted: card_id={}", card_id)
return BaseResponse(result=Card.model_validate(result))
async def reject_card(self, card_id: int) -> BaseResponse:
"""Rejects a card: deletes it from the database."""
result = await self.game_data.delete_one({"card_id": card_id})
if result.deleted_count == 0:
logger.debug("Attempted to reject non-existent card: card_id={}", card_id)
return BaseResponse(result="Card doesn't exist", error=True)
logger.debug("Card rejected and deleted: card_id={}", card_id)
return BaseResponse(result=f"Card {card_id} rejected and deleted")
async def select_choice(self, card_id: int, choice: str) -> BaseResponse:
@@ -206,6 +214,7 @@ class MongoWorker:
{"user_id": user_id},
{"$push": {"liked_card_ids": card_id}},
)
logger.debug("Card liked: card_id={}, user_id={}", card_id, user_id)
return BaseResponse(result=True, error=False)
async def dislike_card(self, card_id: int, user_id: int) -> BaseResponse:
@@ -227,6 +236,7 @@ class MongoWorker:
{"user_id": user_id},
{"$push": {"disliked_card_ids": card_id}},
)
logger.debug("Card disliked: card_id={}, user_id={}", card_id, user_id)
return BaseResponse(result=True, error=False)
@@ -253,6 +263,7 @@ class MongoWorker:
if not updated_user:
return BaseResponse(result="Difficulty adding comment_id to user", error=True)
logger.debug("Comment added: comment_id={}, card_id={}, author_id={}", new_comment.comment_id, card_id, user_id)
return BaseResponse(result=new_comment)
async def get_comments(self, card_id: int) -> BaseResponse:
+6 -6
View File
@@ -1,7 +1,6 @@
import os
import json
import asyncio
import logging
from typing import Callable, Awaitable
import aio_pika
@@ -9,9 +8,7 @@ from aio_pika.abc import AbstractIncomingMessage
from dotenv import load_dotenv
from schemas.base_schemas import Card
logger = logging.getLogger(__name__)
from logger import logger
class RabbitWorker:
@@ -21,8 +18,10 @@ class RabbitWorker:
f"amqp://{os.getenv('RABBIT_USER')}:{os.getenv('RABBIT_PASS')}"
f"@{os.getenv('RABBIT_HOST')}:{os.getenv('RABBIT_PORT')}"
)
logger.info("RabbitWorker connection established.")
async def send_to_moderation(self, card: Card) -> None:
logger.debug("Preparing to send card {} to moderation queue...", card.card_id)
connection = await aio_pika.connect_robust(self.url)
async with connection:
channel = await connection.channel()
@@ -34,7 +33,8 @@ class RabbitWorker:
),
routing_key="moderation",
)
logger.info("Card %s sent to moderation queue", card.card_id)
logger.debug("Card {} successfully published to moderation queue", card.card_id)
logger.info("Card {} sent to moderation queue", card.card_id)
async def consume_moderation(
self,
@@ -55,7 +55,7 @@ class RabbitWorker:
card = Card.model_validate(card_data)
await callback(card)
except Exception as exc:
logger.error("Error processing moderation message: %s", exc)
logger.error("Error processing moderation message: {}", exc)
await queue.consume(on_message)
+1
View File
@@ -6,3 +6,4 @@ uvicorn==0.34.0
aio-pika==10.0.1
aiogram==3.18.0
aiohttp==3.11.18
loguru==0.7.3
+6 -8
View File
@@ -11,7 +11,6 @@ When a button is pressed, the bot calls protected endpoints
import os
import json
import asyncio
import logging
from datetime import datetime
import aiohttp
@@ -21,18 +20,17 @@ from aiogram.types import CallbackQuery, InlineKeyboardButton, InlineKeyboardMar
from schemas.base_schemas import Card
from rabbit_worker import RabbitWorker
from logger import logger, setup_logging
load_dotenv()
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
setup_logging()
# ── Configuration ──────────────────────────────────────────────
BOT_TOKEN = os.getenv("TG_BOT_TOKEN")
ADMIN_CHAT_ID = int(os.getenv("TG_ADMIN_CHAT_ID", "0"))
API_BASE_URL = os.getenv("API_BASE_URL", "http://localhost:5000")
MODERATION_SECRET = os.getenv("MODERATION_SECRET", "change-me-in-production")
bot = Bot(token=BOT_TOKEN)
dp = Dispatcher()
rabbit = RabbitWorker()
@@ -52,7 +50,7 @@ async def _get_author_username(author_id: int) -> str:
if username:
return f"@{username}"
except Exception as exc:
logger.warning("Failed to fetch username for %s: %s", author_id, exc)
logger.warning("Failed to fetch username for {}: {}", author_id, exc)
return str(author_id)
@@ -97,7 +95,7 @@ async def send_card_to_admin(card: Card) -> None:
reply_markup=keyboard,
parse_mode="HTML",
)
logger.info("Sent card %s to admin chat", card.card_id)
logger.info("Sent card {} to admin chat", card.card_id)
# ── Calling protected API endpoints ──────────────────────────
@@ -131,7 +129,7 @@ async def on_accept(callback: CallbackQuery) -> None:
parse_mode="HTML",
)
await callback.answer("Карточка принята!")
logger.info("Card %s accepted by admin", card_id)
logger.info("Card {} accepted by admin", card_id)
@dp.callback_query(F.data.startswith("reject:"))
@@ -148,7 +146,7 @@ async def on_reject(callback: CallbackQuery) -> None:
parse_mode="HTML",
)
await callback.answer("Карточка отклонена!")
logger.info("Card %s rejected by admin", card_id)
logger.info("Card {} rejected by admin", card_id)
# ── aiogram Lifecycle hooks ────────────────────────────────────
_rabbit_task: asyncio.Task | None = None