Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 60 additions & 1 deletion __mocks__/better-sqlite3.js
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,42 @@ class Statement {
if (!this._db._events.find((e) => e.tx_hash === txHash)) {
this._db._events.push({ type, ledger, tx_hash: txHash, payload });
}
} else if (sql.startsWith('INSERT INTO INDEXER_STATE') || sql.startsWith('INSERT OR REPLACE INTO INDEXER_STATE')) {
return { changes: 1, lastInsertRowid: 0 };
}

if (sql.startsWith('INSERT INTO INDEXER_STATE') || sql.startsWith('INSERT OR REPLACE INTO INDEXER_STATE')) {
const [key, value] = args;
this._db._state.set(key, value);
return { changes: 1, lastInsertRowid: 0 };
}

if (sql.startsWith('INSERT INTO IDEMPOTENCY_KEYS')) {
const [key, expiresAt, requestHash, method, path, statusCode, responseBody, createdAt] = args;
this._db._idempotencyRows.set(key, {
key,
expires_at: expiresAt,
request_hash: requestHash,
method,
path,
status_code: statusCode,
response_body: responseBody,
created_at: createdAt,
});
return { changes: 1, lastInsertRowid: 0 };
}

if (sql.startsWith('DELETE FROM IDEMPOTENCY_KEYS')) {
const threshold = args[0];
let deleted = 0;
for (const [key, row] of Array.from(this._db._idempotencyRows.entries())) {
if (row.expires_at <= threshold) {
this._db._idempotencyRows.delete(key);
deleted += 1;
}
}
return { changes: deleted, lastInsertRowid: 0 };
}

return { changes: 1, lastInsertRowid: 0 };
}

Expand All @@ -31,6 +63,25 @@ class Statement {
const value = this._db._state.get(key);
return value !== undefined ? { value } : undefined;
}

if (sql.includes('FROM IDEMPOTENCY_KEYS')) {
const [key, now] = args;
const row = this._db._idempotencyRows.get(key);
if (row && row.expires_at > now) {
return {
key: row.key,
expiresAt: row.expires_at,
requestHash: row.request_hash,
method: row.method,
path: row.path,
statusCode: row.status_code,
responseBody: row.response_body,
createdAt: row.created_at,
};
}
return undefined;
}

return undefined;
}

Expand All @@ -42,6 +93,13 @@ class Statement {
}
return [...this._db._events];
}

if (sql.includes('FROM IDEMPOTENCY_KEYS')) {
return Array.from(this._db._idempotencyRows.values()).map((row) => ({
key: row.key,
}));
}

return [];
}
}
Expand All @@ -50,6 +108,7 @@ class Database {
constructor(_path) {
this._events = [];
this._state = new Map();
this._idempotencyRows = new Map();
}

exec(_sql) {
Expand Down
10 changes: 10 additions & 0 deletions db/003_idempotency_keys.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
CREATE TABLE IF NOT EXISTS idempotency_keys (
key TEXT PRIMARY KEY,
expires_at INTEGER NOT NULL,
request_hash TEXT NOT NULL,
method TEXT NOT NULL,
path TEXT NOT NULL,
status_code INTEGER NOT NULL DEFAULT 0,
response_body TEXT NOT NULL DEFAULT '',
created_at INTEGER NOT NULL
);
3 changes: 2 additions & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,8 @@
"build": "tsc",
"start": "node dist/index.js",
"test": "jest --runInBand",
"lint": "eslint 'src/**/*.ts' 'tests/**/*.ts' --ext .ts"
"lint": "eslint 'src/**/*.ts' 'tests/**/*.ts' --ext .ts",
"purge:idempotency-keys": "npm run build && node dist/scripts/purgeIdempotencyKeys.js"
},
"dependencies": {
"@stellar/stellar-sdk": "12.1.0",
Expand Down
13 changes: 12 additions & 1 deletion src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,14 @@ const ConfigSchema = z.object({
dbPath: z.string().default('scout-off.db'),
});

function required(key: string): string {
const value = process.env[key];
if (!value) {
throw new Error(`Missing required environment variable: ${key}`);
}
return value;
}

const config = {
port: parseInt(process.env.PORT ?? '4000', 10),
network: (process.env.NETWORK ?? 'testnet') as 'testnet' | 'mainnet',
Expand Down Expand Up @@ -46,7 +54,6 @@ const config = {
xFrameOptions: process.env.SECURITY_X_FRAME_OPTIONS ?? 'DENY',
referrerPolicy: process.env.SECURITY_REFERRER_POLICY ?? 'no-referrer',
},
logLevel: (process.env.LOG_LEVEL ?? 'info') as 'debug' | 'info' | 'warn' | 'error',
webhook: {
enabled: process.env.WEBHOOK_ENABLED === 'true',
url: process.env.WEBHOOK_URL ?? ''
Expand All @@ -56,6 +63,10 @@ const config = {
windowMs: parseInt(process.env.RATE_LIMIT_WINDOW_MS ?? '60000', 10),
max: parseInt(process.env.RATE_LIMIT_MAX ?? '60', 10),
},
idempotency: {
ttlSeconds: parseInt(process.env.IDEMPOTENCY_TTL_SECONDS ?? '86400', 10),
purgeIntervalMs: parseInt(process.env.IDEMPOTENCY_PURGE_INTERVAL_MS ?? '60000', 10),
},
};

export default config;
4 changes: 4 additions & 0 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@ import adminRoutes from './routes/admin';
import { errorHandler } from './middleware/errorHandler';
import { securityHeaders } from './middleware/securityHeaders';
import { correlationId } from './middleware/correlationId';
import { responseTime } from './middleware/responseTime';
import { idempotencyMiddleware, startIdempotencyPurgeJob } from './middleware/idempotency';
import { indexEvents } from './services/indexer';
import { logger } from './utils/logger';
import { stellarHealth } from './services/stellar';
Expand All @@ -21,6 +23,7 @@ app.use(correlationId);
app.use(securityHeaders);
app.use(responseTime);
app.use(express.json());
app.use(idempotencyMiddleware);

app.get('/health', async (_req, res) => {
const healthStatus: Record<string, 'ok' | 'error' | 'disabled'> = {};
Expand Down Expand Up @@ -118,6 +121,7 @@ app.listen(config.port, () => {

poll();
setInterval(poll, 5_000);
startIdempotencyPurgeJob();
});

export default app;
144 changes: 144 additions & 0 deletions src/middleware/idempotency.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,144 @@
import Database from 'better-sqlite3';
import fs from 'fs';
import path from 'path';
import type { NextFunction, Request, Response } from 'express';
import config from '../config';

export interface IdempotencyRecord {
key: string;
expiresAt: number;
requestHash: string;
method: string;
path: string;
statusCode: number;
responseBody: string;
createdAt: number;
}

interface IdempotencyRequest extends Request {
idempotencyKey?: string;
idempotencyReplay?: boolean;
}

const DEFAULT_TTL_SECONDS = 86_400;
const DEFAULT_PURGE_INTERVAL_MS = 60_000;

let db: InstanceType<typeof Database> | undefined;

function getDatabase(): InstanceType<typeof Database> {
if (!db) {
db = new Database(config.dbPath);
initializeDatabase();
}

return db;
}

function initializeDatabase(): void {
const migrationPath = path.resolve(__dirname, '../../db/003_idempotency_keys.sql');
const migrationSql = fs.readFileSync(migrationPath, 'utf8');

getDatabase().exec(migrationSql);
getDatabase().exec(`
CREATE INDEX IF NOT EXISTS idx_idempotency_keys_expires_at
ON idempotency_keys (expires_at);
`);
}

export function getIdempotencyDatabase(): InstanceType<typeof Database> {
return getDatabase();
}

function getTtlSeconds(): number {
const parsed = Number.parseInt(String(process.env.IDEMPOTENCY_TTL_SECONDS ?? ''), 10);
if (Number.isFinite(parsed) && parsed > 0) {
return parsed;
}
return config.idempotency?.ttlSeconds && config.idempotency.ttlSeconds > 0
? config.idempotency.ttlSeconds
: DEFAULT_TTL_SECONDS;
}

function getPurgeIntervalMs(): number {
const parsed = Number.parseInt(String(process.env.IDEMPOTENCY_PURGE_INTERVAL_MS ?? ''), 10);
if (Number.isFinite(parsed) && parsed > 0) {
return parsed;
}
return config.idempotency?.purgeIntervalMs && config.idempotency.purgeIntervalMs > 0
? config.idempotency.purgeIntervalMs
: DEFAULT_PURGE_INTERVAL_MS;
}

function isMutatingMethod(method: string): boolean {
return ['POST', 'PUT', 'PATCH', 'DELETE'].includes(method.toUpperCase());
}

function buildRequestHash(req: Request): string {
return `${req.method}:${req.originalUrl || req.url}`;
}

export function getIdempotencyRecord(key: string, now: number = Math.floor(Date.now() / 1000)): IdempotencyRecord | undefined {
return getDatabase()
.prepare(
'SELECT key, expires_at AS expiresAt, request_hash AS requestHash, method, path, status_code AS statusCode, response_body AS responseBody, created_at AS createdAt FROM idempotency_keys WHERE key = ? AND expires_at > ?'
)
.get(key, now) as IdempotencyRecord | undefined;
}

export function recordIdempotencyKey(
key: string,
req: Request,
expiresAt: number,
now: number = Math.floor(Date.now() / 1000)
): void {
getDatabase()
.prepare(
`INSERT INTO idempotency_keys (key, expires_at, request_hash, method, path, status_code, response_body, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)`
)
.run(key, expiresAt, buildRequestHash(req), req.method, req.originalUrl || req.url, 0, '', now);
}

export function purgeExpiredIdempotencyKeys(now: number = Math.floor(Date.now() / 1000)): number {
return getDatabase().prepare('DELETE FROM idempotency_keys WHERE expires_at <= ?').run(now).changes;
}

export function startIdempotencyPurgeJob(intervalMs: number = getPurgeIntervalMs()): NodeJS.Timeout | undefined {
if (intervalMs <= 0) {
return undefined;
}

return setInterval(() => {
purgeExpiredIdempotencyKeys();
}, intervalMs);
}

export function idempotencyMiddleware(req: IdempotencyRequest, _res: Response, next: NextFunction): void {
if (!isMutatingMethod(req.method)) {
next();
return;
}

const key = req.get('Idempotency-Key') || req.get('idempotency-key');
if (!key) {
next();
return;
}

const now = Math.floor(Date.now() / 1000);
const record = getIdempotencyRecord(key, now);

if (record) {
req.idempotencyKey = key;
req.idempotencyReplay = true;
next();
return;
}

const expiresAt = now + getTtlSeconds();
recordIdempotencyKey(key, req, expiresAt, now);

req.idempotencyKey = key;
req.idempotencyReplay = false;
next();
}
4 changes: 4 additions & 0 deletions src/scripts/purgeIdempotencyKeys.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
import { purgeExpiredIdempotencyKeys } from '../middleware/idempotency';

const deleted = purgeExpiredIdempotencyKeys();
console.log(`Deleted ${deleted} expired idempotency key rows.`);
39 changes: 39 additions & 0 deletions tests/middleware/idempotency.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
import fs from 'fs';
import os from 'os';
import path from 'path';

describe('idempotency cleanup', () => {
let dbPath: string;

beforeEach(() => {
dbPath = path.join(fs.mkdtempSync(path.join(os.tmpdir(), 'idempotency-')), 'db.sqlite');
process.env.DB_PATH = dbPath;
jest.resetModules();
});

afterEach(() => {
fs.rmSync(path.dirname(dbPath), { recursive: true, force: true });
});

it('removes expired rows while preserving unexpired ones', () => {
const { getIdempotencyDatabase, purgeExpiredIdempotencyKeys } = require('../../src/middleware/idempotency');
const db = getIdempotencyDatabase();

db.prepare(
`INSERT INTO idempotency_keys (key, expires_at, request_hash, method, path, status_code, response_body, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)`
).run('expired-key', 10, 'hash-1', 'POST', '/api/players', 200, '{"ok":true}', 1);

db.prepare(
`INSERT INTO idempotency_keys (key, expires_at, request_hash, method, path, status_code, response_body, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)`
).run('active-key', 1000, 'hash-2', 'POST', '/api/players', 200, '{"ok":true}', 1);

const deleted = purgeExpiredIdempotencyKeys(100);

expect(deleted).toBe(1);
expect(db.prepare('SELECT key FROM idempotency_keys ORDER BY key').all()).toEqual([
{ key: 'active-key' },
]);
});
});