From 12a8782627101925d912763c713be221ff16aec7 Mon Sep 17 00:00:00 2001 From: izakdvlpr Date: Thu, 9 Jul 2026 18:09:31 -0300 Subject: [PATCH 1/8] feat: add activity and notification table --- .../migration.sql | 171 +++++++++++- prisma/schema.prisma | 254 +++++++++++++----- prisma/seed.ts | 14 +- 3 files changed, 357 insertions(+), 82 deletions(-) rename prisma/migrations/{20260707215819_init => 20260709183715_init}/migration.sql (82%) diff --git a/prisma/migrations/20260707215819_init/migration.sql b/prisma/migrations/20260709183715_init/migration.sql similarity index 82% rename from prisma/migrations/20260707215819_init/migration.sql rename to prisma/migrations/20260709183715_init/migration.sql index 079637f..fb24fe9 100644 --- a/prisma/migrations/20260707215819_init/migration.sql +++ b/prisma/migrations/20260709183715_init/migration.sql @@ -8,10 +8,13 @@ CREATE TYPE "UserTier" AS ENUM ('Tracker', 'Archivist', 'ArchiveMaster'); CREATE TYPE "CommentType" AS ENUM ('Anime', 'Manga', 'TVShow', 'Movie', 'Game', 'Book', 'Profile'); -- CreateEnum -CREATE TYPE "ReactionType" AS ENUM ('Comment', 'FeedEvent', 'GameReview', 'AnimeReview', 'MangaReview', 'TvShowReview', 'MovieReview', 'BookReview'); +CREATE TYPE "ReactionType" AS ENUM ('Comment', 'Activity', 'GameReview', 'AnimeReview', 'MangaReview', 'TvShowReview', 'MovieReview', 'BookReview'); -- CreateEnum -CREATE TYPE "FeedEventType" AS ENUM ('NewFollower', 'NewFavorite', 'NewList', 'NewListItem', 'NewReview', 'NewWatch', 'NewProgress'); +CREATE TYPE "NotificationType" AS ENUM ('System', 'CommentOnProfile', 'ReactionOnComment', 'ReactionOnActivity', 'ReactionOnAnimeReview', 'ReactionOnMangaReview', 'ReactionOnTvShowReview', 'ReactionOnMovieReview', 'ReactionOnGameReview', 'ReactionOnBookReview'); + +-- CreateEnum +CREATE TYPE "ActivityType" AS ENUM ('AccountCreated', 'ListCreated', 'ListItemAdded', 'FavoriteAdded', 'ReviewAdded', 'ProgressStarted', 'ProgressCompleted', 'Watched', 'Followed', 'MedalEarned'); -- CreateEnum CREATE TYPE "WatchEpisodeStatus" AS ENUM ('NotWatched', 'Watching', 'Completed', 'Paused', 'Dropped', 'Planning'); @@ -171,7 +174,7 @@ CREATE TABLE "Reaction" ( "type" "ReactionType" NOT NULL, "userId" TEXT NOT NULL, "commentId" TEXT, - "feedEventId" TEXT, + "activityId" TEXT, "gameReviewId" TEXT, "animeReviewId" TEXT, "mangaReviewId" TEXT, @@ -184,16 +187,56 @@ CREATE TABLE "Reaction" ( ); -- CreateTable -CREATE TABLE "FeedEvent" ( +CREATE TABLE "Notification" ( "id" TEXT NOT NULL, - "type" "FeedEventType" NOT NULL, - "userId" TEXT NOT NULL, + "type" "NotificationType" NOT NULL, + "recipientId" TEXT NOT NULL, + "actorId" TEXT, "metadata" JSONB, - "entityIds" TEXT[], - "count" INTEGER NOT NULL DEFAULT 1, + "readAt" TIMESTAMP(3), "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "profileId" TEXT, + "commentId" TEXT, + "reactionId" TEXT, + "activityId" TEXT, + "animeReviewId" TEXT, + "mangaReviewId" TEXT, + "tvShowReviewId" TEXT, + "movieReviewId" TEXT, + "gameReviewId" TEXT, + "bookReviewId" TEXT, - CONSTRAINT "FeedEvent_pkey" PRIMARY KEY ("id") + CONSTRAINT "Notification_pkey" PRIMARY KEY ("id") +); + +-- CreateTable +CREATE TABLE "Activity" ( + "id" TEXT NOT NULL, + "type" "ActivityType" NOT NULL, + "userId" TEXT NOT NULL, + "metadata" JSONB, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "listId" TEXT, + "listItemId" TEXT, + "favoriteId" TEXT, + "animeReviewId" TEXT, + "mangaReviewId" TEXT, + "tvShowReviewId" TEXT, + "movieReviewId" TEXT, + "gameReviewId" TEXT, + "bookReviewId" TEXT, + "animeProgressId" TEXT, + "mangaProgressId" TEXT, + "tvShowProgressId" TEXT, + "movieProgressId" TEXT, + "gameProgressId" TEXT, + "bookProgressId" TEXT, + "animeEpisodeWatchId" TEXT, + "tvShowEpisodeWatchId" TEXT, + "followingId" TEXT, + "userMedalId" TEXT, + + CONSTRAINT "Activity_pkey" PRIMARY KEY ("id") ); -- CreateTable @@ -842,7 +885,7 @@ CREATE INDEX "Comment_profileId_idx" ON "Comment"("profileId"); CREATE UNIQUE INDEX "Reaction_userId_commentId_key" ON "Reaction"("userId", "commentId"); -- CreateIndex -CREATE UNIQUE INDEX "Reaction_userId_feedEventId_key" ON "Reaction"("userId", "feedEventId"); +CREATE UNIQUE INDEX "Reaction_userId_activityId_key" ON "Reaction"("userId", "activityId"); -- CreateIndex CREATE UNIQUE INDEX "Reaction_userId_gameReviewId_key" ON "Reaction"("userId", "gameReviewId"); @@ -863,7 +906,16 @@ CREATE UNIQUE INDEX "Reaction_userId_movieReviewId_key" ON "Reaction"("userId", CREATE UNIQUE INDEX "Reaction_userId_bookReviewId_key" ON "Reaction"("userId", "bookReviewId"); -- CreateIndex -CREATE INDEX "FeedEvent_userId_idx" ON "FeedEvent"("userId"); +CREATE INDEX "Notification_recipientId_readAt_idx" ON "Notification"("recipientId", "readAt"); + +-- CreateIndex +CREATE INDEX "Notification_recipientId_createdAt_idx" ON "Notification"("recipientId", "createdAt"); + +-- CreateIndex +CREATE INDEX "Activity_userId_idx" ON "Activity"("userId"); + +-- CreateIndex +CREATE INDEX "Activity_type_idx" ON "Activity"("type"); -- CreateIndex CREATE UNIQUE INDEX "Game_igdbId_key" ON "Game"("igdbId"); @@ -1037,7 +1089,7 @@ ALTER TABLE "Reaction" ADD CONSTRAINT "Reaction_userId_fkey" FOREIGN KEY ("userI ALTER TABLE "Reaction" ADD CONSTRAINT "Reaction_commentId_fkey" FOREIGN KEY ("commentId") REFERENCES "Comment"("id") ON DELETE CASCADE ON UPDATE CASCADE; -- AddForeignKey -ALTER TABLE "Reaction" ADD CONSTRAINT "Reaction_feedEventId_fkey" FOREIGN KEY ("feedEventId") REFERENCES "FeedEvent"("id") ON DELETE CASCADE ON UPDATE CASCADE; +ALTER TABLE "Reaction" ADD CONSTRAINT "Reaction_activityId_fkey" FOREIGN KEY ("activityId") REFERENCES "Activity"("id") ON DELETE CASCADE ON UPDATE CASCADE; -- AddForeignKey ALTER TABLE "Reaction" ADD CONSTRAINT "Reaction_gameReviewId_fkey" FOREIGN KEY ("gameReviewId") REFERENCES "GameReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; @@ -1058,7 +1110,100 @@ ALTER TABLE "Reaction" ADD CONSTRAINT "Reaction_movieReviewId_fkey" FOREIGN KEY ALTER TABLE "Reaction" ADD CONSTRAINT "Reaction_bookReviewId_fkey" FOREIGN KEY ("bookReviewId") REFERENCES "BookReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; -- AddForeignKey -ALTER TABLE "FeedEvent" ADD CONSTRAINT "FeedEvent_userId_fkey" FOREIGN KEY ("userId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE; +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_recipientId_fkey" FOREIGN KEY ("recipientId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_actorId_fkey" FOREIGN KEY ("actorId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_profileId_fkey" FOREIGN KEY ("profileId") REFERENCES "Profile"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_commentId_fkey" FOREIGN KEY ("commentId") REFERENCES "Comment"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_reactionId_fkey" FOREIGN KEY ("reactionId") REFERENCES "Reaction"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_activityId_fkey" FOREIGN KEY ("activityId") REFERENCES "Activity"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_animeReviewId_fkey" FOREIGN KEY ("animeReviewId") REFERENCES "AnimeReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_mangaReviewId_fkey" FOREIGN KEY ("mangaReviewId") REFERENCES "MangaReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_tvShowReviewId_fkey" FOREIGN KEY ("tvShowReviewId") REFERENCES "TVShowReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_movieReviewId_fkey" FOREIGN KEY ("movieReviewId") REFERENCES "MovieReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_gameReviewId_fkey" FOREIGN KEY ("gameReviewId") REFERENCES "GameReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Notification" ADD CONSTRAINT "Notification_bookReviewId_fkey" FOREIGN KEY ("bookReviewId") REFERENCES "BookReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_userId_fkey" FOREIGN KEY ("userId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_listId_fkey" FOREIGN KEY ("listId") REFERENCES "List"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_listItemId_fkey" FOREIGN KEY ("listItemId") REFERENCES "ListItem"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_favoriteId_fkey" FOREIGN KEY ("favoriteId") REFERENCES "Favorite"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_animeReviewId_fkey" FOREIGN KEY ("animeReviewId") REFERENCES "AnimeReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_mangaReviewId_fkey" FOREIGN KEY ("mangaReviewId") REFERENCES "MangaReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_tvShowReviewId_fkey" FOREIGN KEY ("tvShowReviewId") REFERENCES "TVShowReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_movieReviewId_fkey" FOREIGN KEY ("movieReviewId") REFERENCES "MovieReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_gameReviewId_fkey" FOREIGN KEY ("gameReviewId") REFERENCES "GameReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_bookReviewId_fkey" FOREIGN KEY ("bookReviewId") REFERENCES "BookReview"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_animeProgressId_fkey" FOREIGN KEY ("animeProgressId") REFERENCES "AnimeProgress"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_mangaProgressId_fkey" FOREIGN KEY ("mangaProgressId") REFERENCES "MangaProgress"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_tvShowProgressId_fkey" FOREIGN KEY ("tvShowProgressId") REFERENCES "TVShowProgress"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_movieProgressId_fkey" FOREIGN KEY ("movieProgressId") REFERENCES "MovieProgress"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_gameProgressId_fkey" FOREIGN KEY ("gameProgressId") REFERENCES "GameProgress"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_bookProgressId_fkey" FOREIGN KEY ("bookProgressId") REFERENCES "BookProgress"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_animeEpisodeWatchId_fkey" FOREIGN KEY ("animeEpisodeWatchId") REFERENCES "AnimeEpisodeWatch"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_tvShowEpisodeWatchId_fkey" FOREIGN KEY ("tvShowEpisodeWatchId") REFERENCES "TVShowEpisodeWatch"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_followingId_fkey" FOREIGN KEY ("followingId") REFERENCES "Following"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_userMedalId_fkey" FOREIGN KEY ("userMedalId") REFERENCES "UserMedal"("id") ON DELETE CASCADE ON UPDATE CASCADE; -- AddForeignKey ALTER TABLE "AnimeEpisodeWatch" ADD CONSTRAINT "AnimeEpisodeWatch_userId_fkey" FOREIGN KEY ("userId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE; diff --git a/prisma/schema.prisma b/prisma/schema.prisma index e815a6a..6b20e4b 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -43,7 +43,7 @@ model User { reactions Reaction[] followers Following[] @relation("UserFollowing") following Following[] @relation("UserFollowers") - feedEvents FeedEvent[] + activities Activity[] animeEpisodesWatches AnimeEpisodeWatch[] lists List[] tvshowEpisodesWatches TvShowEpisodeWatch[] @@ -62,6 +62,8 @@ model User { favorites Favorite[] userMedals UserMedal[] payments Payment[] + notificationsReceived Notification[] @relation("NotificationRecipient") + notificationsSent Notification[] @relation("NotificationActor") @@unique([email]) @@unique([username]) @@ -126,8 +128,9 @@ model Profile { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - comments Comment[] + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + comments Comment[] + notifications Notification[] } model Medal { @@ -149,8 +152,9 @@ model UserMedal { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - medal Medal @relation(fields: [medalId], references: [id], onDelete: Cascade) + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + medal Medal @relation(fields: [medalId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([userId, medalId]) } @@ -162,8 +166,9 @@ model Following { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - follower User @relation("UserFollowers", fields: [followerId], references: [id], onDelete: Cascade) - following User @relation("UserFollowing", fields: [followingId], references: [id], onDelete: Cascade) + follower User @relation("UserFollowers", fields: [followerId], references: [id], onDelete: Cascade) + following User @relation("UserFollowing", fields: [followingId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([followerId, followingId]) @@index([followingId]) @@ -204,6 +209,8 @@ model Comment { profile Profile? @relation(fields: [profileId], references: [id], onDelete: Cascade) reactions Reaction[] + notifications Notification[] + @@index([animeId]) @@index([mangaId]) @@index([tvShowId]) @@ -215,7 +222,7 @@ model Comment { enum ReactionType { Comment - FeedEvent + Activity GameReview AnimeReview MangaReview @@ -230,7 +237,7 @@ model Reaction { type ReactionType userId String commentId String? - feedEventId String? + activityId String? gameReviewId String? animeReviewId String? mangaReviewId String? @@ -241,7 +248,7 @@ model Reaction { user User @relation(fields: [userId], references: [id], onDelete: Cascade) comment Comment? @relation(fields: [commentId], references: [id], onDelete: Cascade) - feedEvent FeedEvent? @relation(fields: [feedEventId], references: [id], onDelete: Cascade) + activity Activity? @relation(fields: [activityId], references: [id], onDelete: Cascade) gameReview GameReview? @relation(fields: [gameReviewId], references: [id], onDelete: Cascade) animeReview AnimeReview? @relation(fields: [animeReviewId], references: [id], onDelete: Cascade) mangaReview MangaReview? @relation(fields: [mangaReviewId], references: [id], onDelete: Cascade) @@ -249,8 +256,10 @@ model Reaction { movieReview MovieReview? @relation(fields: [movieReviewId], references: [id], onDelete: Cascade) bookReview BookReview? @relation(fields: [bookReviewId], references: [id], onDelete: Cascade) + notifications Notification[] + @@unique([userId, commentId]) - @@unique([userId, feedEventId]) + @@unique([userId, activityId]) @@unique([userId, gameReviewId]) @@unique([userId, animeReviewId]) @@unique([userId, mangaReviewId]) @@ -259,29 +268,123 @@ model Reaction { @@unique([userId, bookReviewId]) } -enum FeedEventType { - NewFollower - NewFavorite - NewList - NewListItem - NewReview - NewWatch - NewProgress +enum NotificationType { + System + CommentOnProfile + ReactionOnComment + ReactionOnActivity + ReactionOnAnimeReview + ReactionOnMangaReview + ReactionOnTvShowReview + ReactionOnMovieReview + ReactionOnGameReview + ReactionOnBookReview +} + +model Notification { + id String @id @default(uuid(7)) + type NotificationType + recipientId String + // Null for System notifications: nobody triggered them. + actorId String? + // Free-form payload for System notifications (title, description, url). + metadata Json? + readAt DateTime? + createdAt DateTime @default(now()) + profileId String? + commentId String? + reactionId String? + activityId String? + animeReviewId String? + mangaReviewId String? + tvShowReviewId String? + movieReviewId String? + gameReviewId String? + bookReviewId String? + + recipient User @relation("NotificationRecipient", fields: [recipientId], references: [id], onDelete: Cascade) + actor User? @relation("NotificationActor", fields: [actorId], references: [id], onDelete: Cascade) + profile Profile? @relation(fields: [profileId], references: [id], onDelete: Cascade) + comment Comment? @relation(fields: [commentId], references: [id], onDelete: Cascade) + reaction Reaction? @relation(fields: [reactionId], references: [id], onDelete: Cascade) + activity Activity? @relation(fields: [activityId], references: [id], onDelete: Cascade) + animeReview AnimeReview? @relation(fields: [animeReviewId], references: [id], onDelete: Cascade) + mangaReview MangaReview? @relation(fields: [mangaReviewId], references: [id], onDelete: Cascade) + tvShowReview TvShowReview? @relation(fields: [tvShowReviewId], references: [id], onDelete: Cascade) + movieReview MovieReview? @relation(fields: [movieReviewId], references: [id], onDelete: Cascade) + gameReview GameReview? @relation(fields: [gameReviewId], references: [id], onDelete: Cascade) + bookReview BookReview? @relation(fields: [bookReviewId], references: [id], onDelete: Cascade) + + @@index([recipientId, readAt]) + @@index([recipientId, createdAt]) } -model FeedEvent { - id String @id @default(uuid(7)) - type FeedEventType +enum ActivityType { + AccountCreated + ListCreated + ListItemAdded + FavoriteAdded + ReviewAdded + ProgressStarted + ProgressCompleted + Watched + Followed + MedalEarned +} + +model Activity { + id String @id @default(uuid(7)) + type ActivityType userId String metadata Json? - entityIds String[] - count Int @default(1) - createdAt DateTime @default(now()) + createdAt DateTime @default(now()) - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - reactions Reaction[] + listId String? + listItemId String? + favoriteId String? + animeReviewId String? + mangaReviewId String? + tvShowReviewId String? + movieReviewId String? + gameReviewId String? + bookReviewId String? + animeProgressId String? + mangaProgressId String? + tvShowProgressId String? + movieProgressId String? + gameProgressId String? + bookProgressId String? + animeEpisodeWatchId String? + tvShowEpisodeWatchId String? + followingId String? + userMedalId String? + + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + reactions Reaction[] + list List? @relation(fields: [listId], references: [id], onDelete: Cascade) + listItem ListItem? @relation(fields: [listItemId], references: [id], onDelete: Cascade) + favorite Favorite? @relation(fields: [favoriteId], references: [id], onDelete: Cascade) + animeReview AnimeReview? @relation(fields: [animeReviewId], references: [id], onDelete: Cascade) + mangaReview MangaReview? @relation(fields: [mangaReviewId], references: [id], onDelete: Cascade) + tvShowReview TvShowReview? @relation(fields: [tvShowReviewId], references: [id], onDelete: Cascade) + movieReview MovieReview? @relation(fields: [movieReviewId], references: [id], onDelete: Cascade) + gameReview GameReview? @relation(fields: [gameReviewId], references: [id], onDelete: Cascade) + bookReview BookReview? @relation(fields: [bookReviewId], references: [id], onDelete: Cascade) + animeProgress AnimeProgress? @relation(fields: [animeProgressId], references: [id], onDelete: Cascade) + mangaProgress MangaProgress? @relation(fields: [mangaProgressId], references: [id], onDelete: Cascade) + tvShowProgress TvShowProgress? @relation(fields: [tvShowProgressId], references: [id], onDelete: Cascade) + movieProgress MovieProgress? @relation(fields: [movieProgressId], references: [id], onDelete: Cascade) + gameProgress GameProgress? @relation(fields: [gameProgressId], references: [id], onDelete: Cascade) + bookProgress BookProgress? @relation(fields: [bookProgressId], references: [id], onDelete: Cascade) + animeEpisodeWatch AnimeEpisodeWatch? @relation(fields: [animeEpisodeWatchId], references: [id], onDelete: Cascade) + tvShowEpisodeWatch TvShowEpisodeWatch? @relation(fields: [tvShowEpisodeWatchId], references: [id], onDelete: Cascade) + following Following? @relation(fields: [followingId], references: [id], onDelete: Cascade) + userMedal UserMedal? @relation(fields: [userMedalId], references: [id], onDelete: Cascade) + + notifications Notification[] @@index([userId]) + @@index([type]) } model Game { @@ -591,8 +694,9 @@ model AnimeEpisodeWatch { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - anime Anime @relation(fields: [animeId], references: [id], onDelete: Cascade) + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + anime Anime @relation(fields: [animeId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([userId, animeId, episode]) @@index([userId, status]) @@ -608,8 +712,9 @@ model TvShowEpisodeWatch { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - tvShow TvShow @relation(fields: [tvShowId], references: [id], onDelete: Cascade) + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + tvShow TvShow @relation(fields: [tvShowId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([userId, tvShowId, season, episode]) @@index([userId, status]) @@ -640,8 +745,9 @@ model AnimeProgress { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - anime Anime @relation(fields: [animeId], references: [id], onDelete: Cascade) + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + anime Anime @relation(fields: [animeId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([userId, animeId]) } @@ -658,8 +764,9 @@ model MangaProgress { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - manga Manga @relation(fields: [mangaId], references: [id], onDelete: Cascade) + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + manga Manga @relation(fields: [mangaId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([userId, mangaId]) } @@ -676,8 +783,9 @@ model TvShowProgress { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - tvShow TvShow @relation(fields: [tvShowId], references: [id], onDelete: Cascade) + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + tvShow TvShow @relation(fields: [tvShowId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([userId, tvShowId]) @@map("TVShowProgress") @@ -694,8 +802,9 @@ model MovieProgress { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - movie Movie @relation(fields: [movieId], references: [id], onDelete: Cascade) + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + movie Movie @relation(fields: [movieId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([userId, movieId]) } @@ -711,8 +820,9 @@ model GameProgress { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - game Game @relation(fields: [gameId], references: [id], onDelete: Cascade) + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + game Game @relation(fields: [gameId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([userId, gameId]) } @@ -729,8 +839,9 @@ model BookProgress { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - book Book @relation(fields: [bookId], references: [id], onDelete: Cascade) + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + book Book @relation(fields: [bookId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([userId, bookId]) } @@ -753,9 +864,11 @@ model AnimeReview { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - anime Anime @relation(fields: [animeId], references: [id], onDelete: Cascade) - reactions Reaction[] + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + anime Anime @relation(fields: [animeId], references: [id], onDelete: Cascade) + reactions Reaction[] + activities Activity[] + notifications Notification[] @@unique([userId, animeId]) } @@ -775,9 +888,11 @@ model MangaReview { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - manga Manga @relation(fields: [mangaId], references: [id], onDelete: Cascade) - reactions Reaction[] + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + manga Manga @relation(fields: [mangaId], references: [id], onDelete: Cascade) + reactions Reaction[] + activities Activity[] + notifications Notification[] @@unique([userId, mangaId]) } @@ -797,9 +912,11 @@ model TvShowReview { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - tvShow TvShow @relation(fields: [tvShowId], references: [id], onDelete: Cascade) - reactions Reaction[] + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + tvShow TvShow @relation(fields: [tvShowId], references: [id], onDelete: Cascade) + reactions Reaction[] + activities Activity[] + notifications Notification[] @@unique([userId, tvShowId]) @@map("TVShowReview") @@ -820,9 +937,11 @@ model MovieReview { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - movie Movie @relation(fields: [movieId], references: [id], onDelete: Cascade) - reactions Reaction[] + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + movie Movie @relation(fields: [movieId], references: [id], onDelete: Cascade) + reactions Reaction[] + activities Activity[] + notifications Notification[] @@unique([userId, movieId]) } @@ -847,6 +966,8 @@ model GameReview { game Game @relation(fields: [gameId], references: [id], onDelete: Cascade) gameReviewScreenshots GameReviewScreenshot[] reactions Reaction[] + activities Activity[] + notifications Notification[] @@unique([userId, gameId]) } @@ -879,9 +1000,11 @@ model BookReview { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - book Book @relation(fields: [bookId], references: [id], onDelete: Cascade) - reactions Reaction[] + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + book Book @relation(fields: [bookId], references: [id], onDelete: Cascade) + reactions Reaction[] + activities Activity[] + notifications Notification[] @@unique([userId, bookId]) } @@ -904,8 +1027,9 @@ model List { createdAt DateTime @default(now()) updatedAt DateTime @updatedAt - user User @relation(fields: [userId], references: [id], onDelete: Cascade) - listItems ListItem[] + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + listItems ListItem[] + activities Activity[] @@unique([userId, name]) } @@ -926,8 +1050,9 @@ model ListItem { manga Manga? @relation(fields: [mangaId], references: [id], onDelete: Cascade) tvShow TvShow? @relation(fields: [tvShowId], references: [id], onDelete: Cascade) movie Movie? @relation(fields: [movieId], references: [id], onDelete: Cascade) - game Game? @relation(fields: [gameId], references: [id], onDelete: Cascade) - book Book? @relation(fields: [bookId], references: [id], onDelete: Cascade) + game Game? @relation(fields: [gameId], references: [id], onDelete: Cascade) + book Book? @relation(fields: [bookId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([listId, animeId]) @@unique([listId, mangaId]) @@ -965,8 +1090,9 @@ model Favorite { manga Manga? @relation(fields: [mangaId], references: [id], onDelete: Cascade) tvShow TvShow? @relation(fields: [tvShowId], references: [id], onDelete: Cascade) movie Movie? @relation(fields: [movieId], references: [id], onDelete: Cascade) - game Game? @relation(fields: [gameId], references: [id], onDelete: Cascade) - book Book? @relation(fields: [bookId], references: [id], onDelete: Cascade) + game Game? @relation(fields: [gameId], references: [id], onDelete: Cascade) + book Book? @relation(fields: [bookId], references: [id], onDelete: Cascade) + activities Activity[] @@unique([userId, animeId]) @@unique([userId, mangaId]) diff --git a/prisma/seed.ts b/prisma/seed.ts index 73774e6..255ea8d 100644 --- a/prisma/seed.ts +++ b/prisma/seed.ts @@ -46,8 +46,10 @@ export async function populateMedals(prisma: PrismaClient) { ], skipDuplicates: true, }); - - console.log(`Inserted ${medals.count} medals.`); + + if (medals.count > 0) { + console.log(`Inserted ${medals.count} medals.`); + } } async function createFirstUser(prisma: PrismaClient) { @@ -76,7 +78,7 @@ async function createFirstUser(prisma: PrismaClient) { }); if (!userExists) { - await prisma.user.create({ + const createdUser = await prisma.user.create({ data: { id: userData.id, email: userData.email, @@ -103,8 +105,10 @@ async function createFirstUser(prisma: PrismaClient) { insertedCount++; } } - - console.log(`Inserted ${insertedCount} users.`); + + if (insertedCount > 0) { + console.log(`Inserted ${insertedCount} users.`); + } } main() From 4cbb058a121a287cd33e2fd37daab394944bb12a Mon Sep 17 00:00:00 2001 From: izakdvlpr Date: Thu, 9 Jul 2026 18:09:41 -0300 Subject: [PATCH 2/8] feat: create add modal script --- scripts/add-medal.tsx | 74 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 74 insertions(+) create mode 100644 scripts/add-medal.tsx diff --git a/scripts/add-medal.tsx b/scripts/add-medal.tsx new file mode 100644 index 0000000..df0e1ec --- /dev/null +++ b/scripts/add-medal.tsx @@ -0,0 +1,74 @@ +import "dotenv/config"; + +import { PrismaPg } from "@prisma/adapter-pg"; +import { Pool } from "pg"; + +import { ActivityType, PrismaClient } from "../prisma/generated/client"; + +// Usage: bun run scripts/add-medal.tsx +const [userId, medalId] = process.argv.slice(2); + +if (!medalId || !userId) { + console.log("Usage: bun run scripts/add-medal.tsx "); + process.exit(1); +} + +const connectionString = `${process.env.DATABASE_URL}`; +const pool = new Pool({ connectionString }); + +const adapter = new PrismaPg(pool); +const prisma = new PrismaClient({ adapter }); + +async function main() { + const medal = await prisma.medal.findUnique({ where: { id: medalId } }); + + if (!medal) { + console.error(`Medal not found: ${medalId}`); + process.exit(1); + } + + const user = await prisma.user.findUnique({ where: { id: userId } }); + + if (!user) { + console.error(`User not found: ${userId}`); + process.exit(1); + } + + const alreadyHasMedal = await prisma.userMedal.findFirst({ + where: { userId, medalId }, + }); + + if (alreadyHasMedal) { + console.log(`User "${user.name}" already has medal "${medal.name}".`); + return; + } + + const userMedal = await prisma.userMedal.create({ + data: { userId, medalId }, + }); + + await prisma.activity.create({ + data: { + type: ActivityType.MedalEarned, + userId, + userMedalId: userMedal.id, + metadata: { id: userMedal.id, medal: { ...medal } }, + }, + }); + + console.log(`Medal "${medal.name}" granted to user "${user.name}".`); +} + +main() + .then(async () => { + await prisma.$disconnect(); + await pool.end(); + }) + .catch(async (err) => { + console.error(err); + + await prisma.$disconnect(); + await pool.end(); + + process.exit(1); + }); From 4740e837d7e8ff14fbd9e923f7a142836bed6b83 Mon Sep 17 00:00:00 2001 From: izakdvlpr Date: Thu, 9 Jul 2026 18:10:39 -0300 Subject: [PATCH 3/8] feat: create activity module --- src/app.module.ts | 6 +- src/modules/activity/activity.module.ts | 11 + src/modules/activity/activity.utils.ts | 18 + .../controller/activity.controller.ts | 48 +++ src/modules/activity/dto/activity.dto.ts | 100 ++++++ .../get-activities-by-user-following.dto.ts | 5 + .../dto/get-activities-by-user.dto.ts | 5 + .../dto/get-activities.dto.ts} | 2 +- .../activity/service/activity.service.ts | 336 ++++++++++++++++++ .../controller/feed-event.controller.ts | 30 -- src/modules/feed-event/dto/feed-event.dto.ts | 40 --- .../dto/get-feed-events-by-user.dto.ts | 7 - src/modules/feed-event/feed-event.module.ts | 11 - .../feed-event/service/feed-event.service.ts | 109 ------ .../queue/processors/activity.processor.ts | 58 +++ .../queue/processors/feed-event.processor.ts | 115 ------ .../activity.controller.spec.ts} | 2 +- .../queue/activity.processor.spec.ts} | 2 +- ...rvice.spec.ts => activity.service.spec.ts} | 2 +- 19 files changed, 589 insertions(+), 318 deletions(-) create mode 100644 src/modules/activity/activity.module.ts create mode 100644 src/modules/activity/activity.utils.ts create mode 100644 src/modules/activity/controller/activity.controller.ts create mode 100644 src/modules/activity/dto/activity.dto.ts create mode 100644 src/modules/activity/dto/get-activities-by-user-following.dto.ts create mode 100644 src/modules/activity/dto/get-activities-by-user.dto.ts rename src/modules/{feed-event/dto/get-feed-events.dto.ts => activity/dto/get-activities.dto.ts} (60%) create mode 100644 src/modules/activity/service/activity.service.ts delete mode 100644 src/modules/feed-event/controller/feed-event.controller.ts delete mode 100644 src/modules/feed-event/dto/feed-event.dto.ts delete mode 100644 src/modules/feed-event/dto/get-feed-events-by-user.dto.ts delete mode 100644 src/modules/feed-event/feed-event.module.ts delete mode 100644 src/modules/feed-event/service/feed-event.service.ts create mode 100644 src/shared/infra/queue/processors/activity.processor.ts delete mode 100644 src/shared/infra/queue/processors/feed-event.processor.ts rename test/unit/{infra/queue/feed-event.processor.spec.ts => controllers/activity.controller.spec.ts} (74%) rename test/unit/{controllers/feed-event.controller.spec.ts => infra/queue/activity.processor.spec.ts} (74%) rename test/unit/services/{feed-event.service.spec.ts => activity.service.spec.ts} (75%) diff --git a/src/app.module.ts b/src/app.module.ts index 734b60c..561a7a2 100644 --- a/src/app.module.ts +++ b/src/app.module.ts @@ -4,16 +4,17 @@ import { ConfigModule } from "@nestjs/config"; import { APP_GUARD, APP_INTERCEPTOR } from "@nestjs/core"; import { JwtModule } from "@nestjs/jwt"; import { ThrottlerModule } from "@nestjs/throttler"; +import { ActivityModule } from "./modules/activity/activity.module"; import { AnimeModule } from "./modules/anime/anime.module"; import { AuthModule } from "./modules/auth/auth.module"; import { BookModule } from "./modules/book/book.module"; import { CommentModule } from "./modules/comment/comment.module"; import { FavoriteModule } from "./modules/favorite/favorite.module"; -import { FeedEventModule } from "./modules/feed-event/feed-event.module"; import { GameModule } from "./modules/game/game.module"; import { ListModule } from "./modules/list/list.module"; import { MangaModule } from "./modules/manga/manga.module"; import { MovieModule } from "./modules/movie/movie.module"; +import { NotificationModule } from "./modules/notification/notification.module"; import { PaymentModule } from "./modules/payment/payment.module"; import { ProfileModule } from "./modules/profile/profile.module"; import { ReactionModule } from "./modules/reaction/reaction.module"; @@ -47,7 +48,7 @@ import { MetricsInterceptor } from "./shared/interceptors/metrics.interceptor"; DatabaseModule, HealthModule, AuthModule, - FeedEventModule, + ActivityModule, CacheModule, IntegrationsModule, UploadModule, @@ -55,6 +56,7 @@ import { MetricsInterceptor } from "./shared/interceptors/metrics.interceptor"; ProfileModule, CommentModule, ReactionModule, + NotificationModule, GameModule, MovieModule, TVShowModule, diff --git a/src/modules/activity/activity.module.ts b/src/modules/activity/activity.module.ts new file mode 100644 index 0000000..6d1f91e --- /dev/null +++ b/src/modules/activity/activity.module.ts @@ -0,0 +1,11 @@ +import { Module } from "@nestjs/common"; +import { ActivityController } from "./controller/activity.controller"; +import { ActivityService } from "./service/activity.service"; + +@Module({ + imports: [], + controllers: [ActivityController], + providers: [ActivityService], + exports: [ActivityService], +}) +export class ActivityModule {} diff --git a/src/modules/activity/activity.utils.ts b/src/modules/activity/activity.utils.ts new file mode 100644 index 0000000..405f8a7 --- /dev/null +++ b/src/modules/activity/activity.utils.ts @@ -0,0 +1,18 @@ +import { ActivityType, ProgressStatus } from "@prisma/generated/enums"; + +const STARTED_STATUSES: ProgressStatus[] = [ProgressStatus.Watching, ProgressStatus.Playing, ProgressStatus.Reading]; + +// Maps a progress status to the activity it should emit. +// Only "started" (Watching/Playing/Reading) and Completed generate activities; +// every other status returns null (no activity). +export function activityTypeFromProgressStatus(status: ProgressStatus): ActivityType | null { + if (status === ProgressStatus.Completed) { + return ActivityType.ProgressCompleted; + } + + if (STARTED_STATUSES.includes(status)) { + return ActivityType.ProgressStarted; + } + + return null; +} diff --git a/src/modules/activity/controller/activity.controller.ts b/src/modules/activity/controller/activity.controller.ts new file mode 100644 index 0000000..afd467b --- /dev/null +++ b/src/modules/activity/controller/activity.controller.ts @@ -0,0 +1,48 @@ +import { Controller, Get, Param, Query, UseGuards } from "@nestjs/common"; +import { ApiTags } from "@nestjs/swagger"; +import { AuthGuard, Session, type UserSession } from "@thallesp/nestjs-better-auth"; +import { GetActivitiesDto } from "../dto/get-activities.dto"; +import { GetActivitiesByUserDto } from "../dto/get-activities-by-user.dto"; +import { ActivityService } from "../service/activity.service"; +import { GetActivitiesByUserFollowingDto } from '../dto/get-activities-by-user-following.dto'; + +@ApiTags("Activity") +@Controller("/activities") +export class ActivityController { + constructor(private readonly activityService: ActivityService) {} + + @Get("/global") + async getActivities(@Query() query: GetActivitiesDto) { + const activities = await this.activityService.getActivities(query); + + return { activities }; + } + + @Get("/user/:userId") + async getActivitiesByUserId(@Param("userId") userId: string, @Query() query: GetActivitiesByUserDto) { + const activities = await this.activityService.getActivitiesByUserId({ + ...query, + userId, + }); + + return { activities }; + } + + @Get("/user/:userId/calendar") + async getUserActivityCalendarById(@Param("userId") userId: string) { + const activityCalendar = await this.activityService.getUserActivityCalendarById(userId); + + return { activityCalendar }; + } + + @Get("/following") + @UseGuards(AuthGuard) + async getActivitiesByUserFollowing(@Session() session: UserSession, @Query() query: GetActivitiesByUserFollowingDto) { + const activities = await this.activityService.getActivitiesByUserFollowing({ + ...query, + userId: session.user.id, + }); + + return { activities }; + } +} diff --git a/src/modules/activity/dto/activity.dto.ts b/src/modules/activity/dto/activity.dto.ts new file mode 100644 index 0000000..c89ac87 --- /dev/null +++ b/src/modules/activity/dto/activity.dto.ts @@ -0,0 +1,100 @@ +import { ApiProperty, ApiPropertyOptional } from "@nestjs/swagger"; +import { ActivityType } from "@prisma/generated/enums"; +import { IsEnum, IsNotEmpty, IsOptional, IsUUID } from "class-validator"; + +export interface ActivityMetadata { + readonly id: string; + readonly [key: string]: any; +} + +export class CreateActivityDto { + @IsEnum(ActivityType) + @IsNotEmpty() + @ApiProperty({ enum: ActivityType }) + readonly type: ActivityType; + + @IsNotEmpty() + @IsUUID() + @ApiProperty({ type: "string", format: "uuid" }) + readonly userId: string; + + @IsOptional() + @IsUUID("7") + readonly listId?: string; + + @IsOptional() + @IsUUID("7") + readonly listItemId?: string; + + @IsOptional() + @IsUUID("7") + readonly favoriteId?: string; + + @IsOptional() + @IsUUID("7") + readonly animeReviewId?: string; + + @IsOptional() + @IsUUID("7") + readonly mangaReviewId?: string; + + @IsOptional() + @IsUUID("7") + readonly tvShowReviewId?: string; + + @IsOptional() + @IsUUID("7") + readonly movieReviewId?: string; + + @IsOptional() + @IsUUID("7") + readonly gameReviewId?: string; + + @IsOptional() + @IsUUID("7") + readonly bookReviewId?: string; + + @IsOptional() + @IsUUID("7") + readonly animeProgressId?: string; + + @IsOptional() + @IsUUID("7") + readonly mangaProgressId?: string; + + @IsOptional() + @IsUUID("7") + readonly tvShowProgressId?: string; + + @IsOptional() + @IsUUID("7") + readonly movieProgressId?: string; + + @IsOptional() + @IsUUID("7") + readonly gameProgressId?: string; + + @IsOptional() + @IsUUID("7") + readonly bookProgressId?: string; + + @IsOptional() + @IsUUID("7") + readonly animeEpisodeWatchId?: string; + + @IsOptional() + @IsUUID("7") + readonly tvShowEpisodeWatchId?: string; + + @IsOptional() + @IsUUID("7") + readonly followingId?: string; + + @IsOptional() + @IsUUID("7") + readonly userMedalId?: string; + + @IsOptional() + @ApiPropertyOptional({ type: "object", additionalProperties: true }) + readonly metadata?: ActivityMetadata; +} diff --git a/src/modules/activity/dto/get-activities-by-user-following.dto.ts b/src/modules/activity/dto/get-activities-by-user-following.dto.ts new file mode 100644 index 0000000..78ec999 --- /dev/null +++ b/src/modules/activity/dto/get-activities-by-user-following.dto.ts @@ -0,0 +1,5 @@ +import { OffsetPaginationParamsDto } from "@/shared/infra/database/dtos/offset-pagination.dto"; + +export class GetActivitiesByUserFollowingDto extends OffsetPaginationParamsDto { + readonly userId: string; +} diff --git a/src/modules/activity/dto/get-activities-by-user.dto.ts b/src/modules/activity/dto/get-activities-by-user.dto.ts new file mode 100644 index 0000000..85256bc --- /dev/null +++ b/src/modules/activity/dto/get-activities-by-user.dto.ts @@ -0,0 +1,5 @@ +import { OffsetPaginationParamsDto } from "@/shared/infra/database/dtos/offset-pagination.dto"; + +export class GetActivitiesByUserDto extends OffsetPaginationParamsDto { + readonly userId: string; +} diff --git a/src/modules/feed-event/dto/get-feed-events.dto.ts b/src/modules/activity/dto/get-activities.dto.ts similarity index 60% rename from src/modules/feed-event/dto/get-feed-events.dto.ts rename to src/modules/activity/dto/get-activities.dto.ts index bf054be..4e3e3a8 100644 --- a/src/modules/feed-event/dto/get-feed-events.dto.ts +++ b/src/modules/activity/dto/get-activities.dto.ts @@ -1,3 +1,3 @@ import { OffsetPaginationParamsDto } from "@/shared/infra/database/dtos/offset-pagination.dto"; -export class GetFeedEventsDto extends OffsetPaginationParamsDto {} +export class GetActivitiesDto extends OffsetPaginationParamsDto {} diff --git a/src/modules/activity/service/activity.service.ts b/src/modules/activity/service/activity.service.ts new file mode 100644 index 0000000..4d6905f --- /dev/null +++ b/src/modules/activity/service/activity.service.ts @@ -0,0 +1,336 @@ +import { Injectable } from "@nestjs/common"; +import { ActivityFindManyArgs } from "@prisma/generated/models"; +import { ERROR_CODES } from "@/shared/constants/error-codes"; +import { AppException } from "@/shared/exceptions/app.exceptions"; +import { DatabaseService } from "@/shared/infra/database/database.service"; +import { CreateActivityDto } from "../dto/activity.dto"; +import { GetActivitiesDto } from "../dto/get-activities.dto"; +import { GetActivitiesByUserDto } from "../dto/get-activities-by-user.dto"; +import { GetActivitiesByUserFollowingDto } from '../dto/get-activities-by-user-following.dto'; + +const SOURCE_FIELDS = [ + "listId", + "listItemId", + "favoriteId", + "animeReviewId", + "mangaReviewId", + "tvShowReviewId", + "movieReviewId", + "gameReviewId", + "bookReviewId", + "animeProgressId", + "mangaProgressId", + "tvShowProgressId", + "movieProgressId", + "gameProgressId", + "bookProgressId", + "animeEpisodeWatchId", + "tvShowEpisodeWatchId", + "followingId", + "userMedalId", +] as const; + +// Polymorphic media selects — same field mapping used by favorite.service.ts: +// anime/manga -> title + imageUrl, tvShow -> name + posterUrl, +// movie -> title + posterUrl, game -> name + coverUrl, book -> title + imageUrl. +const ANIME_SELECT = { select: { id: true, malId: true, title: true, imageUrl: true } }; +const MANGA_SELECT = { select: { id: true, malId: true, title: true, imageUrl: true } }; +const TVSHOW_SELECT = { select: { id: true, tmdbId: true, name: true, posterUrl: true } }; +const MOVIE_SELECT = { select: { id: true, tmdbId: true, title: true, posterUrl: true } }; +const GAME_SELECT = { select: { id: true, igdbId: true, name: true, coverUrl: true } }; +const BOOK_SELECT = { select: { id: true, hardcoverId: true, title: true, imageUrl: true } }; + +// Every one of the 6 media relations (favorite/listItem are polymorphic — one is non-null). +const POLY_MEDIA = { + anime: ANIME_SELECT, + manga: MANGA_SELECT, + tvShow: TVSHOW_SELECT, + movie: MOVIE_SELECT, + game: GAME_SELECT, + book: BOOK_SELECT, +}; + +const USER_SELECT = { + select: { + id: true, + name: true, + username: true, + profile: { select: { avatarUrl: true } }, + }, +}; + +const INCLUDE = { + _count: { + select: { + reactions: true, + }, + }, + reactions: { + orderBy: { createdAt: "desc" as const }, + select: { + id: true, + emoji: true, + createdAt: true, + user: { + select: { + id: true, + username: true, + }, + }, + }, + }, + // Actor. + user: USER_SELECT, + // Reviews — each with its own criteria fields + media relation. + animeReview: { + select: { + id: true, + overall: true, + story: true, + characters: true, + animation: true, + sound: true, + enjoyment: true, + summary: true, + anime: ANIME_SELECT, + }, + }, + mangaReview: { + select: { + id: true, + overall: true, + art: true, + worldbuilding: true, + story: true, + characters: true, + summary: true, + manga: MANGA_SELECT, + }, + }, + tvShowReview: { + select: { + id: true, + overall: true, + direction: true, + production: true, + acting: true, + story: true, + summary: true, + tvShow: TVSHOW_SELECT, + }, + }, + movieReview: { + select: { + id: true, + overall: true, + direction: true, + production: true, + acting: true, + story: true, + summary: true, + movie: MOVIE_SELECT, + }, + }, + gameReview: { + select: { + id: true, + overall: true, + graphics: true, + sound: true, + story: true, + gameplay: true, + summary: true, + game: GAME_SELECT, + }, + }, + bookReview: { + select: { + id: true, + overall: true, + characters: true, + language: true, + theme: true, + summary: true, + book: BOOK_SELECT, + }, + }, + // Progress — status + media. + animeProgress: { select: { id: true, status: true, anime: ANIME_SELECT } }, + mangaProgress: { select: { id: true, status: true, manga: MANGA_SELECT } }, + tvShowProgress: { select: { id: true, status: true, tvShow: TVSHOW_SELECT } }, + movieProgress: { select: { id: true, status: true, movie: MOVIE_SELECT } }, + gameProgress: { select: { id: true, status: true, game: GAME_SELECT } }, + bookProgress: { select: { id: true, status: true, book: BOOK_SELECT } }, + // Episode watches. + animeEpisodeWatch: { select: { id: true, status: true, episode: true, anime: ANIME_SELECT } }, + tvShowEpisodeWatch: { + select: { id: true, status: true, season: true, episode: true, tvShow: TVSHOW_SELECT }, + }, + // Lists / favorites (polymorphic media). + list: { select: { id: true, name: true } }, + listItem: { select: { id: true, ...POLY_MEDIA } }, + favorite: { select: { id: true, type: true, ...POLY_MEDIA } }, + // Social. + following: { select: { id: true, following: USER_SELECT } }, + userMedal: { select: { id: true, medal: { select: { id: true, name: true, imageUrl: true } } } }, +}; + +@Injectable() +export class ActivityService { + constructor(private readonly databaseService: DatabaseService) {} + + async createActivity(createActivityDto: CreateActivityDto) { + const { type, userId, metadata } = createActivityDto; + + const source = SOURCE_FIELDS.reduce>((acc, field) => { + const value = createActivityDto[field]; + + if (value) { + acc[field] = value; + } + + return acc; + }, {}); + + // 1 activity per source: replace any existing activity referencing the same + // source row (progress status transitions, watch upsert re-fires, etc.). + if (Object.keys(source).length > 0) { + await this.databaseService.activity.deleteMany({ where: source }); + } + + await this.databaseService.activity.create({ + data: { + type, + userId, + ...source, + ...(metadata && { metadata: { ...metadata } }), + }, + }); + } + + async getActivitiesByUserId(getActivitiesByUserIdDto: GetActivitiesByUserDto, friendIds: string[] = []) { + const userExists = await this.databaseService.user.findUnique({ + where: { id: getActivitiesByUserIdDto.userId }, + }); + + if (!userExists) { + throw new AppException(ERROR_CODES.USER_NOT_FOUND); + } + + const pagination = await this.databaseService.offsetPagination({ + model: "activity", + page: getActivitiesByUserIdDto.page, + itemsPerPage: getActivitiesByUserIdDto.itemsPerPage, + where: { + userId: { + in: [...friendIds, getActivitiesByUserIdDto.userId], + }, + }, + include: INCLUDE, + orderBy: { createdAt: "desc" }, + }); + + return { ...pagination, items: this.groupActivities(pagination.items) }; + } + + async getActivitiesByUserFollowing(getActivitiesByUserIdDto: GetActivitiesByUserFollowingDto) { + const following = await this.databaseService.following.findMany({ + where: { followerId: getActivitiesByUserIdDto.userId }, + select: { followingId: true }, + }); + + const friendIds = following.map((item) => item.followingId); + + return this.getActivitiesByUserId({ ...getActivitiesByUserIdDto, userId: getActivitiesByUserIdDto.userId }, friendIds); + } + + async getActivities(getActivitiesDto: GetActivitiesDto) { + const pagination = await this.databaseService.offsetPagination({ + model: "activity", + orderBy: { createdAt: "desc" }, + page: getActivitiesDto.page, + itemsPerPage: getActivitiesDto.itemsPerPage, + include: INCLUDE, + }); + + return { ...pagination, items: this.groupActivities(pagination.items) }; + } + + // Collapse consecutive activities of the same (userId, type) inside a 1h window + // into a single group with a count, so the feed can render "favorited 3 animes". + private groupActivities(items: any[]) { + const WINDOW_MS = 60 * 60 * 1000; + const groups: any[] = []; + + for (const item of items) { + const last = groups[groups.length - 1]; + const sameGroup = + last && + last.type === item.type && + last.userId === item.userId && + new Date(last.createdAt).getTime() - new Date(item.createdAt).getTime() <= WINDOW_MS; + + if (sameGroup) { + last.items.push(item); + last.count += 1; + } else { + groups.push({ + type: item.type, + userId: item.userId, + createdAt: item.createdAt, + count: 1, + items: [item], + }); + } + } + + return groups; + } + + async getUserActivityCalendarById(userId: string) { + const userExists = await this.databaseService.user.findUnique({ + where: { id: userId }, + select: { id: true }, + }); + + if (!userExists) { + throw new AppException(ERROR_CODES.USER_NOT_FOUND); + } + + const startDate = new Date(); + startDate.setHours(0, 0, 0, 0); + startDate.setDate(startDate.getDate() - 364); + + const activities = await this.databaseService.activity.findMany({ + where: { + userId, + createdAt: { + gte: startDate, + }, + }, + select: { + createdAt: true, + }, + orderBy: { + createdAt: "asc", + }, + }); + + const groupedByDate = new Map(); + + for (const { createdAt } of activities) { + const date = createdAt.toISOString().split("T")[0]; + const currentCount = groupedByDate.get(date) ?? 0; + + groupedByDate.set(date, currentCount + 1); + } + + return { + total: activities.length, + items: [...groupedByDate.entries()].map(([date, count]) => ({ + date, + count, + })), + }; + } +} diff --git a/src/modules/feed-event/controller/feed-event.controller.ts b/src/modules/feed-event/controller/feed-event.controller.ts deleted file mode 100644 index b61db7a..0000000 --- a/src/modules/feed-event/controller/feed-event.controller.ts +++ /dev/null @@ -1,30 +0,0 @@ -import { Controller, Get, Query, UseGuards } from "@nestjs/common"; -import { ApiTags } from "@nestjs/swagger"; -import { AuthGuard, Session, type UserSession } from "@thallesp/nestjs-better-auth"; -import { GetFeedEventsDto } from "../dto/get-feed-events.dto"; -import { GetFeedEventsByUserDto } from "../dto/get-feed-events-by-user.dto"; -import { FeedEventService } from "../service/feed-event.service"; - -@ApiTags("Feed Event") -@Controller("/feed") -export class FeedEventController { - constructor(private readonly feedEventService: FeedEventService) {} - - @Get("/global") - async getFeedEvents(@Query() query: GetFeedEventsDto) { - const feedEvents = await this.feedEventService.getFeedEvents(query); - - return { feedEvents }; - } - - @Get("/user") - @UseGuards(AuthGuard) - async getFeedEventsByUserId(@Session() session: UserSession, @Query() query: GetFeedEventsByUserDto) { - const feedEvents = await this.feedEventService.getFeedEventsByUserId({ - ...query, - userId: session.user.id, - }); - - return { feedEvents }; - } -} diff --git a/src/modules/feed-event/dto/feed-event.dto.ts b/src/modules/feed-event/dto/feed-event.dto.ts deleted file mode 100644 index c83c433..0000000 --- a/src/modules/feed-event/dto/feed-event.dto.ts +++ /dev/null @@ -1,40 +0,0 @@ -import { ApiProperty, ApiPropertyOptional } from "@nestjs/swagger"; -import { FeedEventType } from "@prisma/generated/enums"; -import { IsArray, IsEnum, IsInt, IsNotEmpty, IsOptional, IsUUID } from "class-validator"; - -export interface FeedEventMetadata { - readonly id: string; - readonly [key: string]: any; -} - -export class FeedEventDto { - @IsEnum(FeedEventType) - @IsNotEmpty() - @ApiProperty({ enum: FeedEventType }) - readonly type: FeedEventType; - - @IsNotEmpty() - @IsUUID() - @ApiProperty({ type: "string", format: "uuid" }) - readonly userId: string; - - @IsArray() - @IsOptional() - @IsUUID("7", { each: true }) - @ApiPropertyOptional({ type: [String], format: "uuid" }) - readonly entityIds?: string[] = []; - - @IsInt() - @IsOptional() - @ApiPropertyOptional({ type: "integer", default: 1 }) - readonly count?: number = 1; - - @IsOptional() - @ApiPropertyOptional({ - oneOf: [ - { type: "object", additionalProperties: true }, - { type: "array", items: { type: "object", additionalProperties: true } }, - ], - }) - readonly metadata: FeedEventMetadata | FeedEventMetadata[]; -} diff --git a/src/modules/feed-event/dto/get-feed-events-by-user.dto.ts b/src/modules/feed-event/dto/get-feed-events-by-user.dto.ts deleted file mode 100644 index 12ff143..0000000 --- a/src/modules/feed-event/dto/get-feed-events-by-user.dto.ts +++ /dev/null @@ -1,7 +0,0 @@ -import { ApiProperty } from "@nestjs/swagger"; -import { OffsetPaginationParamsDto } from "@/shared/infra/database/dtos/offset-pagination.dto"; - -export class GetFeedEventsByUserDto extends OffsetPaginationParamsDto { - @ApiProperty({ type: "string", format: "uuid" }) - readonly userId: string; -} diff --git a/src/modules/feed-event/feed-event.module.ts b/src/modules/feed-event/feed-event.module.ts deleted file mode 100644 index 876b50b..0000000 --- a/src/modules/feed-event/feed-event.module.ts +++ /dev/null @@ -1,11 +0,0 @@ -import { Module } from "@nestjs/common"; -import { FeedEventController } from "./controller/feed-event.controller"; -import { FeedEventService } from "./service/feed-event.service"; - -@Module({ - imports: [], - controllers: [FeedEventController], - providers: [FeedEventService], - exports: [FeedEventService], -}) -export class FeedEventModule {} diff --git a/src/modules/feed-event/service/feed-event.service.ts b/src/modules/feed-event/service/feed-event.service.ts deleted file mode 100644 index 902d100..0000000 --- a/src/modules/feed-event/service/feed-event.service.ts +++ /dev/null @@ -1,109 +0,0 @@ -import { Injectable } from "@nestjs/common"; -import { FeedEventFindManyArgs } from "@prisma/generated/models"; -import { ERROR_CODES } from "@/shared/constants/error-codes"; -import { AppException } from "@/shared/exceptions/app.exceptions"; -import { DatabaseService } from "@/shared/infra/database/database.service"; -import { FeedEventDto } from "../dto/feed-event.dto"; -import { GetFeedEventsDto } from "../dto/get-feed-events.dto"; -import { GetFeedEventsByUserDto } from "../dto/get-feed-events-by-user.dto"; - -@Injectable() -export class FeedEventService { - constructor(private readonly databaseService: DatabaseService) {} - - async createFeedEvent(feedEventDto: FeedEventDto) { - const { type, userId, metadata } = feedEventDto; - - await this.databaseService.feedEvent.create({ - data: { - type, - userId, - metadata: { ...metadata }, - }, - }); - } - - async getFeedEventsByUserId(getFeedEventsByUserIdDto: GetFeedEventsByUserDto) { - const userExists = await this.databaseService.user.findUnique({ - where: { id: getFeedEventsByUserIdDto.userId }, - }); - - if (!userExists) { - throw new AppException(ERROR_CODES.USER_NOT_FOUND); - } - - const following = await this.databaseService.following.findMany({ - where: { followerId: getFeedEventsByUserIdDto.userId }, - select: { followingId: true }, - }); - - const friendIds = following.map((item) => item.followingId); - - const pagination = await this.databaseService.offsetPagination({ - model: "feedEvent", - page: getFeedEventsByUserIdDto.page, - itemsPerPage: getFeedEventsByUserIdDto.itemsPerPage, - where: { - userId: { - in: [...friendIds, getFeedEventsByUserIdDto.userId], - }, - }, - include: { - _count: { - select: { - reactions: true, - }, - }, - reactions: { - take: 3, - orderBy: { createdAt: "desc" }, - select: { - id: true, - emoji: true, - createdAt: true, - user: { - select: { - username: true, - }, - }, - }, - }, - }, - orderBy: { createdAt: "desc" }, - }); - - return pagination; - } - - async getFeedEvents(getFeedEventsDto: GetFeedEventsDto) { - const pagination = await this.databaseService.offsetPagination({ - model: "feedEvent", - orderBy: { createdAt: "desc" }, - page: getFeedEventsDto.page, - itemsPerPage: getFeedEventsDto.itemsPerPage, - include: { - _count: { - select: { - reactions: true, - }, - }, - reactions: { - take: 3, - orderBy: { createdAt: "desc" }, - select: { - id: true, - emoji: true, - createdAt: true, - user: { - select: { - username: true, - }, - }, - }, - }, - }, - }); - - return pagination; - } -} diff --git a/src/shared/infra/queue/processors/activity.processor.ts b/src/shared/infra/queue/processors/activity.processor.ts new file mode 100644 index 0000000..a58c5d1 --- /dev/null +++ b/src/shared/infra/queue/processors/activity.processor.ts @@ -0,0 +1,58 @@ +import { OnWorkerEvent, Processor, WorkerHost } from "@nestjs/bullmq"; +import { Logger } from "@nestjs/common"; +import { Job } from "bullmq"; +import { CreateActivityDto } from "@/modules/activity/dto/activity.dto"; +import { ActivityService } from "@/modules/activity/service/activity.service"; +import { ACTIVITY_JOB } from "@/shared/constants/job"; +import { ACTIVITY_QUEUE } from "@/shared/constants/queue"; + +export type ActivityJobData = CreateActivityDto; + +@Processor(ACTIVITY_QUEUE, { concurrency: 10 }) +export class ActivityProcessor extends WorkerHost { + private readonly logger = new Logger(ActivityProcessor.name); + + constructor(private readonly activityService: ActivityService) { + super(); + } + + async process(job: Job) { + if (job.name === ACTIVITY_JOB) { + await this.activityService.createActivity(job.data as ActivityJobData); + + return; + } + + throw new Error(`Unsupported activity job name: ${job.name}`); + } + + @OnWorkerEvent("active") + onActive(job: Job) { + this.logger.log( + `Processing job [${ACTIVITY_QUEUE}] | job=${job.id} name=${job.name} attempt=${job.attemptsMade + 1}`, + ); + } + + @OnWorkerEvent("completed") + onCompleted(job: Job) { + this.logger.log(`Job completed [${ACTIVITY_QUEUE}] | job=${job.id} name=${job.name}`); + } + + @OnWorkerEvent("failed") + onFailed(job: Job | undefined, error: Error) { + if (!job) return; + + const maxAttempts = job.opts?.attempts ?? 1; + const willRetry = job.attemptsMade < maxAttempts; + + if (willRetry) { + this.logger.warn( + `Job failed, retrying [${ACTIVITY_QUEUE}] | job=${job.id} name=${job.name} attempt=${job.attemptsMade}/${maxAttempts} error=${error.message}`, + ); + } else { + this.logger.error( + `Job removed from queue after max attempts [${ACTIVITY_QUEUE}] | job=${job.id} name=${job.name} attempts=${job.attemptsMade}/${maxAttempts} error=${error.message}`, + ); + } + } +} diff --git a/src/shared/infra/queue/processors/feed-event.processor.ts b/src/shared/infra/queue/processors/feed-event.processor.ts deleted file mode 100644 index 3bb9ce6..0000000 --- a/src/shared/infra/queue/processors/feed-event.processor.ts +++ /dev/null @@ -1,115 +0,0 @@ -import { OnWorkerEvent, Processor, WorkerHost } from "@nestjs/bullmq"; -import { Logger } from "@nestjs/common"; -import { Job } from "bullmq"; -import { FeedEventDto, FeedEventMetadata } from "@/modules/feed-event/dto/feed-event.dto"; -import { FeedEventService } from "@/modules/feed-event/service/feed-event.service"; -import { FEED_EVENT_FLUSH_AGGREGATION_JOB, FEED_EVENT_JOB } from "@/shared/constants/job"; -import { FEED_EVENT_QUEUE } from "@/shared/constants/queue"; -import { CacheService } from "../../cache/cache.service"; -import { QueueService } from "../queue.service"; - -export type FeedEventJobData = FeedEventDto; - -export interface FeedEventFlushAggregationJobData { - aggKey: string; - windowsMs: number; -} - -@Processor(FEED_EVENT_QUEUE, { concurrency: 10 }) -export class FeedEventProcessor extends WorkerHost { - private readonly logger = new Logger(FeedEventProcessor.name); - - constructor( - private readonly queueService: QueueService, - private readonly feedEventService: FeedEventService, - private readonly cacheService: CacheService, - ) { - super(); - } - - async process(job: Job) { - if (job.name === FEED_EVENT_JOB) { - const { type, userId, metadata } = job.data as FeedEventJobData; - - const aggKey = `feed:agg:${userId}:${type}`; - const lockKey = `${aggKey}:lock`; - - await this.cacheService.redis.lPush(aggKey, JSON.stringify({ userId, type, metadata })); - await this.cacheService.redis.expire(aggKey, 600); - - const windowsMs = 5 * 60 * 1000; - - const isLeader = await this.cacheService.redis.set(lockKey, "1", { NX: true, PX: windowsMs }); - - if (isLeader) { - await this.queueService.toFeedEventFlushAggregationJob({ aggKey, windowsMs }); - - this.logger.log(`Aggregation window opened: ${aggKey}`); - } - - return; - } - - if (job.name === FEED_EVENT_FLUSH_AGGREGATION_JOB) { - const { aggKey } = job.data as FeedEventFlushAggregationJobData; - - const raw = await this.cacheService.redis.lRange(aggKey, 0, -1); - - await this.cacheService.redis.del(aggKey); - - if (!raw.length) return; - - const events = raw.map((r) => JSON.parse(r)); - - const first = events[0]; - - const entityIds = events.map((e) => (e.metadata as FeedEventMetadata)?.id).filter(Boolean); - - await this.feedEventService.createFeedEvent({ - type: first.type, - userId: first.userId, - count: events.length, - entityIds, - metadata: { - ...(events.length === 1 - ? (first.metadata as FeedEventMetadata) - : (events.map((e) => e.metadata) as FeedEventMetadata[])), - }, - }); - - return; - } - - throw new Error(`Unsupported feed event job name: ${job.name}`); - } - - @OnWorkerEvent("active") - onActive(job: Job) { - this.logger.log( - `Processing job [${FEED_EVENT_QUEUE}] | job=${job.id} name=${job.name} attempt=${job.attemptsMade + 1}`, - ); - } - - @OnWorkerEvent("completed") - onCompleted(job: Job) { - this.logger.log(`Job completed [${FEED_EVENT_QUEUE}] | job=${job.id} name=${job.name}`); - } - - @OnWorkerEvent("failed") - onFailed(job: Job | undefined, error: Error) { - if (!job) return; - - const maxAttempts = job.opts?.attempts ?? 1; - const willRetry = job.attemptsMade < maxAttempts; - - if (willRetry) { - this.logger.warn( - `Job failed, retrying [${FEED_EVENT_QUEUE}] | job=${job.id} name=${job.name} attempt=${job.attemptsMade}/${maxAttempts} error=${error.message}`, - ); - } else { - this.logger.error( - `Job removed from queue after max attempts [${FEED_EVENT_QUEUE}] | job=${job.id} name=${job.name} attempts=${job.attemptsMade}/${maxAttempts} error=${error.message}`, - ); - } - } -} diff --git a/test/unit/infra/queue/feed-event.processor.spec.ts b/test/unit/controllers/activity.controller.spec.ts similarity index 74% rename from test/unit/infra/queue/feed-event.processor.spec.ts rename to test/unit/controllers/activity.controller.spec.ts index f3291b1..840040f 100644 --- a/test/unit/infra/queue/feed-event.processor.spec.ts +++ b/test/unit/controllers/activity.controller.spec.ts @@ -1,6 +1,6 @@ import { describe, expect, it } from "vitest"; -describe("FeedEventProcessor", () => { +describe("ActivityController", () => { it("true is true", () => { expect(true).toBe(true); }); diff --git a/test/unit/controllers/feed-event.controller.spec.ts b/test/unit/infra/queue/activity.processor.spec.ts similarity index 74% rename from test/unit/controllers/feed-event.controller.spec.ts rename to test/unit/infra/queue/activity.processor.spec.ts index 6820717..eb2fdb0 100644 --- a/test/unit/controllers/feed-event.controller.spec.ts +++ b/test/unit/infra/queue/activity.processor.spec.ts @@ -1,6 +1,6 @@ import { describe, expect, it } from "vitest"; -describe("FeedEventController", () => { +describe("ActivityProcessor", () => { it("true is true", () => { expect(true).toBe(true); }); diff --git a/test/unit/services/feed-event.service.spec.ts b/test/unit/services/activity.service.spec.ts similarity index 75% rename from test/unit/services/feed-event.service.spec.ts rename to test/unit/services/activity.service.spec.ts index c2b0739..e77407e 100644 --- a/test/unit/services/feed-event.service.spec.ts +++ b/test/unit/services/activity.service.spec.ts @@ -1,6 +1,6 @@ import { describe, expect, it } from "vitest"; -describe("FeedEventService", () => { +describe("ActivityService", () => { it("true is true", () => { expect(true).toBe(true); }); From 10b258c6f7849998b965c17b89c463d5be27f091 Mon Sep 17 00:00:00 2001 From: izakdvlpr Date: Thu, 9 Jul 2026 18:10:49 -0300 Subject: [PATCH 4/8] feat: create notification module --- .../controller/notification.controller.ts | 87 +++++++ .../notification/dto/get-notifications.dto.ts | 16 ++ .../notification/dto/notification.dto.ts | 43 ++++ .../notification/notification.module.ts | 11 + .../service/notification.service.ts | 236 ++++++++++++++++++ 5 files changed, 393 insertions(+) create mode 100644 src/modules/notification/controller/notification.controller.ts create mode 100644 src/modules/notification/dto/get-notifications.dto.ts create mode 100644 src/modules/notification/dto/notification.dto.ts create mode 100644 src/modules/notification/notification.module.ts create mode 100644 src/modules/notification/service/notification.service.ts diff --git a/src/modules/notification/controller/notification.controller.ts b/src/modules/notification/controller/notification.controller.ts new file mode 100644 index 0000000..ad209de --- /dev/null +++ b/src/modules/notification/controller/notification.controller.ts @@ -0,0 +1,87 @@ +import { + Controller, + Delete, + Get, + HttpCode, + HttpStatus, + Param, + ParseUUIDPipe, + Post, + Query, + UseGuards, +} from "@nestjs/common"; +import { ApiTags } from "@nestjs/swagger"; +import { AuthGuard, Session, type UserSession } from "@thallesp/nestjs-better-auth"; +import { GetNotificationsDto } from "../dto/get-notifications.dto"; +import { NotificationService } from "../service/notification.service"; + +// Literal paths are declared before `/:notificationId`, otherwise Nest matches "read"/"unread" as +// the param and ParseUUIDPipe rejects the request. +@ApiTags("Notification") +@Controller("/notifications") +@UseGuards(AuthGuard) +export class NotificationController { + constructor(private readonly notificationService: NotificationService) {} + + @Get("/") + async getNotifications(@Session() session: UserSession, @Query() query: GetNotificationsDto) { + const notifications = await this.notificationService.getNotifications({ + ...query, + userId: session.user.id, + }); + + return { notifications }; + } + + @Get("/unread/count") + async getUnreadCount(@Session() session: UserSession) { + const count = await this.notificationService.getUnreadCount(session.user.id); + + return { count }; + } + + @Post("/read/all") + @HttpCode(HttpStatus.NO_CONTENT) + async markAllAsRead(@Session() session: UserSession) { + await this.notificationService.markAllAsRead(session.user.id); + } + + @Delete("/read/all") + @HttpCode(HttpStatus.NO_CONTENT) + async markAllAsUnread(@Session() session: UserSession) { + await this.notificationService.markAllAsUnread(session.user.id); + } + + @Delete("/all") + @HttpCode(HttpStatus.NO_CONTENT) + async deleteAllNotifications(@Session() session: UserSession) { + await this.notificationService.deleteAllNotifications(session.user.id); + } + + @Post("/:notificationId/read") + @HttpCode(HttpStatus.NO_CONTENT) + async markAsRead( + @Session() session: UserSession, + @Param("notificationId", new ParseUUIDPipe()) notificationId: string, + ) { + await this.notificationService.markAsRead(session.user.id, notificationId); + } + + @Delete("/:notificationId/read") + @HttpCode(HttpStatus.NO_CONTENT) + async markAsUnread( + @Session() session: UserSession, + @Param("notificationId", new ParseUUIDPipe()) notificationId: string, + ) { + await this.notificationService.markAsUnread(session.user.id, notificationId); + } + + @Delete("/:notificationId") + @HttpCode(HttpStatus.NO_CONTENT) + async deleteNotification( + @Session() session: UserSession, + @Param("notificationId", new ParseUUIDPipe()) notificationId: string, + ) { + await this.notificationService.deleteNotification(session.user.id, notificationId); + } +} diff --git a/src/modules/notification/dto/get-notifications.dto.ts b/src/modules/notification/dto/get-notifications.dto.ts new file mode 100644 index 0000000..33a3100 --- /dev/null +++ b/src/modules/notification/dto/get-notifications.dto.ts @@ -0,0 +1,16 @@ +import { ApiPropertyOptional } from "@nestjs/swagger"; +import { Transform } from "class-transformer"; +import { IsBoolean, IsOptional } from "class-validator"; +import { OffsetPaginationParamsDto } from "@/shared/infra/database/dtos/offset-pagination.dto"; + +export class GetNotificationsDto extends OffsetPaginationParamsDto { + @IsOptional() + @Transform(({ value }) => (value === "true" ? true : value === "false" ? false : value)) + @IsBoolean() + @ApiPropertyOptional({ type: "boolean", description: "Filter by read state. Omit to return every notification." }) + readonly read?: boolean; +} + +export class GetNotificationsByUserDto extends GetNotificationsDto { + readonly userId: string; +} diff --git a/src/modules/notification/dto/notification.dto.ts b/src/modules/notification/dto/notification.dto.ts new file mode 100644 index 0000000..0aa3c56 --- /dev/null +++ b/src/modules/notification/dto/notification.dto.ts @@ -0,0 +1,43 @@ +import { ApiProperty } from "@nestjs/swagger"; +import { ArrayNotEmpty, IsArray, IsNotEmpty, IsObject, IsUUID } from "class-validator"; + +// Shape of Notification.metadata when type is System. Kept open on purpose — what the system sends +// is not settled yet, so anything beyond these fields rides along untouched. +// +// Prefer titleKey/descriptionKey: the web app resolves them through i18n at render time, so the +// notification follows the reader's language. The literal title/description are the escape hatch for +// content that has no translation (a one-off announcement, a user-supplied string). +export interface SystemNotificationMetadata { + readonly title?: string; + readonly titleKey?: string; + readonly description?: string; + readonly descriptionKey?: string; + readonly url?: string; + readonly [key: string]: any; +} + +export class CreateSystemNotificationDto { + @IsArray() + @ArrayNotEmpty() + @ApiProperty({ type: "array", items: { type: "string" }, description: "Users that receive the notification." }) + readonly recipientIds: string[]; + + @IsObject() + @IsNotEmpty() + @ApiProperty({ type: "object", additionalProperties: true }) + readonly metadata: SystemNotificationMetadata; +} + +export class CreateCommentNotificationDto { + @IsNotEmpty() + @IsUUID("7") + @ApiProperty({ type: "string", format: "uuid" }) + readonly commentId: string; +} + +export class CreateReactionNotificationDto { + @IsNotEmpty() + @IsUUID("7") + @ApiProperty({ type: "string", format: "uuid" }) + readonly reactionId: string; +} diff --git a/src/modules/notification/notification.module.ts b/src/modules/notification/notification.module.ts new file mode 100644 index 0000000..2b804d9 --- /dev/null +++ b/src/modules/notification/notification.module.ts @@ -0,0 +1,11 @@ +import { Module } from "@nestjs/common"; +import { NotificationController } from "./controller/notification.controller"; +import { NotificationService } from "./service/notification.service"; + +@Module({ + imports: [], + controllers: [NotificationController], + providers: [NotificationService], + exports: [NotificationService], +}) +export class NotificationModule {} diff --git a/src/modules/notification/service/notification.service.ts b/src/modules/notification/service/notification.service.ts new file mode 100644 index 0000000..145d917 --- /dev/null +++ b/src/modules/notification/service/notification.service.ts @@ -0,0 +1,236 @@ +import { Injectable } from "@nestjs/common"; +import { CommentType, NotificationType, ReactionType } from "@prisma/generated/enums"; +import { NotificationFindManyArgs } from "@prisma/generated/models"; +import { ERROR_CODES } from "@/shared/constants/error-codes"; +import { AppException } from "@/shared/exceptions/app.exceptions"; +import { DatabaseService } from "@/shared/infra/database/database.service"; +import { GetNotificationsByUserDto } from "../dto/get-notifications.dto"; +import { + CreateCommentNotificationDto, + CreateReactionNotificationDto, + CreateSystemNotificationDto, +} from "../dto/notification.dto"; + +const REACTION_TARGETS = { + [ReactionType.Comment]: { type: NotificationType.ReactionOnComment, relation: "comment", sourceField: "commentId" }, + [ReactionType.Activity]: { + type: NotificationType.ReactionOnActivity, + relation: "activity", + sourceField: "activityId", + }, + [ReactionType.AnimeReview]: { + type: NotificationType.ReactionOnAnimeReview, + relation: "animeReview", + sourceField: "animeReviewId", + }, + [ReactionType.MangaReview]: { + type: NotificationType.ReactionOnMangaReview, + relation: "mangaReview", + sourceField: "mangaReviewId", + }, + [ReactionType.TvShowReview]: { + type: NotificationType.ReactionOnTvShowReview, + relation: "tvShowReview", + sourceField: "tvShowReviewId", + }, + [ReactionType.MovieReview]: { + type: NotificationType.ReactionOnMovieReview, + relation: "movieReview", + sourceField: "movieReviewId", + }, + [ReactionType.GameReview]: { + type: NotificationType.ReactionOnGameReview, + relation: "gameReview", + sourceField: "gameReviewId", + }, + [ReactionType.BookReview]: { + type: NotificationType.ReactionOnBookReview, + relation: "bookReview", + sourceField: "bookReviewId", + }, +} as const; + +const OWNER_SELECT = { select: { userId: true } }; + +const ACTOR_SELECT = { + select: { + id: true, + name: true, + username: true, + profile: { + select: { + id: true, + avatarUrl: true, + }, + }, + }, +}; + +@Injectable() +export class NotificationService { + constructor(private readonly databaseService: DatabaseService) {} + + async createSystemNotification({ recipientIds, metadata }: CreateSystemNotificationDto) { + await this.databaseService.notification.createMany({ + data: recipientIds.map((recipientId) => ({ + type: NotificationType.System, + recipientId, + metadata: { ...metadata }, + })), + }); + } + + async createFromComment({ commentId }: CreateCommentNotificationDto) { + const comment = await this.databaseService.comment.findUnique({ + where: { id: commentId }, + include: { profile: OWNER_SELECT }, + }); + + if (!comment || comment.type !== CommentType.Profile || !comment.profile) { + return; + } + + const recipientId = comment.profile.userId; + + if (recipientId === comment.userId) { + return; + } + + await this.databaseService.notification.create({ + data: { + type: NotificationType.CommentOnProfile, + recipientId, + actorId: comment.userId, + commentId: comment.id, + profileId: comment.profileId, + }, + }); + } + + async createFromReaction({ reactionId }: CreateReactionNotificationDto) { + const reaction = await this.databaseService.reaction.findUnique({ + where: { id: reactionId }, + include: { + comment: OWNER_SELECT, + activity: OWNER_SELECT, + animeReview: OWNER_SELECT, + mangaReview: OWNER_SELECT, + tvShowReview: OWNER_SELECT, + movieReview: OWNER_SELECT, + gameReview: OWNER_SELECT, + bookReview: OWNER_SELECT, + }, + }); + + if (!reaction) { + return; + } + + const target = REACTION_TARGETS[reaction.type]; + const owner = reaction[target.relation]; + const sourceId = reaction[target.sourceField]; + + if (!owner || !sourceId || owner.userId === reaction.userId) { + return; + } + + await this.databaseService.notification.create({ + data: { + type: target.type, + recipientId: owner.userId, + actorId: reaction.userId, + reactionId: reaction.id, + [target.sourceField]: sourceId, + }, + }); + } + + async getNotifications({ userId, read, page, itemsPerPage }: GetNotificationsByUserDto) { + const pagination = await this.databaseService.offsetPagination({ + model: "notification", + where: { + recipientId: userId, + ...(read !== undefined && { readAt: read ? { not: null } : null }), + }, + page, + itemsPerPage, + orderBy: { createdAt: "desc" }, + include: { + actor: ACTOR_SELECT, + comment: { + select: { + id: true, + type: true, + content: true, + }, + }, + reaction: { + select: { + id: true, + type: true, + emoji: true, + }, + }, + }, + }); + + return pagination; + } + + async getUnreadCount(userId: string) { + const count = await this.databaseService.notification.count({ + where: { recipientId: userId, readAt: null }, + }); + + return count; + } + + async markAsRead(userId: string, notificationId: string) { + await this.updateReadAt(userId, notificationId, new Date()); + } + + async markAsUnread(userId: string, notificationId: string) { + await this.updateReadAt(userId, notificationId, null); + } + + async markAllAsRead(userId: string) { + await this.databaseService.notification.updateMany({ + where: { recipientId: userId, readAt: null }, + data: { readAt: new Date() }, + }); + } + + async markAllAsUnread(userId: string) { + await this.databaseService.notification.updateMany({ + where: { recipientId: userId, readAt: { not: null } }, + data: { readAt: null }, + }); + } + + async deleteNotification(userId: string, notificationId: string) { + const { count } = await this.databaseService.notification.deleteMany({ + where: { id: notificationId, recipientId: userId }, + }); + + if (count === 0) { + throw new AppException(ERROR_CODES.NOTIFICATION_NOT_FOUND); + } + } + + async deleteAllNotifications(userId: string) { + await this.databaseService.notification.deleteMany({ + where: { recipientId: userId }, + }); + } + + private async updateReadAt(userId: string, notificationId: string, readAt: Date | null) { + const { count } = await this.databaseService.notification.updateMany({ + where: { id: notificationId, recipientId: userId }, + data: { readAt }, + }); + + if (count === 0) { + throw new AppException(ERROR_CODES.NOTIFICATION_NOT_FOUND); + } + } +} From f923b7f18df0c7edab27a42e9afc20436a15fbea Mon Sep 17 00:00:00 2001 From: izakdvlpr Date: Thu, 9 Jul 2026 18:11:03 -0300 Subject: [PATCH 5/8] feat: implement activity and notification module --- .../anime/dto/create-anime-review.dto.ts | 2 +- .../service/anime-episode-watch.service.ts | 40 +++++++++++++++++-- .../anime/service/anime-progress.service.ts | 18 ++++++--- .../anime/service/anime-review.service.ts | 7 ++-- src/modules/auth/config/auth.config.ts | 18 +++++++++ .../book/dto/create-book-review.dto.ts | 2 +- .../book/service/book-progress.service.ts | 17 +++++--- .../book/service/book-review.service.ts | 7 ++-- .../comment/service/comment.service.ts | 10 ++++- .../favorite/service/favorite.service.ts | 7 ++-- .../game/dto/create-game-review.dto.ts | 2 +- .../game/service/game-progress.service.ts | 17 +++++--- .../game/service/game-review.service.ts | 7 ++-- src/modules/list/service/list.service.ts | 12 +++--- .../manga/dto/create-manga-review.dto.ts | 2 +- .../manga/service/manga-progress.service.ts | 17 +++++--- .../manga/service/manga-review.service.ts | 7 ++-- .../movie/dto/create-movie-review.dto.ts | 2 +- .../movie/service/movie-progress.service.ts | 17 +++++--- .../movie/service/movie-review.service.ts | 7 ++-- src/modules/payment/service/stripe.service.ts | 11 ++++- .../reaction/dto/create-reaction.dto.ts | 6 +-- src/modules/reaction/dto/get-reactions.dto.ts | 6 +-- .../reaction/service/reaction.service.ts | 14 +++++-- .../tv-show/dto/create-tv-show-review.dto.ts | 2 +- .../service/tv-show-episode-watch.service.ts | 40 +++++++++++++++++-- .../service/tv-show-progress.service.ts | 18 ++++++--- .../tv-show/service/tv-show-review.service.ts | 7 ++-- src/modules/user/service/user.service.ts | 7 ++-- 29 files changed, 236 insertions(+), 93 deletions(-) diff --git a/src/modules/anime/dto/create-anime-review.dto.ts b/src/modules/anime/dto/create-anime-review.dto.ts index c5882a8..a5ceb79 100644 --- a/src/modules/anime/dto/create-anime-review.dto.ts +++ b/src/modules/anime/dto/create-anime-review.dto.ts @@ -44,7 +44,7 @@ export class CreateAnimeReviewDto { readonly enjoyment?: number; @IsOptional() - @MaxLength(250) + @MaxLength(500) @ApiPropertyOptional({ type: "string", maxLength: 250 }) readonly summary?: string; diff --git a/src/modules/anime/service/anime-episode-watch.service.ts b/src/modules/anime/service/anime-episode-watch.service.ts index 41f7f9a..2d754ce 100644 --- a/src/modules/anime/service/anime-episode-watch.service.ts +++ b/src/modules/anime/service/anime-episode-watch.service.ts @@ -1,17 +1,37 @@ import { Injectable } from "@nestjs/common"; -import { WatchEpisodeStatus } from "@prisma/generated/enums"; +import { AnimeEpisodeWatch } from "@prisma/generated/client"; +import { ActivityType, WatchEpisodeStatus } from "@prisma/generated/enums"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; +import { QueueService } from "@/shared/infra/queue/queue.service"; import { CreateOrUpdateAnimeEpisodeWatchDto } from "../dto/create-or-update-anime-episode-watch.dto"; import { DeleteAllAnimeEpisodeWatchDto } from "../dto/delete-all-anime-episode-watch.dto"; import { DeleteAnimeEpisodeWatchDto } from "../dto/delete-anime-episode-watch.dto"; import { GetAnimeEpisodeWatchDto } from "../dto/get-anime-episode-watch.dto"; import { WatchAllAnimeEpisodesDto } from "../dto/watch-all-anime-episodes.dto"; +const WATCHED_STATUSES: WatchEpisodeStatus[] = [WatchEpisodeStatus.Watching, WatchEpisodeStatus.Completed]; + @Injectable() export class AnimeEpisodeWatchService { - constructor(private readonly databaseService: DatabaseService) {} + constructor( + private readonly databaseService: DatabaseService, + private readonly queueService: QueueService, + ) {} + + private async emitWatchedActivities(userId: string, watches: AnimeEpisodeWatch[]) { + for (const watch of watches) { + if (!WATCHED_STATUSES.includes(watch.status)) continue; + + await this.queueService.toActivityJob({ + type: ActivityType.Watched, + userId, + animeEpisodeWatchId: watch.id, + metadata: { ...watch }, + }); + } + } async createOrUpdateAnimeEpisodeWatch(createOrUpdateAnimeEpisodeWatchDto: CreateOrUpdateAnimeEpisodeWatchDto) { const { animeId, userId, episodes } = createOrUpdateAnimeEpisodeWatchDto; @@ -35,10 +55,12 @@ export class AnimeEpisodeWatchService { const batchSize = 50; + const watches: AnimeEpisodeWatch[] = []; + for (let i = 0; i < episodes.length; i += batchSize) { const batch = episodes.slice(i, i + batchSize); - await Promise.all( + const results = await Promise.all( batch.map(({ episode, status }) => this.databaseService.animeEpisodeWatch.upsert({ where: { @@ -58,7 +80,11 @@ export class AnimeEpisodeWatchService { }), ), ); + + watches.push(...results); } + + await this.emitWatchedActivities(userId, watches); } async watchAllAnimeEpisodes({ animeId, userId }: WatchAllAnimeEpisodesDto) { @@ -78,10 +104,12 @@ export class AnimeEpisodeWatchService { const episodeNumbers = Array.from({ length: anime.numberOfEpisodes }, (_, i) => i + 1); const batchSize = 50; + const watches: AnimeEpisodeWatch[] = []; + for (let i = 0; i < episodeNumbers.length; i += batchSize) { const batch = episodeNumbers.slice(i, i + batchSize); - await Promise.all( + const results = await Promise.all( batch.map((episode) => this.databaseService.animeEpisodeWatch.upsert({ where: { @@ -101,7 +129,11 @@ export class AnimeEpisodeWatchService { }), ), ); + + watches.push(...results); } + + await this.emitWatchedActivities(userId, watches); } async deleteAnimeEpisodeWatch({ userId, animeId, episode }: DeleteAnimeEpisodeWatchDto) { diff --git a/src/modules/anime/service/anime-progress.service.ts b/src/modules/anime/service/anime-progress.service.ts index 2a1b7aa..7b30f18 100644 --- a/src/modules/anime/service/anime-progress.service.ts +++ b/src/modules/anime/service/anime-progress.service.ts @@ -1,6 +1,7 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType, ProgressStatus } from "@prisma/generated/enums"; +import { ProgressStatus } from "@prisma/generated/enums"; import { AnimeProgressFindManyArgs } from "@prisma/generated/models"; +import { activityTypeFromProgressStatus } from "@/modules/activity/activity.utils"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; @@ -66,11 +67,16 @@ export class AnimeProgressService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewProgress, - userId, - metadata: { ...animeProgress }, - }); + const activityType = activityTypeFromProgressStatus(status); + + if (activityType) { + await this.queueService.toActivityJob({ + type: activityType, + userId, + animeProgressId: animeProgress.id, + metadata: { ...animeProgress }, + }); + } if (status === ProgressStatus.Completed) { await this.animeEpisodeWatchService.watchAllAnimeEpisodes({ animeId, userId }); diff --git a/src/modules/anime/service/anime-review.service.ts b/src/modules/anime/service/anime-review.service.ts index 4a6a9e7..7cf5f32 100644 --- a/src/modules/anime/service/anime-review.service.ts +++ b/src/modules/anime/service/anime-review.service.ts @@ -1,5 +1,5 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; +import { ActivityType } from "@prisma/generated/enums"; import { AnimeReviewFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; @@ -59,9 +59,10 @@ export class AnimeReviewService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewReview, + await this.queueService.toActivityJob({ + type: ActivityType.ReviewAdded, userId: createAnimeReviewDto.userId, + animeReviewId: animeReview.id, metadata: { ...animeReview }, }); } diff --git a/src/modules/auth/config/auth.config.ts b/src/modules/auth/config/auth.config.ts index 18c94d6..2920ecf 100644 --- a/src/modules/auth/config/auth.config.ts +++ b/src/modules/auth/config/auth.config.ts @@ -1,6 +1,7 @@ import type { BetterAuthOptions } from "@better-auth/core"; import { Logger } from "@nestjs/common"; import { ConfigService } from "@nestjs/config"; +import { ActivityType } from "@prisma/generated/enums"; import * as bcrypt from "bcrypt"; import { prismaAdapter } from "better-auth/adapters/prisma"; import { bearer, customSession, lastLoginMethod, magicLink, openAPI, username } from "better-auth/plugins"; @@ -208,6 +209,23 @@ export function getAuthConfig(params: AuthConfigParams) { userId: user.id, avatarUrl: user.image, }); + + await queueService.toActivityJob({ + type: ActivityType.AccountCreated, + userId: user.id, + metadata: { + id: user.id, + name: user.name, + }, + }); + + await queueService.toSystemNotificationJob({ + recipientIds: [user.id], + metadata: { + titleKey: "notifications:welcome.title", + descriptionKey: "notifications:welcome.description", + }, + }); }, }, }, diff --git a/src/modules/book/dto/create-book-review.dto.ts b/src/modules/book/dto/create-book-review.dto.ts index d8f697e..45d9d46 100644 --- a/src/modules/book/dto/create-book-review.dto.ts +++ b/src/modules/book/dto/create-book-review.dto.ts @@ -46,7 +46,7 @@ export class CreateBookReviewDto { readonly theme?: number; @IsOptional() - @MaxLength(250) + @MaxLength(500) @ApiPropertyOptional({ type: "string", maxLength: 250, diff --git a/src/modules/book/service/book-progress.service.ts b/src/modules/book/service/book-progress.service.ts index 9b56fad..80aea0c 100644 --- a/src/modules/book/service/book-progress.service.ts +++ b/src/modules/book/service/book-progress.service.ts @@ -1,6 +1,6 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; import { BookProgressFindManyArgs } from "@prisma/generated/models"; +import { activityTypeFromProgressStatus } from "@/modules/activity/activity.utils"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; @@ -66,11 +66,16 @@ export class BookProgressService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewProgress, - userId, - metadata: { ...bookProgress }, - }); + const activityType = activityTypeFromProgressStatus(status); + + if (activityType) { + await this.queueService.toActivityJob({ + type: activityType, + userId, + bookProgressId: bookProgress.id, + metadata: { ...bookProgress }, + }); + } } async deleteBookProgress(bookProgressId: string, userId: string) { diff --git a/src/modules/book/service/book-review.service.ts b/src/modules/book/service/book-review.service.ts index b4a2a75..eb1bea1 100644 --- a/src/modules/book/service/book-review.service.ts +++ b/src/modules/book/service/book-review.service.ts @@ -1,5 +1,5 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; +import { ActivityType } from "@prisma/generated/enums"; import { BookReviewFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; @@ -55,9 +55,10 @@ export class BookReviewService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewReview, + await this.queueService.toActivityJob({ + type: ActivityType.ReviewAdded, userId: createBookReviewDto.userId, + bookReviewId: bookReview.id, metadata: { ...bookReview }, }); } diff --git a/src/modules/comment/service/comment.service.ts b/src/modules/comment/service/comment.service.ts index 961b779..66d98b5 100644 --- a/src/modules/comment/service/comment.service.ts +++ b/src/modules/comment/service/comment.service.ts @@ -3,18 +3,22 @@ import { CommentFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; +import { QueueService } from "@/shared/infra/queue/queue.service"; import { CreateCommentDto } from "../dto/create-comment.dto"; import { DeleteCommentDto } from "../dto/delete-comment.dto"; import { GetCommentsDto } from "../dto/get-comments.dto"; @Injectable() export class CommentService { - constructor(private readonly databaseService: DatabaseService) {} + constructor( + private readonly databaseService: DatabaseService, + private readonly queueService: QueueService, + ) {} async createComment(createCommentDto: CreateCommentDto) { const { type, content, userId, ...entityIds } = createCommentDto; - await this.databaseService.comment.create({ + const comment = await this.databaseService.comment.create({ data: { type, content, @@ -22,6 +26,8 @@ export class CommentService { ...entityIds, }, }); + + await this.queueService.toCommentNotificationJob({ commentId: comment.id }); } async deleteComment(deleteCommentDto: DeleteCommentDto) { diff --git a/src/modules/favorite/service/favorite.service.ts b/src/modules/favorite/service/favorite.service.ts index e72b5cb..53a8cf2 100644 --- a/src/modules/favorite/service/favorite.service.ts +++ b/src/modules/favorite/service/favorite.service.ts @@ -1,5 +1,5 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; +import { ActivityType } from "@prisma/generated/enums"; import { FavoriteFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; @@ -104,9 +104,10 @@ export class FavoriteService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewFavorite, + await this.queueService.toActivityJob({ + type: ActivityType.FavoriteAdded, userId, + favoriteId: favorite.id, metadata: { ...favorite }, }); } diff --git a/src/modules/game/dto/create-game-review.dto.ts b/src/modules/game/dto/create-game-review.dto.ts index e179eb8..75a52cd 100644 --- a/src/modules/game/dto/create-game-review.dto.ts +++ b/src/modules/game/dto/create-game-review.dto.ts @@ -73,7 +73,7 @@ export class CreateGameReviewDto { readonly platform?: string; @IsOptional() - @MaxLength(250) + @MaxLength(500) @ApiPropertyOptional({ type: "string", maxLength: 250, diff --git a/src/modules/game/service/game-progress.service.ts b/src/modules/game/service/game-progress.service.ts index 1c2acb7..96f504d 100644 --- a/src/modules/game/service/game-progress.service.ts +++ b/src/modules/game/service/game-progress.service.ts @@ -1,6 +1,6 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; import { GameProgressFindManyArgs } from "@prisma/generated/models"; +import { activityTypeFromProgressStatus } from "@/modules/activity/activity.utils"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; @@ -64,11 +64,16 @@ export class GameProgressService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewProgress, - userId, - metadata: { ...gameProgress }, - }); + const activityType = activityTypeFromProgressStatus(status); + + if (activityType) { + await this.queueService.toActivityJob({ + type: activityType, + userId, + gameProgressId: gameProgress.id, + metadata: { ...gameProgress }, + }); + } } async deleteGameProgress(gameProgressId: string, userId: string) { diff --git a/src/modules/game/service/game-review.service.ts b/src/modules/game/service/game-review.service.ts index c73dddb..583698c 100644 --- a/src/modules/game/service/game-review.service.ts +++ b/src/modules/game/service/game-review.service.ts @@ -1,5 +1,5 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; +import { ActivityType } from "@prisma/generated/enums"; import { GameReviewFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; @@ -57,9 +57,10 @@ export class GameReviewService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewReview, + await this.queueService.toActivityJob({ + type: ActivityType.ReviewAdded, userId: createGameReviewDto.userId, + gameReviewId: gameReview.id, metadata: { ...gameReview }, }); } diff --git a/src/modules/list/service/list.service.ts b/src/modules/list/service/list.service.ts index 703ba5b..602e4aa 100644 --- a/src/modules/list/service/list.service.ts +++ b/src/modules/list/service/list.service.ts @@ -1,5 +1,5 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; +import { ActivityType } from "@prisma/generated/enums"; import { ListFindManyArgs, ListItemFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; @@ -51,9 +51,10 @@ export class ListService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewList, + await this.queueService.toActivityJob({ + type: ActivityType.ListCreated, userId, + listId: list.id, metadata: { ...list }, }); } @@ -160,9 +161,10 @@ export class ListService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewListItem, + await this.queueService.toActivityJob({ + type: ActivityType.ListItemAdded, userId, + listItemId: listItem.id, metadata: { ...listItem }, }); } diff --git a/src/modules/manga/dto/create-manga-review.dto.ts b/src/modules/manga/dto/create-manga-review.dto.ts index d84b7b9..55def34 100644 --- a/src/modules/manga/dto/create-manga-review.dto.ts +++ b/src/modules/manga/dto/create-manga-review.dto.ts @@ -19,7 +19,7 @@ export class CreateMangaReviewDto { readonly worldbuilding?: number; @IsOptional() - @MaxLength(250) + @MaxLength(500) readonly summary?: string; @IsOptional() diff --git a/src/modules/manga/service/manga-progress.service.ts b/src/modules/manga/service/manga-progress.service.ts index 778d420..1128229 100644 --- a/src/modules/manga/service/manga-progress.service.ts +++ b/src/modules/manga/service/manga-progress.service.ts @@ -1,6 +1,6 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; import { MangaProgressFindManyArgs } from "@prisma/generated/models"; +import { activityTypeFromProgressStatus } from "@/modules/activity/activity.utils"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; @@ -78,11 +78,16 @@ export class MangaProgressService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewProgress, - userId, - metadata: { ...mangaProgress }, - }); + const activityType = activityTypeFromProgressStatus(status); + + if (activityType) { + await this.queueService.toActivityJob({ + type: activityType, + userId, + mangaProgressId: mangaProgress.id, + metadata: { ...mangaProgress }, + }); + } } async deleteMangaProgress(mangaProgressId: string, userId: string) { diff --git a/src/modules/manga/service/manga-review.service.ts b/src/modules/manga/service/manga-review.service.ts index e9a512d..e0d367a 100644 --- a/src/modules/manga/service/manga-review.service.ts +++ b/src/modules/manga/service/manga-review.service.ts @@ -1,5 +1,5 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; +import { ActivityType } from "@prisma/generated/enums"; import { MangaReviewFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; @@ -56,9 +56,10 @@ export class MangaReviewService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewReview, + await this.queueService.toActivityJob({ + type: ActivityType.ReviewAdded, userId: createMangaReviewDto.userId, + mangaReviewId: mangaReview.id, metadata: { ...mangaReview }, }); } diff --git a/src/modules/movie/dto/create-movie-review.dto.ts b/src/modules/movie/dto/create-movie-review.dto.ts index 695559e..7e8edb7 100644 --- a/src/modules/movie/dto/create-movie-review.dto.ts +++ b/src/modules/movie/dto/create-movie-review.dto.ts @@ -46,7 +46,7 @@ export class CreateMovieReviewDto { readonly acting?: number; @IsOptional() - @MaxLength(250) + @MaxLength(500) @ApiPropertyOptional({ type: "string", maxLength: 250, diff --git a/src/modules/movie/service/movie-progress.service.ts b/src/modules/movie/service/movie-progress.service.ts index f48e07d..0ea25c4 100644 --- a/src/modules/movie/service/movie-progress.service.ts +++ b/src/modules/movie/service/movie-progress.service.ts @@ -1,6 +1,6 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; import { MovieProgressFindManyArgs } from "@prisma/generated/models"; +import { activityTypeFromProgressStatus } from "@/modules/activity/activity.utils"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; @@ -59,11 +59,16 @@ export class MovieProgressService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewProgress, - userId, - metadata: { ...movieProgress }, - }); + const activityType = activityTypeFromProgressStatus(status); + + if (activityType) { + await this.queueService.toActivityJob({ + type: activityType, + userId, + movieProgressId: movieProgress.id, + metadata: { ...movieProgress }, + }); + } } async deleteMovieProgress(movieProgressId: string, userId: string) { diff --git a/src/modules/movie/service/movie-review.service.ts b/src/modules/movie/service/movie-review.service.ts index f82746c..8160788 100644 --- a/src/modules/movie/service/movie-review.service.ts +++ b/src/modules/movie/service/movie-review.service.ts @@ -1,5 +1,5 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; +import { ActivityType } from "@prisma/generated/enums"; import { MovieReviewFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; @@ -57,9 +57,10 @@ export class MovieReviewService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewReview, + await this.queueService.toActivityJob({ + type: ActivityType.ReviewAdded, userId: createMovieReviewDto.userId, + movieReviewId: movieReview.id, metadata: { ...movieReview }, }); } diff --git a/src/modules/payment/service/stripe.service.ts b/src/modules/payment/service/stripe.service.ts index c050690..1978eee 100644 --- a/src/modules/payment/service/stripe.service.ts +++ b/src/modules/payment/service/stripe.service.ts @@ -1,7 +1,7 @@ import { HttpService } from "@nestjs/axios"; import { forwardRef, Inject, Injectable, Logger } from "@nestjs/common"; import { ConfigService } from "@nestjs/config"; -import { PaymentStatus } from "@prisma/generated/enums"; +import { ActivityType, PaymentStatus } from "@prisma/generated/enums"; import { firstValueFrom } from "rxjs"; import Stripe from "stripe"; import { CACHE_KEYS } from "@/shared/constants/cache"; @@ -391,12 +391,19 @@ export class StripeService { }); if (!alreadyHasMedal) { - await this.databaseService.userMedal.create({ + const userMedal = await this.databaseService.userMedal.create({ data: { userId: payment.userId, medalId: contributorMedal.id, }, }); + + await this.queueService.toActivityJob({ + type: ActivityType.MedalEarned, + userId: payment.userId, + userMedalId: userMedal.id, + metadata: { id: userMedal.id, medal: { ...contributorMedal } }, + }); } } diff --git a/src/modules/reaction/dto/create-reaction.dto.ts b/src/modules/reaction/dto/create-reaction.dto.ts index c9a73d7..e3f639e 100644 --- a/src/modules/reaction/dto/create-reaction.dto.ts +++ b/src/modules/reaction/dto/create-reaction.dto.ts @@ -51,14 +51,14 @@ export class CreateReactionDto { readonly commentId?: string; @IsOptional() - @ReactionRequiredForType(ReactionType.FeedEvent) + @ReactionRequiredForType(ReactionType.Activity) @ApiProperty({ type: "string", format: "uuid", required: false, - description: "Required when type is FeedEvent", + description: "Required when type is Activity", }) - readonly feedEventId?: string; + readonly activityId?: string; @IsOptional() @ReactionRequiredForType(ReactionType.GameReview) diff --git a/src/modules/reaction/dto/get-reactions.dto.ts b/src/modules/reaction/dto/get-reactions.dto.ts index 005fa2f..c90108e 100644 --- a/src/modules/reaction/dto/get-reactions.dto.ts +++ b/src/modules/reaction/dto/get-reactions.dto.ts @@ -20,14 +20,14 @@ export class GetReactionsDto extends OffsetPaginationParamsDto { readonly commentId?: string; @IsOptional() - @ReactionRequiredForType(ReactionType.FeedEvent) + @ReactionRequiredForType(ReactionType.Activity) @ApiProperty({ type: "string", format: "uuid", required: false, - description: "Required when type is FeedEvent", + description: "Required when type is Activity", }) - readonly feedEventId?: string; + readonly activityId?: string; @IsOptional() @ReactionRequiredForType(ReactionType.GameReview) diff --git a/src/modules/reaction/service/reaction.service.ts b/src/modules/reaction/service/reaction.service.ts index cc9fe3c..a13651f 100644 --- a/src/modules/reaction/service/reaction.service.ts +++ b/src/modules/reaction/service/reaction.service.ts @@ -3,18 +3,22 @@ import { ReactionFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; +import { QueueService } from "@/shared/infra/queue/queue.service"; import { CreateReactionDto } from "../dto/create-reaction.dto"; import { DeleteReactionDto } from "../dto/delete-reaction.dto"; import { GetReactionsDto } from "../dto/get-reactions.dto"; @Injectable() export class ReactionService { - constructor(private readonly databaseService: DatabaseService) {} + constructor( + private readonly databaseService: DatabaseService, + private readonly queueService: QueueService, + ) {} async createReaction(createReactionDto: CreateReactionDto) { const { emoji, userId, type, ...entityIds } = createReactionDto; - await this.databaseService.reaction.create({ + const reaction = await this.databaseService.reaction.create({ data: { type, emoji, @@ -22,6 +26,8 @@ export class ReactionService { ...entityIds, }, }); + + await this.queueService.toReactionNotificationJob({ reactionId: reaction.id }); } async deleteReaction(deleteReactionDto: DeleteReactionDto) { @@ -44,7 +50,7 @@ export class ReactionService { where: { type: getReactionsDto.type, commentId: getReactionsDto.commentId, - feedEventId: getReactionsDto.feedEventId, + activityId: getReactionsDto.activityId, gameReviewId: getReactionsDto.gameReviewId, animeReviewId: getReactionsDto.animeReviewId, mangaReviewId: getReactionsDto.mangaReviewId, @@ -56,7 +62,7 @@ export class ReactionService { itemsPerPage: getReactionsDto.itemsPerPage, include: { comment: true, - feedEvent: true, + activity: true, user: { select: { id: true, diff --git a/src/modules/tv-show/dto/create-tv-show-review.dto.ts b/src/modules/tv-show/dto/create-tv-show-review.dto.ts index 3bb9230..317daec 100644 --- a/src/modules/tv-show/dto/create-tv-show-review.dto.ts +++ b/src/modules/tv-show/dto/create-tv-show-review.dto.ts @@ -46,7 +46,7 @@ export class CreateTVShowReviewDto { readonly acting?: number; @IsOptional() - @MaxLength(250) + @MaxLength(500) @ApiPropertyOptional({ type: "string", maxLength: 250, diff --git a/src/modules/tv-show/service/tv-show-episode-watch.service.ts b/src/modules/tv-show/service/tv-show-episode-watch.service.ts index ed54be1..3be3a7d 100644 --- a/src/modules/tv-show/service/tv-show-episode-watch.service.ts +++ b/src/modules/tv-show/service/tv-show-episode-watch.service.ts @@ -1,18 +1,38 @@ import { Injectable } from "@nestjs/common"; -import { WatchEpisodeStatus } from "@prisma/generated/enums"; +import { TvShowEpisodeWatch } from "@prisma/generated/client"; +import { ActivityType, WatchEpisodeStatus } from "@prisma/generated/enums"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; import { TMDBTVShowSeason } from "@/shared/infra/integrations/tmdb.service"; +import { QueueService } from "@/shared/infra/queue/queue.service"; import { CreateOrUpdateTVShowEpisodeWatchDto } from "../dto/create-or-update-tv-show-episode-watch.dto"; import { DeleteAllTVShowEpisodeWatchDto } from "../dto/delete-all-tv-show-episode-watch.dto"; import { DeleteTVShowEpisodeWatchDto } from "../dto/delete-tv-show-episode-watch.dto"; import { GetTVShowEpisodeWatchDto } from "../dto/get-tv-show-episode-watch.dto"; import { WatchAllEpisodesOfTVShowDto } from "../dto/watch-all-episodes-of-tv-show.dto"; +const WATCHED_STATUSES: WatchEpisodeStatus[] = [WatchEpisodeStatus.Watching, WatchEpisodeStatus.Completed]; + @Injectable() export class TVShowEpisodeWatchService { - constructor(private readonly databaseService: DatabaseService) {} + constructor( + private readonly databaseService: DatabaseService, + private readonly queueService: QueueService, + ) {} + + private async emitWatchedActivities(userId: string, watches: TvShowEpisodeWatch[]) { + for (const watch of watches) { + if (!WATCHED_STATUSES.includes(watch.status)) continue; + + await this.queueService.toActivityJob({ + type: ActivityType.Watched, + userId, + tvShowEpisodeWatchId: watch.id, + metadata: { ...watch }, + }); + } + } async createOrUpdateTVShowEpisodeWatch(createOrUpdateTVShowEpisodeWatchDto: CreateOrUpdateTVShowEpisodeWatchDto) { const { tvShowId, userId, episodes } = createOrUpdateTVShowEpisodeWatchDto; @@ -36,10 +56,12 @@ export class TVShowEpisodeWatchService { const batchSize = 50; + const watches: TvShowEpisodeWatch[] = []; + for (let i = 0; i < episodes.length; i += batchSize) { const batch = episodes.slice(i, i + batchSize); - await Promise.all( + const results = await Promise.all( batch.map(({ season, episode, status }) => this.databaseService.tvShowEpisodeWatch.upsert({ where: { @@ -61,7 +83,11 @@ export class TVShowEpisodeWatchService { }), ), ); + + watches.push(...results); } + + await this.emitWatchedActivities(userId, watches); } async watchAllEpisodesOfTVShow({ tvShowId, userId }: WatchAllEpisodesOfTVShowDto) { @@ -81,13 +107,15 @@ export class TVShowEpisodeWatchService { const seasons = tvShow.seasons as unknown as TMDBTVShowSeason[]; const batchSize = 50; + const watches: TvShowEpisodeWatch[] = []; + for (const season of seasons) { const episodeNumbers = Array.from({ length: season.numberOfEpisodes }, (_, i) => i + 1); for (let i = 0; i < episodeNumbers.length; i += batchSize) { const batch = episodeNumbers.slice(i, i + batchSize); - await Promise.all( + const results = await Promise.all( batch.map((episode) => this.databaseService.tvShowEpisodeWatch.upsert({ where: { @@ -109,8 +137,12 @@ export class TVShowEpisodeWatchService { }), ), ); + + watches.push(...results); } } + + await this.emitWatchedActivities(userId, watches); } async deleteTVShowEpisodeWatch({ userId, tvShowId, season, episode }: DeleteTVShowEpisodeWatchDto) { diff --git a/src/modules/tv-show/service/tv-show-progress.service.ts b/src/modules/tv-show/service/tv-show-progress.service.ts index 35e9506..06903b6 100644 --- a/src/modules/tv-show/service/tv-show-progress.service.ts +++ b/src/modules/tv-show/service/tv-show-progress.service.ts @@ -1,6 +1,7 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType, ProgressStatus } from "@prisma/generated/enums"; +import { ProgressStatus } from "@prisma/generated/enums"; import { TvShowProgressFindManyArgs } from "@prisma/generated/models"; +import { activityTypeFromProgressStatus } from "@/modules/activity/activity.utils"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; @@ -68,11 +69,16 @@ export class TVShowProgressService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewReview, - userId, - metadata: { ...tvShowProgress }, - }); + const activityType = activityTypeFromProgressStatus(status); + + if (activityType) { + await this.queueService.toActivityJob({ + type: activityType, + userId, + tvShowProgressId: tvShowProgress.id, + metadata: { ...tvShowProgress }, + }); + } if (status === ProgressStatus.Completed) { await this.tvShowEpisodeWatchService.watchAllEpisodesOfTVShow({ tvShowId, userId }); diff --git a/src/modules/tv-show/service/tv-show-review.service.ts b/src/modules/tv-show/service/tv-show-review.service.ts index d54d26f..5a46cc0 100644 --- a/src/modules/tv-show/service/tv-show-review.service.ts +++ b/src/modules/tv-show/service/tv-show-review.service.ts @@ -1,5 +1,5 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType } from "@prisma/generated/enums"; +import { ActivityType } from "@prisma/generated/enums"; import { TvShowReviewFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; @@ -56,9 +56,10 @@ export class TVShowReviewService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewReview, + await this.queueService.toActivityJob({ + type: ActivityType.ReviewAdded, userId: createTVShowReviewDto.userId, + tvShowReviewId: tvShowReview.id, metadata: { ...tvShowReview }, }); } diff --git a/src/modules/user/service/user.service.ts b/src/modules/user/service/user.service.ts index c90e428..e30d33f 100644 --- a/src/modules/user/service/user.service.ts +++ b/src/modules/user/service/user.service.ts @@ -1,5 +1,5 @@ import { Injectable } from "@nestjs/common"; -import { FeedEventType, ProgressStatus } from "@prisma/generated/enums"; +import { ActivityType, ProgressStatus } from "@prisma/generated/enums"; import { FollowingFindManyArgs, UserFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; @@ -319,9 +319,10 @@ export class UserService { }, }); - await this.queueService.toFeedEventJob({ - type: FeedEventType.NewFollower, + await this.queueService.toActivityJob({ + type: ActivityType.Followed, userId, + followingId: following.id, metadata: { ...following }, }); } From 810ca1fd16be42415303a19734c7248bd237b1bf Mon Sep 17 00:00:00 2001 From: izakdvlpr Date: Thu, 9 Jul 2026 18:11:44 -0300 Subject: [PATCH 6/8] feat: create activity and notification queue --- src/shared/constants/error-codes.ts | 4 + src/shared/constants/job.ts | 8 +- src/shared/constants/queue.ts | 4 +- src/shared/infra/email/email.service.ts | 2 +- .../processors/notification.processor.ts | 78 +++++++++++++++++++ src/shared/infra/queue/queue.module.ts | 16 ++-- src/shared/infra/queue/queue.service.ts | 50 ++++++++---- 7 files changed, 136 insertions(+), 26 deletions(-) create mode 100644 src/shared/infra/queue/processors/notification.processor.ts diff --git a/src/shared/constants/error-codes.ts b/src/shared/constants/error-codes.ts index d52f5fd..f1df79a 100644 --- a/src/shared/constants/error-codes.ts +++ b/src/shared/constants/error-codes.ts @@ -75,6 +75,10 @@ export const ERROR_CODES = { code: "REACTION_NOT_FOUND", status: 404, }, + NOTIFICATION_NOT_FOUND: { + code: "NOTIFICATION_NOT_FOUND", + status: 404, + }, TMDB_SERVICE_UNAVAILABLE: { code: "TMDB_SERVICE_UNAVAILABLE", status: 503, diff --git a/src/shared/constants/job.ts b/src/shared/constants/job.ts index 6bb7265..689b516 100644 --- a/src/shared/constants/job.ts +++ b/src/shared/constants/job.ts @@ -1,6 +1,10 @@ -export const FEED_EVENT_JOB = "feed-event-job"; +export const ACTIVITY_JOB = "activity-job"; -export const FEED_EVENT_FLUSH_AGGREGATION_JOB = "feed-event-flush-aggregation-job"; +export const NOTIFICATION_SYSTEM_JOB = "notification-system-job"; + +export const NOTIFICATION_COMMENT_JOB = "notification-comment-job"; + +export const NOTIFICATION_REACTION_JOB = "notification-reaction-job"; export const MAGIC_LINK_JOB = "magic-link-job"; diff --git a/src/shared/constants/queue.ts b/src/shared/constants/queue.ts index 2b9706e..5d53a4f 100644 --- a/src/shared/constants/queue.ts +++ b/src/shared/constants/queue.ts @@ -1,3 +1,5 @@ -export const FEED_EVENT_QUEUE = "feed-event-queue"; +export const ACTIVITY_QUEUE = "activity-queue"; + +export const NOTIFICATION_QUEUE = "notification-queue"; export const EMAIL_QUEUE = "email-queue"; diff --git a/src/shared/infra/email/email.service.ts b/src/shared/infra/email/email.service.ts index 6740032..c114eaf 100644 --- a/src/shared/infra/email/email.service.ts +++ b/src/shared/infra/email/email.service.ts @@ -66,7 +66,7 @@ export class EmailService { await this.resendService.send({ from: this.configService.get("RESEND_FROM")!, to: paymentFailedEmailDto.userEmail, - subject: "Payment failed – we'll retry automatically", + subject: "Payment failed - we'll retry automatically", html: this.getHtmlTemplate("payment-failed-email", paymentFailedEmailDto), }); } diff --git a/src/shared/infra/queue/processors/notification.processor.ts b/src/shared/infra/queue/processors/notification.processor.ts new file mode 100644 index 0000000..9fe4ac0 --- /dev/null +++ b/src/shared/infra/queue/processors/notification.processor.ts @@ -0,0 +1,78 @@ +import { OnWorkerEvent, Processor, WorkerHost } from "@nestjs/bullmq"; +import { Logger } from "@nestjs/common"; +import { Job } from "bullmq"; +import { + CreateCommentNotificationDto, + CreateReactionNotificationDto, + CreateSystemNotificationDto, +} from "@/modules/notification/dto/notification.dto"; +import { NotificationService } from "@/modules/notification/service/notification.service"; +import { NOTIFICATION_COMMENT_JOB, NOTIFICATION_REACTION_JOB, NOTIFICATION_SYSTEM_JOB } from "@/shared/constants/job"; +import { NOTIFICATION_QUEUE } from "@/shared/constants/queue"; + +export type CommentNotificationJobData = CreateCommentNotificationDto; + +export type ReactionNotificationJobData = CreateReactionNotificationDto; + +export type SystemNotificationJobData = CreateSystemNotificationDto; + +@Processor(NOTIFICATION_QUEUE, { concurrency: 10 }) +export class NotificationProcessor extends WorkerHost { + private readonly logger = new Logger(NotificationProcessor.name); + + constructor(private readonly notificationService: NotificationService) { + super(); + } + + async process(job: Job) { + if (job.name === NOTIFICATION_SYSTEM_JOB) { + await this.notificationService.createSystemNotification(job.data as SystemNotificationJobData); + + return; + } + + if (job.name === NOTIFICATION_COMMENT_JOB) { + await this.notificationService.createFromComment(job.data as CommentNotificationJobData); + + return; + } + + if (job.name === NOTIFICATION_REACTION_JOB) { + await this.notificationService.createFromReaction(job.data as ReactionNotificationJobData); + + return; + } + + throw new Error(`Unsupported notification job name: ${job.name}`); + } + + @OnWorkerEvent("active") + onActive(job: Job) { + this.logger.log( + `Processing job [${NOTIFICATION_QUEUE}] | job=${job.id} name=${job.name} attempt=${job.attemptsMade + 1}`, + ); + } + + @OnWorkerEvent("completed") + onCompleted(job: Job) { + this.logger.log(`Job completed [${NOTIFICATION_QUEUE}] | job=${job.id} name=${job.name}`); + } + + @OnWorkerEvent("failed") + onFailed(job: Job | undefined, error: Error) { + if (!job) return; + + const maxAttempts = job.opts?.attempts ?? 1; + const willRetry = job.attemptsMade < maxAttempts; + + if (willRetry) { + this.logger.warn( + `Job failed, retrying [${NOTIFICATION_QUEUE}] | job=${job.id} name=${job.name} attempt=${job.attemptsMade}/${maxAttempts} error=${error.message}`, + ); + } else { + this.logger.error( + `Job removed from queue after max attempts [${NOTIFICATION_QUEUE}] | job=${job.id} name=${job.name} attempts=${job.attemptsMade}/${maxAttempts} error=${error.message}`, + ); + } + } +} diff --git a/src/shared/infra/queue/queue.module.ts b/src/shared/infra/queue/queue.module.ts index 41c8c92..ed6707e 100644 --- a/src/shared/infra/queue/queue.module.ts +++ b/src/shared/infra/queue/queue.module.ts @@ -1,18 +1,21 @@ import { BullModule } from "@nestjs/bullmq"; import { Global, Module } from "@nestjs/common"; import { ConfigService } from "@nestjs/config"; -import { FeedEventModule } from "@/modules/feed-event/feed-event.module"; -import { EMAIL_QUEUE, FEED_EVENT_QUEUE } from "@/shared/constants/queue"; +import { ActivityModule } from "@/modules/activity/activity.module"; +import { NotificationModule } from "@/modules/notification/notification.module"; +import { ACTIVITY_QUEUE, EMAIL_QUEUE, NOTIFICATION_QUEUE } from "@/shared/constants/queue"; import { EmailModule } from "../email/email.module"; +import { ActivityProcessor } from "./processors/activity.processor"; import { EmailProcessor } from "./processors/email.processor"; -import { FeedEventProcessor } from "./processors/feed-event.processor"; +import { NotificationProcessor } from "./processors/notification.processor"; import { QueueService } from "./queue.service"; @Global() @Module({ imports: [ EmailModule, - FeedEventModule, + ActivityModule, + NotificationModule, BullModule.forRootAsync({ inject: [ConfigService], useFactory: async (configService: ConfigService) => ({ @@ -30,11 +33,12 @@ import { QueueService } from "./queue.service"; }, }), }), - BullModule.registerQueue({ name: FEED_EVENT_QUEUE }), + BullModule.registerQueue({ name: ACTIVITY_QUEUE }), + BullModule.registerQueue({ name: NOTIFICATION_QUEUE }), BullModule.registerQueue({ name: EMAIL_QUEUE }), ], controllers: [], - providers: [QueueService, FeedEventProcessor, EmailProcessor], + providers: [QueueService, ActivityProcessor, NotificationProcessor, EmailProcessor], exports: [QueueService], }) export class QueueModule {} diff --git a/src/shared/infra/queue/queue.service.ts b/src/shared/infra/queue/queue.service.ts index 9727a6a..e585462 100644 --- a/src/shared/infra/queue/queue.service.ts +++ b/src/shared/infra/queue/queue.service.ts @@ -1,30 +1,39 @@ import { InjectQueue } from "@nestjs/bullmq"; import { Injectable, Logger } from "@nestjs/common"; import { JobsOptions, Queue } from "bullmq"; -import { FeedEventDto } from "@/modules/feed-event/dto/feed-event.dto"; +import { CreateActivityDto } from "@/modules/activity/dto/activity.dto"; import { - FEED_EVENT_FLUSH_AGGREGATION_JOB, - FEED_EVENT_JOB, + CreateCommentNotificationDto, + CreateReactionNotificationDto, + CreateSystemNotificationDto, +} from "@/modules/notification/dto/notification.dto"; +import { + ACTIVITY_JOB, MAGIC_LINK_JOB, + NOTIFICATION_COMMENT_JOB, + NOTIFICATION_REACTION_JOB, + NOTIFICATION_SYSTEM_JOB, PAYMENT_FAILED_JOB, PAYMENT_SUCCESS_JOB, RESET_PASSWORD_JOB, SUBSCRIPTION_CANCELLED_JOB, } from "@/shared/constants/job"; -import { EMAIL_QUEUE, FEED_EVENT_QUEUE } from "@/shared/constants/queue"; +import { ACTIVITY_QUEUE, EMAIL_QUEUE, NOTIFICATION_QUEUE } from "@/shared/constants/queue"; import { MagicLinkEmailDto } from "../email/dto/magic-link-email.dto"; import { PaymentFailedEmailDto } from "../email/dto/payment-failed-email.dto"; import { PaymentSuccessEmailDto } from "../email/dto/payment-success-email.dto"; import { ResetPasswordEmailDto } from "../email/dto/reset-password-email.dto"; import { SubscriptionCancelledEmailDto } from "../email/dto/subscription-cancelled-email.dto"; -type QueueName = typeof EMAIL_QUEUE | typeof FEED_EVENT_QUEUE; +type QueueName = typeof EMAIL_QUEUE | typeof ACTIVITY_QUEUE | typeof NOTIFICATION_QUEUE; type JobName = - | typeof FEED_EVENT_JOB + | typeof ACTIVITY_JOB + | typeof NOTIFICATION_SYSTEM_JOB + | typeof NOTIFICATION_COMMENT_JOB + | typeof NOTIFICATION_REACTION_JOB | typeof MAGIC_LINK_JOB | typeof RESET_PASSWORD_JOB - | typeof FEED_EVENT_FLUSH_AGGREGATION_JOB | typeof PAYMENT_SUCCESS_JOB | typeof PAYMENT_FAILED_JOB | typeof SUBSCRIPTION_CANCELLED_JOB; @@ -36,15 +45,18 @@ export class QueueService { constructor( @InjectQueue(EMAIL_QUEUE) private readonly emailQueue: Queue, - @InjectQueue(FEED_EVENT_QUEUE) - private readonly feedEventQueue: Queue, + @InjectQueue(ACTIVITY_QUEUE) + private readonly activityQueue: Queue, + @InjectQueue(NOTIFICATION_QUEUE) + private readonly notificationQueue: Queue, ) {} private async addJob(queueName: QueueName, jobName: JobName, data: unknown, options?: JobsOptions) { try { const queues = { [EMAIL_QUEUE]: this.emailQueue, - [FEED_EVENT_QUEUE]: this.feedEventQueue, + [ACTIVITY_QUEUE]: this.activityQueue, + [NOTIFICATION_QUEUE]: this.notificationQueue, } as const; const job = await queues[queueName].add(jobName, data, options); @@ -55,14 +67,20 @@ export class QueueService { } } - async toFeedEventJob(feedEventDto: FeedEventDto) { - await this.addJob(FEED_EVENT_QUEUE, FEED_EVENT_JOB, feedEventDto); + async toActivityJob(createActivityDto: CreateActivityDto) { + await this.addJob(ACTIVITY_QUEUE, ACTIVITY_JOB, createActivityDto); + } + + async toSystemNotificationJob(createSystemNotificationDto: CreateSystemNotificationDto) { + await this.addJob(NOTIFICATION_QUEUE, NOTIFICATION_SYSTEM_JOB, createSystemNotificationDto); + } + + async toCommentNotificationJob(createCommentNotificationDto: CreateCommentNotificationDto) { + await this.addJob(NOTIFICATION_QUEUE, NOTIFICATION_COMMENT_JOB, createCommentNotificationDto); } - async toFeedEventFlushAggregationJob(feedEventFlushAggregationDto: { aggKey: string; windowsMs: number }) { - await this.addJob(FEED_EVENT_QUEUE, FEED_EVENT_FLUSH_AGGREGATION_JOB, feedEventFlushAggregationDto, { - delay: feedEventFlushAggregationDto.windowsMs, - }); + async toReactionNotificationJob(createReactionNotificationDto: CreateReactionNotificationDto) { + await this.addJob(NOTIFICATION_QUEUE, NOTIFICATION_REACTION_JOB, createReactionNotificationDto); } async toMagicLinkJob(magicLinkEmailDto: MagicLinkEmailDto) { From 118fb3ead41ffd671b9dba32363a2865bbf828ff Mon Sep 17 00:00:00 2001 From: izakdvlpr Date: Thu, 9 Jul 2026 18:12:40 -0300 Subject: [PATCH 7/8] chore(biome): fix format --- prisma/seed.ts | 6 +++--- .../activity/controller/activity.controller.ts | 6 +++--- src/modules/activity/service/activity.service.ts | 11 +++++++---- 3 files changed, 13 insertions(+), 10 deletions(-) diff --git a/prisma/seed.ts b/prisma/seed.ts index 255ea8d..9629b3b 100644 --- a/prisma/seed.ts +++ b/prisma/seed.ts @@ -46,7 +46,7 @@ export async function populateMedals(prisma: PrismaClient) { ], skipDuplicates: true, }); - + if (medals.count > 0) { console.log(`Inserted ${medals.count} medals.`); } @@ -78,7 +78,7 @@ async function createFirstUser(prisma: PrismaClient) { }); if (!userExists) { - const createdUser = await prisma.user.create({ + await prisma.user.create({ data: { id: userData.id, email: userData.email, @@ -105,7 +105,7 @@ async function createFirstUser(prisma: PrismaClient) { insertedCount++; } } - + if (insertedCount > 0) { console.log(`Inserted ${insertedCount} users.`); } diff --git a/src/modules/activity/controller/activity.controller.ts b/src/modules/activity/controller/activity.controller.ts index afd467b..cf9a57a 100644 --- a/src/modules/activity/controller/activity.controller.ts +++ b/src/modules/activity/controller/activity.controller.ts @@ -3,8 +3,8 @@ import { ApiTags } from "@nestjs/swagger"; import { AuthGuard, Session, type UserSession } from "@thallesp/nestjs-better-auth"; import { GetActivitiesDto } from "../dto/get-activities.dto"; import { GetActivitiesByUserDto } from "../dto/get-activities-by-user.dto"; +import { GetActivitiesByUserFollowingDto } from "../dto/get-activities-by-user-following.dto"; import { ActivityService } from "../service/activity.service"; -import { GetActivitiesByUserFollowingDto } from '../dto/get-activities-by-user-following.dto'; @ApiTags("Activity") @Controller("/activities") @@ -27,14 +27,14 @@ export class ActivityController { return { activities }; } - + @Get("/user/:userId/calendar") async getUserActivityCalendarById(@Param("userId") userId: string) { const activityCalendar = await this.activityService.getUserActivityCalendarById(userId); return { activityCalendar }; } - + @Get("/following") @UseGuards(AuthGuard) async getActivitiesByUserFollowing(@Session() session: UserSession, @Query() query: GetActivitiesByUserFollowingDto) { diff --git a/src/modules/activity/service/activity.service.ts b/src/modules/activity/service/activity.service.ts index 4d6905f..6283156 100644 --- a/src/modules/activity/service/activity.service.ts +++ b/src/modules/activity/service/activity.service.ts @@ -6,7 +6,7 @@ import { DatabaseService } from "@/shared/infra/database/database.service"; import { CreateActivityDto } from "../dto/activity.dto"; import { GetActivitiesDto } from "../dto/get-activities.dto"; import { GetActivitiesByUserDto } from "../dto/get-activities-by-user.dto"; -import { GetActivitiesByUserFollowingDto } from '../dto/get-activities-by-user-following.dto'; +import { GetActivitiesByUserFollowingDto } from "../dto/get-activities-by-user-following.dto"; const SOURCE_FIELDS = [ "listId", @@ -232,7 +232,7 @@ export class ActivityService { return { ...pagination, items: this.groupActivities(pagination.items) }; } - + async getActivitiesByUserFollowing(getActivitiesByUserIdDto: GetActivitiesByUserFollowingDto) { const following = await this.databaseService.following.findMany({ where: { followerId: getActivitiesByUserIdDto.userId }, @@ -241,7 +241,10 @@ export class ActivityService { const friendIds = following.map((item) => item.followingId); - return this.getActivitiesByUserId({ ...getActivitiesByUserIdDto, userId: getActivitiesByUserIdDto.userId }, friendIds); + return this.getActivitiesByUserId( + { ...getActivitiesByUserIdDto, userId: getActivitiesByUserIdDto.userId }, + friendIds, + ); } async getActivities(getActivitiesDto: GetActivitiesDto) { @@ -286,7 +289,7 @@ export class ActivityService { return groups; } - + async getUserActivityCalendarById(userId: string) { const userExists = await this.databaseService.user.findUnique({ where: { id: userId }, From f16dc52954be1677fdec345b2f6d8fc1411bacdc Mon Sep 17 00:00:00 2001 From: izakdvlpr Date: Fri, 10 Jul 2026 14:02:10 -0300 Subject: [PATCH 8/8] feat: activity watch episode --- .../migration.sql | 16 ++++- prisma/schema.prisma | 10 ++++ src/modules/activity/activity.utils.ts | 12 +++- .../activity/dto/sync-watched-activity.dto.ts | 20 +++++++ .../activity/service/activity.service.ts | 60 +++++++++++++++++-- .../service/anime-episode-watch.service.ts | 42 +++++-------- .../service/tv-show-episode-watch.service.ts | 42 +++++-------- src/shared/constants/job.ts | 2 + .../queue/processors/activity.processor.ts | 9 ++- src/shared/infra/queue/queue.service.ts | 7 +++ 10 files changed, 155 insertions(+), 65 deletions(-) rename prisma/migrations/{20260709183715_init => 20260710165614_init}/migration.sql (98%) create mode 100644 src/modules/activity/dto/sync-watched-activity.dto.ts diff --git a/prisma/migrations/20260709183715_init/migration.sql b/prisma/migrations/20260710165614_init/migration.sql similarity index 98% rename from prisma/migrations/20260709183715_init/migration.sql rename to prisma/migrations/20260710165614_init/migration.sql index fb24fe9..e57fd75 100644 --- a/prisma/migrations/20260709183715_init/migration.sql +++ b/prisma/migrations/20260710165614_init/migration.sql @@ -14,7 +14,7 @@ CREATE TYPE "ReactionType" AS ENUM ('Comment', 'Activity', 'GameReview', 'AnimeR CREATE TYPE "NotificationType" AS ENUM ('System', 'CommentOnProfile', 'ReactionOnComment', 'ReactionOnActivity', 'ReactionOnAnimeReview', 'ReactionOnMangaReview', 'ReactionOnTvShowReview', 'ReactionOnMovieReview', 'ReactionOnGameReview', 'ReactionOnBookReview'); -- CreateEnum -CREATE TYPE "ActivityType" AS ENUM ('AccountCreated', 'ListCreated', 'ListItemAdded', 'FavoriteAdded', 'ReviewAdded', 'ProgressStarted', 'ProgressCompleted', 'Watched', 'Followed', 'MedalEarned'); +CREATE TYPE "ActivityType" AS ENUM ('AccountCreated', 'ListCreated', 'ListItemAdded', 'FavoriteAdded', 'ReviewAdded', 'ProgressStarted', 'ProgressCompleted', 'ProgressPaused', 'ProgressDropped', 'Watched', 'Followed', 'MedalEarned'); -- CreateEnum CREATE TYPE "WatchEpisodeStatus" AS ENUM ('NotWatched', 'Watching', 'Completed', 'Paused', 'Dropped', 'Planning'); @@ -235,6 +235,8 @@ CREATE TABLE "Activity" ( "tvShowEpisodeWatchId" TEXT, "followingId" TEXT, "userMedalId" TEXT, + "animeId" TEXT, + "tvShowId" TEXT, CONSTRAINT "Activity_pkey" PRIMARY KEY ("id") ); @@ -917,6 +919,12 @@ CREATE INDEX "Activity_userId_idx" ON "Activity"("userId"); -- CreateIndex CREATE INDEX "Activity_type_idx" ON "Activity"("type"); +-- CreateIndex +CREATE INDEX "Activity_userId_type_animeId_idx" ON "Activity"("userId", "type", "animeId"); + +-- CreateIndex +CREATE INDEX "Activity_userId_type_tvShowId_idx" ON "Activity"("userId", "type", "tvShowId"); + -- CreateIndex CREATE UNIQUE INDEX "Game_igdbId_key" ON "Game"("igdbId"); @@ -1205,6 +1213,12 @@ ALTER TABLE "Activity" ADD CONSTRAINT "Activity_followingId_fkey" FOREIGN KEY (" -- AddForeignKey ALTER TABLE "Activity" ADD CONSTRAINT "Activity_userMedalId_fkey" FOREIGN KEY ("userMedalId") REFERENCES "UserMedal"("id") ON DELETE CASCADE ON UPDATE CASCADE; +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_animeId_fkey" FOREIGN KEY ("animeId") REFERENCES "Anime"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "Activity" ADD CONSTRAINT "Activity_tvShowId_fkey" FOREIGN KEY ("tvShowId") REFERENCES "TVShow"("id") ON DELETE CASCADE ON UPDATE CASCADE; + -- AddForeignKey ALTER TABLE "AnimeEpisodeWatch" ADD CONSTRAINT "AnimeEpisodeWatch_userId_fkey" FOREIGN KEY ("userId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE; diff --git a/prisma/schema.prisma b/prisma/schema.prisma index 6b20e4b..34891db 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -327,6 +327,8 @@ enum ActivityType { ReviewAdded ProgressStarted ProgressCompleted + ProgressPaused + ProgressDropped Watched Followed MedalEarned @@ -358,6 +360,8 @@ model Activity { tvShowEpisodeWatchId String? followingId String? userMedalId String? + animeId String? + tvShowId String? user User @relation(fields: [userId], references: [id], onDelete: Cascade) reactions Reaction[] @@ -380,11 +384,15 @@ model Activity { tvShowEpisodeWatch TvShowEpisodeWatch? @relation(fields: [tvShowEpisodeWatchId], references: [id], onDelete: Cascade) following Following? @relation(fields: [followingId], references: [id], onDelete: Cascade) userMedal UserMedal? @relation(fields: [userMedalId], references: [id], onDelete: Cascade) + anime Anime? @relation(fields: [animeId], references: [id], onDelete: Cascade) + tvShow TvShow? @relation(fields: [tvShowId], references: [id], onDelete: Cascade) notifications Notification[] @@index([userId]) @@index([type]) + @@index([userId, type, animeId]) + @@index([userId, type, tvShowId]) } model Game { @@ -543,6 +551,7 @@ model Anime { favorites Favorite[] listItems ListItem[] comments Comment[] + activities Activity[] } model Manga { @@ -630,6 +639,7 @@ model TvShow { favorites Favorite[] listItems ListItem[] comments Comment[] + activities Activity[] @@map("TVShow") } diff --git a/src/modules/activity/activity.utils.ts b/src/modules/activity/activity.utils.ts index 405f8a7..0a89362 100644 --- a/src/modules/activity/activity.utils.ts +++ b/src/modules/activity/activity.utils.ts @@ -3,13 +3,21 @@ import { ActivityType, ProgressStatus } from "@prisma/generated/enums"; const STARTED_STATUSES: ProgressStatus[] = [ProgressStatus.Watching, ProgressStatus.Playing, ProgressStatus.Reading]; // Maps a progress status to the activity it should emit. -// Only "started" (Watching/Playing/Reading) and Completed generate activities; -// every other status returns null (no activity). +// Started (Watching/Playing/Reading), Completed, Paused and Dropped generate +// activities; every other status returns null (no activity). export function activityTypeFromProgressStatus(status: ProgressStatus): ActivityType | null { if (status === ProgressStatus.Completed) { return ActivityType.ProgressCompleted; } + if (status === ProgressStatus.Paused) { + return ActivityType.ProgressPaused; + } + + if (status === ProgressStatus.Dropped) { + return ActivityType.ProgressDropped; + } + if (STARTED_STATUSES.includes(status)) { return ActivityType.ProgressStarted; } diff --git a/src/modules/activity/dto/sync-watched-activity.dto.ts b/src/modules/activity/dto/sync-watched-activity.dto.ts new file mode 100644 index 0000000..75f7dc9 --- /dev/null +++ b/src/modules/activity/dto/sync-watched-activity.dto.ts @@ -0,0 +1,20 @@ +import { ArrayNotEmpty, IsArray, IsInt, IsPositive, IsUUID, ValidateIf } from "class-validator"; + +export class SyncWatchedActivityDto { + @IsUUID("7") + readonly userId: string; + + @ValidateIf((dto: SyncWatchedActivityDto) => !dto.tvShowId) + @IsUUID("7") + readonly animeId?: string; + + @ValidateIf((dto: SyncWatchedActivityDto) => !dto.animeId) + @IsUUID("7") + readonly tvShowId?: string; + + @IsArray() + @ArrayNotEmpty() + @IsInt({ each: true }) + @IsPositive({ each: true }) + readonly episodes: number[]; +} diff --git a/src/modules/activity/service/activity.service.ts b/src/modules/activity/service/activity.service.ts index 6283156..192ead0 100644 --- a/src/modules/activity/service/activity.service.ts +++ b/src/modules/activity/service/activity.service.ts @@ -1,4 +1,5 @@ import { Injectable } from "@nestjs/common"; +import { ActivityType } from "@prisma/generated/enums"; import { ActivityFindManyArgs } from "@prisma/generated/models"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; @@ -7,6 +8,7 @@ import { CreateActivityDto } from "../dto/activity.dto"; import { GetActivitiesDto } from "../dto/get-activities.dto"; import { GetActivitiesByUserDto } from "../dto/get-activities-by-user.dto"; import { GetActivitiesByUserFollowingDto } from "../dto/get-activities-by-user-following.dto"; +import { SyncWatchedActivityDto } from "../dto/sync-watched-activity.dto"; const SOURCE_FIELDS = [ "listId", @@ -161,11 +163,10 @@ const INCLUDE = { movieProgress: { select: { id: true, status: true, movie: MOVIE_SELECT } }, gameProgress: { select: { id: true, status: true, game: GAME_SELECT } }, bookProgress: { select: { id: true, status: true, book: BOOK_SELECT } }, - // Episode watches. - animeEpisodeWatch: { select: { id: true, status: true, episode: true, anime: ANIME_SELECT } }, - tvShowEpisodeWatch: { - select: { id: true, status: true, season: true, episode: true, tvShow: TVSHOW_SELECT }, - }, + // Watched episodes are now a single per-series activity holding a range in + // metadata; the media is referenced directly by the anime/tvShow relation. + anime: ANIME_SELECT, + tvShow: TVSHOW_SELECT, // Lists / favorites (polymorphic media). list: { select: { id: true, name: true } }, listItem: { select: { id: true, ...POLY_MEDIA } }, @@ -208,6 +209,52 @@ export class ActivityService { }); } + // One Watched activity per (user, series) per 1h window, carrying an episode + // range in metadata. Watching more episodes of the same series within the + // window extends the existing activity in place (keeps its id/reactions); + // after the window closes, the next episode starts a fresh activity. + async syncWatchedActivity(syncWatchedActivityDto: SyncWatchedActivityDto) { + const WINDOW_MS = 60 * 60 * 1000; + const { userId, animeId, tvShowId, episodes } = syncWatchedActivityDto; + + if (episodes.length === 0) return; + + const min = Math.min(...episodes); + const max = Math.max(...episodes); + const seriesWhere = animeId ? { animeId } : { tvShowId }; + + const recent = await this.databaseService.activity.findFirst({ + where: { userId, type: ActivityType.Watched, ...seriesWhere }, + orderBy: { createdAt: "desc" }, + }); + + if (recent && Date.now() - recent.createdAt.getTime() <= WINDOW_MS) { + const meta = (recent.metadata ?? {}) as { from?: number; to?: number; count?: number }; + + await this.databaseService.activity.update({ + where: { id: recent.id }, + data: { + metadata: { + from: Math.min(meta.from ?? min, min), + to: Math.max(meta.to ?? max, max), + count: (meta.count ?? 0) + episodes.length, + }, + }, + }); + + return; + } + + await this.databaseService.activity.create({ + data: { + type: ActivityType.Watched, + userId, + ...seriesWhere, + metadata: { from: min, to: max, count: episodes.length }, + }, + }); + } + async getActivitiesByUserId(getActivitiesByUserIdDto: GetActivitiesByUserDto, friendIds: string[] = []) { const userExists = await this.databaseService.user.findUnique({ where: { id: getActivitiesByUserIdDto.userId }, @@ -261,6 +308,8 @@ export class ActivityService { // Collapse consecutive activities of the same (userId, type) inside a 1h window // into a single group with a count, so the feed can render "favorited 3 animes". + // Watched is never merged here — it's already a single per-series activity + // (see syncWatchedActivity), so each one stands as its own group. private groupActivities(items: any[]) { const WINDOW_MS = 60 * 60 * 1000; const groups: any[] = []; @@ -269,6 +318,7 @@ export class ActivityService { const last = groups[groups.length - 1]; const sameGroup = last && + item.type !== "Watched" && last.type === item.type && last.userId === item.userId && new Date(last.createdAt).getTime() - new Date(item.createdAt).getTime() <= WINDOW_MS; diff --git a/src/modules/anime/service/anime-episode-watch.service.ts b/src/modules/anime/service/anime-episode-watch.service.ts index 2d754ce..ef73006 100644 --- a/src/modules/anime/service/anime-episode-watch.service.ts +++ b/src/modules/anime/service/anime-episode-watch.service.ts @@ -1,6 +1,5 @@ import { Injectable } from "@nestjs/common"; -import { AnimeEpisodeWatch } from "@prisma/generated/client"; -import { ActivityType, WatchEpisodeStatus } from "@prisma/generated/enums"; +import { WatchEpisodeStatus } from "@prisma/generated/enums"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; @@ -20,19 +19,6 @@ export class AnimeEpisodeWatchService { private readonly queueService: QueueService, ) {} - private async emitWatchedActivities(userId: string, watches: AnimeEpisodeWatch[]) { - for (const watch of watches) { - if (!WATCHED_STATUSES.includes(watch.status)) continue; - - await this.queueService.toActivityJob({ - type: ActivityType.Watched, - userId, - animeEpisodeWatchId: watch.id, - metadata: { ...watch }, - }); - } - } - async createOrUpdateAnimeEpisodeWatch(createOrUpdateAnimeEpisodeWatchDto: CreateOrUpdateAnimeEpisodeWatchDto) { const { animeId, userId, episodes } = createOrUpdateAnimeEpisodeWatchDto; @@ -55,12 +41,10 @@ export class AnimeEpisodeWatchService { const batchSize = 50; - const watches: AnimeEpisodeWatch[] = []; - for (let i = 0; i < episodes.length; i += batchSize) { const batch = episodes.slice(i, i + batchSize); - const results = await Promise.all( + await Promise.all( batch.map(({ episode, status }) => this.databaseService.animeEpisodeWatch.upsert({ where: { @@ -80,11 +64,17 @@ export class AnimeEpisodeWatchService { }), ), ); - - watches.push(...results); } - await this.emitWatchedActivities(userId, watches); + // Manual watching feeds a single per-series Watched activity (range in + // metadata); bulk "watch all" / series completion stays silent. + const watchedEpisodes = episodes + .filter(({ status }) => WATCHED_STATUSES.includes(status)) + .map(({ episode }) => episode); + + if (watchedEpisodes.length > 0) { + await this.queueService.toWatchedActivityJob({ userId, animeId, episodes: watchedEpisodes }); + } } async watchAllAnimeEpisodes({ animeId, userId }: WatchAllAnimeEpisodesDto) { @@ -104,12 +94,12 @@ export class AnimeEpisodeWatchService { const episodeNumbers = Array.from({ length: anime.numberOfEpisodes }, (_, i) => i + 1); const batchSize = 50; - const watches: AnimeEpisodeWatch[] = []; - + // Bulk mark (series completion / "watch all"): intentionally emits no + // Watched activities — those only come from manual, one-by-one watching. for (let i = 0; i < episodeNumbers.length; i += batchSize) { const batch = episodeNumbers.slice(i, i + batchSize); - const results = await Promise.all( + await Promise.all( batch.map((episode) => this.databaseService.animeEpisodeWatch.upsert({ where: { @@ -129,11 +119,7 @@ export class AnimeEpisodeWatchService { }), ), ); - - watches.push(...results); } - - await this.emitWatchedActivities(userId, watches); } async deleteAnimeEpisodeWatch({ userId, animeId, episode }: DeleteAnimeEpisodeWatchDto) { diff --git a/src/modules/tv-show/service/tv-show-episode-watch.service.ts b/src/modules/tv-show/service/tv-show-episode-watch.service.ts index 3be3a7d..dd3dba8 100644 --- a/src/modules/tv-show/service/tv-show-episode-watch.service.ts +++ b/src/modules/tv-show/service/tv-show-episode-watch.service.ts @@ -1,6 +1,5 @@ import { Injectable } from "@nestjs/common"; -import { TvShowEpisodeWatch } from "@prisma/generated/client"; -import { ActivityType, WatchEpisodeStatus } from "@prisma/generated/enums"; +import { WatchEpisodeStatus } from "@prisma/generated/enums"; import { ERROR_CODES } from "@/shared/constants/error-codes"; import { AppException } from "@/shared/exceptions/app.exceptions"; import { DatabaseService } from "@/shared/infra/database/database.service"; @@ -21,19 +20,6 @@ export class TVShowEpisodeWatchService { private readonly queueService: QueueService, ) {} - private async emitWatchedActivities(userId: string, watches: TvShowEpisodeWatch[]) { - for (const watch of watches) { - if (!WATCHED_STATUSES.includes(watch.status)) continue; - - await this.queueService.toActivityJob({ - type: ActivityType.Watched, - userId, - tvShowEpisodeWatchId: watch.id, - metadata: { ...watch }, - }); - } - } - async createOrUpdateTVShowEpisodeWatch(createOrUpdateTVShowEpisodeWatchDto: CreateOrUpdateTVShowEpisodeWatchDto) { const { tvShowId, userId, episodes } = createOrUpdateTVShowEpisodeWatchDto; @@ -56,12 +42,10 @@ export class TVShowEpisodeWatchService { const batchSize = 50; - const watches: TvShowEpisodeWatch[] = []; - for (let i = 0; i < episodes.length; i += batchSize) { const batch = episodes.slice(i, i + batchSize); - const results = await Promise.all( + await Promise.all( batch.map(({ season, episode, status }) => this.databaseService.tvShowEpisodeWatch.upsert({ where: { @@ -83,11 +67,17 @@ export class TVShowEpisodeWatchService { }), ), ); - - watches.push(...results); } - await this.emitWatchedActivities(userId, watches); + // Manual watching feeds a single per-series Watched activity (range in + // metadata); bulk "watch all" / series completion stays silent. + const watchedEpisodes = episodes + .filter(({ status }) => WATCHED_STATUSES.includes(status)) + .map(({ episode }) => episode); + + if (watchedEpisodes.length > 0) { + await this.queueService.toWatchedActivityJob({ userId, tvShowId, episodes: watchedEpisodes }); + } } async watchAllEpisodesOfTVShow({ tvShowId, userId }: WatchAllEpisodesOfTVShowDto) { @@ -107,15 +97,15 @@ export class TVShowEpisodeWatchService { const seasons = tvShow.seasons as unknown as TMDBTVShowSeason[]; const batchSize = 50; - const watches: TvShowEpisodeWatch[] = []; - + // Bulk mark (series completion / "watch all"): intentionally emits no + // Watched activities — those only come from manual, one-by-one watching. for (const season of seasons) { const episodeNumbers = Array.from({ length: season.numberOfEpisodes }, (_, i) => i + 1); for (let i = 0; i < episodeNumbers.length; i += batchSize) { const batch = episodeNumbers.slice(i, i + batchSize); - const results = await Promise.all( + await Promise.all( batch.map((episode) => this.databaseService.tvShowEpisodeWatch.upsert({ where: { @@ -137,12 +127,8 @@ export class TVShowEpisodeWatchService { }), ), ); - - watches.push(...results); } } - - await this.emitWatchedActivities(userId, watches); } async deleteTVShowEpisodeWatch({ userId, tvShowId, season, episode }: DeleteTVShowEpisodeWatchDto) { diff --git a/src/shared/constants/job.ts b/src/shared/constants/job.ts index 689b516..3512fe5 100644 --- a/src/shared/constants/job.ts +++ b/src/shared/constants/job.ts @@ -1,5 +1,7 @@ export const ACTIVITY_JOB = "activity-job"; +export const WATCHED_ACTIVITY_JOB = "watched-activity-job"; + export const NOTIFICATION_SYSTEM_JOB = "notification-system-job"; export const NOTIFICATION_COMMENT_JOB = "notification-comment-job"; diff --git a/src/shared/infra/queue/processors/activity.processor.ts b/src/shared/infra/queue/processors/activity.processor.ts index a58c5d1..0203223 100644 --- a/src/shared/infra/queue/processors/activity.processor.ts +++ b/src/shared/infra/queue/processors/activity.processor.ts @@ -2,8 +2,9 @@ import { OnWorkerEvent, Processor, WorkerHost } from "@nestjs/bullmq"; import { Logger } from "@nestjs/common"; import { Job } from "bullmq"; import { CreateActivityDto } from "@/modules/activity/dto/activity.dto"; +import { SyncWatchedActivityDto } from "@/modules/activity/dto/sync-watched-activity.dto"; import { ActivityService } from "@/modules/activity/service/activity.service"; -import { ACTIVITY_JOB } from "@/shared/constants/job"; +import { ACTIVITY_JOB, WATCHED_ACTIVITY_JOB } from "@/shared/constants/job"; import { ACTIVITY_QUEUE } from "@/shared/constants/queue"; export type ActivityJobData = CreateActivityDto; @@ -23,6 +24,12 @@ export class ActivityProcessor extends WorkerHost { return; } + if (job.name === WATCHED_ACTIVITY_JOB) { + await this.activityService.syncWatchedActivity(job.data as SyncWatchedActivityDto); + + return; + } + throw new Error(`Unsupported activity job name: ${job.name}`); } diff --git a/src/shared/infra/queue/queue.service.ts b/src/shared/infra/queue/queue.service.ts index e585462..4eed3c3 100644 --- a/src/shared/infra/queue/queue.service.ts +++ b/src/shared/infra/queue/queue.service.ts @@ -2,6 +2,7 @@ import { InjectQueue } from "@nestjs/bullmq"; import { Injectable, Logger } from "@nestjs/common"; import { JobsOptions, Queue } from "bullmq"; import { CreateActivityDto } from "@/modules/activity/dto/activity.dto"; +import { SyncWatchedActivityDto } from "@/modules/activity/dto/sync-watched-activity.dto"; import { CreateCommentNotificationDto, CreateReactionNotificationDto, @@ -17,6 +18,7 @@ import { PAYMENT_SUCCESS_JOB, RESET_PASSWORD_JOB, SUBSCRIPTION_CANCELLED_JOB, + WATCHED_ACTIVITY_JOB, } from "@/shared/constants/job"; import { ACTIVITY_QUEUE, EMAIL_QUEUE, NOTIFICATION_QUEUE } from "@/shared/constants/queue"; import { MagicLinkEmailDto } from "../email/dto/magic-link-email.dto"; @@ -29,6 +31,7 @@ type QueueName = typeof EMAIL_QUEUE | typeof ACTIVITY_QUEUE | typeof NOTIFICATIO type JobName = | typeof ACTIVITY_JOB + | typeof WATCHED_ACTIVITY_JOB | typeof NOTIFICATION_SYSTEM_JOB | typeof NOTIFICATION_COMMENT_JOB | typeof NOTIFICATION_REACTION_JOB @@ -71,6 +74,10 @@ export class QueueService { await this.addJob(ACTIVITY_QUEUE, ACTIVITY_JOB, createActivityDto); } + async toWatchedActivityJob(syncWatchedActivityDto: SyncWatchedActivityDto) { + await this.addJob(ACTIVITY_QUEUE, WATCHED_ACTIVITY_JOB, syncWatchedActivityDto); + } + async toSystemNotificationJob(createSystemNotificationDto: CreateSystemNotificationDto) { await this.addJob(NOTIFICATION_QUEUE, NOTIFICATION_SYSTEM_JOB, createSystemNotificationDto); }