diff --git a/config/max-lines-baseline.txt b/config/max-lines-baseline.txt index 384318e9ed6..468845ca6ab 100644 --- a/config/max-lines-baseline.txt +++ b/config/max-lines-baseline.txt @@ -14,4 +14,3 @@ mobile-config app/h/*/source-control/*.tsx mobile-config app/h/*/tasks.tsx mobile-config app/index.tsx mobile-config scripts/mock-server.ts -mobile-config src/transport/rpc-client.ts diff --git a/mobile/.oxlintrc.json b/mobile/.oxlintrc.json index d58682609a3..0e7a2e18cce 100644 --- a/mobile/.oxlintrc.json +++ b/mobile/.oxlintrc.json @@ -33,12 +33,6 @@ "max-lines": ["error", { "max": 1422, "skipBlankLines": true, "skipComments": true }] } }, - { - "files": ["src/transport/rpc-client.ts"], - "rules": { - "max-lines": ["error", { "max": 1074, "skipBlankLines": true, "skipComments": true }] - } - }, { "files": ["scripts/mock-server.ts"], "rules": { diff --git a/mobile/src/browser/mobile-browser-rpc-client.ts b/mobile/src/browser/mobile-browser-rpc-client.ts index a006fad9aa0..492adfa8c60 100644 --- a/mobile/src/browser/mobile-browser-rpc-client.ts +++ b/mobile/src/browser/mobile-browser-rpc-client.ts @@ -1,6 +1,6 @@ -import type { RpcClient, RpcClientSendRequestOptions } from '../transport/rpc-client' +import type { RpcClient, SendRequestOptions } from '../transport/rpc-client' +import type { RpcStreamSubscribeOptions } from '../transport/rpc-client-stream-registry' import type { RpcResponse } from '../transport/types' -import type { RpcClientSubscribeOptions } from '../transport/rpc-client-subscribe-options' import type { HostSessionBrowserOperations } from '../session/host-session-browser-operations' export type MobileBrowserRpcClient = Pick @@ -46,7 +46,7 @@ export function createMobileBrowserRpcClient( async function sendRequest( method: string, rawParams: unknown = {}, - _options?: RpcClientSendRequestOptions + _options?: SendRequestOptions ): Promise { const params = (rawParams ?? {}) as BrowserParams try { @@ -113,7 +113,7 @@ export function createMobileBrowserRpcClient( method: string, rawParams: unknown, onData: (result: unknown) => void, - options?: RpcClientSubscribeOptions + options?: RpcStreamSubscribeOptions ): () => void { if (method !== 'browser.screencast') { onData({ type: 'error', message: `Unsupported browser subscription: ${method}` }) diff --git a/mobile/src/transport/direct-rpc-client.ts b/mobile/src/transport/direct-rpc-client.ts new file mode 100644 index 00000000000..c92654e8c82 --- /dev/null +++ b/mobile/src/transport/direct-rpc-client.ts @@ -0,0 +1,309 @@ +import type { ConnectOptions, RpcClient, SendRequestOptions } from './rpc-client' +import { DirectConnectionLog } from './direct-connection-log' +import { RpcClientAuthenticationRetry } from './rpc-client-authentication-retry' +import { RpcClientConnectionState } from './rpc-client-connection-state' +import { + RpcClientReconnectSchedule, + RPC_RECONNECT_ATTEMPT_LIMIT +} from './rpc-client-reconnect-schedule' +import { RpcClientRequestTracker } from './rpc-client-request-tracker' +import { RpcClientSocketCloseController } from './rpc-client-socket-close-controller' +import { RpcClientSocketFactory } from './rpc-client-socket-factory' +import type { RpcClientSocketSession } from './rpc-client-socket-session' +import { + RpcClientStreamRegistry, + type RpcStreamingListener, + type RpcStreamSubscribeOptions +} from './rpc-client-stream-registry' +import { nudgeRpcClientForeground } from './rpc-client-foreground-nudge' +import { isLivenessProbeResponseId, RpcClientLivenessSession } from './rpc-client-liveness-session' +import type { TerminalStreamFrame } from './terminal-stream-protocol' +import type { ConnectionState, ForegroundNudgeReason, RpcResponse } from './types' + +export class DirectRpcClient implements RpcClient { + private socketSession: RpcClientSocketSession | null = null + private readonly connectionState: RpcClientConnectionState + private readonly reconnect: RpcClientReconnectSchedule + private readonly requests: RpcClientRequestTracker + private readonly streams: RpcClientStreamRegistry + private readonly liveness: RpcClientLivenessSession + private readonly socketFactory: RpcClientSocketFactory + private readonly authenticationRetry: RpcClientAuthenticationRetry + private readonly socketClose: RpcClientSocketCloseController + private requestCounter = 0 + private readonly connectionLog: DirectConnectionLog + private intentionallyClosed = false + private authenticationGeneration = 0 + + constructor( + private readonly endpoint: string, + private readonly deviceToken: string, + serverPublicKeyB64: string, + private readonly options: ConnectOptions + ) { + this.connectionLog = new DirectConnectionLog(endpoint, options.onLog) + this.reconnect = new RpcClientReconnectSchedule({ + openConnection: () => this.openConnection(), + rejectConnectWaiters: (reason) => this.connectionState.rejectWaiters(reason), + emitLog: (message, detail) => + this.connectionLog.emit('info', message, detail, { code: 'retry-scheduled' }) + }) + this.connectionState = new RpcClientConnectionState({ + endpoint, + initialListener: options.onStateChange, + getReconnectAttempt: () => this.reconnect.getAttempt(), + isClosed: () => this.intentionallyClosed + }) + this.streams = new RpcClientStreamRegistry({ + nextId: () => this.nextId(), + deviceToken, + getState: () => this.connectionState.get(), + sendEncrypted: (request) => this.sendEncrypted(request) + }) + this.requests = new RpcClientRequestTracker({ + nextId: () => this.nextId(), + deviceToken, + getState: () => this.connectionState.get(), + waitForConnected: (timeoutMs) => this.waitForConnected(timeoutMs), + sendEncrypted: (request) => this.sendEncrypted(request) + }) + this.liveness = new RpcClientLivenessSession({ + sendProbe: (probeId) => this.sendLivenessProbe(probeId), + terminate: (session) => { + if (this.socketSession === session) { + this.socketClose.forceClose(session) + } + }, + onTimeout: this.connectionLog.livenessTimeout, + nextId: () => this.nextId() + }) + this.socketFactory = new RpcClientSocketFactory({ + endpoint, + deviceToken, + serverPublicKeyB64, + getCurrentSocket: () => this.socketSession?.socket ?? null, + getState: () => this.getState(), + getReconnectAttempt: () => this.getReconnectAttempt(), + getLastConnectedAt: () => this.getLastConnectedAt(), + isIntentionallyClosed: () => this.intentionallyClosed, + emitLog: this.connectionLog.emit, + onHandshakeStarted: () => this.connectionState.publish('handshaking'), + onAuthenticated: (session) => this.handleAuthenticated(session), + onAuthRejected: (reason) => this.authenticationRetry.reject(reason), + onRpcResponse: (response) => this.handleRpcResponse(response), + onBinary: (bytes) => this.streams.handleBinary(bytes), + onAuthenticatedInbound: (session) => this.liveness.noteInbound(session), + onClosed: (session, closeCode) => this.socketClose.handle(session, closeCode), + onForcedClose: (session) => this.socketClose.forceClose(session) + }) + this.authenticationRetry = new RpcClientAuthenticationRetry({ + endpoint, + stopLiveness: () => this.liveness.stop(), + emitWarning: (message, detail) => + this.connectionLog.emit('warn', message, detail, { code: 'authentication-rejected' }), + retry: (reason) => this.retryAuthentication(reason), + latchFailure: (reason) => this.latchAuthenticationFailure(reason) + }) + this.socketClose = new RpcClientSocketCloseController({ + connectionState: this.connectionState, + reconnect: this.reconnect, + requests: this.requests, + streams: this.streams, + socketFactory: this.socketFactory, + authenticationRetry: this.authenticationRetry, + getCurrentSession: () => this.socketSession, + clearCurrentSession: () => (this.socketSession = null), + getAuthenticationGeneration: () => this.authenticationGeneration, + isIntentionallyClosed: () => this.intentionallyClosed, + stopLiveness: (session) => this.liveness.stop(session), + emitWarning: (message, detail, evidence) => + this.connectionLog.emit('warn', message, detail, evidence) + }) + this.openConnection() + } + + sendRequest( + method: string, + params?: unknown, + options?: SendRequestOptions + ): Promise { + return this.requests.sendRequest(method, params, options) + } + + subscribe( + method: string, + params: unknown, + onData: RpcStreamingListener, + options?: RpcStreamSubscribeOptions + ): () => void { + return this.streams.subscribe(method, params, onData, options) + } + + updateTerminalSubscriptionViewport( + terminal: string, + viewport: { cols: number; rows: number } + ): void { + this.streams.updateTerminalViewport(terminal, viewport) + } + + sendTerminalBinaryFrame(frame: TerminalStreamFrame): boolean { + return this.socketSession?.sendTerminalBinaryFrame(frame) ?? false + } + + getState(): ConnectionState { + return this.connectionState.get() + } + + getReconnectAttempt(): number { + return this.reconnect.getAttempt() + } + + getLastConnectedAt(): number | null { + return this.connectionState.getLastConnectedAt() + } + + getLastInboundAt = (): number | null => this.liveness.getLastInboundAt() + + onStateChange(listener: (state: ConnectionState) => void): () => void { + return this.connectionState.addListener(listener) + } + + notifyForeground(_reason?: ForegroundNudgeReason): void { + if (this.intentionallyClosed) { + return + } + nudgeRpcClientForeground({ + getState: () => this.getState(), + getReconnectAttempt: () => this.getReconnectAttempt(), + dialAgeMs: Date.now() - this.socketFactory.getDialStartedAt(), + hasReconnectTimer: () => this.reconnect.hasTimer(), + probeLiveness: () => this.liveness.probeNow(), + abandonDial: () => { + const dialing = this.socketSession + if (!dialing) { + return false + } + this.socketClose.forceClose(dialing) + return true + }, + redial: (keepTimer) => this.reconnect.redialNow(keepTimer) + }) + } + + close(): void { + this.intentionallyClosed = true + this.reconnect.cancel() + const session = this.socketSession + session?.dispose() + session?.clearTimers() + this.liveness.stop() + session?.close() + this.socketSession = null + session?.clearKey() + this.connectionState.publish('disconnected') + this.requests.rejectAll('Client closed', { deliveryUnknown: true }) + } + + private openConnection(): void { + if (this.intentionallyClosed) { + return + } + // Why: retired socket buffers still hold process memory; dialing again would stack another. + if (!this.socketFactory.canOpen()) { + this.connectionLog.emit( + 'error', + 'WebSocket reconnect deferred', + 'Retired socket buffers are still draining' + ) + this.connectionState.publish('reconnecting') + this.reconnect.schedule() + return + } + this.connectionState.publish('connecting') + this.socketSession = this.socketFactory.open() + } + + private handleAuthenticated(session: RpcClientSocketSession): void { + console.log('[net] e2ee_authenticated — connected', { streamCount: this.streams.size() }) + this.liveness.start(session) + this.authenticationGeneration++ + this.reconnect.authenticated() + this.authenticationRetry.accepted() + this.connectionState.publish('connected') + this.connectionLog.emit('success', 'Authenticated', 'Channel ready for RPC', { + code: 'direct-connected' + }) + this.streams.replayAfterAuthentication() + } + + private handleRpcResponse(response: RpcResponse): void { + if (isLivenessProbeResponseId(response.id)) { + return + } + if (!response.ok && response.error.code === 'unauthorized') { + this.authenticationRetry.reject('Unauthorized — pairing may be revoked') + return + } + if (!this.streams.handleResponse(response)) { + this.requests.resolve(response) + } + } + + private retryAuthentication(reason: string): void { + const closing = this.socketSession + closing?.dispose() + this.socketSession = null + closing?.clearKey() + this.streams.markForReplay() + this.requests.rejectAll(reason) + closing?.close() + this.connectionState.publish('reconnecting') + this.reconnect.schedule() + } + + private latchAuthenticationFailure(reason: string): void { + this.intentionallyClosed = true + this.socketSession?.dispose() + this.socketSession?.close() + this.socketSession = null + this.connectionState.publish('auth-failed') + this.requests.rejectAll(reason) + } + + private sendEncrypted(request: unknown): boolean { + if (this.socketSession) { + return this.socketSession.sendEncrypted(request) + } + console.log('[net] sendEncrypted FAILED — channel not ready', { + hasWs: false, + hasKey: false, + state: this.getState() + }) + return false + } + + private sendLivenessProbe(probeId: string): boolean { + if (this.getState() !== 'connected') { + return false + } + return this.sendEncrypted({ + id: probeId, + deviceToken: this.deviceToken, + method: 'status.get' + }) + } + + private waitForConnected(timeoutMs?: number): Promise { + if ( + this.getState() === 'reconnecting' && + this.getReconnectAttempt() >= RPC_RECONNECT_ATTEMPT_LIMIT + ) { + return Promise.reject(new Error('Connection retry limit reached')) + } + return this.connectionState.waitForConnected(timeoutMs) + } + + private nextId(): string { + return `rpc-${++this.requestCounter}-${Date.now()}` + } +} diff --git a/mobile/src/transport/mobile-terminal-binary-sender.ts b/mobile/src/transport/mobile-terminal-binary-sender.ts index 82c7a7cf3bb..cb3d896be50 100644 --- a/mobile/src/transport/mobile-terminal-binary-sender.ts +++ b/mobile/src/transport/mobile-terminal-binary-sender.ts @@ -1,5 +1,5 @@ -import { encryptedTerminalMultiplexFrame } from './rpc-client-terminal-multiplex' -import type { TerminalStreamFrame } from './terminal-stream-protocol' +import { encryptBytes } from './e2ee' +import { encodeTerminalStreamFrame, type TerminalStreamFrame } from './terminal-stream-protocol' export function sendMobileTerminalBinaryFrame(args: { frame: TerminalStreamFrame @@ -12,7 +12,7 @@ export function sendMobileTerminalBinaryFrame(args: { return false } try { - args.socket.send(encryptedTerminalMultiplexFrame(args.frame, args.sharedKey)) + args.socket.send(encryptBytes(encodeTerminalStreamFrame(args.frame), args.sharedKey)) return true } catch { if (args.isConnected) { diff --git a/mobile/src/transport/rpc-client-authentication-retry.ts b/mobile/src/transport/rpc-client-authentication-retry.ts new file mode 100644 index 00000000000..1e4ca6c6f0b --- /dev/null +++ b/mobile/src/transport/rpc-client-authentication-retry.ts @@ -0,0 +1,46 @@ +import { redactedWebSocketEndpoint } from './redacted-websocket-endpoint' + +const AUTH_RETRY_BUDGET = 3 + +type AuthenticationRetryOptions = { + endpoint: string + stopLiveness: () => void + emitWarning: (message: string, detail: string) => void + retry: (reason: string) => void + latchFailure: (reason: string) => void +} + +export class RpcClientAuthenticationRetry { + private rejectionCount = 0 + + constructor(private readonly options: AuthenticationRetryOptions) {} + + accepted(): void { + this.rejectionCount = 0 + } + + reject(reason: string, preserveRecovery = false): void { + this.options.stopLiveness() + this.rejectionCount++ + if (this.rejectionCount < AUTH_RETRY_BUDGET) { + console.log('[net] auth rejected — retrying handshake', { + attempt: this.rejectionCount, + budget: AUTH_RETRY_BUDGET, + endpoint: redactedWebSocketEndpoint(this.options.endpoint) + }) + this.options.emitWarning( + 'Authentication rejected', + `Retrying (${this.rejectionCount}/${AUTH_RETRY_BUDGET})` + ) + if (!preserveRecovery) { + this.options.retry(reason) + } + return + } + console.log('[net] auth rejected — budget exhausted, latching auth-failed', { + attempt: this.rejectionCount, + endpoint: redactedWebSocketEndpoint(this.options.endpoint) + }) + this.options.latchFailure(reason) + } +} diff --git a/mobile/src/transport/rpc-client-browser-stream-slot.ts b/mobile/src/transport/rpc-client-browser-stream-slot.ts new file mode 100644 index 00000000000..275b9ad125b --- /dev/null +++ b/mobile/src/transport/rpc-client-browser-stream-slot.ts @@ -0,0 +1,49 @@ +// Why: only one browser.screencast stream may own the binary channel; a replacement +// must be parked until the server acknowledges it, or frames route to a dead listener. +export class RpcClientBrowserStreamSlot { + private activeRequestId: string | null = null + private pendingRequestId: string | null = null + + getActiveRequestId(): string | null { + return this.activeRequestId + } + + /** Returns the request ids this replacement supersedes. */ + replaceWith(id: string): string[] { + const superseded = [this.activeRequestId, this.pendingRequestId].filter( + (candidate): candidate is string => candidate !== null && candidate !== id + ) + this.pendingRequestId = id + this.activeRequestId = null + return superseded + } + + markPending(id: string): void { + this.pendingRequestId = id + this.activeRequestId = null + } + + /** False when the server acknowledged a stream this client no longer wants. */ + acknowledge(id: string): boolean { + if (this.pendingRequestId !== id && this.activeRequestId !== id) { + return false + } + this.pendingRequestId = null + this.activeRequestId = id + return true + } + + clear(id: string): void { + if (this.activeRequestId === id) { + this.activeRequestId = null + } + if (this.pendingRequestId === id) { + this.pendingRequestId = null + } + } + + clearAll(): void { + this.activeRequestId = null + this.pendingRequestId = null + } +} diff --git a/mobile/src/transport/rpc-client-connection-state.ts b/mobile/src/transport/rpc-client-connection-state.ts new file mode 100644 index 00000000000..772dd0412e5 --- /dev/null +++ b/mobile/src/transport/rpc-client-connection-state.ts @@ -0,0 +1,114 @@ +import { redactedWebSocketEndpoint } from './redacted-websocket-endpoint' +import type { ConnectionState } from './types' + +type ConnectWaiter = { + resolve: () => void + reject: (error: Error) => void + timeout: ReturnType | null +} + +type ConnectionStateOptions = { + endpoint: string + initialListener?: (state: ConnectionState) => void + getReconnectAttempt: () => number + isClosed: () => boolean +} + +export class RpcClientConnectionState { + private state: ConnectionState = 'disconnected' + private lastConnectedAt: number | null = null + private stateEnteredAt = Date.now() + private readonly listeners = new Set<(state: ConnectionState) => void>() + private readonly waiters: ConnectWaiter[] = [] + + constructor(private readonly options: ConnectionStateOptions) { + if (options.initialListener) { + this.listeners.add(options.initialListener) + } + } + + get(): ConnectionState { + return this.state + } + + getLastConnectedAt(): number | null { + return this.lastConnectedAt + } + + publish(next: ConnectionState): void { + if (this.state === next) { + return + } + const previous = this.state + const dweltMs = Date.now() - this.stateEnteredAt + this.state = next + this.stateEnteredAt = Date.now() + console.log('[net] state', { + from: previous, + to: next, + dweltMs, + attempt: this.options.getReconnectAttempt(), + endpoint: redactedWebSocketEndpoint(this.options.endpoint) + }) + if (next === 'connected') { + this.lastConnectedAt = Date.now() + this.resolveWaiters() + } else if (next === 'disconnected' || next === 'auth-failed') { + this.rejectWaiters( + next === 'auth-failed' ? 'Unauthorized — pairing may be revoked' : 'Connection closed' + ) + } + for (const listener of this.listeners) { + listener(next) + } + } + + waitForConnected(timeoutMs?: number): Promise { + if (this.state === 'connected') { + return Promise.resolve() + } + if (this.options.isClosed()) { + return Promise.reject(new Error('Client closed')) + } + return new Promise((resolve, reject) => { + const waiter: ConnectWaiter = { resolve, reject, timeout: null } + if (timeoutMs !== undefined) { + waiter.timeout = setTimeout( + () => { + const index = this.waiters.indexOf(waiter) + if (index !== -1) { + this.waiters.splice(index, 1) + } + reject(new Error('Timed out while connecting to the remote Orca runtime.')) + }, + Math.max(0, timeoutMs) + ) + } + this.waiters.push(waiter) + }) + } + + rejectWaiters(reason: string): void { + const error = new Error(reason) + for (const waiter of this.waiters.splice(0)) { + if (waiter.timeout) { + clearTimeout(waiter.timeout) + } + waiter.reject(error) + } + } + + addListener(listener: (state: ConnectionState) => void): () => void { + this.listeners.add(listener) + return () => this.listeners.delete(listener) + } + + private resolveWaiters(): void { + for (const waiter of this.waiters.splice(0)) { + if (waiter.timeout) { + clearTimeout(waiter.timeout) + } + waiter.resolve() + } + } +} diff --git a/mobile/src/transport/rpc-client-contract.ts b/mobile/src/transport/rpc-client-contract.ts deleted file mode 100644 index 4087059016f..00000000000 --- a/mobile/src/transport/rpc-client-contract.ts +++ /dev/null @@ -1,37 +0,0 @@ -import type { TerminalStreamFrame } from './terminal-stream-protocol' -import type { RpcClientSubscribeOptions } from './rpc-client-subscribe-options' -import type { ConnectionState, ForegroundNudgeReason, RpcResponse } from './types' - -export type RpcClientSendRequestOptions = { - timeoutMs?: number - /** Spend the timeout across connection wait and acknowledgement. */ - budgetSpansConnect?: boolean - /** Reject immediately instead of replaying stale terminal input after reconnect. */ - failWhenDisconnected?: boolean -} - -export type RpcClient = { - sendRequest: ( - method: string, - params?: unknown, - options?: RpcClientSendRequestOptions - ) => Promise - subscribe: ( - method: string, - params: unknown, - onData: (result: unknown) => void, - options?: RpcClientSubscribeOptions - ) => () => void - updateTerminalSubscriptionViewport: ( - terminal: string, - viewport: { cols: number; rows: number } - ) => void - sendTerminalBinaryFrame: (frame: TerminalStreamFrame) => boolean - getState: () => ConnectionState - getReconnectAttempt: () => number - getLastConnectedAt: () => number | null - getLastInboundAt?: () => number | null - onStateChange: (listener: (state: ConnectionState) => void) => () => void - notifyForeground: (reason?: ForegroundNudgeReason) => void - close: () => void -} diff --git a/mobile/src/transport/rpc-client-foreground-nudge.ts b/mobile/src/transport/rpc-client-foreground-nudge.ts new file mode 100644 index 00000000000..6c117b4209e --- /dev/null +++ b/mobile/src/transport/rpc-client-foreground-nudge.ts @@ -0,0 +1,35 @@ +import { isStaleForegroundDial } from './rpc-stale-dial' +import type { ConnectionState } from './types' + +// Why: iOS/Android can kill the TCP path while backgrounded, so a resume probes or redials. +export function nudgeRpcClientForeground(args: { + // Why: abandoning a stale dial moves the state, so the redial branch must re-read it. + getState: () => ConnectionState + getReconnectAttempt: () => number + dialAgeMs: number + hasReconnectTimer: () => boolean + probeLiveness: () => void + abandonDial: () => boolean + redial: (keepTimer: boolean) => void +}): void { + if (args.getState() === 'connected') { + console.log('[net] foreground — probing live connection') + args.probeLiveness() + return + } + let abandoned = false + if (isStaleForegroundDial(args.getState(), args.dialAgeMs)) { + console.log('[net] foreground — abandoning stale dial', { + state: args.getState(), + dialAgeMs: args.dialAgeMs + }) + abandoned = args.abandonDial() + } + if (args.getState() === 'reconnecting') { + console.log('[net] foreground — restarting reconnect loop', { + attempt: args.getReconnectAttempt(), + hadTimer: args.hasReconnectTimer() + }) + args.redial(!abandoned) + } +} diff --git a/mobile/src/transport/rpc-client-liveness-session.ts b/mobile/src/transport/rpc-client-liveness-session.ts new file mode 100644 index 00000000000..362e9d2d652 --- /dev/null +++ b/mobile/src/transport/rpc-client-liveness-session.ts @@ -0,0 +1,68 @@ +import { + RpcSessionLivenessWatchdog, + type LivenessTimeoutEvidence +} from './rpc-session-liveness-watchdog' +import type { RpcClientSocketSession } from './rpc-client-socket-session' + +export const LIVENESS_REQUEST_ID_PREFIX = 'mobile-liveness-' + +export function isLivenessProbeResponseId(id: string): boolean { + return id.startsWith(LIVENESS_REQUEST_ID_PREFIX) +} + +// Why: the watchdog is per authenticated socket; a probe or terminate on a retired +// session would kill the socket that replaced it. +export class RpcClientLivenessSession { + private readonly watchdog: RpcSessionLivenessWatchdog + private session: RpcClientSocketSession | null = null + + constructor( + private readonly options: { + sendProbe: (probeId: string) => boolean + terminate: (session: RpcClientSocketSession) => void + onTimeout: (evidence: LivenessTimeoutEvidence) => void + nextId: () => string + } + ) { + this.watchdog = new RpcSessionLivenessWatchdog({ + transport: 'direct', + sendProbe: (identity) => + identity === this.session + ? this.options.sendProbe(`${LIVENESS_REQUEST_ID_PREFIX}${this.options.nextId()}`) + : false, + terminate: (identity) => { + if (identity === this.session && this.session) { + this.options.terminate(this.session) + } + }, + onTimeout: this.options.onTimeout + }) + } + + start(session: RpcClientSocketSession): void { + this.session = session + this.watchdog.start(session) + } + + stop(session?: RpcClientSocketSession): void { + if (!this.session || (session && this.session !== session)) { + return + } + this.watchdog.stop(this.session) + this.session = null + } + + probeNow(): void { + if (this.session) { + this.watchdog.probeNow(this.session) + } + } + + noteInbound(session: RpcClientSocketSession): void { + this.watchdog.noteAuthenticatedInbound(session) + } + + getLastInboundAt(): number | null { + return this.watchdog.getLastInboundAt() || null + } +} diff --git a/mobile/src/transport/rpc-client-reconnect-schedule.ts b/mobile/src/transport/rpc-client-reconnect-schedule.ts new file mode 100644 index 00000000000..7497b776ebc --- /dev/null +++ b/mobile/src/transport/rpc-client-reconnect-schedule.ts @@ -0,0 +1,68 @@ +const RECONNECT_DELAYS = [500, 1000, 2000, 4000, 8000, 15_000, 30_000, 60_000] +export const RPC_RECONNECT_ATTEMPT_LIMIT = 12 +const TRICKLE_RECONNECT_DELAY_MS = 90_000 + +type ReconnectScheduleOptions = { + openConnection: () => void + rejectConnectWaiters: (reason: string) => void + emitLog: (message: string, detail: string) => void +} + +export class RpcClientReconnectSchedule { + private attempt = 0 + private timer: ReturnType | null = null + + constructor(private readonly options: ReconnectScheduleOptions) {} + + getAttempt(): number { + return this.attempt + } + + authenticated(): void { + this.attempt = 0 + } + + schedule(): void { + const trickle = this.attempt >= RPC_RECONNECT_ATTEMPT_LIMIT + let delayMs: number + if (trickle) { + delayMs = TRICKLE_RECONNECT_DELAY_MS + this.options.rejectConnectWaiters('Connection retry limit reached') + } else { + delayMs = RECONNECT_DELAYS[Math.min(this.attempt, RECONNECT_DELAYS.length - 1)]! + this.attempt++ + } + console.log('[net] scheduleReconnect', { + delayMs, + attempt: this.attempt, + trickle + }) + this.options.emitLog( + `Reconnect scheduled in ${delayMs}ms`, + trickle ? `Attempt ${this.attempt} (slow retry)` : `Attempt ${this.attempt}` + ) + this.timer = setTimeout(() => { + this.timer = null + this.options.openConnection() + }, delayMs) + } + + redialNow(resetAttempts: boolean): void { + this.cancel() + if (resetAttempts) { + this.attempt = 0 + } + this.options.openConnection() + } + + hasTimer(): boolean { + return this.timer !== null + } + + cancel(): void { + if (this.timer) { + clearTimeout(this.timer) + this.timer = null + } + } +} diff --git a/mobile/src/transport/rpc-client-request-state.ts b/mobile/src/transport/rpc-client-request-state.ts deleted file mode 100644 index 075cbd8a6e6..00000000000 --- a/mobile/src/transport/rpc-client-request-state.ts +++ /dev/null @@ -1,12 +0,0 @@ -import type { RpcResponse } from './types' - -export type RpcPendingRequest = { - resolve: (response: RpcResponse) => void - reject: (error: Error) => void -} - -export type RpcConnectWaiter = { - resolve: () => void - reject: (error: Error) => void - timeout: ReturnType | null -} diff --git a/mobile/src/transport/rpc-client-request-tracker.ts b/mobile/src/transport/rpc-client-request-tracker.ts new file mode 100644 index 00000000000..34edd2631e2 --- /dev/null +++ b/mobile/src/transport/rpc-client-request-tracker.ts @@ -0,0 +1,105 @@ +import { markRpcDeliveryUnknown } from './rpc-delivery-ambiguity' +import { openRpcRequestBudget, resolvePostConnectRequestTimeout } from './rpc-request-budget' +import type { SendRequestOptions } from './rpc-client' +import type { ConnectionState, RpcResponse } from './types' + +const REQUEST_TIMEOUT_MS = 30_000 + +type PendingRequest = { + resolve: (response: RpcResponse) => void + reject: (error: Error) => void +} + +type RequestTrackerOptions = { + nextId: () => string + getState: () => ConnectionState + waitForConnected: (timeoutMs?: number) => Promise + sendEncrypted: (request: unknown) => boolean + deviceToken: string +} + +export class RpcClientRequestTracker { + private readonly pending = new Map() + + constructor(private readonly options: RequestTrackerOptions) {} + + async sendRequest( + method: string, + params?: unknown, + requestOptions?: SendRequestOptions + ): Promise { + const budget = openRpcRequestBudget(requestOptions) + const waitStart = budget.startedAt + const wasConnected = this.options.getState() === 'connected' + if (requestOptions?.failWhenDisconnected && !wasConnected) { + throw new Error(`Not connected: ${method}`) + } + await this.options.waitForConnected(requestOptions?.timeoutMs) + if (!wasConnected) { + console.log('[net] sendRequest waited for connect', { + method, + waitedMs: Date.now() - waitStart + }) + } + + return new Promise((resolve, reject) => { + const id = this.options.nextId() + const timeoutMs = resolvePostConnectRequestTimeout(budget, REQUEST_TIMEOUT_MS) + const timeout = setTimeout(() => { + this.pending.delete(id) + console.log('[net] sendRequest TIMEOUT', { + method, + timeoutMs, + state: this.options.getState() + }) + reject(markRpcDeliveryUnknown(new Error(`Request timed out: ${method}`))) + }, timeoutMs) + this.pending.set(id, { + resolve: (response) => { + clearTimeout(timeout) + resolve(response) + }, + reject: (error) => { + clearTimeout(timeout) + reject(error) + } + }) + if ( + !this.options.sendEncrypted({ + id, + deviceToken: this.options.deviceToken, + method, + params + }) + ) { + this.pending.delete(id) + clearTimeout(timeout) + reject(new Error('Connection interrupted')) + } + }) + } + + resolve(response: RpcResponse): boolean { + const request = this.pending.get(response.id) + if (!request) { + return false + } + this.pending.delete(response.id) + request.resolve(response) + return true + } + + rejectAll(reason: string, options?: { deliveryUnknown?: boolean }): void { + const error = options?.deliveryUnknown + ? markRpcDeliveryUnknown(new Error(reason)) + : new Error(reason) + for (const [id, request] of this.pending) { + this.pending.delete(id) + queueMicrotask(() => request.reject(error)) + } + } + + size(): number { + return this.pending.size + } +} diff --git a/mobile/src/transport/rpc-client-runtime-events.test.ts b/mobile/src/transport/rpc-client-runtime-events.test.ts new file mode 100644 index 00000000000..d18739423f7 --- /dev/null +++ b/mobile/src/transport/rpc-client-runtime-events.test.ts @@ -0,0 +1,135 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { connect, type RpcClient } from './rpc-client' + +vi.mock('./e2ee', () => ({ + generateKeyPair: () => ({ + publicKey: new Uint8Array(32), + secretKey: new Uint8Array(32) + }), + deriveSharedKey: () => new Uint8Array(32), + publicKeyFromBase64: () => new Uint8Array(32), + publicKeyToBase64: () => 'client-public-key', + encrypt: (plaintext: string) => `encrypted:${plaintext}`, + decrypt: (raw: string) => raw.replace(/^encrypted:/, ''), + decryptBytes: (bytes: Uint8Array) => bytes +})) + +class RuntimeEventTestSocket { + static CONNECTING = 0 + static OPEN = 1 + static CLOSING = 2 + static CLOSED = 3 + + readonly CONNECTING = RuntimeEventTestSocket.CONNECTING + readonly OPEN = RuntimeEventTestSocket.OPEN + readonly CLOSING = RuntimeEventTestSocket.CLOSING + readonly CLOSED = RuntimeEventTestSocket.CLOSED + + readyState = RuntimeEventTestSocket.CONNECTING + onopen: (() => void) | null = null + onclose: (() => void) | null = null + onmessage: ((event: { data: unknown }) => void) | null = null + onerror: (() => void) | null = null + sent: string[] = [] + + constructor(readonly endpoint: string) { + sockets.push(this) + } + + send(payload: string): void { + this.sent.push(payload) + } + + close(): void { + this.readyState = RuntimeEventTestSocket.CLOSED + this.onclose?.() + } + + open(): void { + this.readyState = RuntimeEventTestSocket.OPEN + this.onopen?.() + } + + receive(payload: unknown): void { + this.onmessage?.({ data: payload }) + } +} + +type SentRequest = { id: string; method: string; params?: unknown } + +const sockets: RuntimeEventTestSocket[] = [] +const originalWebSocket = globalThis.WebSocket + +function sentRequests(socket: RuntimeEventTestSocket, method: string): SentRequest[] { + return socket.sent + .map((payload) => JSON.parse(payload.replace(/^encrypted:/, '')) as SentRequest) + .filter((request) => request.method === method) +} + +function connectReadyClient(): { client: RpcClient; socket: RuntimeEventTestSocket } { + const client = connect('ws://desktop.invalid', 'token', 'server-key') + const socket = sockets[0]! + socket.open() + socket.receive(JSON.stringify({ type: 'e2ee_ready' })) + socket.receive('encrypted:{"type":"e2ee_authenticated"}') + return { client, socket } +} + +function emitReady( + socket: RuntimeEventTestSocket, + requestId: string, + subscriptionId: string +): void { + socket.receive( + `encrypted:${JSON.stringify({ + id: requestId, + ok: true, + streaming: true, + result: { type: 'ready', subscriptionId }, + _meta: { runtimeId: 'r1' } + })}` + ) +} + +describe('runtime client-event stream disposal', () => { + beforeEach(() => { + vi.useFakeTimers() + sockets.length = 0 + globalThis.WebSocket = RuntimeEventTestSocket as unknown as typeof WebSocket + }) + + afterEach(() => { + globalThis.WebSocket = originalWebSocket + vi.useRealTimers() + }) + + it('unsubscribes a stream disposed after ready', () => { + const { client, socket } = connectReadyClient() + const unsubscribe = client.subscribe('runtime.clientEvents.subscribe', null, () => {}) + const request = sentRequests(socket, 'runtime.clientEvents.subscribe')[0]! + emitReady(socket, request.id, 'runtime-events:test') + + unsubscribe() + + expect(sentRequests(socket, 'runtime.clientEvents.unsubscribe')).toEqual([ + expect.objectContaining({ params: { subscriptionId: 'runtime-events:test' } }) + ]) + client.close() + }) + + it('keeps a tombstone and unsubscribes a stream disposed before ready', () => { + const { client, socket } = connectReadyClient() + const listener = vi.fn() + const unsubscribe = client.subscribe('runtime.clientEvents.subscribe', null, listener) + const request = sentRequests(socket, 'runtime.clientEvents.subscribe')[0]! + + unsubscribe() + emitReady(socket, request.id, 'runtime-events:late') + + expect(listener).not.toHaveBeenCalled() + expect(sentRequests(socket, 'runtime.clientEvents.unsubscribe')).toEqual([ + expect.objectContaining({ params: { subscriptionId: 'runtime-events:late' } }) + ]) + client.close() + }) +}) diff --git a/mobile/src/transport/rpc-client-socket-close-controller.ts b/mobile/src/transport/rpc-client-socket-close-controller.ts new file mode 100644 index 00000000000..a7b0cecdc82 --- /dev/null +++ b/mobile/src/transport/rpc-client-socket-close-controller.ts @@ -0,0 +1,98 @@ +import type { RpcClientAuthenticationRetry } from './rpc-client-authentication-retry' +import type { RpcClientConnectionState } from './rpc-client-connection-state' +import type { RpcClientReconnectSchedule } from './rpc-client-reconnect-schedule' +import type { RpcClientRequestTracker } from './rpc-client-request-tracker' +import type { RpcClientSocketFactory } from './rpc-client-socket-factory' +import type { RpcClientSocketSession } from './rpc-client-socket-session' +import type { RpcClientStreamRegistry } from './rpc-client-stream-registry' +import { RpcSynthesizedCloseIndex } from './rpc-socket-close-evidence' +import type { ConnectionLogEntry } from './types' + +const UNAUTHORIZED_CLOSE_CODE = 4001 + +type SocketCloseControllerOptions = { + connectionState: RpcClientConnectionState + reconnect: RpcClientReconnectSchedule + requests: RpcClientRequestTracker + streams: RpcClientStreamRegistry + socketFactory: RpcClientSocketFactory + authenticationRetry: RpcClientAuthenticationRetry + getCurrentSession: () => RpcClientSocketSession | null + clearCurrentSession: () => void + getAuthenticationGeneration: () => number + isIntentionallyClosed: () => boolean + stopLiveness: (session: RpcClientSocketSession) => void + emitWarning: ( + message: string, + detail: string, + evidence?: Pick + ) => void +} + +export class RpcClientSocketCloseController { + private readonly synthesizedCloses = new RpcSynthesizedCloseIndex() + + constructor(private readonly options: SocketCloseControllerOptions) {} + + forceClose(session: RpcClientSocketSession): void { + session.close() + if (this.options.getCurrentSession() === session) { + this.synthesizedCloses.remember(session.socket, this.options.getAuthenticationGeneration()) + this.handle(session) + } + } + + handle(session: RpcClientSocketSession, closeCode?: number): void { + if (this.options.getCurrentSession() !== session) { + if ( + this.synthesizedCloses.takeUnauthorized( + session.socket, + closeCode, + this.options.getAuthenticationGeneration(), + UNAUTHORIZED_CLOSE_CODE + ) + ) { + this.options.authenticationRetry.reject('Unauthorized — pairing may be revoked', true) + return + } + console.log('[net] handleSocketClosed STALE — ignoring (ws already swapped)', { + state: this.options.connectionState.get(), + attempt: this.options.reconnect.getAttempt() + }) + return + } + this.options.socketFactory.noteClosed() + session.dispose() + session.clearTimers() + session.clearKey() + this.options.clearCurrentSession() + this.options.streams.markForReplay() + this.options.stopLiveness(session) + if (this.options.isIntentionallyClosed()) { + console.log('[net] handleSocketClosed — intentional close') + this.options.connectionState.publish('disconnected') + this.options.requests.rejectAll('Connection closed', { deliveryUnknown: true }) + return + } + if (closeCode === UNAUTHORIZED_CLOSE_CODE) { + console.log('[net] handleSocketClosed — unauthorized close code', { + attempt: this.options.reconnect.getAttempt() + }) + this.options.authenticationRetry.reject('Unauthorized — pairing may be revoked') + return + } + console.log('[net] handleSocketClosed → reconnect', { + pendingCount: this.options.requests.size(), + streamCount: this.options.streams.size(), + attempt: this.options.reconnect.getAttempt() + }) + this.options.emitWarning( + 'WebSocket closed', + `${closeCode == null ? 'Close code unavailable' : `Close code ${closeCode}`}; reconnect scheduled`, + { code: 'socket-closed' } + ) + this.options.requests.rejectAll('Connection interrupted', { deliveryUnknown: true }) + this.options.connectionState.publish('reconnecting') + this.options.reconnect.schedule() + } +} diff --git a/mobile/src/transport/rpc-client-socket-factory.ts b/mobile/src/transport/rpc-client-socket-factory.ts new file mode 100644 index 00000000000..60dc8c8e592 --- /dev/null +++ b/mobile/src/transport/rpc-client-socket-factory.ts @@ -0,0 +1,90 @@ +import { publicKeyFromBase64 } from './e2ee' +import { RpcClientSocketSession } from './rpc-client-socket-session' +import { processMobileOutboundMemoryBudget } from './mobile-outbound-memory-budget' +import { redactedWebSocketEndpoint } from './redacted-websocket-endpoint' +import type { ConnectionLogEmitter, ConnectionState, RpcResponse } from './types' + +type SocketFactoryOptions = { + endpoint: string + deviceToken: string + serverPublicKeyB64: string + getCurrentSocket: () => WebSocket | null + getState: () => ConnectionState + getReconnectAttempt: () => number + getLastConnectedAt: () => number | null + isIntentionallyClosed: () => boolean + emitLog: ConnectionLogEmitter + onHandshakeStarted: () => void + onAuthenticated: (session: RpcClientSocketSession) => void + onAuthRejected: (reason: string) => void + onRpcResponse: (response: RpcResponse) => void + onBinary: (bytes: Uint8Array) => void + onAuthenticatedInbound: (session: RpcClientSocketSession) => void + onClosed: (session: RpcClientSocketSession, closeCode?: number) => void + onForcedClose: (session: RpcClientSocketSession) => void +} + +export class RpcClientSocketFactory { + private readonly serverPublicKey: Uint8Array + private lastInboundAt: number | null = null + private lastSocketClosedAt: number | null = null + private constructionCount = 0 + private dialStartedAt = 0 + + constructor(private readonly options: SocketFactoryOptions) { + this.serverPublicKey = publicKeyFromBase64(options.serverPublicKeyB64) + } + + canOpen(): boolean { + return processMobileOutboundMemoryBudget.canRegisterBufferedAmount() + } + + open(): RpcClientSocketSession { + const now = Date.now() + const lastConnectedAt = this.options.getLastConnectedAt() + this.constructionCount++ + console.log('[net] openConnection', { + attempt: this.options.getReconnectAttempt(), + endpoint: redactedWebSocketEndpoint(this.options.endpoint), + wsCount: this.constructionCount, + msSinceLastConnected: lastConnectedAt !== null ? now - lastConnectedAt : null, + msSinceLastClose: this.lastSocketClosedAt !== null ? now - this.lastSocketClosedAt : null, + msSinceLastInbound: this.lastInboundAt !== null ? now - this.lastInboundAt : null + }) + this.dialStartedAt = now + this.options.emitLog( + 'info', + this.options.getReconnectAttempt() > 0 + ? `Reconnecting (attempt ${this.options.getReconnectAttempt() + 1})` + : 'Opening WebSocket', + redactedWebSocketEndpoint(this.options.endpoint) + ) + return new RpcClientSocketSession({ + endpoint: this.options.endpoint, + deviceToken: this.options.deviceToken, + serverPublicKey: this.serverPublicKey, + getCurrentSocket: this.options.getCurrentSocket, + getState: this.options.getState, + getReconnectAttempt: this.options.getReconnectAttempt, + isIntentionallyClosed: this.options.isIntentionallyClosed, + emitLog: this.options.emitLog, + onHandshakeStarted: this.options.onHandshakeStarted, + onAuthenticated: this.options.onAuthenticated, + onAuthRejected: this.options.onAuthRejected, + onRpcResponse: this.options.onRpcResponse, + onBinary: this.options.onBinary, + onAnyInbound: (receivedAt) => (this.lastInboundAt = receivedAt), + onAuthenticatedInbound: this.options.onAuthenticatedInbound, + onClosed: this.options.onClosed, + onForcedClose: this.options.onForcedClose + }) + } + + getDialStartedAt(): number { + return this.dialStartedAt + } + + noteClosed(): void { + this.lastSocketClosedAt = Date.now() + } +} diff --git a/mobile/src/transport/rpc-client-socket-session.ts b/mobile/src/transport/rpc-client-socket-session.ts new file mode 100644 index 00000000000..0b22227c0d0 --- /dev/null +++ b/mobile/src/transport/rpc-client-socket-session.ts @@ -0,0 +1,265 @@ +import { decrypt, deriveSharedKey, generateKeyPair, publicKeyToBase64 } from './e2ee' +import { createMobileDirectRpcOutbound } from './mobile-direct-rpc-outbound' +import { createMobileDirectRpcSender } from './mobile-direct-rpc-sender' +import { + createMobileInboundFrameQueue, + MOBILE_INBOUND_BUFFER_OVERFLOW_MESSAGE, + MOBILE_INBOUND_FRAME_TOO_LARGE_MESSAGE, + mobileInboundFrameLogDetail, + type MobileInboundFrameQueue +} from './mobile-inbound-frame-queue' +import { tryParseMobileJsonTextWithinLimits } from './mobile-json-text-admission' +import { handleMobileRpcSocketBinaryMessage } from './mobile-rpc-binary-frame-handler' +import { sendMobileTerminalBinaryFrame } from './mobile-terminal-binary-sender' +import { redactedWebSocketEndpoint } from './redacted-websocket-endpoint' +import { isRpcResponse } from './rpc-response-shape' +import { RpcClientSocketTimeouts } from './rpc-client-socket-timeouts' +import { isStaleRpcSocketEvent, logRpcSocketClose } from './rpc-socket-close-evidence' +import { describeSocketEvent } from './socket-event-debug' +import type { TerminalStreamFrame } from './terminal-stream-protocol' +import type { ConnectionLogEmitter, ConnectionState, RpcResponse } from './types' + +type SocketSessionOptions = { + endpoint: string + deviceToken: string + serverPublicKey: Uint8Array + getCurrentSocket: () => WebSocket | null + getState: () => ConnectionState + getReconnectAttempt: () => number + isIntentionallyClosed: () => boolean + emitLog: ConnectionLogEmitter + onHandshakeStarted: () => void + onAuthenticated: (session: RpcClientSocketSession) => void + onAuthRejected: (reason: string) => void + onRpcResponse: (response: RpcResponse) => void + onBinary: (bytes: Uint8Array) => void + onAnyInbound: (receivedAt: number) => void + onAuthenticatedInbound: (session: RpcClientSocketSession) => void + onClosed: (session: RpcClientSocketSession, closeCode?: number) => void + onForcedClose: (session: RpcClientSocketSession) => void +} + +export class RpcClientSocketSession { + readonly socket: WebSocket + readonly constructedAt = Date.now() + private sharedKey: Uint8Array | null = null + private authenticated = false + private lastInboundAt: number | null = null + private readonly timeouts: RpcClientSocketTimeouts + private readonly outbound: ReturnType + private readonly inboundQueue: MobileInboundFrameQueue + readonly sendEncrypted: (request: unknown) => boolean + + constructor(private readonly options: SocketSessionOptions) { + this.socket = new WebSocket(options.endpoint) + this.outbound = createMobileDirectRpcOutbound({ + socket: this.socket, + isActive: () => this.options.getCurrentSocket() === this.socket, + onOverflow: () => this.closeForOverload('Outbound', 'Mobile RPC outbound buffer overflow') + }) + this.inboundQueue = createMobileInboundFrameQueue({ + process: (raw) => this.handleMessage(raw), + onError: (error) => this.closeForOverload('Inbound', mobileInboundFrameLogDetail(error)), + overflowMessage: MOBILE_INBOUND_BUFFER_OVERFLOW_MESSAGE, + frameTooLargeMessage: MOBILE_INBOUND_FRAME_TOO_LARGE_MESSAGE + }) + this.sendEncrypted = createMobileDirectRpcSender({ + getOutbound: () => this.outbound, + getSharedKey: () => this.sharedKey, + getSocket: () => this.socket, + getState: () => this.options.getState(), + onSocketDesync: () => this.options.onForcedClose(this) + }) + this.timeouts = new RpcClientSocketTimeouts({ + emitLog: options.emitLog, + getReconnectAttempt: options.getReconnectAttempt, + expire: () => this.options.onForcedClose(this) + }) + this.attachHandlers() + this.timeouts.armConnect( + () => this.options.getCurrentSocket() === this.socket, + () => this.socket.readyState + ) + } + + sendTerminalBinaryFrame(frame: TerminalStreamFrame): boolean { + return sendMobileTerminalBinaryFrame({ + frame, + socket: this.socket, + sharedKey: this.sharedKey, + isConnected: this.options.getState() === 'connected', + onSocketClosed: () => this.options.onForcedClose(this) + }) + } + + close(): void { + this.socket.close() + } + + dispose(): void { + this.inboundQueue.dispose() + this.outbound.dispose() + } + + private closeForOverload(direction: 'Inbound' | 'Outbound', detail: string): void { + this.options.emitLog('error', `${direction} WebSocket overload`, detail) + this.options.onForcedClose(this) + } + + clearTimers(): void { + this.timeouts.clearAll() + } + + clearKey(): void { + this.sharedKey = null + } + + private attachHandlers(): void { + this.socket.onopen = () => { + if (this.isStale('open')) { + return + } + console.log('[net] ws.onopen', { attempt: this.options.getReconnectAttempt() }) + this.timeouts.clearConnect() + this.options.onHandshakeStarted() + this.options.emitLog('success', 'WebSocket open', 'Starting E2EE handshake') + const ephemeral = generateKeyPair() + const hello = JSON.stringify({ + type: 'e2ee_hello', + publicKeyB64: publicKeyToBase64(ephemeral.publicKey) + }) + try { + this.socket.send(hello) + } catch { + this.options.onForcedClose(this) + return + } + this.options.emitLog('info', 'Sent e2ee_hello', 'Awaiting server e2ee_ready') + this.sharedKey = deriveSharedKey(ephemeral.secretKey, this.options.serverPublicKey) + this.timeouts.armHandshake( + () => this.options.getCurrentSocket() === this.socket && !this.authenticated + ) + } + this.socket.onmessage = (event) => { + if (!this.isStale('message')) { + void this.inboundQueue.enqueue(event.data) + } + } + this.socket.onclose = (event) => { + this.inboundQueue.dispose() + this.outbound.socketClosed() + const closeCode = logRpcSocketClose({ + event, + state: this.options.getState(), + attempt: this.options.getReconnectAttempt(), + intentionallyClosed: this.options.isIntentionallyClosed(), + endpoint: redactedWebSocketEndpoint(this.options.endpoint), + constructedAt: this.constructedAt, + authenticated: this.authenticated, + lastInboundAt: this.lastInboundAt + }) + this.options.onClosed(this, closeCode) + } + this.socket.onerror = (event) => { + if (this.isStale('error')) { + return + } + console.log('[net] ws.onerror', { + state: this.options.getState(), + attempt: this.options.getReconnectAttempt(), + eventFields: describeSocketEvent(event).fields + }) + } + } + + private handleMessage(rawData: unknown): Promise | void { + const receivedAt = Date.now() + this.lastInboundAt = receivedAt + this.options.onAnyInbound(receivedAt) + const raw = typeof rawData === 'string' ? rawData : null + if (!this.authenticated) { + this.handleHandshakeMessage(raw) + return + } + const key = this.sharedKey + if (!key || key.length !== 32) { + return + } + if (raw === null) { + return handleMobileRpcSocketBinaryMessage({ + rawData, + key, + isCurrent: () => this.options.getCurrentSocket() === this.socket, + onFrame: (plaintext) => { + this.options.onAuthenticatedInbound(this) + this.options.onBinary(plaintext) + } + }) + } + const plaintext = decrypt(raw, key) + if (plaintext === null) { + return + } + this.options.onAuthenticatedInbound(this) + const response = tryParseMobileJsonTextWithinLimits(plaintext) + if (!isRpcResponse(response)) { + return + } + this.outbound.acknowledge(response.id) + this.options.onRpcResponse(response) + } + + private handleHandshakeMessage(raw: string | null): void { + if (raw === null) { + return + } + const plaintextControl = tryParseMobileJsonTextWithinLimits<{ type?: unknown }>(raw) + if (plaintextControl?.type === 'e2ee_ready') { + this.options.emitLog('success', 'Received e2ee_ready', 'Sending device token') + this.sendEncrypted({ type: 'e2ee_auth', deviceToken: this.options.deviceToken }) + return + } + if (!this.sharedKey || this.sharedKey.length !== 32) { + return + } + const plaintext = decrypt(raw, this.sharedKey) + if (plaintext === null) { + return + } + const message = tryParseMobileJsonTextWithinLimits<{ + type?: unknown + ok?: unknown + error?: { code?: unknown } + }>(plaintext) + if (!message) { + return + } + if (message.type === 'e2ee_authenticated') { + this.outbound.acknowledgeAuthentication() + this.timeouts.clearHandshake() + this.authenticated = true + this.options.onAuthenticated(this) + } else if ( + message.type === 'e2ee_error' || + (message.ok === false && message.error?.code === 'unauthorized') + ) { + this.outbound.acknowledgeAuthentication() + // Why: the failure signal is loggable; the server error body is not. + console.log('[net] e2ee auth FAILED', { + signal: message.type === 'e2ee_error' ? 'e2ee_error' : 'unauthorized_response' + }) + this.timeouts.clearHandshake() + this.options.onAuthRejected('Unauthorized — pairing may be revoked') + } + } + + private isStale(eventName: string): boolean { + return isStaleRpcSocketEvent( + this.options.getCurrentSocket(), + this.socket, + eventName, + this.options.getState(), + this.options.getReconnectAttempt() + ) + } +} diff --git a/mobile/src/transport/rpc-client-socket-timeouts.ts b/mobile/src/transport/rpc-client-socket-timeouts.ts new file mode 100644 index 00000000000..5a96447ad1e --- /dev/null +++ b/mobile/src/transport/rpc-client-socket-timeouts.ts @@ -0,0 +1,79 @@ +import type { ConnectionLogEmitter } from './types' + +export const RPC_SOCKET_CONNECT_TIMEOUT_MS = 12_000 +export const RPC_SOCKET_HANDSHAKE_TIMEOUT_MS = 5_000 +const WEBSOCKET_CONNECTING_STATE = 0 + +type TimeoutHandle = ReturnType | null + +// Why: RN can leave a dial or handshake pending forever on flaky handoffs; force a reconnect instead. +export class RpcClientSocketTimeouts { + private connectTimer: TimeoutHandle = null + private handshakeTimer: TimeoutHandle = null + + constructor( + private readonly options: { + emitLog: ConnectionLogEmitter + getReconnectAttempt: () => number + expire: () => void + } + ) {} + + armConnect(isCurrentSocket: () => boolean, readyState: () => number): void { + this.connectTimer = setTimeout(() => { + this.connectTimer = null + if (!isCurrentSocket() || readyState() !== WEBSOCKET_CONNECTING_STATE) { + return + } + console.log('[net] connect-timeout fired (onopen never arrived)', { + attempt: this.options.getReconnectAttempt(), + timeoutMs: RPC_SOCKET_CONNECT_TIMEOUT_MS + }) + this.options.emitLog( + 'error', + 'WebSocket connect timeout', + `No TCP/WS handshake within ${RPC_SOCKET_CONNECT_TIMEOUT_MS / 1000}s — endpoint unreachable?`, + { code: 'connect-timeout' } + ) + this.options.expire() + }, RPC_SOCKET_CONNECT_TIMEOUT_MS) + } + + armHandshake(isPending: () => boolean): void { + this.handshakeTimer = setTimeout(() => { + this.handshakeTimer = null + if (!isPending()) { + return + } + console.log('[net] handshake-timeout fired (e2ee_authenticated never arrived)', { + timeoutMs: RPC_SOCKET_HANDSHAKE_TIMEOUT_MS + }) + this.options.emitLog( + 'error', + 'Handshake timeout', + `No e2ee_ready/e2ee_authenticated within ${RPC_SOCKET_HANDSHAKE_TIMEOUT_MS / 1000}s`, + { code: 'handshake-timeout' } + ) + this.options.expire() + }, RPC_SOCKET_HANDSHAKE_TIMEOUT_MS) + } + + clearConnect(): void { + if (this.connectTimer) { + clearTimeout(this.connectTimer) + this.connectTimer = null + } + } + + clearHandshake(): void { + if (this.handshakeTimer) { + clearTimeout(this.handshakeTimer) + this.handshakeTimer = null + } + } + + clearAll(): void { + this.clearConnect() + this.clearHandshake() + } +} diff --git a/mobile/src/transport/rpc-client-stream-registry.test.ts b/mobile/src/transport/rpc-client-stream-registry.test.ts new file mode 100644 index 00000000000..0e7d561be6e --- /dev/null +++ b/mobile/src/transport/rpc-client-stream-registry.test.ts @@ -0,0 +1,101 @@ +import { describe, expect, it } from 'vitest' +import { RpcClientStreamRegistry } from './rpc-client-stream-registry' +import { encodeTerminalStreamFrame, TerminalStreamOpcode } from './terminal-stream-protocol' +import type { ConnectionState, RpcResponse } from './types' + +type SentRequest = { + id: string + method: string + params?: unknown +} + +function createRegistry(initialState: ConnectionState = 'connected') { + const sent: SentRequest[] = [] + let state = initialState + let id = 0 + const registry = new RpcClientStreamRegistry({ + nextId: () => `rpc-${++id}`, + deviceToken: 'device-token', + getState: () => state, + sendEncrypted: (request) => { + sent.push(request as SentRequest) + return true + } + }) + return { + registry, + sent, + setState(next: ConnectionState) { + state = next + } + } +} + +function streamingResponse(id: string, result: unknown): RpcResponse { + return { + id, + ok: true, + streaming: true, + result, + _meta: { runtimeId: 'runtime-1' } + } +} + +function terminalOutput(streamId: number, chunk: string): Uint8Array { + return encodeTerminalStreamFrame({ + opcode: TerminalStreamOpcode.Output, + streamId, + seq: 1, + payload: new TextEncoder().encode(chunk) + }) +} + +describe('RpcClientStreamRegistry', () => { + it('replays the latest terminal viewport without retaining stale stream routing', () => { + const { registry, sent } = createRegistry() + const events: unknown[] = [] + registry.subscribe( + 'terminal.subscribe', + { terminal: 'term-1', viewport: { cols: 45, rows: 20 } }, + (event) => events.push(event) + ) + const first = sent[0]! + registry.handleResponse(streamingResponse(first.id, { type: 'subscribed', streamId: 7 })) + registry.handleBinary(terminalOutput(7, 'before')) + + registry.updateTerminalViewport('term-1', { cols: 60, rows: 24 }) + registry.markForReplay() + registry.replayAfterAuthentication() + registry.handleBinary(terminalOutput(7, 'stale')) + + expect(sent[1]).toMatchObject({ + id: first.id, + method: 'terminal.subscribe', + params: { terminal: 'term-1', viewport: { cols: 60, rows: 24 } } + }) + expect(events).toEqual([ + { type: 'subscribed', streamId: 7 }, + { type: 'data', streamId: 7, chunk: 'before' } + ]) + }) + + it('keeps a disposed browser tombstone until ready can be unsubscribed', () => { + const { registry, sent } = createRegistry() + const dispose = registry.subscribe('browser.screencast', { page: 'page-1' }, () => {}) + const request = sent[0]! + + dispose() + expect(sent).toHaveLength(1) + + registry.handleResponse( + streamingResponse(request.id, { + type: 'ready', + subscriptionId: 'browser-screencast:page-1:test' + }) + ) + expect(sent[1]).toMatchObject({ + method: 'browser.screencast.unsubscribe', + params: { subscriptionId: 'browser-screencast:page-1:test' } + }) + }) +}) diff --git a/mobile/src/transport/rpc-client-stream-registry.ts b/mobile/src/transport/rpc-client-stream-registry.ts new file mode 100644 index 00000000000..ad5476b7ce4 --- /dev/null +++ b/mobile/src/transport/rpc-client-stream-registry.ts @@ -0,0 +1,305 @@ +import { + decodeBrowserScreencastFrame, + type BrowserScreencastFrame +} from './browser-screencast-protocol' +import { + buildStreamUnsubscribe, + buildTerminalUnsubscribeParams, + updateTerminalSubscriptionViewport +} from './rpc-client-terminal-subscription' +import { buildReadyStreamUnsubscribe } from './rpc-client-server-subscription' +import { + isStreamingSubscriptionReadyResult, + isTerminalSubscribedResult +} from './rpc-subscription-result-shapes' +import { RpcClientBrowserStreamSlot } from './rpc-client-browser-stream-slot' +import { routeTerminalMultiplexFrame } from './rpc-client-terminal-multiplex' +import { RpcClientTerminalStreamRouter } from './rpc-client-terminal-stream-router' +import type { TerminalStreamFrame } from './terminal-stream-protocol' +import type { ConnectionState, RpcResponse, RpcSuccess } from './types' + +export type RpcStreamingListener = (result: unknown) => void + +export type RpcStreamSubscribeOptions = { + onBinaryFrame?: (frame: BrowserScreencastFrame) => void + onTerminalBinaryFrame?: (frame: TerminalStreamFrame) => boolean +} + +type StreamRequest = { + method: string + params: unknown + listener: RpcStreamingListener + onBinaryFrame?: (frame: BrowserScreencastFrame) => void + onTerminalBinaryFrame?: (frame: TerminalStreamFrame) => boolean + subscriptionId?: string + cancelled?: boolean + sent?: boolean +} + +type StreamRegistryOptions = { + nextId: () => string + deviceToken: string + getState: () => ConnectionState + sendEncrypted: (request: unknown) => boolean +} + +export class RpcClientStreamRegistry { + private readonly streams = new Map() + private readonly terminalRouter = new RpcClientTerminalStreamRouter() + private readonly browserSlot = new RpcClientBrowserStreamSlot() + + constructor(private readonly options: StreamRegistryOptions) {} + + subscribe( + method: string, + params: unknown, + listener: RpcStreamingListener, + subscribeOptions?: RpcStreamSubscribeOptions + ): () => void { + const id = this.options.nextId() + const stream: StreamRequest = { + method, + params, + listener, + onBinaryFrame: subscribeOptions?.onBinaryFrame, + onTerminalBinaryFrame: subscribeOptions?.onTerminalBinaryFrame + } + this.streams.set(id, stream) + if (method === 'browser.screencast') { + this.replaceBrowserStream(id) + } + if (this.options.getState() === 'connected') { + if (this.send(id, stream)) { + stream.sent = true + } else { + this.emitError(stream, 'Connection interrupted') + this.remove(id) + } + } else { + console.log('[net] subscribe queued — waiting for connected', { + method, + state: this.options.getState() + }) + } + return () => this.dispose(id) + } + + replayAfterAuthentication(): void { + for (const [id, stream] of this.streams) { + if (stream.cancelled) { + this.remove(id) + continue + } + if (stream.sent) { + continue + } + if (stream.method === 'browser.screencast') { + this.browserSlot.markPending(id) + } + this.resetTerminalRouting(id) + if (this.send(id, stream)) { + stream.sent = true + } else { + this.markForReplay() + break + } + } + } + + markForReplay(): void { + this.browserSlot.clearAll() + for (const [id, stream] of this.streams) { + stream.sent = false + this.resetTerminalRouting(id) + } + } + + handleResponse(response: RpcResponse): boolean { + if (response.ok && response.streaming === true) { + this.handleStreamingResponse(response) + return true + } + const stream = this.streams.get(response.id) + if (response.ok) { + const result = (response as RpcSuccess).result as Record | null + if (stream && result?.type === 'end') { + if (!stream.cancelled) { + stream.listener(result) + } + this.remove(response.id) + return true + } + if (stream && result?.type === 'scrollback') { + stream.listener(result) + return true + } + } + if (!stream) { + return false + } + this.emitError( + stream, + response.ok ? 'Streaming request ended before it was ready.' : response.error.message, + response.ok ? undefined : response.error + ) + this.remove(response.id) + return true + } + + handleBinary(bytes: Uint8Array): void { + const browserFrame = decodeBrowserScreencastFrame(bytes) + if (browserFrame) { + this.handleBrowserFrame(browserFrame) + return + } + if (routeTerminalMultiplexFrame(bytes, this.streams.values())) { + return + } + this.terminalRouter.handle(bytes) + } + + updateTerminalViewport(terminal: string, viewport: { cols: number; rows: number }): void { + updateTerminalSubscriptionViewport(this.streams.values(), terminal, viewport) + } + + size(): number { + return this.streams.size + } + + private handleStreamingResponse(response: RpcSuccess): void { + const stream = this.streams.get(response.id) + if (!stream) { + return + } + const result = response.result + if (isStreamingSubscriptionReadyResult(result)) { + stream.subscriptionId = result.subscriptionId + if (stream.cancelled) { + this.sendServerSubscriptionUnsubscribe(stream) + this.remove(response.id) + return + } + if (stream.method === 'browser.screencast' && !this.browserSlot.acknowledge(response.id)) { + this.sendBrowserUnsubscribe(result.subscriptionId) + this.remove(response.id) + return + } + } + if (isTerminalSubscribedResult(result)) { + this.terminalRouter.register(response.id, result.streamId, stream.listener) + } + if (!stream.cancelled) { + stream.listener(result) + } + } + + private dispose(id: string): void { + const stream = this.streams.get(id) + if (stream?.method === 'browser.screencast') { + stream.cancelled = true + this.browserSlot.clear(id) + this.disposeServerSubscription(id, stream) + return + } + if ( + stream?.method === 'runtime.clientEvents.subscribe' || + stream?.method === 'accounts.subscribe' || + stream?.method === 'files.watch' + ) { + this.disposeServerSubscription(id, stream) + return + } + if (stream?.method === 'terminal.subscribe') { + const params = buildTerminalUnsubscribeParams(stream.params) + if (params) { + this.sendRpc('terminal.unsubscribe', params) + } + } else { + const unsubscribe = buildStreamUnsubscribe(stream?.method, stream?.params) + if (unsubscribe) { + this.sendRpc(unsubscribe.method, unsubscribe.params) + } + } + this.remove(id) + } + + private replaceBrowserStream(id: string): void { + for (const superseded of this.browserSlot.replaceWith(id)) { + this.dispose(superseded) + } + } + + private disposeServerSubscription(id: string, stream: StreamRequest): void { + stream.cancelled = true + if (stream.subscriptionId) { + this.sendServerSubscriptionUnsubscribe(stream) + this.remove(id) + } else if (!stream.sent) { + this.remove(id) + } + } + + private sendServerSubscriptionUnsubscribe(stream: StreamRequest): void { + if (!stream.subscriptionId) { + return + } + const unsubscribe = buildReadyStreamUnsubscribe(stream.method, stream.subscriptionId) + if (unsubscribe) { + this.sendRpc(unsubscribe.method, unsubscribe.params) + } + } + + private sendBrowserUnsubscribe(subscriptionId: string): void { + this.sendRpc('browser.screencast.unsubscribe', { subscriptionId }) + } + + private sendRpc(method: string, params: unknown): boolean { + return this.options.sendEncrypted({ + id: this.options.nextId(), + deviceToken: this.options.deviceToken, + method, + params + }) + } + + private send(id: string, stream: StreamRequest): boolean { + return this.options.sendEncrypted({ + id, + deviceToken: this.options.deviceToken, + method: stream.method, + params: stream.params + }) + } + + private remove(id: string): void { + const stream = this.streams.get(id) + this.streams.delete(id) + this.browserSlot.clear(id) + this.resetTerminalRouting(id) + if (stream?.method === 'browser.screencast') { + stream.cancelled = true + } + } + + private resetTerminalRouting(id: string): void { + this.terminalRouter.reset(id) + } + + private handleBrowserFrame(frame: BrowserScreencastFrame): void { + const activeId = this.browserSlot.getActiveRequestId() + if (!activeId) { + return + } + const stream = this.streams.get(activeId) + if (!stream || stream.cancelled || stream.method !== 'browser.screencast') { + return + } + stream.onBinaryFrame?.(frame) + } + + private emitError(stream: StreamRequest, message: string, error?: unknown): void { + if (!stream.cancelled) { + stream.listener({ type: 'error', message, error }) + } + } +} diff --git a/mobile/src/transport/rpc-client-subscribe-options.ts b/mobile/src/transport/rpc-client-subscribe-options.ts deleted file mode 100644 index c10299ad519..00000000000 --- a/mobile/src/transport/rpc-client-subscribe-options.ts +++ /dev/null @@ -1,7 +0,0 @@ -import type { BrowserScreencastFrame } from './browser-screencast-protocol' -import type { TerminalStreamFrame } from './terminal-stream-protocol' - -export type RpcClientSubscribeOptions = { - onBinaryFrame?: (frame: BrowserScreencastFrame) => void - onTerminalBinaryFrame?: (frame: TerminalStreamFrame) => boolean -} diff --git a/mobile/src/transport/rpc-client-terminal-multiplex.ts b/mobile/src/transport/rpc-client-terminal-multiplex.ts index 2116d896331..90fb4ff47b0 100644 --- a/mobile/src/transport/rpc-client-terminal-multiplex.ts +++ b/mobile/src/transport/rpc-client-terminal-multiplex.ts @@ -1,9 +1,4 @@ -import { encryptBytes } from './e2ee' -import { - decodeTerminalStreamFrame, - encodeTerminalStreamFrame, - type TerminalStreamFrame -} from './terminal-stream-protocol' +import { decodeTerminalStreamFrame, type TerminalStreamFrame } from './terminal-stream-protocol' type MultiplexStream = { method: string @@ -30,10 +25,3 @@ export function routeTerminalMultiplexFrame( } return false } - -export function encryptedTerminalMultiplexFrame( - frame: TerminalStreamFrame, - sharedKey: Uint8Array -): Uint8Array { - return encryptBytes(encodeTerminalStreamFrame(frame), sharedKey) -} diff --git a/mobile/src/transport/rpc-client-terminal-stream-router.ts b/mobile/src/transport/rpc-client-terminal-stream-router.ts new file mode 100644 index 00000000000..0ac4d26c3fc --- /dev/null +++ b/mobile/src/transport/rpc-client-terminal-stream-router.ts @@ -0,0 +1,38 @@ +import { + handleTerminalBinaryFrame, + type TerminalSnapshotState +} from './rpc-client-terminal-binary-frame' + +type TerminalListener = (result: unknown) => void + +export class RpcClientTerminalStreamRouter { + private readonly listeners = new Map() + private readonly idsByRequest = new Map>() + private readonly snapshots = new Map() + + register(requestId: string, streamId: number, listener: TerminalListener): void { + const ids = this.idsByRequest.get(requestId) ?? new Set() + this.idsByRequest.set(requestId, ids) + ids.add(streamId) + this.listeners.set(streamId, listener) + } + + reset(requestId: string): void { + const streamIds = this.idsByRequest.get(requestId) + if (!streamIds) { + return + } + for (const streamId of streamIds) { + this.listeners.delete(streamId) + this.snapshots.delete(streamId) + } + this.idsByRequest.delete(requestId) + } + + handle(bytes: Uint8Array): void { + handleTerminalBinaryFrame(bytes, { + terminalSnapshots: this.snapshots, + getListener: (streamId) => this.listeners.get(streamId) + }) + } +} diff --git a/mobile/src/transport/rpc-client.ts b/mobile/src/transport/rpc-client.ts index bb166295c1d..a3e7d6102ec 100644 --- a/mobile/src/transport/rpc-client.ts +++ b/mobile/src/transport/rpc-client.ts @@ -1,111 +1,52 @@ -import type { - RpcResponse, - RpcSuccess, - ConnectionState, - ConnectionLogLevel, - ConnectionLogSink, - ForegroundNudgeReason -} from './types' -import { - generateKeyPair, - deriveSharedKey, - publicKeyFromBase64, - publicKeyToBase64, - decrypt -} from './e2ee' -import { - handleTerminalBinaryFrame, - type TerminalSnapshotState -} from './rpc-client-terminal-binary-frame' -import { - decodeBrowserScreencastFrame, - type BrowserScreencastFrame -} from './browser-screencast-protocol' -import type { RpcClientSubscribeOptions } from './rpc-client-subscribe-options' -import { - buildStreamUnsubscribe, - buildTerminalUnsubscribeParams, - updateTerminalSubscriptionViewport as updateCachedTerminalSubscriptionViewport -} from './rpc-client-terminal-subscription' -import { describeSocketEvent } from './socket-event-debug' -import { - isStaleRpcSocketEvent, - logRpcSocketClose, - RpcSynthesizedCloseIndex -} from './rpc-socket-close-evidence' -import { markRpcDeliveryUnknown } from './rpc-delivery-ambiguity' -import { openRpcRequestBudget, resolvePostConnectRequestTimeout } from './rpc-request-budget' -import { isRpcResponse } from './rpc-response-shape' -import { - isStreamingSubscriptionReadyResult, - isTerminalSubscribedResult -} from './rpc-subscription-result-shapes' -import { - RpcSessionLivenessWatchdog, - type RpcSessionIdentity -} from './rpc-session-liveness-watchdog' -import { isStaleForegroundDial } from './rpc-stale-dial' -import { - createMobileInboundFrameQueue, - MOBILE_INBOUND_BUFFER_OVERFLOW_MESSAGE, - MOBILE_INBOUND_FRAME_TOO_LARGE_MESSAGE, - mobileInboundFrameLogDetail -} from './mobile-inbound-frame-queue' -import { createMobileDirectRpcOutbound } from './mobile-direct-rpc-outbound' -import { createMobileDirectRpcSender } from './mobile-direct-rpc-sender' -import { sendMobileTerminalBinaryFrame } from './mobile-terminal-binary-sender' -import { handleMobileRpcSocketBinaryMessage } from './mobile-rpc-binary-frame-handler' -import { processMobileOutboundMemoryBudget } from './mobile-outbound-memory-budget' -import { redactedWebSocketEndpoint } from './redacted-websocket-endpoint' -import { tryParseMobileJsonTextWithinLimits } from './mobile-json-text-admission' +import { DirectRpcClient } from './direct-rpc-client' +import type { RpcStreamSubscribeOptions } from './rpc-client-stream-registry' import type { TerminalStreamFrame } from './terminal-stream-protocol' -import { buildServerSubscriptionUnsubscribe } from './rpc-client-server-subscription' -import { routeTerminalMultiplexFrame } from './rpc-client-terminal-multiplex' -import type { RpcClient, RpcClientSendRequestOptions } from './rpc-client-contract' -import type { RpcConnectWaiter, RpcPendingRequest } from './rpc-client-request-state' +import type { + ConnectionLogSink, + ConnectionState, + ForegroundNudgeReason, + RpcResponse +} from './types' -export type { - RpcClient, - RpcClientSendRequestOptions, - RpcClientSendRequestOptions as SendRequestOptions -} from './rpc-client-contract' +export type SendRequestOptions = { + timeoutMs?: number + /** Include the connect wait in the caller's timeout budget. */ + budgetSpansConnect?: boolean + /** Reject instead of replaying the request after reconnect. */ + failWhenDisconnected?: boolean +} type StreamingListener = (result: unknown) => void -type StreamRequest = { - method: string - params: unknown - listener: StreamingListener - onBinaryFrame?: (frame: BrowserScreencastFrame) => void - onTerminalBinaryFrame?: (frame: TerminalStreamFrame) => boolean - subscriptionId?: string - cancelled?: boolean - sent?: boolean +export type RpcClient = { + sendRequest: ( + method: string, + params?: unknown, + options?: SendRequestOptions + ) => Promise + subscribe: ( + method: string, + params: unknown, + onData: StreamingListener, + options?: RpcStreamSubscribeOptions + ) => () => void + updateTerminalSubscriptionViewport: ( + terminal: string, + viewport: { cols: number; rows: number } + ) => void + // Why: the hosted bridge multiplexes terminal bytes over one binary channel. + sendTerminalBinaryFrame: (frame: TerminalStreamFrame) => boolean + getState: () => ConnectionState + getReconnectAttempt: () => number + getLastConnectedAt: () => number | null + getLastInboundAt?: () => number | null + onStateChange: (listener: (state: ConnectionState) => void) => () => void + notifyForeground: (reason?: ForegroundNudgeReason) => void + close: () => void } -// Why: tiered backoff — fast early entries recover blips; the slow tail avoids burning a SYN every 4s on an unreachable desktop. -const RECONNECT_DELAYS = [500, 1000, 2000, 4000, 8000, 15_000, 30_000, 60_000] -// Why: ≈6 min of failure before the re-pair banner; MUST stay aligned with connection-health.ts UNREACHABLE_ATTEMPTS. -const GIVE_UP_AFTER_ATTEMPTS = 12 -// Why: never park past the cap — a wedged VPN fires no AppState/network nudge to revive it, so trickle-dial every 90s to self-heal. -const TRICKLE_RECONNECT_DELAY_MS = 90_000 -// Why: one unauthorized isn't proof the pairing is dead (issue #5200) — retry the handshake this many times before latching auth-failed. -const AUTH_RETRY_BUDGET = 3 -// Why: a desktop that regenerated its E2EE keypair sends an e2ee_error we can't decrypt — the 4001 close code is the only surviving auth-failure signal. -const UNAUTHORIZED_CLOSE_CODE = 4001 -const REQUEST_TIMEOUT_MS = 30_000 -// Why: an explicit `timeoutMs` is one budget for the whole call. If the connect wait -// ate nearly all of it, still give the written frame a moment to be answered rather -// than arming a 1ms timer. -const CONNECT_TIMEOUT_MS = 12_000 -const HANDSHAKE_TIMEOUT_MS = 5_000 -// Why: RN may not expose WebSocket.readyState constants, but the CONNECTING protocol value (0) is stable across runtimes. -const WEBSOCKET_CONNECTING_STATE = 0 -const LIVENESS_REQUEST_ID_PREFIX = 'mobile-liveness-' - export type ConnectOptions = { onStateChange?: (state: ConnectionState) => void - // Fires for every lifecycle event so the UI can show where 'Connecting…' is stuck (e.g. broken Tailscale route). onLog?: ConnectionLogSink } @@ -115,1088 +56,9 @@ export function connect( serverPublicKeyB64: string, optionsOrLegacy?: ConnectOptions | ((state: ConnectionState) => void) ): RpcClient { - // Why: keep backward-compat with callers that pass a bare onStateChange fn. const options: ConnectOptions = typeof optionsOrLegacy === 'function' ? { onStateChange: optionsOrLegacy } : (optionsOrLegacy ?? {}) - const onStateChange = options.onStateChange - const onLog = options.onLog - let logCounter = 0 - function emitLog(level: ConnectionLogLevel, message: string, detail?: string) { - if (!onLog) { - return - } - onLog({ - id: `log-${++logCounter}-${Date.now()}`, - ts: Date.now(), - level, - message, - detail - }) - } - let ws: WebSocket | null = null - const synthesizedCloses = new RpcSynthesizedCloseIndex() - let outbound: ReturnType | null = null - let state: ConnectionState = 'disconnected' - let requestCounter = 0 - let reconnectAttempt = 0 - let reconnectTimer: ReturnType | null = null - let connectTimer: ReturnType | null = null - let handshakeTimer: ReturnType | null = null - let intentionallyClosed = false - // Consecutive auth rejections; tolerate up to AUTH_RETRY_BUDGET (issue #5200) before latching to avoid a needless re-pair. - let authRejectionCount = 0 - let authenticationGeneration = 0 - let lastConnectedAt: number | null = null - // Why: cheap diagnostics for RN/OkHttp process-state poisoning (retry cadence, inbound traffic, close timing). - let lastInboundAt: number | null = null - let livenessIdentity: (RpcSessionIdentity & { socket: WebSocket }) | null = null - let lastWsClosedAt: number | null = null - let wsConstructionCounter = 0 - let dialStartedAt = 0 - - // Why: fresh ephemeral keypair per connection provides forward secrecy. - let sharedKey: Uint8Array | null = null - const serverPublicKey = publicKeyFromBase64(serverPublicKeyB64) - - const pending = new Map() - const streamListeners = new Map() - const terminalStreamListeners = new Map() - const terminalStreamIdsByRequest = new Map>() - const terminalSnapshots = new Map() - let activeBrowserScreencastRequestId: string | null = null - let pendingBrowserScreencastRequestId: string | null = null - const stateListeners = new Set<(state: ConnectionState) => void>() - const connectWaiters: RpcConnectWaiter[] = [] - - if (onStateChange) { - stateListeners.add(onStateChange) - } - - // Diagnostic: dwell time in the current state, for spotting "stuck in connecting/reconnecting". - let stateEnteredAt = Date.now() - - function rejectConnectWaiters(reason: string) { - const error = new Error(reason) - for (const waiter of connectWaiters.splice(0)) { - if (waiter.timeout) { - clearTimeout(waiter.timeout) - } - waiter.reject(error) - } - } - - function setState(next: ConnectionState) { - if (state === next) { - return - } - const prev = state - const dwelt = Date.now() - stateEnteredAt - state = next - stateEnteredAt = Date.now() - console.log('[net] state', { - from: prev, - to: next, - dweltMs: dwelt, - attempt: reconnectAttempt, - endpoint: redactedWebSocketEndpoint(endpoint) - }) - if (next === 'connected') { - lastConnectedAt = Date.now() - authenticationGeneration++ - // Why: only a completed E2EE handshake proves the path is healthy (issue #10119). - reconnectAttempt = 0 - // Why: a clean handshake proves the token is valid — reset the auth retry budget. - authRejectionCount = 0 - for (const waiter of connectWaiters.splice(0)) { - if (waiter.timeout) { - clearTimeout(waiter.timeout) - } - waiter.resolve() - } - } else if (next === 'disconnected' || next === 'auth-failed') { - const reason = - next === 'auth-failed' ? 'Unauthorized — pairing may be revoked' : 'Connection closed' - rejectConnectWaiters(reason) - } - for (const listener of stateListeners) { - listener(next) - } - } - - function waitForConnected(timeoutMs?: number): Promise { - if (state === 'connected') { - return Promise.resolve() - } - if (intentionallyClosed) { - return Promise.reject(new Error('Client closed')) - } - if (state === 'reconnecting' && reconnectAttempt >= GIVE_UP_AFTER_ATTEMPTS) { - // Why: past the cap the loop only trickles every 90s — fail fast instead of hanging on a long-unreachable host. - return Promise.reject(new Error('Connection retry limit reached')) - } - return new Promise((resolve, reject) => { - const waiter: RpcConnectWaiter = { resolve, reject, timeout: null } - if (timeoutMs !== undefined) { - // Why: per-request timeouts must cover offline/reconnect waiting, not just the RPC after connect. - waiter.timeout = setTimeout( - () => { - const index = connectWaiters.indexOf(waiter) - if (index !== -1) { - connectWaiters.splice(index, 1) - } - reject(new Error('Timed out while connecting to the remote Orca runtime.')) - }, - Math.max(0, timeoutMs) - ) - } - connectWaiters.push(waiter) - }) - } - - function nextId(): string { - return `rpc-${++requestCounter}-${Date.now()}` - } - - function disposeActiveOutbound(): void { - outbound?.dispose() - outbound = null - } - - function openConnection() { - if (intentionallyClosed) { - return - } - - const now = Date.now() - wsConstructionCounter++ - console.log('[net] openConnection', { - attempt: reconnectAttempt, - endpoint: redactedWebSocketEndpoint(endpoint), - // Why: diagnostic for RN/OkHttp pool corruption — high wsCount + repeated 1006 closes means process-state stuck. - wsCount: wsConstructionCounter, - msSinceLastConnected: lastConnectedAt != null ? now - lastConnectedAt : null, - msSinceLastClose: lastWsClosedAt != null ? now - lastWsClosedAt : null, - msSinceLastInbound: lastInboundAt != null ? now - lastInboundAt : null - }) - setState('connecting') - dialStartedAt = now - sharedKey = null - - emitLog( - 'info', - reconnectAttempt > 0 ? `Reconnecting (attempt ${reconnectAttempt + 1})` : 'Opening WebSocket', - redactedWebSocketEndpoint(endpoint) - ) - - if (!processMobileOutboundMemoryBudget.canRegisterBufferedAmount()) { - emitLog('error', 'WebSocket reconnect deferred', 'Retired socket buffers are still draining') - setState('reconnecting') - scheduleReconnect() - return - } - - ws = new WebSocket(endpoint) - const openingWs = ws - let openingWsAuthenticated = false - let openingWsLastInboundAt: number | null = null - let openingLivenessIdentity: (RpcSessionIdentity & { socket: WebSocket }) | null = null - const closeForOverload = (direction: 'Inbound' | 'Outbound', detail: string): void => { - emitLog('error', `${direction} WebSocket overload`, detail) - openingWs.close() - if (ws === openingWs) { - handleSocketClosed(openingWs) - } - } - const openingOutbound = createMobileDirectRpcOutbound({ - socket: openingWs, - isActive: () => ws === openingWs, - onOverflow: () => closeForOverload('Outbound', 'Mobile RPC outbound buffer overflow') - }) - outbound = openingOutbound - const inboundQueue = createMobileInboundFrameQueue({ - process: handleSocketMessage, - onError: (error) => closeForOverload('Inbound', mobileInboundFrameLogDetail(error)), - overflowMessage: MOBILE_INBOUND_BUFFER_OVERFLOW_MESSAGE, - frameTooLargeMessage: MOBILE_INBOUND_FRAME_TOO_LARGE_MESSAGE - }) - // Why: RN can leave opens pending forever on flaky handoffs — force reconnect if onopen never arrives. - connectTimer = setTimeout(() => { - connectTimer = null - if (ws === openingWs && openingWs.readyState === WEBSOCKET_CONNECTING_STATE) { - console.log('[net] connect-timeout fired (onopen never arrived)', { - attempt: reconnectAttempt, - timeoutMs: CONNECT_TIMEOUT_MS - }) - emitLog( - 'error', - 'WebSocket connect timeout', - `No TCP/WS handshake within ${CONNECT_TIMEOUT_MS / 1000}s — endpoint unreachable?` - ) - closeAndSynthesize(openingWs) - } - }, CONNECT_TIMEOUT_MS) - - ws.onopen = () => { - if (isStaleRpcSocketEvent(ws, openingWs, 'open', state, reconnectAttempt)) { - return - } - console.log('[net] ws.onopen', { attempt: reconnectAttempt }) - clearConnectTimer() - // Why: no reconnectAttempt reset here — an open socket isn't a healthy session - // until e2ee_authenticated. Resetting pre-handshake pinned the counter at 0↔1, - // so a handshake-stall loop never escalated past "Connecting…" (issue #10119). - setState('handshaking') - emitLog('success', 'WebSocket open', 'Starting E2EE handshake') - - // Why: fresh ephemeral keypair per connection provides forward secrecy. - const ephemeral = generateKeyPair() - const hello = JSON.stringify({ - type: 'e2ee_hello', - publicKeyB64: publicKeyToBase64(ephemeral.publicKey) - }) - try { - openingWs.send(hello) - } catch { - closeAndSynthesize(openingWs) - return - } - emitLog('info', 'Sent e2ee_hello', 'Awaiting server e2ee_ready') - - sharedKey = deriveSharedKey(ephemeral.secretKey, serverPublicKey) - - handshakeTimer = setTimeout(() => { - handshakeTimer = null - if (ws !== openingWs || state !== 'handshaking') { - return - } - console.log('[net] handshake-timeout fired (e2ee_authenticated never arrived)', { - timeoutMs: HANDSHAKE_TIMEOUT_MS - }) - emitLog( - 'error', - 'Handshake timeout', - `No e2ee_ready/e2ee_authenticated within ${HANDSHAKE_TIMEOUT_MS / 1000}s` - ) - closeAndSynthesize(openingWs) - }, HANDSHAKE_TIMEOUT_MS) - } - - ws.onmessage = (event) => { - if (isStaleRpcSocketEvent(ws, openingWs, 'message', state, reconnectAttempt)) { - return - } - void inboundQueue.enqueue(event.data) - } - - function handleSocketMessage(rawData: unknown): Promise | void { - const receivedAt = Date.now() - lastInboundAt = receivedAt - openingWsLastInboundAt = receivedAt - const raw = typeof rawData === 'string' ? rawData : null - - // Why: e2ee_ready is plaintext (precedes encrypted auth); e2ee_authenticated/e2ee_error are encrypted. - if (state === 'handshaking') { - if (raw === null) { - return - } - const plaintextControl = tryParseMobileJsonTextWithinLimits>(raw) - if (plaintextControl?.type === 'e2ee_ready') { - emitLog('success', 'Received e2ee_ready', 'Sending device token') - sendEncrypted({ type: 'e2ee_auth', deviceToken }) - return - } - - if (!sharedKey || sharedKey.length !== 32) { - return - } - - const plaintext = decrypt(raw, sharedKey) - if (plaintext === null) { - return - } - - const msg = tryParseMobileJsonTextWithinLimits>(plaintext) - if (msg) { - if (msg.type === 'e2ee_authenticated') { - openingOutbound.acknowledgeAuthentication() - if (handshakeTimer) { - clearTimeout(handshakeTimer) - handshakeTimer = null - } - console.log('[net] e2ee_authenticated — connected', { - streamCount: streamListeners.size - }) - openingWsAuthenticated = true - openingLivenessIdentity = { socket: openingWs } - livenessIdentity = openingLivenessIdentity - livenessWatchdog.start(openingLivenessIdentity) - setState('connected') - emitLog('success', 'Authenticated', 'Channel ready for RPC') - for (const [id, stream] of streamListeners) { - if (stream.cancelled) { - removeStreamListener(id) - continue - } - // Why: a UI listener notified synchronously by setState('connected') may already have sent this stream — skip it. - if (stream.sent) { - continue - } - if (stream.method === 'browser.screencast') { - pendingBrowserScreencastRequestId = id - activeBrowserScreencastRequestId = null - } - resetTerminalStreamRoutingForRequest(id) - if ( - sendEncrypted({ id, deviceToken, method: stream.method, params: stream.params }) - ) { - stream.sent = true - } else { - // The failed write already starts recovery; retain every stream for replay. - markStreamsForReplay() - break - } - } - } else if ( - msg.type === 'e2ee_error' || - (!msg.ok && (msg.error as { code?: unknown } | undefined)?.code === 'unauthorized') - ) { - openingOutbound.acknowledgeAuthentication() - console.log('[net] e2ee auth FAILED', { - signal: msg.type === 'e2ee_error' ? 'e2ee_error' : 'unauthorized_response' - }) - if (handshakeTimer) { - clearTimeout(handshakeTimer) - handshakeTimer = null - } - handleAuthRejection('Unauthorized — pairing may be revoked') - } - } - return - } - - // Why: sharedKey can be null after destroy() or a reconnect race — don't decrypt with an invalid key. - if (!sharedKey || sharedKey.length !== 32) { - return - } - - if (raw === null) { - return handleMobileRpcSocketBinaryMessage({ - rawData, - key: sharedKey, - isCurrent: () => ws === openingWs, - onFrame: (frame) => { - if (openingLivenessIdentity) { - livenessWatchdog.noteAuthenticatedInbound(openingLivenessIdentity) - } - handleBinaryFrame(frame) - } - }) - } - - const plaintext = decrypt(raw, sharedKey) - if (plaintext === null) { - return - } - if (openingLivenessIdentity) { - livenessWatchdog.noteAuthenticatedInbound(openingLivenessIdentity) - } - - const response = tryParseMobileJsonTextWithinLimits(plaintext) - if (response === null) { - return - } - if (!isRpcResponse(response)) { - return - } - openingOutbound.acknowledge(response.id) - if (response.id.startsWith(LIVENESS_REQUEST_ID_PREFIX)) { - return - } - // Why: a mid-session unauthorized may be transient (issue #5200) — handleAuthRejection retries before latching auth-failed. - if (!response.ok && response.error.code === 'unauthorized') { - handleAuthRejection('Unauthorized — pairing may be revoked') - return - } - - const isStreaming = response.ok && (response as RpcSuccess).streaming === true - - if (isStreaming) { - const stream = streamListeners.get(response.id) - if (stream && response.ok) { - const result = (response as RpcSuccess).result - if (isStreamingSubscriptionReadyResult(result)) { - stream.subscriptionId = result.subscriptionId - if (stream.cancelled) { - sendServerSubscriptionUnsubscribe(stream) - removeStreamListener(response.id) - return - } - if (stream.method === 'browser.screencast') { - if ( - pendingBrowserScreencastRequestId !== response.id && - activeBrowserScreencastRequestId !== response.id - ) { - sendBrowserScreencastUnsubscribe(result.subscriptionId) - removeStreamListener(response.id) - return - } - pendingBrowserScreencastRequestId = null - activeBrowserScreencastRequestId = response.id - } - } - if (isTerminalSubscribedResult(result)) { - let ids = terminalStreamIdsByRequest.get(response.id) - if (!ids) { - ids = new Set() - terminalStreamIdsByRequest.set(response.id, ids) - } - ids.add(result.streamId) - terminalStreamListeners.set(result.streamId, stream.listener) - } - if (!stream.cancelled) { - stream.listener(result) - } - } - return - } - - if (response.ok) { - const result = (response as RpcSuccess).result as Record | null - if (result && result.type === 'end') { - const stream = streamListeners.get(response.id) - if (stream) { - if (!stream.cancelled) { - stream.listener(result) - } - removeStreamListener(response.id) - return - } - } - if (result && result.type === 'scrollback') { - const stream = streamListeners.get(response.id) - if (stream) { - stream.listener(result) - return - } - } - } - - const stream = streamListeners.get(response.id) - if (stream) { - if (!response.ok) { - emitStreamError(stream, response.error.message, response.error) - } else { - emitStreamError(stream, 'Streaming request ended before it was ready.') - } - removeStreamListener(response.id) - return - } - - const req = pending.get(response.id) - if (req) { - pending.delete(response.id) - req.resolve(response) - } - } - - ws.onclose = (event) => { - inboundQueue.dispose() - openingOutbound.socketClosed() - if (outbound === openingOutbound) { - disposeActiveOutbound() - } - const closeCode = logRpcSocketClose({ - event, - state, - attempt: reconnectAttempt, - intentionallyClosed, - endpoint: redactedWebSocketEndpoint(endpoint), - constructedAt: now, - authenticated: openingWsAuthenticated, - lastInboundAt: openingWsLastInboundAt - }) - handleSocketClosed(openingWs, { closeCode }) - } - - ws.onerror = (event) => { - if (isStaleRpcSocketEvent(ws, openingWs, 'error', state, reconnectAttempt)) { - return - } - const errEvent = describeSocketEvent(event) - console.log('[net] ws.onerror', { - state, - attempt: reconnectAttempt, - eventFields: errEvent.fields - }) - } - } - - function handleSocketClosed( - closedWs: WebSocket, - opts: { timedOut?: boolean; closeCode?: number } = {} - ) { - if (ws !== closedWs) { - if ( - synthesizedCloses.takeUnauthorized( - closedWs, - opts.closeCode, - authenticationGeneration, - UNAUTHORIZED_CLOSE_CODE - ) - ) { - handleAuthRejection('Unauthorized — pairing may be revoked', true) - return - } - console.log('[net] handleSocketClosed STALE — ignoring (ws already swapped)', { - state, - attempt: reconnectAttempt - }) - return - } - lastWsClosedAt = Date.now() - clearConnectTimer() - disposeActiveOutbound() - ws = null - sharedKey = null - activeBrowserScreencastRequestId = null - pendingBrowserScreencastRequestId = null - markStreamsForReplay() - clearHandshakeTimer() - if (livenessIdentity?.socket === closedWs) { - livenessWatchdog.stop(livenessIdentity) - livenessIdentity = null - } - if (intentionallyClosed) { - console.log('[net] handleSocketClosed — intentional close') - setState('disconnected') - rejectAllPending('Connection closed', { deliveryUnknown: true }) - return - } - // Why: a bare 4001 close means the desktop rejected our pairing but the encrypted - // e2ee_error never arrived (or was undecryptable) — count it against the auth - // retry budget instead of looping the generic reconnect forever. - if (opts.closeCode === UNAUTHORIZED_CLOSE_CODE) { - console.log('[net] handleSocketClosed — unauthorized close code', { - attempt: reconnectAttempt - }) - handleAuthRejection('Unauthorized — pairing may be revoked') - return - } - console.log('[net] handleSocketClosed → reconnect', { - timedOut: !!opts.timedOut, - pendingCount: pending.size, - streamCount: streamListeners.size, - attempt: reconnectAttempt - }) - emitLog('warn', 'WebSocket closed', 'Will attempt to reconnect') - rejectAllPending('Connection interrupted', { deliveryUnknown: true }) - setState('reconnecting') - scheduleReconnect() - } - - // Why: an auth rejection may be transient (issue #5200) — retry up to AUTH_RETRY_BUDGET times before latching auth-failed. - function handleAuthRejection(reason: string, preserveRecovery = false): void { - if (livenessIdentity) { - livenessWatchdog.stop(livenessIdentity) - livenessIdentity = null - } - authRejectionCount++ - if (authRejectionCount < AUTH_RETRY_BUDGET) { - console.log('[net] auth rejected — retrying handshake', { - attempt: authRejectionCount, - budget: AUTH_RETRY_BUDGET, - endpoint: redactedWebSocketEndpoint(endpoint) - }) - emitLog( - 'warn', - 'Authentication rejected', - `Retrying (${authRejectionCount}/${AUTH_RETRY_BUDGET})` - ) - if (preserveRecovery) { - return - } - activeBrowserScreencastRequestId = null - pendingBrowserScreencastRequestId = null - // Why: close without setting intentionallyClosed so handleSocketClosed routes to reconnect and retries the handshake. - const closing = ws - disposeActiveOutbound() - ws = null - sharedKey = null - // Why: close cleanup stale-bails here, so mark active streams for replay. - markStreamsForReplay() - rejectAllPending(reason) - if (closing) { - closing.close() - } - setState('reconnecting') - scheduleReconnect() - return - } - activeBrowserScreencastRequestId = null - pendingBrowserScreencastRequestId = null - console.log('[net] auth rejected — budget exhausted, latching auth-failed', { - attempt: authRejectionCount, - endpoint: redactedWebSocketEndpoint(endpoint) - }) - intentionallyClosed = true - disposeActiveOutbound() - ws?.close() - ws = null - setState('auth-failed') - rejectAllPending(reason) - } - - function scheduleReconnect() { - // Why: past the cap, trickle (never park) — a parked loop only revives on a network transition a wedged VPN never produces. - const pastGiveUpCap = reconnectAttempt >= GIVE_UP_AFTER_ATTEMPTS - let delay: number - if (pastGiveUpCap) { - // Why: hold the counter at the cap — connection-health's "Can't reach desktop" verdict keys off attempts >= 12. - delay = TRICKLE_RECONNECT_DELAY_MS - rejectConnectWaiters('Connection retry limit reached') - } else { - delay = RECONNECT_DELAYS[Math.min(reconnectAttempt, RECONNECT_DELAYS.length - 1)]! - reconnectAttempt++ - } - console.log('[net] scheduleReconnect', { - delayMs: delay, - attempt: reconnectAttempt, - trickle: pastGiveUpCap - }) - emitLog( - 'info', - `Reconnect scheduled in ${delay}ms`, - pastGiveUpCap ? `Attempt ${reconnectAttempt} (slow retry)` : `Attempt ${reconnectAttempt}` - ) - reconnectTimer = setTimeout(() => { - reconnectTimer = null - openConnection() - }, delay) - } - - function clearConnectTimer() { - if (connectTimer) { - clearTimeout(connectTimer) - connectTimer = null - } - } - - function clearHandshakeTimer() { - if (handshakeTimer) { - clearTimeout(handshakeTimer) - handshakeTimer = null - } - } - - // Why: a revival signal dials at once rather than waiting out the armed backoff. - // Only the caller knows whether the attempt that led here should be forgiven. - function redialNow(resetAttempts: boolean) { - if (reconnectTimer) { - clearTimeout(reconnectTimer) - reconnectTimer = null - } - if (resetAttempts) { - reconnectAttempt = 0 - } - openConnection() - } - - // Why: React Native can omit onclose for a wedged iOS transport, so every forced - // close has to synthesize the close event it may never deliver. - function closeAndSynthesize(socket: WebSocket) { - socket.close() - if (ws === socket) { - synthesizedCloses.remember(socket, authenticationGeneration) - handleSocketClosed(socket, { timedOut: true }) - } - } - - function rejectAllPending(reason: string, options?: { deliveryUnknown?: boolean }) { - // Why: pending entries only exist after a successful socket write, so a close - // here means the host may have processed them — mark the ambiguity for callers. - const error = options?.deliveryUnknown - ? markRpcDeliveryUnknown(new Error(reason)) - : new Error(reason) - for (const [id, req] of pending) { - pending.delete(id) - queueMicrotask(() => req.reject(error)) - } - } - - function removeStreamListener(id: string): void { - const stream = streamListeners.get(id) - streamListeners.delete(id) - if (activeBrowserScreencastRequestId === id) { - activeBrowserScreencastRequestId = null - } - if (pendingBrowserScreencastRequestId === id) { - pendingBrowserScreencastRequestId = null - } - const terminalStreamIds = terminalStreamIdsByRequest.get(id) - if (terminalStreamIds) { - for (const streamId of terminalStreamIds) { - terminalStreamListeners.delete(streamId) - terminalSnapshots.delete(streamId) - } - terminalStreamIdsByRequest.delete(id) - } - if (stream?.method === 'browser.screencast') { - stream.cancelled = true - } - } - - function markStreamsForReplay(): void { - for (const [id, stream] of streamListeners) { - stream.sent = false - resetTerminalStreamRoutingForRequest(id) - } - } - - function resetTerminalStreamRoutingForRequest(id: string): void { - const terminalStreamIds = terminalStreamIdsByRequest.get(id) - if (!terminalStreamIds) { - return - } - for (const streamId of terminalStreamIds) { - terminalStreamListeners.delete(streamId) - terminalSnapshots.delete(streamId) - } - terminalStreamIdsByRequest.delete(id) - } - - function emitStreamError(stream: StreamRequest, message: string, error?: unknown): void { - if (stream.cancelled) { - return - } - stream.listener({ type: 'error', message, error }) - } - - function disposeBrowserScreencastStream(id: string): void { - const stream = streamListeners.get(id) - if (!stream || stream.method !== 'browser.screencast') { - return - } - stream.cancelled = true - if (activeBrowserScreencastRequestId === id) { - activeBrowserScreencastRequestId = null - } - if (pendingBrowserScreencastRequestId === id) { - pendingBrowserScreencastRequestId = null - } - disposeServerSubscriptionStream(id, stream) - } - - function disposeServerSubscriptionStream(id: string, stream: StreamRequest): void { - stream.cancelled = true - if (stream.subscriptionId) { - sendServerSubscriptionUnsubscribe(stream) - removeStreamListener(id) - return - } - // Why: a sent stream may still reply `ready`; keep the tombstone to unsubscribe it (queued streams never reached the desktop). - if (!stream.sent) { - removeStreamListener(id) - } - } - - function handleBinaryFrame(bytes: Uint8Array): void { - const browserFrame = decodeBrowserScreencastFrame(bytes) - if (browserFrame) { - handleBrowserBinaryFrame(browserFrame) - return - } - if (routeTerminalMultiplexFrame(bytes, streamListeners.values())) { - return - } - handleTerminalBinaryFrame(bytes, { - terminalSnapshots, - getListener: (streamId) => terminalStreamListeners.get(streamId) - }) - } - - function handleBrowserBinaryFrame(frame: BrowserScreencastFrame) { - if (!activeBrowserScreencastRequestId) { - return - } - const stream = streamListeners.get(activeBrowserScreencastRequestId) - if (!stream || stream.cancelled || stream.method !== 'browser.screencast') { - return - } - stream.onBinaryFrame?.(frame) - } - - const sendEncrypted = createMobileDirectRpcSender({ - getOutbound: () => outbound, - getSharedKey: () => sharedKey, - getSocket: () => ws, - getState: () => state, - onSocketDesync: (socket) => { - synthesizedCloses.remember(socket, authenticationGeneration) - handleSocketClosed(socket, { timedOut: false }) - } - }) - - function sendBrowserScreencastUnsubscribe(subscriptionId: string): void { - sendEncrypted({ - id: nextId(), - deviceToken, - method: 'browser.screencast.unsubscribe', - params: { subscriptionId } - }) - } - - function sendServerSubscriptionUnsubscribe(stream: StreamRequest): void { - if (!stream.subscriptionId) { - return - } - const unsubscribe = buildServerSubscriptionUnsubscribe(stream.method, stream.subscriptionId) - if (unsubscribe) { - sendEncrypted({ - id: nextId(), - deviceToken, - method: unsubscribe.method, - params: unsubscribe.params - }) - } - } - - const livenessWatchdog = new RpcSessionLivenessWatchdog({ - transport: 'direct', - sendProbe: (identity) => { - if (identity !== livenessIdentity || state !== 'connected') { - return false - } - return sendEncrypted({ - id: `${LIVENESS_REQUEST_ID_PREFIX}${nextId()}`, - deviceToken, - method: 'status.get' - }) - }, - terminate: (identity) => { - if (identity === livenessIdentity && livenessIdentity.socket === ws) { - closeAndSynthesize(livenessIdentity.socket) - } - } - }) - - openConnection() - - return { - async sendRequest( - method: string, - params?: unknown, - options?: RpcClientSendRequestOptions - ): Promise { - const budget = openRpcRequestBudget(options) - const waitStart = budget.startedAt - const wasConnected = state === 'connected' - if (options?.failWhenDisconnected && !wasConnected) { - throw new Error(`Not connected: ${method}`) - } - await waitForConnected(options?.timeoutMs) - if (!wasConnected) { - console.log('[net] sendRequest waited for connect', { - method, - waitedMs: Date.now() - waitStart - }) - } - - return new Promise((resolve, reject) => { - const id = nextId() - const timeoutMs = resolvePostConnectRequestTimeout(budget, REQUEST_TIMEOUT_MS) - const timeout = setTimeout(() => { - pending.delete(id) - console.log('[net] sendRequest TIMEOUT', { - method, - timeoutMs, - state - }) - // Why: the frame was written 30s ago — the host may have processed it. - reject(markRpcDeliveryUnknown(new Error(`Request timed out: ${method}`))) - }, timeoutMs) - - pending.set(id, { - resolve: (response) => { - clearTimeout(timeout) - resolve(response) - }, - reject: (error) => { - clearTimeout(timeout) - reject(error) - } - }) - - if (!sendEncrypted({ id, deviceToken, method, params })) { - pending.delete(id) - clearTimeout(timeout) - reject(new Error('Connection interrupted')) - } - }) - }, - - subscribe( - method: string, - params: unknown, - onData: StreamingListener, - options?: RpcClientSubscribeOptions - ): () => void { - const id = nextId() - const stream: StreamRequest = { - method, - params, - listener: onData, - onBinaryFrame: options?.onBinaryFrame, - onTerminalBinaryFrame: options?.onTerminalBinaryFrame - } - streamListeners.set(id, stream) - if (method === 'browser.screencast') { - if (activeBrowserScreencastRequestId && activeBrowserScreencastRequestId !== id) { - disposeBrowserScreencastStream(activeBrowserScreencastRequestId) - } - if (pendingBrowserScreencastRequestId && pendingBrowserScreencastRequestId !== id) { - disposeBrowserScreencastStream(pendingBrowserScreencastRequestId) - } - // Why: screencast frames carry no stream id, so route only after the new stream's ready to drop stale old-page pixels. - pendingBrowserScreencastRequestId = id - activeBrowserScreencastRequestId = null - } - - if (state === 'connected') { - if (sendEncrypted({ id, deviceToken, method, params })) { - stream.sent = true - } else { - emitStreamError(stream, 'Connection interrupted') - removeStreamListener(id) - } - } else { - // Registered now; the outbound subscribe is (re-)sent once the channel reaches 'connected'. - console.log('[net] subscribe queued — waiting for connected', { method, state }) - } - - return () => { - const stream = streamListeners.get(id) - if (stream?.method === 'browser.screencast') { - disposeBrowserScreencastStream(id) - return - } - if ( - stream?.method === 'runtime.clientEvents.subscribe' || - stream?.method === 'accounts.subscribe' || - stream?.method === 'files.watch' - ) { - disposeServerSubscriptionStream(id, stream) - return - } - if (stream?.method === 'terminal.subscribe') { - // Why: server keys cleanup by composite `${terminal}:${clientId}` so two phones don't evict each other. See docs/mobile-presence-lock.md. - const unsubscribeParams = buildTerminalUnsubscribeParams(stream.params) - if (unsubscribeParams) { - sendEncrypted({ - id: nextId(), - deviceToken, - method: 'terminal.unsubscribe', - params: unsubscribeParams - }) - } - } else { - const unsub = buildStreamUnsubscribe(stream?.method, stream?.params) - if (unsub) { - sendEncrypted({ id: nextId(), deviceToken, method: unsub.method, params: unsub.params }) - } - } - removeStreamListener(id) - } - }, - updateTerminalSubscriptionViewport( - terminal: string, - viewport: { cols: number; rows: number } - ): void { - updateCachedTerminalSubscriptionViewport(streamListeners.values(), terminal, viewport) - }, - sendTerminalBinaryFrame(frame: TerminalStreamFrame): boolean { - return sendMobileTerminalBinaryFrame({ - frame, - socket: ws, - sharedKey, - isConnected: state === 'connected', - onSocketClosed: handleSocketClosed - }) - }, - getState(): ConnectionState { - return state - }, - getReconnectAttempt(): number { - return reconnectAttempt - }, - getLastConnectedAt(): number | null { - return lastConnectedAt - }, - getLastInboundAt(): number | null { - return livenessWatchdog.getLastInboundAt() || null - }, - onStateChange(listener: (state: ConnectionState) => void): () => void { - stateListeners.add(listener) - return () => stateListeners.delete(listener) - }, - notifyForeground(_reason?: ForegroundNudgeReason): void { - if (intentionallyClosed) { - return - } - if (state === 'connected') { - // Why: resume probes now; three fair misses detect a half-open socket within 24s. - console.log('[net] foreground — probing live connection') - if (livenessIdentity) { - livenessWatchdog.probeNow(livenessIdentity) - } - return - } - const dialing = ws - const dialAgeMs = Date.now() - dialStartedAt - let abandoned = false - if (dialing && isStaleForegroundDial(state, dialAgeMs)) { - console.log('[net] foreground — abandoning stale dial', { state, dialAgeMs }) - closeAndSynthesize(dialing) - abandoned = true - } - if (state === 'reconnecting') { - // Why: foreground is a strong user signal — restart immediately instead of waiting out a 60s/90s backoff timer. - console.log('[net] foreground — restarting reconnect loop', { - attempt: reconnectAttempt, - hadTimer: !!reconnectTimer - }) - // Why: an abandoned dial keeps the failure it already represents. It never - // authenticated, so it is the same failure the connect timeout would have - // booked had we waited it out — we skip the wait, we don't pardon it. Zeroing - // there would let a resume (or a flapping network) reset the counter faster - // than it climbs, pinning the card at "Connecting…" through a real outage - // (issue #10119). A redial with no dial to abandon is a genuinely fresh start. - redialNow(!abandoned) - } - }, - - close() { - intentionallyClosed = true - if (reconnectTimer) { - clearTimeout(reconnectTimer) - reconnectTimer = null - } - clearConnectTimer() - clearHandshakeTimer() - disposeActiveOutbound() - if (livenessIdentity) { - livenessWatchdog.stop(livenessIdentity) - livenessIdentity = null - } - if (ws) { - ws.close() - ws = null - } - sharedKey = null - setState('disconnected') - // Why: closing the client cannot retract request frames already written. - rejectAllPending('Client closed', { deliveryUnknown: true }) - } - } + return new DirectRpcClient(endpoint, deviceToken, serverPublicKeyB64, options) }