import type { IncomingMessage } from 'node:http'; import { BaseMessage, HandlerRequestData, InvokePrefixes, ISrpcServerOptions, RequestData, ResponseData, SrpcDisconnectCause, SrpcConnection, SrpcMeta, SrpcStream, TSrpcMessageHandlerFnOrClass } from './types'; export declare class SrpcServer { protected readonly options: ISrpcServerOptions; private readonly logger; private readonly wsServer; private readonly streamConnectionHandlers; private readonly streamDisconnectionHandlers; private readonly streamMessageHandlers; private readonly broadcastHandlers; private readonly blockedClientRequests; private readonly pendingClientRequests; private readonly pendingClientRequestBytes; private readonly inFlightClientRequests; private readonly inFlightClientRequestBytes; private readonly pendingServerRequestBytes; private readonly lateReplyTombstonesByStream; private readonly backpressuredByteStreams; private readonly backpressuredByteStreamBytes; private readonly publishedStreams; private readonly authReplayNoncesByPrincipal; private authNonceConsumer?; private readonly cleanupUpgradeHandler?; private readonly inactivityCheckInterval; private pendingHandshakeCount; private clientAuthorizer?; private clientKeyFetcher?; readonly streamsById: Map>; readonly streamsByClientId: Map>; protected readonly pendingStreamsByClientId: Map>; constructor(options: ISrpcServerOptions); private readonly verifyClient; private attachConnection; private createStream; private handleStreamEstablished; protected postEstablishCheck(_stream: SrpcStream): Promise; /** * Runs after client authentication succeeds but before an sRPC stream is * created. Subclasses may use this to fence admission on external state. */ protected beforeClientAdmission(): Promise; /** * Hook for private/cluster CAS work after the initial protocol frame is * queued but before user connection handlers and public local publication. */ protected onStreamWillActivate(_stream: SrpcStream): void | Promise; /** * Allows a transport extension to recognize a locally-owned sender whose * state is maintained outside the core SrpcByteStream sender map. * Core remains strict by default. */ protected hasExternalByteStreamSender(_stream: SrpcStream, _streamId: number): boolean; protected getCurrentStreamByClientId(clientId: string): SrpcStream | undefined; protected isCurrentStream(stream: SrpcStream): boolean; private activateStream; private handleWsMessage; private handleStreamDataReceived; private openClientRequests; protected handleByteSubstreamOperation(stream: SrpcStream, op: NonNullable, requestId?: string): void; private clearByteStreamBackpressureFlag; private clearByteStreamState; private updateByteStreamBufferedBytes; private totalBackpressuredByteStreamBytes; private handleClientRequest; protected runMessageHandler(handler: TSrpcMessageHandlerFnOrClass, unknown, unknown>, stream: SrpcStream, data: unknown): Promise; protected onStreamConnected(stream: SrpcStream): Promise; protected onStreamActivated(_stream: SrpcStream): void | Promise; protected onStreamDisconnected(stream: SrpcStream, cause: SrpcDisconnectCause): void; private revokeStream; protected cleanupStream(stream: SrpcStream, cause?: SrpcDisconnectCause, remoteClose?: { code?: number; reason?: string; }): void; private closeStreamGracefully; private handleStreamDisconnected; private handleStreamError; private terminateInactiveStreams; private validateClientAuth; private consumeAuthReplayToken; private fetchClientKey; private getRemoteAddress; protected writeToStream(stream: SrpcStream, data: TServerOutput): boolean; private writeEncodedToStream; protected writeToStreamAsync(stream: SrpcStream, data: TServerOutput): Promise; private logTraffic; private canWriteToStream; private closeStreamWithError; private isStreamDispatchAvailable; setClientAuthorizer(authorizer: (metadata: Record, req: IncomingMessage) => Promise> | boolean | Partial): void; setClientKeyFetcher(fetcher: (clientId: string) => Promise | false | string): void; /** * Installs an atomic service-wide authentication replay-token consumer. It * must return true only for the first `(principal, token)` consumption * through `expiresAt`; false rejects a replay. Auth-v2 uses its nonce and * legacy auth-v1 uses its signed stream ID. Core retains its fair local guard. */ setAuthNonceConsumer(consumer: (principal: string, nonce: string, expiresAt: number) => boolean | Promise): void; registerConnectionHandler(handler: (stream: SrpcStream) => void | Promise): void; registerMessageHandler

>(prefix: P, handler: TSrpcMessageHandlerFnOrClass, HandlerRequestData, ResponseData>): void; registerDisconnectHandler(handler: (stream: SrpcStream, cause: SrpcDisconnectCause) => void): void; registerBroadcastHandler(type: string, handler: (data: any, senderInstanceId: number) => void | Promise): void; broadcast(type: string, data: any, options?: { skipSelf?: boolean; }): Promise; resolveClient(clientId: string): Promise | undefined>; listClients(): Promise[]>; /** List streams physically connected to this server process. */ getLocalStreams(): Promise[]>; disconnectClient(streamOrClientId: SrpcStream | string, reason?: string): Promise; updateClientMetadata(streamOrClientId: SrpcStream | string, metadata: Partial): Promise; invoke

>(stream: SrpcStream, prefix: P, data: RequestData, timeoutMs?: number): Promise>; protected invokeWithRequestId

>(stream: SrpcStream, prefix: P, data: RequestData, timeoutMs: number, requestId: string): Promise>; protected encodeMeshInvokeRequest(prefix: string, data: unknown): Uint8Array; protected decodeMeshInvokeRequest(prefix: string, data: Uint8Array): unknown; protected encodeMeshInvokeResponse(prefix: string, data: unknown): Uint8Array; protected decodeMeshInvokeResponse(prefix: string, data: Uint8Array): unknown; protected reserveByteStreamSenderIds(stream: SrpcStream, count: number): number[]; protected writeByteStreamOperation(stream: SrpcStream, operation: NonNullable): Promise; private addLateReplyTombstone; private isLateReply; private pruneLateReplyTombstones; private get lateReplyTombstoneTtlMs(); private get maxPendingClientRequests(); private get maxPendingClientRequestBytes(); private get maxInFlightClientRequests(); private get maxInFlightClientRequestBytes(); private get maxBufferedBytes(); private get maxMessageBytes(); private get maxPendingServerRequests(); private get maxPendingServerRequestBytes(); private get maxPendingHandshakes(); private get maxActiveStreams(); private get maxClientIdBytes(); private get maxClientMetadataBytes(); private get maxAuthReplayPrincipals(); private validRemoteByteStreamOperation; private isLocalAuthNonceReplay; private consumeLocalAuthNonce; close(): void; static createInvoke(instanceFn: () => SrpcServer): SrpcServer['invoke']; } //# sourceMappingURL=SrpcServer.d.ts.map