import net from "node:net"; interface ForgeMessage { type: string; from?: string; fromWorkerId?: number; payload?: unknown; requestId?: string; hmac?: string; error?: { message: string; code?: string; statusCode?: number; } | string | null; [key: string]: unknown; } interface PendingRequest { resolve: (value: unknown) => void; reject: (reason: Error) => void; timer: ReturnType; socketKey: string | null; settled: boolean; } interface WorkerChannelManagerOptions { channels?: string[]; onChannelDead?: (key: string, attempts: number) => void; } interface SupervisorMessage { type: string; [key: string]: unknown; } interface SocketRegistry { [key: string]: string; } interface TopologyEntry { connections: number; keys: string[]; } interface ForgeSocket extends net.Socket { _drainWarned?: boolean; } export declare class WorkerChannelManager { serviceName: string; workerId: number; _dependencies: Set | null; _onChannelDead: ((key: string, attempts: number) => void) | null; _server: net.Server | null; _socketPath: string | null; outbound: Map; inbound: Map; serviceConnections: Map; rrIndex: Map; pendingRequests: Map; requestCounter: number; _workerId: number; onMessage: ((from: string, payload: unknown) => void) | null; onRequest: ((from: string, payload: unknown) => unknown) | null; _supervisorSend: ((msg: SupervisorMessage) => void) | null; _registry: SocketRegistry; _reconnectAttempts: Map; _reconnectTimers: Map>; _socketKeyMap: Map; backpressureEvents: number; _reconnectQueue: Map; constructor(serviceName: string, workerId: number, options?: WorkerChannelManagerOptions); /** * Initialize — set up supervisor IPC listener. */ init(supervisorSend: (msg: SupervisorMessage) => void): void; /** * Start our UDS server so other workers can connect to us. */ _startServer(socketDir: string, serviceName: string, workerId: number): void; /** * Update our knowledge of the socket registry and connect to new peers. */ _updateRegistry(registry: SocketRegistry): void; /** * Establish an outbound connection to another worker's UDS server. */ _connectTo(key: string, socketPath: string): void; /** * Send a length-prefixed JSON frame over a socket. * P6: Buffer.from(json) avoids double-scan of the string. * @returns true if the write was accepted into the kernel buffer */ _sendFrame(socket: ForgeSocket, msg: ForgeMessage): boolean; /** * Send a pre-built frame buffer over a socket (for broadcast optimization). * P6: Serialize once, send to all recipients. */ _sendRawFrame(socket: ForgeSocket, frameBuffer: Buffer): boolean; /** * Build a length-prefixed frame buffer from a message object. * P6: Used by broadcast to serialize once and send to all. */ _buildFrame(msg: ForgeMessage): Buffer; /** * Handle a message from another worker (inbound or outbound socket). */ _handleIncomingMessage(socket: net.Socket, msg: ForgeMessage): void; /** * Send a fire-and-forget message. Direct UDS path if available. */ send(target: string, payload: unknown): void; /** * Broadcast to all workers of a target service. */ broadcast(target: string, payload: unknown): void; /** * Request/response over direct UDS. */ request(target: string, payload: unknown, timeoutMs?: number): Promise; /** * Pick a socket to a target service (round-robin). */ _pickSocket(target: string): ForgeSocket | null; /** * P19: O(1) socket → key lookup via reverse Map (replaces linear scan). */ _socketKey(socket: net.Socket): string; /** * Reject all pending requests that were sent over a specific socket. */ _rejectPendingForSocket(deadKey: string): void; hasDirectConnection(target: string): boolean; topology(): Record; destroy(): Promise; } export {}; //# sourceMappingURL=WorkerChannelManager.d.ts.map