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
1 change: 0 additions & 1 deletion .github/workflows/build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -115,4 +115,3 @@ jobs:
run: npm test --if-present

# Workflow run retention settings
retention-days: 30
1 change: 0 additions & 1 deletion .github/workflows/contract-release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -29,4 +29,3 @@ jobs:
release_token: ${{ secrets.GITHUB_TOKEN }}

# Workflow run retention settings
retention-days: 90
1 change: 0 additions & 1 deletion .github/workflows/dapp-ipfs.yml
Original file line number Diff line number Diff line change
Expand Up @@ -67,4 +67,3 @@ jobs:
echo "- URL: ${{ steps.storacha.outputs.url }}" >> "$GITHUB_STEP_SUMMARY"

# Workflow run retention settings
retention-days: 30
1 change: 0 additions & 1 deletion .github/workflows/secrets-check.yml
Original file line number Diff line number Diff line change
Expand Up @@ -22,4 +22,3 @@ jobs:
run: ./scripts/check-k8s-secrets.sh

# Workflow run retention settings
retention-days: 30
16 changes: 16 additions & 0 deletions backend/src/controllers/healthController.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,22 @@ export const redis: Redis | null = config.REDIS_URL
})
: null;

export async function closeHealthDependencies(): Promise<void> {
const closers: Promise<unknown>[] = [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;
Expand Down
91 changes: 61 additions & 30 deletions backend/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand All @@ -18,13 +22,21 @@ 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<typeof setInterval> | undefined;
let idempotencyCleanup: ReturnType<typeof setInterval> | undefined;
let isShuttingDown = false;

// Initialize Socket.IO
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`);
Expand All @@ -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();
Expand All @@ -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();
Expand All @@ -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'));
66 changes: 66 additions & 0 deletions backend/src/services/__tests__/gracefulShutdown.test.ts
Original file line number Diff line number Diff line change
@@ -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();
});
});
3 changes: 3 additions & 0 deletions backend/src/services/contractEventIndexer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand Down
91 changes: 91 additions & 0 deletions backend/src/services/gracefulShutdown.ts
Original file line number Diff line number Diff line change
@@ -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<void>;
closeDependencies: () => Promise<void>;
timeoutMs?: number;
exit?: (code: number) => void;
}

async function drainHttpServer(
server: Server,
logger: ShutdownLogger,
timeoutMs: number
): Promise<void> {
await new Promise<void>((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<void> {
let shutdownPromise: Promise<void> | null = null;

return (signal: NodeJS.Signals): Promise<void> => {
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;
};
}
4 changes: 4 additions & 0 deletions backend/src/services/rateLimitService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -261,3 +261,7 @@ export class RateLimitService {
}

export const rateLimitService = new RateLimitService();

export async function closeRateLimitRedis(): Promise<void> {
await RedisClient.disconnect();
}