diff --git a/app/logger.py b/app/logger.py new file mode 100644 index 0000000..dc392d8 --- /dev/null +++ b/app/logger.py @@ -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 diff --git a/app/main.py b/app/main.py index f3a45e0..c2eda16 100644 --- a/app/main.py +++ b/app/main.py @@ -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() diff --git a/app/middleware.py b/app/middleware.py new file mode 100644 index 0000000..c82bfaf --- /dev/null +++ b/app/middleware.py @@ -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 diff --git a/app/mongo_worker.py b/app/mongo_worker.py index 7de4381..a152c13 100644 --- a/app/mongo_worker.py +++ b/app/mongo_worker.py @@ -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: diff --git a/app/rabbit_worker.py b/app/rabbit_worker.py index 065f805..2c91475 100644 --- a/app/rabbit_worker.py +++ b/app/rabbit_worker.py @@ -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: @@ -20,9 +17,11 @@ class RabbitWorker: self.url = ( 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) diff --git a/app/requirements.txt b/app/requirements.txt index a146643..68f6a0b 100644 --- a/app/requirements.txt +++ b/app/requirements.txt @@ -5,4 +5,5 @@ python-dotenv==1.0.1 uvicorn==0.34.0 aio-pika==10.0.1 aiogram==3.18.0 -aiohttp==3.11.18 \ No newline at end of file +aiohttp==3.11.18 +loguru==0.7.3 \ No newline at end of file diff --git a/app/tg_bot.py b/app/tg_bot.py index 0201c3f..3de8424 100644 --- a/app/tg_bot.py +++ b/app/tg_bot.py @@ -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