import type { BaseMessage, InvokePrefixes, ISrpcServerOptions, RequestData, ResponseData, SrpcConnection, SrpcDisconnectCause, SrpcMeta, SrpcStream } from '../../srpc/types'; import type { Server } from 'node:http'; import type { MeshBroadcastMap, MeshBroadcastOptions, MeshServiceOptions } from '../mesh'; import { SrpcServer } from '../../srpc/SrpcServer'; import { MeshClientRegistry } from './mesh-client-registry'; import type { MeshClientRedisRegistryOptions } from './mesh-client-redis-registry'; import type { MeshSrpcConnection } from './srpc-registry-metadata'; import { type MeshClientRegistryBackend, type RegisteredClient } from './types'; export interface MeshSrpcServerOptions { meshKey: string; meshOptions?: MeshServiceOptions; /** * Register meshStart()/meshStop() with the current App lifecycle. * Defaults to true. Disable this when the application needs bounded or * custom mesh startup/shutdown handling. */ autoLifecycle?: boolean; registryBackend?: MeshClientRegistryBackend; /** Limits for the built-in Redis registry. Ignored with registryBackend. */ registryOptions?: MeshClientRedisRegistryOptions; extractRegistryMetadata?: (stream: SrpcStream) => TRegistryMeta; meshLink?: { advertiseUrl?: string; path?: string; secret?: string; /** * Explicit mesh-link listener. When omitted, mesh upgrades use the * running application HTTP listener, independent of the sRPC client * listener. Before an application listener exists, TSF preserves the * legacy fallback to the sRPC server's supplied httpServer. */ httpServer?: Server; connectTimeoutMs?: number; requestTimeoutMs?: number; idleTimeoutMs?: number; maxFrameBytes?: number; maxBufferedBytes?: number; maxEndpointPins?: number; }; } export declare class MeshSrpcServer extends SrpcServer { private meshClientService; private meshLogger; private extractRegistryMetadataFn?; private readonly meshKey; private readonly meshLinkOptions; private meshLinkRuntime?; private meshLinkController?; private unregisterMeshLinkRoute?; private meshStartPromise?; private meshStopPromise?; private meshFatalCleanupPromise?; private meshLeaseCleanupPromise?; private meshLeaseCleanupRequired; private meshLeaseCleanupRetryTimer?; private meshLeaseCleanupRetryMs; /** Retry a fresh, fenced mesh generation after Redis becomes available again. */ private meshRecoveryTimer?; private meshRecoveryPromise?; private meshRecoveryGeneration; private meshRecoveryAttempt; private meshRecoveryEnabled; /** @internal Overridable by focused recovery tests. */ private meshRecoveryRetryMs; /** @internal Overridable by focused recovery tests. */ private meshRecoveryMaxRetryMs; /** * A start that shutdown abandoned. New starts wait for its eventual * rollback so a late mesh start cannot overlap a replacement lifecycle. */ private meshPendingStartCleanup?; /** A detached controller close from a cancelled, not-yet-ready start. */ private meshPendingLinkClose?; private meshCleanupFailure?; private meshStartGeneration; private meshStartCancelled; private meshStartRollbackFailure?; private meshRunning; private meshLeaseFailure?; private meshStopping; private meshClosed; private meshLinkRequestTimeoutMs; private readonly unregisterLifecycleHandlers; private readonly announcedSenderIds; private readonly pendingMeshInvocationRequestIds; private readonly failedMeshInvocationTombstones; private failedMeshInvocationOverflowUntil; private connectedCallbacks; private disconnectedCallbacks; private orphanedCallbacks; private clientRegistryMetadata; private lifecycleConnectedStreams; /** * Lifecycle disconnects captured before a cancelled startup fully fences * mesh membership. Their callbacks cannot run until rollback has released * the abandoned membership, otherwise consumers can observe a false * offline transition while ownership is still indeterminate. */ private pendingStartDisconnects; private pendingStartDisconnectFenceDeadlines; private pendingStartDisconnectCallbackQueued; private rollbackGatedDisconnects; private meshSupersedeReconcileMs; private meshSupersedeReconcileRetryMs; private clientRegistryChains; private clientCallbackChains; private pendingSyncs; constructor(options: ISrpcServerOptions & MeshSrpcServerOptions); private handleClientSuperseded; /** Fence streams and release the link route after the mesh lease is lost. */ private beginMeshLeaseCleanup; private scheduleMeshLeaseCleanupRetry; private clearMeshLeaseCleanupRetry; /** * Lease loss is a hard split-brain boundary, but not a permanent process * failure. Once the exact old ownership is removed, start a new mesh * generation. Existing streams were already fenced and must reconnect. */ private scheduleMeshRecovery; private attemptMeshRecovery; private canRecoverMesh; private cancelMeshRecovery; private handleMeshLeaseLost; /** * A mesh client is never admitted until this process can reserve its * ownership in the shared registry. This centralizes the readiness gate * rather than requiring each application authorizer to remember it. */ protected beforeClientAdmission(): Promise; /** * Defers stream activation until mesh reservation succeeds. * Installs meta proxy and reserves the client atomically in Redis * using the protocol-version supersession rules, all serialized in the * per-client registry chain. Protocol v1 retains its implicit replacement * behavior; protocol v2+ requires the explicit supersede flag. Reserved * clients remain hidden from lookup/invoke until onStreamActivated promotes * them to active. * * If registration returns a conflict, cleanupStream is called * (which fires onStreamDisconnected and drains the queue) so * no connection handlers or RPCs ever run on the rejected stream. */ protected postEstablishCheck(stream: SrpcStream): Promise; private extractRegistryMetadata; private shouldAllowSupersede; private static readonly PROXIED; /** * Install a Proxy on stream.meta that schedules a microtask-debounced * sync to Redis whenever any property is mutated. * * This means handler code, connection handlers, and external code * (e.g. FreeSwitch controller) can all mutate stream.meta directly * and the mesh registry stays in sync - no manual sync calls needed. * * **Limitation:** Only top-level property mutations are tracked. * Nested mutations (e.g. `stream.meta.user.name = 'Bob'`) do NOT * trigger a sync. For nested metadata, either reassign the top-level * property (`stream.meta.user = { ...stream.meta.user, name: 'Bob' }`) * or call `updateClientMetadata()` explicitly. */ private installMetaProxy; /** * Schedule a microtask-debounced sync for a client. * Multiple synchronous mutations are batched into a single sync. */ private scheduleSyncStreamMeta; protected onStreamConnected(stream: SrpcStream): Promise; protected onStreamWillActivate(stream: SrpcStream): Promise; protected onStreamActivated(stream: SrpcStream): Promise; protected onStreamDisconnected(stream: SrpcStream, cause: SrpcDisconnectCause): void; /** * Sync the current stream.meta to the mesh registry. * Called automatically by the meta proxy's microtask debounce. * Routed through enqueueClientRegistry so updates are serialized * after initial registration (prevents lost updates if registration * hasn't completed yet). */ private syncStreamMeta; private enqueueClientRegistry; private enqueueClientCallback; get meshInstanceId(): number; get clientRegistry(): MeshClientRegistry; get startupState(): 'stopped' | 'starting' | 'ready' | 'draining' | 'failed'; ready(): Promise; /** * Read a generation-fenced registry record without materializing a remote * connection, link capability, or byte-stream sender pool. Use this for * connection-state and metadata decisions; resolveClient() is for an imminent * invoke, byte stream, metadata mutation, or disconnect. */ getRegisteredClient(clientId: string): Promise | undefined>; /** Registry-only counterpart to listClients(); it never creates remote handles. */ listRegisteredClients(): Promise[]>; /** * Update metadata for a client, regardless of which node owns it. * Routes through the mesh to the owning node so that stream.meta * reflects the change immediately and the proxy auto-syncs to Redis. * For local streams, you can also mutate stream.meta directly. */ updateClientMetadata(connectionOrClientId: SrpcConnection | string, metadata: Partial): Promise; onClientConnected(handler: (clientId: string, metadata: TRegistryMeta) => void | Promise): void; onClientDisconnected(handler: (clientId: string, metadata: TRegistryMeta) => void | Promise): void; onNodeClientsOrphaned(handler: (nodeId: number, clients: RegisteredClient[]) => void | Promise): void; registerBroadcastHandler(type: K, handler: (data: TBroadcasts[K], senderInstanceId: number) => void | Promise): void; broadcast(type: K, data: TBroadcasts[K], options?: MeshBroadcastOptions): Promise; resolveClient(clientId: string, deadlineAt?: number): Promise | undefined>; listClients(): Promise[]>; disconnectClient(connectionOrClientId: SrpcConnection | string, reason?: string): Promise; /** * Invoke a client method across any node in the mesh. * Overloaded: when called with a stream, delegates to SrpcServer.invoke. * When called with a clientId string, routes through the mesh. */ invoke

>(connectionOrClientId: MeshSrpcConnection | string, prefix: P, data: RequestData, timeoutMs?: number): Promise>; invoke(connectionOrClientId: SrpcConnection | string, prefix: string, data: unknown, timeoutMs?: number): Promise; meshStart(): Promise; private startMeshLifecycle; private assertMeshStartCurrent; private startMesh; meshStop(): Promise; private stopMesh; private disconnectAllMeshStreams; private capturePendingStartDisconnects; private captureDeferredDisconnect; private dispatchPendingStartDisconnects; private shouldDispatchPendingStartDisconnect; private queuePendingStartDisconnectCallback; private markPendingStartDisconnectFence; private failDeferredDisconnectReconciliation; private releaseMeshLinkResources; private trackDetachedLinkClose; private deferPendingMeshStartRollback; private rollbackMeshStart; close(): void; protected handleByteSubstreamOperation(stream: SrpcStream, operation: NonNullable, requestId?: string): void; private rememberFailedMeshInvocation; private pruneFailedMeshInvocationTombstones; private destroyFailedMeshInvocationSender; protected hasExternalByteStreamSender(stream: SrpcStream, streamId: number): boolean; private createMeshLinkController; /** @internal Isolates the core invoke boundary for focused failure-path verification. */ private invokeMeshClientWithRequestId; private applyMetadataToLocalStream; private assertCurrentMeshStream; private projectMetadataWithoutMutation; private resolveMeshLinkConfig; /** * Normal applications start mesh links after their main HTTP listener is * bound, so omitting meshLink.httpServer intentionally hosts upgrades * there. Standalone servers historically supplied only options.httpServer * and may have no live application listener; retain that safe fallback so * their advertised listener and upgrade listener cannot diverge. */ private resolveMeshLinkHttpServer; } //# sourceMappingURL=mesh-srpc-server.d.ts.map