Files
ruvnet--ruflo/v3/mcp/connection-pool.ts
wehub-resource-sync 23f7624596
ADR-166 MCP Bridge Security Lock / Static-source security lock (push) Failing after 0s
ADR-166 MCP Bridge Security Lock / Compose default binds loopback + Mongo has auth (push) Failing after 2s
CodeQL Advanced / Analyze (rust) (push) Failing after 0s
ADR-166 MCP Bridge Security Lock / plugin-agent-federation bindHost default (push) Failing after 1s
ADR-166 MCP Bridge Security Lock / Runtime behavior — 401 + terminal gate + fail-closed (push) Failing after 4s
business-pods-smoke / smoke (push) Failing after 1s
all-plugins-smoke / smoke-all (push) Failing after 2s
CI/CD Pipeline / Security & Code Quality (push) Failing after 1s
CI/CD Pipeline / Test Suite (ubuntu-latest) (push) Failing after 1s
CI/CD Pipeline / Build & Package (macos-latest) (push) Has been skipped
CI/CD Pipeline / Build & Package (ubuntu-latest) (push) Has been skipped
CI/CD Pipeline / Build & Package (windows-latest) (push) Has been skipped
CI/CD Pipeline / Documentation & Examples (push) Failing after 1s
Clone Tracker (14-day rolling) / Snapshot clones for ruflo ecosystem (push) Failing after 1s
CodeQL Advanced / Analyze (actions) (push) Failing after 1s
CodeQL Advanced / Analyze (javascript-typescript) (push) Failing after 1s
federation-peer-rust / stable-noop (push) Failing after 1s
metaharness-ci / score (push) Failing after 1s
metaharness-ci / router-compat (push) Failing after 0s
metaharness-ci / similarity-tests (push) Failing after 0s
no-agentbbs-smoke / smoke-without-agentbbs (push) Failing after 1s
V3 CI/CD Pipeline / Build V3 (windows-latest) (push) Has been skipped
codex-integration-audit / Codex integration audit (push) Failing after 1s
helpers-manifest-guard / guard (push) Failing after 1s
🔗 Cross-Agent Integration Tests / 🤝 Agent Coordination Tests (push) Has been skipped
🔗 Cross-Agent Integration Tests / 🧠 Memory Sharing Integration (push) Has been skipped
🔗 Cross-Agent Integration Tests / 🛡️ Fault Tolerance Tests (push) Has been skipped
🔗 Cross-Agent Integration Tests / ⚡ Performance Integration Tests (push) Has been skipped
metaharness-ci / mcp-scan (push) Failing after 1s
metaharness-ci / eject-dryrun (push) Failing after 1s
metaharness-ci / metaharness-real-data (push) Failing after 0s
no-cli-optdep-bloat-2561 / guard (push) Failing after 1s
no-metaharness-smoke / smoke-without-metaharness (push) Failing after 1s
no-phantom-agentic-flow-subpath / guard (push) Failing after 1s
🔄 Automated Rollback Manager / 🚨 Failure Detection (push) Failing after 1s
V3 CI/CD Pipeline / Plugin hooks smoke / ubuntu-latest / Node 22 (push) Failing after 1s
V3 CI/CD Pipeline / ruflo-graph-intelligence build + test smoke (#2044, ADR-123) (push) Failing after 1s
CVE Audit Gate / Audit root (critical-blocking) (push) Failing after 2s
cost-tracker-smoke / smoke (push) Failing after 3s
oia-audit-weekly / audit (push) Failing after 2s
ruflo-agent-smoke / ruflo-agent structural smoke (push) Failing after 1s
📊 Status Badges Update / 📊 Update Status Badges (push) Failing after 1s
V3 CI/CD Pipeline / Static regression guards (#2267 YAML + (push) Failing after 1s
V3 CI/CD Pipeline / Test V3 Packages (push) Failing after 0s
V3 CI/CD Pipeline / agent_execute provider routing smoke (#2042) (push) Failing after 0s
CVE Audit Gate / Audit v3 (critical-blocking) (push) Failing after 1s
federation-peer-rust / stable-native (push) Failing after 2s
🔗 Cross-Agent Integration Tests / 🚀 Integration Test Setup (push) Failing after 2s
neural-trader-smoke / runtime-smoke (push) Failing after 1s
V3 CI/CD Pipeline / Build V3 (macos-latest) (push) Has been skipped
V3 CI/CD Pipeline / Build V3 (ubuntu-latest) (push) Has been skipped
V3 CI/CD Pipeline / Type Check V3 (push) Failing after 1s
V3 CI/CD Pipeline / Smoke (no better-sqlite3) / ubuntu-latest / Node 24 (push) Failing after 1s
V3 CI/CD Pipeline / Smoke (no better-sqlite3) / ubuntu-latest / Node 22 (push) Failing after 2s
V3 CI/CD Pipeline / browser rvf create flag smoke (#2015) (push) Failing after 0s
V3 CI/CD Pipeline / Dependency review (#2046) (push) Has been skipped
V3 CI/CD Pipeline / Supply-chain audit (#2046) (push) Failing after 0s
V3 CI/CD Pipeline / witness marker drift smoke (#2021) (push) Failing after 1s
V3 CI/CD Pipeline / neural-trader portfolio CG smoke (#2068, ADR-126 Phase 3) (push) Failing after 1s
V3 CI/CD Pipeline / neural-trader backtest signing smoke (#2068, ADR-126 Phase 4) (push) Failing after 1s
V3 CI/CD Pipeline / kg-extract type-import classification smoke (#2049) (push) Failing after 0s
V3 CI/CD Pipeline / witness verify precondition smoke (#1880) (push) Failing after 2s
V3 CI/CD Pipeline / neural-trader pipeline risk-gate smoke (#2068, ADR-126 Phase 5) (push) Failing after 0s
V3 CI/CD Pipeline / neural-trader feature attribution smoke (#2068, ADR-126 Phase 6) (push) Failing after 0s
V3 CI/CD Pipeline / plugin-registry signature verification smoke (#1922, CWE-347) (push) Failing after 4s
V3 CI/CD Pipeline / memory stats legacy-DB smoke (#2120) (push) Failing after 4s
V3 CI/CD Pipeline / github deprecated actions smoke (#2089, ADR-127 Phase 3) (push) Failing after 1s
V3 CI/CD Pipeline / graph query + pathfinder smoke (ADR-130 P2+P5) (push) Has been skipped
V3 CI/CD Pipeline / graph trajectory hooks smoke (ADR-130 P3) (push) Has been skipped
V3 CI/CD Pipeline / graph plugin adapter smoke (ADR-130 P4) (push) Has been skipped
V3 CI/CD Pipeline / graph benchmark (ADR-130 P6) (push) Has been skipped
V3 CI/CD Pipeline / statusline generator delegation smoke (#2195) (push) Failing after 1s
V3 CI/CD Pipeline / wizard init regression guard (#2206 (push) Failing after 1s
V3 CI/CD Pipeline / memory no-stray-db smoke (ADR-125 P7) (push) Failing after 1s
V3 CI/CD Pipeline / github-safe injection smoke (#2089, ADR-127 Phase 1) (push) Failing after 1s
V3 CI/CD Pipeline / github actions pin smoke (#2089, ADR-127 Phase 1) (push) Failing after 1s
V3 CI/CD Pipeline / github attribution opt-in smoke (#2089, ADR-127 Phase 4) (push) Failing after 1s
V3 CI/CD Pipeline / pre-bash hook safety smoke (#2017) (push) Failing after 1s
V3 CI/CD Pipeline / Memory import smoke / ubuntu-latest (push) Failing after 0s
V3 CI/CD Pipeline / MCP protocol smoke / ubuntu-latest (push) Failing after 2s
V3 CI/CD Pipeline / ruvllm WASM auto-init smoke (#2086) (push) Failing after 4s
V3 CI/CD Pipeline / MCP paired-tool round-trip smoke (#1889) (push) Failing after 1s
V3 CI/CD Pipeline / Plugin package install-safety (#1902/#1903/#1904) (push) Failing after 1s
V3 CI/CD Pipeline / Tool description discoverability (ADR-112) (push) Failing after 3s
V3 CI/CD Pipeline / CLI npx-install smoke (#1147 / (22) (push) Failing after 1s
V3 CI/CD Pipeline / CLI npx-install smoke (#1147 / (24) (push) Failing after 1s
V3 CI/CD Pipeline / Windows hook shim smoke (#2132) / ubuntu-latest (push) Failing after 2s
V3 CI/CD Pipeline / Windows hook execution smoke (#2132) / ubuntu-latest (push) Failing after 1s
V3 CI/CD Pipeline / Windows init hooks smoke (#2132) / ubuntu-latest (push) Failing after 1s
V3 CI/CD Pipeline / Vector-index dimension audit (#1947) (push) Failing after 0s
V3 CI/CD Pipeline / Hook-command install safety (#1921) (push) Failing after 1s
V3 CI/CD Pipeline / ToolOutputGuardrail smoke (ADR-131, (push) Failing after 1s
V3 CI/CD Pipeline / init-bundle invariants smoke (#2095, ADR-128 Phase 5) (push) Failing after 1s
V3 CI/CD Pipeline / wasm provider bridge smoke (ADR-129 P1) (push) Failing after 2s
V3 CI/CD Pipeline / wasm gallery CRUD smoke (ADR-129 P3) (push) Failing after 1s
V3 CI/CD Pipeline / wasm plugin bridge smoke (ADR-129 P4) (push) Failing after 0s
V3 CI/CD Pipeline / wasm compose smoke (ADR-129 P2) (push) Failing after 4s
V3 CI/CD Pipeline / graph schema smoke (ADR-130 P1) (push) Failing after 0s
Validate Marketplace / validate (push) Failing after 1s
🔍 Verification Pipeline / 🚀 Setup Verification (push) Failing after 1s
🔍 Verification Pipeline / 🛡️ Security Verification (push) Has been skipped
🔍 Verification Pipeline / 📝 Code Quality (push) Has been skipped
🔍 Verification Pipeline / 🧪 Test Verification (${{ matrix.os }}, Node ${{ matrix.node }}) (push) Has been skipped
🔍 Verification Pipeline / 🏗️ Build Verification (push) Has been skipped
🔍 Verification Pipeline / 📚 Documentation Verification (push) Has been skipped
CVE Audit Gate / High-severity report (warn only) (push) Has been cancelled
🔄 Automated Rollback Manager / 🔄 Execute Rollback (push) Has been cancelled
🔄 Automated Rollback Manager / ✅ Post-Rollback Verification (push) Has been cancelled
🔄 Automated Rollback Manager / 📊 Rollback Monitoring (push) Has been cancelled
V3 CI/CD Pipeline / Windows init hooks smoke (#2132) / windows-latest (push) Has been cancelled
V3 CI/CD Pipeline / Windows hook execution smoke (#2132) / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Windows hook execution smoke (#2132) / windows-latest (push) Has been cancelled
🔄 Automated Rollback Manager / ⏳ Manual Rollback Approval (push) Has been cancelled
V3 CI/CD Pipeline / MCP protocol smoke / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Memory import smoke / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Windows hook shim smoke (#2132) / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Windows hook shim smoke (#2132) / windows-latest (push) Has been cancelled
V3 CI/CD Pipeline / Windows init hooks smoke (#2132) / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Witness verify (signed manifest) / macos-latest (push) Has been cancelled
V3 CI/CD Pipeline / Witness verify (signed manifest) / ubuntu-latest (push) Has been cancelled
V3 CI/CD Pipeline / Witness verify (signed manifest) / windows-latest (push) Has been cancelled
V3 CI/CD Pipeline / Publish to npm (alpha) (push) Has been cancelled
V3 CI/CD Pipeline / Smoke (no better-sqlite3) / macos-latest / Node 22 (push) Has been cancelled
V3 CI/CD Pipeline / Plugin hooks smoke / macos-latest / Node 22 (push) Has been cancelled
CI/CD Pipeline / Deploy & Release (push) Has been cancelled
CI/CD Pipeline / CI Status (push) Has been cancelled
🔗 Cross-Agent Integration Tests / 📊 Integration Test Report (push) Has been cancelled
🔄 Automated Rollback Manager / 🔍 Pre-Rollback Validation (push) Has been cancelled
🔍 Verification Pipeline / ⚡ Performance Verification (push) Has been cancelled
🔍 Verification Pipeline / 📊 Verification Report (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 12:02:19 +08:00

439 lines
12 KiB
TypeScript

/**
* V3 MCP Connection Pool Manager
*
* High-performance connection pooling for MCP server:
* - Reusable connections to reduce overhead
* - Max connections: 10 (configurable)
* - Idle timeout handling with automatic eviction
* - Connection health monitoring
* - Graceful shutdown support
*
* Performance Targets:
* - Connection acquire: <5ms
* - Connection release: <1ms
*/
import { EventEmitter } from 'events';
import {
PooledConnection,
ConnectionPoolStats,
ConnectionPoolConfig,
ConnectionState,
IConnectionPool,
ILogger,
TransportType,
} from './types.js';
/**
* Default connection pool configuration
*/
const DEFAULT_POOL_CONFIG: ConnectionPoolConfig = {
maxConnections: 10,
minConnections: 2,
idleTimeout: 30000, // 30 seconds
acquireTimeout: 5000, // 5 seconds
maxWaitingClients: 50,
evictionRunInterval: 10000, // 10 seconds
};
/**
* Connection wrapper with lifecycle management
*/
class ManagedConnection implements PooledConnection {
public state: ConnectionState = 'idle';
public lastUsedAt: Date;
public useCount: number = 0;
constructor(
public readonly id: string,
public readonly transport: TransportType,
public readonly createdAt: Date = new Date(),
public metadata?: Record<string, unknown>
) {
this.lastUsedAt = this.createdAt;
}
/**
* Mark connection as busy
*/
acquire(): void {
this.state = 'busy';
this.lastUsedAt = new Date();
this.useCount++;
}
/**
* Mark connection as idle
*/
release(): void {
this.state = 'idle';
this.lastUsedAt = new Date();
}
/**
* Check if connection is expired
*/
isExpired(idleTimeout: number): boolean {
if (this.state !== 'idle') return false;
return Date.now() - this.lastUsedAt.getTime() > idleTimeout;
}
/**
* Check if connection is healthy
*/
isHealthy(): boolean {
return this.state !== 'error' && this.state !== 'closed';
}
}
/**
* Waiting client for connection
*/
interface WaitingClient {
resolve: (connection: PooledConnection) => void;
reject: (error: Error) => void;
timestamp: number;
}
/**
* Connection Pool Manager
*
* Manages a pool of reusable connections for optimal performance
*/
export class ConnectionPool extends EventEmitter implements IConnectionPool {
private readonly config: ConnectionPoolConfig;
private readonly connections: Map<string, ManagedConnection> = new Map();
private readonly waitingClients: WaitingClient[] = [];
private evictionTimer?: NodeJS.Timeout;
private connectionCounter: number = 0;
private isShuttingDown: boolean = false;
// Statistics
private stats = {
totalAcquired: 0,
totalReleased: 0,
totalCreated: 0,
totalDestroyed: 0,
acquireTimeTotal: 0,
acquireCount: 0,
};
constructor(
config: Partial<ConnectionPoolConfig> = {},
private readonly logger: ILogger,
private readonly transportType: TransportType = 'in-process'
) {
super();
this.config = { ...DEFAULT_POOL_CONFIG, ...config };
this.startEvictionTimer();
this.initializeMinConnections();
}
/**
* Initialize minimum number of connections
*/
private async initializeMinConnections(): Promise<void> {
const promises: Promise<void>[] = [];
for (let i = 0; i < this.config.minConnections; i++) {
promises.push(this.createConnection());
}
await Promise.all(promises);
this.logger.debug('Connection pool initialized', {
minConnections: this.config.minConnections,
});
}
/**
* Create a new connection
*/
private async createConnection(): Promise<ManagedConnection> {
const id = `conn-${++this.connectionCounter}-${Date.now()}`;
const connection = new ManagedConnection(id, this.transportType);
this.connections.set(id, connection);
this.stats.totalCreated++;
this.emit('pool:connection:created', { connectionId: id });
this.logger.debug('Connection created', { id, total: this.connections.size });
return connection;
}
/**
* Acquire a connection from the pool
*/
async acquire(): Promise<PooledConnection> {
const startTime = performance.now();
if (this.isShuttingDown) {
throw new Error('Connection pool is shutting down');
}
// Try to find an idle connection
for (const connection of this.connections.values()) {
if (connection.state === 'idle' && connection.isHealthy()) {
connection.acquire();
this.stats.totalAcquired++;
this.recordAcquireTime(startTime);
this.emit('pool:connection:acquired', { connectionId: connection.id });
this.logger.debug('Connection acquired from pool', { id: connection.id });
return connection;
}
}
// Create new connection if under limit
if (this.connections.size < this.config.maxConnections) {
const connection = await this.createConnection();
connection.acquire();
this.stats.totalAcquired++;
this.recordAcquireTime(startTime);
this.emit('pool:connection:acquired', { connectionId: connection.id });
return connection;
}
// Wait for a connection to become available
return this.waitForConnection(startTime);
}
/**
* Wait for a connection to become available
*/
private waitForConnection(startTime: number): Promise<PooledConnection> {
return new Promise((resolve, reject) => {
if (this.waitingClients.length >= this.config.maxWaitingClients) {
reject(new Error('Connection pool exhausted - max waiting clients reached'));
return;
}
const client: WaitingClient = {
resolve: (connection) => {
this.recordAcquireTime(startTime);
resolve(connection);
},
reject,
timestamp: Date.now(),
};
this.waitingClients.push(client);
// Set timeout
setTimeout(() => {
const index = this.waitingClients.indexOf(client);
if (index !== -1) {
this.waitingClients.splice(index, 1);
reject(new Error(`Connection acquire timeout after ${this.config.acquireTimeout}ms`));
}
}, this.config.acquireTimeout);
});
}
/**
* Release a connection back to the pool
*/
release(connection: PooledConnection): void {
const managed = this.connections.get(connection.id);
if (!managed) {
this.logger.warn('Attempted to release unknown connection', { id: connection.id });
return;
}
// Check for waiting clients first
const waitingClient = this.waitingClients.shift();
if (waitingClient) {
managed.acquire();
this.stats.totalAcquired++;
this.emit('pool:connection:acquired', { connectionId: connection.id });
waitingClient.resolve(managed);
return;
}
// Return to pool
managed.release();
this.stats.totalReleased++;
this.emit('pool:connection:released', { connectionId: connection.id });
this.logger.debug('Connection released to pool', { id: connection.id });
}
/**
* Destroy a connection (remove from pool)
*/
destroy(connection: PooledConnection): void {
const managed = this.connections.get(connection.id);
if (!managed) {
return;
}
managed.state = 'closed';
this.connections.delete(connection.id);
this.stats.totalDestroyed++;
this.emit('pool:connection:destroyed', { connectionId: connection.id });
this.logger.debug('Connection destroyed', { id: connection.id });
// Create new connection to maintain minimum if needed
if (this.connections.size < this.config.minConnections && !this.isShuttingDown) {
this.createConnection().catch((err) => {
this.logger.error('Failed to create replacement connection', err);
});
}
}
/**
* Get pool statistics
*/
getStats(): ConnectionPoolStats {
let idleCount = 0;
let busyCount = 0;
for (const connection of this.connections.values()) {
if (connection.state === 'idle') idleCount++;
else if (connection.state === 'busy') busyCount++;
}
return {
totalConnections: this.connections.size,
idleConnections: idleCount,
busyConnections: busyCount,
pendingRequests: this.waitingClients.length,
totalAcquired: this.stats.totalAcquired,
totalReleased: this.stats.totalReleased,
totalCreated: this.stats.totalCreated,
totalDestroyed: this.stats.totalDestroyed,
avgAcquireTime: this.stats.acquireCount > 0
? this.stats.acquireTimeTotal / this.stats.acquireCount
: 0,
};
}
/**
* Drain the pool (wait for all connections to be released)
*/
async drain(): Promise<void> {
this.isShuttingDown = true;
this.logger.info('Draining connection pool');
// Reject all waiting clients
while (this.waitingClients.length > 0) {
const client = this.waitingClients.shift();
client?.reject(new Error('Connection pool is draining'));
}
// Wait for busy connections to be released
const maxWait = 10000; // 10 seconds
const startTime = Date.now();
while (Date.now() - startTime < maxWait) {
let busyCount = 0;
for (const connection of this.connections.values()) {
if (connection.state === 'busy') busyCount++;
}
if (busyCount === 0) break;
await new Promise((resolve) => setTimeout(resolve, 100));
}
this.logger.info('Connection pool drained');
}
/**
* Clear all connections from the pool
*/
async clear(): Promise<void> {
this.stopEvictionTimer();
await this.drain();
// Destroy all remaining connections
for (const connection of this.connections.values()) {
connection.state = 'closed';
}
this.connections.clear();
this.logger.info('Connection pool cleared');
}
/**
* Start the eviction timer
*/
private startEvictionTimer(): void {
this.evictionTimer = setInterval(() => {
this.evictIdleConnections();
}, this.config.evictionRunInterval);
}
/**
* Stop the eviction timer
*/
private stopEvictionTimer(): void {
if (this.evictionTimer) {
clearInterval(this.evictionTimer);
this.evictionTimer = undefined;
}
}
/**
* Evict idle connections that have exceeded the timeout
*/
private evictIdleConnections(): void {
if (this.isShuttingDown) return;
const toEvict: ManagedConnection[] = [];
for (const connection of this.connections.values()) {
if (
connection.isExpired(this.config.idleTimeout) &&
this.connections.size > this.config.minConnections
) {
toEvict.push(connection);
}
}
for (const connection of toEvict) {
this.destroy(connection);
this.logger.debug('Evicted idle connection', { id: connection.id });
}
if (toEvict.length > 0) {
this.logger.info('Evicted idle connections', { count: toEvict.length });
}
}
/**
* Record acquire time for statistics
*/
private recordAcquireTime(startTime: number): void {
const duration = performance.now() - startTime;
this.stats.acquireTimeTotal += duration;
this.stats.acquireCount++;
}
/**
* Get all connections (for debugging/monitoring)
*/
getConnections(): PooledConnection[] {
return Array.from(this.connections.values());
}
/**
* Check if pool is healthy
*/
isHealthy(): boolean {
return !this.isShuttingDown && this.connections.size >= this.config.minConnections;
}
}
/**
* Create a connection pool with default settings
*/
export function createConnectionPool(
config: Partial<ConnectionPoolConfig> = {},
logger: ILogger,
transportType: TransportType = 'in-process'
): ConnectionPool {
return new ConnectionPool(config, logger, transportType);
}