diff --git a/src/database/migrations/0003_add_idempotency_keys.sql b/src/database/migrations/0003_add_idempotency_keys.sql new file mode 100644 index 0000000..dda003e --- /dev/null +++ b/src/database/migrations/0003_add_idempotency_keys.sql @@ -0,0 +1,13 @@ +CREATE TABLE IF NOT EXISTS "idempotency_keys" ( + "key" varchar(64) PRIMARY KEY, + "user_id" uuid NOT NULL REFERENCES "users"("id") ON DELETE CASCADE, + "endpoint" varchar(255) NOT NULL, + "request_hash" varchar(64) NOT NULL, + "response_status" integer, + "response_body" jsonb, + "tx_hash" varchar(64), + "created_at" timestamp with time zone NOT NULL DEFAULT now(), + "expires_at" timestamp with time zone NOT NULL +); + +CREATE INDEX IF NOT EXISTS "idx_idempotency_expires" ON "idempotency_keys" ("expires_at"); diff --git a/src/database/schema.ts b/src/database/schema.ts index a719b2c..019d2f4 100644 --- a/src/database/schema.ts +++ b/src/database/schema.ts @@ -150,3 +150,25 @@ export const credentials = pgTable( ), ] ); + +// ─── Idempotency Keys ───────────────────────────────────────────────────── + +export const idempotencyKeys = pgTable( + "idempotency_keys", + { + key: varchar("key", { length: 64 }).primaryKey(), + userId: uuid("user_id") + .notNull() + .references(() => users.id, { onDelete: "cascade" }), + endpoint: varchar("endpoint", { length: 255 }).notNull(), + requestHash: varchar("request_hash", { length: 64 }).notNull(), + responseStatus: integer("response_status"), + responseBody: jsonb("response_body"), + txHash: varchar("tx_hash", { length: 64 }), + createdAt: timestamp("created_at", { withTimezone: true }) + .notNull() + .defaultNow(), + expiresAt: timestamp("expires_at", { withTimezone: true }).notNull(), + }, + (table) => [index("idx_idempotency_expires").on(table.expiresAt)] +); diff --git a/src/jobs/cleanup-idempotency.ts b/src/jobs/cleanup-idempotency.ts new file mode 100644 index 0000000..548babf --- /dev/null +++ b/src/jobs/cleanup-idempotency.ts @@ -0,0 +1,26 @@ +import { cleanupIdempotencyKeys } from "../middleware/idempotency.js"; +import { logger } from "../utils/logger.js"; + +let cleanupInterval: ReturnType | null = null; + +export function startIdempotencyCleanup(): void { + const CLEANUP_INTERVAL_MS = 60 * 60 * 1000; // 1 hour + + cleanupInterval = setInterval(async () => { + try { + const deleted = await cleanupIdempotencyKeys(); + if (deleted > 0) { + logger.info({ deleted }, "Cleaned up expired idempotency keys"); + } + } catch (err) { + logger.error({ err }, "Failed to clean up idempotency keys"); + } + }, CLEANUP_INTERVAL_MS); +} + +export function stopIdempotencyCleanup(): void { + if (cleanupInterval) { + clearInterval(cleanupInterval); + cleanupInterval = null; + } +} diff --git a/src/middleware/idempotency.ts b/src/middleware/idempotency.ts new file mode 100644 index 0000000..20eb228 --- /dev/null +++ b/src/middleware/idempotency.ts @@ -0,0 +1,66 @@ +import { db } from "../config/database.js"; +import { idempotencyKeys } from "../database/schema.js"; +import { eq, and, lt } from "drizzle-orm"; +import { sha256Hash } from "../utils/crypto.js"; +import { ConflictError } from "../utils/errors.js"; + +export async function checkIdempotency( + key: string, + userId: string, + endpoint: string, + requestBody: unknown +): Promise<{ cached: boolean; response?: { status: number; body: unknown } }> { + const requestHash = sha256Hash(JSON.stringify(requestBody)); + + const existing = await db.query.idempotencyKeys.findFirst({ + where: eq(idempotencyKeys.key, key), + }); + + if (existing) { + if (existing.requestHash === requestHash && existing.responseBody) { + return { + cached: true, + response: { + status: existing.responseStatus ?? 200, + body: existing.responseBody, + }, + }; + } + throw new ConflictError("Idempotency key reused with different request body"); + } + + await db.insert(idempotencyKeys).values({ + key, + userId, + endpoint, + requestHash, + expiresAt: new Date(Date.now() + 24 * 60 * 60 * 1000), + }); + + return { cached: false }; +} + +export async function storeIdempotentResponse( + key: string, + status: number, + body: unknown, + txHash?: string +): Promise { + await db + .update(idempotencyKeys) + .set({ + responseStatus: status, + responseBody: body, + txHash: txHash ?? null, + }) + .where(eq(idempotencyKeys.key, key)); +} + +export async function cleanupIdempotencyKeys(): Promise { + const result = await db + .delete(idempotencyKeys) + .where(lt(idempotencyKeys.expiresAt, new Date())) + .returning({ key: idempotencyKeys.key }); + + return result.length; +} diff --git a/src/modules/credentials/credential.controller.ts b/src/modules/credentials/credential.controller.ts index fcb9823..c19d3c1 100644 --- a/src/modules/credentials/credential.controller.ts +++ b/src/modules/credentials/credential.controller.ts @@ -2,6 +2,10 @@ import type { FastifyRequest, FastifyReply } from "fastify"; import { credentialService } from "./credential.service.js"; import type { AuthenticatedRequest } from "../../middleware/auth.js"; import type { MintCredentialBody } from "./credential.types.js"; +import { + checkIdempotency, + storeIdempotentResponse, +} from "../../middleware/idempotency.js"; export class CredentialController { /** @@ -13,14 +17,50 @@ export class CredentialController { reply: FastifyReply ): Promise { const { authUser } = request as AuthenticatedRequest; - const { courseId, submissionId } = (request as any).validatedBody; - const result = await credentialService.mint( + const { courseId, submissionId, idempotencyKey } = (request as any).validatedBody; + + const { cached, response } = await checkIdempotency( + idempotencyKey, authUser.id, - courseId, - submissionId + "/credentials/mint", + request.body ); - reply.status(201).send({ success: true, data: result }); + if (cached) { + reply.status(response!.status).send(response!.body); + return; + } + + try { + const result = await credentialService.mint( + authUser.id, + courseId, + submissionId + ); + + await storeIdempotentResponse( + idempotencyKey, + 201, + { success: true, data: result }, + result.mintTxHash + ); + + reply.status(201).send({ success: true, data: result }); + } catch (err: unknown) { + const statusCode = + err && typeof err === "object" && "statusCode" in err + ? (err as { statusCode: number }).statusCode + : 500; + const message = + err instanceof Error ? err.message : "Internal server error"; + + await storeIdempotentResponse(idempotencyKey, statusCode, { + success: false, + error: message, + }); + + throw err; + } } /** diff --git a/src/modules/credentials/credential.types.ts b/src/modules/credentials/credential.types.ts index 216d2c5..db5092a 100644 --- a/src/modules/credentials/credential.types.ts +++ b/src/modules/credentials/credential.types.ts @@ -5,6 +5,7 @@ import { z } from "zod"; export const mintCredentialSchema = z.object({ courseId: z.string().uuid("Invalid course ID"), submissionId: z.string().uuid("Invalid submission ID"), + idempotencyKey: z.string().min(16).max(64), }); // ─── Types ────────────────────────────────────────────────────────────────── diff --git a/src/modules/rewards/reward.controller.ts b/src/modules/rewards/reward.controller.ts index dad4208..d3ae599 100644 --- a/src/modules/rewards/reward.controller.ts +++ b/src/modules/rewards/reward.controller.ts @@ -2,6 +2,10 @@ import type { FastifyRequest, FastifyReply } from "fastify"; import { rewardService } from "./reward.service.js"; import type { AuthenticatedRequest } from "../../middleware/auth.js"; import type { ClaimRewardBody } from "./reward.types.js"; +import { + checkIdempotency, + storeIdempotentResponse, +} from "../../middleware/idempotency.js"; export class RewardController { /** @@ -13,10 +17,46 @@ export class RewardController { reply: FastifyReply ): Promise { const { authUser } = request as AuthenticatedRequest; - const { submissionId } = (request as any).validatedBody; - const result = await rewardService.claimReward(authUser.id, submissionId); + const { submissionId, idempotencyKey } = (request as any).validatedBody; - reply.send({ success: true, data: result }); + const { cached, response } = await checkIdempotency( + idempotencyKey, + authUser.id, + "/rewards/claim", + request.body + ); + + if (cached) { + reply.status(response!.status).send(response!.body); + return; + } + + try { + const result = await rewardService.claimReward(authUser.id, submissionId); + + await storeIdempotentResponse( + idempotencyKey, + 200, + { success: true, data: result }, + result.txHash ?? undefined + ); + + reply.send({ success: true, data: result }); + } catch (err: unknown) { + const statusCode = + err && typeof err === "object" && "statusCode" in err + ? (err as { statusCode: number }).statusCode + : 500; + const message = + err instanceof Error ? err.message : "Internal server error"; + + await storeIdempotentResponse(idempotencyKey, statusCode, { + success: false, + error: message, + }); + + throw err; + } } /** diff --git a/src/modules/rewards/reward.types.ts b/src/modules/rewards/reward.types.ts index 5d1a316..a898f44 100644 --- a/src/modules/rewards/reward.types.ts +++ b/src/modules/rewards/reward.types.ts @@ -4,6 +4,7 @@ import { z } from "zod"; export const claimRewardSchema = z.object({ submissionId: z.string().uuid("Invalid submission ID"), + idempotencyKey: z.string().min(16).max(64), }); // ─── Types ────────────────────────────────────────────────────────────────── diff --git a/src/server.ts b/src/server.ts index 7ed3b05..afdcfbd 100644 --- a/src/server.ts +++ b/src/server.ts @@ -20,6 +20,10 @@ import { stopRetryProcessor, type RetryJob, } from "./services/retry-queue.js"; +import { + startIdempotencyCleanup, + stopIdempotencyCleanup, +} from "./jobs/cleanup-idempotency.js"; import { processRewardClaim } from "./modules/rewards/reward.service.js"; import { warmCourseCache } from "./cache/warmer.js"; @@ -161,6 +165,7 @@ async function start() { const app = await buildApp(); startRetryProcessor(processRetryJob); + startIdempotencyCleanup(); try { await warmCourseCache(); @@ -177,6 +182,7 @@ async function start() { const shutdown = async (signal: string) => { logger.info({ signal }, "Received shutdown signal"); stopRetryProcessor(); + stopIdempotencyCleanup(); await app.close(); await closeDatabase(); await closeRedis(); diff --git a/tests/e2e/credentials.test.ts b/tests/e2e/credentials.test.ts index c041223..e144b42 100644 --- a/tests/e2e/credentials.test.ts +++ b/tests/e2e/credentials.test.ts @@ -60,6 +60,7 @@ describe("Credentials API", () => { payload: { courseId: "00000000-0000-0000-0000-000000000001", submissionId: "00000000-0000-0000-0000-000000000002", + idempotencyKey: "test-key-credentials-mint-unauth", }, }); @@ -77,6 +78,7 @@ describe("Credentials API", () => { payload: { courseId: "00000000-0000-0000-0000-000000000001", submissionId: "00000000-0000-0000-0000-000000000001", + idempotencyKey: "test-key-credentials-mint-valid", }, }); @@ -101,6 +103,7 @@ describe("Credentials API", () => { payload: { courseId: "00000000-0000-0000-0000-000000000001", submissionId: "00000000-0000-0000-0000-000000000099", + idempotencyKey: "test-key-credentials-mint-notpassed", }, }); @@ -123,6 +126,7 @@ describe("Credentials API", () => { payload: { courseId: "00000000-0000-0000-0000-000000000001", submissionId: "00000000-0000-0000-0000-000000000001", + idempotencyKey: "test-key-credentials-mint-001", }, }); @@ -135,6 +139,7 @@ describe("Credentials API", () => { payload: { courseId: "00000000-0000-0000-0000-000000000001", submissionId: "00000000-0000-0000-0000-000000000001", + idempotencyKey: "test-key-credentials-mint-002", }, }); @@ -145,5 +150,22 @@ describe("Credentials API", () => { } } }); + + it("should reject request without idempotency key", async () => { + const token = createToken(); + + const response = await app.inject({ + method: "POST", + url: "/api/credentials/mint", + headers: { authorization: `Bearer ${token}` }, + payload: { + courseId: "00000000-0000-0000-0000-000000000001", + submissionId: "00000000-0000-0000-0000-000000000002", + }, + }); + + // Auth may reject (401) or validation may reject missing idempotencyKey (400) + expect([400, 401]).toContain(response.statusCode); + }); }); }); diff --git a/tests/e2e/rewards.test.ts b/tests/e2e/rewards.test.ts index cfdc7b7..421d539 100644 --- a/tests/e2e/rewards.test.ts +++ b/tests/e2e/rewards.test.ts @@ -21,6 +21,7 @@ describe("Rewards API", () => { url: "/api/rewards/claim", payload: { submissionId: "00000000-0000-0000-0000-000000000000", + idempotencyKey: "test-key-rewards-claim-unauth", }, }); @@ -48,6 +49,47 @@ describe("Rewards API", () => { // Auth may reject the token (401) or validation may reject the ID (400) expect([400, 401]).toContain(response.statusCode); }); + + it("should reject request without idempotency key", async () => { + const token = app.jwt.sign({ + sub: "00000000-0000-0000-0000-000000000001", + stellarAddress: + "GALICE0000000000000000000000000000000000000000000000000000000", + }); + + const response = await app.inject({ + method: "POST", + url: "/api/rewards/claim", + headers: { authorization: `Bearer ${token}` }, + payload: { + submissionId: "00000000-0000-0000-0000-000000000000", + }, + }); + + // Auth may reject (401) or validation may reject missing idempotencyKey (400) + expect([400, 401]).toContain(response.statusCode); + }); + + it("should reject idempotency key that is too short", async () => { + const token = app.jwt.sign({ + sub: "00000000-0000-0000-0000-000000000001", + stellarAddress: + "GALICE0000000000000000000000000000000000000000000000000000000", + }); + + const response = await app.inject({ + method: "POST", + url: "/api/rewards/claim", + headers: { authorization: `Bearer ${token}` }, + payload: { + submissionId: "00000000-0000-0000-0000-000000000000", + idempotencyKey: "short", + }, + }); + + // Auth may reject (401) or validation may reject short key (400) + expect([400, 401]).toContain(response.statusCode); + }); }); describe("GET /api/rewards/history", () => { diff --git a/tests/unit/services/idempotency.test.ts b/tests/unit/services/idempotency.test.ts new file mode 100644 index 0000000..4f71d0b --- /dev/null +++ b/tests/unit/services/idempotency.test.ts @@ -0,0 +1,213 @@ +import { describe, it, expect, vi, beforeEach } from "vitest"; + +vi.mock("../../../src/config/database.js", () => { + const mockDb = { + query: { + idempotencyKeys: { + findFirst: vi.fn(), + }, + }, + insert: vi.fn(), + update: vi.fn(), + delete: vi.fn(), + }; + return { db: mockDb }; +}); + +vi.mock("../../../src/utils/logger.js", () => ({ + logger: { info: vi.fn(), error: vi.fn(), warn: vi.fn() }, +})); + +import { db } from "../../../src/config/database.js"; +import { + checkIdempotency, + storeIdempotentResponse, + cleanupIdempotencyKeys, +} from "../../../src/middleware/idempotency.js"; + +const mockDb = vi.mocked(db); + +function makeChainable(result: unknown) { + const chain: Record = {}; + chain.then = (resolve: Function, reject: Function) => + Promise.resolve(result).then(resolve, reject); + chain.where = vi.fn().mockReturnValue(chain); + chain.set = vi.fn().mockReturnValue(chain); + chain.values = vi.fn().mockReturnValue(chain); + chain.returning = vi.fn().mockResolvedValue(result); + return chain; +} + +describe("Idempotency Middleware", () => { + beforeEach(() => { + vi.clearAllMocks(); + }); + + describe("checkIdempotency", () => { + it("should return cached=true when key exists with same request hash and response", async () => { + const existingRecord = { + key: "test-key-1234567890", + requestHash: "abc123", + responseBody: { success: true, data: { id: "1" } }, + responseStatus: 200, + }; + + vi.mocked(mockDb.query.idempotencyKeys.findFirst).mockResolvedValue( + existingRecord as any + ); + + const sha256Hash = (await import("../../../src/utils/crypto.js")).sha256Hash; + const requestBody = { submissionId: "sub-1" }; + const expectedHash = sha256Hash(JSON.stringify(requestBody)); + + existingRecord.requestHash = expectedHash; + + const result = await checkIdempotency( + "test-key-1234567890", + "user-1", + "/rewards/claim", + requestBody + ); + + expect(result.cached).toBe(true); + expect(result.response).toEqual({ + status: 200, + body: { success: true, data: { id: "1" } }, + }); + }); + + it("should return cached=false and insert new key when key does not exist", async () => { + vi.mocked(mockDb.query.idempotencyKeys.findFirst).mockResolvedValue( + undefined + ); + + const insertChain = makeChainable([]); + mockDb.insert.mockReturnValue(insertChain); + + const result = await checkIdempotency( + "new-key-1234567890", + "user-1", + "/rewards/claim", + { submissionId: "sub-1" } + ); + + expect(result.cached).toBe(false); + expect(mockDb.insert).toHaveBeenCalled(); + }); + + it("should throw ConflictError when key is reused with different request body", async () => { + const existingRecord = { + key: "test-key-1234567890", + requestHash: "different-hash", + responseBody: { success: true }, + responseStatus: 200, + }; + + vi.mocked(mockDb.query.idempotencyKeys.findFirst).mockResolvedValue( + existingRecord as any + ); + + await expect( + checkIdempotency( + "test-key-1234567890", + "user-1", + "/rewards/claim", + { submissionId: "sub-2" } + ) + ).rejects.toThrow("Idempotency key reused with different request body"); + }); + + it("should throw ConflictError when key exists with same hash but no response body", async () => { + const sha256Hash = (await import("../../../src/utils/crypto.js")).sha256Hash; + const requestBody = { submissionId: "sub-1" }; + const requestHash = sha256Hash(JSON.stringify(requestBody)); + + const existingRecord = { + key: "test-key-1234567890", + requestHash, + responseBody: null, + responseStatus: null, + }; + + vi.mocked(mockDb.query.idempotencyKeys.findFirst).mockResolvedValue( + existingRecord as any + ); + + await expect( + checkIdempotency( + "test-key-1234567890", + "user-1", + "/rewards/claim", + requestBody + ) + ).rejects.toThrow("Idempotency key reused with different request body"); + }); + }); + + describe("storeIdempotentResponse", () => { + it("should update the idempotency key with response data", async () => { + const updateChain = makeChainable([]); + const whereChain = makeChainable([]); + updateChain.where = vi.fn().mockReturnValue(whereChain); + mockDb.update.mockReturnValue(updateChain); + + await storeIdempotentResponse( + "test-key-1234567890", + 200, + { success: true, data: { id: "1" } }, + "tx-hash-abc" + ); + + expect(mockDb.update).toHaveBeenCalled(); + expect(updateChain.set).toHaveBeenCalledWith({ + responseStatus: 200, + responseBody: { success: true, data: { id: "1" } }, + txHash: "tx-hash-abc", + }); + }); + + it("should set txHash to null when not provided", async () => { + const updateChain = makeChainable([]); + const whereChain = makeChainable([]); + updateChain.where = vi.fn().mockReturnValue(whereChain); + mockDb.update.mockReturnValue(updateChain); + + await storeIdempotentResponse("test-key-1234567890", 400, { + success: false, + error: "Bad request", + }); + + expect(updateChain.set).toHaveBeenCalledWith({ + responseStatus: 400, + responseBody: { success: false, error: "Bad request" }, + txHash: null, + }); + }); + }); + + describe("cleanupIdempotencyKeys", () => { + it("should delete expired keys and return count", async () => { + const deletedRows = [ + { key: "expired-key-1" }, + { key: "expired-key-2" }, + ]; + + const deleteChain = makeChainable(deletedRows); + mockDb.delete.mockReturnValue(deleteChain); + + const count = await cleanupIdempotencyKeys(); + + expect(count).toBe(2); + expect(mockDb.delete).toHaveBeenCalled(); + }); + + it("should return 0 when no expired keys exist", async () => { + const deleteChain = makeChainable([]); + mockDb.delete.mockReturnValue(deleteChain); + + const count = await cleanupIdempotencyKeys(); + + expect(count).toBe(0); + }); + }); +});