Merge host-side review fixes: query-reply guard, files.unwatch, asar dedupe, gzip cache bound

Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb
This commit is contained in:
Jinwoo-H
2026-09-04 02:22:03 -04:00
16 changed files with 509 additions and 25 deletions
+6
View File
@@ -205,6 +205,12 @@ module.exports = {
// Why: out/electron-dev caches `pnpm dev`'s per-branch Electron.app copies (~270MB each).
// CI never creates it, but packaging on a machine that has run dev would pack them all.
'!out/electron-dev{,/**/*}',
// Why: the hosted mobile-web bundle ships through mobileWebExtraResource to
// Resources/mobile-web, which is the only copy resolveMobileWebPackageRoot reads when
// packaged. Packing it here duplicated it, and -export is the pre-hardening Metro
// artifact that still carries the eval/new Function forms the packager strips.
'!out/mobile-web-rnw{,/**/*}',
'!out/mobile-web-rnw-export{,/**/*}',
'!electron.vite.config.{js,ts,mjs,cjs}',
'!{.eslintcache,eslint.config.mjs,.prettierignore,.prettierrc.yaml,CHANGELOG.md,README.md}',
'!{.env,.env.*,.npmrc,pnpm-lock.yaml}',
@@ -84,6 +84,31 @@ describe('electron-builder config', () => {
expect(packs('out/main/examples/index.js')).toBe(true)
})
// Why: `files` is exclusion-only, so both mobile-web trees packed into app.asar on top
// of the Resources/mobile-web copy that resolveMobileWebPackageRoot actually reads.
it('keeps the hosted mobile-web bundle out of app.asar', () => {
const matcher = new FileMatcher('/app', '/dest', (value) => value, electronBuilderConfig.files)
matcher.prependPattern('**/*')
const isPacked = matcher.createFilter()
const packs = (repoPath) => isPacked(join('/app', repoPath), { isDirectory: () => false })
for (const extraResourceOnly of [
'out/mobile-web-rnw/manifest.json',
'out/mobile-web-rnw/assets/index.js',
'out/mobile-web-rnw-export/index.html',
'out/mobile-web-rnw-export/_expo/static/js/web/entry.js'
]) {
expect(packs(extraResourceOnly)).toBe(false)
}
// The extraResources copy is the one the packaged host reads.
for (const platform of ['mac', 'linux', 'win']) {
expect(electronBuilderConfig[platform].extraResources).toContainEqual({
from: 'out/mobile-web-rnw',
to: 'mobile-web'
})
}
})
// Why: out/electron-dev holds `pnpm dev`'s cached Electron.app copies (~270MB per branch).
// CI never creates it, so only a local package would have hit this -- silently, as bulk.
it('keeps cached dev Electron bundles out of app.asar', () => {
@@ -0,0 +1,44 @@
import { readFileSync } from 'node:fs'
import { join, resolve } from 'node:path'
import { describe, expect, it } from 'vitest'
import { parse as parseYaml } from 'yaml'
const projectDir = resolve(import.meta.dirname, '../..')
const SPEC = 'tests/e2e/hosted-mobile-webview-ssh.spec.ts'
// Why: no GitHub runner offers an iOS simulator and a Docker daemon, so this spec is
// manual-only. Pin the facts that keep "unrun" honest rather than leaving them in a YAML
// comment: it is excluded from the lane that would report a green skip, it has a named
// manual entry point, and it refuses to run itself when either prerequisite is missing.
describe('hosted mobile WebView SSH e2e contract', () => {
const spec = readFileSync(join(projectDir, SPEC), 'utf8')
it('stays out of the changed-spec lane that cannot run it', () => {
const workflow = parseYaml(readFileSync(join(projectDir, '.github/workflows/e2e.yml'), 'utf8'))
const changedRun = workflow.jobs['changed-e2e'].steps.find(
(step) => step.name === 'Run changed E2E specs'
)
expect(changedRun.run).toContain(`. != "${SPEC}"`)
})
it('keeps a named manual entry point for both the dev and packaged runs', () => {
const scripts = JSON.parse(readFileSync(join(projectDir, 'package.json'), 'utf8')).scripts
expect(scripts['test:e2e:hosted-mobile-webview:ssh']).toContain(
'run-hosted-mobile-webview-ssh-e2e.mjs'
)
expect(scripts['test:e2e:hosted-mobile-webview:ssh:packaged']).toContain(
'run-packaged-hosted-mobile-webview-ssh-e2e.mjs'
)
})
it('refuses to run itself without a Docker daemon or macOS', () => {
expect(spec).toContain(
"test.skip(!RUN_DOCKER_SSH, 'Set ORCA_E2E_SSH_DOCKER=1 to run Docker-backed SSH tests.')"
)
expect(spec).toContain(
"test.skip(process.platform !== 'darwin', 'Hosted iOS WebView automation requires macOS.')"
)
})
})
@@ -0,0 +1,42 @@
import { describe, expect, it, vi } from 'vitest'
import { handleBrowserScreencastEvent } from './mobile-browser-stream-events'
function handlerArgs() {
return {
busyRef: { current: true },
clearStartupTimer: vi.fn(),
lastZoomResetUrlRef: { current: '' },
resetBrowserZoomState: vi.fn(),
setAddressValue: vi.fn(),
setBusy: vi.fn(),
setDialog: vi.fn(),
setError: vi.fn()
}
}
describe('handleBrowserScreencastEvent', () => {
// Why: the host adds screencast event types without a capability gate, so a client
// that predates one must ignore it rather than throw or disturb the stream's state.
it('ignores an event type it does not recognize', () => {
const args = handlerArgs()
expect(() =>
handleBrowserScreencastEvent({ ...args, event: { type: 'not-a-known-event' } })
).not.toThrow()
expect(args.clearStartupTimer).not.toHaveBeenCalled()
expect(args.setAddressValue).not.toHaveBeenCalled()
expect(args.setBusy).not.toHaveBeenCalled()
expect(args.setDialog).not.toHaveBeenCalled()
expect(args.setError).not.toHaveBeenCalled()
expect(args.busyRef.current).toBe(true)
})
it('still routes a recognized event', () => {
const args = handlerArgs()
handleBrowserScreencastEvent({ ...args, event: { type: 'dialogClosed' } })
expect(args.setDialog).toHaveBeenCalledWith(null)
})
})
@@ -2,6 +2,7 @@ import type {
MobileWebTerminalDeviceInputResult,
MobileWebTerminalRequest
} from '../../../src/shared/mobile-web/terminal-stream-contract'
import { isTerminalQueryReply } from '../../../src/shared/terminal-query-reply'
import type { RpcClient } from '../transport/rpc-client'
import { TerminalStreamOpcode } from '../transport/terminal-stream-protocol'
import { MobileWebBrokerError } from './mobile-web-broker-error'
@@ -27,15 +28,18 @@ export function handleMobileWebTerminalInput(args: {
return enqueueInput(args.record, async () => {
requireLive(args.isLive)
if (args.request.operation === 'input' || args.request.operation === 'queryReply') {
const bytes = atob(args.request.data)
// Why: the page's own label is not evidence; the host drops opcode 18 that fails this same grammar check.
if (args.request.operation === 'queryReply' && !isTerminalQueryReply(bytes)) {
return null
}
sendMobileWebTerminalFrame(
args.client,
args.record,
args.request.operation === 'queryReply'
? args.record.supportsQueryReply
? TerminalStreamOpcode.QueryReply
: TerminalStreamOpcode.Input
args.request.operation === 'queryReply' && args.record.supportsQueryReply
? TerminalStreamOpcode.QueryReply
: TerminalStreamOpcode.Input,
Uint8Array.from(atob(args.request.data), (character) => character.charCodeAt(0))
Uint8Array.from(bytes, (character) => character.charCodeAt(0))
)
return null
}
@@ -268,6 +268,39 @@ describe('MobileWebTerminalStreams', () => {
expect(decodeTerminalStreamText(harness.sentFrames.at(-1)!.payload)).toBe('\x1b[0n')
})
it('drops a page query-reply that is not a reply grammar', async () => {
const harness = createHarness()
await harness.streams.start({
requestId: 'request-bad-query',
subscriptionId: SUBSCRIPTION_ID,
payload: subscribePayload(),
client: harness.client,
isRequestActive: () => true
})
harness.emitMultiplex({ type: 'ready' })
const subscribe = harness.sentFrames.at(-1)!
const hostStreamId = decodeTerminalStreamJson<Record<string, unknown>>(subscribe.payload)!
.streamId as number
harness.emitMultiplex({
type: 'subscribed',
streamId: hostStreamId,
capabilities: { queryReply: 1 }
})
const framesBefore = harness.sentFrames.length
await harness.streams.handle(
{
operation: 'queryReply',
streamId: SUBSCRIPTION_ID,
sequence: 0,
data: Buffer.from(':q!\r').toString('base64')
},
harness.client
)
expect(harness.sentFrames.length).toBe(framesBefore)
})
it('falls back to legacy input when an older host omits query-reply support', async () => {
const harness = createHarness()
await harness.streams.start({
@@ -39,6 +39,8 @@ export type BrowserScreencastSession = {
export type BrowserScreencastEvent =
| { type: 'dialog'; dialogType: string; message: string }
| { type: 'dialogClosed' }
// Rule 3, ungated: clients route screencast results through an unvalidated cast into an
// if/else-if chain with no else, so a build that predates this member ignores it silently.
| {
type: 'navigation'
tab: { url: string; title: string; canGoBack: boolean; canGoForward: boolean }
+35 -5
View File
@@ -51,12 +51,10 @@ const MOBILE_DYNAMIC_RPC_METHODS = [
]
const MOBILE_STREAMING_CLEANUP_RPC_METHODS = [
// Why: shared-control unsubscribe methods are sent from generated cleanup
// paths, so literal mobile source scanning cannot discover every one.
'accounts.unsubscribe',
'browser.screencast.unsubscribe',
// Why: these cleanups are sent from generated paths that literal mobile source
// scanning cannot discover, and they are not in the server-subscription map
// derived below (that map's pairs are asserted from the client's own source).
'notifications.unsubscribe',
'runtime.clientEvents.unsubscribe',
'session.tabs.unsubscribe',
'session.tabs.unsubscribeAll',
'terminal.unsubscribe'
@@ -120,6 +118,25 @@ function mobileRpcAllowlist(): Set<string> {
return new Set([...allowlist[1]!.matchAll(/'([^']+)'/g)].map((match) => match[1]!))
}
/** Subscribe -> unsubscribe pairs read from the client's own teardown map. */
function serverSubscriptionUnsubscribePairs(): [string, string][] {
const source = readFileSync(
join(process.cwd(), 'mobile/src/transport/rpc-client-server-subscription.ts'),
'utf8'
)
const map = source.match(/const unsubscribeMethod = \{([\s\S]*?)\}\[method\]/)
if (!map) {
throw new Error('buildServerSubscriptionUnsubscribe map not found')
}
const pairs = [...map[1]!.matchAll(/'([^']+)':\s*'([^']+)'/g)].map(
(match) => [match[1]!, match[2]!] as [string, string]
)
if (pairs.length === 0) {
throw new Error('buildServerSubscriptionUnsubscribe map is empty')
}
return pairs
}
function registeredRuntimeMethods(): Set<string> {
return new Set(ALL_RPC_METHODS.map((method) => method.name))
}
@@ -150,6 +167,19 @@ describe('mobile RPC allowlist', () => {
expect(missing).toEqual([])
})
it('allows the unsubscribe for every allowlisted server subscription', () => {
// Why: allowlisting a subscribe without its unsubscribe leaks the host-side
// subscription for the socket's life, and the client's cancel swallows the
// forbidden reply. Derive the pairs so the gap class cannot recur.
const allowed = mobileRpcAllowlist()
const registered = registeredRuntimeMethods()
const missing = serverSubscriptionUnsubscribePairs()
.filter(([subscribe]) => allowed.has(subscribe))
.filter(([, unsubscribe]) => !allowed.has(unsubscribe) || !registered.has(unsubscribe))
expect(missing).toEqual([])
})
it('does not grant mobile credentials control over host updates', () => {
const allowed = mobileRpcAllowlist()
expect(
@@ -4,6 +4,7 @@ import {
decodeTerminalStreamText
} from '../../../../../shared/terminal-stream-protocol'
import { isTerminalInputLockedForClient, sendTerminalStreamInput } from './terminal-input-delivery'
import { isAcceptableTerminalQueryReplyFrame } from './terminal-query-reply-guard'
import type { TerminalSubscriptionArgs } from './terminal-legacy-subscription-types'
import { updateViewportForClient } from './terminal-viewport-update'
@@ -52,16 +53,24 @@ export function registerLegacyBinaryControlFrames(
if (isTerminalInputLockedForClient(runtime, ptyId, params.client)) {
return
}
const isQueryReply = frame.opcode === TerminalStreamOpcode.QueryReply
void controls.getDesktopClaimTail().then(async (claimed) => {
if (!claimed || isTerminalInputLockedForClient(runtime, ptyId, params.client)) {
return
}
// Why: opcode 18 skips the mobile input floor, so it needs every guard terminal.send applies; drop otherwise.
if (
isQueryReply &&
!isAcceptableTerminalQueryReplyFrame({ runtime, ptyId, text, client: params.client })
) {
return
}
const outcome = await sendTerminalStreamInput(runtime, {
terminal: params.terminal,
text,
client: params.client,
isMobile,
inputKind: frame.opcode === TerminalStreamOpcode.QueryReply ? 'query-reply' : 'input'
inputKind: isQueryReply ? 'query-reply' : 'input'
})
if (!controls.isClosed() && outcome === 'rejected' && supportsWriteUnavailable) {
controls.sendFrame(TerminalStreamOpcode.WriteUnavailable)
@@ -10,6 +10,7 @@ import {
TerminalMultiplexSourceRangeAckFrame
} from './stream-schemas'
import { isTerminalInputLockedForClient, sendTerminalStreamInput } from './terminal-input-delivery'
import { isAcceptableTerminalQueryReplyFrame } from './terminal-query-reply-guard'
import {
getOutputAfterSnapshotSeq,
normalizeMultiplexSnapshotScrollbackRows
@@ -73,16 +74,29 @@ export function installMultiplexSlotFrames(
}
// Mobile already has the higher-priority floor, so a rejected desktop claim must not suppress later phone input.
const inputClaimTail = stream.isMobile ? Promise.resolve(true) : stream.desktopClaimTail
const isQueryReply = frame.opcode === TerminalStreamOpcode.QueryReply
void inputClaimTail.then(async (claimed) => {
if (!claimed || isTerminalInputLockedForClient(runtime, stream.ptyId, stream.client)) {
return
}
// Why: opcode 18 skips the mobile input floor, so it needs every guard terminal.send applies; drop otherwise.
if (
isQueryReply &&
!isAcceptableTerminalQueryReplyFrame({
runtime,
ptyId: stream.ptyId,
text,
client: stream.client
})
) {
return
}
const outcome = await sendTerminalStreamInput(runtime, {
terminal: stream.terminal,
text,
client: stream.client,
isMobile: stream.isMobile,
inputKind: frame.opcode === TerminalStreamOpcode.QueryReply ? 'query-reply' : 'input'
inputKind: isQueryReply ? 'query-reply' : 'input'
})
state.notifyStreamWriteUnavailable(stream, outcome)
})
@@ -0,0 +1,47 @@
import { isTerminalQueryReply } from '../../../../../shared/terminal-query-reply'
import type { OrcaRuntimeService } from '../../../orca-runtime'
import type { TerminalViewportClient } from './terminal-stream-types'
type TerminalQueryReplyIdentity = {
text: string | undefined
client: TerminalViewportClient | undefined
/** Authenticated device token when the transport carries one; the declared client id must match it. */
connectionClientId?: string | undefined
}
/**
* Client id allowed to author this query reply, or null when the request is not a
* well-formed reply from an identified mobile client. Shared so the JSON
* `terminal.send` path and the binary opcode-18 frame paths cannot drift.
*/
export function resolveTerminalQueryReplyAuthorId(args: TerminalQueryReplyIdentity): string | null {
const authorId = args.connectionClientId ?? args.client?.id
if (
!args.text ||
!isTerminalQueryReply(args.text) ||
args.client?.type !== 'mobile' ||
!authorId ||
(args.connectionClientId !== undefined && args.client.id !== args.connectionClientId)
) {
return null
}
return authorId
}
/**
* Whole guard for a binary query-reply frame: reply grammar, mobile identity, and
* reply authority. A frame that fails is dropped, matching what the JSON path does
* for the same failures (throw on shape, `accepted: false` on authority).
*/
export function isAcceptableTerminalQueryReplyFrame(
args: TerminalQueryReplyIdentity & {
runtime: Pick<OrcaRuntimeService, 'isMobileTerminalQueryReplyAuthority'>
ptyId: string | undefined
}
): boolean {
const authorId = resolveTerminalQueryReplyAuthorId(args)
if (!authorId || !args.ptyId) {
return false
}
return args.runtime.isMobileTerminalQueryReplyAuthority(args.ptyId, authorId)
}
@@ -1,8 +1,8 @@
import { isAgentSessionPtyWriteRefusedError } from '../../../../../shared/agent-session-pty-write-admission'
import { assertLegacyAiVaultResumeCommandAllowed } from '../../../../ai-vault/structured-session-ownership'
import { InvalidArgumentError, defineMethod, type RpcAnyMethod } from '../../core'
import { isTerminalQueryReply } from '../../../../../shared/terminal-query-reply'
import { assertTerminalAgentSendable } from '../../terminal-agent-send-guard'
import { resolveTerminalQueryReplyAuthorId } from './terminal-query-reply-guard'
import { TerminalSend } from './unary-schemas'
import {
assertTerminalSendExactPtyBinding,
@@ -33,18 +33,18 @@ export const TERMINAL_SEND_METHODS: RpcAnyMethod[] = [
runtime.ensureStructuredAgentSessionHost()
)
}
const queryReplyClientId = clientId ?? params.client?.id
const queryReplyClientId = resolveTerminalQueryReplyAuthorId({
text: params.text,
client: params.client,
connectionClientId: clientId
})
if (
params.inputKind === 'query-reply' &&
(!params.text ||
!isTerminalQueryReply(params.text) ||
(!queryReplyClientId ||
params.enter === true ||
params.interrupt === true ||
params.agentPrompt === true ||
params.requireAgentStatus !== undefined ||
params.client?.type !== 'mobile' ||
!queryReplyClientId ||
(clientId !== undefined && params.client.id !== clientId))
params.requireAgentStatus !== undefined)
) {
throw new InvalidArgumentError('Invalid terminal query reply')
}
@@ -33,8 +33,14 @@ type AssetRangeReader = (path: string, offset: number, length: number) => Promis
type MobileWebPackageAssetsOptions = {
resolveRoot?: () => string
readAssetRange?: AssetRangeReader
gzipCacheMaxBytes?: number
}
// Why: `length` is client-chosen, so one source byte can be cached under eight keys and
// the fingerprint that clears the map never changes on a shipped desktop. Bound the
// retained total instead, evicting least-recently-used first.
const GZIP_CACHE_MAX_BYTES = 16 * 1024 * 1024
export class MobileWebPackageAssets {
private readonly resolveRoot: () => string
private readonly readAssetRange: AssetRangeReader
@@ -46,10 +52,13 @@ export class MobileWebPackageAssets {
} | null = null
private readonly readStates = new Map<string, PackageReadState>()
private readonly gzipChunks = new Map<string, Buffer>()
private gzipChunkBytes = 0
private readonly gzipCacheMaxBytes: number
constructor(options: MobileWebPackageAssetsOptions = {}) {
this.resolveRoot = options.resolveRoot ?? resolveMobileWebPackageRoot
this.readAssetRange = options.readAssetRange ?? readAssetRange
this.gzipCacheMaxBytes = options.gzipCacheMaxBytes ?? GZIP_CACHE_MAX_BYTES
}
async getManifest(): Promise<MobileWebPackageManifestResponse> {
@@ -131,6 +140,7 @@ export class MobileWebPackageAssets {
try {
const verified = await promise
this.gzipChunks.clear()
this.gzipChunkBytes = 0
this.cached = { fingerprint, package: verified }
return verified
} finally {
@@ -146,7 +156,7 @@ export class MobileWebPackageAssets {
): Promise<MobileWebPackageGzipAssetChunk> {
const requestedLength = params.length ?? MOBILE_WEB_PACKAGE_CHUNK_BYTES
const key = `${params.buildId}:${params.path}:${params.offset}:${requestedLength}`
const cached = this.gzipChunks.get(key)
const cached = this.readGzipChunk(key)
if (cached) {
throwIfAborted(options.signal)
const verified = await this.getVerifiedPackage()
@@ -163,7 +173,7 @@ export class MobileWebPackageAssets {
}
const read = await this.readVerifiedRange(params, options, requestedLength)
const compressed = gzipSync(read.bytes, { level: 6 })
this.gzipChunks.set(key, compressed)
this.storeGzipChunk(key, compressed)
return this.gzipResponse(
read.buildId,
read.asset,
@@ -173,6 +183,37 @@ export class MobileWebPackageAssets {
)
}
/** Reading marks the entry most-recently-used; Map iteration is insertion-ordered. */
private readGzipChunk(key: string): Buffer | undefined {
const cached = this.gzipChunks.get(key)
if (!cached) {
return undefined
}
this.gzipChunks.delete(key)
this.gzipChunks.set(key, cached)
return cached
}
private storeGzipChunk(key: string, compressed: Buffer): void {
if (compressed.byteLength > this.gzipCacheMaxBytes) {
return
}
const previous = this.gzipChunks.get(key)
if (previous) {
this.gzipChunks.delete(key)
this.gzipChunkBytes -= previous.byteLength
}
this.gzipChunks.set(key, compressed)
this.gzipChunkBytes += compressed.byteLength
for (const [oldest, bytes] of this.gzipChunks) {
if (this.gzipChunkBytes <= this.gzipCacheMaxBytes) {
return
}
this.gzipChunks.delete(oldest)
this.gzipChunkBytes -= bytes.byteLength
}
}
private gzipResponse(
buildId: string,
asset: MobileWebAsset,
@@ -1,5 +1,5 @@
import { createHash } from 'node:crypto'
import { mkdtemp, mkdir, rm, writeFile } from 'node:fs/promises'
import { mkdtemp, mkdir, readFile, rm, writeFile } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { gunzipSync } from 'node:zlib'
@@ -84,6 +84,38 @@ describe('mobile web package ranged gzip reads', () => {
expect(second).toEqual(first)
})
it('evicts the least recently used range once the gzip cache is full', async () => {
const fixture = await createRangeFixture()
const script = fixture.manifest.assets.find((asset) => asset.role === 'script')!
const length = 4 * MOBILE_WEB_PACKAGE_CHUNK_BYTES
const first = { buildId: fixture.manifest.buildId, path: script.path, offset: 0, length }
const second = { ...first, offset: MOBILE_WEB_PACKAGE_CHUNK_BYTES }
// Both ranges are the same run of identical bytes, so one cached entry is the cap.
const sized = await new MobileWebPackageAssets({
resolveRoot: () => fixture.root
}).getAssetGzipChunk(first)
let reads = 0
const assets = new MobileWebPackageAssets({
resolveRoot: () => fixture.root,
gzipCacheMaxBytes: sized.byteLength,
readAssetRange: async (path, offset, byteLength) => {
reads += 1
return readFile(path).then((bytes) => bytes.subarray(offset, offset + byteLength))
}
})
await assets.getAssetGzipChunk(first)
await assets.getAssetGzipChunk(first)
expect(reads).toBe(1)
await assets.getAssetGzipChunk(second)
// Caching the second range evicted the first, so it must be read again.
await assets.getAssetGzipChunk(first)
expect(reads).toBe(3)
// The re-read made the first range the most recent, so it stays cached.
await assets.getAssetGzipChunk(first)
expect(reads).toBe(3)
})
it('keeps a chunk read and a range read at the same offset apart', async () => {
const fixture = await createRangeFixture()
const assets = new MobileWebPackageAssets({ resolveRoot: () => fixture.root })
@@ -139,13 +139,167 @@ describe('terminal query-reply opcode negotiation', () => {
})
})
function queryReplyFrame(streamId: number) {
describe('terminal query-reply opcode guards', () => {
const GUARD_CASES = [
{
name: 'a non-authority phone',
authority: false,
text: QUERY_REPLY,
type: 'mobile' as const
},
{
name: 'non-reply grammar',
authority: true,
text: ':q!\r',
type: 'mobile' as const
},
{
name: 'a desktop client',
authority: true,
text: QUERY_REPLY,
type: 'desktop' as const
}
]
it.each(GUARD_CASES)('drops multiplex opcode 18 from $name', async (testCase) => {
const sendTerminal = vi.fn().mockResolvedValue({ accepted: true })
const harness = startDesktopMultiplexSubscribe({
sendTerminal,
handleMobileSubscribe: vi.fn().mockResolvedValue(undefined),
handleMobileUnsubscribe: vi.fn(),
isMobileTerminalQueryReplyAuthority: vi.fn().mockReturnValue(testCase.authority)
})
await vi.waitFor(() => expect(harness.handlers.has(0)).toBe(true))
harness.handlers.get(0)?.(
subscribeFrame({
streamId: 11,
clientType: testCase.type,
negotiated: true
})
)
await vi.waitFor(() => expect(harness.handlers.has(11)).toBe(true))
harness.handlers.get(11)?.(queryReplyFrame(11, testCase.text))
// Ordinary input on the same stream proves the stream is live and the drop is guard-specific.
harness.handlers.get(11)?.(inputFrame(11))
await vi.waitFor(() => expect(sendTerminal).toHaveBeenCalledTimes(1))
expect(sendTerminal.mock.calls[0]?.[1]).toMatchObject({ text: ':q!\r' })
harness.registry.cleanupSubscription('terminal-multiplex:conn-desktop-first-paint')
await harness.dispatchPromise
})
it.each(GUARD_CASES)('drops direct opcode 18 from $name', async (testCase) => {
const registry = createSubscriptionRegistryDouble()
const messages: string[] = []
const handlers = new Map<
number,
(frame: NonNullable<ReturnType<typeof decodeTerminalStreamFrame>>) => void
>()
const sendTerminal = vi.fn().mockResolvedValue({ accepted: true })
const runtime = stubRuntime({
...directSubscribeRuntime(registry),
isMobileTerminalQueryReplyAuthority: vi.fn().mockReturnValue(testCase.authority),
sendTerminal
})
const dispatcher = new RpcDispatcher({ runtime, methods: TERMINAL_METHODS })
const dispatchPromise = dispatcher.dispatchStreaming(
makeRequest('terminal.subscribe', {
terminal: 'terminal-1',
client: { id: 'client-1', type: testCase.type },
capabilities: { terminalBinaryStream: 1, queryReply: 1 }
}),
(message) => messages.push(message),
{
connectionId: `conn-guard-${testCase.name}`,
sendBinary: vi.fn(),
registerBinaryStreamHandler: (streamId, handler) => {
handlers.set(streamId, handler)
return () => handlers.delete(streamId)
}
}
)
await vi.waitFor(() =>
expect(messages.some((message) => JSON.parse(message).result?.type === 'subscribed')).toBe(
true
)
)
const streamId = messages
.map((message) => JSON.parse(message).result)
.find((event) => event?.type === 'subscribed').streamId
handlers.get(streamId)?.(queryReplyFrame(streamId, testCase.text))
handlers.get(streamId)?.(inputFrame(streamId))
await vi.waitFor(() => expect(sendTerminal).toHaveBeenCalledTimes(1))
expect(sendTerminal.mock.calls[0]?.[1]).toMatchObject({ text: ':q!\r' })
runtime.cleanupSubscription('terminal-1:client-1')
await dispatchPromise
})
})
function directSubscribeRuntime(registry: ReturnType<typeof createSubscriptionRegistryDouble>) {
return {
resolveLeafForHandle: vi.fn().mockReturnValue({ ptyId: 'pty-1' }),
readTerminal: vi.fn().mockResolvedValue({ tail: [], truncated: false }),
serializeTerminalBuffer: vi.fn().mockResolvedValue(null),
getTerminalSize: vi.fn().mockReturnValue({ cols: 80, rows: 24 }),
getMobileDisplayMode: vi.fn().mockReturnValue('auto'),
getDriver: vi.fn().mockReturnValue({ kind: 'idle' }),
getLayout: vi.fn().mockReturnValue({ seq: 1 }),
subscribeToTerminalData: vi.fn().mockReturnValue(vi.fn()),
subscribeToTerminalResize: vi.fn().mockReturnValue(vi.fn()),
subscribeToFitOverrideChanges: vi.fn().mockReturnValue(vi.fn()),
registerOwnedSubscriptionCleanup: vi.fn(registry.registerOwnedSubscriptionCleanup),
cleanupSubscription: vi.fn(registry.cleanupSubscription),
waitForTerminal: vi.fn(() => new Promise<RuntimeTerminalWait>(() => {})),
handleMobileSubscribe: vi.fn().mockResolvedValue(undefined),
handleMobileUnsubscribe: vi.fn()
}
}
function subscribeFrame(options: {
streamId: number
clientType: 'mobile' | 'desktop'
negotiated: boolean
}) {
return decodeTerminalStreamFrame(
encodeTerminalStreamFrame({
opcode: TerminalStreamOpcode.Subscribe,
streamId: 0,
seq: 1,
payload: encodeTerminalStreamJson({
streamId: options.streamId,
terminal: 'terminal-1',
client: { id: 'client-1', type: options.clientType },
capabilities: {
ackOutput: 1,
...(options.negotiated ? { queryReply: 1 } : {})
}
})
})
)!
}
function inputFrame(streamId: number) {
return decodeTerminalStreamFrame(
encodeTerminalStreamFrame({
opcode: TerminalStreamOpcode.Input,
streamId,
seq: 3,
payload: encodeTerminalStreamText(':q!\r')
})
)!
}
function queryReplyFrame(streamId: number, text: string = QUERY_REPLY) {
return decodeTerminalStreamFrame(
encodeTerminalStreamFrame({
opcode: TerminalStreamOpcode.QueryReply,
streamId,
seq: 2,
payload: encodeTerminalStreamText(QUERY_REPLY)
payload: encodeTerminalStreamText(text)
})
)!
}
@@ -47,6 +47,7 @@ export const MOBILE_RPC_METHOD_ALLOWLIST = new Set([
'files.readTerminalArtifactPreview',
'files.resolveTerminalPath',
'files.searchPaths',
'files.unwatch',
'files.watch',
'files.writeIfUnchanged',
'files.writeTerminalArtifact',