mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 16:02:56 +00:00
perf(mobile): pipeline and widen mobile-web package reads
The hybrid app pulled its 9.1 MiB mobile-web package one 48 KiB chunk per RPC round trip, strictly serially, and 99.3% of the package is a single 9.5 MB script — so the whole download was 245 sequential round trips against whatever the relay path's latency happened to be. Measured on the real package over the real mobile relay client with a 120 ms round trip: 31.7 s (0.30 MB/s). Two changes, both negotiated: - the downloader now keeps up to MOBILE_WEB_PACKAGE_MAX_CONCURRENT_READS reads in flight across the whole manifest while still draining to the stager in offset order (the native stage appends at the file's current length), and narrows the window instead of failing when a host answers mobile_web_package_read_limited; - mobileWeb.package.asset.gzip takes an optional `length` so one read answers up to eight chunks. The params schema is strict, so older hosts reject the field — it is gated on the new mobileWeb.package.range.v1 capability, and the client still splits the range into 48 KiB stage writes. Same package, same harness, relay path: 31.7 s -> 2.8 s at 120 ms RTT and 62.9 s -> 5.5 s at 250 ms RTT, with 245 requests down to 76. The download still dies when the app is backgrounded past RELAY_BACKGROUND_GRACE_MS — the session suspends, the refresh effect aborts, and the stage is discarded — so the progress panel now says to keep the app open. Resuming from the staged offset needs a native stage-reopen API on both platforms and is not attempted here. mobile/scripts/measure-mobile-web-package-download.mjs drives the real downloader over both the direct and relay mobile clients with an injectable round trip, which is where every number above comes from. Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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')
|
||||
}
|
||||
@@ -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<string, unknown> }
|
||||
|
||||
type Teardown = (() => Promise<void> | void)[]
|
||||
|
||||
export async function measureMobileWebPackageDownload(
|
||||
options: MeasurementOptions
|
||||
): Promise<MeasurementResult> {
|
||||
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<RpcClient> {
|
||||
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<RpcClient> {
|
||||
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<WebSocket, string>()
|
||||
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<number> {
|
||||
await new Promise<void>((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<void>((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<unknown>, 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<void> {
|
||||
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}`)
|
||||
}
|
||||
@@ -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())
|
||||
}
|
||||
})
|
||||
@@ -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({
|
||||
<Text style={styles.packageProgressBytes}>
|
||||
{formatBytes(progress.completedBytes)} of {formatBytes(progress.totalBytes)}
|
||||
</Text>
|
||||
{progress.phase === 'downloading' ? (
|
||||
<Text style={styles.packageProgressBytes}>{STAY_OPEN_HINT}</Text>
|
||||
) : null}
|
||||
</View>
|
||||
)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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<string, Uint8Array>
|
||||
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<typeof createStager>): void {
|
||||
const stagedByPath = new Map<string, Buffer[]>()
|
||||
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('<!doctype html><title>Orca</title>')
|
||||
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<RpcResponse> => {
|
||||
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')
|
||||
}
|
||||
@@ -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<TCommit>(
|
||||
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<TCommit>(
|
||||
}
|
||||
}
|
||||
|
||||
async function downloadAsset<TCommit>(
|
||||
request: MobileWebPackageRequest,
|
||||
stager: MobileWebPackageStager<TCommit>,
|
||||
manifest: MobileWebManifest,
|
||||
asset: MobileWebAsset,
|
||||
chunkBytes: number,
|
||||
signal: AbortSignal | undefined,
|
||||
useGzip = false,
|
||||
onChunkWritten?: (bytes: number) => void
|
||||
): Promise<void> {
|
||||
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<TCommit>(args: {
|
||||
request: MobileWebPackageRequest
|
||||
stager: MobileWebPackageStager<TCommit>
|
||||
manifest: MobileWebManifest
|
||||
chunkBytes: number
|
||||
signal: AbortSignal | undefined
|
||||
useGzip: boolean
|
||||
rangeBytes: number
|
||||
maxConcurrentRequests: number
|
||||
onChunkWritten: (bytes: number) => void
|
||||
}): Promise<void> {
|
||||
const tasks = planChunkTasks(args.manifest, args.rangeBytes)
|
||||
const inFlight = new Map<number, Promise<SettledChunk>>()
|
||||
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<TCommit>(
|
||||
args: {
|
||||
request: MobileWebPackageRequest
|
||||
stager: MobileWebPackageStager<TCommit>
|
||||
manifest: MobileWebManifest
|
||||
signal: AbortSignal | undefined
|
||||
useGzip: boolean
|
||||
chunkBytes: number
|
||||
rangeBytes: number
|
||||
},
|
||||
task: ChunkTask,
|
||||
onReadLimited: () => void
|
||||
): Promise<SettledChunk> {
|
||||
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<void> {
|
||||
return new Promise((resolve) => setTimeout(resolve, durationMs))
|
||||
}
|
||||
|
||||
function throwIfAborted(signal: AbortSignal | undefined): void {
|
||||
if (signal?.aborted) {
|
||||
throw new MobileWebPackageDownloadError('cancelled')
|
||||
|
||||
@@ -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])
|
||||
}
|
||||
|
||||
@@ -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) => {
|
||||
|
||||
@@ -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"
|
||||
]
|
||||
}
|
||||
|
||||
@@ -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<string, MobileWebAsset>
|
||||
fileStatsByPath: ReadonlyMap<string, { size: number; mtimeMs: number }>
|
||||
}
|
||||
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<Buffer>
|
||||
@@ -67,10 +64,31 @@ export class MobileWebPackageAssets {
|
||||
params: MobileWebPackageAssetParams,
|
||||
options: { connectionId?: string; signal?: AbortSignal } = {}
|
||||
): Promise<MobileWebPackageAssetChunk> {
|
||||
// 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<MobileWebPackageGzipAssetChunk> {
|
||||
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<void> {
|
||||
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<VerifiedMobileWebPackage> {
|
||||
const manifest = parseManifest(manifestBytes)
|
||||
if (sha256(serializeMobileWebManifestForBuildId(manifest)) !== manifest.buildId) {
|
||||
throw new Error('mobile_web_package_build_invalid')
|
||||
}
|
||||
const fileStatsByPath = new Map<string, { size: number; mtimeMs: number }>()
|
||||
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<Buffer> {
|
||||
try {
|
||||
return await readFile(resolve(root, 'manifest.json'))
|
||||
} catch {
|
||||
throw new Error('mobile_web_package_unavailable')
|
||||
}
|
||||
}
|
||||
|
||||
async function assertManifestFingerprint(root: string, expected: string): Promise<void> {
|
||||
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<Buffer> {
|
||||
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()
|
||||
|
||||
@@ -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('<!doctype html><title>Orca</title>', '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')
|
||||
}
|
||||
@@ -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<string, MobileWebAsset>
|
||||
fileStatsByPath: ReadonlyMap<string, { size: number; mtimeMs: number }>
|
||||
}
|
||||
|
||||
export async function verifyPackage(
|
||||
root: string,
|
||||
manifestBytes: Buffer,
|
||||
manifestFingerprint: string
|
||||
): Promise<VerifiedMobileWebPackage> {
|
||||
const manifest = parseManifest(manifestBytes)
|
||||
if (sha256(serializeMobileWebManifestForBuildId(manifest)) !== manifest.buildId) {
|
||||
throw new Error('mobile_web_package_build_invalid')
|
||||
}
|
||||
const fileStatsByPath = new Map<string, { size: number; mtimeMs: number }>()
|
||||
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<Buffer> {
|
||||
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<void> {
|
||||
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<Buffer> {
|
||||
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')
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user