diff --git a/mobile/src/transport/logical-subscription-registry.ts b/mobile/src/transport/logical-subscription-registry.ts deleted file mode 100644 index a696664c6a9..00000000000 --- a/mobile/src/transport/logical-subscription-registry.ts +++ /dev/null @@ -1,121 +0,0 @@ -import type { RpcClient } from './rpc-client' - -type SubscriptionRecord = { - method: string - params: unknown - listener: (result: unknown) => void - options?: Parameters[3] - disposePhysical: (() => void) | null - cancelled: boolean -} - -export type LogicalSubscriptionRegistryContext = { - isClosed: () => boolean - isSuspended: () => boolean - activeSession: () => RpcClient - generation: () => number -} - -/** - * The logical subscriptions of a stable client: each record outlives the physical - * session it is currently attached to, and survives suspend and migration by - * re-attaching rather than by asking the caller to resubscribe. - */ -export class LogicalSubscriptionRegistry { - private readonly records = new Map() - private nextId = 0 - - constructor(private readonly context: LogicalSubscriptionRegistryContext) {} - - add( - method: string, - params: unknown, - listener: (result: unknown) => void, - options?: Parameters[3] - ): () => void { - const id = ++this.nextId - const record: SubscriptionRecord = { - method, - params, - listener, - options, - disposePhysical: null, - cancelled: false - } - this.records.set(id, record) - if (!this.context.isSuspended()) { - this.attach(record, this.context.activeSession(), this.context.generation()) - } - return () => { - if (record.cancelled) { - return - } - record.cancelled = true - record.disposePhysical?.() - record.disposePhysical = null - this.records.delete(id) - } - } - - updateTerminalViewport(terminal: string, viewport: { cols: number; rows: number }): void { - for (const record of this.records.values()) { - if ( - record.params && - typeof record.params === 'object' && - 'terminal' in record.params && - record.params.terminal === terminal - ) { - record.params = { ...record.params, viewport } - } - } - if (!this.context.isSuspended()) { - this.context.activeSession().updateTerminalSubscriptionViewport(terminal, viewport) - } - } - - // Why: attach before disposing the previous physical subscription so no gap opens, - // while callbacks stay fenced until the caller makes nextGeneration current. - replayOnto(nextSession: RpcClient, nextGeneration: number): void { - for (const record of this.records.values()) { - const disposePrevious = record.disposePhysical - this.attach(record, nextSession, nextGeneration) - disposePrevious?.() - } - } - - /** Detach from the physical session but keep the records for a later replay. */ - detachAll(): void { - for (const record of this.records.values()) { - record.disposePhysical?.() - record.disposePhysical = null - } - } - - disposeAll(): void { - for (const record of this.records.values()) { - record.disposePhysical?.() - } - this.records.clear() - } - - private attach( - record: SubscriptionRecord, - session: RpcClient, - subscriptionGeneration: number - ): void { - record.disposePhysical = session.subscribe( - record.method, - record.params, - (result) => { - if ( - !this.context.isClosed() && - !record.cancelled && - this.context.generation() === subscriptionGeneration - ) { - record.listener(result) - } - }, - record.options - ) - } -} diff --git a/mobile/src/transport/stable-logical-rpc-client.ts b/mobile/src/transport/stable-logical-rpc-client.ts index 337fe2d4bd7..fb514128382 100644 --- a/mobile/src/transport/stable-logical-rpc-client.ts +++ b/mobile/src/transport/stable-logical-rpc-client.ts @@ -7,7 +7,6 @@ import { import { waitForAuthenticated } from './replacement-session-authentication' import { projectMobileRpcRequestParams } from './mobile-rpc-request-projection' import { LogicalClientConnectionPath } from './logical-client-connection-path' -import { LogicalSubscriptionRegistry } from './logical-subscription-registry' export type MobileConnectionPath = 'lan' | 'tailscale' | 'relay' @@ -25,6 +24,15 @@ export function isLogicalClientCutoverError(error: unknown): boolean { ) } +type SubscriptionRecord = { + method: string + params: unknown + listener: (result: unknown) => void + options?: Parameters[3] + disposePhysical: (() => void) | null + cancelled: boolean +} + type PendingRequest = { reject: (error: Error) => void } @@ -64,13 +72,9 @@ export function createStableLogicalRpcClient( let generation = 1 let closed = false let suspended = false + let nextSubscriptionId = 0 let activeStateUnsubscribe: (() => void) | null = null - const subscriptions = new LogicalSubscriptionRegistry({ - isClosed: () => closed, - isSuspended: () => suspended, - activeSession: () => activeSession, - generation: () => generation - }) + const subscriptions = new Map() const pendingRequests = new Set() const stateListeners = new Set<(state: ConnectionState) => void>() let state = initialSession.getState() @@ -116,11 +120,44 @@ export function createStableLogicalRpcClient( if (closed) { return () => {} } - return subscriptions.add(method, params, listener, options) + const id = ++nextSubscriptionId + const record: SubscriptionRecord = { + method, + params, + listener, + options, + disposePhysical: null, + cancelled: false + } + subscriptions.set(id, record) + if (!suspended) { + attachSubscription(record, activeSession, generation) + } + return () => { + if (record.cancelled) { + return + } + record.cancelled = true + record.disposePhysical?.() + record.disposePhysical = null + subscriptions.delete(id) + } }, updateTerminalSubscriptionViewport(terminal, viewport) { - subscriptions.updateTerminalViewport(terminal, viewport) + for (const record of subscriptions.values()) { + if ( + record.params && + typeof record.params === 'object' && + 'terminal' in record.params && + record.params.terminal === terminal + ) { + record.params = { ...record.params, viewport } + } + } + if (!suspended) { + activeSession.updateTerminalSubscriptionViewport(terminal, viewport) + } }, getState: () => state, @@ -143,7 +180,10 @@ export function createStableLogicalRpcClient( closed = true activeStateUnsubscribe?.() activeStateUnsubscribe = null - subscriptions.disposeAll() + for (const record of subscriptions.values()) { + record.disposePhysical?.() + } + subscriptions.clear() // Why: let the physical close settle in-flight requests — it knows which // frames were written and marks those delivery-unknown; a blanket local // reject would erase that distinction. @@ -158,7 +198,10 @@ export function createStableLogicalRpcClient( suspended = true activeStateUnsubscribe?.() activeStateUnsubscribe = null - subscriptions.detachAll() + for (const record of subscriptions.values()) { + record.disposePhysical?.() + record.disposePhysical = null + } // Why: let the physical close settle in-flight requests — it knows which // frames were written and marks those delivery-unknown (a suspend can cut // over a half-open relay whose sends may already be delivered). @@ -211,7 +254,11 @@ export function createStableLogicalRpcClient( // Why: replay on the authenticated replacement before closing the old // session, but fence callbacks until the generation becomes current. - subscriptions.replayOnto(nextSession, nextGeneration) + for (const record of subscriptions.values()) { + const disposePrevious = record.disposePhysical + attachSubscription(record, nextSession, nextGeneration) + disposePrevious?.() + } generation = nextGeneration activeSession = nextSession activePath = path @@ -258,6 +305,23 @@ export function createStableLogicalRpcClient( } } + function attachSubscription( + record: SubscriptionRecord, + session: RpcClient, + subscriptionGeneration: number + ): void { + record.disposePhysical = session.subscribe( + record.method, + record.params, + (result) => { + if (!closed && !record.cancelled && generation === subscriptionGeneration) { + record.listener(result) + } + }, + record.options + ) + } + function bindActiveState(session: RpcClient, sessionGeneration: number): void { activeStateUnsubscribe = session.onStateChange((next) => { if (!closed && generation === sessionGeneration && session === activeSession) {