mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
refactor(mobile): restore main's RPC client split under the hybrid transport
Replace the 1,202-line rpc-client.ts closure with main's DirectRpcClient decomposition (12 restored modules) and re-apply the branch's 33-hunk delta onto it: the outbound backpressure queue and process memory budget, the inbound frame queue with overflow/oversize close, JSON structure admission, the extracted binary-frame handler, terminal multiplex routing with sendTerminalBinaryFrame, accounts/files server-subscription teardown, and protocol-only endpoint redaction. Deletes rpc-client-contract.ts, rpc-client-request-state.ts and rpc-client-subscribe-options.ts, which duplicated main's rpc-client.ts and stream registry. Four helpers keep every restored module under the 300-line cap without a suppression, so rpc-client.ts drops its max-lines override. Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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<RpcClient, 'sendRequest' | 'subscribe'>
|
||||
@@ -46,7 +46,7 @@ export function createMobileBrowserRpcClient(
|
||||
async function sendRequest(
|
||||
method: string,
|
||||
rawParams: unknown = {},
|
||||
_options?: RpcClientSendRequestOptions
|
||||
_options?: SendRequestOptions
|
||||
): Promise<RpcResponse> {
|
||||
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}` })
|
||||
|
||||
@@ -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<RpcResponse> {
|
||||
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<void> {
|
||||
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()}`
|
||||
}
|
||||
}
|
||||
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
@@ -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<typeof setTimeout> | 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<void> {
|
||||
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()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<RpcResponse>
|
||||
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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
@@ -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<typeof setTimeout> | 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
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<typeof setTimeout> | null
|
||||
}
|
||||
@@ -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<void>
|
||||
sendEncrypted: (request: unknown) => boolean
|
||||
deviceToken: string
|
||||
}
|
||||
|
||||
export class RpcClientRequestTracker {
|
||||
private readonly pending = new Map<string, PendingRequest>()
|
||||
|
||||
constructor(private readonly options: RequestTrackerOptions) {}
|
||||
|
||||
async sendRequest(
|
||||
method: string,
|
||||
params?: unknown,
|
||||
requestOptions?: SendRequestOptions
|
||||
): Promise<RpcResponse> {
|
||||
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
|
||||
}
|
||||
}
|
||||
@@ -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()
|
||||
})
|
||||
})
|
||||
@@ -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<ConnectionLogEntry, 'code' | 'path'>
|
||||
) => 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()
|
||||
}
|
||||
}
|
||||
@@ -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()
|
||||
}
|
||||
}
|
||||
@@ -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<typeof createMobileDirectRpcOutbound>
|
||||
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> | 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()
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -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<typeof setTimeout> | 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()
|
||||
}
|
||||
}
|
||||
@@ -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' }
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -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<string, StreamRequest>()
|
||||
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<string, unknown> | 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 })
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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<number, TerminalListener>()
|
||||
private readonly idsByRequest = new Map<string, Set<number>>()
|
||||
private readonly snapshots = new Map<number, TerminalSnapshotState>()
|
||||
|
||||
register(requestId: string, streamId: number, listener: TerminalListener): void {
|
||||
const ids = this.idsByRequest.get(requestId) ?? new Set<number>()
|
||||
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)
|
||||
})
|
||||
}
|
||||
}
|
||||
+41
-1179
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user