import { EventEmitter } from "node:events"; interface WorkerLike { id?: number; threadId?: number; process?: { connected: boolean; }; isDead?: () => boolean; isConnected?: () => boolean; send: (message: unknown) => void; on: (event: string, listener: (...args: unknown[]) => void) => this; off: (event: string, listener: (...args: unknown[]) => void) => this; } interface WorkerEntry { id: number; worker: WorkerLike; mode: string; socketPath: string | null; _messageHandler?: (msg: unknown) => void; _errorHandler?: (err: unknown) => void; } interface UnregisterOptions { suppressBroadcast?: boolean; } export interface ServiceStats { workerCount: number; ids: number[]; directSockets: number; } /** * DirectMessageBus — Unix Domain Socket Mesh (Supervisor Side) * * The fundamental problem: cluster.Worker.send() can only transfer * sockets/servers, not MessagePort objects. So we can't use MessageChannel * for direct worker-to-worker communication in cluster mode. * * Solution: Each worker opens a Unix domain socket server. Workers connect * directly to each other via these sockets. The supervisor only tells * workers WHERE to connect (socket paths), then gets out of the way. * * Architecture: * * SETUP (supervisor involved briefly): * 1. Each worker starts a UDS server at /tmp/forge-{pid}/{service}-{worker}.sock * 2. Worker reports its socket path to supervisor via IPC * 3. Supervisor broadcasts the full socket registry to all workers * 4. Workers establish direct connections to each other * * RUNTIME (supervisor NOT involved in message routing): * Worker A ──UDS──► Worker B (direct, length-prefixed JSON) * Worker B ──UDS──► Worker A * * Supervisor only handles: * - Socket path registry distribution (one-time + on new workers) * - Health checks (periodic pull, not per-message) * - Worker lifecycle (restart, scale) */ export declare class DirectMessageBus extends EventEmitter { workers: Map; socketRegistry: Map; _registeredWorkerIds: Set; _workerErrorHandlers: WeakMap void>; _connections: Map>; _socketDir: string; _broadcastTimer: ReturnType | null; constructor(); get socketDir(): string; /** * Register a worker. Tell it to start a UDS server, then * broadcast the updated registry to all workers. */ registerWorker(serviceName: string, worker: WorkerLike, mode?: string): void; _scheduleBroadcast(): void; _sendInitSocket(worker: WorkerLike, serviceName: string, workerId: number): void; _sendRegistryTo(worker: WorkerLike): void; _broadcastRegistry(): void; unregisterWorker(serviceName: string, workerId: number, options?: UnregisterOptions): void; unregisterService(serviceName: string): void; requestHealthChecks(): void; stats(): Record; /** * CR-IPC-9: Get the connection matrix for debugging. */ getConnectionMatrix(): Record; cleanup(): void; } export {}; //# sourceMappingURL=DirectMessageBus.d.ts.map