diff --git a/config/electron-builder.config.cjs b/config/electron-builder.config.cjs index 30d81bcef26..4650933b9e3 100644 --- a/config/electron-builder.config.cjs +++ b/config/electron-builder.config.cjs @@ -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}', diff --git a/config/scripts/electron-builder-config.test.mjs b/config/scripts/electron-builder-config.test.mjs index f72d972675d..b0e7ceddea7 100644 --- a/config/scripts/electron-builder-config.test.mjs +++ b/config/scripts/electron-builder-config.test.mjs @@ -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', () => { diff --git a/config/scripts/hosted-mobile-webview-ssh-e2e-contract.test.mjs b/config/scripts/hosted-mobile-webview-ssh-e2e-contract.test.mjs new file mode 100644 index 00000000000..6089749c5b3 --- /dev/null +++ b/config/scripts/hosted-mobile-webview-ssh-e2e-contract.test.mjs @@ -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.')" + ) + }) +}) diff --git a/mobile/src/browser/mobile-browser-stream-events.test.ts b/mobile/src/browser/mobile-browser-stream-events.test.ts new file mode 100644 index 00000000000..d55ff8cd35d --- /dev/null +++ b/mobile/src/browser/mobile-browser-stream-events.test.ts @@ -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) + }) +}) diff --git a/mobile/src/mobile-web/mobile-web-terminal-input.ts b/mobile/src/mobile-web/mobile-web-terminal-input.ts index daaccb5f009..f619ab6fe31 100644 --- a/mobile/src/mobile-web/mobile-web-terminal-input.ts +++ b/mobile/src/mobile-web/mobile-web-terminal-input.ts @@ -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 } diff --git a/mobile/src/mobile-web/mobile-web-terminal-streams.test.ts b/mobile/src/mobile-web/mobile-web-terminal-streams.test.ts index d5b7409f455..dbc13cd4232 100644 --- a/mobile/src/mobile-web/mobile-web-terminal-streams.test.ts +++ b/mobile/src/mobile-web/mobile-web-terminal-streams.test.ts @@ -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>(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({ diff --git a/src/main/browser/browser-screencast-stream-types.ts b/src/main/browser/browser-screencast-stream-types.ts index 76fbfde6d8f..e3ee950ac04 100644 --- a/src/main/browser/browser-screencast-stream-types.ts +++ b/src/main/browser/browser-screencast-stream-types.ts @@ -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 } diff --git a/src/main/runtime/mobile-rpc-allowlist.test.ts b/src/main/runtime/mobile-rpc-allowlist.test.ts index 6a2f090c0b8..7f75c9b9795 100644 --- a/src/main/runtime/mobile-rpc-allowlist.test.ts +++ b/src/main/runtime/mobile-rpc-allowlist.test.ts @@ -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 { 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 { 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( diff --git a/src/main/runtime/rpc/methods/terminal/terminal-legacy-binary-control-frames.ts b/src/main/runtime/rpc/methods/terminal/terminal-legacy-binary-control-frames.ts index 808ccf63241..5c4ffef70cc 100644 --- a/src/main/runtime/rpc/methods/terminal/terminal-legacy-binary-control-frames.ts +++ b/src/main/runtime/rpc/methods/terminal/terminal-legacy-binary-control-frames.ts @@ -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) diff --git a/src/main/runtime/rpc/methods/terminal/terminal-multiplex-slot-frames.ts b/src/main/runtime/rpc/methods/terminal/terminal-multiplex-slot-frames.ts index b848a8ddeb3..bd30759f86e 100644 --- a/src/main/runtime/rpc/methods/terminal/terminal-multiplex-slot-frames.ts +++ b/src/main/runtime/rpc/methods/terminal/terminal-multiplex-slot-frames.ts @@ -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) }) diff --git a/src/main/runtime/rpc/methods/terminal/terminal-query-reply-guard.ts b/src/main/runtime/rpc/methods/terminal/terminal-query-reply-guard.ts new file mode 100644 index 00000000000..9fc6160b3ad --- /dev/null +++ b/src/main/runtime/rpc/methods/terminal/terminal-query-reply-guard.ts @@ -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 + ptyId: string | undefined + } +): boolean { + const authorId = resolveTerminalQueryReplyAuthorId(args) + if (!authorId || !args.ptyId) { + return false + } + return args.runtime.isMobileTerminalQueryReplyAuthority(args.ptyId, authorId) +} diff --git a/src/main/runtime/rpc/methods/terminal/terminal-send-method.ts b/src/main/runtime/rpc/methods/terminal/terminal-send-method.ts index 6bafeb93966..200217f4661 100644 --- a/src/main/runtime/rpc/methods/terminal/terminal-send-method.ts +++ b/src/main/runtime/rpc/methods/terminal/terminal-send-method.ts @@ -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') } diff --git a/src/main/runtime/rpc/mobile-web-package-assets.ts b/src/main/runtime/rpc/mobile-web-package-assets.ts index 2da21be10d3..344e4f74d1d 100644 --- a/src/main/runtime/rpc/mobile-web-package-assets.ts +++ b/src/main/runtime/rpc/mobile-web-package-assets.ts @@ -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() private readonly gzipChunks = new Map() + 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 { @@ -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 { 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, diff --git a/src/main/runtime/rpc/mobile-web-package-range-reads.test.ts b/src/main/runtime/rpc/mobile-web-package-range-reads.test.ts index 1c7dfc303c7..dd70f312a89 100644 --- a/src/main/runtime/rpc/mobile-web-package-range-reads.test.ts +++ b/src/main/runtime/rpc/mobile-web-package-range-reads.test.ts @@ -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 }) diff --git a/src/main/runtime/rpc/terminal-query-reply-negotiation.test.ts b/src/main/runtime/rpc/terminal-query-reply-negotiation.test.ts index e64d43f95e7..ff8a94baa3f 100644 --- a/src/main/runtime/rpc/terminal-query-reply-negotiation.test.ts +++ b/src/main/runtime/rpc/terminal-query-reply-negotiation.test.ts @@ -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>) => 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) { + 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(() => {})), + 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) }) )! } diff --git a/src/main/runtime/runtime-rpc/runtime-rpc-mobile-method-allowlist.ts b/src/main/runtime/runtime-rpc/runtime-rpc-mobile-method-allowlist.ts index 952068fce57..0e77974ab62 100644 --- a/src/main/runtime/runtime-rpc/runtime-rpc-mobile-method-allowlist.ts +++ b/src/main/runtime/runtime-rpc/runtime-rpc-mobile-method-allowlist.ts @@ -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',