refactor(mobile): cache the host catalog per connection and unify host ceilings

The desktop's page-method catalog is a module constant, so its grants can only
change when the desktop process restarts, which tears down the socket and makes
the broker replace the client. Reading it over the wire before every forwarded
request cost a full round-trip per host call and, at mount, several identical
reads at once. Cache grants per client with in-flight dedupe, keyed by method so
a sibling's cancellation cannot fail a peer, and clear it wherever the broker
already clears per-connection authority state.

The host-request concurrency ceiling lived twice: a literal 4 in the accounting
module and a grant limit the shell never read. Take it from the grant alone and
raise it to what a product page needs, counting only one-shot forwards, since
subscriptions have their own ledger cap and catalog reads no longer reach the
wire per request.

Byte ceilings disagreed in both directions. The grant schema let a desktop
advertise a size a shipped shell hard-rejects, so cap it at the shared bridge
envelope, and give the read-heavy desktop methods that whole envelope instead of
an arbitrary 512 KiB. Forwarding also serialized each payload up to four times
per direction; one walk now yields both the verdict and the length.

Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb
This commit is contained in:
Jinwoo-H
2026-09-07 04:40:23 -04:00
parent 15d19c308b
commit 86706dddbd
21 changed files with 495 additions and 116 deletions
@@ -2,6 +2,7 @@ import { MobileWebBrowserAuthority } from './mobile-web-browser-authority'
import { MobileWebAgentHistoryAuthority } from './mobile-web-agent-history-authority'
import { MobileWebAgentHistoryPager } from './mobile-web-agent-history-pager'
import { MobileWebAgentHistoryResume } from './mobile-web-agent-history-resume'
import { MobileWebHostCatalogCache } from './mobile-web-host-catalog-cache'
import { MobileWebNativeChatAuthority } from './mobile-web-native-chat-authority'
import { MobileWebSourceControlBranchComparePager } from './mobile-web-source-control-branch-compare-pager'
import { MobileWebTerminalArtifactAuthority } from './mobile-web-terminal-artifact-authority'
@@ -15,6 +16,7 @@ export class MobileWebCapabilityAuthorities {
readonly agentHistoryPager: MobileWebAgentHistoryPager
readonly agentHistoryResume: MobileWebAgentHistoryResume
readonly browser: MobileWebBrowserAuthority
readonly hostCatalog: MobileWebHostCatalogCache
readonly nativeChat: MobileWebNativeChatAuthority
readonly sourceControlBranchCompare: MobileWebSourceControlBranchComparePager
readonly terminalArtifact: MobileWebTerminalArtifactAuthority
@@ -28,6 +30,7 @@ export class MobileWebCapabilityAuthorities {
this.agentHistoryPager = new MobileWebAgentHistoryPager(options.randomBytes)
this.agentHistoryResume = new MobileWebAgentHistoryResume(options.randomBytes)
this.browser = new MobileWebBrowserAuthority()
this.hostCatalog = new MobileWebHostCatalogCache()
this.nativeChat = new MobileWebNativeChatAuthority(options.randomBytes)
this.sourceControlBranchCompare = new MobileWebSourceControlBranchComparePager()
this.terminalArtifact = new MobileWebTerminalArtifactAuthority(options)
@@ -42,6 +45,7 @@ export class MobileWebCapabilityAuthorities {
this.agentHistoryPager.clear()
this.agentHistoryResume.clear()
this.browser.clear()
this.hostCatalog.clear()
this.nativeChat.clear()
this.sourceControlBranchCompare.clear()
this.terminalArtifact.clear()
@@ -255,6 +255,7 @@ export class MobileWebCapabilityBroker {
sourceControlBranchCompare: this.authorities.sourceControlBranchCompare,
speechAuthority: this.speechAuthority,
workspaceSubscriptions: this.subscriptions.workspace,
hostCatalog: this.authorities.hostCatalog,
hostSubscriptions: this.subscriptions.host,
terminalStreams: this.terminalStreams,
commitMessageGeneration: this.commitMessageGeneration,
@@ -211,6 +211,7 @@ async function subscribeBrowser(args: Deps, request: SubscriptionRequest): Promi
async function subscribeWorkspace(args: Deps, request: SubscriptionRequest): Promise<unknown> {
if (request.operation === 'hostSubscribe') {
await args.hostSubscriptions.start({
catalog: args.hostCatalog,
getPageSessionId: args.getPageSessionId,
requestId: request.requestId,
subscriptionId: request.subscriptionId,
@@ -1,3 +1,4 @@
import type { MobileWebHostCatalogCache } from './mobile-web-host-catalog-cache'
import type { MobileWebHostSubscriptions } from './mobile-web-host-subscriptions'
import type { MobileWebBridgePageMessage } from '../../../src/shared/mobile-web/bridge-contract'
import type { RpcClient } from '../transport/rpc-client'
@@ -33,6 +34,7 @@ export type MobileWebCapabilityExecutionDependencies = {
agentHistoryAuthority: MobileWebAgentHistoryAuthority
agentHistoryPager: MobileWebAgentHistoryPager
agentHistoryResume: MobileWebAgentHistoryResume
hostCatalog: MobileWebHostCatalogCache
hostSubscriptions: MobileWebHostSubscriptions
accountSubscriptions: MobileWebAccountSubscriptions
browserStreams: MobileWebBrowserStreams
@@ -0,0 +1,149 @@
import { describe, expect, it, vi } from 'vitest'
import type { RpcClient } from '../transport/rpc-client'
import { createMobileWebBridgeRoundtripFixture } from './mobile-web-bridge-roundtrip-fixture'
import { MobileWebHostCatalogCache } from './mobile-web-host-catalog-cache'
import { MOBILE_WEB_PRODUCTION_GRANTS } from './mobile-web-production-grants'
const grant = {
method: 'mobileWeb.sourceControl.status',
workspaceParam: 'worktree',
maxRequestBytes: 16 * 1024,
maxResponseBytes: 512 * 1024
}
function cacheFixture() {
const sendRequest = vi
.fn<RpcClient['sendRequest']>()
.mockImplementation(async (_method, input) => ({
ok: true,
result: {
grants: (input as { methods: string[] }).methods
.filter((method) => method !== 'future.absent')
.map((method) => ({ ...grant, method }))
}
}))
return { sendRequest, client: { sendRequest } as unknown as RpcClient }
}
describe('host catalog cache', () => {
it('reads the desktop catalog once per method for a connection', async () => {
const { sendRequest, client } = cacheFixture()
const cache = new MobileWebHostCatalogCache()
await expect(cache.grant(client, grant.method)).resolves.toMatchObject(grant)
await expect(cache.grant(client, grant.method)).resolves.toMatchObject(grant)
expect(sendRequest).toHaveBeenCalledOnce()
})
it('collapses concurrent first reads of one method into a single catalog RPC', async () => {
const { sendRequest, client } = cacheFixture()
const cache = new MobileWebHostCatalogCache()
const resolved = await Promise.all([
cache.grant(client, grant.method),
cache.grant(client, grant.method),
cache.grant(client, grant.method)
])
expect(resolved.every((entry) => entry?.method === grant.method)).toBe(true)
expect(sendRequest).toHaveBeenCalledOnce()
})
it('remembers a method the desktop refused to advertise', async () => {
const { sendRequest, client } = cacheFixture()
const cache = new MobileWebHostCatalogCache()
await expect(cache.grant(client, 'future.absent')).resolves.toBeNull()
await expect(cache.grant(client, 'future.absent')).resolves.toBeNull()
expect(sendRequest).toHaveBeenCalledOnce()
})
it('asks only for the methods it has not resolved yet', async () => {
const { sendRequest, client } = cacheFixture()
const cache = new MobileWebHostCatalogCache()
await cache.read(client, { methods: [grant.method, 'future.absent'] })
await expect(cache.read(client, { methods: [grant.method, 'future.other'] })).resolves.toEqual({
grants: [
expect.objectContaining({ method: grant.method }),
{ ...grant, method: 'future.other' }
]
})
expect(sendRequest).toHaveBeenCalledTimes(2)
expect(sendRequest).toHaveBeenLastCalledWith(
'mobileWeb.host.catalog',
{ methods: ['future.other'] },
expect.any(Object)
)
})
it('re-reads after a client swap discards the connection it was read from', async () => {
const { sendRequest, client } = cacheFixture()
const cache = new MobileWebHostCatalogCache()
await cache.grant(client, grant.method)
cache.clear()
await cache.grant(client, grant.method)
expect(sendRequest).toHaveBeenCalledTimes(2)
await cache.grant({ sendRequest } as unknown as RpcClient, grant.method)
expect(sendRequest).toHaveBeenCalledTimes(3)
})
it('lets a peer retry when the request that owned the shared read fails', async () => {
const { sendRequest, client } = cacheFixture()
const failure = { ok: false as const, error: { code: 'internal_error' } }
sendRequest.mockResolvedValueOnce(failure)
const cache = new MobileWebHostCatalogCache()
const owner = cache.grant(client, grant.method)
const peer = cache.grant(client, grant.method)
await expect(owner).rejects.toMatchObject({ code: 'host_error' })
await expect(peer).resolves.toMatchObject(grant)
expect(sendRequest).toHaveBeenCalledTimes(2)
})
})
describe('host catalog reads across a broker client swap', () => {
it('serves forwarded requests from one catalog read until the client is replaced', async () => {
const sendRequest = vi.fn<RpcClient['sendRequest']>().mockImplementation(async (method) => {
if (method === 'worktree.ps') {
return {
ok: true,
result: {
worktrees: [{ worktreeId: 'host-workspace', repo: '/repo', displayName: 'Workspace' }]
}
}
}
if (method === 'mobileWeb.host.catalog') {
return { ok: true, result: { grants: [grant] } }
}
return {
ok: true,
result: {
entries: [],
conflictOperation: 'unknown',
branch: 'main',
totalCount: 0,
truncated: false
}
}
})
const client = { sendRequest } as unknown as RpcClient
const fixture = createMobileWebBridgeRoundtripFixture({
grants: MOBILE_WEB_PRODUCTION_GRANTS,
rpcClient: client
})
const statusPayload = async () => {
const snapshot = await fixture.client.workspaceSnapshot({ limit: 10 })
return { workspaceId: snapshot.workspaces[0]!.id, limit: 10 }
}
const payload = await statusPayload()
const catalogReads = () =>
sendRequest.mock.calls.filter(([method]) => method === 'mobileWeb.host.catalog').length
await Promise.all([
fixture.client.sourceControlStatus(payload),
fixture.client.sourceControlStatus(payload)
])
await fixture.client.sourceControlStatus(payload)
expect(catalogReads()).toBe(1)
// The swap retires every page handle too, so the page rediscovers its workspace first.
fixture.broker.replaceClient(client)
await fixture.client.sourceControlStatus(await statusPayload())
expect(catalogReads()).toBe(2)
})
})
@@ -0,0 +1,132 @@
import {
MobileWebHostCatalogPayloadSchema,
MobileWebHostCatalogResultSchema,
mobileWebHostPayloadWithinBounds,
type MobileWebHostGrant
} from '../../../src/shared/mobile-web/host-rpc-contract'
import type { RpcClient, SendRequestOptions } from '../transport/rpc-client'
import { MobileWebBrokerError, mobileWebBrokerHostRpcError } from './mobile-web-broker-error'
import { MOBILE_WEB_HOST_REQUEST_TIMEOUT_MS } from './mobile-web-host-requests'
/** A grant the desktop advertised, or `null` for a method it refused to advertise. */
type CatalogEntry = MobileWebHostGrant | null
/**
* The desktop's catalog is a module constant, so its grants can only change when the desktop
* process restarts, which tears the socket down and makes the broker replace this client. One
* catalog read per method per connection therefore answers every forwarded request.
*/
export class MobileWebHostCatalogCache {
private client: RpcClient | null = null
private readonly entries = new Map<string, CatalogEntry>()
private readonly inFlight = new Map<string, Promise<CatalogEntry>>()
clear(): void {
this.client = null
this.entries.clear()
this.inFlight.clear()
}
async read(
client: RpcClient,
input: unknown,
options?: SendRequestOptions
): Promise<{ grants: MobileWebHostGrant[] }> {
const payload = MobileWebHostCatalogPayloadSchema.parse(input)
const entries = await this.resolve(client, [...new Set(payload.methods)], options)
return { grants: entries.filter((entry) => entry !== null) }
}
async grant(
client: RpcClient,
method: string,
options?: SendRequestOptions
): Promise<CatalogEntry> {
return (await this.resolve(client, [method], options))[0] ?? null
}
private resolve(
client: RpcClient,
methods: readonly string[],
options: SendRequestOptions | undefined
): Promise<CatalogEntry[]> {
if (this.client !== client) {
this.clear()
this.client = client
}
const missing = methods.filter(
(method) => !this.entries.has(method) && !this.inFlight.has(method)
)
if (missing.length > 0) {
this.track(missing, this.send(client, missing, options))
}
const owned = new Set(missing)
return Promise.all(
methods.map((method) => this.settle(client, method, options, owned.has(method)))
)
}
private settle(
client: RpcClient,
method: string,
options: SendRequestOptions | undefined,
owned: boolean
): CatalogEntry | Promise<CatalogEntry> {
const cached = this.entries.get(method)
if (cached !== undefined) {
return cached
}
const shared = this.inFlight.get(method)!
if (owned) {
return shared
}
// A sibling request cancelling or timing out its own catalog read must not fail this one.
return shared.catch(async () => {
const settled = this.entries.get(method)
return settled !== undefined
? settled
: ((await this.send(client, [method], options)).get(method) ?? null)
})
}
private track(methods: readonly string[], request: Promise<Map<string, CatalogEntry>>): void {
for (const method of methods) {
const entry = request.then((found) => found.get(method) ?? null)
this.inFlight.set(method, entry)
const forget = (): void => {
if (this.inFlight.get(method) === entry) {
this.inFlight.delete(method)
}
}
void entry.then(forget, forget)
}
}
private async send(
client: RpcClient,
methods: readonly string[],
options: SendRequestOptions | undefined
): Promise<Map<string, CatalogEntry>> {
const response = await client.sendRequest(
'mobileWeb.host.catalog',
{ methods },
options ?? { timeoutMs: MOBILE_WEB_HOST_REQUEST_TIMEOUT_MS, budgetSpansConnect: true }
)
if (!response.ok) {
throw mobileWebBrokerHostRpcError(response.error)
}
if (!mobileWebHostPayloadWithinBounds(response.result)) {
throw new MobileWebBrokerError('too_large')
}
const { grants } = MobileWebHostCatalogResultSchema.parse(response.result)
const found = new Map(
methods.map((method) => [method, grants.find((grant) => grant.method === method) ?? null])
)
if (this.client === client) {
for (const [method, entry] of found) {
this.entries.set(method, entry)
}
}
return found
}
}
@@ -1,8 +1,13 @@
import { describe, expect, it, vi } from 'vitest'
import type { RpcClient } from '../transport/rpc-client'
import { MobileWebHostCatalogCache } from './mobile-web-host-catalog-cache'
import { executeMobileWebHostRequest } from './mobile-web-host-requests'
import { MobileWebWorkspaceAuthority } from './mobile-web-workspace-authority'
import { MOBILE_WEB_PRODUCTION_GRANTS } from './mobile-web-production-grants'
import { MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES } from '../../../src/shared/mobile-web/bridge-limits'
import {
MOBILE_WEB_PRODUCTION_GRANT_INDEX,
MOBILE_WEB_PRODUCTION_GRANTS
} from './mobile-web-production-grants'
import { createMobileWebBridgeRoundtripFixture } from './mobile-web-bridge-roundtrip-fixture'
const grant = {
@@ -18,6 +23,7 @@ function fixture() {
const sendRequest = vi.fn<RpcClient['sendRequest']>()
const args = {
authority,
catalog: new MobileWebHostCatalogCache(),
client: { sendRequest } as unknown as RpcClient,
isActive: () => true,
payload: {
@@ -121,17 +127,40 @@ describe('host-advertised unary forwarding', () => {
expect(sendRequest).toHaveBeenCalledTimes(1)
})
it('keeps native hard ceilings even when the trusted host permits a larger result', async () => {
it('refuses a grant advertising more than a shipped shell can deliver', async () => {
const { args, sendRequest } = fixture()
sendRequest.mockResolvedValueOnce({
ok: true,
result: { grants: [{ ...grant, maxResponseBytes: 10_000_000 }] }
})
await expect(executeMobileWebHostRequest(args)).rejects.toMatchObject({ name: 'ZodError' })
expect(sendRequest).toHaveBeenCalledOnce()
})
it('keeps native hard ceilings at the largest grant the desktop may advertise', async () => {
const { args, sendRequest } = fixture()
sendRequest
.mockResolvedValueOnce({
ok: true,
result: { grants: [{ ...grant, maxResponseBytes: 10_000_000 }] }
result: {
grants: [{ ...grant, maxResponseBytes: MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES }]
}
})
.mockResolvedValueOnce({ ok: true, result: { text: 'x'.repeat(640 * 1024) } })
await expect(executeMobileWebHostRequest(args)).rejects.toMatchObject({ code: 'too_large' })
})
it('enforces the advertised response ceiling below the envelope', async () => {
const { args, sendRequest } = fixture()
sendRequest
.mockResolvedValueOnce({
ok: true,
result: { grants: [{ ...grant, maxResponseBytes: 1024 }] }
})
.mockResolvedValueOnce({ ok: true, result: { text: 'x'.repeat(4096) } })
await expect(executeMobileWebHostRequest(args)).rejects.toMatchObject({ code: 'too_large' })
})
it('enforces host request bounds before executing', async () => {
const { args, sendRequest } = fixture()
sendRequest.mockResolvedValueOnce({
@@ -168,7 +197,9 @@ describe('host-advertised unary forwarding', () => {
})
const snapshot = await client.workspaceSnapshot({ limit: 10 })
const payload = { workspaceId: snapshot.workspaces[0]!.id, limit: 10 }
for (let i = 0; i < 4; i++) {
const ceiling =
MOBILE_WEB_PRODUCTION_GRANT_INDEX.get('workspace.hostRequest')!.limits.maxConcurrent
for (let i = 0; i < ceiling; i++) {
const controller = new AbortController()
const pending = client.sourceControlStatus(payload, { signal: controller.signal })
const rejection = expect(pending).rejects.toMatchObject({ code: 'cancelled' })
@@ -178,7 +209,7 @@ describe('host-advertised unary forwarding', () => {
await expect(client.sourceControlStatus(payload)).rejects.toMatchObject({
code: 'rate_limited'
})
expect(finishCatalog).toHaveLength(4)
expect(finishCatalog).toHaveLength(1)
finishCatalog.forEach((finish) => finish())
})
@@ -1,34 +1,27 @@
import {
MobileWebHostCatalogPayloadSchema,
MobileWebHostCatalogResultSchema,
MobileWebHostRequestPayloadSchema,
mobileWebHostPayloadWithinBounds
mobileWebHostPayloadByteLength
} from '../../../src/shared/mobile-web/host-rpc-contract'
import type { RpcClient, SendRequestOptions } from '../transport/rpc-client'
import { MobileWebBrokerError, mobileWebBrokerHostRpcError } from './mobile-web-broker-error'
import { mobileWebEncodedByteLength } from './mobile-web-request-accounting'
import type { MobileWebWorkspaceAuthority } from './mobile-web-workspace-authority'
import type { MobileWebHostCatalogCache } from './mobile-web-host-catalog-cache'
import type {
MobileWebHostWorkspaceId,
MobileWebWorkspaceAuthority
} from './mobile-web-workspace-authority'
const HOST_REQUEST_TIMEOUT_MS = 15_000
export const MOBILE_WEB_HOST_REQUEST_TIMEOUT_MS = 15_000
export async function readMobileWebHostCatalog(
client: RpcClient,
input: unknown,
options: SendRequestOptions = { timeoutMs: HOST_REQUEST_TIMEOUT_MS, budgetSpansConnect: true }
) {
const payload = MobileWebHostCatalogPayloadSchema.parse(input)
const response = await client.sendRequest('mobileWeb.host.catalog', payload, options)
if (!response.ok) {
throw mobileWebBrokerHostRpcError(response.error)
}
if (!mobileWebHostPayloadWithinBounds(response.result)) {
throw new MobileWebBrokerError('too_large')
}
return MobileWebHostCatalogResultSchema.parse(response.result)
/** The page handle a request is scoped to, resolved once and re-checked at every dispatch. Absent
* only for host-scoped grants, which carry no workspace at all. */
export type MobileWebHostRequestScope = {
pageWorkspaceId: string
hostWorkspaceId: MobileWebHostWorkspaceId
}
export type MobileWebHostRequestArguments = {
client: RpcClient
catalog: MobileWebHostCatalogCache
authority: MobileWebWorkspaceAuthority
payload: unknown
isActive: () => boolean
@@ -36,27 +29,31 @@ export type MobileWebHostRequestArguments = {
requestOptions?: () => SendRequestOptions
}
export function assertMobileWebHostRequestScope(
authority: MobileWebWorkspaceAuthority,
scope: MobileWebHostRequestScope | undefined
): void {
if (scope) {
authority.assertHostWorkspaceBinding(scope.pageWorkspaceId, scope.hostWorkspaceId)
}
}
export async function prepareMobileWebHostRequest(
args: MobileWebHostRequestArguments,
mode: 'once' | 'subscription'
) {
const payload = MobileWebHostRequestPayloadSchema.parse(args.payload)
if (!mobileWebHostPayloadWithinBounds(payload.params)) {
throw new MobileWebBrokerError('too_large')
}
const hostWorkspaceId =
const scope =
payload.workspaceId === undefined
? undefined
: args.authority.hostWorkspaceId(payload.workspaceId)
const catalog = await readMobileWebHostCatalog(
args.client,
{ methods: [payload.method] },
args.requestOptions?.()
)
const grant = catalog.grants.find((entry) => entry.method === payload.method)
: {
pageWorkspaceId: payload.workspaceId,
hostWorkspaceId: args.authority.hostWorkspaceId(payload.workspaceId)
}
const grant = await args.catalog.grant(args.client, payload.method, args.requestOptions?.())
if (
!grant ||
(grant.scope === 'host') !== (payload.workspaceId === undefined) ||
(grant.scope === 'host') !== (scope === undefined) ||
(grant.mode ?? 'once') !== mode ||
(mode === 'subscription' && !grant.unsubscribeMethod)
) {
@@ -65,9 +62,7 @@ export async function prepareMobileWebHostRequest(
if (!args.isActive()) {
throw new MobileWebBrokerError('cancelled')
}
if (payload.workspaceId !== undefined && hostWorkspaceId !== undefined) {
args.authority.assertHostWorkspaceBinding(payload.workspaceId, hostWorkspaceId)
}
assertMobileWebHostRequestScope(args.authority, scope)
if (
grant.pageSessionParam &&
(!args.getPageSessionId || grant.pageSessionParam === grant.workspaceParam)
@@ -76,24 +71,21 @@ export async function prepareMobileWebHostRequest(
}
const params = {
...payload.params,
...(grant.workspaceParam && hostWorkspaceId
? { [grant.workspaceParam]: `id:${hostWorkspaceId}` }
: {}),
// Scope agreement above plus the grant schema's refine make workspaceParam present here.
...(scope ? { [grant.workspaceParam!]: `id:${scope.hostWorkspaceId}` } : {}),
...(grant.pageSessionParam ? { [grant.pageSessionParam]: await args.getPageSessionId!() } : {})
}
if (
!mobileWebHostPayloadWithinBounds(params) ||
mobileWebEncodedByteLength(params) > grant.maxRequestBytes
) {
const requestBytes = mobileWebHostPayloadByteLength(params)
if (requestBytes === undefined || requestBytes > grant.maxRequestBytes) {
throw new MobileWebBrokerError('too_large')
}
return { payload, hostWorkspaceId, grant, params }
return { payload, scope, grant, params }
}
export async function executeMobileWebHostRequest(
args: MobileWebHostRequestArguments
): Promise<unknown> {
const deadline = Date.now() + HOST_REQUEST_TIMEOUT_MS
const deadline = Date.now() + MOBILE_WEB_HOST_REQUEST_TIMEOUT_MS
const beforeSend = () => {
if (!args.isActive()) {
throw new MobileWebBrokerError('cancelled')
@@ -106,16 +98,14 @@ export async function executeMobileWebHostRequest(
beforeSend()
return { timeoutMs: deadline - Date.now(), budgetSpansConnect: true, beforeSend }
}
const { payload, hostWorkspaceId, grant, params } = await prepareMobileWebHostRequest(
const { payload, scope, grant, params } = await prepareMobileWebHostRequest(
{ ...args, requestOptions },
'once'
)
const options = requestOptions()
options.beforeSend = () => {
beforeSend()
if (payload.workspaceId !== undefined && hostWorkspaceId !== undefined) {
args.authority.assertHostWorkspaceBinding(payload.workspaceId, hostWorkspaceId)
}
assertMobileWebHostRequestScope(args.authority, scope)
}
const response = await args.client.sendRequest(payload.method, params, options)
if (!response.ok) {
@@ -124,13 +114,9 @@ export async function executeMobileWebHostRequest(
if (!args.isActive()) {
throw new MobileWebBrokerError('cancelled')
}
if (payload.workspaceId !== undefined && hostWorkspaceId !== undefined) {
args.authority.assertHostWorkspaceBinding(payload.workspaceId, hostWorkspaceId)
}
if (
!mobileWebHostPayloadWithinBounds(response.result) ||
mobileWebEncodedByteLength(response.result) > grant.maxResponseBytes
) {
assertMobileWebHostRequestScope(args.authority, scope)
const responseBytes = mobileWebHostPayloadByteLength(response.result)
if (responseBytes === undefined || responseBytes > grant.maxResponseBytes) {
throw new MobileWebBrokerError('too_large')
}
return response.result
@@ -1,5 +1,6 @@
import { describe, expect, it, vi } from 'vitest'
import type { RpcClient } from '../transport/rpc-client'
import { MobileWebHostCatalogCache } from './mobile-web-host-catalog-cache'
import { executeMobileWebHostRequest } from './mobile-web-host-requests'
import { MobileWebWorkspaceAuthority } from './mobile-web-workspace-authority'
import { MOBILE_WEB_PRODUCTION_GRANTS } from './mobile-web-production-grants'
@@ -21,6 +22,7 @@ function fixture() {
sendRequest,
args: {
authority: new MobileWebWorkspaceAuthority((length) => new Uint8Array(length)),
catalog: new MobileWebHostCatalogCache(),
client: { sendRequest } as unknown as RpcClient,
isActive: () => true,
payload: { method: grant.method, params: { enabled: false } }
@@ -1,5 +1,6 @@
import { describe, expect, it, vi } from 'vitest'
import type { RpcClient } from '../transport/rpc-client'
import { MobileWebHostCatalogCache } from './mobile-web-host-catalog-cache'
import { MobileWebHostSubscriptions } from './mobile-web-host-subscriptions'
import { MobileWebWorkspaceAuthority } from './mobile-web-workspace-authority'
@@ -30,6 +31,7 @@ function fixture() {
postClosed
})
const args = {
catalog: new MobileWebHostCatalogCache(),
requestId: 'request',
subscriptionId: 'stream',
isActive: () => true,
@@ -1,23 +1,20 @@
import { mobileWebHostPayloadWithinBounds } from '../../../src/shared/mobile-web/host-rpc-contract'
import { mobileWebHostPayloadByteLength } from '../../../src/shared/mobile-web/host-rpc-contract'
import { MobileWebBrokerError } from './mobile-web-broker-error'
import {
assertMobileWebHostRequestScope,
prepareMobileWebHostRequest,
type MobileWebHostRequestArguments
type MobileWebHostRequestArguments,
type MobileWebHostRequestScope
} from './mobile-web-host-requests'
import { mobileWebEncodedByteLength } from './mobile-web-request-accounting'
import {
MobileWebSubscriptionLedger,
type MobileWebSubscriptionLedgerConfig,
type MobileWebSubscriptionRecord
} from './mobile-web-subscription-ledger'
import type {
MobileWebHostWorkspaceId,
MobileWebWorkspaceAuthority
} from './mobile-web-workspace-authority'
import type { MobileWebWorkspaceAuthority } from './mobile-web-workspace-authority'
type HostStreamRecord = MobileWebSubscriptionRecord & {
pageWorkspaceId: string | undefined
hostWorkspaceId: MobileWebHostWorkspaceId | undefined
scope: MobileWebHostRequestScope | undefined
maxEventBytes: number
closing: boolean
}
@@ -41,7 +38,7 @@ export class MobileWebHostSubscriptions extends MobileWebSubscriptionLedger<
}
): Promise<void> {
this.admit(args.subscriptionId)
const { payload, hostWorkspaceId, grant, params } = await prepareMobileWebHostRequest(
const { payload, scope, grant, params } = await prepareMobileWebHostRequest(
{
...args,
authority: this.config.workspaceAuthority
@@ -53,8 +50,7 @@ export class MobileWebHostSubscriptions extends MobileWebSubscriptionLedger<
}
const record: HostStreamRecord = {
...this.newRecord(args.requestId),
pageWorkspaceId: payload.workspaceId,
hostWorkspaceId,
scope,
maxEventBytes: grant.maxResponseBytes,
closing: false
}
@@ -70,12 +66,7 @@ export class MobileWebHostSubscriptions extends MobileWebSubscriptionLedger<
protected override canDeliver(subscriptionId: string, record: HostStreamRecord): boolean {
try {
if (record.pageWorkspaceId !== undefined && record.hostWorkspaceId !== undefined) {
this.config.workspaceAuthority.assertHostWorkspaceBinding(
record.pageWorkspaceId,
record.hostWorkspaceId
)
}
assertMobileWebHostRequestScope(this.config.workspaceAuthority, record.scope)
return true
} catch {
this.cancel(subscriptionId, { code: 'not_found', retryable: false })
@@ -91,10 +82,8 @@ export class MobileWebHostSubscriptions extends MobileWebSubscriptionLedger<
) {
return
}
if (
!mobileWebHostPayloadWithinBounds(event) ||
mobileWebEncodedByteLength(event) > record.maxEventBytes
) {
const eventBytes = mobileWebHostPayloadByteLength(event)
if (eventBytes === undefined || eventBytes > record.maxEventBytes) {
this.cancel(subscriptionId, { code: 'too_large', retryable: false })
return
}
@@ -8,19 +8,21 @@ import {
const SHELL_DIR = resolve(__dirname)
const REAUTHORIZATION =
/\.(?:assertHostWorkspaceBinding|assertHostRepoBinding|assertHostedTarget)\(/g
/(?:\.(?:assertHostWorkspaceBinding|assertHostRepoBinding|assertHostedTarget)|\bassertMobileWebHostRequestScope)\(/g
const HANDLE_RESOLUTION =
/\.(?:hostWorkspaceId|hostRepoId|hostConnectionId|resolveGitHub|resolveGitLab|resolveLinear)\(/
/** Every module that reauthorizes an opaque handle, and how many times. Pinned so deleting a
* reauthorization arm fails here even when the surrounding module keeps others. This counts
* sites; it does not prove each one sits after the awaited read it guards. */
* reauthorization arm fails here even when the surrounding module keeps others. Host forwarding
* reauthorizes through `assertMobileWebHostRequestScope`, so its declaration, its own assert and
* every call site count. This counts sites; it does not prove each one sits after the awaited
* read it guards. */
const REAUTHORIZATION_SITES: Record<string, number> = {
'mobile-web-agent-history-resume.ts': 1,
'mobile-web-browser-resource-binding.ts': 1,
'mobile-web-file-operations.ts': 1,
'mobile-web-file-write.ts': 1,
'mobile-web-host-requests.ts': 3,
'mobile-web-host-requests.ts': 5,
'mobile-web-host-subscriptions.ts': 1,
'mobile-web-markdown-operations.ts': 2,
'mobile-web-native-chat-binding.ts': 2,
@@ -17,7 +17,7 @@ export const MOBILE_WEB_PRODUCTION_GRANTS = [
...capabilityGrants('workspace', {
hostSubscribe: grantLimits(600 * 1024, 1024, 8, 8, 2),
hostCatalog: grantLimits(8 * 1024, 32 * 1024, 2, 8, 2),
hostRequest: grantLimits(600 * 1024, 600 * 1024, 4, 12, 4),
hostRequest: grantLimits(600 * 1024, 600 * 1024, 16, 32, 4),
snapshot: grantLimits(1 * 1024, 128 * 1024, 2, 4, 1),
repositories: grantLimits(256, 128 * 1024, 2, 4, 1),
subscribe: grantLimits(256, 1 * 1024, 1, 4, 1),
@@ -116,7 +116,7 @@ export function mobileWebRequestAtCapacity(args: {
}): boolean {
const key = mobileWebOperationKey(args.request)
return (
(args.isHostRequest && args.hostRequestsInFlight >= 4) ||
(args.isHostRequest && args.hostRequestsInFlight >= args.maxConcurrent) ||
args.pending.size >= MOBILE_WEB_BRIDGE_MAX_PENDING_REQUESTS ||
(args.request.mode === 'subscription' &&
mobileWebSubscriptionCount(args.pending.values(), args.ledgers) >=
@@ -127,12 +127,11 @@ export function mobileWebRequestAtCapacity(args: {
)
}
// Only one-shot forwards hold host work past a page cancellation; subscriptions are capped by
// their own ledger and catalog reads are served from a per-connection cache.
export function mobileWebIsHostRequest(request: {
capability: string
operation: string
}): boolean {
return (
request.capability === 'workspace' &&
['hostRequest', 'hostCatalog', 'hostSubscribe'].includes(request.operation)
)
return request.capability === 'workspace' && request.operation === 'hostRequest'
}
@@ -1,6 +1,7 @@
import { describe, expect, it } from 'vitest'
import { MOBILE_WEB_BRIDGE_MAX_SUBSCRIPTIONS } from '../../../src/shared/mobile-web/bridge-limits'
import { mobileWebRequestAtCapacity } from './mobile-web-request-accounting'
import { MOBILE_WEB_PRODUCTION_GRANT_INDEX } from './mobile-web-production-grants'
import { mobileWebIsHostRequest, mobileWebRequestAtCapacity } from './mobile-web-request-accounting'
describe('aggregate subscription admission', () => {
it('counts pending generic streams alongside active legacy and terminal streams', () => {
@@ -24,8 +25,8 @@ describe('aggregate subscription admission', () => {
: 0
}
],
isHostRequest: true,
hostRequestsInFlight: 1,
isHostRequest: false,
hostRequestsInFlight: 0,
maxConcurrent: 8
}
expect(mobileWebRequestAtCapacity(args)).toBe(true)
@@ -34,15 +35,37 @@ describe('aggregate subscription admission', () => {
})
it('retains the actual host work ceiling after page cancellation removes pending state', () => {
expect(
const args = {
pending: new Map(),
request: { mode: 'once' as const, capability: 'workspace', operation: 'hostRequest' },
ledgers: [],
isHostRequest: true,
maxConcurrent: 8
}
expect(mobileWebRequestAtCapacity({ ...args, hostRequestsInFlight: 8 })).toBe(true)
expect(mobileWebRequestAtCapacity({ ...args, hostRequestsInFlight: 7 })).toBe(false)
})
it('takes the host-request ceiling from the granted budget, not a private literal', () => {
const grant = MOBILE_WEB_PRODUCTION_GRANT_INDEX.get('workspace.hostRequest')
expect(grant).toBeDefined()
const atCapacity = (hostRequestsInFlight: number) =>
mobileWebRequestAtCapacity({
pending: new Map(),
request: { mode: 'subscription', capability: 'workspace', operation: 'hostSubscribe' },
request: { mode: 'once', capability: 'workspace', operation: 'hostRequest' },
ledgers: [],
isHostRequest: true,
hostRequestsInFlight: 4,
maxConcurrent: 8
hostRequestsInFlight,
maxConcurrent: grant!.limits.maxConcurrent
})
).toBe(true)
expect(atCapacity(grant!.limits.maxConcurrent - 1)).toBe(false)
expect(atCapacity(grant!.limits.maxConcurrent)).toBe(true)
})
it('leaves catalog reads and host streams off the one-shot host ceiling', () => {
expect(mobileWebIsHostRequest({ capability: 'workspace', operation: 'hostRequest' })).toBe(true)
for (const operation of ['hostCatalog', 'hostSubscribe']) {
expect(mobileWebIsHostRequest({ capability: 'workspace', operation })).toBe(false)
}
})
})
@@ -1,7 +1,7 @@
import type { MobileWebBridgePageMessage } from '../../../src/shared/mobile-web/bridge-contract'
import type { MobileWebCapabilityExecutionDependencies } from './mobile-web-capability-execution-dependencies'
import { MobileWebBrokerError } from './mobile-web-broker-error'
import { executeMobileWebHostRequest, readMobileWebHostCatalog } from './mobile-web-host-requests'
import { executeMobileWebHostRequest } from './mobile-web-host-requests'
import { executeMobileWebWorkspaceOperation } from './mobile-web-workspace-operations'
type OnceRequest = Extract<
@@ -14,11 +14,12 @@ export async function executeWorkspace(
request: OnceRequest
): Promise<unknown> {
if (request.operation === 'hostCatalog') {
return readMobileWebHostCatalog(args.connectedClient(), request.payload)
return args.hostCatalog.read(args.connectedClient(), request.payload)
}
if (request.operation === 'hostRequest') {
return executeMobileWebHostRequest({
client: args.connectedClient(),
catalog: args.hostCatalog,
authority: args.workspaceAuthority,
getPageSessionId: args.getPageSessionId,
payload: request.payload,
@@ -415,6 +415,7 @@ describe('hosted mobile bridge over cloud Relay transport', () => {
source: 'transcript'
}
])
// One catalog read per method per connection: bind and read share the second one.
expect(observedMethods).toEqual([
'pairing.getEndpoints',
'runtime.clientCapabilities.update',
@@ -423,9 +424,7 @@ describe('hosted mobile bridge over cloud Relay transport', () => {
'mobileWeb.page.subscribe',
'mobileWeb.session.snapshot',
'mobileWeb.host.catalog',
'mobileWeb.host.catalog',
'mobileWeb.nativeChat.bind',
'mobileWeb.host.catalog',
'mobileWeb.nativeChat.read'
])
expect(JSON.stringify({ sessionSnapshot, transcript })).not.toContain('relay-provider-session')
@@ -1,4 +1,5 @@
import { describe, expect, it } from 'vitest'
import { MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES } from '../../../../shared/mobile-web/bridge-limits'
import { MOBILE_WEB_HOST_CATALOG_METHOD } from './mobile-web-host-catalog'
import type { RpcContext } from '../core'
import { ALL_RPC_METHODS } from './index'
@@ -35,7 +36,14 @@ describe('mobile web host catalog', () => {
method,
workspaceParam: 'worktree',
maxRequestBytes: 16 * 1024,
maxResponseBytes: 512 * 1024
// Directory listings, file reads and diffs get the whole bridge envelope.
maxResponseBytes: [
'mobileWeb.files.readDir',
'mobileWeb.files.read',
'mobileWeb.sourceControl.diff'
].includes(method)
? MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES
: 512 * 1024
}))
})
for (const method of [
@@ -1,9 +1,20 @@
import { MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES } from '../../../../shared/mobile-web/bridge-limits'
import {
MobileWebHostCatalogPayloadSchema,
type MobileWebHostGrant
} from '../../../../shared/mobile-web/host-rpc-contract'
import { defineMethod } from '../core'
// Reads whose result is a directory, a file, a diff or a transcript page, so the only honest
// ceiling is the bridge envelope the shell can actually deliver.
const READ_HEAVY_METHODS = new Set([
'mobileWeb.files.readDir',
'mobileWeb.files.read',
'mobileWeb.sourceControl.diff',
'mobileWeb.session.snapshot',
'mobileWeb.nativeChat.read'
])
// Only page-safe results belong here; transport credentials never enter this catalog.
const PAGE_METHODS = new Map<string, MobileWebHostGrant>(
[
@@ -40,8 +51,13 @@ const PAGE_METHODS = new Map<string, MobileWebHostGrant>(
method.startsWith('mobileWeb.session.')
? { pageSessionParam: 'pageSession' }
: {}),
maxRequestBytes: method === 'mobileWeb.nativeChat.mutate' ? 600 * 1024 : 16 * 1024,
maxResponseBytes: 512 * 1024
maxRequestBytes:
method === 'mobileWeb.nativeChat.mutate'
? MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES
: 16 * 1024,
maxResponseBytes: READ_HEAVY_METHODS.has(method)
? MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES
: 512 * 1024
}
])
)
@@ -1,7 +1,9 @@
import { describe, expect, it } from 'vitest'
import { MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES } from './bridge-limits'
import {
MobileWebHostRequestPayloadSchema,
MobileWebHostGrantSchema,
mobileWebHostPayloadByteLength,
mobileWebHostPayloadWithinBounds
} from './host-rpc-contract'
@@ -29,6 +31,32 @@ describe('generic host payload transport', () => {
.success
).toBe(false)
})
it('refuses a grant advertising more than a shipped shell can deliver', () => {
const grant = { method: 'future.read', scope: 'host' as const, maxRequestBytes: 1024 }
const envelope = MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES
expect(
MobileWebHostGrantSchema.safeParse({ ...grant, maxResponseBytes: envelope }).success
).toBe(true)
expect(
MobileWebHostGrantSchema.safeParse({ ...grant, maxResponseBytes: envelope + 1 }).success
).toBe(false)
expect(
MobileWebHostGrantSchema.safeParse({
...grant,
maxRequestBytes: envelope + 1,
maxResponseBytes: 1024
}).success
).toBe(false)
})
it('reports the encoded length once for callers that also need the verdict', () => {
expect(mobileWebHostPayloadByteLength({ future: 'value' })).toBe(
new TextEncoder().encode(JSON.stringify({ future: 'value' })).byteLength
)
expect(mobileWebHostPayloadByteLength('é'.repeat(310 * 1024))).toBeUndefined()
expect(mobileWebHostPayloadByteLength(() => {})).toBeUndefined()
})
it('bounds depth, node count and encoded bytes independently of domain shape', () => {
let nested: unknown = null
for (let i = 0; i < 34; i++) {
+13 -9
View File
@@ -37,8 +37,8 @@ export const MobileWebHostGrantSchema = z
.max(80)
.regex(/^[A-Za-z][A-Za-z0-9]*$/)
.optional(),
maxRequestBytes: z.number().int().positive(),
maxResponseBytes: z.number().int().positive()
maxRequestBytes: z.number().int().positive().max(MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES),
maxResponseBytes: z.number().int().positive().max(MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES)
})
.refine(
(grant) =>
@@ -57,29 +57,33 @@ export type MobileWebHostGrant = z.infer<typeof MobileWebHostGrantSchema>
export type MobileWebHostRequestPayload = z.infer<typeof MobileWebHostRequestPayloadSchema>
export function mobileWebHostPayloadWithinBounds(value: unknown): boolean {
return mobileWebHostPayloadByteLength(value) !== undefined
}
/** Encoded length, or undefined when the value cannot cross the bridge at all. Callers that need
* both the verdict and the size must use this, not a second serialization. */
export function mobileWebHostPayloadByteLength(value: unknown): number | undefined {
const pending = [{ value, depth: 0 }]
let nodes = 0
while (pending.length > 0) {
const entry = pending.pop()!
if (++nodes > 40_000 || entry.depth > 32) {
return false
return undefined
}
if (entry.value !== null && typeof entry.value === 'object') {
for (const child of Object.values(entry.value)) {
pending.push({ value: child, depth: entry.depth + 1 })
if (pending.length > 40_000) {
return false
return undefined
}
}
} else if (
entry.value !== null &&
!['string', 'boolean', 'number'].includes(typeof entry.value)
) {
return false
return undefined
}
}
return (
new TextEncoder().encode(JSON.stringify(value)).byteLength <=
MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES
)
const bytes = new TextEncoder().encode(JSON.stringify(value)).byteLength
return bytes <= MOBILE_WEB_BRIDGE_MAX_OPERATION_BYTES ? bytes : undefined
}