diff --git a/backend/src/controllers/healthController.ts b/backend/src/controllers/healthController.ts index 2f11c9b8..709946ed 100644 --- a/backend/src/controllers/healthController.ts +++ b/backend/src/controllers/healthController.ts @@ -13,6 +13,22 @@ export const redis: Redis | null = config.REDIS_URL }) : null; +export async function closeHealthDependencies(): Promise { + const closers: Promise[] = [pool.end()]; + if (redis) { + closers.push(redis.quit()); + } + + const results = await Promise.allSettled(closers); + const failures = results + .filter((result): result is PromiseRejectedResult => result.status === 'rejected') + .map((result) => result.reason); + + if (failures.length > 0) { + throw new AggregateError(failures, 'Failed to close health-check dependencies'); + } +} + export interface DependencyStatus { status: 'connected' | 'disconnected' | 'not_configured' | 'unknown'; error?: string; diff --git a/backend/src/index.ts b/backend/src/index.ts index c883ae89..d2b3abe6 100644 --- a/backend/src/index.ts +++ b/backend/src/index.ts @@ -3,6 +3,10 @@ import { createServer } from 'node:http'; import app from './app.js'; import logger from './utils/logger.js'; import config from './config/index.js'; +import pool from './config/database.js'; +import { closeHealthDependencies } from './controllers/healthController.js'; +import { closeRateLimitRedis } from './services/rateLimitService.js'; +import { createGracefulShutdown } from './services/gracefulShutdown.js'; import { initializeSocket } from './services/socketService.js'; import { scheduleExecutor } from './services/scheduleExecutor.js'; import { contractEventIndexer } from './services/contractEventIndexer.js'; @@ -18,6 +22,9 @@ const server = createServer(app); // Part-49 job handles — assigned on server start, cleaned up on shutdown let usageSnapshotJob: { stop(): void }; let integrityCheckJob: { stop(): void }; +let auditCacheCleanup: ReturnType | undefined; +let idempotencyCleanup: ReturnType | undefined; +let isShuttingDown = false; // Initialize Socket.IO initializeSocket(server); @@ -25,6 +32,11 @@ initializeSocket(server); const PORT = config.port || process.env.PORT || 4000; server.listen(PORT, () => { + if (isShuttingDown) { + server.close(); + return; + } + logger.info(`Server running on port ${PORT}`); logger.info(`Environment: ${config.nodeEnv}`); logger.info(`Health check: http://localhost:${PORT}/health`); @@ -48,7 +60,7 @@ server.listen(PORT, () => { logger.info('Part-49 jobs scheduled (usage snapshots + audit integrity)'); // Part 45 — cleanup expired audit cache every hour - setInterval( + auditCacheCleanup = setInterval( async () => { try { const deleted = await auditAnalyticsService.cleanupExpiredCache(); @@ -64,7 +76,7 @@ server.listen(PORT, () => { logger.info('Part-45 audit cache cleanup scheduled'); // Idempotency key cleanup — every hour, remove expired keys - setInterval( + idempotencyCleanup = setInterval( async () => { try { await cleanupExpiredIdempotencyKeys(); @@ -77,35 +89,54 @@ server.listen(PORT, () => { logger.info('Idempotency key cleanup scheduled'); }); -// Graceful shutdown handling -const shutdown = () => { - logger.info('Shutting down gracefully...'); - - // Stop the schedule executor - scheduleExecutor.stop(); - - liquidityAlertChecker.stop(); - - // Stop Part-49 cron jobs - usageSnapshotJob?.stop(); - integrityCheckJob?.stop(); - - // Stop the contract event indexer - contractEventIndexer.stop(); - - // Close the server - server.close(() => { - logger.info('Server closed'); - process.exit(0); - }); +// Stop accepting HTTP before stopping future background work. The shared +// shutdown service drains HTTP for up to 30 seconds before closing dependencies. +const gracefulShutdown = createGracefulShutdown({ + server, + logger, + stopBackgroundWork: async () => { + if (auditCacheCleanup) clearInterval(auditCacheCleanup); + if (idempotencyCleanup) clearInterval(idempotencyCleanup); + + const stops = [ + () => scheduleExecutor.stop(), + () => liquidityAlertChecker.stop(), + () => usageSnapshotJob?.stop(), + () => integrityCheckJob?.stop(), + () => contractEventIndexer.stop(), + ]; + const results = await Promise.allSettled( + stops.map((stop) => Promise.resolve().then(stop)) + ); + const failures = results + .filter((result): result is PromiseRejectedResult => result.status === 'rejected') + .map((result) => result.reason); + if (failures.length > 0) { + throw new AggregateError(failures, 'Failed to stop one or more background jobs'); + } + }, + closeDependencies: async () => { + const results = await Promise.allSettled([ + Promise.resolve().then(() => pool.end()), + Promise.resolve().then(closeHealthDependencies), + Promise.resolve().then(closeRateLimitRedis), + ]); + const failures = results + .filter((result): result is PromiseRejectedResult => result.status === 'rejected') + .map((result) => result.reason); + if (failures.length > 0) { + throw new AggregateError(failures, 'Failed to close one or more backend dependencies'); + } + }, +}); - // Force shutdown after 10 seconds - setTimeout(() => { - logger.error('Forced shutdown after timeout'); +const shutdown = (signal: NodeJS.Signals): void => { + isShuttingDown = true; + void gracefulShutdown(signal).catch((error) => { + logger.error('Graceful shutdown failed', { error }); process.exit(1); - }, 10000); + }); }; -// Listen for termination signals -process.on('SIGTERM', shutdown); -process.on('SIGINT', shutdown); +process.on('SIGTERM', () => shutdown('SIGTERM')); +process.on('SIGINT', () => shutdown('SIGINT')); diff --git a/backend/src/services/__tests__/gracefulShutdown.test.ts b/backend/src/services/__tests__/gracefulShutdown.test.ts new file mode 100644 index 00000000..a443b512 --- /dev/null +++ b/backend/src/services/__tests__/gracefulShutdown.test.ts @@ -0,0 +1,66 @@ +import type { Server } from 'node:http'; +import { createGracefulShutdown } from '../gracefulShutdown'; + +describe('createGracefulShutdown', () => { + it('drains HTTP before dependency cleanup and exits cleanly', async () => { + let onClose: (() => void) | undefined; + const server = { + close: jest.fn((callback: () => void) => { + onClose = callback; + return server; + }), + closeAllConnections: jest.fn(), + } as unknown as Server; + const stopBackgroundWork = jest.fn(); + const closeDependencies = jest.fn().mockResolvedValue(undefined); + const exit = jest.fn(); + const logger = { info: jest.fn(), warn: jest.fn(), error: jest.fn() }; + + const shutdown = createGracefulShutdown({ + server, + logger, + stopBackgroundWork, + closeDependencies, + exit, + }); + + const pending = shutdown('SIGTERM'); + expect(server.close).toHaveBeenCalledTimes(1); + expect(stopBackgroundWork).toHaveBeenCalledTimes(1); + expect(closeDependencies).not.toHaveBeenCalled(); + + onClose?.(); + await pending; + + expect(closeDependencies).toHaveBeenCalledTimes(1); + expect(exit).toHaveBeenCalledWith(0); + }); + + it('uses the 30 second drain bound before cleanup', async () => { + jest.useFakeTimers(); + const server = { + close: jest.fn(() => server), + closeAllConnections: jest.fn(), + } as unknown as Server; + const closeDependencies = jest.fn().mockResolvedValue(undefined); + const exit = jest.fn(); + const logger = { info: jest.fn(), warn: jest.fn(), error: jest.fn() }; + + const pending = createGracefulShutdown({ + server, + logger, + stopBackgroundWork: jest.fn(), + closeDependencies, + exit, + })('SIGINT'); + + await Promise.resolve(); + jest.advanceTimersByTime(30_000); + await pending; + + expect(server.closeAllConnections).toHaveBeenCalledTimes(1); + expect(closeDependencies).toHaveBeenCalledTimes(1); + expect(exit).toHaveBeenCalledWith(0); + jest.useRealTimers(); + }); +}); diff --git a/backend/src/services/contractEventIndexer.ts b/backend/src/services/contractEventIndexer.ts index 7e66c746..d772d7a9 100644 --- a/backend/src/services/contractEventIndexer.ts +++ b/backend/src/services/contractEventIndexer.ts @@ -37,6 +37,9 @@ export class ContractEventIndexer { // Run immediately on startup await this.pollAndIndexEvents(); + + // Shutdown may have stopped the indexer while the initial poll was pending. + if (!this.isRunning) return; // Then poll at regular intervals this.intervalId = setInterval(async () => { diff --git a/backend/src/services/gracefulShutdown.ts b/backend/src/services/gracefulShutdown.ts new file mode 100644 index 00000000..0ae49e77 --- /dev/null +++ b/backend/src/services/gracefulShutdown.ts @@ -0,0 +1,91 @@ +import type { Server } from 'node:http'; + +export interface ShutdownLogger { + info(message: string): unknown; + warn(message: string): unknown; + error(message: string, meta?: unknown): unknown; +} + +export interface GracefulShutdownOptions { + server: Server; + logger: ShutdownLogger; + stopBackgroundWork: () => void | Promise; + closeDependencies: () => Promise; + timeoutMs?: number; + exit?: (code: number) => void; +} + +async function drainHttpServer( + server: Server, + logger: ShutdownLogger, + timeoutMs: number +): Promise { + await new Promise((resolve) => { + let settled = false; + let timeout: NodeJS.Timeout | undefined; + + const finish = () => { + if (settled) return; + settled = true; + if (timeout) clearTimeout(timeout); + resolve(); + }; + + server.close((error) => { + if (error && (error as NodeJS.ErrnoException).code !== 'ERR_SERVER_NOT_RUNNING') { + logger.error('HTTP server close failed', { error }); + } + finish(); + }); + + if (!settled) { + timeout = setTimeout(() => { + logger.warn(`HTTP drain exceeded ${timeoutMs}ms; closing remaining connections`); + server.closeAllConnections(); + finish(); + }, timeoutMs); + } + }); +} + +export function createGracefulShutdown({ + server, + logger, + stopBackgroundWork, + closeDependencies, + timeoutMs = 30_000, + exit = (code) => process.exit(code), +}: GracefulShutdownOptions): (signal: NodeJS.Signals) => Promise { + let shutdownPromise: Promise | null = null; + + return (signal: NodeJS.Signals): Promise => { + if (shutdownPromise) return shutdownPromise; + + shutdownPromise = (async () => { + logger.info(`Received ${signal}; starting graceful shutdown`); + + // server.close() runs immediately inside drainHttpServer, so new HTTP + // connections are refused before background cleanup begins. + const httpDrain = drainHttpServer(server, logger, timeoutMs); + + try { + await stopBackgroundWork(); + } catch (error) { + logger.error('Failed to stop background work cleanly', { error }); + } + + await httpDrain; + + try { + await closeDependencies(); + } catch (error) { + logger.error('Failed to close one or more runtime dependencies', { error }); + } + + logger.info('Graceful shutdown complete'); + exit(0); + })(); + + return shutdownPromise; + }; +} diff --git a/backend/src/services/rateLimitService.ts b/backend/src/services/rateLimitService.ts index f90773e7..4cd5fb64 100644 --- a/backend/src/services/rateLimitService.ts +++ b/backend/src/services/rateLimitService.ts @@ -261,3 +261,7 @@ export class RateLimitService { } export const rateLimitService = new RateLimitService(); + +export async function closeRateLimitRedis(): Promise { + await RedisClient.disconnect(); +}