import { MeshService, type MeshBroadcastMap, type MeshBroadcastOptions, type MeshServiceOptions } from '../mesh'; import { type MeshClientRedisRegistryOptions } from './mesh-client-redis-registry'; import { MeshClientRegistry } from './mesh-client-registry'; import { type MeshClientRegistryBackend, type RegisteredClient } from './types'; /** * Point-to-point MeshClientService operations are transported over an * authenticated mesh link. Redis retains only mesh membership, registry, and * optional broadcast coordination. */ export interface MeshClientRemoteTransport { invokeClient(nodeId: number, request: { clientId: string; connectionId: string; type: string; data: unknown; timeoutMs?: number; deadlineAt?: number; }): Promise; fenceClient(nodeId: number, request: { clientId: string; connectionId: string; reason?: string; timeoutMs?: number; }): Promise; updateClientMetadata(nodeId: number, request: { clientId: string; connectionId: string; metadata: TMeta; }): Promise; } export interface MeshClientServiceOptions { key: string; meshOptions?: MeshServiceOptions; registryBackend?: MeshClientRegistryBackend; /** Limits for the built-in Redis registry. Ignored with registryBackend. */ registryOptions?: MeshClientRedisRegistryOptions; clientInvokeFn: (clientId: string, type: string, data: unknown, timeoutMs?: number, connectionId?: string) => Promise; clientUpdateMetaFn?: (clientId: string, metadata: TMeta) => boolean; clientProjectMetaFn?: (clientId: string, metadata: unknown, connectionId?: string) => { updated: boolean; metadata: TMeta; }; clientApplyMetaFn?: (clientId: string, metadata: unknown, connectionId?: string) => boolean; /** * Bounded fallback for backends without durable orphan delivery. * Each non-empty node cleanup snapshot consumes one item. */ maxPendingOrphanItems?: number; maxPendingOrphanBytes?: number; /** Expiry for finalized in-memory fallback obligations. Defaults to one hour. */ pendingOrphanTtlMs?: number; } export declare class MeshClientService { /** @internal */ readonly mesh: MeshService; private logger; private registry; private backend; private clientInvokeFn; private clientUpdateMetaFn; private clientProjectMetaFn; private clientApplyMetaFn; private remoteTransport?; private nodeCleanedUpCallbacks; private clientSupersededCallbacks; private leaseLostCallbacks; private registryRefreshTimer?; private registryRefreshInFlight?; private orphanRetryTimer?; private readonly pendingOrphanCallbacks; private readonly fallbackDeliveryInFlight; private pendingOrphanBytes; private readonly orphanClaimerId; private orphanDrainInFlight?; private activeDurableOrphan?; private readonly registryRefreshIntervalMs; private readonly maxClientIdBytes; private readonly maxMetadataBytes; private lifecycle; private fallbackCleanupLifecycle; private readonly exactUnregisterObligations; private readonly exactUnregisterClientChains; private ownershipGeneration; private hasStarted; private readonly maxPendingOrphanItems; private readonly maxPendingOrphanBytes; private readonly pendingOrphanTtlMs; /** @internal Overridable by focused ownership-failure tests. */ private ownershipClaimDeadlineMs; /** @internal Overridable by focused ownership-failure tests. */ private ownershipRetryIntervalMs; /** @internal Overridable by focused ownership-failure tests. */ private exactUnregisterDeadlineMs; constructor(options: MeshClientServiceOptions); onNodeClientsOrphaned(cb: (nodeId: number, orphaned: RegisteredClient[]) => void | Promise): void; onClientSuperseded(cb: (clientId: string, connectionId?: string, reason?: string) => boolean | void | Promise): void; /** Fence local stream delivery after this mesh member loses its lease. */ onLeaseLost(cb: (reason?: Error) => void | Promise): void; /** @internal Installed by MeshSrpcServer before the service starts. */ setRemoteTransport(transport: MeshClientRemoteTransport | undefined): void; get instanceId(): number; get clientRegistry(): MeshClientRegistry; /** Whether mesh membership and client ownership operations are active. */ get isRunning(): boolean; getAuthNonceConsumer(): ((principal: string, nonce: string, expiresAt: number) => Promise) | undefined; private running; /** Synchronous shutdown fence that rejects new client admission. */ private admissionFenced; /** Exact registry node ID whose ownership cleanup still must succeed. */ private registryCleanupNodeId?; /** MeshService stop initiated by the synchronous shutdown fence. */ private shutdownFencePromise?; /** * Re-open admission for a new lifecycle after the prior stop has fully * completed. MeshSrpcServer only calls this after any cancelled start's * cleanup barrier has settled. */ prepareStart(): void; /** @internal Reject new ownership admission while preserving graceful drain delivery. */ fenceAdmission(): void; /** @internal Immediately fence admission/delivery and stop mesh membership. */ fenceForShutdown(): void; start(): Promise; private doStart; stop(): Promise; private sweepRegistryOwnership; /** @internal Removes retained registry ownership without stopping MeshService (safe from lease-loss callbacks). */ cleanupRegistryOwnership(): Promise; private doStop; private serializeLifecycle; private stopRegistryTimers; private meshLeaseSafe; private hasDurableOrphanBackend; private cleanupNodeForFallback; private retryOrphanCallbacks; private enqueueFallbackOrphan; private pruneFallbackOrphans; private removeFallbackOrphan; private deliverOrphanCallbacks; private doDeliverOrphanCallbacks; private drainDurableOrphans; /** * Register a client on this node. Returns true if registered, false if * another node owns the client and `allowSupersede` is false (conflict). */ registerClient(clientId: string, metadata: TMeta, allowSupersede: boolean, connectionId: string): Promise; /** * Reserve ownership of a clientId without exposing it for lookup/invoke * until activation completes. */ reserveClient(clientId: string, metadata: TMeta, allowSupersede: boolean, connectionId: string): Promise; private takeOwnership; private createOwnershipClaim; private commitOwnershipClaim; private removeFencedClaimPrevious; private isSameClientGeneration; private isCommittedClaim; private fenceClaimAmbiguity; private fencePreviousOwner; /** * Promote a same-node reservation to an active, discoverable client. */ activateClient(clientId: string, metadata: TMeta, connectionId: string): Promise; unregisterClient(clientId: string, connectionId: string): Promise; private retryExactUnregister; private isOwnershipGenerationSafe; updateClientMetadata(clientId: string, metadata: TMeta, connectionId?: string): Promise; private projectClientMetadata; private applyClientMetadata; registerBroadcastHandler(type: K, handler: (data: TBroadcasts[K], senderInstanceId: number) => void | Promise): void; broadcast(type: K, data: TBroadcasts[K], options?: MeshBroadcastOptions): Promise; invoke(clientId: string, type: string, data: unknown, timeoutMs?: number, connectionId?: string): Promise; disconnectClient(clientId: string, connectionId?: string, reason?: string): Promise; private requireRemoteTransport; private validateClientRecord; private validateClientId; } //# sourceMappingURL=mesh-client-service.d.ts.map