mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 08:01:56 +00:00
test(mobile): prove uncertain input stays fenced across host downgrade
This commit is contained in:
@@ -623,6 +623,7 @@ jobs:
|
||||
tests/e2e/cross-version-wire/release-checkout.unit.test.ts
|
||||
tests/e2e/cross-version-wire/versioned-mobile-terminal-wire.unit.test.ts
|
||||
tests/e2e/cross-version-wire/cross-version-mobile-input.unit.test.ts
|
||||
tests/e2e/cross-version-wire/cross-version-mobile-input-recovery.unit.test.ts
|
||||
tests/e2e/cross-version-wire/cross-version-browser-placement.unit.test.ts
|
||||
tests/e2e/cross-version-wire/cross-version-terminal-wire.unit.test.ts
|
||||
tests/e2e/cross-version-wire/reported-lossy-initial-snapshot.unit.test.ts
|
||||
|
||||
@@ -165,8 +165,14 @@ the released app's full input-routing stack or released writer behavior.
|
||||
|
||||
The harness does **not** cover the session-tab sync channel, legacy agent-session
|
||||
publications, file or Git RPCs, physical mobile keyboards, or the public relay service.
|
||||
The mobile journey does not yet cover reconnect or uncertain delivery. Changes on those
|
||||
paths still need their own reasoning against the three rules above.
|
||||
`cross-version-mobile-input-recovery.unit.test.ts` retains the production logical client
|
||||
through current → pre-ordered-input host → current encrypted sessions. After prefix bytes
|
||||
reach the PTY but their receipt is dropped, input and JSON fallback stay fenced until an
|
||||
explicit compatible recovery. Final exact bytes prove no replay or stray Enter. The legacy
|
||||
host is pinned to v1.4.197 because this scenario specifically requires no ordered-input
|
||||
support. Session assembly and receipt loss are fixture-owned; this is not public-relay
|
||||
reconnection or recovery-UI coverage. Changes on uncovered paths still need their own
|
||||
reasoning against the three rules above.
|
||||
|
||||
## Worked example: `agentWait` on terminal and worker reads
|
||||
|
||||
|
||||
@@ -0,0 +1,123 @@
|
||||
import { randomBytes } from 'node:crypto'
|
||||
import { expect, it, vi } from 'vitest'
|
||||
import { createStableLogicalRpcClient } from '../../../mobile/src/transport/stable-logical-rpc-client'
|
||||
import {
|
||||
createOrderedInputPtyTestRig,
|
||||
inputProofDeadline
|
||||
} from '../../../src/main/runtime/rpc/terminal-ordered-input-pty-test-rig'
|
||||
import { loadMobileTerminalWireBuild } from './versioned-mobile-terminal-wire'
|
||||
import { WORKING_TREE } from './versioned-terminal-wire'
|
||||
import { openMobileInputWireSession } from './mobile-input-wire-session'
|
||||
|
||||
vi.mock('expo-crypto', () => ({
|
||||
getRandomBytes: (length: number) => new Uint8Array(randomBytes(length))
|
||||
}))
|
||||
|
||||
// This scenario requires an actual pre-ordered-input host, unlike the rolling happy-path baseline.
|
||||
const PRE_ORDERED_INPUT_HOST = 'v1.4.197'
|
||||
|
||||
it('fences uncertain input across a released-host downgrade until explicit compatible recovery', async () => {
|
||||
const [current, oldHost] = await Promise.all([
|
||||
loadMobileTerminalWireBuild(WORKING_TREE),
|
||||
loadMobileTerminalWireBuild(PRE_ORDERED_INPUT_HOST)
|
||||
])
|
||||
const expected = Buffer.from('prefix한글fresh\r')
|
||||
const rig = await createOrderedInputPtyTestRig(expected)
|
||||
const sessions: Awaited<ReturnType<typeof openMobileInputWireSession>>[] = []
|
||||
let logical: ReturnType<typeof createStableLogicalRpcClient> | undefined
|
||||
const subscriptions: { capabilities?: { orderedInput?: unknown } }[] = []
|
||||
try {
|
||||
const initial = await openMobileInputWireSession(current, current, rig.runtime, {
|
||||
dropInputReceipts: true,
|
||||
connectionId: 'recovery-initial'
|
||||
})
|
||||
sessions.push(initial)
|
||||
logical = createStableLogicalRpcClient(initial.rpc, 'relay')
|
||||
logical.subscribe(
|
||||
'terminal.subscribe',
|
||||
{
|
||||
terminal: 'terminal-1',
|
||||
client: { id: 'phone', type: 'mobile' },
|
||||
capabilities: { terminalBinaryStream: 1 }
|
||||
},
|
||||
(event) => {
|
||||
const result = event as { type?: string; capabilities?: { orderedInput?: unknown } }
|
||||
if (result.type === 'subscribed') {
|
||||
subscriptions.push(result)
|
||||
}
|
||||
}
|
||||
)
|
||||
await vi.waitFor(() => expect(subscriptions).toHaveLength(1))
|
||||
expect(logical.supportsTerminalStreamInput?.('terminal-1')).toBe(true)
|
||||
const prefix = logical.sendTerminalStreamInput!('terminal-1', 'prefix한글')!
|
||||
await vi.waitFor(() => {
|
||||
expect(initial.counts().droppedReceipts).toBe(1)
|
||||
expect(rig.bytes()).toEqual(Buffer.from('prefix한글'))
|
||||
})
|
||||
logical.suspendActiveSession()
|
||||
expect(await inputProofDeadline(prefix, 'uncertain prefix settlement')).toBe(false)
|
||||
expect(logical.getTerminalStreamInputFailure?.('terminal-1')).toMatchObject({
|
||||
outcome: 'unknown'
|
||||
})
|
||||
|
||||
const legacy = await openMobileInputWireSession(oldHost, current, rig.runtime, {
|
||||
connectionId: 'recovery-legacy'
|
||||
})
|
||||
sessions.push(legacy)
|
||||
const legacyJson = vi.spyOn(legacy.rpc, 'sendRequest')
|
||||
const legacyStream = vi.spyOn(legacy.rpc, 'sendTerminalStreamInput')
|
||||
await logical.migrateTo(legacy.rpc, 'relay')
|
||||
await vi.waitFor(() => expect(subscriptions).toHaveLength(2))
|
||||
expect(subscriptions[1]?.capabilities?.orderedInput, PRE_ORDERED_INPUT_HOST).toBeUndefined()
|
||||
expect(legacy.rpc.supportsTerminalStreamInput?.('terminal-1')).toBe(false)
|
||||
expect(logical.recoverTerminalStreamInput?.('terminal-1')).toBe(false)
|
||||
const enter = logical.sendTerminalStreamInput!('terminal-1', '\r')
|
||||
expect(enter).not.toBeNull()
|
||||
expect(await enter).toBe(false)
|
||||
await expect(
|
||||
logical.sendRequest('terminal.send', { terminal: 'terminal-1', text: '\r' })
|
||||
).rejects.toThrow('Terminal input stopped')
|
||||
expect(legacyJson).not.toHaveBeenCalled()
|
||||
expect(legacyStream).not.toHaveBeenCalled()
|
||||
expect(legacy.counts()).toEqual({ binaryInputs: 0, jsonInputs: 0, droppedReceipts: 0 })
|
||||
|
||||
const compatible = await openMobileInputWireSession(current, current, rig.runtime, {
|
||||
connectionId: 'recovery-compatible'
|
||||
})
|
||||
sessions.push(compatible)
|
||||
const compatibleStream = vi.spyOn(compatible.rpc, 'sendTerminalStreamInput')
|
||||
const compatibleJson = vi.spyOn(compatible.rpc, 'sendRequest')
|
||||
await logical.migrateTo(compatible.rpc, 'relay')
|
||||
await vi.waitFor(() => expect(subscriptions).toHaveLength(3))
|
||||
expect(compatible.rpc.supportsTerminalStreamInput?.('terminal-1')).toBe(true)
|
||||
expect(await logical.sendTerminalStreamInput!('terminal-1', '\r')).toBe(false)
|
||||
expect(compatible.counts().binaryInputs).toBe(0)
|
||||
expect(compatibleStream).not.toHaveBeenCalled()
|
||||
await expect(
|
||||
logical.sendRequest('terminal.send', { terminal: 'terminal-1', text: '\r' })
|
||||
).rejects.toThrow('Terminal input stopped')
|
||||
expect(compatibleJson).not.toHaveBeenCalled()
|
||||
expect(logical.recoverTerminalStreamInput?.('terminal-1')).toBe(true)
|
||||
expect(
|
||||
await inputProofDeadline(
|
||||
logical.sendTerminalStreamInput!('terminal-1', 'fresh\r')!,
|
||||
'fresh receipt'
|
||||
)
|
||||
).toBe(true)
|
||||
await inputProofDeadline(rig.inputDelivered, 'recovered PTY bytes')
|
||||
expect(rig.bytes()).toEqual(expected)
|
||||
expect(initial.counts()).toEqual({ binaryInputs: 1, jsonInputs: 0, droppedReceipts: 1 })
|
||||
expect(compatible.counts()).toEqual({ binaryInputs: 1, jsonInputs: 0, droppedReceipts: 0 })
|
||||
for (const session of sessions) {
|
||||
expect(session.errors).toEqual([])
|
||||
}
|
||||
} finally {
|
||||
try {
|
||||
logical?.close()
|
||||
const cleanup = await Promise.allSettled(sessions.map((session) => session.dispose()))
|
||||
expect(cleanup.filter((result) => result.status === 'rejected')).toEqual([])
|
||||
} finally {
|
||||
await rig.close()
|
||||
}
|
||||
}
|
||||
}, 60_000)
|
||||
@@ -1,153 +1,30 @@
|
||||
import { expect } from 'vitest'
|
||||
import nacl from 'tweetnacl'
|
||||
import WebSocket, { WebSocketServer } from 'ws'
|
||||
import { RelayPendingRequests } from '../../../mobile/src/transport/relay-pending-requests'
|
||||
import {
|
||||
createOrderedInputPtyTestRig,
|
||||
inputProofDeadline
|
||||
} from '../../../src/main/runtime/rpc/terminal-ordered-input-pty-test-rig'
|
||||
import type { TerminalStreamFrame } from './versioned-terminal-wire'
|
||||
import type { MobileTerminalWireBuild as MobileInputWireBuild } from './versioned-mobile-terminal-wire'
|
||||
import type { MobileTerminalWireBuild } from './versioned-mobile-terminal-wire'
|
||||
import { openMobileInputWireSession } from './mobile-input-wire-session'
|
||||
|
||||
const INPUTS = ['BOM:\ufeff한글\nline-2', '\t\x1b[A\x7f\x03', '\r']
|
||||
|
||||
export async function runMobileInputSkewJourney(
|
||||
host: MobileInputWireBuild,
|
||||
client: MobileInputWireBuild
|
||||
host: MobileTerminalWireBuild,
|
||||
client: MobileTerminalWireBuild
|
||||
) {
|
||||
const expected = Buffer.from(INPUTS.join(''))
|
||||
const rig = await createOrderedInputPtyTestRig(expected)
|
||||
const server = new WebSocketServer({ host: '127.0.0.1', port: 0, perMessageDeflate: false })
|
||||
const keys = nacl.box.keyPair()
|
||||
const abort = new AbortController()
|
||||
const pending = new RelayPendingRequests()
|
||||
const dispatches: Promise<unknown>[] = []
|
||||
const channels: InstanceType<MobileInputWireBuild['E2EEChannel']>[] = []
|
||||
const errors: unknown[] = []
|
||||
let socket: WebSocket | undefined
|
||||
let mobile: InstanceType<MobileInputWireBuild['MobileE2EEV2PhysicalChannel']> | undefined
|
||||
let streams: InstanceType<MobileInputWireBuild['MobileRelayRpcStreams']> | undefined
|
||||
let session: Awaited<ReturnType<typeof openMobileInputWireSession>> | undefined
|
||||
let unsubscribe: (() => void) | undefined
|
||||
let binaryInputs = 0
|
||||
let jsonInputs = 0
|
||||
let offered: unknown
|
||||
let echo: unknown
|
||||
try {
|
||||
await inputProofDeadline(
|
||||
new Promise<void>((resolve, reject) => {
|
||||
server.once('listening', resolve)
|
||||
server.once('error', reject)
|
||||
}),
|
||||
'mobile skew listener'
|
||||
)
|
||||
const dispatcher = new host.RpcDispatcher({
|
||||
runtime: rig.runtime,
|
||||
methods: host.TERMINAL_METHODS
|
||||
})
|
||||
server.on('connection', (peer) => {
|
||||
const handlers = new Map<number, (frame: TerminalStreamFrame) => void>()
|
||||
const channel = new host.E2EEChannel(peer, {
|
||||
serverSecretKey: keys.secretKey,
|
||||
resolveAuthenticatedDevice: (token) =>
|
||||
token === 'skew-device'
|
||||
? { deviceId: 'phone', deviceToken: token, scope: 'mobile' }
|
||||
: null,
|
||||
transportContext: { transport: 'relay', relayHostId: 'AbCdEf0123_-xyZ9' },
|
||||
requireV2: true,
|
||||
onReady: () => {},
|
||||
onError: (code, reason) => errors.push({ code, reason })
|
||||
})
|
||||
channels.push(channel)
|
||||
channel.onMessage((text, reply, sendBinary) => {
|
||||
const request = JSON.parse(text)
|
||||
if (request.method === 'terminal.subscribe') {
|
||||
offered = request.params?.capabilities
|
||||
}
|
||||
if (request.method === 'terminal.send') {
|
||||
jsonInputs++
|
||||
}
|
||||
const work = dispatcher.dispatchStreaming(request, reply, {
|
||||
connectionId: 'mobile-skew',
|
||||
signal: abort.signal,
|
||||
clientKind: 'mobile',
|
||||
sendBinary,
|
||||
registerBinaryStreamHandler: (id, handler) => {
|
||||
handlers.set(id, handler)
|
||||
return () => {
|
||||
handlers.delete(id)
|
||||
}
|
||||
}
|
||||
})
|
||||
dispatches.push(work)
|
||||
void work.catch((error) => errors.push(error))
|
||||
})
|
||||
channel.onBinaryMessage((bytes) => {
|
||||
const frame = host.codec.decodeTerminalStreamFrame(bytes)
|
||||
if (!frame) {
|
||||
errors.push(new Error('Host refused a mobile frame'))
|
||||
return
|
||||
}
|
||||
if (frame.opcode === host.codec.TerminalStreamOpcode.Input) {
|
||||
binaryInputs++
|
||||
}
|
||||
handlers.get(frame.streamId)?.(frame)
|
||||
})
|
||||
peer.on('message', (raw, binary) =>
|
||||
channel.handleRawMessage(binary ? new Uint8Array(raw as Buffer) : raw.toString())
|
||||
)
|
||||
peer.on('error', (error) => errors.push(error))
|
||||
})
|
||||
const address = server.address()
|
||||
if (!address || typeof address === 'string') {
|
||||
throw new Error('Missing skew server port')
|
||||
}
|
||||
socket = new WebSocket(`ws://127.0.0.1:${address.port}`)
|
||||
await inputProofDeadline(
|
||||
new Promise<void>((resolve, reject) => {
|
||||
socket!.once('open', resolve)
|
||||
socket!.once('error', reject)
|
||||
}),
|
||||
'mobile skew socket'
|
||||
)
|
||||
let authenticated!: () => void
|
||||
const ready = new Promise<void>((resolve) => {
|
||||
authenticated = resolve
|
||||
})
|
||||
streams = new client.MobileRelayRpcStreams({
|
||||
nextId: () => pending.nextId(),
|
||||
waitForConnected: () => ready,
|
||||
sendFrame: (request) => mobile!.sendText(JSON.stringify(request)),
|
||||
sendBinary: (bytes) => mobile!.sendBinary(bytes)
|
||||
})
|
||||
mobile = new client.MobileE2EEV2PhysicalChannel({
|
||||
session: client.MobileE2EEV2ClientSession.create({
|
||||
desktopPublicKeyB64: Buffer.from(keys.publicKey).toString('base64'),
|
||||
transport: 'relay',
|
||||
relayHostId: 'AbCdEf0123_-xyZ9'
|
||||
}),
|
||||
socket,
|
||||
deviceToken: 'skew-device',
|
||||
decodeBinary: async (raw) => (raw instanceof Uint8Array ? raw : null),
|
||||
onAuthenticated: authenticated,
|
||||
onText: (text) => {
|
||||
const response = JSON.parse(text)
|
||||
if (!pending.settle(response)) {
|
||||
streams!.handleResponse(response)
|
||||
}
|
||||
},
|
||||
onBinary: (bytes) => streams!.handleBinary(bytes),
|
||||
onError: (error) => errors.push(error)
|
||||
})
|
||||
socket.on('message', (raw, binary) => {
|
||||
void mobile!.handleMessage(binary ? new Uint8Array(raw as Buffer) : raw.toString())
|
||||
})
|
||||
mobile.start()
|
||||
await inputProofDeadline(ready, 'mobile skew authentication')
|
||||
session = await openMobileInputWireSession(host, client, rig.runtime)
|
||||
const { rpc } = session
|
||||
let subscribed!: () => void
|
||||
const subscription = new Promise<void>((resolve) => {
|
||||
subscribed = resolve
|
||||
})
|
||||
unsubscribe = streams.subscribe(
|
||||
unsubscribe = rpc.subscribe(
|
||||
'terminal.subscribe',
|
||||
{
|
||||
terminal: 'terminal-1',
|
||||
@@ -163,83 +40,43 @@ export async function runMobileInputSkewJourney(
|
||||
}
|
||||
)
|
||||
await inputProofDeadline(subscription, 'mobile skew subscription')
|
||||
const ordered = streams.supportsTerminalStreamInput?.('terminal-1') === true
|
||||
const ordered = rpc.supportsTerminalStreamInput?.('terminal-1') === true
|
||||
for (const text of INPUTS) {
|
||||
if (ordered) {
|
||||
expect(
|
||||
await inputProofDeadline(
|
||||
streams.sendTerminalStreamInput!('terminal-1', text)!,
|
||||
rpc.sendTerminalStreamInput!('terminal-1', text)!,
|
||||
'mobile skew receipt'
|
||||
)
|
||||
).toBe(true)
|
||||
} else {
|
||||
const id = pending.nextId()
|
||||
const response = new Promise<unknown>((resolve, reject) => {
|
||||
const timer = setTimeout(() => {
|
||||
pending.drop(id)
|
||||
reject(new Error('Mobile skew JSON timeout'))
|
||||
}, 5000)
|
||||
pending.track(id, {
|
||||
resolve,
|
||||
reject,
|
||||
timer
|
||||
})
|
||||
if (
|
||||
!mobile!.sendText(
|
||||
JSON.stringify({
|
||||
id,
|
||||
method: 'terminal.send',
|
||||
params: { terminal: 'terminal-1', text }
|
||||
})
|
||||
)
|
||||
) {
|
||||
clearTimeout(timer)
|
||||
pending.drop(id)
|
||||
reject(new Error('Mobile skew JSON send failed'))
|
||||
}
|
||||
})
|
||||
expect(await response).toMatchObject({ ok: true })
|
||||
expect(
|
||||
await rpc.sendRequest('terminal.send', { terminal: 'terminal-1', text })
|
||||
).toMatchObject({ ok: true })
|
||||
}
|
||||
}
|
||||
await inputProofDeadline(rig.inputDelivered, 'mobile skew PTY delivery')
|
||||
expect(rig.bytes()).toEqual(expected)
|
||||
const { binaryInputs, jsonInputs } = session.counts()
|
||||
expect(binaryInputs).toBe(ordered ? INPUTS.length : 0)
|
||||
expect(jsonInputs).toBe(ordered ? 0 : INPUTS.length)
|
||||
expect(errors).toEqual([])
|
||||
expect(session.errors).toEqual([])
|
||||
return {
|
||||
host: host.revision,
|
||||
client: client.revision,
|
||||
ordered,
|
||||
offered,
|
||||
offered: session.offered(),
|
||||
echo,
|
||||
binaryInputs,
|
||||
jsonInputs,
|
||||
hex: rig.bytes().toString('hex')
|
||||
}
|
||||
} finally {
|
||||
unsubscribe?.()
|
||||
streams?.clear()
|
||||
pending.rejectAll(new Error('Mobile skew cleanup'))
|
||||
abort.abort()
|
||||
mobile?.dispose()
|
||||
socket?.terminate()
|
||||
for (const channel of channels) {
|
||||
channel.destroy()
|
||||
}
|
||||
for (const peer of server.clients) {
|
||||
peer.terminate()
|
||||
}
|
||||
try {
|
||||
await inputProofDeadline(Promise.allSettled(dispatches), 'mobile skew dispatcher cleanup')
|
||||
unsubscribe?.()
|
||||
await session?.dispose()
|
||||
} finally {
|
||||
try {
|
||||
await inputProofDeadline(
|
||||
new Promise<void>((resolve) => server.close(() => resolve())),
|
||||
'mobile skew server cleanup'
|
||||
)
|
||||
} finally {
|
||||
await rig.close()
|
||||
}
|
||||
await rig.close()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,236 @@
|
||||
import type { RpcClient } from '../../../mobile/src/transport/rpc-client'
|
||||
import type { ConnectionState } from '../../../mobile/src/transport/types'
|
||||
import nacl from 'tweetnacl'
|
||||
import WebSocket, { WebSocketServer } from 'ws'
|
||||
import { RelayPendingRequests } from '../../../mobile/src/transport/relay-pending-requests'
|
||||
import { inputProofDeadline } from '../../../src/main/runtime/rpc/terminal-ordered-input-pty-test-rig'
|
||||
import type { TerminalStreamFrame } from './versioned-terminal-wire'
|
||||
import type { MobileTerminalWireBuild as MobileInputWireBuild } from './versioned-mobile-terminal-wire'
|
||||
|
||||
export async function openMobileInputWireSession(
|
||||
host: MobileInputWireBuild,
|
||||
client: MobileInputWireBuild,
|
||||
runtime: unknown,
|
||||
options: { dropInputReceipts?: boolean; connectionId?: string } = {}
|
||||
) {
|
||||
const server = new WebSocketServer({ host: '127.0.0.1', port: 0, perMessageDeflate: false })
|
||||
const keys = nacl.box.keyPair()
|
||||
const abort = new AbortController()
|
||||
const pending = new RelayPendingRequests()
|
||||
const dispatches: Promise<unknown>[] = []
|
||||
const channels: InstanceType<MobileInputWireBuild['E2EEChannel']>[] = []
|
||||
const errors: unknown[] = []
|
||||
let socket: WebSocket | undefined
|
||||
let mobile: InstanceType<MobileInputWireBuild['MobileE2EEV2PhysicalChannel']> | undefined
|
||||
let streams: InstanceType<MobileInputWireBuild['MobileRelayRpcStreams']> | undefined
|
||||
let binaryInputs = 0
|
||||
let jsonInputs = 0
|
||||
let offered: unknown
|
||||
let closed = false
|
||||
let droppedReceipts = 0
|
||||
const listeners = new Set<(state: ConnectionState) => void>()
|
||||
const close = () => {
|
||||
if (closed) {
|
||||
return
|
||||
}
|
||||
closed = true
|
||||
streams?.clear()
|
||||
pending.rejectAll(new Error('Mobile skew cleanup'))
|
||||
abort.abort()
|
||||
mobile?.dispose()
|
||||
socket?.terminate()
|
||||
for (const channel of channels) {
|
||||
channel.destroy()
|
||||
}
|
||||
for (const peer of server.clients) {
|
||||
peer.terminate()
|
||||
}
|
||||
for (const listener of listeners) {
|
||||
listener('disconnected')
|
||||
}
|
||||
}
|
||||
const dispose = async () => {
|
||||
close()
|
||||
try {
|
||||
await inputProofDeadline(Promise.allSettled(dispatches), 'mobile skew dispatcher cleanup')
|
||||
} finally {
|
||||
await inputProofDeadline(
|
||||
new Promise<void>((resolve) => server.close(() => resolve())),
|
||||
'mobile skew server cleanup'
|
||||
)
|
||||
}
|
||||
}
|
||||
try {
|
||||
await inputProofDeadline(
|
||||
new Promise<void>((resolve, reject) => {
|
||||
server.once('listening', resolve)
|
||||
server.once('error', reject)
|
||||
}),
|
||||
'mobile skew listener'
|
||||
)
|
||||
const dispatcher = new host.RpcDispatcher({
|
||||
runtime,
|
||||
methods: host.TERMINAL_METHODS
|
||||
})
|
||||
server.on('connection', (peer) => {
|
||||
const handlers = new Map<number, (frame: TerminalStreamFrame) => void>()
|
||||
const channel = new host.E2EEChannel(peer, {
|
||||
serverSecretKey: keys.secretKey,
|
||||
resolveAuthenticatedDevice: (token) =>
|
||||
token === 'skew-device'
|
||||
? { deviceId: 'phone', deviceToken: token, scope: 'mobile' }
|
||||
: null,
|
||||
transportContext: { transport: 'relay', relayHostId: 'AbCdEf0123_-xyZ9' },
|
||||
requireV2: true,
|
||||
onReady: () => {},
|
||||
onError: (code, reason) => errors.push({ code, reason })
|
||||
})
|
||||
channels.push(channel)
|
||||
channel.onMessage((text, reply, sendBinary) => {
|
||||
const request = JSON.parse(text)
|
||||
if (request.method === 'terminal.subscribe') {
|
||||
offered = request.params?.capabilities
|
||||
}
|
||||
if (request.method === 'terminal.send') {
|
||||
jsonInputs++
|
||||
}
|
||||
const work = dispatcher.dispatchStreaming(request, reply, {
|
||||
connectionId: options.connectionId ?? 'mobile-skew',
|
||||
signal: abort.signal,
|
||||
clientKind: 'mobile',
|
||||
sendBinary: (bytes) => {
|
||||
const frame = host.codec.decodeTerminalStreamFrame(bytes)
|
||||
const metadata =
|
||||
frame?.opcode === host.codec.TerminalStreamOpcode.Metadata
|
||||
? host.codec.decodeTerminalStreamJson<{ inputReceipt?: unknown }>(frame.payload)
|
||||
: null
|
||||
if (options.dropInputReceipts && metadata?.inputReceipt) {
|
||||
droppedReceipts++
|
||||
return true
|
||||
}
|
||||
return sendBinary(bytes)
|
||||
},
|
||||
registerBinaryStreamHandler: (id, handler) => {
|
||||
handlers.set(id, handler)
|
||||
return () => {
|
||||
handlers.delete(id)
|
||||
}
|
||||
}
|
||||
})
|
||||
dispatches.push(work)
|
||||
void work.catch((error) => errors.push(error))
|
||||
})
|
||||
channel.onBinaryMessage((bytes) => {
|
||||
const frame = host.codec.decodeTerminalStreamFrame(bytes)
|
||||
if (!frame) {
|
||||
errors.push(new Error('Host refused a mobile frame'))
|
||||
return
|
||||
}
|
||||
if (frame.opcode === host.codec.TerminalStreamOpcode.Input) {
|
||||
binaryInputs++
|
||||
}
|
||||
handlers.get(frame.streamId)?.(frame)
|
||||
})
|
||||
peer.on('message', (raw, binary) =>
|
||||
channel.handleRawMessage(binary ? new Uint8Array(raw as Buffer) : raw.toString())
|
||||
)
|
||||
peer.on('error', (error) => errors.push(error))
|
||||
})
|
||||
const address = server.address()
|
||||
if (!address || typeof address === 'string') {
|
||||
throw new Error('Missing skew server port')
|
||||
}
|
||||
socket = new WebSocket(`ws://127.0.0.1:${address.port}`)
|
||||
await inputProofDeadline(
|
||||
new Promise<void>((resolve, reject) => {
|
||||
socket!.once('open', resolve)
|
||||
socket!.once('error', reject)
|
||||
}),
|
||||
'mobile skew socket'
|
||||
)
|
||||
let authenticated!: () => void
|
||||
const ready = new Promise<void>((resolve) => {
|
||||
authenticated = resolve
|
||||
})
|
||||
streams = new client.MobileRelayRpcStreams({
|
||||
nextId: () => pending.nextId(),
|
||||
waitForConnected: () => ready,
|
||||
sendFrame: (request) => mobile!.sendText(JSON.stringify(request)),
|
||||
sendBinary: (bytes) => mobile!.sendBinary(bytes)
|
||||
})
|
||||
mobile = new client.MobileE2EEV2PhysicalChannel({
|
||||
session: client.MobileE2EEV2ClientSession.create({
|
||||
desktopPublicKeyB64: Buffer.from(keys.publicKey).toString('base64'),
|
||||
transport: 'relay',
|
||||
relayHostId: 'AbCdEf0123_-xyZ9'
|
||||
}),
|
||||
socket,
|
||||
deviceToken: 'skew-device',
|
||||
decodeBinary: async (raw) => (raw instanceof Uint8Array ? raw : null),
|
||||
onAuthenticated: authenticated,
|
||||
onText: (text) => {
|
||||
const response = JSON.parse(text)
|
||||
if (!pending.settle(response)) {
|
||||
streams!.handleResponse(response)
|
||||
}
|
||||
},
|
||||
onBinary: (bytes) => streams!.handleBinary(bytes),
|
||||
onError: (error) => errors.push(error)
|
||||
})
|
||||
socket.on('message', (raw, binary) => {
|
||||
void mobile!.handleMessage(binary ? new Uint8Array(raw as Buffer) : raw.toString())
|
||||
})
|
||||
mobile.start()
|
||||
await inputProofDeadline(ready, 'mobile skew authentication')
|
||||
const rpc: RpcClient = {
|
||||
supportsTerminalStreamInput: streams.supportsTerminalStreamInput?.bind(streams),
|
||||
sendTerminalStreamInput: streams.sendTerminalStreamInput?.bind(streams),
|
||||
getTerminalStreamInputFailure: streams.getTerminalStreamInputFailure?.bind(streams),
|
||||
recoverTerminalStreamInput: streams.recoverTerminalStreamInput?.bind(streams),
|
||||
cancelTerminalStreamInput: streams.cancelTerminalStreamInput?.bind(streams),
|
||||
fenceTerminalStreamInput: streams.fenceTerminalStreamInput?.bind(streams),
|
||||
subscribe: (...args) => streams!.subscribe(...args),
|
||||
updateTerminalSubscriptionViewport: () => {},
|
||||
getState: () => (closed ? 'disconnected' : 'connected'),
|
||||
getReconnectAttempt: () => 0,
|
||||
getLastConnectedAt: () => null,
|
||||
onStateChange: (listener) => {
|
||||
listeners.add(listener)
|
||||
return () => {
|
||||
listeners.delete(listener)
|
||||
}
|
||||
},
|
||||
notifyForeground: () => {},
|
||||
close,
|
||||
sendRequest: (method, params) => {
|
||||
const id = pending.nextId()
|
||||
return new Promise((resolve, reject) => {
|
||||
const timer = setTimeout(() => {
|
||||
pending.drop(id)
|
||||
reject(new Error('Mobile skew JSON timeout'))
|
||||
}, 5000)
|
||||
pending.track(id, { resolve, reject, timer })
|
||||
try {
|
||||
if (!mobile!.sendText(JSON.stringify({ id, method, params }))) {
|
||||
throw new Error('Mobile skew JSON send failed')
|
||||
}
|
||||
} catch (error) {
|
||||
clearTimeout(timer)
|
||||
pending.drop(id)
|
||||
reject(error)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
return {
|
||||
rpc,
|
||||
dispose,
|
||||
errors,
|
||||
counts: () => ({ binaryInputs, jsonInputs, droppedReceipts }),
|
||||
offered: () => offered
|
||||
}
|
||||
} catch (error) {
|
||||
await dispose()
|
||||
throw error
|
||||
}
|
||||
}
|
||||
@@ -19,7 +19,17 @@ export type MobileWireStreams = Pick<
|
||||
MobileRelayRpcStreams,
|
||||
'subscribe' | 'handleResponse' | 'handleBinary' | 'clear'
|
||||
> &
|
||||
Partial<Pick<MobileRelayRpcStreams, 'supportsTerminalStreamInput' | 'sendTerminalStreamInput'>>
|
||||
Partial<
|
||||
Pick<
|
||||
MobileRelayRpcStreams,
|
||||
| 'supportsTerminalStreamInput'
|
||||
| 'sendTerminalStreamInput'
|
||||
| 'getTerminalStreamInputFailure'
|
||||
| 'recoverTerminalStreamInput'
|
||||
| 'cancelTerminalStreamInput'
|
||||
| 'fenceTerminalStreamInput'
|
||||
>
|
||||
>
|
||||
export type MobileWireHostChannel = Pick<
|
||||
E2EEChannel,
|
||||
'onMessage' | 'onBinaryMessage' | 'handleRawMessage' | 'destroy'
|
||||
|
||||
Reference in New Issue
Block a user