From f49bdfb9af1be6346a0a838078e3615f8139f28f Mon Sep 17 00:00:00 2001 From: Damilola Ogunrotimi <98775983+Fury03@users.noreply.github.com> Date: Sat, 29 Aug 2026 16:53:04 +0000 Subject: [PATCH] feat(analytics): add real-time active users tracking Track concurrent active users in Redis with a sliding inactivity window, updated on every authenticated request via a global middleware, and expose the current count through a new analytics endpoint. - activeUsersService: Redis sorted-set tracking with configurable timeout (ACTIVE_USER_TIMEOUT_SECONDS) and graceful degradation when Redis is down - activityTracker middleware: decodes the Bearer token and records activity fire-and-forget, wired globally in app.js - GET /api/analytics/active-users endpoint (authenticated) --- .env.example | 4 + app.js | 9 + .../analytics/activeUsersController.js | 25 +++ src/middlewares/analytics/activityTracker.js | 33 ++++ src/routes/analytics/activeUsersRoutes.js | 13 ++ src/services/analytics/activeUsersService.js | 110 +++++++++++++ test/activeUsers.test.js | 154 ++++++++++++++++++ 7 files changed, 348 insertions(+) create mode 100644 src/controllers/analytics/activeUsersController.js create mode 100644 src/middlewares/analytics/activityTracker.js create mode 100644 src/routes/analytics/activeUsersRoutes.js create mode 100644 src/services/analytics/activeUsersService.js create mode 100644 test/activeUsers.test.js diff --git a/.env.example b/.env.example index 63d767f..985e201 100644 --- a/.env.example +++ b/.env.example @@ -153,6 +153,10 @@ REDIS_PORT=6379 # REDIS_USERNAME=default # REDIS_PASSWORD=your_password +# Real-time active-user tracking (issue #243): how long (seconds) a user may +# stay idle before they stop counting as active (default 300). +# ACTIVE_USER_TIMEOUT_SECONDS=300 + # Jitsi configuration for video calls (optional) # JITSI_MEET_DOMAIN=your_jitsi_domain # JITSI_APP_ID=your_app_id diff --git a/app.js b/app.js index e174272..d0e95cd 100644 --- a/app.js +++ b/app.js @@ -20,6 +20,7 @@ import { } from "./src/middlewares/security.js"; import { sanitizeInput } from "./src/middlewares/validate.js"; import { rtlMiddleware } from "./src/middlewares/rtl.js"; +import { trackActivity } from "./src/middlewares/analytics/activityTracker.js"; import { errorHandler, notFound, @@ -64,6 +65,7 @@ import courseBundleRoutes from "./src/routes/course-bundle.routes.js"; import certificateRoutes from "./src/routes/certificate.routes.js"; import badgeRoutes from "./src/routes/badge.routes.js"; import achievementRoutes from "./src/routes/api/achievements.js"; +import activeUsersRoutes from "./src/routes/analytics/activeUsersRoutes.js"; import { healthCheck, ping } from "./src/controllers/healthController.js"; import databaseHealthRoutes from "./src/routes/health/database.js"; import databaseMetricsRoutes from "./src/routes/metrics/database.js"; @@ -180,6 +182,10 @@ app.use(hppMiddleware); app.use(sanitizeInput); app.use(rtlMiddleware); +// Issue #243 — Real-time active-user tracking. Runs for every request and +// only does work when a valid Bearer token is present (see the middleware). +app.use(trackActivity); + // ====================== // ROUTES // ====================== @@ -269,6 +275,9 @@ app.use(versionMiddleware); app.use("/api/v1", generousLimiter, v1Router); app.use("/api/v2", generousLimiter, v2Router); +// Issue #243 — Real-time platform analytics (active users). +app.use("/api/analytics", generousLimiter, activeUsersRoutes); + // Issue #212 — Hashtag trending endpoints. app.use("/api/hashtags", generousLimiter, hashtagRoutes); diff --git a/src/controllers/analytics/activeUsersController.js b/src/controllers/analytics/activeUsersController.js new file mode 100644 index 0000000..699df24 --- /dev/null +++ b/src/controllers/analytics/activeUsersController.js @@ -0,0 +1,25 @@ +// controllers/analytics/activeUsersController.js +import logger from "../../config/logger.js"; +import activeUsersService from "../../services/analytics/activeUsersService.js"; + +/** + * GET /api/analytics/active-users + * Return the current number of concurrent active users (unique users seen + * within the configured inactivity window). + */ +export const getActiveUsers = async (req, res) => { + try { + const activeUsers = await activeUsersService.getActiveUserCount(); + res.status(200).json({ + success: true, + activeUsers, + timeoutSeconds: activeUsersService.getTimeoutSeconds(), + }); + } catch (error) { + logger.error("Failed to retrieve active user count:", error); + res.status(500).json({ + success: false, + message: "Failed to retrieve active user count", + }); + } +}; diff --git a/src/middlewares/analytics/activityTracker.js b/src/middlewares/analytics/activityTracker.js new file mode 100644 index 0000000..c9faf1b --- /dev/null +++ b/src/middlewares/analytics/activityTracker.js @@ -0,0 +1,33 @@ +// middlewares/analytics/activityTracker.js +// +// Records the authenticated user's activity on every request so real-time +// active-user counts reflect live usage. Mounted globally in app.js; it only +// does work when the request carries a valid Bearer token — the user id is +// decoded from the JWT (no database lookup) and the Redis write is +// fire-and-forget, so tracking can never slow down or break a request. + +import jwt from "jsonwebtoken"; +import logger from "../../config/logger.js"; +import activeUsersService from "../../services/analytics/activeUsersService.js"; + +const JWT_SECRET = process.env.JWT_SECRET || "deenbridge-temp-secret-key-2024"; + +export const trackActivity = (req, res, next) => { + const authorization = req.headers?.authorization || ""; + if (authorization.startsWith("Bearer ")) { + try { + const decoded = jwt.verify(authorization.slice(7), JWT_SECRET); + if (decoded?.userId) { + activeUsersService + .trackActivity({ userId: decoded.userId }) + .catch((err) => logger.warn("Activity tracking skipped:", err.message)); + } + } catch { + // Invalid/expired token — the route's own auth will reject the request; + // there is nothing meaningful to track here. + } + } + next(); +}; + +export default trackActivity; diff --git a/src/routes/analytics/activeUsersRoutes.js b/src/routes/analytics/activeUsersRoutes.js new file mode 100644 index 0000000..367c6e7 --- /dev/null +++ b/src/routes/analytics/activeUsersRoutes.js @@ -0,0 +1,13 @@ +// routes/analytics/activeUsersRoutes.js +// +// Real-time platform usage endpoints. Mounted at /api/analytics in app.js. +import express from "express"; +import { protect } from "../../middlewares/authMiddleware.js"; +import { getActiveUsers } from "../../controllers/analytics/activeUsersController.js"; + +const router = express.Router(); + +// Current concurrent active user count (any authenticated user may read it). +router.get("/active-users", protect, getActiveUsers); + +export default router; diff --git a/src/services/analytics/activeUsersService.js b/src/services/analytics/activeUsersService.js new file mode 100644 index 0000000..dc6ec6b --- /dev/null +++ b/src/services/analytics/activeUsersService.js @@ -0,0 +1,110 @@ +// services/analytics/activeUsersService.js +// +// Real-time active-user tracking backed by a Redis sorted set. Each +// authenticated request bumps the user's "last seen" score; users whose score +// falls outside the (configurable) activity window are pruned, and the +// concurrent active-user count is simply the size of the set. +// +// Degrades gracefully: when Redis is unavailable every method becomes a no-op +// (tracking is skipped, the count is 0) so the platform keeps working without +// the analytics layer. + +import { getRedisClient, isRedisReady } from "../../config/redis.js"; + +const ACTIVE_USERS_KEY = "analytics:active-users"; +const DEFAULT_TIMEOUT_SECONDS = 300; // 5 minutes + +export class ActiveUsersService { + /** + * @param {object} [options] + * @param {import("redis").RedisClientType|null} [options.redis] - Optional + * injected client (used by tests). Defaults to the app's shared client. + * @param {number|null} [options.timeoutSeconds] - Inactivity timeout override. + */ + constructor({ redis = null, timeoutSeconds = null } = {}) { + this.redis = redis; + this.timeoutSeconds = timeoutSeconds; + } + + /** + * How long (seconds) a user may stay idle before they stop counting as + * active. Reads ACTIVE_USER_TIMEOUT_SECONDS unless overridden (e.g. by a + * test or a caller that wants a different window). + */ + getTimeoutSeconds() { + return ( + this.timeoutSeconds || + parseInt(process.env.ACTIVE_USER_TIMEOUT_SECONDS || String(DEFAULT_TIMEOUT_SECONDS), 10) || + DEFAULT_TIMEOUT_SECONDS + ); + } + + /** @returns {import("redis").RedisClientType|null} The Redis client in use. */ + _client() { + return this.redis || getRedisClient(); + } + + /** @returns {boolean} Whether a usable Redis client is available. */ + _isReady() { + return this.redis ? true : isRedisReady(); + } + + /** + * Record activity for a user (idempotent — one entry per user) and prune + * entries that have been idle longer than the timeout. + * + * @param {object} params + * @param {string|number} params.userId - The authenticated user's id. + * @returns {Promise} 1 when tracked, 0 when Redis is unavailable. + */ + async trackActivity({ userId }) { + if (!userId) return 0; + if (!this._isReady()) return 0; + + const client = this._client(); + const now = Date.now(); + const timeoutMs = this.getTimeoutSeconds() * 1000; + + await client.zAdd(ACTIVE_USERS_KEY, [{ score: now, value: String(userId) }]); + await client.zRemRangeByScore(ACTIVE_USERS_KEY, 0, now - timeoutMs); + return 1; + } + + /** + * Current number of concurrent active users (unique users seen within the + * activity window). + * + * @returns {Promise} The count, or 0 when Redis is unavailable. + */ + async getActiveUserCount() { + if (!this._isReady()) return 0; + + const client = this._client(); + const now = Date.now(); + const timeoutMs = this.getTimeoutSeconds() * 1000; + + await client.zRemRangeByScore(ACTIVE_USERS_KEY, 0, now - timeoutMs); + return client.zCard(ACTIVE_USERS_KEY); + } + + /** + * Test/dependency-injection seam: swap in a Redis-compatible client. + * + * @param {import("redis").RedisClientType|null} client - The client to use. + */ + setRedis(client) { + this.redis = client; + } + + /** + * Override the inactivity timeout (used by tests to simulate expiry). + * + * @param {number} seconds - Timeout in seconds. + */ + setTimeoutSeconds(seconds) { + this.timeoutSeconds = seconds; + } +} + +export const activeUsersService = new ActiveUsersService(); +export default activeUsersService; diff --git a/test/activeUsers.test.js b/test/activeUsers.test.js new file mode 100644 index 0000000..87296af --- /dev/null +++ b/test/activeUsers.test.js @@ -0,0 +1,154 @@ +import request from "supertest"; +import mongoose from "mongoose"; +import { MongoMemoryServer } from "mongodb-memory-server"; + +import app from "../app.js"; +import User from "../src/models/User.js"; +import activeUsersService from "../src/services/analytics/activeUsersService.js"; +import { seedUserAndLogin } from "./helpers/testAuth.js"; + +// Minimal Redis-compatible in-memory client covering exactly the commands the +// active-users service uses (zAdd / zRemRangeByScore / zCard). The real Redis +// connection is unavailable in the test environment, so this fake stands in at +// the client seam — everything above it (middleware -> service -> endpoint) is +// exercised for real. +const createFakeRedis = () => { + const store = new Map(); // member -> score + return { + _store: store, + async zAdd(_key, members) { + for (const member of members) store.set(member.value, member.score); + return members.length; + }, + async zRemRangeByScore(_key, min, max) { + let removed = 0; + for (const [member, score] of store) { + if (score >= min && score <= max) { + store.delete(member); + removed += 1; + } + } + return removed; + }, + async zCard() { + return store.size; + }, + }; +}; + +describe("Real-time active users tracking (#243)", () => { + let mongoServer; + let readerToken; + let authorToken; + let fakeRedis; + + beforeAll(async () => { + if (mongoose.connection.readyState !== 0) { + await mongoose.disconnect(); + } + mongoServer = await MongoMemoryServer.create(); + await mongoose.connect(mongoServer.getUri()); + + const reader = await seedUserAndLogin(app, { + name: "Active Reader", + email: "active-reader@example.com", + }); + readerToken = reader.token; + + const author = await seedUserAndLogin(app, { + name: "Active Author", + email: "active-author@example.com", + role: "mentor", + }); + authorToken = author.token; + }); + + afterAll(async () => { + if (mongoose.connection.readyState !== 0) { + await mongoose.disconnect(); + } + if (mongoServer) { + await mongoServer.stop(); + } + }); + + beforeEach(() => { + // NOTE: seeded users are intentionally kept — their login tokens from + // beforeAll must stay valid. Only the Redis state is reset per test. + fakeRedis = createFakeRedis(); + activeUsersService.setRedis(fakeRedis); + activeUsersService.setTimeoutSeconds(300); + }); + + afterEach(() => { + activeUsersService.setRedis(null); + }); + + it("requires authentication to read the active-user count", async () => { + const res = await request(app).get("/api/analytics/active-users"); + + expect(res.status).toBe(401); + }); + + it("counts the requesting user via the activity middleware", async () => { + const res = await request(app) + .get("/api/analytics/active-users") + .set("Authorization", `Bearer ${readerToken}`); + + expect(res.status).toBe(200); + expect(res.body.success).toBe(true); + expect(res.body.timeoutSeconds).toBe(300); + // The middleware tracked this very request before the handler counted. + expect(res.body.activeUsers).toBeGreaterThanOrEqual(1); + expect(fakeRedis._store.size).toBeGreaterThanOrEqual(1); + }); + + it("counts each unique user once and ignores repeated activity", async () => { + await request(app) + .get("/api/analytics/active-users") + .set("Authorization", `Bearer ${readerToken}`); + await request(app) + .get("/api/analytics/active-users") + .set("Authorization", `Bearer ${authorToken}`); + + const res = await request(app) + .get("/api/analytics/active-users") + .set("Authorization", `Bearer ${readerToken}`); + + expect(res.body.activeUsers).toBe(2); + }); + + it("expires users who have been idle longer than the timeout", async () => { + // Seed one fresh entry (via the service) and one stale entry directly. + await activeUsersService.trackActivity({ userId: "fresh-user" }); + const staleScore = Date.now() - activeUsersService.getTimeoutSeconds() * 1000 - 1000; + fakeRedis._store.set("stale-user", staleScore); + + const count = await activeUsersService.getActiveUserCount(); + + expect(count).toBe(1); + expect(fakeRedis._store.has("stale-user")).toBe(false); + }); + + it("respects a shorter timeout via the environment override seam", async () => { + activeUsersService.setTimeoutSeconds(60); + await activeUsersService.trackActivity({ userId: "fresh-user" }); + const staleScore = Date.now() - 61 * 1000; + fakeRedis._store.set("stale-user", staleScore); + + const count = await activeUsersService.getActiveUserCount(); + + expect(count).toBe(1); + }); + + it("returns 0 without error when Redis is unavailable", async () => { + activeUsersService.setRedis(null); + + const res = await request(app) + .get("/api/analytics/active-users") + .set("Authorization", `Bearer ${readerToken}`); + + expect(res.status).toBe(200); + expect(res.body.activeUsers).toBe(0); + }); +});