diff --git a/docs/reference/remote-wire-compatibility.md b/docs/reference/remote-wire-compatibility.md index 13741c03789..0adbd351592 100644 --- a/docs/reference/remote-wire-compatibility.md +++ b/docs/reference/remote-wire-compatibility.md @@ -75,6 +75,28 @@ Treat these as wire changes even though nothing in the codec moves: If old clients cannot interpret the new projection correctly, gate it behind a runtime capability the same way Rule 2 gates an opcode. +## Worked example — a new optional param on a strict schema + +`mobileWeb.package.asset.gzip` gained an optional `length` so a phone can pull several +48 KiB chunks per round trip. That is not Rule 1 free: the params schema is zod +`.strict()`, so a host that predates the field answers `invalid_argument` rather than +ignoring it. The negotiation is therefore Rule 2 shaped even though nothing about the +framing changed: + +- the host advertises `MOBILE_WEB_PACKAGE_RANGE_RUNTIME_CAPABILITY` + (`mobileWeb.package.range.v1`) in `RUNTIME_CAPABILITIES`; +- `useMobileWebPackageCapability` reads it off the same probe that already gates + `mobileWeb.package.gzip.v1`, and only then does the downloader send `length`; +- a host that answers without the capability keeps getting one-chunk reads, and the + response it returns still validates against the old client's chunk schema. + +The response side is Rule 1 in the safe direction: `sourceByteLength`'s ceiling moved up +to `MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES`, but a host only returns more than one chunk to a +client that asked for it, so an old client never receives a chunk its schema would refuse. +Read the ceiling before widening it again — 384 KiB of incompressible bytes is still +inside the relay's 1 MiB control-lane frame after gzip, base64, and the E2EE layer's own +base64, and a larger range is not. + ## Enforcement `tests/e2e/cross-version-wire/cross-version-terminal-wire.unit.test.ts` runs the real diff --git a/mobile/scripts/measure-mobile-web-package-download.mjs b/mobile/scripts/measure-mobile-web-package-download.mjs new file mode 100644 index 00000000000..ceff5b9bc03 --- /dev/null +++ b/mobile/scripts/measure-mobile-web-package-download.mjs @@ -0,0 +1,112 @@ +#!/usr/bin/env node +// Measures mobile-web package download throughput over the real mobile RPC transports. +// +// node mobile/scripts/measure-mobile-web-package-download.mjs \ +// --root out/mobile-web-rnw --paths direct,relay --rtt 0,40,120 --concurrency 1,4 +// +// --rtt values are full round-trip milliseconds injected on the loopback socket (half per +// hop), so a cloud-relay round trip can be modelled without a production relay. +import { mkdtempSync, rmSync, writeFileSync } from 'node:fs' +import { join, resolve } from 'node:path' +import { pathToFileURL } from 'node:url' +import * as esbuild from 'esbuild' + +const scriptDirectory = import.meta.dirname +const repositoryRoot = resolve(scriptDirectory, '..', '..') + +const options = parseArguments(process.argv.slice(2)) +// Why: the bundle keeps `ws` external, so it has to live inside the repo to resolve it. +const buildDirectory = mkdtempSync(join(repositoryRoot, 'node_modules', '.orca-measure-')) +try { + const { measureMobileWebPackageDownload } = await bundleHarness(buildDirectory) + const results = [] + for (const path of options.paths) { + for (const rtt of options.rtts) { + for (const concurrency of options.concurrencies) { + for (const gzip of options.gzips) { + const result = await measureMobileWebPackageDownload({ + path, + packageRoot: options.packageRoot, + oneWayDelayMs: rtt / 2, + gzip, + rangeBytes: options.rangeBytes, + maxConcurrentRequests: concurrency + }) + results.push(result) + console.log(formatRow(result)) + } + } + } + } + console.log('') + console.log(formatTable(results)) +} finally { + rmSync(buildDirectory, { recursive: true, force: true }) +} + +async function bundleHarness(outputDirectory) { + const shimPath = join(outputDirectory, 'expo-crypto-shim.mjs') + writeFileSync( + shimPath, + "import { randomBytes } from 'node:crypto'\n" + + 'export function getRandomBytes(length) { return new Uint8Array(randomBytes(length)) }\n' + + 'export default { getRandomBytes }\n' + ) + const outfile = join(outputDirectory, 'harness.mjs') + await esbuild.build({ + entryPoints: [join(scriptDirectory, 'mobile-web-package-download-measurement.ts')], + outfile, + bundle: true, + platform: 'node', + format: 'esm', + target: 'node22', + logLevel: 'error', + absWorkingDir: repositoryRoot, + alias: { 'expo-crypto': shimPath }, + external: ['ws', 'electron'], + banner: { + js: "import { createRequire as __createRequire } from 'node:module'\nconst require = __createRequire(import.meta.url)\n" + } + }) + return import(pathToFileURL(outfile).href) +} + +function parseArguments(argv) { + const flags = new Map() + for (let index = 0; index < argv.length; index += 1) { + const flag = argv[index] + if (flag.startsWith('--')) { + flags.set(flag.slice(2), argv[index + 1]) + index += 1 + } + } + return { + packageRoot: resolve(repositoryRoot, flags.get('root') ?? 'out/mobile-web-rnw'), + paths: (flags.get('paths') ?? 'direct,relay').split(','), + rtts: (flags.get('rtt') ?? '0').split(',').map(Number), + concurrencies: (flags.get('concurrency') ?? '1').split(',').map(Number), + gzips: (flags.get('gzip') ?? 'true').split(',').map((value) => value === 'true'), + rangeBytes: Number(flags.get('range-kib') ?? '48') * 1024 + } +} + +function formatRow(result) { + return [ + result.path.padEnd(6), + `rtt=${String(result.oneWayDelayMs * 2).padStart(3)}ms`, + `conc=${String(result.maxConcurrentRequests).padStart(2)}`, + `gzip=${result.gzip ? 'on ' : 'off'}`, + `range=${String(result.rangeBytes / 1024).padStart(3)}KiB`, + `${(result.wallMs / 1000).toFixed(1).padStart(7)}s`, + `${(result.bytesPerSecond / 1_000_000).toFixed(2).padStart(6)} MB/s`, + `chunks=${String(result.chunkRequests).padStart(4)}`, + `wire=${(result.wireBytesToPhone / 1_048_576).toFixed(1).padStart(6)} MiB`, + `p50=${result.latencyMedianMs.toFixed(1).padStart(7)}ms`, + `p95=${result.latencyP95Ms.toFixed(1).padStart(7)}ms`, + `peak=${result.peakInFlight}` + ].join(' ') +} + +function formatTable(results) { + return results.map(formatRow).join('\n') +} diff --git a/mobile/scripts/mobile-web-package-download-measurement.ts b/mobile/scripts/mobile-web-package-download-measurement.ts new file mode 100644 index 00000000000..3f09b256c27 --- /dev/null +++ b/mobile/scripts/mobile-web-package-download-measurement.ts @@ -0,0 +1,461 @@ +// Harness behind scripts/measure-mobile-web-package-download.mjs. It drives the real mobile +// package downloader over the real mobile RPC clients (direct WebSocket and the cloud-relay +// session) against the real desktop MobileWebPackageAssets reader, with an injectable per-hop +// delay so a cloud round trip can be modelled on loopback. +import { mkdtempSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import nacl from 'tweetnacl' +import WebSocketClient, { WebSocketServer, type RawData, type WebSocket } from 'ws' +import { connect as connectDirectRpcClient, type RpcClient } from '../src/transport/rpc-client' +import { connectMobileRelayRpcSession } from '../src/transport/mobile-relay-rpc-session' +import { + downloadMobileWebPackage, + type MobileWebPackageStager +} from '../src/mobile-web/mobile-web-package-downloader' +import { MOBILE_WEB_BRIDGE_PROTOCOL_VERSION } from '../../src/shared/mobile-web/bridge-contract' +import type { MobileWebPackageAssetParams } from '../../src/shared/mobile-web/package-rpc-contract' +import type { MobileRelayEndpoint } from '../../src/shared/mobile-relay-credential-contract' +import { DeviceRegistry } from '../../src/main/runtime/device-registry' +import { + MobileSocketWiring, + type MobileSocketTransport +} from '../../src/main/runtime/rpc/mobile-socket-wiring' +import { CloudRelayTransport } from '../../src/main/runtime/rpc/relay-transport' +import { deriveRelayHostId } from '../../src/main/runtime/relay/relay-http-client' +import { MobileWebPackageAssets } from '../../src/main/runtime/rpc/mobile-web-package-assets' + +export type MeasurementOptions = { + path: 'direct' | 'relay' + packageRoot: string + oneWayDelayMs: number + gzip: boolean + rangeBytes: number + maxConcurrentRequests: number +} + +export type MeasurementResult = MeasurementOptions & { + totalBytes: number + wallMs: number + bytesPerSecond: number + chunkRequests: number + wireBytesToPhone: number + latencyMedianMs: number + latencyP95Ms: number + peakInFlight: number +} + +type MeasurementRpcRequest = { id: string; method: string; params?: Record } + +type Teardown = (() => Promise | void)[] + +export async function measureMobileWebPackageDownload( + options: MeasurementOptions +): Promise { + const teardown: Teardown = [] + try { + const packageAssets = new MobileWebPackageAssets({ resolveRoot: () => options.packageRoot }) + const userDataPath = mkdtempSync(join(tmpdir(), 'orca-package-measure-')) + teardown.push(() => rmSync(userDataPath, { recursive: true, force: true })) + const registry = new DeviceRegistry(userDataPath) + const device = registry.addDevice('Measurement phone', 'mobile') + const desktopKeys = nacl.box.keyPair() + const desktopPublicKeyB64 = Buffer.from(desktopKeys.publicKey).toString('base64') + const relayHostId = deriveRelayHostId(desktopKeys.publicKey) + const relayEndpoint: MobileRelayEndpoint = { + v: 1, + directorUrl: 'https://relay.measure.test', + cellUrl: 'https://relay-c1.measure.test', + assignmentEpoch: 1, + relayHostId, + e2eeFraming: 2 + } + + let wireBytesToPhone = 0 + const countInbound = (bytes: number): void => { + wireBytesToPhone += bytes + } + const wiring = new MobileSocketWiring({ + deviceRegistry: registry, + e2eeKeypair: { + publicKey: desktopKeys.publicKey, + secretKey: desktopKeys.secretKey, + publicKeyB64: desktopPublicKeyB64 + }, + onText: (socket, plaintext, reply) => { + const request = JSON.parse(plaintext) as MeasurementRpcRequest + if (request.method === 'pairing.getEndpoints') { + reply( + rpcSuccess(request.id, { + v: 1, + relay: relayEndpoint, + resumeConfirmation: { + v: 1, + reqId: request.params?.resumeConfirmReqId, + currentVersion: 3, + acceptedAs: 'current', + renewed: true, + resumeExpiresAt: Date.now() + 900_000 + } + }) + ) + return + } + const params = request.params as unknown as MobileWebPackageAssetParams + const readOptions = { connectionId: socket.connectionId } + if (request.method === 'mobileWeb.package.manifest') { + settle(request.id, packageAssets.getManifest(), reply) + return + } + if (request.method === 'mobileWeb.package.asset') { + settle(request.id, packageAssets.getAssetChunk(params, readOptions), reply) + return + } + if (request.method === 'mobileWeb.package.asset.gzip') { + settle(request.id, packageAssets.getAssetGzipChunk(params, readOptions), reply) + return + } + reply(rpcSuccess(request.id, {})) + }, + onBinary: () => {}, + onClose: () => {} + }) + + const client = + options.path === 'relay' + ? await openRelayClient({ + wiring, + relayEndpoint, + relayHostId, + relayDeviceId: device.deviceId, + deviceToken: device.token, + desktopPublicKeyB64, + oneWayDelayMs: options.oneWayDelayMs, + teardown, + countInbound + }) + : await openDirectClient({ + wiring, + deviceToken: device.token, + desktopPublicKeyB64, + oneWayDelayMs: options.oneWayDelayMs, + teardown, + countInbound + }) + + const latencies: number[] = [] + let chunkRequests = 0 + let inFlight = 0 + let peakInFlight = 0 + const startedAt = performance.now() + const downloaded = await downloadMobileWebPackage( + async (method, params) => { + inFlight += 1 + peakInFlight = Math.max(peakInFlight, inFlight) + const requestStartedAt = performance.now() + try { + return await client.sendRequest(method, params, { timeoutMs: 180_000 }) + } finally { + inFlight -= 1 + if (method !== 'mobileWeb.package.manifest') { + chunkRequests += 1 + latencies.push(performance.now() - requestStartedAt) + } + } + }, + createDiscardingStager(), + { + shellBridgeVersion: MOBILE_WEB_BRIDGE_PROTOCOL_VERSION, + useGzip: options.gzip, + rangeBytes: options.rangeBytes, + maxConcurrentRequests: options.maxConcurrentRequests + } + ) + const wallMs = performance.now() - startedAt + client.close() + return { + ...options, + totalBytes: downloaded.manifest.totalBytes, + wallMs, + bytesPerSecond: downloaded.manifest.totalBytes / (wallMs / 1000), + chunkRequests, + wireBytesToPhone, + latencyMedianMs: percentile(latencies, 0.5), + latencyP95Ms: percentile(latencies, 0.95), + peakInFlight + } + } finally { + for (const dispose of teardown.toReversed()) { + await dispose() + } + } +} + +async function openRelayClient(args: { + wiring: MobileSocketWiring + relayEndpoint: MobileRelayEndpoint + relayHostId: string + relayDeviceId: string + deviceToken: string + desktopPublicKeyB64: string + oneWayDelayMs: number + teardown: Teardown + countInbound: (bytes: number) => void +}): Promise { + const relayServer = new WebSocketServer({ host: '127.0.0.1', port: 0, perMessageDeflate: false }) + const port = await listen(relayServer, args.teardown) + + let hostSocket: WebSocket | null = null + let phoneSocket: WebSocket | null = null + const splice = (): void => { + if (!hostSocket || !phoneSocket) { + return + } + const host = hostSocket + const phone = phoneSocket + host.on('message', (raw, isBinary) => { + args.countInbound(rawByteLength(raw)) + forward(phone, raw, isBinary, args.oneWayDelayMs) + }) + phone.on('message', (raw, isBinary) => forward(host, raw, isBinary, args.oneWayDelayMs)) + phone.send( + JSON.stringify({ + type: 'relay-hello', + ok: true, + credentialKind: 'resume', + leaseExpiresAt: Date.now() + 900_000, + acceptedCredentialVersion: 3, + acceptedAs: 'current', + resumeExpiresAt: Date.now() + 900_000 + }) + ) + } + relayServer.on('connection', (socket, request) => { + socket.once('message', () => { + if (request.url === '/v1/host/data/connection-1') { + hostSocket = socket + } else { + phoneSocket = socket + } + splice() + }) + }) + + const transport = new CloudRelayTransport({ + cellUrl: `http://127.0.0.1:${port}`, + relayHostId: args.relayHostId, + generation: 1 + }) + args.teardown.push(() => transport.stop()) + args.wiring.attachTransport(transport, (socket) => transport.metadataFor(socket)) + await transport.start() + await transport.openConnection({ + connId: 'connection-1', + connTicket: 'A'.repeat(43), + kind: 'resume', + relayDeviceId: args.relayDeviceId, + attachDeadlineMs: 15_000 + }) + + const session = connectMobileRelayRpcSession({ + relay: args.relayEndpoint, + resumeToken: 'B'.repeat(43), + resumeCredentialVersion: 3, + resumeConfirmReqId: 'measure-confirm', + deviceToken: args.deviceToken, + desktopPublicKeyB64: args.desktopPublicKeyB64, + requestTimeoutMs: 180_000, + createSocket: () => + new WebSocketClient(`ws://127.0.0.1:${port}/v1/connect/${args.relayHostId}`, { + perMessageDeflate: false, + maxPayload: 16 * 1024 * 1024 + }) as unknown as globalThis.WebSocket + }) + args.teardown.push(() => session.close()) + await waitFor(() => session.getState() === 'connected', 30_000, 'relay session connect') + return session +} + +async function openDirectClient(args: { + wiring: MobileSocketWiring + deviceToken: string + desktopPublicKeyB64: string + oneWayDelayMs: number + teardown: Teardown + countInbound: (bytes: number) => void +}): Promise { + const server = new WebSocketServer({ host: '127.0.0.1', port: 0, perMessageDeflate: false }) + const port = await listen(server, args.teardown) + args.wiring.attachTransport( + createDirectLoopbackTransport(server, args.oneWayDelayMs, args.countInbound) + ) + const client = connectDirectRpcClient( + `ws://127.0.0.1:${port}`, + args.deviceToken, + args.desktopPublicKeyB64 + ) + args.teardown.push(() => client.close()) + await waitFor(() => client.getState() === 'connected', 30_000, 'direct client connect') + return client +} + +function createDirectLoopbackTransport( + server: WebSocketServer, + oneWayDelayMs: number, + countInbound: (bytes: number) => void +): MobileSocketTransport { + type MessageHandler = ( + message: string | Uint8Array, + reply: (response: string) => void, + ws: WebSocket + ) => void + const messageHandlers: MessageHandler[] = [] + const closeHandlers: ((clientId: string | null, ws: WebSocket, other: boolean) => void)[] = [] + const clientIds = new Map() + server.on('connection', (socket) => { + // Why: the desktop replies through E2EEChannel's own ws.send, so the outbound hop is + // delayed by wrapping send rather than by a forwarding splice like the relay path has. + delaySocketSend(socket, oneWayDelayMs, countInbound) + socket.on('message', (raw, isBinary) => { + const message = isBinary ? toBytes(raw) : raw.toString('utf8') + after(oneWayDelayMs, () => { + for (const handler of messageHandlers) { + handler(message, () => {}, socket) + } + }) + }) + socket.on('close', () => { + const clientId = clientIds.get(socket) ?? null + clientIds.delete(socket) + for (const handler of closeHandlers) { + handler(clientId, socket, false) + } + }) + }) + return { + onMessage: (handler) => messageHandlers.push(handler), + onConnectionClose: (handler) => closeHandlers.push(handler), + setClientId: (ws, clientId) => clientIds.set(ws, clientId), + terminateClientConnections: (clientId) => { + let terminated = 0 + for (const [socket, id] of clientIds) { + if (id === clientId) { + socket.terminate() + terminated += 1 + } + } + return terminated + } + } +} + +function delaySocketSend( + socket: WebSocket, + oneWayDelayMs: number, + countInbound: (bytes: number) => void +): void { + const send = socket.send.bind(socket) + socket.send = ((data: string | Uint8Array, ...rest: unknown[]) => { + countInbound(typeof data === 'string' ? Buffer.byteLength(data) : data.byteLength) + after(oneWayDelayMs, () => { + if (socket.readyState === socket.OPEN) { + ;(send as (...values: unknown[]) => void)(data, ...rest) + } + }) + }) as typeof socket.send +} + +function forward(socket: WebSocket, raw: RawData, isBinary: boolean, oneWayDelayMs: number): void { + after(oneWayDelayMs, () => { + if (socket.readyState === socket.OPEN) { + socket.send(raw, { binary: isBinary }) + } + }) +} + +function after(delayMs: number, run: () => void): void { + if (delayMs <= 0) { + queueMicrotask(run) + return + } + setTimeout(run, delayMs) +} + +async function listen(server: WebSocketServer, teardown: Teardown): Promise { + await new Promise((resolve) => server.once('listening', resolve)) + const address = server.address() + if (!address || typeof address === 'string') { + throw new Error('expected a TCP address') + } + teardown.push( + () => + new Promise((resolve) => { + for (const socket of server.clients) { + socket.terminate() + } + server.close(() => resolve()) + }) + ) + return address.port +} + +function createDiscardingStager(): MobileWebPackageStager<{ buildId: string }> { + return { + async begin() {}, + async writeAssetChunk() {}, + async finishAsset() {}, + async commit(manifest) { + return { buildId: manifest.buildId } + }, + async abort() {} + } +} + +function settle(id: string, operation: Promise, reply: (response: string) => void): void { + void operation.then( + (result) => reply(rpcSuccess(id, result)), + (error: unknown) => + reply( + JSON.stringify({ + id, + ok: false, + error: { code: 'invalid_argument', message: String((error as Error).message) }, + _meta: { runtimeId: 'measure-runtime' } + }) + ) + ) +} + +function rpcSuccess(id: string, result: unknown): string { + return JSON.stringify({ id, ok: true, result, _meta: { runtimeId: 'measure-runtime' } }) +} + +function rawByteLength(raw: RawData): number { + return typeof raw === 'string' ? Buffer.byteLength(raw) : toBytes(raw).byteLength +} + +function toBytes(raw: RawData): Uint8Array { + if (Array.isArray(raw)) { + return Buffer.concat(raw) + } + return raw instanceof ArrayBuffer ? new Uint8Array(raw) : new Uint8Array(raw as Buffer) +} + +function percentile(values: number[], fraction: number): number { + if (values.length === 0) { + return 0 + } + const sorted = [...values].sort((left, right) => left - right) + return sorted[Math.min(sorted.length - 1, Math.floor(fraction * (sorted.length - 1)))]! +} + +async function waitFor(predicate: () => boolean, timeoutMs: number, label: string): Promise { + const deadline = Date.now() + timeoutMs + while (Date.now() < deadline) { + if (predicate()) { + return + } + await new Promise((resolve) => setTimeout(resolve, 25)) + } + throw new Error(`timed out waiting for ${label}`) +} diff --git a/mobile/src/mobile-web/MobileWebPackageProgress.test.tsx b/mobile/src/mobile-web/MobileWebPackageProgress.test.tsx new file mode 100644 index 00000000000..18469403feb --- /dev/null +++ b/mobile/src/mobile-web/MobileWebPackageProgress.test.tsx @@ -0,0 +1,51 @@ +import { createElement } from 'react' +import { act, create, type ReactTestRenderer } from 'react-test-renderer' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { MobileWebPackageProgress } from './MobileWebPackageProgress' + +vi.mock('react-native', () => ({ + Text: 'Text', + View: 'View', + StyleSheet: { create: (styles: unknown) => styles } +})) + +describe('mobile web package progress', () => { + let renderer: ReactTestRenderer | null = null + + beforeEach(() => { + globalThis.IS_REACT_ACT_ENVIRONMENT = true + }) + + afterEach(() => { + act(() => renderer?.unmount()) + renderer = null + }) + + it('tells the user to stay in the app while bytes are still arriving', () => { + act(() => { + renderer = create( + createElement(MobileWebPackageProgress, { + progress: { phase: 'downloading', completedBytes: 30_720, totalBytes: 9_548_835 } + }) + ) + }) + + expect(renderedText()).toContain('Keep Orca open until this finishes.') + }) + + it('drops the stay-open warning once the transfer is done', () => { + act(() => { + renderer = create( + createElement(MobileWebPackageProgress, { + progress: { phase: 'verifying', completedBytes: 9_548_835, totalBytes: 9_548_835 } + }) + ) + }) + + expect(renderedText()).not.toContain('Keep Orca open') + }) + + function renderedText(): string { + return JSON.stringify(renderer!.toJSON()) + } +}) diff --git a/mobile/src/mobile-web/MobileWebPackageProgress.tsx b/mobile/src/mobile-web/MobileWebPackageProgress.tsx index c8d14814f70..7725b1c4fc1 100644 --- a/mobile/src/mobile-web/MobileWebPackageProgress.tsx +++ b/mobile/src/mobile-web/MobileWebPackageProgress.tsx @@ -2,6 +2,10 @@ import { Text, View } from 'react-native' import { hybridShellStyles as styles } from './hybrid-shell-styles' import type { MobileWebPackageDownloadProgress } from './mobile-web-package-downloader' +// Backgrounding suspends the connection after MobileRelayBackgroundGrace's window and the +// staged download is discarded, so the transfer restarts from zero on return. +const STAY_OPEN_HINT = 'Keep Orca open until this finishes.' + export function MobileWebPackageProgress({ progress }: { @@ -34,6 +38,9 @@ export function MobileWebPackageProgress({ {formatBytes(progress.completedBytes)} of {formatBytes(progress.totalBytes)} + {progress.phase === 'downloading' ? ( + {STAY_OPEN_HINT} + ) : null} ) } diff --git a/mobile/src/mobile-web/mobile-web-package-chunk-decoder.ts b/mobile/src/mobile-web/mobile-web-package-chunk-decoder.ts index 74eb9e7e557..b1c571ab5f6 100644 --- a/mobile/src/mobile-web/mobile-web-package-chunk-decoder.ts +++ b/mobile/src/mobile-web/mobile-web-package-chunk-decoder.ts @@ -2,10 +2,10 @@ import { Buffer } from 'buffer/' import { sha256 } from '@noble/hashes/sha256' import { gunzipSync } from 'fflate' import { + MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES, MobileWebPackageAssetChunkSchema, MobileWebPackageGzipAssetChunkSchema } from '../../../src/shared/mobile-web/package-rpc-contract' -import { MOBILE_WEB_PACKAGE_CHUNK_BYTES } from '../../../src/shared/mobile-web/manifest-contract' export function decodeRawMobileWebPackageChunk( result: unknown, @@ -52,7 +52,7 @@ export function decodeGzipMobileWebPackageChunk( ) { return null } - if (expectedLength > MOBILE_WEB_PACKAGE_CHUNK_BYTES) { + if (expectedLength > MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES) { return null } const compressed = decodeCanonicalBase64(chunk.data.dataBase64) diff --git a/mobile/src/mobile-web/mobile-web-package-download-pipelining.test.ts b/mobile/src/mobile-web/mobile-web-package-download-pipelining.test.ts new file mode 100644 index 00000000000..4ba3251f2e2 --- /dev/null +++ b/mobile/src/mobile-web/mobile-web-package-download-pipelining.test.ts @@ -0,0 +1,329 @@ +import { Buffer } from 'buffer/' +import { gzipSync } from 'node:zlib' +import { sha256 } from '@noble/hashes/sha256' +import { describe, expect, it, vi } from 'vitest' +import { + MOBILE_WEB_MANIFEST_SCHEMA_VERSION, + MOBILE_WEB_PACKAGE_CHUNK_BYTES, + serializeMobileWebManifestForBuildId, + type MobileWebAsset, + type MobileWebManifest +} from '../../../src/shared/mobile-web/manifest-contract' +import { + MOBILE_WEB_PACKAGE_MAX_CONCURRENT_READS, + MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES +} from '../../../src/shared/mobile-web/package-rpc-contract' +import type { RpcResponse } from '../transport/types' +import { + downloadMobileWebPackage, + type MobileWebPackageRequest, + type MobileWebPackageStager +} from './mobile-web-package-downloader' + +type AssetParams = { buildId: string; path: string; offset: number; length?: number } + +type Fixture = { + manifest: MobileWebManifest + bytesByPath: Map + request: MobileWebPackageRequest + paramsByCall: AssetParams[] + peakInFlight: () => number +} + +describe('mobile web package download pipelining', () => { + it('overlaps chunk reads while staging them in ascending offset order', async () => { + const fixture = createFixture() + const stager = createStager() + + await downloadMobileWebPackage(fixture.request, stager, { + shellBridgeVersion: 1, + maxConcurrentRequests: 4 + }) + + expect(fixture.peakInFlight()).toBe(4) + expectStagedInOrder(fixture, stager) + }) + + it('defaults to the host per-connection read budget', async () => { + const fixture = createFixture() + + await downloadMobileWebPackage(fixture.request, createStager(), { shellBridgeVersion: 1 }) + + expect(fixture.peakInFlight()).toBe(MOBILE_WEB_PACKAGE_MAX_CONCURRENT_READS) + }) + + it('never exceeds the host read budget even when asked for more', async () => { + const fixture = createFixture() + + await downloadMobileWebPackage(fixture.request, createStager(), { + shellBridgeVersion: 1, + maxConcurrentRequests: 64 + }) + + expect(fixture.peakInFlight()).toBe(MOBILE_WEB_PACKAGE_MAX_CONCURRENT_READS) + }) + + it('reads one chunk at a time when concurrency is disabled', async () => { + const fixture = createFixture() + + await downloadMobileWebPackage(fixture.request, createStager(), { + shellBridgeVersion: 1, + maxConcurrentRequests: 1 + }) + + expect(fixture.peakInFlight()).toBe(1) + }) + + it('narrows the pipeline instead of failing when the host limits reads', async () => { + const fixture = createFixture({ readLimitedCalls: 3 }) + const stager = createStager() + + await downloadMobileWebPackage(fixture.request, stager, { + shellBridgeVersion: 1, + maxConcurrentRequests: 4 + }) + + expect(fixture.paramsByCall.length).toBeGreaterThan(chunkCount(fixture.manifest)) + expectStagedInOrder(fixture, stager) + }) + + it('gives up on a host that never stops limiting reads', async () => { + const fixture = createFixture({ readLimitedCalls: Number.POSITIVE_INFINITY }) + const stager = createStager() + + await expect( + downloadMobileWebPackage(fixture.request, stager, { shellBridgeVersion: 1 }) + ).rejects.toMatchObject({ code: 'mobile_web_package_read_limited' }) + expect(stager.abort).toHaveBeenCalledOnce() + }) + + it('aborts a pipelined download without committing the stage', async () => { + const fixture = createFixture() + const stager = createStager() + const controller = new AbortController() + stager.writeAssetChunk.mockImplementation(async () => controller.abort()) + + await expect( + downloadMobileWebPackage(fixture.request, stager, { + shellBridgeVersion: 1, + signal: controller.signal, + maxConcurrentRequests: 4 + }) + ).rejects.toMatchObject({ code: 'cancelled' }) + expect(stager.commit).not.toHaveBeenCalled() + expect(stager.abort).toHaveBeenCalledOnce() + + // Backgrounding aborts the same way, and nothing survives it: the retry re-reads offset 0. + fixture.paramsByCall.length = 0 + await downloadMobileWebPackage(fixture.request, createStager(), { shellBridgeVersion: 1 }) + expect(fixture.paramsByCall[0]?.offset).toBe(0) + }) + + it('requests a multi-chunk gzip range and stages it one chunk at a time', async () => { + const fixture = createFixture() + const stager = createStager() + + await downloadMobileWebPackage(fixture.request, stager, { + shellBridgeVersion: 1, + useGzip: true, + rangeBytes: MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + }) + + expect( + fixture.paramsByCall.every((params) => params.length === MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES) + ).toBe(true) + expect(fixture.paramsByCall.length).toBeLessThan(chunkCount(fixture.manifest)) + expectStagedInOrder(fixture, stager) + expect( + stager.writeAssetChunk.mock.calls.every( + (call) => call[2].byteLength <= MOBILE_WEB_PACKAGE_CHUNK_BYTES + ) + ).toBe(true) + }) + + it('omits the range length unless the host advertised range reads', async () => { + const fixture = createFixture() + + await downloadMobileWebPackage(fixture.request, createStager(), { + shellBridgeVersion: 1, + useGzip: true + }) + + expect(fixture.paramsByCall.every((params) => params.length === undefined)).toBe(true) + }) + + it('keeps ranged reads off the raw chunk method', async () => { + const fixture = createFixture() + + await downloadMobileWebPackage(fixture.request, createStager(), { + shellBridgeVersion: 1, + rangeBytes: MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + }) + + expect(fixture.paramsByCall.every((params) => params.length === undefined)).toBe(true) + expect(fixture.paramsByCall.length).toBe(chunkCount(fixture.manifest)) + }) + + it('rejects a ranged asset whose bytes do not hash to the manifest entry', async () => { + const fixture = createFixture({ alterLastRangeByte: true }) + const stager = createStager() + + await expect( + downloadMobileWebPackage(fixture.request, stager, { + shellBridgeVersion: 1, + useGzip: true, + rangeBytes: MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + }) + ).rejects.toMatchObject({ code: 'asset_integrity_failed' }) + expect(stager.abort).toHaveBeenCalledOnce() + }) +}) + +function expectStagedInOrder(fixture: Fixture, stager: ReturnType): void { + const stagedByPath = new Map() + for (const [asset, offset, bytes] of stager.writeAssetChunk.mock.calls) { + const staged = stagedByPath.get(asset.path) ?? [] + expect(offset).toBe(staged.reduce((total, part) => total + part.byteLength, 0)) + staged.push(Buffer.from(bytes)) + stagedByPath.set(asset.path, staged) + } + for (const [path, staged] of stagedByPath) { + expect(Buffer.concat(staged)).toEqual(Buffer.from(fixture.bytesByPath.get(path)!)) + } + expect([...stagedByPath.keys()].sort()).toEqual( + fixture.manifest.assets.map((asset) => asset.path).sort() + ) +} + +function chunkCount(manifest: MobileWebManifest): number { + return manifest.assets.reduce( + (total, asset) => total + Math.ceil(asset.byteLength / MOBILE_WEB_PACKAGE_CHUNK_BYTES), + 0 + ) +} + +function createFixture( + options: { readLimitedCalls?: number; alterLastRangeByte?: boolean } = {} +): Fixture { + const document = Buffer.from('Orca') + const script = Buffer.alloc( + MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + MOBILE_WEB_PACKAGE_CHUNK_BYTES + 11, + 0x61 + ) + const assets: MobileWebAsset[] = [ + asset('index.html', document, 'text/html; charset=utf-8', 'document'), + asset(`assets/${sha256Hex(script)}.js`, script, 'text/javascript; charset=utf-8', 'script') + ].sort((left, right) => left.path.localeCompare(right.path)) + const seed: MobileWebManifest = { + schemaVersion: MOBILE_WEB_MANIFEST_SCHEMA_VERSION, + buildId: '0'.repeat(64), + bridge: { minimum: 1, testedThrough: 2 }, + entrypoint: 'index.html', + totalBytes: assets.reduce((total, candidate) => total + candidate.byteLength, 0), + assets + } + const manifest = { + ...seed, + buildId: sha256Hex(Buffer.from(serializeMobileWebManifestForBuildId(seed))) + } + const bytesByPath = new Map([ + ['index.html', document as unknown as Uint8Array], + [assets.find((candidate) => candidate.role === 'script')!.path, script as unknown as Uint8Array] + ]) + const paramsByCall: AssetParams[] = [] + let readLimitedRemaining = options.readLimitedCalls ?? 0 + let inFlight = 0 + let peakInFlight = 0 + + const request = vi.fn(async (method: string, params?: unknown): Promise => { + if (method === 'mobileWeb.package.manifest') { + return success({ manifest, chunkBytes: MOBILE_WEB_PACKAGE_CHUNK_BYTES }) + } + const assetParams = params as AssetParams + paramsByCall.push(assetParams) + inFlight += 1 + peakInFlight = Math.max(peakInFlight, inFlight) + try { + // Why: a real host answers overlapping reads out of order, so the drain path has to + // reorder rather than rely on the request order. + await new Promise((resolve) => setTimeout(resolve, 5 - (paramsByCall.length % 5))) + if (readLimitedRemaining > 0) { + readLimitedRemaining -= 1 + return failure('mobile_web_package_read_limited') + } + const bytes = bytesByPath.get(assetParams.path)! + const requested = assetParams.length ?? MOBILE_WEB_PACKAGE_CHUNK_BYTES + const source = Buffer.from(bytes).subarray( + assetParams.offset, + Math.min(assetParams.offset + requested, bytes.byteLength) + ) + const eof = assetParams.offset + source.byteLength === bytes.byteLength + const payload = + options.alterLastRangeByte && eof && source.byteLength > 1 + ? Buffer.from(source).fill(0x62, source.byteLength - 1) + : source + if (method !== 'mobileWeb.package.asset.gzip') { + return success({ + buildId: assetParams.buildId, + path: assetParams.path, + offset: assetParams.offset, + byteLength: payload.byteLength, + sha256: sha256Hex(payload), + dataBase64: Buffer.from(payload).toString('base64'), + eof + }) + } + const encoded = gzipSync(Buffer.from(payload), { mtime: 0 }) + return success({ + buildId: assetParams.buildId, + path: assetParams.path, + offset: assetParams.offset, + sourceByteLength: payload.byteLength, + byteLength: encoded.byteLength, + sha256: sha256Hex(encoded), + dataBase64: Buffer.from(encoded).toString('base64'), + eof, + encoding: 'gzip' + }) + } finally { + inFlight -= 1 + } + }) + return { manifest, bytesByPath, request, paramsByCall, peakInFlight: () => peakInFlight } +} + +function createStager() { + return { + begin: vi.fn(async () => {}), + writeAssetChunk: vi.fn(async () => {}), + finishAsset: vi.fn(async () => {}), + commit: vi.fn(async (manifest: MobileWebManifest) => ({ generation: manifest.buildId })), + abort: vi.fn(async () => {}) + } satisfies MobileWebPackageStager<{ generation: string }> +} + +function asset( + path: string, + bytes: Uint8Array, + contentType: MobileWebAsset['contentType'], + role: MobileWebAsset['role'] +): MobileWebAsset { + return { path, sha256: sha256Hex(bytes), byteLength: bytes.byteLength, contentType, role } +} + +function success(result: unknown): RpcResponse { + return { id: 'request', ok: true, result, _meta: { runtimeId: 'runtime' } } +} + +function failure(message: string): RpcResponse { + return { + id: 'request', + ok: false, + error: { code: 'invalid_argument', message }, + _meta: { runtimeId: 'runtime' } + } +} + +function sha256Hex(bytes: Uint8Array): string { + return Buffer.from(sha256(bytes)).toString('hex') +} diff --git a/mobile/src/mobile-web/mobile-web-package-downloader.ts b/mobile/src/mobile-web/mobile-web-package-downloader.ts index 4ab6d1926c7..d9f8f612181 100644 --- a/mobile/src/mobile-web/mobile-web-package-downloader.ts +++ b/mobile/src/mobile-web/mobile-web-package-downloader.ts @@ -1,6 +1,8 @@ import { Buffer } from 'buffer/' import { sha256 } from '@noble/hashes/sha256' import { + MOBILE_WEB_PACKAGE_MAX_CONCURRENT_READS, + MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES, MobileWebPackageManifestResponseSchema, isMobileWebPackageErrorCode, type MobileWebPackageErrorCode @@ -19,6 +21,9 @@ import { decodeRawMobileWebPackageChunk } from './mobile-web-package-chunk-decoder' +const MOBILE_WEB_PACKAGE_READ_LIMITED_RETRIES = 4 +const MOBILE_WEB_PACKAGE_READ_LIMITED_BACKOFF_MS = 50 + export const MOBILE_WEB_PACKAGE_DOWNLOAD_ERROR_CODES = [ 'cancelled', 'host_error', @@ -68,6 +73,10 @@ type DownloadMobileWebPackageOptions = { shellBridgeVersion: number signal?: AbortSignal useGzip?: boolean + /** Chunk reads kept in flight at once. One round trip per 48 KiB otherwise caps throughput. */ + maxConcurrentRequests?: number + /** Bytes per gzip read. Only send above chunkBytes when the host advertises range reads. */ + rangeBytes?: number onProgress?: (progress: MobileWebPackageDownloadProgress) => void } @@ -138,25 +147,25 @@ export async function downloadMobileWebPackage( await stager.begin(manifest) stagingStarted = true let completedBytes = 0 - for (const asset of manifest.assets) { - await downloadAsset( - request, - stager, - manifest, - asset, - chunkBytes, - options.signal, - options.useGzip, - (writtenBytes) => { - completedBytes += writtenBytes - options.onProgress?.({ - phase: 'downloading', - completedBytes, - totalBytes: manifest.totalBytes - }) - } - ) - } + await downloadAssetChunks({ + request, + stager, + manifest, + chunkBytes, + signal: options.signal, + useGzip: options.useGzip ?? false, + rangeBytes: options.useGzip ? clampRangeBytes(chunkBytes, options.rangeBytes) : chunkBytes, + maxConcurrentRequests: + options.maxConcurrentRequests ?? MOBILE_WEB_PACKAGE_MAX_CONCURRENT_READS, + onChunkWritten: (writtenBytes) => { + completedBytes += writtenBytes + options.onProgress?.({ + phase: 'downloading', + completedBytes, + totalBytes: manifest.totalBytes + }) + } + }) throwIfAborted(options.signal) options.onProgress?.({ phase: 'verifying', @@ -182,58 +191,149 @@ export async function downloadMobileWebPackage( } } -async function downloadAsset( - request: MobileWebPackageRequest, - stager: MobileWebPackageStager, - manifest: MobileWebManifest, - asset: MobileWebAsset, - chunkBytes: number, - signal: AbortSignal | undefined, - useGzip = false, - onChunkWritten?: (bytes: number) => void -): Promise { - const assetHash = sha256.create() - for (let offset = 0; offset < asset.byteLength; offset += chunkBytes) { - throwIfAborted(signal) - const result = await requestResult( - request, - useGzip ? 'mobileWeb.package.asset.gzip' : 'mobileWeb.package.asset', - { - buildId: manifest.buildId, - path: asset.path, - offset +type ChunkTask = { asset: MobileWebAsset; offset: number; expectedLength: number } +type SettledChunk = { bytes: Uint8Array } | { failure: unknown } + +// The native stage appends each asset chunk at the file's current length, so chunks must reach +// the stager in offset order even though the reads themselves overlap. +async function downloadAssetChunks(args: { + request: MobileWebPackageRequest + stager: MobileWebPackageStager + manifest: MobileWebManifest + chunkBytes: number + signal: AbortSignal | undefined + useGzip: boolean + rangeBytes: number + maxConcurrentRequests: number + onChunkWritten: (bytes: number) => void +}): Promise { + const tasks = planChunkTasks(args.manifest, args.rangeBytes) + const inFlight = new Map>() + let window = clampWindow(args.maxConcurrentRequests) + let issued = 0 + let assetHash = sha256.create() + + const shrinkWindow = (): void => { + window = Math.max(1, window - 1) + } + for (let drained = 0; drained < tasks.length; drained += 1) { + while (issued < tasks.length && issued - drained < window) { + const task = tasks[issued]! + inFlight.set(issued, fetchChunk(args, task, shrinkWindow)) + issued += 1 + } + const settled = await inFlight.get(drained)! + inFlight.delete(drained) + if ('failure' in settled) { + throw settled.failure + } + throwIfAborted(args.signal) + const task = tasks[drained]! + assetHash.update(settled.bytes) + // A ranged read answers several stage chunks at once; the native stage still appends + // one 48 KiB chunk at a time. + for (let written = 0; written < settled.bytes.byteLength; written += args.chunkBytes) { + const slice = settled.bytes.subarray(written, written + args.chunkBytes) + await args.stager.writeAssetChunk(task.asset, task.offset + written, slice) + args.onChunkWritten(slice.byteLength) + } + if (task.offset + task.expectedLength === task.asset.byteLength) { + if (Buffer.from(assetHash.digest()).toString('hex') !== task.asset.sha256) { + throw new MobileWebPackageDownloadError('asset_integrity_failed') } - ) - throwIfAborted(signal) - const expectedLength = Math.min(chunkBytes, asset.byteLength - offset) - const bytes = useGzip - ? decodeGzipMobileWebPackageChunk( - result, - manifest.buildId, - asset.path, - offset, - expectedLength, - asset.byteLength - ) - : decodeRawMobileWebPackageChunk( - result, - manifest.buildId, - asset.path, - offset, - expectedLength, - asset.byteLength - ) - if (!bytes) { - throw new MobileWebPackageDownloadError('invalid_chunk') + assetHash = sha256.create() + await args.stager.finishAsset(task.asset) } - assetHash.update(bytes) - await stager.writeAssetChunk(asset, offset, bytes) - onChunkWritten?.(bytes.byteLength) } - if (Buffer.from(assetHash.digest()).toString('hex') !== asset.sha256) { - throw new MobileWebPackageDownloadError('asset_integrity_failed') +} + +function planChunkTasks(manifest: MobileWebManifest, rangeBytes: number): ChunkTask[] { + const tasks: ChunkTask[] = [] + for (const asset of manifest.assets) { + for (let offset = 0; offset < asset.byteLength; offset += rangeBytes) { + tasks.push({ + asset, + offset, + expectedLength: Math.min(rangeBytes, asset.byteLength - offset) + }) + } + } + return tasks +} + +// Ranged reads ride mobileWeb.package.asset.gzip's optional `length`, which strict-schema +// hosts predating MOBILE_WEB_PACKAGE_RANGE_RUNTIME_CAPABILITY reject. +function clampRangeBytes(chunkBytes: number, rangeBytes: number | undefined): number { + if (rangeBytes === undefined || rangeBytes <= chunkBytes) { + return chunkBytes + } + const chunks = Math.min( + Math.floor(MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES / chunkBytes), + Math.floor(rangeBytes / chunkBytes) + ) + return Math.max(1, chunks) * chunkBytes +} + +function clampWindow(maxConcurrentRequests: number): number { + return Number.isInteger(maxConcurrentRequests) + ? Math.min(MOBILE_WEB_PACKAGE_MAX_CONCURRENT_READS, Math.max(1, maxConcurrentRequests)) + : 1 +} + +// Never rejects: a queued read that fails while an earlier one is still draining would +// otherwise surface as an unhandled rejection before the drain loop reaches it. +async function fetchChunk( + args: { + request: MobileWebPackageRequest + stager: MobileWebPackageStager + manifest: MobileWebManifest + signal: AbortSignal | undefined + useGzip: boolean + chunkBytes: number + rangeBytes: number + }, + task: ChunkTask, + onReadLimited: () => void +): Promise { + for (let attempt = 0; ; attempt += 1) { + try { + throwIfAborted(args.signal) + const result = await requestResult( + args.request, + args.useGzip ? 'mobileWeb.package.asset.gzip' : 'mobileWeb.package.asset', + { + buildId: args.manifest.buildId, + path: task.asset.path, + offset: task.offset, + ...(args.rangeBytes > args.chunkBytes ? { length: args.rangeBytes } : {}) + } + ) + throwIfAborted(args.signal) + const decode = args.useGzip ? decodeGzipMobileWebPackageChunk : decodeRawMobileWebPackageChunk + const bytes = decode( + result, + args.manifest.buildId, + task.asset.path, + task.offset, + task.expectedLength, + task.asset.byteLength + ) + return bytes ? { bytes } : { failure: new MobileWebPackageDownloadError('invalid_chunk') } + } catch (error) { + // Why: a host whose per-connection read budget is narrower than ours must slow the + // pipeline down, not fail the whole package. + if ( + error instanceof MobileWebPackageDownloadError && + error.code === 'mobile_web_package_read_limited' && + attempt < MOBILE_WEB_PACKAGE_READ_LIMITED_RETRIES + ) { + onReadLimited() + await sleep(MOBILE_WEB_PACKAGE_READ_LIMITED_BACKOFF_MS * (attempt + 1)) + continue + } + return { failure: error } + } } - await stager.finishAsset(asset) } async function requestResult( @@ -279,6 +379,10 @@ function mobileWebPackageHostFailureCode(code: string): MobileWebPackageDownload } } +function sleep(durationMs: number): Promise { + return new Promise((resolve) => setTimeout(resolve, durationMs)) +} + function throwIfAborted(signal: AbortSignal | undefined): void { if (signal?.aborted) { throw new MobileWebPackageDownloadError('cancelled') diff --git a/mobile/src/mobile-web/use-mobile-web-package-capability.ts b/mobile/src/mobile-web/use-mobile-web-package-capability.ts index c6ec94889f7..9dbc6665be8 100644 --- a/mobile/src/mobile-web/use-mobile-web-package-capability.ts +++ b/mobile/src/mobile-web/use-mobile-web-package-capability.ts @@ -1,6 +1,7 @@ import { useEffect, useMemo, useState } from 'react' import { MOBILE_WEB_PACKAGE_GZIP_RUNTIME_CAPABILITY, + MOBILE_WEB_PACKAGE_RANGE_RUNTIME_CAPABILITY, MOBILE_WEB_PACKAGE_RUNTIME_CAPABILITY } from '../../../src/shared/protocol-version' import type { RpcClient } from '../transport/rpc-client' @@ -16,6 +17,7 @@ export type MobileWebPackageCapabilityStatus = export type MobileWebPackageCapability = { status: MobileWebPackageCapabilityStatus gzip: boolean + range: boolean } type ResolvedPackageCapability = { @@ -24,6 +26,7 @@ type ResolvedPackageCapability = { connectionId: number | null supported: boolean gzip: boolean + range: boolean } function currentConnectionId(client: RpcClient | null): number | null { @@ -54,23 +57,25 @@ export function useMobileWebPackageCapability(args: { hostId: requestHostId, connectionId: requestConnectionId, supported: capabilities.includes(MOBILE_WEB_PACKAGE_RUNTIME_CAPABILITY), - gzip: capabilities.includes(MOBILE_WEB_PACKAGE_GZIP_RUNTIME_CAPABILITY) + gzip: capabilities.includes(MOBILE_WEB_PACKAGE_GZIP_RUNTIME_CAPABILITY), + range: capabilities.includes(MOBILE_WEB_PACKAGE_RANGE_RUNTIME_CAPABILITY) }) }) }, [client, connectionId, hostId, state]) const capability = state !== 'connected' - ? { status: 'offline' as const, gzip: false } + ? { status: 'offline' as const, gzip: false, range: false } : !client || !hostId || resolved?.client !== client || resolved.hostId !== hostId || resolved.connectionId !== connectionId - ? { status: 'pending' as const, gzip: false } + ? { status: 'pending' as const, gzip: false, range: false } : { status: resolved.supported ? ('supported' as const) : ('update-required' as const), - gzip: resolved.gzip + gzip: resolved.gzip, + range: resolved.range } - return useMemo(() => capability, [capability.gzip, capability.status]) + return useMemo(() => capability, [capability.gzip, capability.range, capability.status]) } diff --git a/mobile/src/mobile-web/use-mobile-web-package-refresh.ts b/mobile/src/mobile-web/use-mobile-web-package-refresh.ts index 500db6f64ce..e3704fbd55a 100644 --- a/mobile/src/mobile-web/use-mobile-web-package-refresh.ts +++ b/mobile/src/mobile-web/use-mobile-web-package-refresh.ts @@ -1,6 +1,7 @@ import { useEffect, type MutableRefObject } from 'react' import ExpoMobileWebShell, { type MobileWebShellSession } from '@orca/expo-mobile-web-shell' import { MOBILE_WEB_BRIDGE_PROTOCOL_VERSION } from '../../../src/shared/mobile-web/bridge-contract' +import { MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES } from '../../../src/shared/mobile-web/package-rpc-contract' import type { RpcClient } from '../transport/rpc-client' import type { ConnectionState, HostProfile } from '../transport/types' import { createMobileWebNativeStager } from './mobile-web-native-stager' @@ -75,6 +76,7 @@ export function useMobileWebPackageRefresh(args: { { shellBridgeVersion: MOBILE_WEB_BRIDGE_PROTOCOL_VERSION, useGzip: packageCapability.gzip, + ...(packageCapability.range ? { rangeBytes: MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES } : {}), signal: controller.signal, onProgress: setPackageProgress, reuseVerifiedBuild: async (buildId) => { diff --git a/mobile/tsconfig.json b/mobile/tsconfig.json index 46b7cffcd5e..7ee087ca785 100644 --- a/mobile/tsconfig.json +++ b/mobile/tsconfig.json @@ -13,5 +13,11 @@ } }, "include": ["**/*.ts", "**/*.tsx"], - "exclude": ["**/*.test.ts", "**/*.test.tsx"] + // The measurement harness drives the desktop main-process reader directly; that code is + // typed against Node's lib, which this React Native program cannot model. + "exclude": [ + "**/*.test.ts", + "**/*.test.tsx", + "scripts/mobile-web-package-download-measurement.ts" + ] } diff --git a/src/main/runtime/rpc/mobile-web-package-assets.ts b/src/main/runtime/rpc/mobile-web-package-assets.ts index 15c1c1b991c..2da21be10d3 100644 --- a/src/main/runtime/rpc/mobile-web-package-assets.ts +++ b/src/main/runtime/rpc/mobile-web-package-assets.ts @@ -1,7 +1,5 @@ -import { createHash } from 'node:crypto' -import { open, readFile, stat } from 'node:fs/promises' +import { stat } from 'node:fs/promises' import { gzipSync } from 'node:zlib' -import { isAbsolute, relative, resolve } from 'node:path' import { MOBILE_WEB_PACKAGE_MAX_CONCURRENT_READS, MOBILE_WEB_PACKAGE_MAX_IN_FLIGHT_BYTES, @@ -15,20 +13,19 @@ import { } from '../../../shared/mobile-web/package-rpc-contract' import { MOBILE_WEB_PACKAGE_CHUNK_BYTES, - MobileWebManifestSchema, - serializeMobileWebManifestForBuildId, - type MobileWebAsset, - type MobileWebManifest + type MobileWebAsset } from '../../../shared/mobile-web/manifest-contract' import { resolveMobileWebPackageRoot } from './mobile-web-package-root' - -type VerifiedMobileWebPackage = { - root: string - manifestFingerprint: string - manifest: MobileWebManifest - assetsByPath: ReadonlyMap - fileStatsByPath: ReadonlyMap -} +import { + assertManifestFingerprint, + readAssetRange, + readManifestBytes, + resolveDeclaredAssetPath, + sha256, + throwIfAborted, + verifyPackage, + type VerifiedMobileWebPackage +} from './mobile-web-package-verification' type PackageReadState = { count: number; bytes: number } type AssetRangeReader = (path: string, offset: number, length: number) => Promise @@ -67,10 +64,31 @@ export class MobileWebPackageAssets { params: MobileWebPackageAssetParams, options: { connectionId?: string; signal?: AbortSignal } = {} ): Promise { + // Ranged reads exist only on the gzip method; the raw chunk schema caps at one chunk. + if (params.length !== undefined && params.length !== MOBILE_WEB_PACKAGE_CHUNK_BYTES) { + throw new Error('mobile_web_package_offset_invalid') + } + const read = await this.readVerifiedRange(params, options, MOBILE_WEB_PACKAGE_CHUNK_BYTES) + return MobileWebPackageAssetChunkSchema.parse({ + buildId: read.buildId, + path: read.asset.path, + offset: params.offset, + byteLength: read.bytes.byteLength, + sha256: sha256(read.bytes), + dataBase64: read.bytes.toString('base64'), + eof: read.eof + }) + } + + private async readVerifiedRange( + params: MobileWebPackageAssetParams, + options: { connectionId?: string; signal?: AbortSignal }, + requestedLength: number + ): Promise<{ buildId: string; asset: MobileWebAsset; bytes: Buffer; eof: boolean }> { throwIfAborted(options.signal) const verified = await this.getVerifiedPackage() const asset = this.validateAssetParams(verified, params) - const byteLength = Math.min(MOBILE_WEB_PACKAGE_CHUNK_BYTES, asset.byteLength - params.offset) + const byteLength = Math.min(requestedLength, asset.byteLength - params.offset) const release = this.acquireRead(options.connectionId ?? 'in-process', byteLength) try { const path = resolveDeclaredAssetPath(verified.root, asset.path) @@ -87,15 +105,12 @@ export class MobileWebPackageAssets { } await assertManifestFingerprint(verified.root, verified.manifestFingerprint) throwIfAborted(options.signal) - return MobileWebPackageAssetChunkSchema.parse({ + return { buildId: verified.manifest.buildId, - path: asset.path, - offset: params.offset, - byteLength, - sha256: sha256(bytes), - dataBase64: bytes.toString('base64'), + asset, + bytes, eof: params.offset + byteLength === asset.byteLength - }) + } } finally { release() } @@ -129,52 +144,69 @@ export class MobileWebPackageAssets { params: MobileWebPackageAssetParams, options: { connectionId?: string; signal?: AbortSignal } = {} ): Promise { - throwIfAborted(options.signal) - const verified = await this.getVerifiedPackage() - const asset = this.validateAssetParams(verified, params) - const key = `${params.buildId}:${params.path}:${params.offset}` - let compressed = this.gzipChunks.get(key) - let raw: MobileWebPackageAssetChunk - if (!compressed) { - raw = await this.getAssetChunk(params, options) - compressed = gzipSync(Buffer.from(raw.dataBase64, 'base64'), { level: 6 }) - this.gzipChunks.set(key, compressed) - } else { - const path = resolveDeclaredAssetPath(verified.root, asset.path) - const expectedStat = verified.fileStatsByPath.get(asset.path)! - const currentStat = await stat(path) - if (currentStat.size !== expectedStat.size || currentStat.mtimeMs !== expectedStat.mtimeMs) { - this.cached = null - throw new Error('mobile_web_package_asset_changed') - } - await assertManifestFingerprint(verified.root, verified.manifestFingerprint) - const rawByteLength = Math.min( - MOBILE_WEB_PACKAGE_CHUNK_BYTES, - asset.byteLength - params.offset + const requestedLength = params.length ?? MOBILE_WEB_PACKAGE_CHUNK_BYTES + const key = `${params.buildId}:${params.path}:${params.offset}:${requestedLength}` + const cached = this.gzipChunks.get(key) + if (cached) { + throwIfAborted(options.signal) + const verified = await this.getVerifiedPackage() + const asset = this.validateAssetParams(verified, params) + await this.assertAssetUnchanged(verified, asset) + const sourceByteLength = Math.min(requestedLength, asset.byteLength - params.offset) + return this.gzipResponse( + verified.manifest.buildId, + asset, + params.offset, + sourceByteLength, + cached ) - raw = { - buildId: verified.manifest.buildId, - path: asset.path, - offset: params.offset, - byteLength: rawByteLength, - sha256: '', - dataBase64: '', - eof: params.offset + rawByteLength === asset.byteLength - } } + const read = await this.readVerifiedRange(params, options, requestedLength) + const compressed = gzipSync(read.bytes, { level: 6 }) + this.gzipChunks.set(key, compressed) + return this.gzipResponse( + read.buildId, + read.asset, + params.offset, + read.bytes.byteLength, + compressed + ) + } + + private gzipResponse( + buildId: string, + asset: MobileWebAsset, + offset: number, + sourceByteLength: number, + compressed: Buffer + ): MobileWebPackageGzipAssetChunk { return MobileWebPackageGzipAssetChunkSchema.parse({ - buildId: raw.buildId, - path: raw.path, - offset: raw.offset, - sourceByteLength: raw.byteLength, + buildId, + path: asset.path, + offset, + sourceByteLength, byteLength: compressed.byteLength, sha256: sha256(compressed), dataBase64: compressed.toString('base64'), - eof: raw.eof, + eof: offset + sourceByteLength === asset.byteLength, encoding: 'gzip' }) } + private async assertAssetUnchanged( + verified: VerifiedMobileWebPackage, + asset: MobileWebAsset + ): Promise { + const path = resolveDeclaredAssetPath(verified.root, asset.path) + const expectedStat = verified.fileStatsByPath.get(asset.path)! + const currentStat = await stat(path) + if (currentStat.size !== expectedStat.size || currentStat.mtimeMs !== expectedStat.mtimeMs) { + this.cached = null + throw new Error('mobile_web_package_asset_changed') + } + await assertManifestFingerprint(verified.root, verified.manifestFingerprint) + } + private validateAssetParams( verified: VerifiedMobileWebPackage, params: MobileWebPackageAssetParams @@ -213,89 +245,4 @@ export class MobileWebPackageAssets { } } -async function verifyPackage( - root: string, - manifestBytes: Buffer, - manifestFingerprint: string -): Promise { - const manifest = parseManifest(manifestBytes) - if (sha256(serializeMobileWebManifestForBuildId(manifest)) !== manifest.buildId) { - throw new Error('mobile_web_package_build_invalid') - } - const fileStatsByPath = new Map() - for (const asset of manifest.assets) { - const path = resolveDeclaredAssetPath(root, asset.path) - const beforeRead = await stat(path) - const bytes = await readFile(path) - const afterRead = await stat(path) - if (bytes.byteLength !== asset.byteLength || sha256(bytes) !== asset.sha256) { - throw new Error('mobile_web_package_asset_invalid') - } - if (beforeRead.size !== afterRead.size || beforeRead.mtimeMs !== afterRead.mtimeMs) { - throw new Error('mobile_web_package_asset_changed') - } - fileStatsByPath.set(asset.path, { size: afterRead.size, mtimeMs: afterRead.mtimeMs }) - } - await assertManifestFingerprint(root, manifestFingerprint) - return { - root, - manifestFingerprint, - manifest, - assetsByPath: new Map(manifest.assets.map((asset) => [asset.path, asset])), - fileStatsByPath - } -} - -function parseManifest(manifestBytes: Buffer): MobileWebManifest { - try { - return MobileWebManifestSchema.parse(JSON.parse(manifestBytes.toString('utf8'))) - } catch { - throw new Error('mobile_web_package_build_invalid') - } -} - -async function readManifestBytes(root: string): Promise { - try { - return await readFile(resolve(root, 'manifest.json')) - } catch { - throw new Error('mobile_web_package_unavailable') - } -} - -async function assertManifestFingerprint(root: string, expected: string): Promise { - if (sha256(await readManifestBytes(root)) !== expected) { - throw new Error('mobile_web_package_build_changed') - } -} - -function resolveDeclaredAssetPath(root: string, assetPath: string): string { - const candidate = resolve(root, ...assetPath.split('/')) - const child = relative(root, candidate) - if (child.startsWith('..') || isAbsolute(child)) { - throw new Error('mobile_web_package_asset_path_invalid') - } - return candidate -} - -async function readAssetRange(path: string, offset: number, length: number): Promise { - const file = await open(path, 'r') - try { - const buffer = Buffer.alloc(length) - const { bytesRead } = await file.read(buffer, 0, length, offset) - return buffer.subarray(0, bytesRead) - } finally { - await file.close() - } -} - -function throwIfAborted(signal: AbortSignal | undefined): void { - if (signal?.aborted) { - throw new Error('mobile_web_package_cancelled') - } -} - -function sha256(value: string | Uint8Array): string { - return createHash('sha256').update(value).digest('hex') -} - export const mobileWebPackageAssets = new MobileWebPackageAssets() 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 new file mode 100644 index 00000000000..1c7dfc303c7 --- /dev/null +++ b/src/main/runtime/rpc/mobile-web-package-range-reads.test.ts @@ -0,0 +1,196 @@ +import { createHash } from 'node:crypto' +import { mkdtemp, mkdir, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { gunzipSync } from 'node:zlib' +import { afterEach, describe, expect, it } from 'vitest' +import { MOBILE_WEB_BRIDGE_PROTOCOL_VERSION } from '../../../shared/mobile-web/bridge-contract' +import { + MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES, + MobileWebPackageAssetParamsSchema +} from '../../../shared/mobile-web/package-rpc-contract' +import { + MOBILE_WEB_MANIFEST_SCHEMA_VERSION, + MOBILE_WEB_PACKAGE_CHUNK_BYTES, + serializeMobileWebManifestForBuildId, + type MobileWebManifest +} from '../../../shared/mobile-web/manifest-contract' +import { MobileWebPackageAssets } from './mobile-web-package-assets' + +const temporaryRoots: string[] = [] + +afterEach(async () => { + await Promise.all( + temporaryRoots.splice(0).map((root) => rm(root, { recursive: true, force: true })) + ) +}) + +describe('mobile web package ranged gzip reads', () => { + it('answers a whole multi-chunk range in one response', async () => { + const fixture = await createRangeFixture() + const assets = new MobileWebPackageAssets({ resolveRoot: () => fixture.root }) + const script = fixture.manifest.assets.find((asset) => asset.role === 'script')! + + const chunk = await assets.getAssetGzipChunk( + { + buildId: fixture.manifest.buildId, + path: script.path, + offset: 0, + length: MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + }, + { connectionId: 'connection-1' } + ) + + expect(chunk.sourceByteLength).toBe(MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES) + expect(chunk.eof).toBe(false) + const bytes = gunzipSync(Buffer.from(chunk.dataBase64, 'base64')) + expect(bytes).toEqual(fixture.scriptBytes.subarray(0, MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES)) + }) + + it('clamps the last range to the end of the asset', async () => { + const fixture = await createRangeFixture() + const assets = new MobileWebPackageAssets({ resolveRoot: () => fixture.root }) + const script = fixture.manifest.assets.find((asset) => asset.role === 'script')! + const offset = MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + + const chunk = await assets.getAssetGzipChunk({ + buildId: fixture.manifest.buildId, + path: script.path, + offset, + length: MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + }) + + expect(chunk.sourceByteLength).toBe(script.byteLength - offset) + expect(chunk.eof).toBe(true) + expect(gunzipSync(Buffer.from(chunk.dataBase64, 'base64'))).toEqual( + fixture.scriptBytes.subarray(offset) + ) + }) + + it('reproduces a cached range byte for byte', async () => { + const fixture = await createRangeFixture() + const assets = new MobileWebPackageAssets({ resolveRoot: () => fixture.root }) + const script = fixture.manifest.assets.find((asset) => asset.role === 'script')! + const params = { + buildId: fixture.manifest.buildId, + path: script.path, + offset: 0, + length: MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + } + + const first = await assets.getAssetGzipChunk(params) + const second = await assets.getAssetGzipChunk(params) + + expect(second).toEqual(first) + }) + + 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 }) + const script = fixture.manifest.assets.find((asset) => asset.role === 'script')! + const address = { buildId: fixture.manifest.buildId, path: script.path, offset: 0 } + + const single = await assets.getAssetGzipChunk(address) + const ranged = await assets.getAssetGzipChunk({ + ...address, + length: MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + }) + + expect(single.sourceByteLength).toBe(MOBILE_WEB_PACKAGE_CHUNK_BYTES) + expect(ranged.sourceByteLength).toBe(MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES) + }) + + it('refuses a ranged read on the uncompressed chunk method', async () => { + const fixture = await createRangeFixture() + const assets = new MobileWebPackageAssets({ resolveRoot: () => fixture.root }) + const script = fixture.manifest.assets.find((asset) => asset.role === 'script')! + + await expect( + assets.getAssetChunk({ + buildId: fixture.manifest.buildId, + path: script.path, + offset: 0, + length: MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + }) + ).rejects.toThrow('mobile_web_package_offset_invalid') + }) + + it('only accepts chunk-aligned range lengths within the cap', async () => { + const address = { buildId: '0'.repeat(64), path: 'index.html', offset: 0 } + + expect( + MobileWebPackageAssetParamsSchema.safeParse({ + ...address, + length: MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + }).success + ).toBe(true) + expect(MobileWebPackageAssetParamsSchema.safeParse(address).success).toBe(true) + expect( + MobileWebPackageAssetParamsSchema.safeParse({ + ...address, + length: MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + MOBILE_WEB_PACKAGE_CHUNK_BYTES + }).success + ).toBe(false) + expect( + MobileWebPackageAssetParamsSchema.safeParse({ + ...address, + length: MOBILE_WEB_PACKAGE_CHUNK_BYTES + 1 + }).success + ).toBe(false) + expect(MobileWebPackageAssetParamsSchema.safeParse({ ...address, length: 0 }).success).toBe( + false + ) + }) +}) + +async function createRangeFixture(): Promise<{ + root: string + manifest: MobileWebManifest + scriptBytes: Buffer +}> { + const root = await mkdtemp(join(tmpdir(), 'orca-mobile-web-range-')) + temporaryRoots.push(root) + const scriptBytes = Buffer.alloc( + MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + MOBILE_WEB_PACKAGE_CHUNK_BYTES + 23, + 0x61 + ) + const documentBytes = Buffer.from('Orca', 'utf8') + const scriptHash = sha256(scriptBytes) + const assets = [ + { + path: `assets/${scriptHash}.js`, + sha256: scriptHash, + byteLength: scriptBytes.byteLength, + contentType: 'text/javascript; charset=utf-8' as const, + role: 'script' as const + }, + { + path: 'index.html', + sha256: sha256(documentBytes), + byteLength: documentBytes.byteLength, + contentType: 'text/html; charset=utf-8' as const, + role: 'document' as const + } + ] + const seed: MobileWebManifest = { + schemaVersion: MOBILE_WEB_MANIFEST_SCHEMA_VERSION, + buildId: '0'.repeat(64), + bridge: { + minimum: MOBILE_WEB_BRIDGE_PROTOCOL_VERSION, + testedThrough: MOBILE_WEB_BRIDGE_PROTOCOL_VERSION + }, + entrypoint: 'index.html', + totalBytes: assets.reduce((total, asset) => total + asset.byteLength, 0), + assets + } + const manifest = { ...seed, buildId: sha256(serializeMobileWebManifestForBuildId(seed)) } + await mkdir(join(root, 'assets'), { recursive: true }) + await writeFile(join(root, ...assets[0]!.path.split('/')), scriptBytes) + await writeFile(join(root, 'index.html'), documentBytes) + await writeFile(join(root, 'manifest.json'), JSON.stringify(manifest)) + return { root, manifest, scriptBytes } +} + +function sha256(value: string | Uint8Array): string { + return createHash('sha256').update(value).digest('hex') +} diff --git a/src/main/runtime/rpc/mobile-web-package-verification.ts b/src/main/runtime/rpc/mobile-web-package-verification.ts new file mode 100644 index 00000000000..5dd8478b509 --- /dev/null +++ b/src/main/runtime/rpc/mobile-web-package-verification.ts @@ -0,0 +1,106 @@ +import { createHash } from 'node:crypto' +import { open, readFile, stat } from 'node:fs/promises' +import { isAbsolute, relative, resolve } from 'node:path' +import { + MobileWebManifestSchema, + serializeMobileWebManifestForBuildId, + type MobileWebAsset, + type MobileWebManifest +} from '../../../shared/mobile-web/manifest-contract' + +export type VerifiedMobileWebPackage = { + root: string + manifestFingerprint: string + manifest: MobileWebManifest + assetsByPath: ReadonlyMap + fileStatsByPath: ReadonlyMap +} + +export async function verifyPackage( + root: string, + manifestBytes: Buffer, + manifestFingerprint: string +): Promise { + const manifest = parseManifest(manifestBytes) + if (sha256(serializeMobileWebManifestForBuildId(manifest)) !== manifest.buildId) { + throw new Error('mobile_web_package_build_invalid') + } + const fileStatsByPath = new Map() + for (const asset of manifest.assets) { + const path = resolveDeclaredAssetPath(root, asset.path) + const beforeRead = await stat(path) + const bytes = await readFile(path) + const afterRead = await stat(path) + if (bytes.byteLength !== asset.byteLength || sha256(bytes) !== asset.sha256) { + throw new Error('mobile_web_package_asset_invalid') + } + if (beforeRead.size !== afterRead.size || beforeRead.mtimeMs !== afterRead.mtimeMs) { + throw new Error('mobile_web_package_asset_changed') + } + fileStatsByPath.set(asset.path, { size: afterRead.size, mtimeMs: afterRead.mtimeMs }) + } + await assertManifestFingerprint(root, manifestFingerprint) + return { + root, + manifestFingerprint, + manifest, + assetsByPath: new Map(manifest.assets.map((asset) => [asset.path, asset])), + fileStatsByPath + } +} + +function parseManifest(manifestBytes: Buffer): MobileWebManifest { + try { + return MobileWebManifestSchema.parse(JSON.parse(manifestBytes.toString('utf8'))) + } catch { + throw new Error('mobile_web_package_build_invalid') + } +} + +export async function readManifestBytes(root: string): Promise { + try { + return await readFile(resolve(root, 'manifest.json')) + } catch { + throw new Error('mobile_web_package_unavailable') + } +} + +export async function assertManifestFingerprint(root: string, expected: string): Promise { + if (sha256(await readManifestBytes(root)) !== expected) { + throw new Error('mobile_web_package_build_changed') + } +} + +export function resolveDeclaredAssetPath(root: string, assetPath: string): string { + const candidate = resolve(root, ...assetPath.split('/')) + const child = relative(root, candidate) + if (child.startsWith('..') || isAbsolute(child)) { + throw new Error('mobile_web_package_asset_path_invalid') + } + return candidate +} + +export async function readAssetRange( + path: string, + offset: number, + length: number +): Promise { + const file = await open(path, 'r') + try { + const buffer = Buffer.alloc(length) + const { bytesRead } = await file.read(buffer, 0, length, offset) + return buffer.subarray(0, bytesRead) + } finally { + await file.close() + } +} + +export function throwIfAborted(signal: AbortSignal | undefined): void { + if (signal?.aborted) { + throw new Error('mobile_web_package_cancelled') + } +} + +export function sha256(value: string | Uint8Array): string { + return createHash('sha256').update(value).digest('hex') +} diff --git a/src/shared/mobile-web/package-rpc-contract.ts b/src/shared/mobile-web/package-rpc-contract.ts index 1c46ad423cc..89dd9cbb9f1 100644 --- a/src/shared/mobile-web/package-rpc-contract.ts +++ b/src/shared/mobile-web/package-rpc-contract.ts @@ -7,12 +7,18 @@ import { import { isMobileWebBase64, isMobileWebSha256 } from './protocol-token-contract' export const MOBILE_WEB_PACKAGE_MAX_CONCURRENT_READS = 4 +// A ranged gzip read answers this many 48 KiB chunks in one round trip. 384 KiB of +// incompressible bytes still fits the relay's 1 MiB control-lane frame after gzip, +// base64, and the E2EE layer's own base64 (a ~1.78x expansion in total). +export const MOBILE_WEB_PACKAGE_RANGE_CHUNKS = 8 +export const MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES = + MOBILE_WEB_PACKAGE_RANGE_CHUNKS * MOBILE_WEB_PACKAGE_CHUNK_BYTES export const MOBILE_WEB_PACKAGE_MAX_IN_FLIGHT_BYTES = - MOBILE_WEB_PACKAGE_MAX_CONCURRENT_READS * MOBILE_WEB_PACKAGE_CHUNK_BYTES + MOBILE_WEB_PACKAGE_MAX_CONCURRENT_READS * MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES export const MOBILE_WEB_PACKAGE_CHUNK_BASE64_CHARS = Math.ceil(MOBILE_WEB_PACKAGE_CHUNK_BYTES / 3) * 4 // Deflate can add a small stored-block header to incompressible chunks. -export const MOBILE_WEB_PACKAGE_GZIP_CHUNK_BYTES = MOBILE_WEB_PACKAGE_CHUNK_BYTES + 64 +export const MOBILE_WEB_PACKAGE_GZIP_CHUNK_BYTES = MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES + 64 export const MOBILE_WEB_PACKAGE_GZIP_CHUNK_BASE64_CHARS = Math.ceil(MOBILE_WEB_PACKAGE_GZIP_CHUNK_BYTES / 3) * 4 @@ -45,7 +51,17 @@ export const MobileWebPackageAssetParamsSchema = z .object({ buildId: z.string().refine(isMobileWebSha256), path: AssetPathSchema, - offset: z.number().int().nonnegative() + offset: z.number().int().nonnegative(), + // Optional, and only on mobileWeb.package.asset.gzip: hosts predating + // MOBILE_WEB_PACKAGE_RANGE_RUNTIME_CAPABILITY reject it (the schema is strict), so + // clients must not send it before the capability probe reports the host supports it. + length: z + .number() + .int() + .positive() + .max(MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES) + .refine((length) => length % MOBILE_WEB_PACKAGE_CHUNK_BYTES === 0) + .optional() }) .strict() @@ -75,7 +91,7 @@ export const MobileWebPackageGzipAssetChunkSchema = z buildId: z.string().refine(isMobileWebSha256), path: AssetPathSchema, offset: z.number().int().nonnegative(), - sourceByteLength: z.number().int().positive().max(MOBILE_WEB_PACKAGE_CHUNK_BYTES), + sourceByteLength: z.number().int().positive().max(MOBILE_WEB_PACKAGE_MAX_RANGE_BYTES), byteLength: z.number().int().positive().max(MOBILE_WEB_PACKAGE_GZIP_CHUNK_BYTES), sha256: z.string().refine(isMobileWebSha256), dataBase64: z diff --git a/src/shared/protocol-version.ts b/src/shared/protocol-version.ts index be4a1782546..f5becb24905 100644 --- a/src/shared/protocol-version.ts +++ b/src/shared/protocol-version.ts @@ -110,6 +110,10 @@ export const CODEX_RESET_CREDIT_RUNTIME_CAPABILITY = 'accounts.codex-reset-credi export const ACCOUNT_IMPORT_RUNTIME_CAPABILITY = 'accounts.import-host-credentials.v1' as const export const MOBILE_WEB_PACKAGE_RUNTIME_CAPABILITY = 'mobileWeb.package.v1' as const export const MOBILE_WEB_PACKAGE_GZIP_RUNTIME_CAPABILITY = 'mobileWeb.package.gzip.v1' as const +// Why: mobileWeb.package.asset.gzip's params schema is strict, so an optional `length` +// reaches older hosts as invalid_argument. Clients may only request a multi-chunk range +// once the host advertises it. +export const MOBILE_WEB_PACKAGE_RANGE_RUNTIME_CAPABILITY = 'mobileWeb.package.range.v1' as const // Why: older hosts cannot reconcile terminal.create's mutation after losing the reply, so clients may only retry unknown outcomes when advertised. export const TERMINAL_CREATE_IDEMPOTENCY_RUNTIME_CAPABILITY = 'terminal.create-idempotency.v2' as const @@ -238,6 +242,7 @@ export const RUNTIME_CAPABILITIES = [ CODEX_RESET_CREDIT_RUNTIME_CAPABILITY, MOBILE_WEB_PACKAGE_RUNTIME_CAPABILITY, MOBILE_WEB_PACKAGE_GZIP_RUNTIME_CAPABILITY, + MOBILE_WEB_PACKAGE_RANGE_RUNTIME_CAPABILITY, SKILL_INSTALL_CAPABILITY, SKILL_BUNDLE_INSTALL_CAPABILITY, SKILL_INSTALL_CANCEL_CAPABILITY,