From 92ceda0b3de1eb00d58d7dc182eaba4beeb8345b Mon Sep 17 00:00:00 2001 From: Harkinkunmi <158626915+TheHalalHunter@users.noreply.github.com> Date: Wed, 30 Sep 2026 16:44:38 +0100 Subject: [PATCH] fix: resolve issues #521 #522 #523 #526 #521 - Separate health check readiness from Stellar dependency checks - /health/live: always 200 (liveness probe) - /health/ready: DB + Redis only, Stellar removed (readiness probe) - /health: full status including Stellar Horizon + Soroban RPC (monitoring) - Add getSorobanServer() accessor to StellarClient - Update README with three-tier health check documentation #522 - Optimise recommendation queries (6 queries -> <=3) - Implement getRecommendedCourses with merged user-context JOIN query - Peer collaborative signal collapsed to single subquery JOIN, cached 24h - Candidate courses + enrollment counts in one query via correlated subquery - Add migration 0006 with 3 missing indexes for recommendation query paths - Add GET /api/courses/recommended route with scoring and reason tags #523 - Make audit logging non-blocking with buffered writes - Replace fire-and-forget db.insert with AuditLogger class - In-memory buffer flushed periodically (AUDIT_FLUSH_INTERVAL_MS) or on capacity (AUDIT_BUFFER_SIZE); flushing guard prevents concurrent races - Timer is unref'd so it does not block clean process exit - stopAuditLogger() awaits final flush before DB connection closes - Add auth.login and auth.login_failed audit events to auth.service.ts - Add AUDIT_BUFFER_SIZE and AUDIT_FLUSH_INTERVAL_MS to config #526 - Fix TOCTOU race condition in deductCredits (negative balances) - Build admin module: adminGuard, types, service, controller, routes - deductCredits uses single atomic UPDATE...WHERE credits >= amount RETURNING eliminating the read-check-write gap that allowed negative balances - On zero rows: EXISTS check distinguishes 404 (no user) vs 422 (low balance) - grantCredits uses same single-statement pattern - adminGuard uses crypto.timingSafeEqual to prevent timing oracle attacks - Add ADMIN_API_KEY to config (min 32 chars, placeholder-rejected) - Both operations emit structured audit events and invalidate user cache --- README.md | 32 ++- src/audit/index.ts | 154 ++++++++++- src/config/index.ts | 19 ++ .../0006_add_recommendation_indexes.sql | 26 ++ src/middleware/auth.ts | 34 +++ src/modules/admin/admin-users.controller.ts | 45 ++++ src/modules/admin/admin-users.routes.ts | 40 +++ src/modules/admin/admin-users.service.ts | 142 ++++++++++ src/modules/admin/admin-users.types.ts | 27 ++ src/modules/auth/auth.service.ts | 19 ++ src/modules/courses/course.controller.ts | 25 +- src/modules/courses/course.routes.ts | 24 +- src/modules/courses/course.service.ts | 255 +++++++++++++++++- src/modules/courses/course.types.ts | 30 +++ src/routes/v1/index.ts | 2 + src/server.ts | 94 ++++--- src/stellar/client.ts | 5 + 17 files changed, 933 insertions(+), 40 deletions(-) create mode 100644 src/database/migrations/0006_add_recommendation_indexes.sql create mode 100644 src/modules/admin/admin-users.controller.ts create mode 100644 src/modules/admin/admin-users.routes.ts create mode 100644 src/modules/admin/admin-users.service.ts create mode 100644 src/modules/admin/admin-users.types.ts diff --git a/README.md b/README.md index 3c3a4c8..b76a778 100644 --- a/README.md +++ b/README.md @@ -127,9 +127,39 @@ The API will be available at `http://localhost:3000`. ### Health +The API exposes three health check endpoints with distinct responsibilities: + | Method | Path | Description | |--------|------|-------------| -| `GET` | `/health` | Health check | +| `GET` | `/health/live` | **Liveness** — always returns `200 { status: "ok" }`. Wire to Kubernetes `livenessProbe` or any "is the process alive?" check. A failing liveness probe triggers a container restart. | +| `GET` | `/health/ready` | **Readiness** — returns `200` when PostgreSQL and Redis are reachable, `503` otherwise. Wire to Kubernetes `readinessProbe` and load balancer health gates. Stellar is deliberately excluded: a blockchain outage must not remove healthy API instances from rotation, since course browsing, quiz taking, and other non-Stellar features continue to work. | +| `GET` | `/health` | **Full health** — checks all four dependencies (DB, Redis, Stellar Horizon, Soroban RPC). Returns `200 healthy` or `503 degraded`. Intended for monitoring dashboards and alerting only — **do not** wire this to probes that restart containers or pull instances from the load balancer. | + +**Example readiness response (healthy):** +```json +{ + "status": "ready", + "checks": { + "database": "ok", + "redis": "ok" + } +} +``` + +**Example full health response (Soroban degraded):** +```json +{ + "status": "degraded", + "timestamp": "2026-09-30T12:00:00.000Z", + "uptime": 3600, + "checks": { + "database": "ok", + "redis": "ok", + "stellar_horizon": "ok", + "stellar_soroban": "error" + } +} +``` ## Database Schema diff --git a/src/audit/index.ts b/src/audit/index.ts index ca2b7b4..304b100 100644 --- a/src/audit/index.ts +++ b/src/audit/index.ts @@ -1,17 +1,22 @@ import { logger } from "../utils/logger.js"; import { db } from "../config/database.js"; import { auditLogs } from "../database/schema.js"; +import { config } from "../config/index.js"; -type AuditEvent = +// ─── Types ─────────────────────────────────────────────────────────────────── + +export type AuditEvent = | "quiz.submitted" | "reward.claimed" | "reward.queued" | "reward.pending_confirmation" | "credential.minted" | "auth.login" - | "auth.login_failed"; + | "auth.login_failed" + | "admin.credits.deducted" + | "admin.credits.granted"; -interface AuditFields { +export interface AuditFields { userId?: string; submissionId?: string; credentialId?: string; @@ -24,9 +29,148 @@ interface AuditFields { queued?: boolean; ip?: string; userAgent?: string; + stellarAddress?: string; + previousCredits?: number; + newCredits?: number; + reason?: string; + adminNote?: string; +} + +interface AuditEntry { + event: AuditEvent; + fields: AuditFields; + timestamp: Date; } +// ─── AuditLogger ───────────────────────────────────────────────────────────── + +/** + * Non-blocking, buffered audit logger. + * + * Entries are accumulated in an in-memory buffer. The buffer is flushed to the + * database either: + * a) periodically, every AUDIT_FLUSH_INTERVAL_MS milliseconds, or + * b) immediately when the buffer reaches AUDIT_BUFFER_SIZE entries. + * + * The public `log()` method is synchronous and returns immediately — callers + * are never blocked or slowed by database latency. Flush errors are logged as + * warnings and never propagate to callers. + * + * Lifecycle: + * `start()` — arms the periodic flush timer (called once at server start) + * `stop()` — clears the timer and awaits a final flush of buffered entries + * (called during graceful shutdown before the DB connection closes) + */ +class AuditLogger { + private buffer: AuditEntry[] = []; + private timer: ReturnType | null = null; + /** + * Guards against concurrent flush calls racing to drain the same buffer. + * Without this, two simultaneous flushes (e.g. timer fires while a + * capacity-triggered flush is in progress) could each splice the same + * slice and attempt to insert duplicate rows. + */ + private flushing = false; + + // ── Lifecycle ───────────────────────────────────────────────────────── + + start(): void { + if (this.timer !== null) return; // idempotent + this.timer = setInterval( + () => void this.flush(), + config.AUDIT_FLUSH_INTERVAL_MS, + ); + // setInterval keeps the event loop alive by design; unref() lets Node.js + // exit normally even if the timer is still armed — the explicit stop() + // call in the shutdown path will drain buffered entries first. + this.timer.unref(); + logger.info( + { + flushIntervalMs: config.AUDIT_FLUSH_INTERVAL_MS, + bufferSize: config.AUDIT_BUFFER_SIZE, + }, + "Audit logger started", + ); + } + + async stop(): Promise { + if (this.timer !== null) { + clearInterval(this.timer); + this.timer = null; + } + await this.flush(); + logger.info("Audit logger stopped"); + } + + // ── Public API ──────────────────────────────────────────────────────── + + /** + * Enqueue an audit entry. Synchronous and non-blocking — returns immediately. + * Also emits a structured log line so the entry is visible in log streams + * even before it is persisted. + */ + log(event: AuditEvent, fields: AuditFields): void { + logger.info({ audit: true, event, ...fields }, `audit: ${event}`); + + this.buffer.push({ event, fields, timestamp: new Date() }); + + if (this.buffer.length >= config.AUDIT_BUFFER_SIZE) { + // Trigger an early flush without awaiting — callers must not block. + void this.flush(); + } + } + + // ── Internal ────────────────────────────────────────────────────────── + + /** + * Drain the buffer and write all pending entries to the database in a + * single INSERT. Safe to call concurrently — re-entrant calls are dropped + * while a flush is already in progress. + */ + async flush(): Promise { + if (this.flushing || this.buffer.length === 0) return; + + this.flushing = true; + // Splice the entire buffer atomically so new entries queued during an + // in-progress flush land in the next cycle rather than being lost. + const entries = this.buffer.splice(0); + + try { + await db.insert(auditLogs).values( + entries.map((e) => ({ event: e.event, fields: e.fields })), + ); + } catch (err) { + logger.warn( + { err, count: entries.length }, + "Audit log flush failed — entries dropped", + ); + // Entries are intentionally not re-queued: re-queuing on failure risks + // unbounded buffer growth if the database is persistently unavailable, + // and audit log loss is preferable to OOM or cascading failures in the + // main application path. The structured log line emitted by log() above + // provides a secondary record in the log stream. + } finally { + this.flushing = false; + } + } +} + +// ─── Singleton & exports ────────────────────────────────────────────────────── + +export const auditLogger = new AuditLogger(); + +/** + * Convenience wrapper used by all call sites in the codebase. + * Delegates to the singleton logger — synchronous and non-blocking. + */ export function auditLog(event: AuditEvent, fields: AuditFields): void { - logger.info({ audit: true, event, ...fields }, `audit: ${event}`); - db.insert(auditLogs).values({ event, fields }).catch((err) => logger.error({ err }, "Failed to persist audit log")); + auditLogger.log(event, fields); +} + +export function startAuditLogger(): void { + auditLogger.start(); +} + +export async function stopAuditLogger(): Promise { + await auditLogger.stop(); } diff --git a/src/config/index.ts b/src/config/index.ts index 261f2c6..170f9f2 100644 --- a/src/config/index.ts +++ b/src/config/index.ts @@ -37,6 +37,24 @@ const envSchema = z.object({ RATE_LIMIT_MAX: z.coerce.number().default(100), RATE_LIMIT_WINDOW_MS: z.coerce.number().default(60_000), + // Admin API — a static bearer token that must be supplied on every request + // to /api/admin/*. Kept separate from JWT_SECRET so admin credentials can + // be rotated independently of user-facing auth. + ADMIN_API_KEY: z + .string() + .min(32, "ADMIN_API_KEY must be at least 32 characters") + .refine( + (val) => val !== "admin-secret" && !val.includes("change-in-production"), + "ADMIN_API_KEY must be a real secret, not a placeholder" + ), + + // Audit log buffer — entries are accumulated in memory and flushed to the + // database periodically or when the buffer reaches capacity. + // AUDIT_BUFFER_SIZE: max entries before an early flush is triggered. + // AUDIT_FLUSH_INTERVAL_MS: how often the background timer fires (ms). + AUDIT_BUFFER_SIZE: z.coerce.number().int().min(1).default(100), + AUDIT_FLUSH_INTERVAL_MS: z.coerce.number().int().min(100).default(5_000), + // AI service (chainlearn-ai) used for quiz generation AI_SERVICE_URL: z.string().url().default("http://localhost:8000"), AI_TIMEOUT_MS: z.coerce.number().default(30_000), @@ -68,6 +86,7 @@ function loadConfig(): Env { STELLAR_QUIZ_CONTRACT_ID: process.env.STELLAR_QUIZ_CONTRACT_ID || "test", STELLAR_REWARD_CONTRACT_ID: process.env.STELLAR_REWARD_CONTRACT_ID || "test", STELLAR_CREDENTIAL_CONTRACT_ID: process.env.STELLAR_CREDENTIAL_CONTRACT_ID || "test", + ADMIN_API_KEY: process.env.ADMIN_API_KEY || "test-admin-api-key-that-is-at-least-32-chars", }); } console.error( diff --git a/src/database/migrations/0006_add_recommendation_indexes.sql b/src/database/migrations/0006_add_recommendation_indexes.sql new file mode 100644 index 0000000..791a320 --- /dev/null +++ b/src/database/migrations/0006_add_recommendation_indexes.sql @@ -0,0 +1,26 @@ +-- Indexes to support the getRecommendedCourses query pattern. +-- +-- Query 1 (user context): fetches all of a user's enrollments with +-- completed_at and joins to credentials. The existing unique index +-- idx_enrollments_user_course covers (user_id, course_id) but does not +-- include completed_at, so a partial scan is needed to filter completed rows. +-- This composite index lets the planner satisfy +-- WHERE user_id = ? +-- and cover completed_at without a heap fetch. +CREATE INDEX IF NOT EXISTS idx_enrollments_user_completed + ON enrollments (user_id, completed_at); + +-- Query 2 (peer collaborative filtering): finds peers who share any of the +-- current user's enrolled courses, then aggregates their other enrollments. +-- The join condition is WHERE course_id = ANY(?) which requires an index +-- on course_id alone. The leading-column of the unique index is user_id, +-- so it is not used for course-first lookups on all planner configurations. +CREATE INDEX IF NOT EXISTS idx_enrollments_course_id + ON enrollments (course_id); + +-- Query 3 (candidate courses): every recommendation query filters by +-- is_active = true and optionally by difficulty. A composite covering index +-- lets the planner satisfy both predicates without visiting the table heap +-- for the filter pass. +CREATE INDEX IF NOT EXISTS idx_courses_active_difficulty + ON courses (is_active, difficulty); diff --git a/src/middleware/auth.ts b/src/middleware/auth.ts index f1cd326..8715c1b 100644 --- a/src/middleware/auth.ts +++ b/src/middleware/auth.ts @@ -77,3 +77,37 @@ export interface AuthUser { export interface AuthenticatedRequest extends FastifyRequest { authUser: AuthUser; } + +/** + * Admin guard — verifies the request carries the static ADMIN_API_KEY in the + * Authorization header as `Bearer `. Intentionally separate from the + * user JWT flow so admin credentials can be rotated independently. + * + * Timing-safe comparison via `crypto.timingSafeEqual` prevents timing attacks + * that could be used to brute-force the key character-by-character. + */ +import crypto from "node:crypto"; +import { config } from "../config/index.js"; + +export async function adminGuard( + request: FastifyRequest, + _reply: FastifyReply, +): Promise { + const authHeader = request.headers.authorization ?? ""; + const token = authHeader.startsWith("Bearer ") + ? authHeader.slice(7) + : ""; + + // Always run the comparison even when token is empty to prevent early-exit + // timing differences from leaking whether the key exists. + const expected = Buffer.from(config.ADMIN_API_KEY, "utf8"); + const provided = Buffer.from(token, "utf8"); + + const valid = + provided.length === expected.length && + crypto.timingSafeEqual(provided, expected); + + if (!valid) { + throw new UnauthorizedError("Invalid or missing admin API key"); + } +} diff --git a/src/modules/admin/admin-users.controller.ts b/src/modules/admin/admin-users.controller.ts new file mode 100644 index 0000000..78f9ee6 --- /dev/null +++ b/src/modules/admin/admin-users.controller.ts @@ -0,0 +1,45 @@ +import type { FastifyRequest, FastifyReply } from "fastify"; +import { adminUsersService } from "./admin-users.service.js"; +import type { UserIdParams, CreditAdjustmentBody } from "./admin-users.types.js"; + +export class AdminUsersController { + /** + * POST /api/admin/users/:userId/credits/deduct + * Atomically deduct credits from a user. Returns the before/after balances. + */ + async deductCredits( + request: FastifyRequest<{ + Params: UserIdParams; + Body: CreditAdjustmentBody; + }>, + reply: FastifyReply, + ): Promise { + const { userId } = request.params; + const { amount, reason } = request.body; + + const result = await adminUsersService.deductCredits(userId, amount, reason); + + reply.send({ success: true, data: result }); + } + + /** + * POST /api/admin/users/:userId/credits/grant + * Atomically grant credits to a user. Returns the before/after balances. + */ + async grantCredits( + request: FastifyRequest<{ + Params: UserIdParams; + Body: CreditAdjustmentBody; + }>, + reply: FastifyReply, + ): Promise { + const { userId } = request.params; + const { amount, reason } = request.body; + + const result = await adminUsersService.grantCredits(userId, amount, reason); + + reply.send({ success: true, data: result }); + } +} + +export const adminUsersController = new AdminUsersController(); diff --git a/src/modules/admin/admin-users.routes.ts b/src/modules/admin/admin-users.routes.ts new file mode 100644 index 0000000..a55ff0d --- /dev/null +++ b/src/modules/admin/admin-users.routes.ts @@ -0,0 +1,40 @@ +import type { FastifyInstance, FastifySchema } from "fastify"; +import { adminUsersController } from "./admin-users.controller.js"; +import { adminGuard } from "../../middleware/auth.js"; +import { validate } from "../../middleware/validation.js"; +import { userIdParamsSchema, creditAdjustmentSchema } from "./admin-users.types.js"; +import type { UserIdParams, CreditAdjustmentBody } from "./admin-users.types.js"; + +export async function adminUsersRoutes(app: FastifyInstance): Promise { + // All routes in this plugin are protected by adminGuard — the hook is + // registered once here rather than on each individual route. + app.addHook("preHandler", adminGuard); + + app.post<{ Params: UserIdParams; Body: CreditAdjustmentBody }>( + "/:userId/credits/deduct", + { + preHandler: [ + validate({ params: userIdParamsSchema, body: creditAdjustmentSchema }), + ], + schema: { + description: "Deduct credits from a user (atomic, race-condition-free)", + tags: ["admin"], + } as FastifySchema, + }, + (request, reply) => adminUsersController.deductCredits(request, reply), + ); + + app.post<{ Params: UserIdParams; Body: CreditAdjustmentBody }>( + "/:userId/credits/grant", + { + preHandler: [ + validate({ params: userIdParamsSchema, body: creditAdjustmentSchema }), + ], + schema: { + description: "Grant credits to a user", + tags: ["admin"], + } as FastifySchema, + }, + (request, reply) => adminUsersController.grantCredits(request, reply), + ); +} diff --git a/src/modules/admin/admin-users.service.ts b/src/modules/admin/admin-users.service.ts new file mode 100644 index 0000000..29543df --- /dev/null +++ b/src/modules/admin/admin-users.service.ts @@ -0,0 +1,142 @@ +import { eq, sql } from "drizzle-orm"; +import { db } from "../../config/database.js"; +import { users } from "../../database/schema.js"; +import { NotFoundError, ValidationError } from "../../utils/errors.js"; +import { auditLog } from "../../audit/index.js"; +import { cacheDel, cacheKey } from "../../cache/index.js"; +import type { CreditAdjustmentResult } from "./admin-users.types.js"; + +export class AdminUsersService { + /** + * Atomically deduct credits from a user's balance. + * + * ## Why this is race-condition-free + * + * A naive implementation would: + * 1. SELECT credits FROM users WHERE id = ? -- read + * 2. if (credits < amount) throw ValidationError -- check + * 3. UPDATE users SET credits = credits - amount -- write + * + * Between steps 1 and 3 any concurrent reward claim or admin grant can + * change the balance, so two concurrent deductions of 80 against a balance + * of 100 could both pass the check and produce -60. This is the classic + * TOCTOU (Time-Of-Check-Time-Of-Use) race condition. + * + * This implementation collapses the check and write into a single atomic + * UPDATE statement: + * + * UPDATE users + * SET credits = credits - amount + * WHERE id = ? + * AND credits >= amount ← balance check is inside the same statement + * RETURNING id, credits + * + * PostgreSQL evaluates the WHERE clause and applies the SET in the same row + * lock, so no concurrent transaction can slip a write in between. If the + * balance is insufficient the WHERE clause matches zero rows and `updated` + * is undefined — we can then distinguish "user not found" from "insufficient + * balance" with a second lightweight EXISTS check. + */ + async deductCredits( + userId: string, + amount: number, + reason: string, + adminNote?: string, + ): Promise { + const [updated] = await db + .update(users) + .set({ + credits: sql`${users.credits} - ${amount}`, + updatedAt: new Date(), + }) + .where( + sql`${users.id} = ${userId} + AND ${users.credits} >= ${amount}`, + ) + .returning({ id: users.id, credits: users.credits }); + + if (!updated) { + // Distinguish "user doesn't exist" from "insufficient balance" so the + // caller gets the correct HTTP status (404 vs 422). + const exists = await db.query.users.findFirst({ + where: eq(users.id, userId), + columns: { id: true }, + }); + + if (!exists) { + throw new NotFoundError("User"); + } + + // User exists but the WHERE credits >= amount condition failed. + throw new ValidationError({ + amount: [ + "Insufficient credits: user does not have enough credits for this deduction", + ], + }); + } + + const newCredits = updated.credits; + const previousCredits = newCredits + amount; + + auditLog("admin.credits.deducted", { + userId, + amount, + previousCredits, + newCredits, + reason, + ...(adminNote ? { adminNote } : {}), + }); + + await cacheDel(cacheKey("user", "profile", userId)); + await cacheDel(cacheKey("user", "progress", userId)); + + return { userId, previousCredits, newCredits, delta: -amount }; + } + + /** + * Atomically grant credits to a user's balance. + * + * Uses `credits + amount` unconditionally — there is no upper bound today, + * so the only safety requirement is that the user actually exists. The + * returning clause gives us the new balance in the same round-trip, + * letting us derive previousCredits without a pre-read. + */ + async grantCredits( + userId: string, + amount: number, + reason: string, + adminNote?: string, + ): Promise { + const [updated] = await db + .update(users) + .set({ + credits: sql`${users.credits} + ${amount}`, + updatedAt: new Date(), + }) + .where(eq(users.id, userId)) + .returning({ id: users.id, credits: users.credits }); + + if (!updated) { + throw new NotFoundError("User"); + } + + const newCredits = updated.credits; + const previousCredits = newCredits - amount; + + auditLog("admin.credits.granted", { + userId, + amount, + previousCredits, + newCredits, + reason, + ...(adminNote ? { adminNote } : {}), + }); + + await cacheDel(cacheKey("user", "profile", userId)); + await cacheDel(cacheKey("user", "progress", userId)); + + return { userId, previousCredits, newCredits, delta: amount }; + } +} + +export const adminUsersService = new AdminUsersService(); diff --git a/src/modules/admin/admin-users.types.ts b/src/modules/admin/admin-users.types.ts new file mode 100644 index 0000000..b13e727 --- /dev/null +++ b/src/modules/admin/admin-users.types.ts @@ -0,0 +1,27 @@ +import { z } from "zod"; + +// ─── Request schemas ────────────────────────────────────────────────────────── + +export const userIdParamsSchema = z.object({ + userId: z.string().uuid("Invalid user ID"), +}); + +export const creditAdjustmentSchema = z.object({ + amount: z + .number() + .int("Amount must be an integer") + .positive("Amount must be greater than zero"), + reason: z.string().min(1).max(255), +}); + +// ─── Types ──────────────────────────────────────────────────────────────────── + +export type UserIdParams = z.infer; +export type CreditAdjustmentBody = z.infer; + +export interface CreditAdjustmentResult { + userId: string; + previousCredits: number; + newCredits: number; + delta: number; +} diff --git a/src/modules/auth/auth.service.ts b/src/modules/auth/auth.service.ts index 8e8e218..4536f87 100644 --- a/src/modules/auth/auth.service.ts +++ b/src/modules/auth/auth.service.ts @@ -7,6 +7,7 @@ import { getNetworkPassphrase } from "../../config/stellar.js"; import { UnauthorizedError } from "../../utils/errors.js"; import { logger } from "../../utils/logger.js"; import { eq } from "drizzle-orm"; +import { auditLog } from "../../audit/index.js"; import type { ChallengeResponse, AuthResponse } from "./auth.types.js"; const CHALLENGE_TTL_SECONDS = 300; // 5 minutes @@ -81,6 +82,22 @@ export class AuthService { stellarAddress: string, challengeId: string, signedChallenge: string + ): Promise { + try { + return await this._verifyChallenge(stellarAddress, challengeId, signedChallenge); + } catch (err) { + // Audit every authentication failure in one place so individual + // throw sites don't each need their own auditLog call. Re-throw + // unchanged so the HTTP layer still returns the correct status. + auditLog("auth.login_failed", { stellarAddress }); + throw err; + } + } + + private async _verifyChallenge( + stellarAddress: string, + challengeId: string, + signedChallenge: string ): Promise { // Atomically retrieve and delete the challenge (single-use, must not be consumed yet). // The key includes the per-request challengeId so it cannot be guessed or clobbered. @@ -205,6 +222,8 @@ export class AuthService { logger.info({ stellarAddress, userId: user.id }, "New user created"); } + auditLog("auth.login", { userId: user.id, stellarAddress }); + return { token: "", // Will be set by controller user: { diff --git a/src/modules/courses/course.controller.ts b/src/modules/courses/course.controller.ts index c09505a..173cbb1 100644 --- a/src/modules/courses/course.controller.ts +++ b/src/modules/courses/course.controller.ts @@ -1,7 +1,7 @@ import type { FastifyRequest, FastifyReply } from "fastify"; import { courseService } from "./course.service.js"; import type { AuthenticatedRequest } from "../../middleware/auth.js"; -import type { ListCoursesQuery, CourseIdParams } from "./course.types.js"; +import type { ListCoursesQuery, CourseIdParams, RecommendationsQuery } from "./course.types.js"; export class CourseController { /** @@ -59,6 +59,29 @@ export class CourseController { message: "Enrolled successfully", }); } + + /** + * GET /api/courses/recommended + * Return a personalised ranked list of courses for the authenticated user. + * Requires auth — recommendations are user-specific. + */ + async getRecommendations( + request: FastifyRequest<{ Querystring: RecommendationsQuery }>, + reply: FastifyReply + ): Promise { + const { authUser } = request as AuthenticatedRequest; + const { limit } = request.query; + const result = await courseService.getRecommendedCourses(authUser.id, limit); + + reply.send({ + success: true, + data: result.courses, + meta: { + inferredDifficulty: result.inferredDifficulty, + count: result.courses.length, + }, + }); + } } export const courseController = new CourseController(); diff --git a/src/modules/courses/course.routes.ts b/src/modules/courses/course.routes.ts index 04f4dcb..a4eca66 100644 --- a/src/modules/courses/course.routes.ts +++ b/src/modules/courses/course.routes.ts @@ -2,7 +2,11 @@ import type { FastifyInstance, FastifySchema } from "fastify"; import { courseController } from "./course.controller.js"; import { authGuard, optionalAuth } from "../../middleware/auth.js"; import { validate } from "../../middleware/validation.js"; -import { listCoursesSchema, courseIdParamsSchema } from "./course.types.js"; +import { + listCoursesSchema, + courseIdParamsSchema, + recommendationsQuerySchema, +} from "./course.types.js"; export async function courseRoutes(app: FastifyInstance): Promise { app.get<{ Querystring: import("./course.types.js").ListCoursesQuery }>( @@ -17,6 +21,24 @@ export async function courseRoutes(app: FastifyInstance): Promise { (request, reply) => courseController.list(request, reply) ); + // Must be registered before /:id so the literal segment "recommended" is + // not swallowed by the UUID param pattern. + app.get<{ Querystring: import("./course.types.js").RecommendationsQuery }>( + "/recommended", + { + preHandler: [ + authGuard, + validate({ querystring: recommendationsQuerySchema }), + ], + schema: { + description: + "Get personalised course recommendations for the authenticated user", + tags: ["courses"], + } as FastifySchema, + }, + (request, reply) => courseController.getRecommendations(request, reply) + ); + app.get<{ Params: { id: string } }>( "/:id", { diff --git a/src/modules/courses/course.service.ts b/src/modules/courses/course.service.ts index e75e88b..0407b81 100644 --- a/src/modules/courses/course.service.ts +++ b/src/modules/courses/course.service.ts @@ -1,6 +1,6 @@ -import { eq, and, count, desc, inArray } from "drizzle-orm"; +import { eq, and, count, desc, inArray, notInArray, avg, sql } from "drizzle-orm"; import { db } from "../../config/database.js"; -import { courses, enrollments, quizzes } from "../../database/schema.js"; +import { courses, enrollments, credentials, quizSubmissions, quizzes } from "../../database/schema.js"; import { NotFoundError, ConflictError } from "../../utils/errors.js"; import { withLock } from "../../utils/lock.js"; import { @@ -15,6 +15,8 @@ import type { ListCoursesQuery, CourseSummary, CourseDetail, + RecommendedCourse, + GetRecommendationsResult, } from "./course.types.js"; export class CourseService { @@ -225,6 +227,255 @@ export class CourseService { await cacheDel(cacheKey("user", "progress", userId)); }); } + + /** + * Returns a personalised list of recommended courses for the given user. + * + * ## Query plan (≤ 3 DB round-trips per request) + * + * **Query 1 — user context** (always runs, never cached individually) + * A single query joining enrollments LEFT JOIN credentials LEFT JOIN a + * quiz_submissions aggregate subquery. Returns: + * - every course the user is already enrolled in (→ exclusion list) + * - which of those they completed (completedAt IS NOT NULL) + * - whether they have a credential (credentialCourseId IS NOT NULL) + * - their average quiz score across all submissions + * This merges original queries 1, 2, and 3 into one round-trip. + * + * **Query 2 — peer collaborative signal** (cached 24 h per user) + * Finds other users who share ≥1 enrolled course with the current user + * (capped at 500 peers), then aggregates the other courses those peers + * enrolled in. Expensive for large datasets; the 24-hour cache means the + * full join only re-runs once per day per user. + * This replaces original query 4. + * + * **Query 3 — candidate courses + enrollment counts** (result cached 1 h) + * Fetches active courses the user is NOT enrolled in, with their total + * enrollment counts included via a lateral subquery, in a single pass. + * This merges original queries 5 and 6 into one round-trip. + * + * ## Scoring + * Each candidate is scored as: + * peerCount × 3 (collaborative signal — strongest indicator) + * + difficultyBonus (2 if matches inferred level, 1 if adjacent) + * + log(enrolledCount + 1) (popularity fallback for new users) + * + * Results are sorted descending by score, capped at `limit`. + */ + async getRecommendedCourses( + userId: string, + limit = 10, + ): Promise { + const RESULT_CACHE_TTL = 60 * 60; // 1 hour + const PEER_CACHE_TTL = 60 * 60 * 24; // 24 hours + const MAX_PEERS = 500; + const DIFFICULTY_ORDER = ["beginner", "intermediate", "advanced"] as const; + + const resultCacheKey = cacheKey("courses", "recommended", userId); + const peerCacheKey = cacheKey("courses", "recommended", "peers", userId); + + // ── Full result cache ──────────────────────────────────────────────── + const cached = await cacheGet( + "courses", + resultCacheKey, + ); + if (cached) return cached; + + // ── Query 1: user context ──────────────────────────────────────────── + // Joins enrollments → credentials (LEFT) → per-user avg score subquery. + // One round-trip replaces the original 3 sequential queries. + const userScoreSubquery = db + .select({ + userId: quizSubmissions.userId, + avgScore: avg(quizSubmissions.score).as("avg_score"), + }) + .from(quizSubmissions) + .where(eq(quizSubmissions.userId, userId)) + .groupBy(quizSubmissions.userId) + .as("user_scores"); + + const userEnrollmentRows = await db + .select({ + courseId: enrollments.courseId, + completedAt: enrollments.completedAt, + credentialCourseId: credentials.courseId, + avgScore: userScoreSubquery.avgScore, + }) + .from(enrollments) + .leftJoin( + credentials, + and( + eq(credentials.userId, userId), + eq(credentials.courseId, enrollments.courseId), + ), + ) + .leftJoin(userScoreSubquery, eq(userScoreSubquery.userId, userId)) + .where(eq(enrollments.userId, userId)); + + const enrolledCourseIds = userEnrollmentRows.map((r) => r.courseId); + const completedCourseIds = new Set( + userEnrollmentRows + .filter((r) => r.completedAt !== null) + .map((r) => r.courseId), + ); + + // Infer preferred difficulty from completed courses + const rawAvgScore = userEnrollmentRows[0]?.avgScore ?? null; + const avgScore = rawAvgScore !== null ? parseFloat(String(rawAvgScore)) : null; + + // Map avg score → difficulty bucket + // ≥ 80 → ready for the next level; < 50 → suggest easier; otherwise stay + const completedCount = completedCourseIds.size; + let inferredDifficulty: (typeof DIFFICULTY_ORDER)[number] | null = null; + if (completedCount > 0 && avgScore !== null) { + // Most-common difficulty among completed courses would require another + // query; instead use score as a proxy (sufficient without extra DB call) + if (avgScore >= 80) { + inferredDifficulty = "intermediate"; // default upward step + } else if (avgScore < 50) { + inferredDifficulty = "beginner"; + } else { + inferredDifficulty = "intermediate"; + } + } else if (completedCount === 0) { + inferredDifficulty = "beginner"; + } + + // ── Query 2: peer collaborative signal (24 h cache) ────────────────── + // Returns a map of courseId → number of peers who enrolled in that course. + let peerCourseSignal = await cacheGet>( + "courses", + peerCacheKey, + ); + + if (!peerCourseSignal) { + peerCourseSignal = {}; + + if (enrolledCourseIds.length > 0) { + // Single query: use a subquery to find peer user IDs inline, then + // aggregate their other enrollments — one round-trip instead of two. + const peersSubquery = db + .selectDistinct({ peerId: enrollments.userId }) + .from(enrollments) + .where( + and( + inArray(enrollments.courseId, enrolledCourseIds), + sql`${enrollments.userId} != ${userId}`, + ), + ) + .limit(MAX_PEERS) + .as("peers"); + + const peerEnrollmentCounts = await db + .select({ + courseId: enrollments.courseId, + peerCount: count().as("peer_count"), + }) + .from(enrollments) + .innerJoin(peersSubquery, eq(enrollments.userId, peersSubquery.peerId)) + .where( + enrolledCourseIds.length > 0 + ? notInArray(enrollments.courseId, enrolledCourseIds) + : sql`true`, + ) + .groupBy(enrollments.courseId); + + for (const row of peerEnrollmentCounts) { + peerCourseSignal[row.courseId] = row.peerCount; + } + } + + await cacheSet(peerCacheKey, peerCourseSignal, PEER_CACHE_TTL); + } + + // ── Query 3: candidate courses + enrollment counts ──────────────────── + // Single query: active courses the user isn't in, with enrollment count + // included via a correlated subquery. Merges original queries 5 and 6. + const candidateQuery = db + .select({ + id: courses.id, + title: courses.title, + description: courses.description, + difficulty: courses.difficulty, + isActive: courses.isActive, + // Inline correlated subquery so enrollment counts come back in the + // same round-trip rather than a separate aggregation query. + enrolledCount: sql`( + SELECT count(*)::int + FROM enrollments e2 + WHERE e2.course_id = ${courses.id} + )`.as("enrolled_count"), + }) + .from(courses) + .where( + and( + eq(courses.isActive, true), + enrolledCourseIds.length > 0 + ? notInArray(courses.id, enrolledCourseIds) + : sql`true`, + ), + ); + + const candidates = await candidateQuery; + + // ── Scoring ────────────────────────────────────────────────────────── + const difficultyIndex = inferredDifficulty + ? DIFFICULTY_ORDER.indexOf(inferredDifficulty) + : -1; + + const scored = candidates.map((course) => { + const peerCount = peerCourseSignal![course.id] ?? 0; + const popularity = Math.log(course.enrolledCount + 1); + + // Difficulty affinity bonus + let difficultyBonus = 0; + if (difficultyIndex >= 0) { + const courseIdx = DIFFICULTY_ORDER.indexOf( + course.difficulty as (typeof DIFFICULTY_ORDER)[number], + ); + if (courseIdx === difficultyIndex) { + difficultyBonus = 2; // exact match + } else if (Math.abs(courseIdx - difficultyIndex) === 1) { + difficultyBonus = 1; // adjacent level + } + } + + const score = peerCount * 3 + difficultyBonus + popularity; + + // Determine primary reason for recommendation + let reason: RecommendedCourse["reason"]; + if (peerCount > 0) { + reason = "peer"; + } else if (difficultyBonus > 0) { + reason = "difficulty"; + } else { + reason = "popular"; + } + + return { course, score, reason }; + }); + + scored.sort((a, b) => b.score - a.score); + const top = scored.slice(0, limit); + + const result: GetRecommendationsResult = { + courses: top.map(({ course, score, reason }) => ({ + id: course.id, + title: course.title, + description: course.description, + difficulty: course.difficulty, + isActive: course.isActive, + enrolledCount: course.enrolledCount, + recommendationScore: Math.round(score * 100) / 100, + reason, + })), + inferredDifficulty, + }; + + await cacheSet(resultCacheKey, result, RESULT_CACHE_TTL); + + return result; + } } export const courseService = new CourseService(); diff --git a/src/modules/courses/course.types.ts b/src/modules/courses/course.types.ts index fcc6a83..1fff406 100644 --- a/src/modules/courses/course.types.ts +++ b/src/modules/courses/course.types.ts @@ -12,10 +12,15 @@ export const courseIdParamsSchema = z.object({ id: z.string().uuid("Invalid course ID"), }); +export const recommendationsQuerySchema = z.object({ + limit: z.coerce.number().int().min(1).max(20).default(10), +}); + // ─── Types ────────────────────────────────────────────────────────────────── export type ListCoursesQuery = z.infer; export type CourseIdParams = z.infer; +export type RecommendationsQuery = z.infer; export interface CourseSummary { id: string; @@ -38,3 +43,28 @@ export interface CourseModule { title: string; order: number; } + +// ─── Recommendations ───────────────────────────────────────────────────────── + +/** + * A single recommended course, extending CourseSummary with a score that + * reflects how strongly the course is recommended for this user (higher = better). + * The score is computed from collaborative peer signal and difficulty affinity, + * and is exposed so clients can show relative relevance if desired. + */ +export interface RecommendedCourse extends Omit { + recommendationScore: number; + /** + * Why this course was recommended. Helps the client display a reason badge. + * "peer" — peers who share your courses also enrolled in this one + * "difficulty" — matches your demonstrated difficulty level + * "popular" — highly enrolled course you haven't started yet + */ + reason: "peer" | "difficulty" | "popular"; +} + +export interface GetRecommendationsResult { + courses: RecommendedCourse[]; + /** Average difficulty level inferred from the user's completed courses. */ + inferredDifficulty: string | null; +} diff --git a/src/routes/v1/index.ts b/src/routes/v1/index.ts index e3c9805..f23f1ee 100644 --- a/src/routes/v1/index.ts +++ b/src/routes/v1/index.ts @@ -6,6 +6,7 @@ import { courseRoutes } from "../../modules/courses/course.routes.js"; import { quizRoutes } from "../../modules/quizzes/quiz.routes.js"; import { rewardRoutes } from "../../modules/rewards/reward.routes.js"; import { credentialRoutes } from "../../modules/credentials/credential.routes.js"; +import { adminUsersRoutes } from "../../modules/admin/admin-users.routes.js"; export async function registerV1Routes(app: FastifyInstance) { await app.register(authRoutes, { prefix: "/auth" }); @@ -14,4 +15,5 @@ export async function registerV1Routes(app: FastifyInstance) { await app.register(quizRoutes, { prefix: "/quizzes" }); await app.register(rewardRoutes, { prefix: "/rewards" }); await app.register(credentialRoutes, { prefix: "/credentials" }); + await app.register(adminUsersRoutes, { prefix: "/admin/users" }); } diff --git a/src/server.ts b/src/server.ts index e077b06..1583df3 100644 --- a/src/server.ts +++ b/src/server.ts @@ -27,6 +27,7 @@ import { } from "./jobs/cleanup-idempotency.js"; import { processRewardClaim } from "./modules/rewards/reward.service.js"; import { warmCourseCache } from "./cache/warmer.js"; +import { startAuditLogger, stopAuditLogger } from "./audit/index.js"; // Versioned route modules import { registerVersionedRoutes } from "./routes/versioning.js"; @@ -93,59 +94,90 @@ async function buildApp() { registerErrorHandler(app); // ─── Health Check ─────────────────────────────────────────────────────── - app.get("/health", async (_request, reply) => { - const [dbCheck, redisCheck, stellarCheck] = await Promise.allSettled([ + // + // Three tiers of health checks: + // + // GET /health/live — Liveness: is the process alive? + // Always returns 200. Used by Kubernetes liveness probes + // to decide whether to restart the container. + // + // GET /health/ready — Readiness: can the API serve requests? + // Returns 200 only when DB and Redis are reachable. + // Used by load balancers / Kubernetes readiness probes. + // Deliberately excludes Stellar — an outage of the + // blockchain network must not pull the API out of rotation + // since course browsing, quiz taking, and other non-Stellar + // features remain fully functional. + // + // GET /health — Full health: detailed status for monitoring dashboards. + // Checks DB, Redis, Stellar Horizon, and Soroban RPC. + // Returns 503 when any dependency is degraded so that + // dashboards and alerting tools can surface the issue, + // but this endpoint should NOT be wired up to probes that + // trigger restarts or traffic removal. + + /** Liveness probe — returns 200 as long as the process is running. */ + app.get("/health/live", async () => ({ status: "ok" })); + + /** + * Readiness probe — returns 200 only when core infrastructure (DB + Redis) + * is reachable. Stellar is intentionally excluded so that blockchain outages + * do not cause load balancers to stop routing traffic to healthy instances. + */ + app.get("/health/ready", async (_request, reply) => { + const [dbCheck, redisCheck] = await Promise.allSettled([ db.execute(sql`SELECT 1`), redis.ping(), - stellarClient.getHorizonServer().root(), ]); - const allHealthy = [dbCheck, redisCheck, stellarCheck].every( - (c) => c.status === "fulfilled", - ); + const ready = + dbCheck.status === "fulfilled" && redisCheck.status === "fulfilled"; - const status = allHealthy ? "healthy" : "degraded"; - - return reply.status(allHealthy ? 200 : 503).send({ - status, - timestamp: new Date().toISOString(), - uptime: process.uptime(), + return reply.status(ready ? 200 : 503).send({ + status: ready ? "ready" : "not_ready", checks: { database: dbCheck.status === "fulfilled" ? "ok" : "error", redis: redisCheck.status === "fulfilled" ? "ok" : "error", - stellar: stellarCheck.status === "fulfilled" ? "ok" : "error", }, }); }); - app.get("/metrics", { preHandler: authGuard }, async (_request, reply) => { - reply.header("Content-Type", registry.contentType); - return reply.send(await registry.metrics()); - }); - - app.get("/health/live", async () => ({ status: "ok" })); - - app.get("/health/ready", async (_request, reply) => { - const [dbCheck, redisCheck, stellarCheck] = await Promise.allSettled([ - db.execute(sql`SELECT 1`), - redis.ping(), - stellarClient.getHorizonServer().root(), - ]); - - const allHealthy = [dbCheck, redisCheck, stellarCheck].every( + /** + * Full health — detailed status including external Stellar dependencies. + * Intended for monitoring dashboards and alerting, not for automated probes + * that restart containers or remove instances from rotation. + */ + app.get("/health", async (_request, reply) => { + const [dbCheck, redisCheck, horizonCheck, sorobanCheck] = + await Promise.allSettled([ + db.execute(sql`SELECT 1`), + redis.ping(), + stellarClient.getHorizonServer().root(), + stellarClient.getSorobanServer().getHealth(), + ]); + + const allHealthy = [dbCheck, redisCheck, horizonCheck, sorobanCheck].every( (c) => c.status === "fulfilled", ); return reply.status(allHealthy ? 200 : 503).send({ - status: allHealthy ? "ready" : "not_ready", + status: allHealthy ? "healthy" : "degraded", + timestamp: new Date().toISOString(), + uptime: process.uptime(), checks: { database: dbCheck.status === "fulfilled" ? "ok" : "error", redis: redisCheck.status === "fulfilled" ? "ok" : "error", - stellar: stellarCheck.status === "fulfilled" ? "ok" : "error", + stellar_horizon: horizonCheck.status === "fulfilled" ? "ok" : "error", + stellar_soroban: sorobanCheck.status === "fulfilled" ? "ok" : "error", }, }); }); + app.get("/metrics", { preHandler: authGuard }, async (_request, reply) => { + reply.header("Content-Type", registry.contentType); + return reply.send(await registry.metrics()); + }); + // ─── API Routes ───────────────────────────────────────────────────────── await registerVersionedRoutes(app); @@ -157,6 +189,7 @@ async function start() { startRetryProcessor(processRetryJob); startIdempotencyCleanup(); + startAuditLogger(); let cacheWarmInterval: ReturnType | null = null; @@ -180,6 +213,7 @@ async function start() { clearInterval(cacheWarmInterval); } await app.close(); + await stopAuditLogger(); await closeDatabase(); await closeRedis(); await shutdownTracing(); diff --git a/src/stellar/client.ts b/src/stellar/client.ts index c3d21a0..e6a5332 100644 --- a/src/stellar/client.ts +++ b/src/stellar/client.ts @@ -113,6 +113,11 @@ export class StellarClient { getHorizonServer(): StellarSdk.Horizon.Server { return this.horizon; } + + /** Expose Soroban RPC server for health checks. */ + getSorobanServer(): StellarSdk.rpc.Server { + return this.soroban; + } } export const stellarClient = new StellarClient();