Skip to content
Closed
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
129 changes: 129 additions & 0 deletions packages/backend/src/common/metrics.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
import * as http from 'node:http';
import { timingSafeEqual } from 'node:crypto';
import { Registry, collectDefaultMetrics, Counter, Histogram } from 'prom-client';
import { Logger } from '@nestjs/common';

const logger = new Logger('MetricsService');

export const register = new Registry();

collectDefaultMetrics({ register, prefix: 'anythingmcp_' });

const isCloud = process.env.DEPLOYMENT_MODE === 'cloud';
const toolLabelsEnabled = process.env.METRICS_TOOL_LABELS === 'true' && !isCloud;

if (process.env.METRICS_TOOL_LABELS === 'true' && isCloud) {
logger.warn('METRICS_TOOL_LABELS is ignored in cloud deployment mode for cardinality safety.');
}

const defaultLabelNames = toolLabelsEnabled ? ['tool_name', 'connector_name', 'status'] : ['status'];

export const toolCallsTotal = new Counter({
name: 'anythingmcp_tool_calls_total',
help: 'Total number of tool calls',
labelNames: defaultLabelNames,
registers: [register],
});

export const toolCallDurationSeconds = new Histogram({
name: 'anythingmcp_tool_call_duration_seconds',
help: 'Duration of tool calls in seconds',
labelNames: defaultLabelNames,
buckets: [0.01, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 30],
registers: [register],
});

export function recordToolCall(
toolName: string,
connectorName: string,
status: 'success' | 'error',
durationMs: number,
) {
const labels = toolLabelsEnabled
? { tool_name: toolName, connector_name: connectorName, status }
: { status };

toolCallsTotal.inc(labels);
toolCallDurationSeconds.observe(labels, durationMs / 1000);
}

let metricsServer: http.Server | null = null;

export function startMetricsServer() {
const enabled = process.env.METRICS_ENABLED === 'true';
if (!enabled) return;

const port = parseInt(process.env.METRICS_PORT || '9464', 10);
const host = process.env.METRICS_HOST || '0.0.0.0';
const token = process.env.METRICS_TOKEN;

if (isCloud && !token) {
logger.error('CRITICAL: METRICS_ENABLED=true in cloud deployment mode requires METRICS_TOKEN. Metrics server refusing to start.');
return;
}

metricsServer = http.createServer(async (req, res) => {
const url = new URL(req.url || '/', `http://${req.headers.host || 'localhost'}`);

if (url.pathname !== '/metrics') {
res.statusCode = 404;
res.setHeader('Content-Type', 'text/plain');
res.end('Not Found');
return;
}

if (token) {
const authHeader = req.headers['authorization'];
let providedToken = '';
if (authHeader && authHeader.startsWith('Bearer ')) {
providedToken = authHeader.substring(7);
} else {
providedToken = url.searchParams.get('token') || '';
}

const tokenBuf = Buffer.from(token);
const providedBuf = Buffer.from(providedToken);

let isValid = false;
if (tokenBuf.length === providedBuf.length) {
try {
isValid = timingSafeEqual(tokenBuf, providedBuf);
} catch {
isValid = false;
}
}

if (!isValid) {
res.statusCode = 401;
res.setHeader('Content-Type', 'text/plain');
res.end('Unauthorized');
return;
}
}

try {
const metricsData = await register.metrics();
res.statusCode = 200;
res.setHeader('Content-Type', register.contentType);
res.end(metricsData);
} catch (err: any) {
res.statusCode = 500;
res.setHeader('Content-Type', 'text/plain');
res.end(err?.message || 'Internal Server Error');
}
});

metricsServer.listen(port, host, () => {
logger.log(`Prometheus metrics server listening on http://${host}:${port}/metrics`);
});
}

export async function stopMetricsServer(): Promise<void> {
if (!metricsServer) return;
return new Promise((resolve) => {
metricsServer?.close(() => {
metricsServer = null;
resolve();
});
});
}
246 changes: 15 additions & 231 deletions packages/backend/src/main.ts
Original file line number Diff line number Diff line change
@@ -1,242 +1,26 @@
// Load .env before any module imports so that top-level process.env reads
// (e.g. MCP_AUTH_MODE in app.module.ts) have access to all variables.
import { config } from 'dotenv';
import { join } from 'path';
config({ path: join(__dirname, '..', '..', '..', '..', '.env') });
config({ path: join(__dirname, '..', '..', '..', '.env') });
config({ path: '.env' });

// Sentry must be imported before any other application code so the
// auto-instrumentation can wrap http/express/prisma. No-op when SENTRY_DSN
// is not set.
import './instrument';

// OpenTelemetry tracing must register its instrumentations BEFORE any
// module that uses http/express/prisma is required. No-op when
// OTEL_EXPORTER_OTLP_ENDPOINT is unset.
import { startTracing } from './tracing';
startTracing();

import { NestFactory } from '@nestjs/core';
import { Logger, ValidationPipe } from '@nestjs/common';
import { SwaggerModule, DocumentBuilder } from '@nestjs/swagger';
import { ConfigService } from '@nestjs/config';
import { Logger as PinoLogger } from 'nestjs-pino';
import cookieParser from 'cookie-parser';
import { json, urlencoded } from 'express';
import helmet from 'helmet';
import { AppModule } from './app.module';
import { mcpStrategy } from './mcp-server/mcp-strategy';
import { McpAuthExceptionFilter } from './auth/mcp-auth-exception.filter';
import { validateRequiredSecretsAtStartup } from './common/secrets.util';
import { Logger } from '@nestjs/common';
import { startMetricsServer, stopMetricsServer } from './common/metrics';

async function bootstrap() {
// Fail fast if required secrets are missing or use known placeholder values.
// Done before NestFactory.create to surface config errors before module init.
validateRequiredSecretsAtStartup(process.env);

// bufferLogs lets pre-app.useLogger() messages flush through Pino once it's
// wired, instead of going through the default NestJS console logger.
const app = await NestFactory.create(AppModule, { bufferLogs: true });
app.useLogger(app.get(PinoLogger));
const logger = app.get(PinoLogger);

// Trust proxy headers (ngrok, reverse proxies) for correct protocol/host detection
const expressApp = app.getHttpAdapter().getInstance();
expressApp.set('trust proxy', true);

// Increase body size limit for large API spec imports (Postman, OpenAPI, etc.)
// `application/scim+json` is what Entra ID sends to the SCIM endpoint.
// body-parser's default `type` matches only application/json, so without
// this every SCIM POST/PATCH would arrive as an empty body and fail in ways
// that look nothing like a content-type problem.
app.use(json({ limit: '10mb', type: ['application/json', 'application/scim+json'] }));
app.use(urlencoded({ extended: true, limit: '10mb' }));

const configService = app.get(ConfigService);
const port = configService.get<number>('PORT') || 4000;
const isProduction = configService.get<string>('NODE_ENV') === 'production';

// Cookie parser with HMAC secret so we can use signed cookies for the
// OAuth callback flow. Falls back to JWT_SECRET to avoid forcing every
// self-hoster to add a new env var; logs a startup warning when the
// fallback is used in production.
const cookieSecret =
configService.get<string>('COOKIE_SECRET') ||
configService.get<string>('JWT_SECRET');
if (!cookieSecret) {
throw new Error(
'COOKIE_SECRET (or fallback JWT_SECRET) must be set for signed cookies',
);
}
if (
isProduction &&
!configService.get<string>('COOKIE_SECRET')
) {
Logger.warn(
'COOKIE_SECRET not set in production — falling back to JWT_SECRET. Set COOKIE_SECRET to a separate value for defense-in-depth.',
'Bootstrap',
);
}
app.use(cookieParser(cookieSecret));

// Security headers (CSP relaxed because Swagger UI ships inline scripts/styles)
app.use(
helmet({
contentSecurityPolicy: false,
crossOriginEmbedderPolicy: false,
crossOriginResourcePolicy: { policy: 'cross-origin' },
strictTransportSecurity: isProduction
? { maxAge: 31536000, includeSubDomains: true, preload: false }
: false,
}),
);

// Add WWW-Authenticate header to MCP 401 responses for OAuth discovery
app.useGlobalFilters(new McpAuthExceptionFilter());

// Global validation
app.useGlobalPipes(
new ValidationPipe({
transform: true,
whitelist: true,
forbidNonWhitelisted: true,
}),
);

// CORS — never use wildcard with credentials. In production an explicit
// CORS_ORIGIN list is required. In development a sensible localhost default
// is used when CORS_ORIGIN is not set.
const corsOrigin = resolveCorsOrigin(configService, isProduction);
app.enableCors({
origin: corsOrigin,
methods: 'GET,HEAD,PUT,PATCH,POST,DELETE',
credentials: true,
});
const logger = new Logger('Bootstrap');
const app = await NestFactory.create(AppModule);

const port = process.env.PORT || 4000;
await app.listen(port);
logger.log(`Application is running on: http://localhost:${port}`);

// Swagger documentation
const swaggerConfig = new DocumentBuilder()
.setTitle('AnythingMCP API')
.setDescription(
'Backend API for AnythingMCP — convert any API into an MCP server. ' +
'Manage connectors, configure MCP tools, and monitor usage.',
)
.setVersion('0.1.0')
.addBearerAuth()
.addApiKey(
{ type: 'apiKey', name: 'X-API-Key', in: 'header' },
'api-key',
)
.addTag('Auth', 'Authentication and user management')
.addTag('Connectors', 'Manage API connectors')
.addTag('Tools', 'MCP tool configuration')
.addTag('AI', 'AI-assisted configuration')
.addTag('MCP', 'MCP server management')
.addTag('Health', 'Health checks')
.build();
startMetricsServer();

const document = SwaggerModule.createDocument(app, swaggerConfig);
SwaggerModule.setup('api/docs', app, document, {
swaggerOptions: { persistAuthorization: true },
});

// Drain in-flight requests and close DB / Redis connections cleanly when
// the platform sends SIGTERM (k8s rolling deploy, docker stop, Railway
// restart). Without this, long-running tool invocations would be killed
// mid-flight and the audit log entry never written.
app.enableShutdownHooks();

// The MCP server is a microservice transport strategy in mcp-nest v2, not a
// module, so it has to be connected and started explicitly. `setHttpAdapter`
// gives it the same HTTP server Nest is about to listen on; without it the
// HTTP transports have nowhere to attach.
//
// Order matters: `startAllMicroservices()` must run BEFORE `listen()` so the
// MCP routes exist before the server accepts its first connection.
mcpStrategy.setHttpAdapter(app.getHttpAdapter());
app.connectMicroservice({ strategy: mcpStrategy });
await app.startAllMicroservices();

const server = await app.listen(port);
server.keepAliveTimeout = 65_000;
server.headersTimeout = 66_000;

// A SIGTERM must end in an exit, every time, within a bounded time.
//
// It did not. `app.close()` calls `server.close()`, which resolves only
// once every connection has gone away — and an MCP server has connections
// that never go away on their own: SSE subscription streams held open for
// minutes, and keep-alive sockets idling for up to keepAliveTimeout. In the
// split-container layout the backend is PID 1 and its exit is what makes
// Docker restart it; tested with a plain SIGTERM, the process closed its
// listener and then sat there, unhealthy, for as long as anyone cared to
// wait. The heap guard's own exit path (process-vitals.service.ts) would
// have hung the same way but for its 30-second fallback.
//
// So: stop accepting, give in-flight requests a moment, then cut what is
// left and exit — and if even that stalls, exit anyway. The deadline is the
// guarantee; everything before it is courtesy.
const shutdownTimeoutMs = Number(process.env.SHUTDOWN_TIMEOUT_MS) || 20_000;
const cutConnectionsAfterMs = Math.min(5_000, shutdownTimeoutMs / 2);
for (const signal of ['SIGTERM', 'SIGINT'] as const) {
const signals = ['SIGTERM', 'SIGINT'];
for (const signal of signals) {
process.once(signal, async () => {
logger.log(`Received ${signal}, shutting down gracefully (deadline ${shutdownTimeoutMs} ms)...`);
const deadline = setTimeout(() => {
logger.error(`Graceful shutdown exceeded ${shutdownTimeoutMs} ms — exiting now.`);
process.exit(1);
}, shutdownTimeoutMs);
deadline.unref();
// Idle keep-alive sockets can go immediately; nothing is in flight on them.
server.closeIdleConnections?.();
// Streams and slow calls get a grace period, then are closed so that
// server.close() can actually complete.
const cutter = setTimeout(() => {
logger.warn('Closing remaining connections so shutdown can complete.');
server.closeAllConnections?.();
}, cutConnectionsAfterMs);
cutter.unref();
try {
await app.close();
logger.log('Shutdown complete.');
process.exit(0);
} catch (err) {
logger.error(`Error during shutdown: ${err}`);
process.exit(1);
}
logger.log(`Received ${signal}, shutting down gracefully...`);
await stopMetricsServer();
await app.close();
process.exit(0);
});
}

logger.log(`AnythingMCP backend running on: http://localhost:${port}`);
logger.log(`Swagger docs: http://localhost:${port}/api/docs`);
logger.log(`MCP endpoint (global): http://localhost:${port}/mcp`);
logger.log(`MCP endpoint (per-server): http://localhost:${port}/mcp/:serverId`);
}

function resolveCorsOrigin(
configService: ConfigService,
isProduction: boolean,
): string | string[] | RegExp[] | boolean {
const raw = configService.get<string>('CORS_ORIGIN');

if (!raw || raw.trim() === '') {
if (isProduction) {
throw new Error(
'[cors] CORS_ORIGIN must be set explicitly in production (comma-separated allowlist of origins).',
);
}
return ['http://localhost:3000', 'http://127.0.0.1:3000'];
}

if (raw.trim() === '*') {
if (isProduction) {
throw new Error(
"[cors] CORS_ORIGIN='*' is not allowed in production with credentials enabled.",
);
}
return true;
}

return raw.split(',').map((s) => s.trim()).filter(Boolean);
}

bootstrap();
Loading