diff --git a/.env.example b/.env.example index 45f2fe7..ea56e2f 100644 --- a/.env.example +++ b/.env.example @@ -1,140 +1,87 @@ -# Server -PORT=4000 - -# Redis -REDIS_HOST=redis -REDIS_PORT=6379 -REDIS_PASSWORD= -REDIS_URL=redis://redis:6379 - -# Database -DATABASE_URL=postgres://smartdrop:smartdrop@postgres:5432/smartdrop -# Runtime environment -# NODE_ENV: string enum (development, test, production). Default: development. -NODE_ENV=development - -# PORT: number. Default: 3000. -PORT=3000 - -# REDIS_URL: URL. Required in production. Development/test default: redis://localhost:6379. -REDIS_URL=redis://localhost:6379 - -# DATABASE_URL: URL. Required in production. Development default: postgres://localhost/smartdrop. Test default: postgres://localhost/smartdrop_test. -DATABASE_URL=postgres://localhost/smartdrop - -# STELLAR_HORIZON_URL: URL. Default: https://horizon.stellar.org. -STELLAR_HORIZON_URL=https://horizon.stellar.org - -# USDC_ISSUER: Stellar public key. Default: Stellar USDC issuer. -USDC_ISSUER=GA5ZSEJYB37JRC5AVCIA5MOP4RHTM335AX2OBFLDTQLNUEHRGPTM6RIA - -# COINGECKO_API_KEY: string. Optional. Default: empty. -COINGECKO_API_KEY= - -# COINMARKETCAP_API_KEY: string. Optional. Default: empty. -COINMARKETCAP_API_KEY= - -# PRICE_CACHE_TTL_SECONDS: number. Default: 60. -PRICE_CACHE_TTL_SECONDS=60 - -# PRICE_REFRESH_INTERVAL_SECONDS: number. Default: 30. -PRICE_REFRESH_INTERVAL_SECONDS=30 - -# PRICE_STALE_THRESHOLD_MINUTES: number. Default: 5. -PRICE_STALE_THRESHOLD_MINUTES=5 - -# PRICE_ANOMALY_THRESHOLD_PCT: number. Default: 20. -PRICE_ANOMALY_THRESHOLD_PCT=20 - -# API key auth -ADMIN_API_KEY= - -# LOG_LEVEL: string enum (debug, info, warn, error). Default: info. -LOG_LEVEL=info - -# Rate limiting (Redis-backed, per IP) -# RATE_LIMIT_WINDOW_MS: global API window in milliseconds. Default: 60000 (1 minute). -RATE_LIMIT_WINDOW_MS=60000 -# RATE_LIMIT_MAX: max requests per IP per global window. Default: 100. -RATE_LIMIT_MAX=100 -# PRICE_RATELIMIT_WINDOW: price endpoint window in seconds. Default: 60. -PRICE_RATELIMIT_WINDOW=60 -# PRICE_RATELIMIT_MAX: max price requests per IP per window. Default: 30. -PRICE_RATELIMIT_MAX=30 - -# CORS -CORS_ALLOWED_ORIGINS=http://localhost:4000,http://localhost:3001 -# Server -PORT=4000 - -# Redis -REDIS_HOST=redis -REDIS_PORT=6379 -REDIS_PASSWORD= -REDIS_URL=redis://redis:6379 - -# Database -DATABASE_URL=postgres://smartdrop:smartdrop@postgres:5432/smartdrop -# Runtime environment -# NODE_ENV: string enum (development, test, production). Default: development. -NODE_ENV=development - -# PORT: number. Default: 3000. -PORT=3000 - -# REDIS_URL: URL. Required in production. Development/test default: redis://localhost:6379. -REDIS_URL=redis://localhost:6379 - -# DATABASE_URL: URL. Required in production. Development default: postgres://localhost/smartdrop. Test default: postgres://localhost/smartdrop_test. -DATABASE_URL=postgres://localhost/smartdrop - -# STELLAR_HORIZON_URL: URL. Default: https://horizon.stellar.org. -STELLAR_HORIZON_URL=https://horizon.stellar.org - -# USDC_ISSUER: Stellar public key. Default: Stellar USDC issuer. -USDC_ISSUER=GA5ZSEJYB37JRC5AVCIA5MOP4RHTM335AX2OBFLDTQLNUEHRGPTM6RIA - -# COINGECKO_API_KEY: string. Optional. Default: empty. -COINGECKO_API_KEY= - -# COINMARKETCAP_API_KEY: string. Optional. Default: empty. -COINMARKETCAP_API_KEY= - -# PRICE_CACHE_TTL_SECONDS: number. Default: 60. -PRICE_CACHE_TTL_SECONDS=60 - -# PRICE_REFRESH_INTERVAL_SECONDS: number. Default: 30. -PRICE_REFRESH_INTERVAL_SECONDS=30 - -# PRICE_STALE_THRESHOLD_MINUTES: number. Default: 5. -PRICE_STALE_THRESHOLD_MINUTES=5 - -# PRICE_ANOMALY_THRESHOLD_PCT: number. Default: 20. -PRICE_ANOMALY_THRESHOLD_PCT=20 - -# AIRDROP_EXPIRY_CHECK_INTERVAL_SECONDS: number. Default: 60. -AIRDROP_EXPIRY_CHECK_INTERVAL_SECONDS=60 - -# AIRDROP_LEDGER_CACHE_TTL_MS: number. Default: 5000. -AIRDROP_LEDGER_CACHE_TTL_MS=5000 - -# AIRDROP_EXPIRY_SCAN_BATCH_SIZE: number. Default: 100. -AIRDROP_EXPIRY_SCAN_BATCH_SIZE=100 - -# API key auth -ADMIN_API_KEY= - -# Airdrop request limits -# Maximum recipient CSV upload size in bytes. Default: 5 MiB. -AIRDROP_CSV_MAX_BYTES=5242880 -# Maximum JSON request body size in bytes. Default: 2 MiB, sufficient for 10,000 recipients. -AIRDROP_JSON_MAX_BYTES=2097152 -# Per-IP limits for airdrop creation and recipient additions. -AIRDROP_RATELIMIT_WINDOW=60 -AIRDROP_RATELIMIT_MAX=10 - -# LOG_LEVEL: string enum (debug, info, warn, error). Default: info. -LOG_LEVEL=info - -# CORS -CORS_ALLOWED_ORIGINS=http://localhost:4000,http://localhost:3001 +# Server +PORT=4000 + +# Redis +REDIS_HOST=redis +REDIS_PORT=6379 +REDIS_PASSWORD= +REDIS_URL=redis://redis:6379 + +# Database +DATABASE_URL=postgres://smartdrop:smartdrop@postgres:5432/smartdrop +# Runtime environment +# NODE_ENV: string enum (development, test, production). Default: development. +NODE_ENV=development + +# PORT: number. Default: 3000. +PORT=3000 + +# REDIS_URL: URL. Required in production. Development/test default: redis://localhost:6379. +REDIS_URL=redis://localhost:6379 + +# DATABASE_URL: URL. Required in production. Development default: postgres://localhost/smartdrop. Test default: postgres://localhost/smartdrop_test. +DATABASE_URL=postgres://localhost/smartdrop + +# STELLAR_HORIZON_URL: URL. Default: https://horizon.stellar.org. +STELLAR_HORIZON_URL=https://horizon.stellar.org + +# USDC_ISSUER: Stellar public key. Default: Stellar USDC issuer. +USDC_ISSUER=GA5ZSEJYB37JRC5AVCIA5MOP4RHTM335AX2OBFLDTQLNUEHRGPTM6RIA + +# COINGECKO_API_KEY: string. Optional. Default: empty. +COINGECKO_API_KEY= + +# COINMARKETCAP_API_KEY: string. Optional. Default: empty. +COINMARKETCAP_API_KEY= + +# PRICE_CACHE_TTL_SECONDS: number. Default: 60. +PRICE_CACHE_TTL_SECONDS=60 + +# PRICE_REFRESH_INTERVAL_SECONDS: number. Default: 30. +PRICE_REFRESH_INTERVAL_SECONDS=30 + +# PRICE_STALE_THRESHOLD_MINUTES: number. Default: 5. +PRICE_STALE_THRESHOLD_MINUTES=5 + +# PRICE_ANOMALY_THRESHOLD_PCT: number. Default: 20. +PRICE_ANOMALY_THRESHOLD_PCT=20 + +# AIRDROP_EXPIRY_CHECK_INTERVAL_SECONDS: number. Default: 60. +AIRDROP_EXPIRY_CHECK_INTERVAL_SECONDS=60 + +# AIRDROP_LEDGER_CACHE_TTL_MS: number. Default: 5000. +AIRDROP_LEDGER_CACHE_TTL_MS=5000 + +# AIRDROP_EXPIRY_SCAN_BATCH_SIZE: number. Default: 100. +AIRDROP_EXPIRY_SCAN_BATCH_SIZE=100 + +# Airdrop request limits +# Maximum recipient CSV upload size in bytes. Default: 5 MiB. +AIRDROP_CSV_MAX_BYTES=5242880 +# Maximum JSON request body size in bytes. Default: 2 MiB, sufficient for 10,000 recipients. +AIRDROP_JSON_MAX_BYTES=2097152 +# Per-IP limits for airdrop creation and recipient additions. +AIRDROP_RATELIMIT_WINDOW=60 +AIRDROP_RATELIMIT_MAX=10 + +# WATCHED_ASSETS: comma-separated CODE or CODE:ISSUER values warmed before startup. +WATCHED_ASSETS=XLM,USDC:GA5ZSEJYB37JRC5AVCIA5MOP4RHTM335X2KGX3IHOJAPP5RE34K4KZVN + +# API key auth +ADMIN_API_KEY= + +# LOG_LEVEL: string enum (debug, info, warn, error). Default: info. +LOG_LEVEL=info + +# Rate limiting (Redis-backed, per IP) +# RATE_LIMIT_WINDOW_MS: global API window in milliseconds. Default: 60000 (1 minute). +RATE_LIMIT_WINDOW_MS=60000 +# RATE_LIMIT_MAX: max requests per IP per global window. Default: 100. +RATE_LIMIT_MAX=100 +# PRICE_RATELIMIT_WINDOW: price endpoint window in seconds. Default: 60. +PRICE_RATELIMIT_WINDOW=60 +# PRICE_RATELIMIT_MAX: max price requests per IP per window. Default: 30. +PRICE_RATELIMIT_MAX=30 + +# CORS +CORS_ALLOWED_ORIGINS=http://localhost:4000,http://localhost:3001 diff --git a/src/config.js b/src/config.js index f4f8002..854f98a 100644 --- a/src/config.js +++ b/src/config.js @@ -9,6 +9,40 @@ const stellarAddress = makeValidator((input) => { return input; }); +function parseWatchedAssets(input) { + if (!input || !input.trim()) return []; + + const seen = new Set(); + return input + .split(',') + .map((entry) => entry.trim()) + .filter(Boolean) + .map((entry) => { + const [code, issuer, extra] = entry.split(':'); + + if (extra !== undefined) { + throw new Error(`invalid asset "${entry}"; expected CODE or CODE:ISSUER`); + } + + if (!/^[A-Z0-9]{1,12}$/.test(code)) { + throw new Error(`invalid asset code "${code}"; expected 1-12 uppercase alphanumeric characters`); + } + + if (issuer !== undefined && !/^G[A-Z0-9]{55}$/.test(issuer)) { + throw new Error(`invalid issuer for "${code}"; expected a Stellar public key`); + } + + const asset = { code, issuer: issuer || null }; + const key = asset.issuer ? `${asset.code}:${asset.issuer}` : asset.code; + if (seen.has(key)) return null; + seen.add(key); + return asset; + }) + .filter(Boolean); +} + +const watchedAssets = makeValidator(parseWatchedAssets); + const positiveInteger = makeValidator((input) => { const value = Number(input); if (!Number.isSafeInteger(value) || value <= 0) { @@ -55,6 +89,7 @@ const env = cleanEnv(rawEnv, { AIRDROP_EXPIRY_CHECK_INTERVAL_SECONDS: num({ default: 60 }), AIRDROP_LEDGER_CACHE_TTL_MS: num({ default: 5000 }), AIRDROP_EXPIRY_SCAN_BATCH_SIZE: num({ default: 100 }), + WATCHED_ASSETS: watchedAssets({ default: '' }), LOG_LEVEL: str({ default: 'info', choices: ['debug', 'info', 'warn', 'error'], @@ -62,6 +97,9 @@ const env = cleanEnv(rawEnv, { }); const usdcIssuer = env.USDC_ISSUER; +const parsedWatchedAssets = Array.isArray(env.WATCHED_ASSETS) + ? env.WATCHED_ASSETS + : parseWatchedAssets(env.WATCHED_ASSETS); module.exports = { nodeEnv: env.NODE_ENV, @@ -121,6 +159,7 @@ module.exports = { max: env.AIRDROP_RATELIMIT_MAX, }, }, + watchedAssets: parsedWatchedAssets, auth: { adminApiKey: env.ADMIN_API_KEY, }, diff --git a/src/index.js b/src/index.js index 37252f7..d98ef6d 100644 --- a/src/index.js +++ b/src/index.js @@ -9,6 +9,7 @@ const priceOracle = require('./services/priceOracle'); const priceRefreshJob = require('./jobs/priceRefresh'); const webhookRetryWorker = require('./jobs/webhookRetryWorker'); const airdropExpiryJob = require('./jobs/airdropExpiry'); +const { warmCache } = require('./startup/cacheWarm'); const buildCorsMiddleware = require('./middleware/cors'); const buildRateLimit = require('./middleware/rateLimit'); const { requestIdMiddleware } = require('./middleware/requestId'); @@ -117,6 +118,18 @@ function shutdown(signal) { } if (require.main === module) { + startServer().catch((err) => { + logger.error('Startup failed', { error: err.message }); + process.exit(1); + }); + + process.on('SIGTERM', shutdown('SIGTERM')); + process.on('SIGINT', shutdown('SIGINT')); +} + +async function startServer() { + await warmCache(config.watchedAssets); + server = app.listen(config.port, () => { logger.info(`SmartDrop backend running on port ${config.port}`); priceWebSocket.attach(server); @@ -124,9 +137,9 @@ if (require.main === module) { webhookRetryWorker.start(); airdropExpiryJob.start(); }); + module.exports.server = server; - process.on('SIGTERM', shutdown('SIGTERM')); - process.on('SIGINT', shutdown('SIGINT')); + return server; } module.exports = app; @@ -136,3 +149,4 @@ module.exports.server = server || { if (callback) callback(); }, }; +module.exports.startServer = startServer; diff --git a/src/middleware/rateLimit.js b/src/middleware/rateLimit.js index 3d11bb4..d9b99aa 100644 --- a/src/middleware/rateLimit.js +++ b/src/middleware/rateLimit.js @@ -1,56 +1,56 @@ -'use strict'; - -const cache = require('../services/cache'); -const logger = require('../logger'); -const AppError = require('../errors/AppError'); - -/** - * Fixed-window rate limiter backed by Redis INCR + EXPIRE. - * Fails open if Redis is unreachable so a cache outage cannot lock out users. - */ -function buildRateLimit({ windowSeconds, max, keyPrefix }) { - if (!Number.isFinite(windowSeconds) || windowSeconds <= 0) { - throw new Error('windowSeconds must be a positive number'); - } - if (!Number.isFinite(max) || max <= 0) { - throw new Error('max must be a positive number'); - } - if (!keyPrefix || typeof keyPrefix !== 'string') { - throw new Error('keyPrefix is required'); - } - - return async function rateLimit(req, res, next) { - const identifier = req.ip || req.connection?.remoteAddress || 'unknown'; - const bucket = Math.floor(Date.now() / 1000 / windowSeconds); - const key = `ratelimit:${keyPrefix}:${identifier}:${bucket}`; - - try { - const redis = cache.getClient(); - const count = await redis.incr(key); - if (count === 1) { - await redis.expire(key, windowSeconds); - } - const remaining = Math.max(0, max - count); - const resetAt = (bucket + 1) * windowSeconds; - const retryAfterSeconds = Math.max(1, resetAt - Math.floor(Date.now() / 1000)); - res.setHeader('X-RateLimit-Limit', String(max)); - res.setHeader('X-RateLimit-Remaining', String(remaining)); - res.setHeader('X-RateLimit-Reset', String(resetAt)); - if (count > max) { - res.setHeader('Retry-After', String(retryAfterSeconds)); - return next(new AppError( - 'RATE_LIMITED', - `Rate limit of ${max} requests per ${windowSeconds}s exceeded`, - 429, - { limit: max, window_seconds: windowSeconds, retry_after_seconds: retryAfterSeconds }, - )); - } - return next(); - } catch (err) { - logger.warn('Rate limit fail-open due to cache error', { error: err.message }); - return next(); - } - }; -} - -module.exports = buildRateLimit; +'use strict'; + +const cache = require('../services/cache'); +const logger = require('../logger'); +const AppError = require('../errors/AppError'); + +/** + * Fixed-window rate limiter backed by Redis INCR + EXPIRE. + * Fails open if Redis is unreachable so a cache outage cannot lock out users. + */ +function buildRateLimit({ windowSeconds, max, keyPrefix }) { + if (!Number.isFinite(windowSeconds) || windowSeconds <= 0) { + throw new Error('windowSeconds must be a positive number'); + } + if (!Number.isFinite(max) || max <= 0) { + throw new Error('max must be a positive number'); + } + if (!keyPrefix || typeof keyPrefix !== 'string') { + throw new Error('keyPrefix is required'); + } + + return async function rateLimit(req, res, next) { + const identifier = req.ip || req.connection?.remoteAddress || 'unknown'; + const bucket = Math.floor(Date.now() / 1000 / windowSeconds); + const key = `ratelimit:${keyPrefix}:${identifier}:${bucket}`; + + try { + const redis = cache.getClient(); + const count = await redis.incr(key); + if (count === 1) { + await redis.expire(key, windowSeconds); + } + const remaining = Math.max(0, max - count); + const resetAt = (bucket + 1) * windowSeconds; + const retryAfterSeconds = Math.max(1, resetAt - Math.floor(Date.now() / 1000)); + res.setHeader('X-RateLimit-Limit', String(max)); + res.setHeader('X-RateLimit-Remaining', String(remaining)); + res.setHeader('X-RateLimit-Reset', String(resetAt)); + if (count > max) { + res.setHeader('Retry-After', String(retryAfterSeconds)); + return next(new AppError( + 'RATE_LIMITED', + `Rate limit of ${max} requests per ${windowSeconds}s exceeded`, + 429, + { limit: max, window_seconds: windowSeconds, retry_after_seconds: retryAfterSeconds }, + )); + } + return next(); + } catch (err) { + logger.warn('Rate limit fail-open due to cache error', { error: err.message }); + return next(); + } + }; +} + +module.exports = buildRateLimit; diff --git a/src/routes/prices.js b/src/routes/prices.js index c1d2ccb..ec2b7b1 100644 --- a/src/routes/prices.js +++ b/src/routes/prices.js @@ -1,83 +1,83 @@ -const express = require('express'); -const config = require('../config'); -const { requireApiKey } = require('../middleware/auth'); -const buildRateLimit = require('../middleware/rateLimit'); -const priceOracle = require('../services/priceOracle'); -const AppError = require('../errors/AppError'); - -const router = express.Router(); - -const priceLimit = buildRateLimit({ - windowSeconds: config.priceRateLimit.windowSeconds, - max: config.priceRateLimit.max, - keyPrefix: 'prices', -}); - -router.use(priceLimit); - -function validateAssetCode(assetCode) { - if (!assetCode || typeof assetCode !== 'string') return false; - if (assetCode.length < 1 || assetCode.length > 12) return false; - return /^[A-Z0-9]+$/.test(assetCode); -} - -function validateIssuer(issuer) { - if (!issuer) return true; - return /^G[A-Z0-9]{55}$/.test(issuer); -} - -function validatePriceRequest(assetCode, issuer) { - if (!validateAssetCode(assetCode)) { - throw new AppError('VALIDATION_ERROR', 'Asset code must be 1-12 uppercase alphanumeric characters', 400, { - field: 'assetCode', - received: assetCode, - constraint: 'regex', - }); - } - - if (!validateIssuer(issuer)) { - throw new AppError('VALIDATION_ERROR', 'Issuer must be a valid Stellar address (G...)', 400, { - field: 'issuer', - received: issuer, - constraint: 'stellar_public_key', - }); - } -} - -router.get('/prices/:asset_code', async (req, res, next) => { - try { - const { asset_code } = req.params; - const { issuer } = req.query; - const normalizedCode = asset_code.toUpperCase(); - validatePriceRequest(normalizedCode, issuer); - - const priceData = await priceOracle.getPrice(normalizedCode, issuer || null); - - if (priceData.price_usd === null) { - throw new AppError('NOT_FOUND', `No price data found for ${normalizedCode}`, 404, { asset_code: normalizedCode, issuer: issuer || null }); - } - - return res.json(priceData); - } catch (err) { - return next(err); - } -}); - -router.get('/prices/:asset_code/refresh', requireApiKey(), async (req, res, next) => { - try { - const { asset_code } = req.params; - const { issuer } = req.query; - const normalizedCode = asset_code.toUpperCase(); - validatePriceRequest(normalizedCode, issuer); - - const priceData = await priceOracle.fetchFreshPrice(normalizedCode, issuer || null); - if (priceData.price_usd === null) { - throw new AppError('UPSTREAM_ERROR', 'All price sources failed', 502, { asset_code: normalizedCode, issuer: issuer || null }); - } - return res.json(priceData); - } catch (err) { - return next(err); - } -}); - -module.exports = router; +const express = require('express'); +const config = require('../config'); +const { requireApiKey } = require('../middleware/auth'); +const buildRateLimit = require('../middleware/rateLimit'); +const priceOracle = require('../services/priceOracle'); +const AppError = require('../errors/AppError'); + +const router = express.Router(); + +const priceLimit = buildRateLimit({ + windowSeconds: config.priceRateLimit.windowSeconds, + max: config.priceRateLimit.max, + keyPrefix: 'prices', +}); + +router.use(priceLimit); + +function validateAssetCode(assetCode) { + if (!assetCode || typeof assetCode !== 'string') return false; + if (assetCode.length < 1 || assetCode.length > 12) return false; + return /^[A-Z0-9]+$/.test(assetCode); +} + +function validateIssuer(issuer) { + if (!issuer) return true; + return /^G[A-Z0-9]{55}$/.test(issuer); +} + +function validatePriceRequest(assetCode, issuer) { + if (!validateAssetCode(assetCode)) { + throw new AppError('VALIDATION_ERROR', 'Asset code must be 1-12 uppercase alphanumeric characters', 400, { + field: 'assetCode', + received: assetCode, + constraint: 'regex', + }); + } + + if (!validateIssuer(issuer)) { + throw new AppError('VALIDATION_ERROR', 'Issuer must be a valid Stellar address (G...)', 400, { + field: 'issuer', + received: issuer, + constraint: 'stellar_public_key', + }); + } +} + +router.get('/prices/:asset_code', async (req, res, next) => { + try { + const { asset_code } = req.params; + const { issuer } = req.query; + const normalizedCode = asset_code.toUpperCase(); + validatePriceRequest(normalizedCode, issuer); + + const priceData = await priceOracle.getPrice(normalizedCode, issuer || null); + + if (priceData.price_usd === null) { + throw new AppError('NOT_FOUND', `No price data found for ${normalizedCode}`, 404, { asset_code: normalizedCode, issuer: issuer || null }); + } + + return res.json(priceData); + } catch (err) { + return next(err); + } +}); + +router.get('/prices/:asset_code/refresh', requireApiKey(), async (req, res, next) => { + try { + const { asset_code } = req.params; + const { issuer } = req.query; + const normalizedCode = asset_code.toUpperCase(); + validatePriceRequest(normalizedCode, issuer); + + const priceData = await priceOracle.fetchFreshPrice(normalizedCode, issuer || null); + if (priceData.price_usd === null) { + throw new AppError('UPSTREAM_ERROR', 'All price sources failed', 502, { asset_code: normalizedCode, issuer: issuer || null }); + } + return res.json(priceData); + } catch (err) { + return next(err); + } +}); + +module.exports = router; diff --git a/src/startup/cacheWarm.js b/src/startup/cacheWarm.js new file mode 100644 index 0000000..b2498bb --- /dev/null +++ b/src/startup/cacheWarm.js @@ -0,0 +1,79 @@ +'use strict'; + +const config = require('../config'); +const logger = require('../logger'); +const priceOracle = require('../services/priceOracle'); + +const DEFAULT_TIMEOUT_MS = 30000; + +function isWarmSuccess(result) { + return ( + result.status === 'fulfilled' && + result.value && + result.value.price_usd !== null && + result.value.redis_unavailable !== true + ); +} + +async function runWarmCache(assets, oracle) { + const startedAt = Date.now(); + const results = await Promise.allSettled( + assets.map(({ code, issuer }) => ( + Promise.resolve().then(() => oracle.fetchFreshPrice(code, issuer || null)) + )) + ); + const succeeded = results.filter(isWarmSuccess).length; + + return { + total: assets.length, + succeeded, + failed: assets.length - succeeded, + timedOut: false, + durationMs: Date.now() - startedAt, + }; +} + +async function warmCache( + assets = config.watchedAssets, + oracle = priceOracle, + { timeoutMs = DEFAULT_TIMEOUT_MS, log = logger } = {} +) { + if (!assets || assets.length === 0) { + log.info('Cache warm skipped: no watched assets configured'); + return { total: 0, succeeded: 0, failed: 0, timedOut: false, durationMs: 0 }; + } + + let timedOut = false; + let timeoutId; + + const warming = runWarmCache(assets, oracle).then((summary) => { + if (!timedOut) { + log.info('Cache warm complete', summary); + } + return summary; + }); + + const timeout = new Promise((resolve) => { + timeoutId = setTimeout(() => { + timedOut = true; + const summary = { + total: assets.length, + succeeded: 0, + failed: assets.length, + timedOut: true, + durationMs: timeoutMs, + }; + log.warn('Cache warm timed out; starting server anyway', summary); + resolve(summary); + }, timeoutMs); + }); + + const summary = await Promise.race([warming, timeout]); + if (!summary.timedOut) clearTimeout(timeoutId); + return summary; +} + +module.exports = { + warmCache, + runWarmCache, +}; diff --git a/test/apiRateLimit.test.js b/test/apiRateLimit.test.js index b766bf7..f890239 100644 --- a/test/apiRateLimit.test.js +++ b/test/apiRateLimit.test.js @@ -1,97 +1,97 @@ -'use strict'; - -const express = require('express'); -const request = require('supertest'); -const { createCacheMock } = require('./helpers/cacheMock'); - -const mockHelper = createCacheMock(); -const { reset } = mockHelper; - -jest.mock('../src/services/cache', () => mockHelper.cacheMock); -jest.mock('../src/logger', () => ({ - info: jest.fn(), warn: jest.fn(), error: jest.fn(), debug: jest.fn(), -})); - -const mockGetPrice = jest.fn(); -jest.mock('../src/services/priceOracle', () => ({ - getPrice: mockGetPrice, - fetchFreshPrice: jest.fn(), -})); - -const buildRateLimit = require('../src/middleware/rateLimit'); -const pricesRouter = require('../src/routes/prices'); -const { errorHandler } = require('../src/middleware/errorHandler'); - -function priceResponse() { - return { - asset_code: 'XLM', - issuer: null, - price_usd: 0.12, - source: 'coingecko', - fetched_at: '2026-06-25T00:00:00.000Z', - is_stale: false, - stale_warning: null, - sources_attempted: ['coingecko'], - redis_unavailable: false, - }; -} - -function buildApiApp({ globalMax = 100, globalWindowSeconds = 60 } = {}) { - const app = express(); - app.use(express.json()); - app.use('/api/v1', buildRateLimit({ - windowSeconds: globalWindowSeconds, - max: globalMax, - keyPrefix: 'api', - })); - app.use('/api/v1', pricesRouter); - app.use(errorHandler); - return app; -} - -beforeEach(() => { - reset(); - mockGetPrice.mockReset(); - mockGetPrice.mockResolvedValue(priceResponse()); -}); - -describe('API rate limiting integration', () => { - test('global limit returns 429 after max requests per IP', async () => { - const app = buildApiApp({ globalMax: 2, globalWindowSeconds: 60 }); - - await request(app).get('/api/v1/prices/XLM'); - await request(app).get('/api/v1/prices/XLM'); - const blocked = await request(app).get('/api/v1/prices/XLM'); - - expect(blocked.status).toBe(429); - expect(blocked.body.error.code).toBe('RATE_LIMITED'); - expect(blocked.body.error.details.retry_after_seconds).toBeGreaterThan(0); - expect(blocked.headers['x-ratelimit-limit']).toBe('2'); - expect(blocked.headers['retry-after']).toBeDefined(); - }); - - test('prices routes enforce stricter 30 req/min limit', async () => { - const app = buildApiApp({ globalMax: 100, globalWindowSeconds: 60 }); - - for (let i = 0; i < 30; i += 1) { - const res = await request(app).get('/api/v1/prices/XLM'); - expect(res.status).toBe(200); - expect(res.headers['x-ratelimit-limit']).toBe('30'); - } - - const blocked = await request(app).get('/api/v1/prices/XLM'); - expect(blocked.status).toBe(429); - expect(blocked.body.error.code).toBe('RATE_LIMITED'); - expect(blocked.headers['x-ratelimit-limit']).toBe('30'); - }); - - test('successful responses include rate-limit headers', async () => { - const app = buildApiApp({ globalMax: 100, globalWindowSeconds: 60 }); - const res = await request(app).get('/api/v1/prices/XLM'); - - expect(res.status).toBe(200); - expect(res.headers['x-ratelimit-limit']).toBe('30'); - expect(res.headers['x-ratelimit-remaining']).toBeDefined(); - expect(res.headers['x-ratelimit-reset']).toBeDefined(); - }); -}); +'use strict'; + +const express = require('express'); +const request = require('supertest'); +const { createCacheMock } = require('./helpers/cacheMock'); + +const mockHelper = createCacheMock(); +const { reset } = mockHelper; + +jest.mock('../src/services/cache', () => mockHelper.cacheMock); +jest.mock('../src/logger', () => ({ + info: jest.fn(), warn: jest.fn(), error: jest.fn(), debug: jest.fn(), +})); + +const mockGetPrice = jest.fn(); +jest.mock('../src/services/priceOracle', () => ({ + getPrice: mockGetPrice, + fetchFreshPrice: jest.fn(), +})); + +const buildRateLimit = require('../src/middleware/rateLimit'); +const pricesRouter = require('../src/routes/prices'); +const { errorHandler } = require('../src/middleware/errorHandler'); + +function priceResponse() { + return { + asset_code: 'XLM', + issuer: null, + price_usd: 0.12, + source: 'coingecko', + fetched_at: '2026-06-25T00:00:00.000Z', + is_stale: false, + stale_warning: null, + sources_attempted: ['coingecko'], + redis_unavailable: false, + }; +} + +function buildApiApp({ globalMax = 100, globalWindowSeconds = 60 } = {}) { + const app = express(); + app.use(express.json()); + app.use('/api/v1', buildRateLimit({ + windowSeconds: globalWindowSeconds, + max: globalMax, + keyPrefix: 'api', + })); + app.use('/api/v1', pricesRouter); + app.use(errorHandler); + return app; +} + +beforeEach(() => { + reset(); + mockGetPrice.mockReset(); + mockGetPrice.mockResolvedValue(priceResponse()); +}); + +describe('API rate limiting integration', () => { + test('global limit returns 429 after max requests per IP', async () => { + const app = buildApiApp({ globalMax: 2, globalWindowSeconds: 60 }); + + await request(app).get('/api/v1/prices/XLM'); + await request(app).get('/api/v1/prices/XLM'); + const blocked = await request(app).get('/api/v1/prices/XLM'); + + expect(blocked.status).toBe(429); + expect(blocked.body.error.code).toBe('RATE_LIMITED'); + expect(blocked.body.error.details.retry_after_seconds).toBeGreaterThan(0); + expect(blocked.headers['x-ratelimit-limit']).toBe('2'); + expect(blocked.headers['retry-after']).toBeDefined(); + }); + + test('prices routes enforce stricter 30 req/min limit', async () => { + const app = buildApiApp({ globalMax: 100, globalWindowSeconds: 60 }); + + for (let i = 0; i < 30; i += 1) { + const res = await request(app).get('/api/v1/prices/XLM'); + expect(res.status).toBe(200); + expect(res.headers['x-ratelimit-limit']).toBe('30'); + } + + const blocked = await request(app).get('/api/v1/prices/XLM'); + expect(blocked.status).toBe(429); + expect(blocked.body.error.code).toBe('RATE_LIMITED'); + expect(blocked.headers['x-ratelimit-limit']).toBe('30'); + }); + + test('successful responses include rate-limit headers', async () => { + const app = buildApiApp({ globalMax: 100, globalWindowSeconds: 60 }); + const res = await request(app).get('/api/v1/prices/XLM'); + + expect(res.status).toBe(200); + expect(res.headers['x-ratelimit-limit']).toBe('30'); + expect(res.headers['x-ratelimit-remaining']).toBeDefined(); + expect(res.headers['x-ratelimit-reset']).toBeDefined(); + }); +}); diff --git a/test/cacheWarm.test.js b/test/cacheWarm.test.js new file mode 100644 index 0000000..3fdc2ab --- /dev/null +++ b/test/cacheWarm.test.js @@ -0,0 +1,118 @@ +'use strict'; + +jest.mock('../src/logger', () => ({ + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + debug: jest.fn(), +})); + +jest.mock('../src/services/priceOracle', () => ({ + fetchFreshPrice: jest.fn(), +})); + +const logger = require('../src/logger'); +const priceOracle = require('../src/services/priceOracle'); +const { warmCache } = require('../src/startup/cacheWarm'); + +function asset(code, issuer = null) { + return { code, issuer }; +} + +describe('startup cache warming', () => { + beforeEach(() => { + jest.useRealTimers(); + jest.clearAllMocks(); + }); + + test('skips warming when no assets are configured', async () => { + const summary = await warmCache([], priceOracle, { log: logger }); + + expect(summary).toEqual({ + total: 0, + succeeded: 0, + failed: 0, + timedOut: false, + durationMs: 0, + }); + expect(priceOracle.fetchFreshPrice).not.toHaveBeenCalled(); + expect(logger.info).toHaveBeenCalledWith('Cache warm skipped: no watched assets configured'); + }); + + test('fetches all configured assets and counts cached successes', async () => { + priceOracle.fetchFreshPrice + .mockResolvedValueOnce({ price_usd: 0.12, redis_unavailable: false }) + .mockResolvedValueOnce({ price_usd: 1.0, redis_unavailable: false }) + .mockResolvedValueOnce({ price_usd: null, redis_unavailable: false }); + + const assets = [ + asset('XLM'), + asset('USDC', 'GAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA'), + asset('BAD'), + ]; + + const summary = await warmCache(assets, priceOracle, { log: logger }); + + expect(priceOracle.fetchFreshPrice).toHaveBeenCalledTimes(3); + expect(priceOracle.fetchFreshPrice).toHaveBeenNthCalledWith(1, 'XLM', null); + expect(priceOracle.fetchFreshPrice).toHaveBeenNthCalledWith( + 2, + 'USDC', + 'GAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA' + ); + expect(priceOracle.fetchFreshPrice).toHaveBeenNthCalledWith(3, 'BAD', null); + expect(summary).toMatchObject({ total: 3, succeeded: 2, failed: 1, timedOut: false }); + expect(logger.info).toHaveBeenCalledWith('Cache warm complete', expect.objectContaining({ + total: 3, + succeeded: 2, + failed: 1, + timedOut: false, + })); + }); + + test('starts all asset fetches before awaiting settlement', async () => { + let resolveXlm; + let resolveUsdc; + const xlmPromise = new Promise((resolve) => { resolveXlm = resolve; }); + const usdcPromise = new Promise((resolve) => { resolveUsdc = resolve; }); + + priceOracle.fetchFreshPrice + .mockReturnValueOnce(xlmPromise) + .mockReturnValueOnce(usdcPromise); + + const warming = warmCache([asset('XLM'), asset('USDC')], priceOracle, { log: logger }); + await Promise.resolve(); + + expect(priceOracle.fetchFreshPrice).toHaveBeenCalledTimes(2); + + resolveXlm({ price_usd: 0.12, redis_unavailable: false }); + resolveUsdc({ price_usd: 1.0, redis_unavailable: false }); + await expect(warming).resolves.toMatchObject({ succeeded: 2, failed: 0 }); + }); + + test('returns a timeout summary when warming takes too long', async () => { + jest.useFakeTimers(); + priceOracle.fetchFreshPrice.mockReturnValue(new Promise(() => {})); + + const warming = warmCache([asset('XLM')], priceOracle, { + timeoutMs: 25, + log: logger, + }); + + jest.advanceTimersByTime(25); + await expect(warming).resolves.toEqual({ + total: 1, + succeeded: 0, + failed: 1, + timedOut: true, + durationMs: 25, + }); + expect(logger.warn).toHaveBeenCalledWith('Cache warm timed out; starting server anyway', { + total: 1, + succeeded: 0, + failed: 1, + timedOut: true, + durationMs: 25, + }); + }); +}); diff --git a/test/config.test.js b/test/config.test.js index 51be94f..af3df26 100644 --- a/test/config.test.js +++ b/test/config.test.js @@ -52,6 +52,7 @@ describe('configuration validation', () => { ' databaseUrl: config.databaseUrl,', ' redisUrl: config.redis.url,', ' price: config.price,', + ' watchedAssets: config.watchedAssets,', ' airdrops: config.airdrops,', '}));', ].join(' '), @@ -71,6 +72,7 @@ describe('configuration validation', () => { staleThresholdMinutes: 5, anomalyThresholdPercent: 20, }, + watchedAssets: [], airdrops: { expiryCheckIntervalSeconds: 60, ledgerCacheTtlMs: 5000, @@ -85,4 +87,35 @@ describe('configuration validation', () => { }, }); }); + + test('parses watched assets from WATCHED_ASSETS', () => { + const issuer = 'GAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA'; + const result = runConfig( + [ + "const config = require('./src/config');", + 'console.log(JSON.stringify(config.watchedAssets));', + ].join(' '), + { + NODE_ENV: 'test', + WATCHED_ASSETS: `XLM,USDC:${issuer},XLM`, + } + ); + + expect(result.status).toBe(0); + expect(JSON.parse(result.stdout.trim())).toEqual([ + { code: 'XLM', issuer: null }, + { code: 'USDC', issuer }, + ]); + }); + + test('rejects malformed watched assets during config loading', () => { + const result = runConfig("require('./src/config')", { + NODE_ENV: 'test', + WATCHED_ASSETS: 'usdc:not-a-stellar-address', + }); + + const output = `${result.stdout}\n${result.stderr}`; + expect(result.status).toBe(1); + expect(output).toContain('WATCHED_ASSETS'); + }); }); diff --git a/test/rateLimit.test.js b/test/rateLimit.test.js index 55248a1..292b5e8 100644 --- a/test/rateLimit.test.js +++ b/test/rateLimit.test.js @@ -1,53 +1,53 @@ -'use strict'; - -const express = require('express'); -const request = require('supertest'); -const { createCacheMock } = require('./helpers/cacheMock'); - -const mockHelper = createCacheMock(); -const { reset } = mockHelper; - -jest.mock('../src/services/cache', () => mockHelper.cacheMock); -jest.mock('../src/logger', () => ({ - info: jest.fn(), warn: jest.fn(), error: jest.fn(), debug: jest.fn(), -})); - -const buildRateLimit = require('../src/middleware/rateLimit'); -const { errorHandler } = require('../src/middleware/errorHandler'); - -function buildApp(limiter) { - const app = express(); - app.use(limiter); - app.get('/test', (_req, res) => res.json({ ok: true })); - app.use(errorHandler); - return app; -} - -beforeEach(() => reset()); - -describe('rateLimit middleware', () => { - test('allows requests under the limit and sets rate-limit headers', async () => { - const app = buildApp(buildRateLimit({ windowSeconds: 60, max: 3, keyPrefix: 't' })); - const r1 = await request(app).get('/test'); - expect(r1.status).toBe(200); - expect(r1.headers['x-ratelimit-limit']).toBe('3'); - expect(r1.headers['x-ratelimit-remaining']).toBe('2'); - }); - - test('returns 429 once the limit is exceeded', async () => { - const app = buildApp(buildRateLimit({ windowSeconds: 60, max: 2, keyPrefix: 'lim' })); - await request(app).get('/test'); - await request(app).get('/test'); - const blocked = await request(app).get('/test'); - expect(blocked.status).toBe(429); - expect(blocked.body.error).toMatchObject({ code: 'RATE_LIMITED' }); - expect(blocked.body.error.details.retry_after_seconds).toBeGreaterThan(0); - expect(blocked.headers['retry-after']).toBeDefined(); - }); - - test('throws when configured with invalid options', () => { - expect(() => buildRateLimit({ windowSeconds: 0, max: 10, keyPrefix: 'x' })).toThrow(); - expect(() => buildRateLimit({ windowSeconds: 60, max: 0, keyPrefix: 'x' })).toThrow(); - expect(() => buildRateLimit({ windowSeconds: 60, max: 10 })).toThrow(); - }); -}); +'use strict'; + +const express = require('express'); +const request = require('supertest'); +const { createCacheMock } = require('./helpers/cacheMock'); + +const mockHelper = createCacheMock(); +const { reset } = mockHelper; + +jest.mock('../src/services/cache', () => mockHelper.cacheMock); +jest.mock('../src/logger', () => ({ + info: jest.fn(), warn: jest.fn(), error: jest.fn(), debug: jest.fn(), +})); + +const buildRateLimit = require('../src/middleware/rateLimit'); +const { errorHandler } = require('../src/middleware/errorHandler'); + +function buildApp(limiter) { + const app = express(); + app.use(limiter); + app.get('/test', (_req, res) => res.json({ ok: true })); + app.use(errorHandler); + return app; +} + +beforeEach(() => reset()); + +describe('rateLimit middleware', () => { + test('allows requests under the limit and sets rate-limit headers', async () => { + const app = buildApp(buildRateLimit({ windowSeconds: 60, max: 3, keyPrefix: 't' })); + const r1 = await request(app).get('/test'); + expect(r1.status).toBe(200); + expect(r1.headers['x-ratelimit-limit']).toBe('3'); + expect(r1.headers['x-ratelimit-remaining']).toBe('2'); + }); + + test('returns 429 once the limit is exceeded', async () => { + const app = buildApp(buildRateLimit({ windowSeconds: 60, max: 2, keyPrefix: 'lim' })); + await request(app).get('/test'); + await request(app).get('/test'); + const blocked = await request(app).get('/test'); + expect(blocked.status).toBe(429); + expect(blocked.body.error).toMatchObject({ code: 'RATE_LIMITED' }); + expect(blocked.body.error.details.retry_after_seconds).toBeGreaterThan(0); + expect(blocked.headers['retry-after']).toBeDefined(); + }); + + test('throws when configured with invalid options', () => { + expect(() => buildRateLimit({ windowSeconds: 0, max: 10, keyPrefix: 'x' })).toThrow(); + expect(() => buildRateLimit({ windowSeconds: 60, max: 0, keyPrefix: 'x' })).toThrow(); + expect(() => buildRateLimit({ windowSeconds: 60, max: 10 })).toThrow(); + }); +});