fix: add service disposal to graceful shutdown

This commit is contained in:
Brad Groux 2026-01-28 17:13:49 -06:00
parent 279221ee38
commit 1e0842ce4e
6 changed files with 164 additions and 46 deletions

View file

@ -3590,3 +3590,5 @@
{"type":"task.status_changed","taskId":"task_20260128_3X46QP","project":"veritas-kanban","status":"in-progress","previousStatus":"todo","id":"evt_ZvJmZLewWSeq","timestamp":"2026-01-28T23:11:16.711Z"}
{"type":"task.status_changed","taskId":"task_20260128_uUydZN","project":"veritas-kanban","status":"done","previousStatus":"in-progress","id":"evt_hqwd3Ec6SsJY","timestamp":"2026-01-28T23:11:21.221Z"}
{"type":"task.status_changed","taskId":"task_20260128_nhfrcW","project":"veritas-kanban","status":"in-progress","previousStatus":"todo","id":"evt_h4DH9ZZp9jPt","timestamp":"2026-01-28T23:11:54.652Z"}
{"type":"task.status_changed","taskId":"task_20260128_3X46QP","project":"veritas-kanban","status":"done","previousStatus":"in-progress","id":"evt_sOAaMhqeJX84","timestamp":"2026-01-28T23:12:25.180Z"}
{"type":"task.status_changed","taskId":"task_20260128_gYodue","project":"veritas-kanban","status":"in-progress","previousStatus":"todo","id":"evt__V2Fj9ltjzJA","timestamp":"2026-01-28T23:13:04.737Z"}

View file

@ -1,4 +1,37 @@
[
{
"id": "activity_1769641984737_iexbxcytg",
"type": "status_changed",
"taskId": "task_20260128_gYodue",
"taskTitle": "PERF: Add service disposal to graceful shutdown",
"details": {
"from": "todo",
"status": "in-progress"
},
"timestamp": "2026-01-28T23:13:04.737Z"
},
{
"id": "activity_1769641953788_zrhhxvv9q",
"type": "comment_added",
"taskId": "task_20260128_3X46QP",
"taskTitle": "QUALITY: Create CI pipeline (.github/workflows)",
"details": {
"author": "Veritas",
"preview": "Created .github/workflows/ci.yml with 4 parallel j..."
},
"timestamp": "2026-01-28T23:12:33.788Z"
},
{
"id": "activity_1769641945180_ybvr2ju43",
"type": "status_changed",
"taskId": "task_20260128_3X46QP",
"taskTitle": "QUALITY: Create CI pipeline (.github/workflows)",
"details": {
"from": "in-progress",
"status": "done"
},
"timestamp": "2026-01-28T23:12:25.180Z"
},
{
"id": "activity_1769641914696_o9slvvwd6",
"type": "status_changed",

View file

@ -8,12 +8,14 @@ import { WebSocketServer, WebSocket } from 'ws';
import { createServer } from 'http';
import path from 'path';
import { fileURLToPath } from 'url';
import { createLogger } from './lib/logger.js';
import { v1Router } from './routes/v1/index.js';
import { agentService } from './routes/agents.js';
import { syncSettingsToServices } from './routes/settings.js';
import { initAgentStatus } from './routes/agent-status.js';
import { getTelemetryService } from './services/telemetry-service.js';
import { ConfigService } from './services/config-service.js';
import { disposeTaskService } from './services/task-service.js';
import { initBroadcast } from './services/broadcast-service.js';
import { runStartupMigrations } from './services/migration-service.js';
import { errorHandler } from './middleware/error-handler.js';
@ -33,6 +35,8 @@ import attachmentRoutes from './routes/attachments.js';
import { configRoutes } from './routes/config.js';
import { agentRoutes } from './routes/agents.js';
const log = createLogger('server');
const app = express();
const PORT = process.env.PORT || 3001;
@ -122,7 +126,7 @@ const corsOptions: cors.CorsOptions = {
if (ALLOWED_ORIGINS.includes(origin)) {
callback(null, true);
} else {
console.warn(`CORS: Blocked request from origin: ${origin}`);
log.warn({ origin }, 'CORS blocked request from disallowed origin');
callback(new Error('Not allowed by CORS'));
}
},
@ -245,6 +249,9 @@ if (process.env.NODE_ENV === 'production') {
// Error handling middleware (must be last)
app.use(errorHandler);
// Module-level config service instance (shared with shutdown handler)
let configService: ConfigService | null = null;
// Initialize services on startup
(async () => {
try {
@ -252,12 +259,12 @@ app.use(errorHandler);
await runStartupMigrations();
// Initialize telemetry service and sync with feature settings
const configService = new ConfigService();
configService = new ConfigService();
const featureSettings = await configService.getFeatureSettings();
syncSettingsToServices(featureSettings);
await getTelemetryService().init();
} catch (err) {
console.error('Failed to initialize services:', err);
log.error({ err }, 'Failed to initialize services');
}
})();
@ -275,7 +282,7 @@ const wss = new WebSocketServer({
const result = validateWebSocketOrigin(origin, ALLOWED_ORIGINS);
if (!result.allowed) {
console.warn(`WebSocket origin rejected: ${origin} — ${result.reason}`);
log.warn({ origin, reason: result.reason }, 'WebSocket origin rejected');
callback(false, 403, 'Forbidden: origin not allowed');
return;
}
@ -298,7 +305,7 @@ wss.on('connection', (ws: AuthenticatedWebSocket, req) => {
const authResult = authenticateWebSocket(req);
if (!authResult.authenticated) {
console.log('WebSocket connection rejected: ' + authResult.error);
log.warn({ error: authResult.error }, 'WebSocket connection rejected');
ws.close(4001, authResult.error || 'Authentication required');
return;
}
@ -310,7 +317,7 @@ wss.on('connection', (ws: AuthenticatedWebSocket, req) => {
isLocalhost: authResult.isLocalhost,
};
console.log(`WebSocket client connected (role: ${authResult.role}, localhost: ${authResult.isLocalhost})`);
log.info({ role: authResult.role, localhost: authResult.isLocalhost }, 'WebSocket client connected');
let subscribedTaskId: string | null = null;
@ -396,12 +403,12 @@ wss.on('connection', (ws: AuthenticatedWebSocket, req) => {
}));
}
} catch (error) {
console.error('WebSocket message error:', error);
log.error({ err: error }, 'WebSocket message error');
}
});
ws.on('close', () => {
console.log('WebSocket client disconnected');
log.info('WebSocket client disconnected');
// Clean up subscriptions
if (subscribedTaskId) {
@ -420,24 +427,55 @@ wss.on('connection', (ws: AuthenticatedWebSocket, req) => {
export { wss };
// Graceful shutdown handler
function gracefulShutdown(signal: string) {
console.log(`\n${signal} received. Shutting down gracefully...`);
async function gracefulShutdown(signal: string) {
log.info({ signal }, 'Shutting down gracefully');
// Close WebSocket connections
console.log('Closing WebSocket connections...');
// 1. Close WebSocket connections first (stop accepting new messages)
log.info({ clients: wss.clients.size }, 'Closing WebSocket connections');
wss.clients.forEach((client) => {
client.close(1000, 'Server shutting down');
});
// Close HTTP server
// Close the WebSocket server itself (stop accepting new connections)
await new Promise<void>((resolve) => {
wss.close((err) => {
if (err) log.error({ err }, 'Error closing WebSocket server');
else log.info('WebSocket server closed');
resolve();
});
});
// 2. Dispose services (release file watchers, flush buffers)
try {
log.info('Disposing services');
// Flush pending telemetry writes
await getTelemetryService().flush();
log.info('Telemetry flushed');
// Dispose task service (closes file watchers, clears cache)
disposeTaskService();
log.info('Task service disposed');
// Dispose config service (closes file watcher, clears cache)
if (configService) {
configService.dispose();
configService = null;
log.info('Config service disposed');
}
} catch (err) {
log.error({ err }, 'Error during service disposal');
}
// 3. Close HTTP server last
server.close(() => {
console.log('HTTP server closed.');
log.info('HTTP server closed');
process.exit(0);
});
// Force exit after 10 seconds
setTimeout(() => {
console.error('Forced shutdown after timeout.');
log.fatal('Forced shutdown after timeout');
process.exit(1);
}, 10000);
}
@ -457,34 +495,29 @@ server.listen(PORT, () => {
: 'Auth: OFF (dev mode)';
const corsLine = `CORS: ${ALLOWED_ORIGINS.length} origins`;
console.log(`
╔═══════════════════════════════════════════════╗
║ Veritas Kanban Server ║
╠═══════════════════════════════════════════════╣
║ API: http://localhost:${PORT} ║
║ WebSocket: ws://localhost:${PORT}/ws ║
║ Health: http://localhost:${PORT}/health ║
║ ${authLine.padEnd(42)}║
║ ${corsLine.padEnd(42)}║
║ Helmet: ON (CSP + security headers) ║
║ Compress: ON (gzip, threshold 1KB) ║
║ Rate Limit: 100 req/min ║
║ Body Limit: 1MB ║
╚═══════════════════════════════════════════════╝
`);
log.info({
port: PORT,
api: `http://localhost:${PORT}`,
ws: `ws://localhost:${PORT}/ws`,
auth: authLine,
cors: corsLine,
helmet: true,
compression: true,
rateLimit: '100 req/min',
bodyLimit: '1MB',
}, 'Veritas Kanban Server started');
// Security warnings for localhost bypass
if (authStatus.localhostBypass) {
if (authStatus.localhostRole === 'admin') {
console.warn(
'⚠️ WARNING: Localhost bypass is active with ADMIN role.\n' +
' Any local process can read, modify, or delete all data without authentication.\n' +
' Set VERITAS_AUTH_LOCALHOST_ROLE=read-only or disable bypass for production.\n'
log.warn(
{ localhostRole: 'admin' },
'Localhost bypass is active with ADMIN role — any local process has full access without authentication'
);
} else {
console.log(
`ℹ️ Localhost bypass active (role: ${authStatus.localhostRole}). ` +
'Local connections can read data without authentication.'
log.info(
{ localhostRole: authStatus.localhostRole },
'Localhost bypass active — local connections can read data without authentication'
);
}
}

44
server/src/lib/logger.ts Normal file
View file

@ -0,0 +1,44 @@
import pino from 'pino';
const isDev = process.env.NODE_ENV !== 'production';
/**
* Structured logger using pino.
*
* - JSON output in production (machine-readable for log aggregators)
* - pino-pretty in development (human-readable, colorized)
* - Log level controlled via LOG_LEVEL env var (default: 'info')
* - Includes timestamp, pid, hostname in every log line
*/
export const logger = pino({
level: process.env.LOG_LEVEL || 'info',
// pino includes pid and hostname by default in JSON mode.
// In dev we use pino-pretty as a transport for colorized output.
...(isDev
? {
transport: {
target: 'pino-pretty',
options: {
colorize: true,
translateTime: 'HH:MM:ss.l',
ignore: 'pid,hostname',
},
},
}
: {
// Production: JSON with explicit ISO timestamp
timestamp: pino.stdTimeFunctions.isoTime,
}),
});
/**
* Create a child logger scoped to a specific module / component.
* Usage:
* const log = createLogger('auth');
* log.info({ role: 'admin' }, 'User authenticated');
*/
export function createLogger(component: string) {
return logger.child({ component });
}
export default logger;

View file

@ -1,4 +1,7 @@
import { Request, Response, NextFunction } from 'express';
import { createLogger } from '../lib/logger.js';
const log = createLogger('error-handler');
// Custom error classes
export class AppError extends Error {
@ -54,7 +57,7 @@ export function errorHandler(
return res.status(err.statusCode).json(response);
}
console.error('Unhandled error:', err);
log.error({ err, requestId, method: req.method, path: req.path }, 'Unhandled error');
res.status(500).json({
error: 'Internal server error',
...(requestId ? { requestId } : {}),

View file

@ -5,6 +5,9 @@ import matter from 'gray-matter';
import { nanoid } from 'nanoid';
import type { Task, CreateTaskInput, UpdateTaskInput, ReviewComment, Subtask, TaskTelemetryEvent, TimeTracking } from '@veritas-kanban/shared';
import { getTelemetryService, type TelemetryService } from './telemetry-service.js';
import { createLogger } from '../lib/logger.js';
const log = createLogger('task-cache');
/**
* Task ID format validation
@ -81,7 +84,7 @@ export class TaskService {
this.cacheLoading = null;
this.cacheInitialized = true;
this.startWatcher();
console.debug(`[TaskCache] Initialized with ${this.cache.size} tasks`);
log.debug({ count: this.cache.size }, 'Cache initialized');
}
/** Read every .md file in tasksDir and populate the cache */
@ -110,7 +113,7 @@ export class TaskService {
const content = await fs.readFile(filepath, 'utf-8');
const task = this.parseTaskFile(content, filename);
if (task) {
console.debug(`[TaskCache] Reloaded ${task.id} from disk`);
log.debug({ taskId: task.id }, 'Reloaded from disk');
this.cache.set(task.id, task);
}
} catch {
@ -126,7 +129,7 @@ export class TaskService {
if (idMatch) {
const id = idMatch[1];
if (this.cache.delete(id)) {
console.debug(`[TaskCache] Invalidated ${id} (file removed)`);
log.debug({ taskId: id }, 'Invalidated (file removed)');
}
}
}
@ -135,7 +138,7 @@ export class TaskService {
private cacheInvalidate(id: string): boolean {
const deleted = this.cache.delete(id);
if (deleted) {
console.debug(`[TaskCache] Invalidated ${id}`);
log.debug({ taskId: id }, 'Invalidated');
}
return deleted;
}
@ -145,10 +148,10 @@ export class TaskService {
const task = this.cache.get(id);
if (task) {
this.cacheStats.hits++;
console.debug(`[TaskCache] HIT ${id} (hits=${this.cacheStats.hits})`);
log.trace({ taskId: id, hits: this.cacheStats.hits }, 'Cache HIT');
} else {
this.cacheStats.misses++;
console.debug(`[TaskCache] MISS ${id} (misses=${this.cacheStats.misses})`);
log.trace({ taskId: id, misses: this.cacheStats.misses }, 'Cache MISS');
}
return task;
}
@ -175,10 +178,10 @@ export class TaskService {
// Ignore events caused by our own writes
if (Date.now() - this.lastWriteTime < WRITE_DEBOUNCE_MS) return;
console.debug(`[TaskCache] File change detected: ${eventType} ${filename}`);
log.debug({ eventType, filename }, 'File change detected');
// Re-read the changed file (or remove from cache if deleted)
this.reloadFile(filename).catch(err =>
console.error(`[TaskCache] Error reloading ${filename}:`, err),
log.error({ err, filename }, 'Error reloading file'),
);
});
} catch (err) {