diff --git a/src/database/migrations/20260818000000_create_leaderboard_and_streaks_tables.ts b/src/database/migrations/20260818000000_create_leaderboard_and_streaks_tables.ts new file mode 100644 index 00000000..aa9a45a5 --- /dev/null +++ b/src/database/migrations/20260818000000_create_leaderboard_and_streaks_tables.ts @@ -0,0 +1,101 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +/** + * Migration: create leaderboard_snapshots and user_activity_streaks tables + * + * leaderboard_snapshots — stores pre-computed leaderboard results from the + * nightly job so the API can respond in < 50ms from a single indexed lookup + * instead of running complex multi-join SQL on every request. + * + * user_activity_streaks — tracks consecutive-day activity (at least one + * milestone completion per day) so streak_days can be served from the DB + * and incremented via the daily streak cron job. + */ +export class CreateLeaderboardAndStreaksTables1755561600000 + implements MigrationInterface +{ + public async up(queryRunner: QueryRunner): Promise { + // ── leaderboard_snapshots ──────────────────────────────────────────────── + await queryRunner.query(` + CREATE TABLE IF NOT EXISTS leaderboard_snapshots ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + type VARCHAR(20) NOT NULL CHECK (type IN ('milestone','path','global')), + target_id UUID NULL, + period VARCHAR(20) NOT NULL CHECK (period IN ('week','month','quarter','all')), + entries JSONB NOT NULL DEFAULT '[]', + computed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() + ); + `); + + /* Each (type, target_id, period) combination has at most one active + snapshot. A NULL target_id is used for global leaderboards so we use + COALESCE to make the unique constraint work with NULLs. */ + await queryRunner.query(` + CREATE UNIQUE INDEX IF NOT EXISTS uq_leaderboard_snapshots_key + ON leaderboard_snapshots (type, COALESCE(target_id::text, ''), period); + `); + + /* Fast lookup by computed_at for pruning stale snapshots */ + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS idx_leaderboard_snapshots_computed_at + ON leaderboard_snapshots (computed_at DESC); + `); + + // ── user_activity_streaks ──────────────────────────────────────────────── + await queryRunner.query(` + CREATE TABLE IF NOT EXISTS user_activity_streaks ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE, + current_streak INT NOT NULL DEFAULT 0, + longest_streak INT NOT NULL DEFAULT 0, + last_active_date DATE NULL, + streak_started_at DATE NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + CONSTRAINT uq_user_activity_streaks_user_id UNIQUE (user_id) + ); + `); + + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS idx_user_activity_streaks_user_id + ON user_activity_streaks (user_id); + `); + + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS idx_user_activity_streaks_last_active + ON user_activity_streaks (last_active_date DESC); + `); + + // ── peer_review_votes ─────────────────────────────────────────────────── + /* Tracks which users "liked" a peer review so the anti-gaming check can + verify that a review received at least one liked vote before counting + toward the reviewer's helpfulReviews score. */ + await queryRunner.query(` + CREATE TABLE IF NOT EXISTS peer_review_votes ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + review_id UUID NOT NULL, + voter_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + CONSTRAINT uq_peer_review_votes UNIQUE (review_id, voter_id) + ); + `); + + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS idx_peer_review_votes_review_id + ON peer_review_votes (review_id); + `); + + await queryRunner.query(` + CREATE INDEX IF NOT EXISTS idx_peer_review_votes_voter_id + ON peer_review_votes (voter_id); + `); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.query(`DROP TABLE IF EXISTS peer_review_votes;`); + await queryRunner.query(`DROP TABLE IF EXISTS user_activity_streaks;`); + await queryRunner.query(`DROP TABLE IF EXISTS leaderboard_snapshots;`); + } +} diff --git a/src/jobs/leaderboardPrecompute.job.ts b/src/jobs/leaderboardPrecompute.job.ts new file mode 100644 index 00000000..311cf88b --- /dev/null +++ b/src/jobs/leaderboardPrecompute.job.ts @@ -0,0 +1,170 @@ +/** + * Nightly Leaderboard Pre-computation Job + * + * Computes leaderboard snapshots for all active learning paths and stores the + * results in leaderboard_snapshots. The leaderboard API then reads from this + * table (< 50 ms) instead of running the full multi-join SQL on every request. + * + * Acceptance criteria: + * - Runs within 30 minutes for 10,000 active users + * - Each user appears exactly once per leaderboard (GROUP BY bug fixed) + * - Only peer reviews with at least 1 liked vote count toward helpfulReviews + * + * Scheduled: nightly at 02:30 UTC (after the analytics refresh window) + */ + +import pool from '../config/database'; +import { CollaborativeLearningService } from '../services/collaborative-learning.service'; +import { logger } from '../utils/logger.utils'; +import { CacheService } from '../services/cache.service'; + +type LeaderboardType = 'milestone' | 'path' | 'global'; +type LeaderboardPeriod = 'week' | 'month' | 'quarter' | 'all'; + +const ALL_PERIODS: LeaderboardPeriod[] = ['week', 'month', 'quarter', 'all']; +const ALL_TYPES: LeaderboardType[] = ['milestone', 'path', 'global']; + +/** + * Run the nightly leaderboard pre-computation. + * Computes snapshots for every (type × period) combination and upserts them + * into leaderboard_snapshots. Path and milestone types additionally compute + * per-entity snapshots for each active learning path / milestone. + */ +export async function runLeaderboardPrecompute(): Promise<{ + elapsed: number; + snapshots: number; + errors: number; +}> { + const startTime = Date.now(); + let snapshotCount = 0; + let errorCount = 0; + + logger.info('[LeaderboardPrecompute] Starting nightly leaderboard pre-computation'); + + try { + // ── 1. Global leaderboards ───────────────────────────────────────────── + for (const period of ALL_PERIODS) { + try { + await upsertSnapshot('global', undefined, period); + snapshotCount++; + } catch (err) { + errorCount++; + logger.error('[LeaderboardPrecompute] Failed to compute global snapshot', { + period, + error: (err as Error).message, + }); + } + } + + // ── 2. Per-path leaderboards ─────────────────────────────────────────── + const { rows: activePaths } = await pool.query<{ id: string }>( + `SELECT DISTINCT id FROM learning_paths WHERE is_published = true` + ); + + logger.info('[LeaderboardPrecompute] Computing path-level snapshots', { + pathCount: activePaths.length, + }); + + for (const { id: pathId } of activePaths) { + for (const period of ALL_PERIODS) { + try { + await upsertSnapshot('path', pathId, period); + snapshotCount++; + } catch (err) { + errorCount++; + logger.error('[LeaderboardPrecompute] Failed path snapshot', { + pathId, + period, + error: (err as Error).message, + }); + } + } + } + + // ── 3. Per-milestone leaderboards ────────────────────────────────────── + // Limit to milestones that have had at least 5 completions to avoid + // sparse snapshots for newly created milestones. + const { rows: activeMilestones } = await pool.query<{ id: string }>( + `SELECT m.id + FROM milestones m + JOIN learning_paths lp ON m.learning_path_id = lp.id + WHERE lp.is_published = true + AND ( + SELECT COUNT(*) FROM milestone_progress mp + WHERE mp.milestone_id = m.id AND mp.status = 'completed' + ) >= 5` + ); + + logger.info('[LeaderboardPrecompute] Computing milestone-level snapshots', { + milestoneCount: activeMilestones.length, + }); + + for (const { id: milestoneId } of activeMilestones) { + for (const period of ALL_PERIODS) { + try { + await upsertSnapshot('milestone', milestoneId, period); + snapshotCount++; + } catch (err) { + errorCount++; + logger.error('[LeaderboardPrecompute] Failed milestone snapshot', { + milestoneId, + period, + error: (err as Error).message, + }); + } + } + } + + // ── 4. Invalidate short-lived caches so next request picks up fresh data + await CacheService.invalidatePattern('leaderboard:*'); + + const elapsed = Date.now() - startTime; + logger.info('[LeaderboardPrecompute] Completed', { + elapsedMs: elapsed, + snapshots: snapshotCount, + errors: errorCount, + }); + + return { elapsed, snapshots: snapshotCount, errors: errorCount }; + } catch (err) { + const elapsed = Date.now() - startTime; + logger.error('[LeaderboardPrecompute] Fatal error', { + error: (err as Error).message, + elapsedMs: elapsed, + }); + throw err; + } +} + +/** + * Compute a single leaderboard using the live SQL and upsert into snapshots. + * Uses ON CONFLICT to keep the table at exactly one row per (type, target_id, period). + */ +async function upsertSnapshot( + type: LeaderboardType, + targetId: string | undefined, + period: LeaderboardPeriod +): Promise { + const leaderboard = await CollaborativeLearningService.computeLeaderboardLive( + type, + targetId, + period + ); + + await pool.query( + `INSERT INTO leaderboard_snapshots + (type, target_id, period, entries, computed_at, updated_at) + VALUES ($1, $2, $3, $4::jsonb, NOW(), NOW()) + ON CONFLICT ON CONSTRAINT uq_leaderboard_snapshots_key + DO UPDATE SET + entries = EXCLUDED.entries, + computed_at = EXCLUDED.computed_at, + updated_at = EXCLUDED.updated_at`, + [ + type, + targetId ?? null, + period, + JSON.stringify(leaderboard.entries), + ] + ); +} diff --git a/src/jobs/streakTracking.job.ts b/src/jobs/streakTracking.job.ts new file mode 100644 index 00000000..614119a9 --- /dev/null +++ b/src/jobs/streakTracking.job.ts @@ -0,0 +1,234 @@ +/** + * Streak Tracking Job + * + * Runs daily at 00:05 UTC. For every user who had at least one milestone + * completion or activity event yesterday it: + * 1. Increments current_streak and updates longest_streak in + * user_activity_streaks + * 2. Writes the current streak value into Redis (streak:current:) + * so the leaderboard API can apply real-time increments without waiting + * for the next nightly snapshot + * + * Users with no activity yesterday have their current_streak reset to 0 and + * their Redis key cleared. + */ + +import pool from '../config/database'; +import { redis } from '../config/redis'; +import { logger } from '../utils/logger.utils'; + +/** 30-day TTL for streak Redis keys — they are refreshed every day they are active */ +const STREAK_REDIS_TTL_SECONDS = 30 * 24 * 60 * 60; + +export interface StreakTrackingResult { + elapsed: number; + usersProcessed: number; + streaksIncremented: number; + streaksReset: number; + errors: number; +} + +/** + * Main entry point — called by the scheduler or maintenance worker. + */ +export async function runStreakTracking(): Promise { + const startTime = Date.now(); + let usersProcessed = 0; + let streaksIncremented = 0; + let streaksReset = 0; + let errors = 0; + + logger.info('[StreakTracking] Starting daily streak tracking job'); + + try { + // ── 1. Find all users who had milestone activity yesterday ───────────── + const { rows: activeYesterday } = await pool.query<{ user_id: string }>( + `SELECT DISTINCT pe.student_id AS user_id + FROM milestone_progress mp + JOIN path_enrollments pe ON mp.enrollment_id = pe.id + WHERE mp.completed_at::date = (CURRENT_DATE - INTERVAL '1 day')::date + AND mp.status = 'completed'` + ); + + const activeUserIds = new Set(activeYesterday.map((r) => r.user_id)); + + // ── 2. Load all existing streak records ──────────────────────────────── + const { rows: existingStreaks } = await pool.query<{ + user_id: string; + current_streak: number; + longest_streak: number; + last_active_date: string | null; + }>( + `SELECT user_id, current_streak, longest_streak, last_active_date + FROM user_activity_streaks` + ); + + const streakMap = new Map( + existingStreaks.map((r) => [r.user_id, r]) + ); + + // ── 3. Collect all enrolled users (for reset processing) ─────────────── + const { rows: allEnrolled } = await pool.query<{ student_id: string }>( + `SELECT DISTINCT student_id FROM path_enrollments WHERE status = 'active'` + ); + + // ── 4. Process each enrolled user ────────────────────────────────────── + const pipeline = redis.pipeline(); + const upserts: Array<{ + userId: string; + currentStreak: number; + longestStreak: number; + lastActiveDate: string | null; + }> = []; + + for (const { student_id: userId } of allEnrolled) { + usersProcessed++; + try { + const existing = streakMap.get(userId); + const wasActiveYesterday = activeUserIds.has(userId); + + let currentStreak: number; + let longestStreak: number; + const yesterday = new Date(); + yesterday.setUTCDate(yesterday.getUTCDate() - 1); + const yesterdayStr = yesterday.toISOString().split('T')[0]; + + if (wasActiveYesterday) { + // Increment streak + const prev = existing?.current_streak ?? 0; + const lastActive = existing?.last_active_date; + + // Streak continues only if last_active_date was the day before yesterday + // (already one day gap = streak reset, two consecutive days = continue) + const lastActiveDayBefore = + lastActive === getDateNDaysAgo(2) || prev === 0; + + if (lastActiveDayBefore || !lastActive) { + currentStreak = prev + 1; + } else { + // More than one day gap — restart streak at 1 + currentStreak = 1; + } + + longestStreak = Math.max(currentStreak, existing?.longest_streak ?? 0); + + upserts.push({ + userId, + currentStreak, + longestStreak, + lastActiveDate: yesterdayStr, + }); + + // Write to Redis for real-time leaderboard reads + pipeline.set( + `streak:current:${userId}`, + String(currentStreak), + 'EX', + STREAK_REDIS_TTL_SECONDS + ); + + streaksIncremented++; + } else { + // No activity yesterday — reset streak + currentStreak = 0; + longestStreak = existing?.longest_streak ?? 0; + + upserts.push({ + userId, + currentStreak, + longestStreak, + lastActiveDate: existing?.last_active_date ?? null, + }); + + // Remove Redis key so leaderboard reflects 0 streak + pipeline.del(`streak:current:${userId}`); + + streaksReset++; + } + } catch (err) { + errors++; + logger.error('[StreakTracking] Error processing user streak', { + userId, + error: (err as Error).message, + }); + } + } + + // ── 5. Flush Redis pipeline ──────────────────────────────────────────── + await pipeline.exec(); + + // ── 6. Bulk upsert streak records to DB ─────────────────────────────── + // Process in batches of 500 to avoid huge parameter lists + const BATCH_SIZE = 500; + for (let i = 0; i < upserts.length; i += BATCH_SIZE) { + const batch = upserts.slice(i, i + BATCH_SIZE); + await upsertStreaksBatch(batch); + } + + const elapsed = Date.now() - startTime; + logger.info('[StreakTracking] Completed', { + elapsedMs: elapsed, + usersProcessed, + streaksIncremented, + streaksReset, + errors, + }); + + return { elapsed, usersProcessed, streaksIncremented, streaksReset, errors }; + } catch (err) { + const elapsed = Date.now() - startTime; + logger.error('[StreakTracking] Fatal error', { + error: (err as Error).message, + elapsedMs: elapsed, + }); + throw err; + } +} + +/** + * Bulk upsert a batch of streak records. + * Uses a single multi-row INSERT ... ON CONFLICT for efficiency. + */ +async function upsertStreaksBatch( + batch: Array<{ + userId: string; + currentStreak: number; + longestStreak: number; + lastActiveDate: string | null; + }> +): Promise { + if (batch.length === 0) return; + + // Build parameterised multi-row VALUES clause + const values: any[] = []; + const placeholders = batch.map((row, i) => { + const base = i * 4; + values.push( + row.userId, + row.currentStreak, + row.longestStreak, + row.lastActiveDate + ); + return `($${base + 1}, $${base + 2}, $${base + 3}, $${base + 4})`; + }); + + await pool.query( + `INSERT INTO user_activity_streaks + (user_id, current_streak, longest_streak, last_active_date, updated_at) + VALUES ${placeholders.join(', ')} + ON CONFLICT (user_id) + DO UPDATE SET + current_streak = EXCLUDED.current_streak, + longest_streak = GREATEST(user_activity_streaks.longest_streak, EXCLUDED.longest_streak), + last_active_date = COALESCE(EXCLUDED.last_active_date, user_activity_streaks.last_active_date), + updated_at = NOW()`, + values + ); +} + +/** Return a date string (YYYY-MM-DD) for N days ago in UTC */ +function getDateNDaysAgo(n: number): string { + const d = new Date(); + d.setUTCDate(d.getUTCDate() - n); + return d.toISOString().split('T')[0]; +} diff --git a/src/routes/collaborative-learning.routes.ts b/src/routes/collaborative-learning.routes.ts new file mode 100644 index 00000000..ab2e2f90 --- /dev/null +++ b/src/routes/collaborative-learning.routes.ts @@ -0,0 +1,268 @@ +import { Router, Request, Response } from 'express'; +import { authenticate } from '../middleware/auth.middleware'; +import { createLimiter } from '../middleware/rate-limit.middleware'; +import { CollaborativeLearningService } from '../services/collaborative-learning.service'; +import { logger } from '../utils/logger.utils'; + +const router = Router(); + +/** + * Rate limiter for peer review submission. + * Max 5 peer reviews per user per 24-hour window. + * The service layer enforces the per-path-per-24h limit; this middleware adds a + * per-user global guard to catch coordinated abuse across many paths. + */ +const peerReviewLimiter = createLimiter({ + profile: { + windowMs: 24 * 60 * 60 * 1000, // 24 hours + max: 5, + message: + 'Rate limit exceeded: you may submit at most 5 peer reviews per 24 hours per learning path.', + }, + keyStrategy: 'user', +}); + +// All collaborative learning routes require authentication +router.use(authenticate); + +// ── Discussion Forums ────────────────────────────────────────────────────────── + +/** + * POST /api/v1/collaborative-learning/forums + * Create a discussion forum for a milestone + */ +router.post('/forums', async (req: Request, res: Response) => { + try { + const { milestoneId, title, description } = req.body; + const creatorId = (req as any).user.id; + + if (!milestoneId) { + return res.status(400).json({ status: 'error', message: 'milestoneId is required' }); + } + + const forum = await CollaborativeLearningService.createMilestoneForum( + milestoneId, + creatorId, + { title, description } + ); + + return res.status(201).json({ status: 'success', data: forum }); + } catch (err: any) { + logger.error('POST /forums error', { error: err.message }); + return res.status(err.status ?? 500).json({ status: 'error', message: err.message }); + } +}); + +/** + * POST /api/v1/collaborative-learning/forums/:forumId/messages + * Post a message in a forum + */ +router.post('/forums/:forumId/messages', async (req: Request, res: Response) => { + try { + const { forumId } = req.params; + const userId = (req as any).user.id; + const { content, parentMessageId } = req.body; + + if (!content) { + return res.status(400).json({ status: 'error', message: 'content is required' }); + } + + const message = await CollaborativeLearningService.postForumMessage( + forumId, + userId, + { content, parentMessageId } + ); + + return res.status(201).json({ status: 'success', data: message }); + } catch (err: any) { + logger.error('POST /forums/:forumId/messages error', { error: err.message }); + return res.status(err.status ?? 500).json({ status: 'error', message: err.message }); + } +}); + +/** + * GET /api/v1/collaborative-learning/forums/:forumId/messages + * Get paginated messages for a forum + */ +router.get('/forums/:forumId/messages', async (req: Request, res: Response) => { + try { + const { forumId } = req.params; + const userId = (req as any).user.id; + const page = parseInt(req.query.page as string, 10) || 1; + const limit = Math.min(parseInt(req.query.limit as string, 10) || 20, 100); + const includeReplies = req.query.includeReplies === 'true'; + + const result = await CollaborativeLearningService.getForumMessages( + forumId, + userId, + { page, limit, includeReplies } + ); + + return res.status(200).json({ status: 'success', data: result }); + } catch (err: any) { + logger.error('GET /forums/:forumId/messages error', { error: err.message }); + return res.status(err.status ?? 500).json({ status: 'error', message: err.message }); + } +}); + +// ── Study Groups ─────────────────────────────────────────────────────────────── + +/** + * POST /api/v1/collaborative-learning/study-groups + * Create a study group for a learning path + */ +router.post('/study-groups', async (req: Request, res: Response) => { + try { + const creatorId = (req as any).user.id; + const { learningPathId, name, description, maxMembers, isPublic, meetingSchedule, communicationChannel } = req.body; + + if (!learningPathId || !name || !description) { + return res.status(400).json({ + status: 'error', + message: 'learningPathId, name, and description are required', + }); + } + + const group = await CollaborativeLearningService.createStudyGroup( + learningPathId, + creatorId, + { name, description, maxMembers, isPublic, meetingSchedule, communicationChannel } + ); + + return res.status(201).json({ status: 'success', data: group }); + } catch (err: any) { + logger.error('POST /study-groups error', { error: err.message }); + return res.status(err.status ?? 500).json({ status: 'error', message: err.message }); + } +}); + +/** + * POST /api/v1/collaborative-learning/study-groups/:groupId/join + * Join a study group + */ +router.post('/study-groups/:groupId/join', async (req: Request, res: Response) => { + try { + const { groupId } = req.params; + const userId = (req as any).user.id; + + const member = await CollaborativeLearningService.joinStudyGroup(groupId, userId); + + return res.status(200).json({ status: 'success', data: member }); + } catch (err: any) { + logger.error('POST /study-groups/:groupId/join error', { error: err.message }); + return res.status(err.status ?? 500).json({ status: 'error', message: err.message }); + } +}); + +// ── Peer Reviews ─────────────────────────────────────────────────────────────── + +/** + * POST /api/v1/collaborative-learning/peer-reviews + * Submit a peer review for a milestone submission. + * + * Rate limited: max 5 submissions per user per 24 hours. + * The service additionally enforces max 5 per reviewer per learning path per 24 h. + */ +router.post( + '/peer-reviews', + peerReviewLimiter, + async (req: Request, res: Response) => { + try { + const reviewerId = (req as any).user.id; + const { milestoneId, submissionId, rating, feedback, criteria, isAnonymous } = req.body; + + if (!milestoneId || !submissionId || rating == null || !feedback) { + return res.status(400).json({ + status: 'error', + message: 'milestoneId, submissionId, rating, and feedback are required', + }); + } + + const review = await CollaborativeLearningService.createPeerReview( + milestoneId, + submissionId, + reviewerId, + { rating, feedback, criteria: criteria ?? [], isAnonymous } + ); + + return res.status(201).json({ status: 'success', data: review }); + } catch (err: any) { + logger.error('POST /peer-reviews error', { error: err.message }); + return res.status(err.status ?? 500).json({ status: 'error', message: err.message }); + } + } +); + +/** + * POST /api/v1/collaborative-learning/peer-reviews/:reviewId/vote + * Vote (like) a peer review. + * A review must receive at least 1 vote to count toward the reviewer's + * helpfulReviews score on the leaderboard. + */ +router.post('/peer-reviews/:reviewId/vote', async (req: Request, res: Response) => { + try { + const { reviewId } = req.params; + const voterId = (req as any).user.id; + + const result = await CollaborativeLearningService.votePeerReview(reviewId, voterId); + + return res.status(200).json({ status: 'success', data: result }); + } catch (err: any) { + logger.error('POST /peer-reviews/:reviewId/vote error', { error: err.message }); + return res.status(err.status ?? 500).json({ status: 'error', message: err.message }); + } +}); + +// ── Leaderboard ──────────────────────────────────────────────────────────────── + +/** + * GET /api/v1/collaborative-learning/leaderboard + * Fetch the leaderboard from pre-computed snapshots. + * Responds in < 50ms when a snapshot is available. + * + * Query params: + * type - 'milestone' | 'path' | 'global' (default: 'global') + * id - UUID of the path or milestone (required when type != 'global') + * period - 'week' | 'month' | 'quarter' | 'all' (default: 'month') + */ +router.get('/leaderboard', async (req: Request, res: Response) => { + try { + const type = (req.query.type as any) || 'global'; + const targetId = req.query.id as string | undefined; + const period = (req.query.period as any) || 'month'; + + if (!['milestone', 'path', 'global'].includes(type)) { + return res.status(400).json({ + status: 'error', + message: "type must be 'milestone', 'path', or 'global'", + }); + } + + if (!['week', 'month', 'quarter', 'all'].includes(period)) { + return res.status(400).json({ + status: 'error', + message: "period must be 'week', 'month', 'quarter', or 'all'", + }); + } + + if ((type === 'milestone' || type === 'path') && !targetId) { + return res.status(400).json({ + status: 'error', + message: `id is required when type is '${type}'`, + }); + } + + const leaderboard = await CollaborativeLearningService.getLeaderboard( + type, + targetId, + period + ); + + return res.status(200).json({ status: 'success', data: leaderboard }); + } catch (err: any) { + logger.error('GET /leaderboard error', { error: err.message }); + return res.status(err.status ?? 500).json({ status: 'error', message: err.message }); + } +}); + +export default router; diff --git a/src/routes/index.ts b/src/routes/index.ts index 2b09ca6d..61feec8a 100644 --- a/src/routes/index.ts +++ b/src/routes/index.ts @@ -50,6 +50,7 @@ import { JwksController } from "../controllers/jwks.controller"; import { metricsRegistry } from "../config/metrics"; import { monitoringConfig } from "../config/monitoring.config"; import adminAuditRoutes from "./admin/audit-logs.routes"; +import collaborativeLearningRoutes from "./collaborative-learning.routes"; const router = Router(); @@ -103,6 +104,7 @@ router.use("/subscriptions", subscriptionRoutes); router.use("/tax", taxRoutes); router.use("/oracle", oracleRoutes); router.use("/admin/audit-logs", adminAuditRoutes); +router.use("/collaborative-learning", collaborativeLearningRoutes); router.use("/", exportRoutes); // JWKS public endpoint — no auth required diff --git a/src/services/collaborative-learning.service.ts b/src/services/collaborative-learning.service.ts index 8b39326a..9f07589b 100644 --- a/src/services/collaborative-learning.service.ts +++ b/src/services/collaborative-learning.service.ts @@ -1,5 +1,6 @@ import pool from "../config/database"; import { CacheService } from "./cache.service"; +import { redis } from "../config/redis"; import { CacheKeys, CacheTTL } from "../utils/cache-key.utils"; import { logger } from "../utils/logger.utils"; import { createError } from "../middleware/errorHandler"; @@ -541,7 +542,11 @@ export const CollaborativeLearningService = { }, /** - * Create peer review for milestone submission + * Create peer review for milestone submission. + * + * Anti-gaming: enforces max 5 reviews per reviewer per learning path per + * 24-hour window at the DB level (rate-limiting happens in the route layer + * too, but this service-level guard is the authoritative enforcement point). */ async createPeerReview( milestoneId: string, @@ -556,10 +561,12 @@ export const CollaborativeLearningService = { ): Promise { // Verify milestone and get submission details const { rows: submissionRows } = await pool.query( - `SELECT ms.*, pe.student_id as submitter_id, u.first_name, u.last_name + `SELECT ms.*, pe.student_id as submitter_id, u.first_name, u.last_name, + m.learning_path_id FROM milestone_submissions ms JOIN path_enrollments pe ON ms.enrollment_id = pe.id JOIN users u ON pe.student_id = u.id + JOIN milestones m ON ms.milestone_id = m.id WHERE ms.id = $1 AND ms.milestone_id = $2`, [submissionId, milestoneId] ); @@ -584,8 +591,34 @@ export const CollaborativeLearningService = { throw createError("Reviewer not enrolled in learning path", 403); } + // Prevent self-review + if (submission.submitter_id === reviewerId) { + throw createError("Cannot review your own submission", 403); + } + const reviewer = reviewerRows[0]; + // Anti-gaming: max 5 reviews per reviewer per learning path per 24 hours + const { rows: recentReviews } = await pool.query( + `SELECT COUNT(*) as review_count + FROM peer_reviews pr + JOIN milestone_submissions ms ON pr.submission_id = ms.id + JOIN path_enrollments pe ON ms.enrollment_id = pe.id + JOIN milestones m ON ms.milestone_id = m.id + WHERE pr.reviewer_id = $1 + AND m.learning_path_id = $2 + AND pr.created_at >= NOW() - INTERVAL '24 hours'`, + [reviewerId, submission.learning_path_id] + ); + + const reviewCount = parseInt(recentReviews[0].review_count, 10); + if (reviewCount >= 5) { + throw createError( + "Rate limit exceeded: maximum 5 peer reviews per learning path per 24 hours", + 429 + ); + } + // Check if review already exists const { rows: existingRows } = await pool.query( `SELECT id FROM peer_reviews WHERE submission_id = $1 AND reviewer_id = $2`, @@ -635,7 +668,48 @@ export const CollaborativeLearningService = { }, /** - * Get leaderboard for learning path or milestone + * Vote (like) a peer review. + * A review only counts toward the reviewer's helpfulReviews score once it + * receives at least one liked vote from another user. + */ + async votePeerReview( + reviewId: string, + voterId: string + ): Promise<{ success: boolean }> { + // Verify the review exists and voter is not the reviewer + const { rows: reviewRows } = await pool.query( + `SELECT id, reviewer_id FROM peer_reviews WHERE id = $1`, + [reviewId] + ); + + if (reviewRows.length === 0) { + throw createError("Review not found", 404); + } + + if (reviewRows[0].reviewer_id === voterId) { + throw createError("Cannot vote on your own review", 403); + } + + // Upsert vote — ignore duplicate if already voted + await pool.query( + `INSERT INTO peer_review_votes (review_id, voter_id) + VALUES ($1, $2) + ON CONFLICT (review_id, voter_id) DO NOTHING`, + [reviewId, voterId] + ); + + logger.info("Peer review vote recorded", { reviewId, voterId }); + return { success: true }; + }, + + /** + * Get leaderboard for learning path or milestone. + * + * Serves from leaderboard_snapshots (pre-computed nightly) for fast responses. + * Applies real-time Redis increments to streak_days so leaderboard reflects + * the current day's activity without waiting for the next nightly run. + * + * If no snapshot exists yet, falls back to live query (first-run scenario). */ async getLeaderboard( type: 'milestone' | 'path' | 'global', @@ -643,31 +717,120 @@ export const CollaborativeLearningService = { period: 'week' | 'month' | 'quarter' | 'all' = 'month' ): Promise { const cacheKey = `leaderboard:${type}:${targetId || 'global'}:${period}`; - - // Try cache first + + // ── 1. Try L1 short-lived cache (serves repeated bursts in < 50ms) ────── const cached = await CacheService.get(cacheKey); if (cached !== null) { - return cached; + return this.applyRedisStreakIncrements(cached); + } + + // ── 2. Serve from pre-computed snapshot ────────────────────────────────── + const { rows: snapshotRows } = await pool.query( + `SELECT entries, computed_at + FROM leaderboard_snapshots + WHERE type = $1 + AND COALESCE(target_id::text, '') = $2 + AND period = $3 + ORDER BY computed_at DESC + LIMIT 1`, + [type, targetId ?? '', period] + ); + + if (snapshotRows.length > 0) { + const snapshot = snapshotRows[0]; + const entries: LeaderboardEntry[] = snapshot.entries; + + const leaderboard: Leaderboard = { + type, + period, + entries, + lastUpdated: (snapshot.computed_at as Date).toISOString() + }; + + // Apply live Redis streak increments + const enriched = await this.applyRedisStreakIncrements(leaderboard); + + // Cache for 60 seconds to absorb request bursts + await CacheService.set(cacheKey, enriched, 60); + return enriched; } + // ── 3. Fallback: live query (no snapshot available yet) ────────────────── + logger.warn('No leaderboard snapshot found — computing live', { type, targetId, period }); + const leaderboard = await this.computeLeaderboardLive(type, targetId, period); + + // Cache for 10 minutes (original behaviour) + await CacheService.set(cacheKey, leaderboard, CacheTTL.medium * 2); + return leaderboard; + }, + + /** + * Apply real-time Redis streak increments to leaderboard entries. + * Reads current streak_days overrides from Redis (written by streak cron) + * so the displayed value is always fresh even between nightly snapshots. + */ + async applyRedisStreakIncrements(leaderboard: Leaderboard): Promise { + if (!leaderboard.entries.length) return leaderboard; + + const pipeline = redis.pipeline(); + for (const entry of leaderboard.entries) { + pipeline.get(`streak:current:${entry.userId}`); + } + + const results = await pipeline.exec(); + + const entries = leaderboard.entries.map((entry, i) => { + const redisDays = results?.[i]?.[1]; + return { + ...entry, + streakDays: redisDays != null ? parseInt(redisDays as string, 10) : entry.streakDays + }; + }); + + // Re-sort by score + streak (streak is a secondary signal in display order) + entries.sort((a, b) => b.score - a.score || b.streakDays - a.streakDays); + + // Re-assign ranks after sort + const ranked = entries.map((e, i) => ({ ...e, rank: i + 1 })); + + return { ...leaderboard, entries: ranked }; + }, + + /** + * Live leaderboard computation — used for first-run or admin-triggered + * refresh before the first nightly snapshot exists. + * + * FIX: removed pe.progress_percentage from GROUP BY. + * Uses MAX(pe.progress_percentage) as an aggregate instead so each user + * appears exactly once even if their path progress changed mid-query. + * + * Anti-gaming: helpful_reviews only counts reviews that received at least + * one liked vote (joined via peer_review_votes). + */ + async computeLeaderboardLive( + type: 'milestone' | 'path' | 'global', + targetId?: string, + period: 'week' | 'month' | 'quarter' | 'all' = 'month' + ): Promise { let query = ''; let params: any[] = []; - // Build query based on type switch (type) { case 'milestone': query = ` - SELECT - u.id as user_id, - u.first_name || ' ' || u.last_name as user_name, - COUNT(mp.id) as completed_milestones, - AVG(mp.progress_percentage) as avg_progress, - COUNT(pr.id) as helpful_reviews, - COALESCE(MAX(mp.completed_at) - MIN(mp.started_at), INTERVAL '0') as total_time + SELECT + u.id AS user_id, + u.first_name || ' ' || u.last_name AS user_name, + COUNT(DISTINCT mp.id) AS completed_milestones, + COALESCE(MAX(mp.progress_percentage), 0) AS avg_progress, + COUNT(DISTINCT pr.id) AS helpful_reviews FROM users u - JOIN path_enrollments pe ON u.id = pe.student_id + JOIN path_enrollments pe ON u.id = pe.student_id JOIN milestone_progress mp ON pe.id = mp.enrollment_id - LEFT JOIN peer_reviews pr ON u.id = pr.reviewer_id + LEFT JOIN peer_reviews pr ON u.id = pr.reviewer_id + AND EXISTS ( + SELECT 1 FROM peer_review_votes prv WHERE prv.review_id = pr.id + ) WHERE mp.milestone_id = $1 AND mp.status = 'completed' `; params = [targetId]; @@ -675,35 +838,43 @@ export const CollaborativeLearningService = { case 'path': query = ` - SELECT - u.id as user_id, - u.first_name || ' ' || u.last_name as user_name, - COUNT(mp.id) as completed_milestones, - AVG(mp.progress_percentage) as avg_progress, - COUNT(pr.id) as helpful_reviews, - pe.progress_percentage as path_progress + SELECT + u.id AS user_id, + u.first_name || ' ' || u.last_name AS user_name, + COUNT(DISTINCT mp.id) AS completed_milestones, + COALESCE(MAX(mp.progress_percentage), 0) AS avg_progress, + COUNT(DISTINCT pr.id) AS helpful_reviews, + COALESCE(MAX(pe.progress_percentage), 0) AS path_progress FROM users u - JOIN path_enrollments pe ON u.id = pe.student_id + JOIN path_enrollments pe ON u.id = pe.student_id JOIN milestone_progress mp ON pe.id = mp.enrollment_id - LEFT JOIN peer_reviews pr ON u.id = pr.reviewer_id + LEFT JOIN peer_reviews pr ON u.id = pr.reviewer_id + AND EXISTS ( + SELECT 1 FROM peer_review_votes prv WHERE prv.review_id = pr.id + ) WHERE pe.learning_path_id = $1 `; params = [targetId]; break; case 'global': + default: query = ` - SELECT - u.id as user_id, - u.first_name || ' ' || u.last_name as user_name, - COUNT(mp.id) as completed_milestones, - AVG(mp.progress_percentage) as avg_progress, - COUNT(pr.id) as helpful_reviews, - COUNT(DISTINCT pe.learning_path_id) as paths_enrolled + SELECT + u.id AS user_id, + u.first_name || ' ' || u.last_name AS user_name, + COUNT(DISTINCT mp.id) AS completed_milestones, + COALESCE(AVG(mp.progress_percentage), 0) AS avg_progress, + COUNT(DISTINCT pr.id) AS helpful_reviews, + COUNT(DISTINCT pe.learning_path_id) AS paths_enrolled FROM users u - JOIN path_enrollments pe ON u.id = pe.student_id + JOIN path_enrollments pe ON u.id = pe.student_id JOIN milestone_progress mp ON pe.id = mp.enrollment_id - LEFT JOIN peer_reviews pr ON u.id = pr.reviewer_id + LEFT JOIN peer_reviews pr ON u.id = pr.reviewer_id + AND EXISTS ( + SELECT 1 FROM peer_review_votes prv WHERE prv.review_id = pr.id + ) + WHERE 1=1 `; break; } @@ -715,8 +886,9 @@ export const CollaborativeLearningService = { params.push(timeFilter); } + // GROUP BY: u.id only — no pe.progress_percentage (was the bug) query += ` - GROUP BY u.id, u.first_name, u.last_name, pe.progress_percentage + GROUP BY u.id, u.first_name, u.last_name ORDER BY completed_milestones DESC, avg_progress DESC, helpful_reviews DESC LIMIT 50 `; @@ -728,23 +900,18 @@ export const CollaborativeLearningService = { userId: row.user_id, userName: row.user_name, score: this.calculateLeaderboardScore(row), - achievements: [], // Would be populated from achievements system - streakDays: 0, // Would be calculated from activity data + achievements: [], + streakDays: 0, // populated by applyRedisStreakIncrements completedMilestones: parseInt(row.completed_milestones), helpfulReviews: parseInt(row.helpful_reviews) })); - const leaderboard: Leaderboard = { + return { type, period, entries, lastUpdated: new Date().toISOString() }; - - // Cache for 10 minutes - await CacheService.set(cacheKey, leaderboard, CacheTTL.medium * 2); - - return leaderboard; }, // Private helper methods diff --git a/src/workers/maintenance.worker.ts b/src/workers/maintenance.worker.ts index 8fd415f5..c2a50204 100644 --- a/src/workers/maintenance.worker.ts +++ b/src/workers/maintenance.worker.ts @@ -5,6 +5,8 @@ import { VerificationService } from "../services/verification.service"; import { AuditLogArchivalJob } from "../jobs/auditLog.job"; import keyRotationJob from "../jobs/keyRotation.job"; import recommendationStatsJob from "../jobs/recommendationStats.job"; +import { runLeaderboardPrecompute } from "../jobs/leaderboardPrecompute.job"; +import { runStreakTracking } from "../jobs/streakTracking.job"; import { logger } from "../utils/logger.utils"; async function processMaintenanceJob(job: Job): Promise { @@ -40,6 +42,24 @@ async function processMaintenanceJob(job: Job): Promise { return; } + if (job.name === "leaderboard-precompute-scheduled") { + logger.info("[MaintenanceWorker] Running nightly leaderboard pre-computation", { + jobId: job.id, + }); + const result = await runLeaderboardPrecompute(); + logger.info("[MaintenanceWorker] Leaderboard pre-computation complete", result); + return; + } + + if (job.name === "streak-tracking-scheduled") { + logger.info("[MaintenanceWorker] Running daily streak tracking", { + jobId: job.id, + }); + const result = await runStreakTracking(); + logger.info("[MaintenanceWorker] Streak tracking complete", result); + return; + } + logger.info("[MaintenanceWorker] Running maintenance tasks", { jobId: job.id }); await runMaintenanceTasks(); } diff --git a/src/workers/scheduler.ts b/src/workers/scheduler.ts index 3397ebe3..a94e3299 100644 --- a/src/workers/scheduler.ts +++ b/src/workers/scheduler.ts @@ -15,6 +15,8 @@ import { accountDeletionJob } from "../jobs/accountDeletion.job"; import databaseMaintenanceJob from "../jobs/database-maintenance.job"; import staleDataCleanupJob from "../jobs/stale-data-cleanup.job"; import deprecationMaintenanceJob from "../jobs/deprecation-maintenance.job"; +import { runLeaderboardPrecompute } from "../jobs/leaderboardPrecompute.job"; +import { runStreakTracking } from "../jobs/streakTracking.job"; import { logger } from "../utils/logger.utils"; import config from "../config"; import { AuditLogModel } from "../models/audit-log.model"; @@ -206,8 +208,32 @@ export async function startScheduler(): Promise { }, ); + // Nightly leaderboard pre-computation — daily at 02:30 UTC + // Writes fresh leaderboard_snapshots so the API can respond in < 50 ms. + await addRepeatableJobIfNotExists( + maintenanceQueue, + "leaderboard-precompute-scheduled", + { jobType: "leaderboard-precompute" }, + { + repeat: { pattern: "30 2 * * *" }, // cron: daily 02:30 UTC + jobId: "leaderboard-precompute-recurring", + }, + ); + + // Daily streak tracking — daily at 00:05 UTC + // Increments/resets user_activity_streaks and writes Redis streak keys. + await addRepeatableJobIfNotExists( + maintenanceQueue, + "streak-tracking-scheduled", + { jobType: "streak-tracking" }, + { + repeat: { pattern: "5 0 * * *" }, // cron: daily 00:05 UTC + jobId: "streak-tracking-recurring", + }, + ); + logger.info( - "Job scheduler started — weekly earnings, session reminders, escrow check, notification cleanup, daily maintenance, verification retry, audit log archival, key rotation, and insight generation registered", + "Job scheduler started — weekly earnings, session reminders, escrow check, notification cleanup, daily maintenance, verification retry, audit log archival, key rotation, insight generation, leaderboard pre-computation, and streak tracking registered", ); if (!backgroundCheckPollingTimer) {