import type WebSocket from 'ws'; import type { ClassType } from '../types'; import type { IByteStreamable } from './SrpcByteStream'; export type RequestKeys = keyof T & `${string}Request`; export type ResponseKeys = keyof T & `${string}Response`; export type RequestPrefix = K extends `${infer P}Request` ? P : never; export type ResponsePrefix = K extends `${infer P}Response` ? P : never; type ExtractPrefix = K extends `${infer P}Request` ? (`${P}Response` extends keyof TRes ? P : never) : never; export type InvokePrefixes = ExtractPrefix; export type RequestData = `${P}Request` extends keyof TReq ? NonNullable : never; export type ResponseData = `${P}Response` extends keyof TRes ? NonNullable : never; export type HandlerRequestData = `${P}Request` extends keyof TReq ? NonNullable : never; export type SrpcMeta = object; export interface BaseMessage { requestId?: string; reply?: boolean; error?: string; userError?: boolean; trace?: { traceId: string; spanId: string; traceFlags: number; }; pingPong?: object; byteStreamOperation?: { streamId: number; write?: { chunk: Uint8Array; }; finish?: object; destroy?: { error?: string; }; }; } /** Ensures untrusted envelope tracing data is safe to pass to OpenTelemetry. */ export declare function isValidSrpcTrace(trace: BaseMessage['trace']): trace is NonNullable; export interface SrpcMessageFns { encode(message: T, writer?: unknown): { finish(): Uint8Array; } | Uint8Array; decode(input: Uint8Array, length?: number): T; } export type SrpcDisconnectCause = 'disconnect' | 'conflict' | 'supersede' | 'timeout' | 'badArg'; export interface IQueuedRequest { exp: number; resolve: (value: unknown) => void; reject: (err: unknown) => void; } export declare class SrpcError extends Error { isUserError?: boolean | undefined; constructor(message: string, isUserError?: boolean | undefined); } /** Preserve the explicit sRPC error contract without promoting ordinary errors. */ export declare function serializeSrpcError(error: unknown): { error: string; userError?: boolean; }; export interface ISrpcLogger { info(...messages: unknown[]): void; warn(...messages: unknown[]): void; error(...messages: unknown[]): void; debug(...messages: unknown[]): void; } /** * Controls per-message sRPC traffic logs. `true` logs envelope/message types; * set `bodies` to include the decoded message body. */ export interface SrpcTrafficLoggingOptions { bodies?: boolean; } export type SrpcTrafficLogging = boolean | SrpcTrafficLoggingOptions; /** Returns the application-level message fields carried by an sRPC envelope. */ export declare function srpcMessageTypes(message: BaseMessage): string[]; export interface ISrpcServerOptions { logger: ISrpcLogger; clientMessage: SrpcMessageFns; serverMessage: SrpcMessageFns; wsPath: string; /** * Protocol version assigned to a handshake that omits `pv` (or legacy `_v`). Leave unset to * require an explicit transport version. Use `1` only where legacy * same-client replacement semantics are intentional. */ defaultUnspecifiedProtocolVersion?: 1 | 2 | 3; logTraffic?: SrpcTrafficLogging; httpServer?: import('node:http').Server; /** How long replies for locally abandoned requests are ignored. Defaults to 60 seconds. */ lateReplyTombstoneTtlMs?: number; /** Maximum client requests buffered before a stream is activated. */ maxPendingClientRequests?: number; /** Maximum decoded client-request bytes buffered before a stream is activated. */ maxPendingClientRequestBytes?: number; /** Maximum concurrent client request handlers for one stream. */ maxInFlightClientRequests?: number; /** Maximum decoded client-request bytes executing concurrently for one stream. */ maxInFlightClientRequestBytes?: number; /** Maximum queued WebSocket bytes per stream before the stream is closed. */ maxBufferedBytes?: number; /** Maximum encoded size of one incoming client WebSocket message. */ maxMessageBytes?: number; /** Maximum pending server-to-client RPCs per stream. */ maxPendingServerRequests?: number; /** Maximum encoded pending server-to-client RPC bytes per stream. */ maxPendingServerRequestBytes?: number; /** Maximum WebSocket authentication handshakes awaiting authorization. */ maxPendingHandshakes?: number; /** Maximum live streams, including streams still activating. */ maxActiveStreams?: number; /** Maximum UTF-8 byte length of a client ID. */ maxClientIdBytes?: number; /** Maximum JSON-encoded byte length of merged query/authorization metadata. */ maxClientMetadataBytes?: number; /** Maximum principals retained by the bounded local authentication replay cache. */ maxAuthReplayPrincipals?: number; /** Optional audience expected in v2 credentials. Defaults to `wsPath`. */ authAudience?: string; } export interface SrpcStream extends IByteStreamable { $ws: WebSocket; $queue: Map; readonly id: string; readonly clientStreamId: string; readonly address: string; readonly clientId: string; readonly appVersion: string; readonly configureTs: number; readonly protocolVersion: 1 | 2 | 3; /** Optional client capabilities negotiated during the WebSocket upgrade. */ readonly capabilities?: ReadonlySet; readonly supersede: boolean; readonly meta: T; readonly connectedAt: number; isActivated: boolean; lastPingAt: number; readonly connected: boolean; close(reason?: string): Promise; } /** * Transport-neutral handle for an sRPC client connection. * * A handle is pinned to one connection generation (`id`). Implementations * must not silently retarget it after the same client reconnects. */ export interface SrpcConnection extends IByteStreamable { readonly id: string; readonly clientId: string; readonly meta: T; readonly connectedAt: number; readonly connected: boolean; close(reason?: string): Promise; } export declare class SrpcClientNotFoundError extends Error { constructor(clientId: string); } export declare class SrpcStaleConnectionError extends Error { constructor(clientId: string); } export declare class SrpcOwnerUnavailableError extends Error { constructor(clientId: string, cause?: unknown); } export declare class SrpcIndeterminateDeliveryError extends Error { constructor(clientId: string, cause?: unknown); } export declare class SrpcMeshProtocolError extends Error { constructor(message: string); } export declare class SrpcMeshAuthenticationError extends Error { constructor(message?: string); } export declare class SrpcBackpressureError extends Error { constructor(message?: string); } export declare class SrpcStreamClosedError extends Error { constructor(message?: string); } export type SrpcMessageHandlerFn = (wrappedStream: C, data: I) => Promise | O; export interface ISrpcMessageHandler { handle: SrpcMessageHandlerFn; } export type TSrpcMessageHandlerClass = ClassType>; export type TSrpcMessageHandlerFnOrClass = SrpcMessageHandlerFn | TSrpcMessageHandlerClass; export declare function isSrpcMessageHandlerClass(handler: TSrpcMessageHandlerFnOrClass): handler is TSrpcMessageHandlerClass; export declare function encodeSrpcMessage(codec: SrpcMessageFns, message: T): Buffer; export {}; //# sourceMappingURL=types.d.ts.map