From ab9ef317e0ef236bb3e440f610456b06c7ad54f5 Mon Sep 17 00:00:00 2001 From: oomokaro1 Date: Sat, 18 Jul 2026 12:49:16 +0100 Subject: [PATCH 1/5] feat: add idempotency_keys table to database schema Add new table for tracking idempotency keys on blockchain transaction endpoints. Includes foreign key to users, request hash for body validation, response caching columns, and expiry index for cleanup. --- .../migrations/0003_add_idempotency_keys.sql | 13 +++++++++++ src/database/schema.ts | 22 +++++++++++++++++++ 2 files changed, 35 insertions(+) create mode 100644 src/database/migrations/0003_add_idempotency_keys.sql 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)] +); From 46f5fb63374c0b876226fa8c5e2cd781908ab1f1 Mon Sep 17 00:00:00 2001 From: oomokaro1 Date: Sat, 18 Jul 2026 12:49:22 +0100 Subject: [PATCH 2/5] feat: implement idempotency middleware for transaction endpoints Add checkIdempotency, storeIdempotentResponse, and cleanupIdempotencyKeys functions. Prevents duplicate on-chain transactions by caching responses and detecting reused keys with different request bodies. --- src/middleware/idempotency.ts | 66 +++++++++++++++++++++++++++++++++++ 1 file changed, 66 insertions(+) create mode 100644 src/middleware/idempotency.ts 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; +} From 7ce003665d53456545a625e0cf8f27758f4f60d2 Mon Sep 17 00:00:00 2001 From: oomokaro1 Date: Sat, 18 Jul 2026 12:49:30 +0100 Subject: [PATCH 3/5] feat: integrate idempotency into reward and credential endpoints Add idempotencyKey field to claim reward and mint credential request schemas. Controllers now check idempotency before processing and store responses after completion, enabling safe retry behavior for clients. --- .../credentials/credential.controller.ts | 50 +++++++++++++++++-- src/modules/credentials/credential.types.ts | 1 + src/modules/rewards/reward.controller.ts | 46 +++++++++++++++-- src/modules/rewards/reward.types.ts | 1 + 4 files changed, 90 insertions(+), 8 deletions(-) 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 ────────────────────────────────────────────────────────────────── From 59f86a6adb25f5fabf72c76683d8d74393fa85fd Mon Sep 17 00:00:00 2001 From: oomokaro1 Date: Sat, 18 Jul 2026 12:49:37 +0100 Subject: [PATCH 4/5] feat: add scheduled cleanup job for expired idempotency keys Cleanup runs hourly to purge idempotency keys older than 24 hours. Integrated into server startup and graceful shutdown lifecycle. --- src/jobs/cleanup-idempotency.ts | 26 ++++++++++++++++++++++++++ src/server.ts | 6 ++++++ 2 files changed, 32 insertions(+) create mode 100644 src/jobs/cleanup-idempotency.ts 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/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(); From 5fec2cab774a4e1fe91c672d3fc4bfe4287f5467 Mon Sep 17 00:00:00 2001 From: oomokaro1 Date: Sat, 18 Jul 2026 12:49:43 +0100 Subject: [PATCH 5/5] test: add unit and E2E tests for idempotency system Unit tests cover checkIdempotency, storeIdempotentResponse, and cleanupIdempotencyKeys with mocked database. E2E tests verify idempotency key validation on reward claim and credential mint endpoints. --- tests/e2e/credentials.test.ts | 22 +++ tests/e2e/rewards.test.ts | 42 +++++ tests/unit/services/idempotency.test.ts | 213 ++++++++++++++++++++++++ 3 files changed, 277 insertions(+) create mode 100644 tests/unit/services/idempotency.test.ts 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); + }); + }); +});