Rewrite sync pymongo to async motor

This commit is contained in:
IgorVolochay
2026-08-20 13:40:38 +03:00
parent 3e87607886
commit 122209bdf4
3 changed files with 327 additions and 290 deletions
+192 -166
View File
@@ -1,22 +1,31 @@
import os
import logging
import pymongo
import motor.motor_asyncio
from datetime import datetime
from dotenv import load_dotenv
from typing import Optional
from pymongo import ReturnDocument
from schemas.base_schemas import *
from schemas.api_schemas import *
from schemas.base_schemas import User, Visited, Card, Comment
from schemas.api_schemas import BaseResponse
logger = logging.getLogger(__name__)
class MongoWorker:
def __init__(self):
load_dotenv()
self.client = pymongo.MongoClient(host = os.getenv('MONGO_HOST'),
port = int(os.getenv('MONGO_PORT')),
username = os.getenv('MONGO_USER'),
password = os.getenv('MONGO_PASS'))
self.client = motor.motor_asyncio.AsyncIOMotorClient(
host=os.getenv('MONGO_HOST'),
port=int(os.getenv('MONGO_PORT', 27017)),
username=os.getenv('MONGO_USER'),
password=os.getenv('MONGO_PASS'),
serverSelectionTimeoutMS=5000,
connectTimeoutMS=5000,
)
self.db = self.client["data"]
self.users_data = self.db["users"]
self.visited_data = self.db["visited"]
@@ -24,191 +33,208 @@ class MongoWorker:
self.game_data = self.db["cards"]
self.comments_data = self.db["comments"]
async def create_indexes(self) -> None:
"""Создаёт индексы при старте приложения."""
await self.users_data.create_index("user_id", unique=True)
await self.game_data.create_index("card_id", unique=True)
await self.game_data.create_index("active_status")
await self.visited_data.create_index("user_id", unique=True)
await self.comments_data.create_index("comment_id", unique=True)
logger.info("MongoDB indexes created.")
def check_user(self, user_id: int) -> bool:
if self.users_data.find_one({"user_id": user_id}):
return True
else:
return False
def add_user(self, user_id: int, username: str, first_name: str, last_name: str, photo_url: str) -> User:
new_user = User(user_id=user_id,
username=username,
first_name=first_name,
last_name=last_name,
photo_url=photo_url,
registration_date=datetime.now().isoformat())
try:
self.users_data.insert_one(new_user.model_dump())
return new_user
except Exception as exception:
return new_user
async def check_user(self, user_id: int) -> bool:
document = await self.users_data.find_one({"user_id": user_id}, {"_id": 1})
return document is not None
def get_user(self, user_id: int) -> User:
return User.model_validate(self.users_data.find_one({"user_id": user_id}))
async def add_user(
self, user_id: int, username: str, first_name: str, last_name: str, photo_url: str) -> User:
new_user = User(
user_id=user_id,
username=username,
first_name=first_name,
last_name=last_name,
photo_url=photo_url,
registration_date=datetime.now().isoformat(),
)
await self.users_data.insert_one(new_user.model_dump())
return new_user
def get_and_update_counter(self, counter_name: str) -> int:
counter = self.counters.find_one_and_update(
async def get_user(self, user_id: int) -> User:
document = await self.users_data.find_one({"user_id": user_id})
return User.model_validate(document)
async def get_and_update_counter(self, counter_name: str) -> int:
"""Атомарно инкрементирует счётчик и возвращает новое значение."""
counter = await self.counters.find_one_and_update(
{"counter_name": counter_name},
{"$inc": {"counter": 1}},
upsert=True,
return_document=True)
return_document=ReturnDocument.AFTER,
)
return counter["counter"]
def get_visited_cards(self, user_id: int) -> BaseResponse:
document = self.visited_data.find_one({"user_id": user_id})
async def get_visited_cards(self, user_id: int) -> BaseResponse:
document = await self.visited_data.find_one({"user_id": user_id})
if not document:
check_user = self.check_user(user_id)
if check_user:
return BaseResponse(result=Visited(user_id=user_id,
cards_visited=set()))
else:
return BaseResponse(result="User doesn't exist", error=True)
else:
return BaseResponse(result=Visited.model_validate(document))
def filter_cards(self, random_cards: list[Card], cards_visited: set) -> tuple[list[Card], list[int]]:
filtered_cards = [card for card in random_cards if card.card_id not in cards_visited]
filtered_cards_id = [filtered_card.card_id for filtered_card in filtered_cards]
return filtered_cards, filtered_cards_id
def update_visited_cards(self, user_id: int, visited_card_id: int) -> Visited:
update_visited = self.visited_data.find_one_and_update({"user_id": user_id},
{"$addToSet": {"cards_visited": visited_card_id}},
upsert=True,
return_document=True)
return Visited.model_validate(update_visited)
if await self.check_user(user_id):
return BaseResponse(result=Visited(user_id=user_id, cards_visited=set()))
return BaseResponse(result="User doesn't exist", error=True)
return BaseResponse(result=Visited.model_validate(document))
async def update_visited_cards(self, user_id: int, visited_card_id: int) -> Visited:
updated = await self.visited_data.find_one_and_update(
{"user_id": user_id},
{"$addToSet": {"cards_visited": visited_card_id}},
upsert=True,
return_document=ReturnDocument.AFTER,
)
return Visited.model_validate(updated)
def add_card_by_api(self, choice_A: str, choice_B: str, author_id: int) -> Card:
new_card = Card(card_id=self.get_and_update_counter(counter_name="card"),
choice_A=choice_A,
choice_B=choice_B,
author_id=author_id,
creation_date=datetime.now().isoformat())
try:
self.game_data.insert_one(new_card.model_dump())
return new_card
except Exception as exception:
print(exception)
return new_card
def add_card_by_base_model(self, new_card: Card) -> Optional[Card]:
new_card.card_id = self.get_and_update_counter(counter_name="card")
try:
self.game_data.insert_one(new_card.model_dump())
return new_card
except Exception as exception:
print(exception)
return new_card
def get_card(self, card_id: int) -> Optional[Card]:
document = self.game_data.find_one({"card_id": card_id})
async def get_card(self, card_id: int) -> Optional[Card]:
document = await self.game_data.find_one({"card_id": card_id})
if document:
return Card.model_validate(document)
else:
return None
return None
def get_random_cards(self, amount: int, active_status: bool) -> Optional[list[Card]]:
pipeline = [{"$match": {"active_status": active_status}},
{"$sample": {"size": amount}}]
raw_items = list(self.game_data.aggregate(pipeline))
async def get_random_cards(self, amount: int,active_status: bool,exclude_ids: Optional[set[int]] = None,) -> Optional[list[Card]]:
"""Возвращает случайные карточки, исключая уже просмотренные (одним запросом)."""
match_filter: dict = {"active_status": active_status}
if exclude_ids:
match_filter["card_id"] = {"$nin": list(exclude_ids)}
pipeline = [
{"$match": match_filter},
{"$sample": {"size": amount}},
]
raw_items = await self.game_data.aggregate(pipeline).to_list(length=amount)
if raw_items:
validated_items = [Card.model_validate(item) for item in raw_items]
return validated_items
else:
return None
return [Card.model_validate(item) for item in raw_items]
return None
def select_choice(self, card_id: int, choice: str) -> BaseResponse:
def filter_cards(self, random_cards: list[Card], cards_visited: set) -> tuple[list[Card], list[int]]:
filtered_cards = [card for card in random_cards if card.card_id not in cards_visited]
filtered_cards_id = [card.card_id for card in filtered_cards]
return filtered_cards, filtered_cards_id
async def add_card_by_api(self, choice_A: str, choice_B: str, author_id: int) -> Card:
new_card = Card(
card_id=await self.get_and_update_counter(counter_name="card"),
choice_A=choice_A,
choice_B=choice_B,
author_id=author_id,
creation_date=datetime.now().isoformat(),
)
await self.game_data.insert_one(new_card.model_dump())
return new_card
async def add_card_by_base_model(self, new_card: Card) -> Optional[Card]:
new_card.card_id = await self.get_and_update_counter(counter_name="card")
try:
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)
raise
async def select_choice(self, card_id: int, choice: str) -> BaseResponse:
if choice == "A":
count_choice = "count_choice_A"
count_field = "count_choice_A"
elif choice == "B":
count_choice = "count_choice_B"
count_field = "count_choice_B"
else:
return BaseResponse(result="Wrong choice", error=True)
result = self.game_data.find_one_and_update({"card_id": card_id},
{"$inc": {"count_total": 1, count_choice: 1}})
result = await self.game_data.find_one_and_update(
{"card_id": card_id},
{"$inc": {"count_total": 1, count_field: 1}},
)
if not result:
return BaseResponse(result="Card doesn't exist", error=True)
else:
return BaseResponse(result=result, error=False)
def check_user_reactions(self, user_id: int, card_id: int) -> BaseResponse:
user_info: User = self.get_user(user_id)
liked_card_ids: list = user_info.liked_card_ids
disliked_card_ids: list = user_info.disliked_card_ids
return BaseResponse(result=True, error=False)
if card_id in liked_card_ids:
async def check_user_reactions(self, user_id: int, card_id: int) -> BaseResponse:
user_info: User = await self.get_user(user_id)
if card_id in user_info.liked_card_ids:
return BaseResponse(result="Card already liked", error=True)
elif card_id in disliked_card_ids:
if card_id in user_info.disliked_card_ids:
return BaseResponse(result="Card already disliked", error=True)
else:
return BaseResponse(result="No reactions", error=False)
return BaseResponse(result="No reactions", error=False)
def like_card(self, card_id: int, user_id: int) -> BaseResponse:
if self.check_user(user_id):
user_reaction = self.check_user_reactions(user_id, card_id)
if user_reaction.error:
return user_reaction
update_card_info = self.game_data.find_one_and_update({"card_id": card_id},
{"$inc": {"count_likes": 1}})
if not update_card_info:
return BaseResponse(result="Card doesn't exist", error=True)
add_card_to_user = self.users_data.update_one({'user_id': user_id},
{'$push': {'liked_card_ids': card_id}})
if not add_card_to_user:
return BaseResponse(result="User doesn't exist", error=True)
else:
return BaseResponse(result=True, error=False)
else:
async def like_card(self, card_id: int, user_id: int) -> BaseResponse:
if not await self.check_user(user_id):
return BaseResponse(result="User doesn't exist", error=True)
def dislike_card(self, card_id: int, user_id: int) -> BaseResponse:
if self.check_user(user_id):
user_reaction = self.check_user_reactions(user_id, card_id)
if user_reaction.error:
return user_reaction
update_card_info = self.game_data.find_one_and_update({"card_id": card_id},
{"$inc": {"count_dislikes": 1}})
if not update_card_info:
return BaseResponse(result="Card doesn't exist", error=True)
add_card_to_user = self.users_data.update_one({'user_id': user_id},
{'$push': {'disliked_card_ids': card_id}})
if not add_card_to_user:
return BaseResponse(result="User doesn't exist", error=True)
else:
return BaseResponse(result=True, error=False)
else:
user_reaction = await self.check_user_reactions(user_id, card_id)
if user_reaction.error:
return user_reaction
updated_card = await self.game_data.find_one_and_update(
{"card_id": card_id},
{"$inc": {"count_likes": 1}},
)
if not updated_card:
return BaseResponse(result="Card doesn't exist", error=True)
await self.users_data.update_one(
{"user_id": user_id},
{"$push": {"liked_card_ids": card_id}},
)
return BaseResponse(result=True, error=False)
async def dislike_card(self, card_id: int, user_id: int) -> BaseResponse:
if not await self.check_user(user_id):
return BaseResponse(result="User doesn't exist", error=True)
def add_comment(self, user_id: int, card_id: int, comment_text: str) -> BaseResponse:
if self.check_user(user_id):
if self.get_card(card_id):
new_comment = Comment(comment_id=self.get_and_update_counter(counter_name="comment"),
author_id=user_id,
card_id=card_id,
commet_text=comment_text,
creation_date=datetime.now().isoformat())
result = self.comments_data.insert_one(new_comment.model_dump())
if result:
update_user_comments = self.users_data.find_one_and_update({"user_id": user_id},
{"$addToSet": {"comments_ids": new_comment.comment_id}})
if update_user_comments:
return BaseResponse(result=new_comment)
else:
return BaseResponse(result="Difficulty adding comment_id to user", error=True)
else:
return BaseResponse(result="Add comment error", error=True)
else:
return BaseResponse(result="Card doesn't exist", error=True)
else:
return BaseResponse(result="User doesn't exist", error=True)
user_reaction = await self.check_user_reactions(user_id, card_id)
if user_reaction.error:
return user_reaction
updated_card = await self.game_data.find_one_and_update(
{"card_id": card_id},
{"$inc": {"count_dislikes": 1}},
)
if not updated_card:
return BaseResponse(result="Card doesn't exist", error=True)
await self.users_data.update_one(
{"user_id": user_id},
{"$push": {"disliked_card_ids": card_id}},
)
return BaseResponse(result=True, error=False)
async def add_comment(self, user_id: int, card_id: int, comment_text: str) -> BaseResponse:
if not await self.check_user(user_id):
return BaseResponse(result="User doesn't exist", error=True)
if not await self.get_card(card_id):
return BaseResponse(result="Card doesn't exist", error=True)
new_comment = Comment(
comment_id=await self.get_and_update_counter(counter_name="comment"),
author_id=user_id,
card_id=card_id,
commet_text=comment_text,
creation_date=datetime.now().isoformat(),
)
await self.comments_data.insert_one(new_comment.model_dump())
updated_user = await self.users_data.find_one_and_update(
{"user_id": user_id},
{"$addToSet": {"comments_ids": new_comment.comment_id}},
return_document=ReturnDocument.AFTER,
)
if not updated_user:
return BaseResponse(result="Difficulty adding comment_id to user", error=True)
return BaseResponse(result=new_comment)
async def get_comments(self, card_id: int) -> BaseResponse:
comments = await self.comments_data.find({"card_id": card_id}).sort("creation_date", -1).to_list(length=None)
comments = [Comment.model_validate(comment) for comment in comments]
return BaseResponse(result=comments)