from aiohttp import web from aiohttp_apispec import docs, headers_schema, json_schema, querystring_schema from apollo_utils.service.clients.aiohttp.utils.response import dump_response_schema from apollo_utils.service.exceptions import NotFound from datetime import datetime from sqlalchemy import all_, and_, or_ from server.constants.core import BASE_API_PREFIX from server.constants.messages import TYPE_MAP from server.db.models.messages import FeedMessage from server.db.models.users import Account from server.schemas.headers import AccountHeader from server.schemas.users.messages.feed import UsersMessagesFeed, UsersMessagesFeedView from server.utils.pagination import paginate router = web.RouteTableDef() @router.get(BASE_API_PREFIX + "/users/messages/feed/") @docs( tags=["feed messages"], summary="Endpoint to get all user's feed messages.", ) @headers_schema(AccountHeader) @querystring_schema(UsersMessagesFeed.Request) @dump_response_schema(UsersMessagesFeed.Response, apply=True) @paginate() async def get_feed(request: web.Request) -> web.Response: headers, params = request["headers"], request["querystring"] status = params.get("status") subject, search = params["subject"], params["search"] account = await Account.get(app_slug=headers["app"], user_id=headers["user_id"]) if not account: raise NotFound(f"Account for user {headers['user_id']} and application {headers['app']} is not found.") filters = [] search_filters = [] if status is not None: filters.append(FeedMessage.status == status) if subject: filters.append(FeedMessage.meta["subject"].astext.in_(subject)) if search: search_filters.extend( [ FeedMessage.data["content"]["track"]["name"].astext.ilike(all_([search])), FeedMessage.data["content"]["playlist"]["name"].astext.ilike(all_([search])), FeedMessage.data["content"]["chart"]["name"].astext.ilike(all_([search])), FeedMessage.data["content"]["track"]["artists_names"].astext.ilike(all_([search])), ] ) messages, count = await FeedMessage.list( order_by=( FeedMessage.created_at.desc(), FeedMessage.id.desc(), ), account_id=account.id, filters=[and_(*filters) & or_(*search_filters)], limit=params.get("limit"), offset=params.get("offset"), ) results = [] for item in messages: meta = item.meta data = item.data content = data["content"] topic = TYPE_MAP.get(meta["type"]).lower() created_at = item.created_at.replace(microsecond=0) message = { "id": item.id, "date": created_at, "topic": topic, "vendor": meta["dsp"].lower(), "is_new": item.status, "title": content.get("title", ""), "message": content["body"], "country_code": meta["country_code"], } extra_data = {} track = content.get("track") if track: track_data = { "track_id": track.get("id"), "track_name": track["name"], "isrc": track["isrc"], "artist_name": " ".join([i["name"] for i in track["artists"] if i["name"]]), } extra_data.update(track_data) playlist = content.get("playlist") if playlist and playlist.get("id"): playlist_data = { "playlist_id": playlist.get("id"), "playlist_name": playlist.get("name"), "playlist_image_url": playlist.get("image_url"), } extra_data.update(playlist_data) try: timedelta_obj = datetime.now() - created_at extra_data["time_delta"] = int(timedelta_obj.total_seconds()) except (ValueError, TypeError, AttributeError): extra_data["time_delta"] = None if topic in ["chart_removals", "starred_playlist_removals", "playlist_removals"]: position = content.get("previous_position") else: position = content.get("current_position") if position is not None: extra_data["position"] = position if topic in ["chart_major_moves", "starred_playlist_major_moves", "playlist_major_moves"]: extra_data["change"] = content.get("current_position", 0) - content.get("previous_position", 0) extra_data["target"] = content.get("chart", {}).get("name", content.get("playlist", {}).get("name")) message.update(extra_data) results.append(message) return results, count @router.post(BASE_API_PREFIX + "/users/messages/feed/view/") @docs( tags=[""], summary="Mark the passed one and previous user’s feed messages as read.", ) @headers_schema(AccountHeader) @json_schema(UsersMessagesFeedView.Request) @dump_response_schema(UsersMessagesFeedView.Response, apply=True) async def update_feed(request: web.Request) -> web.Response: headers, data = request["headers"], request["json"] account = await Account.get(app_slug=headers["app"], user_id=headers["user_id"]) if not account: raise NotFound(f"Account for user {headers['user_id']} and application {headers['app']} is not found.") account_id = account.id filters = [] if data.get("message_id"): filters.append(FeedMessage.id == data["message_id"]) elif data.get("related_id"): filters.append(FeedMessage.message_id == data["related_id"]) message = await FeedMessage.get(account_id=account_id, filters=filters) if message: await FeedMessage.update_bulk( {"status": False}, filters=[ or_(FeedMessage.id <= message.id, FeedMessage.created_at < message.created_at), FeedMessage.account_id == account_id, ], ) return web.Response()