import asyncio from aiohttp import web from aiohttp_apispec import docs, response_schema, json_schema from typing import Any, Dict, List import logging from sqlalchemy import delete, update, insert, func from sqlalchemy.future import select from server.artist.decorators import Helper from server.artist.constants import COUNTRIES_MAPPER from server.dna.models import UserSearch, RecentSearches from server.dna.category.models import Category, CategoryEntity from server.dna.models import RecentSearches from server.dna.schemas import ( UserSearchSchema, ProfileInfoResponse, ProfileResetResponse, LabelsWithUsersResponse, ErrorResponseSchema, RecentSearch, EntityTypeSchema, FavoritesRequestShema, ) from server.dna.utils import QueryWrapper, Utils from server.dna.helpers.recent_searches import RecentSearchesHelper from server.track.constants import GENRES_MAPPER log = logging.getLogger(__name__) @docs(tags=["dna", "user_search", "item"], description="Get user_search by id") @response_schema(UserSearchSchema.ResponseGetSchema) @Helper.args_decompose() @Helper.auth_with_user_id async def get_user_search(request: web.Request, id: int) -> web.Response: user_id = request["user_id"] filters = {"id": id, "user_id": user_id} query = select(UserSearch).filter_by(**filters) result = await QueryWrapper.select(query) return web.json_response(UserSearchSchema.ResponseGetSchema().dump(result)) @docs(tags=["dna", "user_search", "list"], description="Save user search") @response_schema(UserSearchSchema.ResponseGetSchema) @Helper.auth_with_user_id async def get_list_user_search(request: web.Request) -> web.Response: user_id = request["user_id"] filters = {"user_id": user_id} query = select(UserSearch).filter_by(**filters) results = await QueryWrapper.select(query, many=True) return web.json_response(UserSearchSchema.ResponseGetSchema().dump(results, many=True)) @docs(tags=["dna", "user_search", "save"], description="Save user search") @json_schema(UserSearchSchema.RequestSchema) @response_schema(UserSearchSchema.ResponseSchema) @Helper.auth_with_user_id async def save_user_search(request: web.Request) -> web.Response: data = await request.json() UserSearchSchema.RequestSchema().load(data=data) data["user_id"] = request["user_id"] query = insert(UserSearch).values(**data) result = await QueryWrapper.insert(query) if "error" in result: return web.json_response(status=web.HTTPBadRequest.status_code, data=ErrorResponseSchema().dump(result)) return web.json_response(status=web.HTTPCreated.status_code, data=UserSearchSchema.ResponseSchema().dump(result)) @docs(tags=["dna", "user_search", "update"], description="Update user_search") @response_schema(UserSearchSchema.ResponseSchema) @json_schema(UserSearchSchema.RequestSchema) @Helper.args_decompose() @Helper.auth_with_user_id async def update_user_search(request: web.Request, id: int) -> web.Response: data = await request.json() UserSearchSchema.RequestSchema().load(data=data) data["id"] = id user_id = request["user_id"] filters = {"id": id, "user_id": user_id} query = update(UserSearch).values(**data).filter_by(**filters) result = await QueryWrapper.update(query) if "error" in result: return web.json_response(status=web.HTTPBadRequest.status_code, data=ErrorResponseSchema().dump(result)) return web.json_response(UserSearchSchema.ResponseSchema().dump(result)) @docs(tags=["dna", "user_search", "archived"], description="Archive user_search") @response_schema(UserSearchSchema.RequestArchiveSchema) @Helper.args_decompose() @Helper.auth_with_user_id async def user_search_archived(request: web.Request, id: int) -> web.Response: data = await request.json() UserSearchSchema.RequestArchiveSchema().load(data=data) data["id"] = id archived = data["archived"] if archived: data["archived"] = func.now() user_id = request["user_id"] filters = {"id": id, "user_id": user_id} query = update(UserSearch).values(**data).filter_by(**filters) result = await QueryWrapper.update(query) return web.json_response(UserSearchSchema.ResponseSchema().dump(result)) @docs(tags=["dna", "user_search", "delete"], description="Delete user_search") @response_schema(UserSearchSchema.ResponseDeleteSchema) @Helper.args_decompose() @Helper.auth_with_user_id async def user_search_delete(request: web.Request, id: int) -> web.Response: user_id = request["user_id"] filters = {"id": id, "user_id": user_id} query = delete(UserSearch).filter_by(**filters) result = await QueryWrapper.delete(query) return web.json_response(UserSearchSchema.ResponseDeleteSchema().dump(result)) @docs(tags=["dna", "profile_info"], description="Get authorized user info") @response_schema(ProfileInfoResponse) @Helper.auth_with_user_id async def get_profile_info(request: web.Request) -> web.Response: user_id = request["user_id"] atlas_api = request.app["atlas_api"] profile_info, label_names = await asyncio.gather( *( asyncio.wait_for(atlas_api.get_profileinfo(user_id), timeout=5), asyncio.wait_for(atlas_api.get_label_names(user_id), timeout=5), ) ) profile_info["user_id"] = user_id profile_info["labels"] = label_names profile_info["roles"] = profile_info.get("dna/role", []) return web.json_response(ProfileInfoResponse().dump(profile_info)) @docs(tags=["dna", "item", "label", "profile"], description="Get label with profiles") @response_schema(LabelsWithUsersResponse) @Helper.auth_with_user_id @Helper.args_decompose() async def get_label_profiles(request: web.Request, id: int) -> web.Response: response = await Utils.prepare_profile_response(request, id) return web.json_response(LabelsWithUsersResponse().dump(response)) @docs(tags=["dna", "list", "labels", "profile"], description="Get list labels with profiles") @response_schema(LabelsWithUsersResponse) @Helper.auth_with_user_id async def get_list_labels_with_profiles(request: web.Request) -> web.Response: response = await Utils.prepare_profile_response(request) return web.json_response(LabelsWithUsersResponse().dump(response, many=True)) @docs(tags=["dna", "profile", "reset"], description="dna profile reset") @response_schema(ProfileResetResponse) @Helper.auth_with_user_id async def profile_reset(request: web.Request) -> web.Response: def form_response(result: Dict[str, Any]) -> web.Response: log.info("profile was reset") return web.json_response(data=ProfileResetResponse().dump(result)) user_id = request["user_id"] filters = {"user_id": user_id} async def delete_user_searches(): query = delete(UserSearch).filter_by(**filters) await QueryWrapper.delete(query) log.info(f"user_searches by {user_id} profle were removed") async def delete_recent_searches(): query = delete(RecentSearches).filter_by(**filters) await QueryWrapper.delete(query) log.info(f"recent_searches by {user_id} profle were removed") async def delete_categories(): data: dict = {"user_id": user_id, "is_deleted": func.now()} category_filters: dict = {"is_deleted": None, "user_id": user_id} query = update(Category).values(**data).filter_by(**category_filters) await QueryWrapper.update(query) log.info(f"categories by {user_id} profle were removed") async def delete_categories_entities(): category_entity_id_subquery = ( select(CategoryEntity.id).join(Category, Category.id == CategoryEntity.category_id).filter_by(**filters) ) query = delete(CategoryEntity).filter(CategoryEntity.id.in_(category_entity_id_subquery)) await QueryWrapper.delete(query) log.info(f"categories entities by profle {user_id} were removed") async def delete_favorites(): genres, countries = await asyncio.gather( *( asyncio.wait_for(Utils.get_account_favorites(request, "genres"), timeout=5), asyncio.wait_for(Utils.get_account_favorites(request, "countries"), timeout=5), ) ) log.debug(f"profle favorites genres {genres}") log.debug(f"profle favorites countries {countries}") genres_tasks = [ asyncio.create_task(Utils.delete_entity_from_account_favorites(request, "genres", genre["code"])) for genre in genres ] countries_tasks = [ asyncio.create_task(Utils.delete_entity_from_account_favorites(request, "countries", country["code"])) for country in countries ] tasks = genres_tasks + countries_tasks if tasks: await asyncio.gather(asyncio.wait(tasks, timeout=5)) log.info(f"profile favorites items were removed") await asyncio.gather( delete_user_searches(), delete_recent_searches(), delete_categories(), delete_categories_entities(), delete_favorites(), ) result = filters return form_response(result) @docs(tags=["dna", "list", "users", "recent_searches"], description="Get list users recent searches") @response_schema(RecentSearch.ResponseSchema) @Helper.auth_with_user_id async def get_list_recent_searches(request: web.Request) -> web.Response: def form_response(recent_searches): return web.json_response(RecentSearch.ResponseSchema().dump(recent_searches, many=True)) user_id = request["user_id"] filters = {"user_id": user_id} recent_searches = await RecentSearchesHelper.get_recent_searches(filters) return form_response(recent_searches) @docs(tags=["dna", "save", "users", "recent_searches"], description="Save users recent search") @json_schema(RecentSearch.RequestSchema) @response_schema(RecentSearch.CreatedRecentSearchResponseSchema) @Helper.auth_with_user_id async def create_recent_search(request: web.Request) -> web.Response: data = await request.json() data = RecentSearch.RequestSchema().load(data=data) data["user_id"] = request["user_id"] query = insert(RecentSearches).values(**data) result = await QueryWrapper.insert(query) if "error" in result: return web.json_response(status=web.HTTPBadRequest.status_code, data=ErrorResponseSchema().dump(result)) return web.json_response( status=web.HTTPCreated.status_code, data=RecentSearch.CreatedRecentSearchResponseSchema().dump(result) ) @docs(tags=["dna", "get", "list", "users", "favorites"], description="Get users list favorites entities") @Helper.args_decompose() @Helper.auth_with_user_id @Helper.check_request_url(EntityTypeSchema) async def list_favorites(request: web.Request, entity_type) -> web.Response: favorites: List[Dict[str, Any]] = await Utils.get_account_favorites(request, entity_type) return web.json_response(favorites) @docs(tags=["dna", "post", "entity", "favorites"], description="Add an entity to account favorites") @Helper.args_decompose() @Helper.auth_with_user_id @Helper.check_request_url(FavoritesRequestShema) async def add_entity_to_account_favorites(request: web.Request, entity_type: str, entity_id: str) -> web.Response: await Utils.add_entity_to_account_favorites(request, entity_type, entity_id) return web.json_response(status=web.HTTPCreated.status_code) @docs(tags=["dna", "delete", "entity", "favorites"], description="Delete an entity from account favorites") @Helper.args_decompose() @Helper.auth_with_user_id @Helper.check_request_url(FavoritesRequestShema) async def delete_entity_from_account_favorites(request: web.Request, entity_type: str, entity_id: str) -> web.Response: await Utils.delete_entity_from_account_favorites(request, entity_type, entity_id) return web.json_response(status=web.HTTPNoContent.status_code) @docs(tags=["track", "artist", "genres"], description="Genres list") @Helper.args_decompose() async def genres_list(request: web.Request) -> web.Response: return web.json_response(GENRES_MAPPER) @docs(tags=["track", "artist", "genres"], description="Countries list") @Helper.args_decompose() async def countries_list(request: web.Request) -> web.Response: return web.json_response(COUNTRIES_MAPPER)