feat(mobile): forward host-advertised subscriptions

Add bounded generic streams with Desktop cleanup metadata, direct and relay cancellation, and page-owned source-control invalidation. Preserve legacy fallback and protocol 2. Required code gates and page export pass; platform journeys remain tracked.
This commit is contained in:
Jinwoo-H
2026-09-06 18:45:47 -04:00
parent 42c031d85e
commit a0eee2fc21
33 changed files with 988 additions and 98 deletions
@@ -68,12 +68,12 @@ folder and SSH workspaces still use their actual execution owner.
## 3. Generic subscriptions and transport lifecycle
- [ ] Remove the transport's static method-to-unsubscribe dependency for generic
- [x] Remove the transport's static method-to-unsubscribe dependency for generic
streams using host-advertised cleanup or a generic host subscription token.
- [ ] Preserve direct/relay setup, ready, unsubscribe and reconnect behavior.
- [ ] Reuse the existing subscription ledger with bounded pending event bytes
- [x] Reuse the existing subscription ledger with bounded pending event bytes
and event count; enforce aggregate subscription ceilings.
- [ ] Forward domain event shapes without APK-owned projections.
- [x] Forward domain event shapes without APK-owned projections (source-control file watch).
- [ ] Migrate native-chat/session/source-control/account feeds.
- [ ] Preserve terminal binary capability negotiation, acknowledgements,
backpressure and resync; never silently substitute JSON stream semantics.
@@ -284,3 +284,20 @@ A failed/malformed `session.tabs.list` no longer revokes native-chat authority;
it reports a retryable host error. Successful snapshots still establish removal.
Next implementation slice: host-advertised generic subscriptions, cleanup
metadata, bounded queues, and page-side source-control invalidation.
### Generic subscription implementation
Implemented `workspace.hostSubscribe` using Desktop catalog mode/cleanup metadata,
opaque workspace binding, and unchanged protocol 2. Source-control invalidation
now consumes Desktop file-watch events in the page, with legacy shell/host fallback.
Direct and relay transports retain arbitrary cleanup routes across early cancellation;
direct reconnect clears old tokens. Generic transport records cap at 128, including
cancelled records awaiting ready. Ledger queues cap at 64 events / 2 MiB; aggregate
admission counts legacy, terminal and generic subscriptions together.
All required gates pass. Mobile: 834 files / 5,515 passed; root: 318 files / 2,700
passed. Additional capacity tests: 3 files / 14 passed. Deliberate dispatch census
228 → 229; reauthorization census includes generic stream delivery and extracted
workspace dispatch. Logs: `/tmp/orca-ota-e2e/generic-stream-gates/`.
Page build: `925397fb68c40b511854af454cce2a2989134aa34a517f10fd494064f36ea895`.
Native-chat/session/account feed migration and platform E2E remain open.
@@ -1,8 +1,7 @@
import { requireMobileWebConnectedClient } from './mobile-web-connected-client'
import {
MOBILE_WEB_BRIDGE_MAX_PENDING_REQUESTS,
type MobileWebBridgePageMessage,
type MobileWebResumeRoute
import type {
MobileWebBridgePageMessage,
MobileWebResumeRoute
} from '../../../src/shared/mobile-web/bridge-contract'
import type { RpcClient } from '../transport/rpc-client'
import {
@@ -26,10 +25,10 @@ import { rememberMobileWebBrokerRoute } from './mobile-web-broker-route-memory'
import { resolveMobileWebHostNavigationRoute } from './mobile-web-host-navigation-route'
import {
mobileWebEncodedByteLength,
mobileWebRequestAtCapacity,
mobileWebRequestSurvivesCancellation,
mobileWebAgentHistoryContinuation,
mobileWebOperationKey,
mobileWebPendingForOperation,
mobileWebPendingRequestForSubscription,
mobileWebRequestExpectsSubscription,
mobileWebWorkspaceSnapshotContinuation
@@ -155,7 +154,7 @@ export class MobileWebCapabilityBroker {
const isHostRequest =
request.capability === 'workspace' &&
(request.operation === 'hostRequest' || request.operation === 'hostCatalog')
['hostRequest', 'hostCatalog', 'hostSubscribe'].includes(request.operation)
const grant = MOBILE_WEB_PRODUCTION_GRANT_INDEX.get(mobileWebOperationKey(request))
const expectsSubscription = mobileWebRequestExpectsSubscription(request)
if (!grant || (request.mode === 'subscription') !== expectsSubscription) {
@@ -174,13 +173,14 @@ export class MobileWebCapabilityBroker {
return
}
if (
(isHostRequest && this.hostRequestsInFlight >= 4) ||
this.pending.size >= MOBILE_WEB_BRIDGE_MAX_PENDING_REQUESTS ||
mobileWebPendingForOperation(this.pending.values(), mobileWebOperationKey(request)) +
this.subscriptions.countForOperation(mobileWebOperationKey(request)) +
this.terminalStreams.countForOperation(mobileWebOperationKey(request)) +
this.speechAuthority.countForOperation(mobileWebOperationKey(request)) >=
grant.limits.maxConcurrent
mobileWebRequestAtCapacity({
pending: this.pending,
request,
isHostRequest,
hostRequestsInFlight: this.hostRequestsInFlight,
ledgers: [this.subscriptions, this.terminalStreams, this.speechAuthority],
maxConcurrent: grant.limits.maxConcurrent
})
) {
await this.messages.error(request.requestId, 'rate_limited', true)
return
@@ -254,6 +254,7 @@ export class MobileWebCapabilityBroker {
sourceControlBranchCompare: this.authorities.sourceControlBranchCompare,
speechAuthority: this.speechAuthority,
workspaceSubscriptions: this.subscriptions.workspace,
hostSubscriptions: this.subscriptions.host,
terminalStreams: this.terminalStreams,
commitMessageGeneration: this.commitMessageGeneration,
browserAuthority: this.authorities.browser,
@@ -49,7 +49,7 @@ describe('mobile web capability dispatch census', () => {
})
expect(unresolved.map(({ capability, operation }) => `${capability}.${operation}`)).toEqual([])
expect(registeredOperations()).toHaveLength(228)
expect(registeredOperations()).toHaveLength(229)
})
it('carries a dispatch arm for exactly the capabilities that own operations of that mode', () => {
@@ -6,7 +6,7 @@ import type { MobileWebBridgePageMessage } from '../../../src/shared/mobile-web/
import type { MobileWebBridgeCapability } from '../../../src/shared/mobile-web/bridge-operation-registry'
import { MobileWebSourceControlSubscribePayloadSchema } from '../../../src/shared/mobile-web/source-control-operation-contract'
import { MobileWebSpeechSubscribePayloadSchema } from '../../../src/shared/mobile-web/speech-operation-contract'
import { executeMobileWebHostRequest, readMobileWebHostCatalog } from './mobile-web-host-requests'
import { executeWorkspace } from './mobile-web-workspace-capability'
import { executeMobileWebAccountCapability } from './mobile-web-account-capability'
import { executeMobileWebAgentHistoryOperation } from './mobile-web-agent-history-operations'
import { MobileWebBrokerError } from './mobile-web-broker-error'
@@ -23,7 +23,6 @@ import { executeMobileWebSessionOperation } from './mobile-web-session-operation
import { executeMobileWebSourceControlOperation } from './mobile-web-source-control-operations'
import { executeMobileWebSpeechOperation } from './mobile-web-speech-operations'
import { executeMobileWebTaskReadOperation } from './mobile-web-task-read-operations'
import { executeMobileWebWorkspaceOperation } from './mobile-web-workspace-operations'
type PageRequest = Extract<MobileWebBridgePageMessage, { type: 'request' }>
type OnceRequest = Extract<PageRequest, { mode: 'once' }>
@@ -70,35 +69,6 @@ async function executeBrowser(args: Deps, request: OnceRequest): Promise<unknown
})
}
async function executeWorkspace(args: Deps, request: OnceRequest): Promise<unknown> {
if (request.operation === 'hostCatalog') {
return readMobileWebHostCatalog(args.connectedClient(), request.payload)
}
if (request.operation === 'hostRequest') {
return executeMobileWebHostRequest({
client: args.connectedClient(),
authority: args.workspaceAuthority,
payload: request.payload,
isActive: args.isRequestActive
})
}
if (request.capability !== 'workspace' && request.capability !== 'settings') {
throw new MobileWebBrokerError('unsupported_capability')
}
const result = await executeMobileWebWorkspaceOperation({
capability: request.capability,
operation: request.operation,
payload: request.payload,
client: args.connectedClient(),
authority: args.workspaceAuthority,
snapshots: args.workspaceSnapshots
})
if (request.capability === 'workspace' && request.operation === 'activate') {
args.terminalArtifactAuthority.clear()
}
return result
}
async function executeSession(args: Deps, request: OnceRequest): Promise<unknown> {
const result = await executeMobileWebSessionOperation({
operation: request.operation,
@@ -242,6 +212,16 @@ 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({
requestId: request.requestId,
subscriptionId: request.subscriptionId,
payload: request.payload,
client: args.connectedClient(),
isActive: args.isRequestActive
})
return null
}
requireSubscribeOperation(request)
MobileWebWorkspaceSubscribePayloadSchema.parse(request.payload)
args.workspaceSubscriptions.start({
@@ -1,3 +1,4 @@
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'
import type { MobileWebAccountSubscriptions } from './mobile-web-account-subscriptions'
@@ -34,6 +35,7 @@ export type MobileWebCapabilityExecutionDependencies = {
agentHistoryAuthority: MobileWebAgentHistoryAuthority
agentHistoryPager: MobileWebAgentHistoryPager
agentHistoryResume: MobileWebAgentHistoryResume
hostSubscriptions: MobileWebHostSubscriptions
accountSubscriptions: MobileWebAccountSubscriptions
browserStreams: MobileWebBrowserStreams
nativeChatSubscriptions: MobileWebNativeChatSubscriptions
@@ -1,3 +1,4 @@
import { mobileWebRequestExpectsSubscription } from './mobile-web-request-accounting'
import { describe, expect, it, vi } from 'vitest'
import type { MobileWebBridgePageMessage } from '../../../src/shared/mobile-web/bridge-contract'
import type { RpcClient } from '../transport/rpc-client'
@@ -94,6 +95,8 @@ function requestFor(
capability: grant.capability,
operation: grant.operation,
payload,
...(grant.operation === 'subscribe' ? { subscriptionId: `s${requestId.slice(1)}` } : {})
...(mobileWebRequestExpectsSubscription(grant)
? { subscriptionId: `s${requestId.slice(1)}` }
: {})
})
}
@@ -1,3 +1,4 @@
import { MobileWebHostSubscriptions } from './mobile-web-host-subscriptions'
import { MobileWebAccountSubscriptions } from './mobile-web-account-subscriptions'
import type { MobileWebSubscriptionClosure } from './mobile-web-subscription-closure'
import type {
@@ -14,6 +15,7 @@ import type { MobileWebWorkspaceAuthority } from './mobile-web-workspace-authori
import { MobileWebWorkspaceSubscriptions } from './mobile-web-workspace-subscriptions'
export class MobileWebCapabilitySubscriptions {
readonly host: MobileWebHostSubscriptions
readonly account: MobileWebAccountSubscriptions
readonly browser: MobileWebBrowserStreams
readonly nativeChat: MobileWebNativeChatSubscriptions
@@ -34,6 +36,10 @@ export class MobileWebCapabilitySubscriptions {
postEvent: args.postEvent,
postClosed: args.postClosed
}
this.host = new MobileWebHostSubscriptions({
...shared,
workspaceAuthority: args.workspaceAuthority
})
this.account = new MobileWebAccountSubscriptions(shared)
this.browser = new MobileWebBrowserStreams({
...shared,
@@ -56,6 +62,7 @@ export class MobileWebCapabilitySubscriptions {
})
this.workspace = new MobileWebWorkspaceSubscriptions(shared)
this.ledgers = [
this.host,
this.account,
this.browser,
this.nativeChat,
@@ -21,12 +21,17 @@ export async function readMobileWebHostCatalog(client: RpcClient, input: unknown
return MobileWebHostCatalogResultSchema.parse(response.result)
}
export async function executeMobileWebHostRequest(args: {
export type MobileWebHostRequestArguments = {
client: RpcClient
authority: MobileWebWorkspaceAuthority
payload: unknown
isActive: () => boolean
}): Promise<unknown> {
}
export async function prepareMobileWebHostRequest(
args: MobileWebHostRequestArguments,
mode: 'once' | 'subscription'
) {
const payload = MobileWebHostRequestPayloadSchema.parse(args.payload)
if (!mobileWebHostPayloadWithinBounds(payload.params)) {
throw new MobileWebBrokerError('too_large')
@@ -34,7 +39,11 @@ export async function executeMobileWebHostRequest(args: {
const hostWorkspaceId = args.authority.hostWorkspaceId(payload.workspaceId)
const catalog = await readMobileWebHostCatalog(args.client, { methods: [payload.method] })
const grant = catalog.grants.find((entry) => entry.method === payload.method)
if (!grant) {
if (
!grant ||
(grant.mode ?? 'once') !== mode ||
(mode === 'subscription' && !grant.unsubscribeMethod)
) {
throw new MobileWebBrokerError('unsupported_capability')
}
if (!args.isActive()) {
@@ -48,6 +57,16 @@ export async function executeMobileWebHostRequest(args: {
) {
throw new MobileWebBrokerError('too_large')
}
return { payload, hostWorkspaceId, grant, params }
}
export async function executeMobileWebHostRequest(
args: MobileWebHostRequestArguments
): Promise<unknown> {
const { payload, hostWorkspaceId, grant, params } = await prepareMobileWebHostRequest(
args,
'once'
)
const response = await args.client.sendRequest(payload.method, params)
if (!response.ok) {
throw mobileWebBrokerHostRpcError(response.error)
@@ -0,0 +1,112 @@
import { describe, expect, it, vi } from 'vitest'
import type { RpcClient } from '../transport/rpc-client'
import { createMobileWebBridgeRoundtripFixture } from './mobile-web-bridge-roundtrip-fixture'
import { MOBILE_WEB_PRODUCTION_GRANTS } from './mobile-web-production-grants'
function fixture(catalogAvailable = true, genericShell = true) {
let emit: (event: unknown) => void = () => {}
const unsubscribe = vi.fn()
const subscribe = vi.fn<RpcClient['subscribe']>((_method, _params, listener) => {
emit = listener
return unsubscribe
})
const sendRequest = vi.fn<RpcClient['sendRequest']>().mockImplementation(async (method) => {
if (method === 'worktree.ps') {
return {
ok: true,
result: {
worktrees: [
{ worktreeId: 'host-workspace', repo: '/private/repo', displayName: 'Workspace' }
]
}
}
}
return catalogAvailable
? {
ok: true,
result: {
grants: [
{
method: 'mobileWeb.files.watch',
mode: 'subscription',
workspaceParam: 'worktree',
unsubscribeMethod: 'files.unwatch',
maxRequestBytes: 1024,
maxResponseBytes: 512 * 1024
},
{
method: 'future.events',
mode: 'subscription',
workspaceParam: 'scope',
unsubscribeMethod: 'future.release',
maxRequestBytes: 1024,
maxResponseBytes: 512 * 1024
}
]
}
}
: { ok: false, error: { code: 'method_not_found', message: 'Old host' } }
})
const bridge = createMobileWebBridgeRoundtripFixture({
grants: MOBILE_WEB_PRODUCTION_GRANTS.filter(
(grant) => genericShell || grant.operation !== 'hostSubscribe'
),
rpcClient: { sendRequest, subscribe } as unknown as RpcClient
})
return { ...bridge, subscribe, unsubscribe, emit: (event: unknown) => emit(event) }
}
describe('generic subscription bridge compatibility', () => {
it.each([
[true, true],
[false, true],
[true, false]
])('source-control catalog=%s shell=%s', async (catalog, shell) => {
const f = fixture(catalog, shell)
const workspace = (await f.client.workspaceSnapshot({ limit: 10 })).workspaces[0]!.id
const onEvent = vi.fn()
const onError = vi.fn()
const subscription = f.client.sourceControlSubscribe(
{ workspaceId: workspace },
onEvent,
onError
)
await subscription.ready
const generic = catalog && shell
expect(f.subscribe.mock.calls[0]?.[0]).toBe(generic ? 'mobileWeb.files.watch' : 'files.watch')
f.emit({
type: 'changed',
...(generic ? {} : { worktree: 'id:host-workspace' }),
events: [],
futureField: 'new'
})
await vi.waitFor(() =>
expect(onEvent).toHaveBeenCalledWith({ workspaceId: workspace, reason: 'changed' })
)
expect(onError).not.toHaveBeenCalled()
subscription.unsubscribe()
expect(f.unsubscribe).toHaveBeenCalledOnce()
})
it('forwards a future method and event without adding a domain operation', async () => {
const f = fixture()
const workspaceId = (await f.client.workspaceSnapshot({ limit: 10 })).workspaces[0]!.id
const onEvent = vi.fn()
const subscription = f.client.hostSubscribe(
{ method: 'future.events', workspaceId, params: { newParam: 42 } },
onEvent,
vi.fn()
)
await subscription.ready
const event = { futureKind: 'unknown-to-shell', fields: { addedLater: true } }
f.emit(event)
await vi.waitFor(() => expect(onEvent).toHaveBeenCalledWith(event))
expect(f.subscribe).toHaveBeenCalledWith(
'future.events',
{ scope: 'id:host-workspace', newParam: 42 },
expect.any(Function),
{ serverUnsubscribeMethod: 'future.release' }
)
subscription.unsubscribe()
})
})
@@ -0,0 +1,148 @@
import { describe, expect, it, vi } from 'vitest'
import type { RpcClient } from '../transport/rpc-client'
import { MobileWebHostSubscriptions } from './mobile-web-host-subscriptions'
import { MobileWebWorkspaceAuthority } from './mobile-web-workspace-authority'
const grant = {
method: 'future.feed',
mode: 'subscription',
workspaceParam: 'worktree',
unsubscribeMethod: 'future.stop',
maxRequestBytes: 1024,
maxResponseBytes: 512 * 1024
}
function fixture() {
const authority = new MobileWebWorkspaceAuthority((length) => new Uint8Array(length))
authority.synchronize([{ workspaceId: 'host-workspace', repoId: 'host-repo' }])
const postEvent = vi.fn().mockResolvedValue(undefined)
const postClosed = vi.fn()
const unsubscribe = vi.fn()
let emit: (event: unknown) => void = () => {}
const sendRequest = vi.fn().mockResolvedValue({ ok: true, result: { grants: [grant] } })
const subscribe = vi.fn<RpcClient['subscribe']>((_method, _params, listener) => {
emit = listener
return unsubscribe
})
const ledger = new MobileWebHostSubscriptions({
workspaceAuthority: authority,
isActive: () => true,
postEvent,
postClosed
})
const args = {
requestId: 'request',
subscriptionId: 'stream',
isActive: () => true,
client: { sendRequest, subscribe } as unknown as RpcClient,
payload: {
method: grant.method,
workspaceId: authority.pageWorkspaceId('host-workspace'),
params: {}
}
}
return {
ledger,
args,
authority,
sendRequest,
subscribe,
unsubscribe,
postEvent,
postClosed,
emit: (event: unknown) => emit(event)
}
}
describe('generic host subscriptions', () => {
it('forwards future events and uses Desktop cleanup metadata', async () => {
const f = fixture()
await f.ledger.start(f.args)
expect(f.subscribe).toHaveBeenCalledWith(
'future.feed',
{ worktree: 'id:host-workspace' },
expect.any(Function),
{ serverUnsubscribeMethod: 'future.stop' }
)
const event = { type: 'future-shape', nested: { additional: true } }
f.emit(event)
await vi.waitFor(() => expect(f.postEvent).toHaveBeenCalledWith('stream', 0, event))
f.ledger.cancel('stream')
expect(f.unsubscribe).toHaveBeenCalledOnce()
})
it('does not open after cancellation during catalog discovery', async () => {
const f = fixture()
let active = true
f.args.isActive = () => active
f.sendRequest.mockImplementationOnce(async () => {
active = false
return { ok: true, result: { grants: [grant] } }
})
await expect(f.ledger.start(f.args)).rejects.toMatchObject({ code: 'cancelled' })
expect(f.subscribe).not.toHaveBeenCalled()
})
it('retires a late host handle after a synchronous oversized event', async () => {
const f = fixture()
f.subscribe.mockImplementationOnce((_method, _params, emit) => {
emit({ payload: 'x'.repeat(600 * 1024) })
return f.unsubscribe
})
await f.ledger.start(f.args)
expect(f.unsubscribe).toHaveBeenCalledOnce()
expect(f.postClosed).toHaveBeenCalledWith('stream', { code: 'too_large', retryable: false })
})
it('bounds retained event bytes while the page is stalled', async () => {
const f = fixture()
let release!: () => void
f.postEvent.mockImplementationOnce(
() =>
new Promise<void>((resolve) => {
release = resolve
})
)
await f.ledger.start(f.args)
f.emit({ payload: 'x'.repeat(400 * 1024) })
await vi.waitFor(() => expect(f.postEvent).toHaveBeenCalledOnce())
for (let index = 0; index < 6; index++) {
f.emit({ payload: 'x'.repeat(400 * 1024) })
}
expect(f.unsubscribe).toHaveBeenCalledOnce()
expect(f.postClosed).toHaveBeenCalledWith('stream', { code: 'rate_limited', retryable: true })
release()
})
it('does not publish after workspace authority is retired', async () => {
const f = fixture()
await f.ledger.start(f.args)
f.authority.clear()
f.emit({ type: 'future' })
expect(f.postEvent).not.toHaveBeenCalled()
expect(f.unsubscribe).toHaveBeenCalledOnce()
})
it('bounds tiny queued events while the page is stalled', async () => {
const f = fixture()
const stalled = Promise.withResolvers<void>()
f.postEvent.mockReturnValueOnce(stalled.promise)
await f.ledger.start(f.args)
f.emit({ value: 0 })
await vi.waitFor(() => expect(f.postEvent).toHaveBeenCalledOnce())
for (let value = 1; value <= 64; value++) {
f.emit({ value })
}
expect(f.unsubscribe).toHaveBeenCalledOnce()
expect(f.postClosed).toHaveBeenCalledWith('stream', { code: 'rate_limited', retryable: true })
stalled.resolve()
})
it('rejects unary grants in the streaming lane', async () => {
const f = fixture()
f.sendRequest.mockResolvedValueOnce({
ok: true,
result: { grants: [{ ...grant, mode: 'once' }] }
})
await expect(f.ledger.start(f.args)).rejects.toMatchObject({ code: 'unsupported_capability' })
expect(f.subscribe).not.toHaveBeenCalled()
})
})
@@ -0,0 +1,104 @@
import { mobileWebHostPayloadWithinBounds } from '../../../src/shared/mobile-web/host-rpc-contract'
import { MobileWebBrokerError } from './mobile-web-broker-error'
import {
prepareMobileWebHostRequest,
type MobileWebHostRequestArguments
} 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'
type HostStreamRecord = MobileWebSubscriptionRecord & {
pageWorkspaceId: string
hostWorkspaceId: MobileWebHostWorkspaceId
maxEventBytes: number
closing: boolean
}
export class MobileWebHostSubscriptions extends MobileWebSubscriptionLedger<
unknown,
HostStreamRecord
> {
constructor(
private readonly config: MobileWebSubscriptionLedgerConfig<unknown> & {
workspaceAuthority: MobileWebWorkspaceAuthority
}
) {
super({ ...config, operationKey: 'workspace.hostSubscribe' })
}
async start(
args: Omit<MobileWebHostRequestArguments, 'authority'> & {
requestId: string
subscriptionId: string
}
): Promise<void> {
this.admit(args.subscriptionId)
const { payload, hostWorkspaceId, grant, params } = await prepareMobileWebHostRequest(
{
...args,
authority: this.config.workspaceAuthority
},
'subscription'
)
if (!args.isActive()) {
throw new MobileWebBrokerError('cancelled')
}
const record: HostStreamRecord = {
...this.newRecord(args.requestId),
pageWorkspaceId: payload.workspaceId,
hostWorkspaceId,
maxEventBytes: grant.maxResponseBytes,
closing: false
}
this.open(args.subscriptionId, record, () =>
args.client.subscribe(
payload.method,
params,
(event) => this.receive(args.subscriptionId, record, event),
{ serverUnsubscribeMethod: grant.unsubscribeMethod }
)
)
}
protected override canDeliver(subscriptionId: string, record: HostStreamRecord): boolean {
try {
this.config.workspaceAuthority.assertHostWorkspaceBinding(
record.pageWorkspaceId,
record.hostWorkspaceId
)
return true
} catch {
this.cancel(subscriptionId, { code: 'not_found', retryable: false })
return false
}
}
private receive(subscriptionId: string, record: HostStreamRecord, event: unknown): void {
if (
record.closing ||
!this.isCurrent(subscriptionId, record) ||
!this.canDeliver(subscriptionId, record)
) {
return
}
if (
!mobileWebHostPayloadWithinBounds(event) ||
mobileWebEncodedByteLength(event) > record.maxEventBytes
) {
this.cancel(subscriptionId, { code: 'too_large', retryable: false })
return
}
const type =
typeof event === 'object' && event !== null && 'type' in event ? event.type : undefined
record.closing = type === 'end' || type === 'error'
this.enqueue(subscriptionId, record, event, record.closing)
}
}
@@ -20,6 +20,7 @@ const REAUTHORIZATION_SITES: Record<string, number> = {
'mobile-web-file-operations.ts': 1,
'mobile-web-file-write.ts': 1,
'mobile-web-host-requests.ts': 2,
'mobile-web-host-subscriptions.ts': 1,
'mobile-web-markdown-operations.ts': 2,
'mobile-web-native-chat-binding.ts': 1,
'mobile-web-provider-review-creation.ts': 2,
@@ -105,7 +106,7 @@ function shellSources(): Map<string, string> {
function dispatchModules(sources: Map<string, string>, operation: string): string[] {
// The generic arm delegates handle resolution to its bounded executor.
if (operation === 'hostRequest') {
expect(sources.get('mobile-web-capability-execution-arms.ts')).toContain(
expect(sources.get('mobile-web-workspace-capability.ts')).toContain(
'executeMobileWebHostRequest({'
)
return ['mobile-web-host-requests.ts']
@@ -15,6 +15,7 @@ export type { MobileWebOperationGrant } from './mobile-web-production-grant-tabl
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),
snapshot: grantLimits(1 * 1024, 128 * 1024, 2, 4, 1),
@@ -1,3 +1,8 @@
import {
MOBILE_WEB_BRIDGE_MAX_PENDING_REQUESTS,
MOBILE_WEB_BRIDGE_MAX_SUBSCRIPTIONS
} from '../../../src/shared/mobile-web/bridge-limits'
import { MOBILE_WEB_BRIDGE_OPERATIONS } from '../../../src/shared/mobile-web/bridge-operation-registry'
export function mobileWebOperationKey(request: { capability: string; operation: string }): string {
return `${request.capability}.${request.operation}`
}
@@ -7,7 +12,8 @@ export function mobileWebRequestExpectsSubscription(request: {
operation: string
}): boolean {
return (
(request.capability === 'workspace' ||
(request.capability === 'workspace' && request.operation === 'hostSubscribe') ||
((request.capability === 'workspace' ||
request.capability === 'account' ||
request.capability === 'session' ||
request.capability === 'sourceControl' ||
@@ -15,7 +21,7 @@ export function mobileWebRequestExpectsSubscription(request: {
request.capability === 'browser' ||
request.capability === 'nativeChat' ||
request.capability === 'speech') &&
request.operation === 'subscribe'
request.operation === 'subscribe')
)
}
@@ -82,3 +88,41 @@ export function mobileWebPendingRequestForSubscription(
export function mobileWebRequestSurvivesCancellation(pending: { operationKey: string }): boolean {
return pending.operationKey === 'native.alert'
}
export function mobileWebSubscriptionCount(
pending: Iterable<{ subscriptionId?: string }>,
ledgers: readonly { countForOperation: (key: string) => number }[]
): number {
let count = Array.from(pending).filter((entry) => entry.subscriptionId !== undefined).length
for (const [capability, operations] of Object.entries(MOBILE_WEB_BRIDGE_OPERATIONS)) {
for (const [operation, kind] of Object.entries(operations)) {
if (kind === 'subscription') {
for (const ledger of ledgers) {
count += ledger.countForOperation(`${capability}.${operation}`)
}
}
}
}
return count
}
export function mobileWebRequestAtCapacity(args: {
pending: ReadonlyMap<string, { operationKey: string; subscriptionId?: string }>
request: { mode: 'once' | 'subscription'; capability: string; operation: string }
ledgers: readonly { countForOperation: (key: string) => number }[]
isHostRequest: boolean
hostRequestsInFlight: number
maxConcurrent: number
}): boolean {
const key = mobileWebOperationKey(args.request)
return (
(args.isHostRequest && args.hostRequestsInFlight >= 4) ||
args.pending.size >= MOBILE_WEB_BRIDGE_MAX_PENDING_REQUESTS ||
(args.request.mode === 'subscription' &&
mobileWebSubscriptionCount(args.pending.values(), args.ledgers) >=
MOBILE_WEB_BRIDGE_MAX_SUBSCRIPTIONS) ||
mobileWebPendingForOperation(args.pending.values(), key) +
args.ledgers.reduce((sum, ledger) => sum + ledger.countForOperation(key), 0) >=
args.maxConcurrent
)
}
@@ -0,0 +1,48 @@
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'
describe('aggregate subscription admission', () => {
it('counts pending generic streams alongside active legacy and terminal streams', () => {
const pending = new Map([
['pending', { operationKey: 'workspace.hostSubscribe', subscriptionId: 'pending-stream' }]
])
const args = {
pending,
request: {
mode: 'subscription' as const,
capability: 'workspace',
operation: 'hostSubscribe'
},
ledgers: [
{
countForOperation: (key: string) =>
key === 'session.subscribe'
? MOBILE_WEB_BRIDGE_MAX_SUBSCRIPTIONS - 2
: key === 'terminal.subscribe'
? 1
: 0
}
],
isHostRequest: true,
hostRequestsInFlight: 1,
maxConcurrent: 8
}
expect(mobileWebRequestAtCapacity(args)).toBe(true)
pending.clear()
expect(mobileWebRequestAtCapacity(args)).toBe(false)
})
it('retains the actual host work ceiling after page cancellation removes pending state', () => {
expect(
mobileWebRequestAtCapacity({
pending: new Map(),
request: { mode: 'subscription', capability: 'workspace', operation: 'hostSubscribe' },
ledgers: [],
isHostRequest: true,
hostRequestsInFlight: 4,
maxConcurrent: 8
})
).toBe(true)
})
})
@@ -3,6 +3,7 @@ import type {
MobileWebSubscriptionClosure
} from './mobile-web-subscription-closure'
import { MOBILE_WEB_BRIDGE_MAX_SUBSCRIPTIONS } from '../../../src/shared/mobile-web/bridge-contract'
import { mobileWebEncodedByteLength } from './mobile-web-request-accounting'
import { MobileWebBrokerError } from './mobile-web-broker-error'
export type MobileWebSubscriptionRecord = {
@@ -42,6 +43,8 @@ export class MobileWebSubscriptionLedger<
TRecord extends MobileWebSubscriptionRecord = MobileWebSubscriptionRecord
> {
protected readonly records = new Map<string, TRecord>()
private queuedEvents = 0
private queuedBytes = 0
constructor(protected readonly options: MobileWebSubscriptionLedgerOptions<TEvent>) {}
@@ -139,15 +142,31 @@ export class MobileWebSubscriptionLedger<
): void {
const sequence = record.sequence
record.sequence += 1
this.enqueueTask(subscriptionId, record, async () => {
await this.options.postEvent(subscriptionId, sequence, event)
if (retireAfterDelivery) {
this.cancel(subscriptionId)
}
})
this.enqueueTask(
subscriptionId,
record,
async () => {
await this.options.postEvent(subscriptionId, sequence, event)
if (retireAfterDelivery) {
this.cancel(subscriptionId)
}
},
mobileWebEncodedByteLength(event)
)
}
protected enqueueTask(subscriptionId: string, record: TRecord, task: () => Promise<void>): void {
protected enqueueTask(
subscriptionId: string,
record: TRecord,
task: () => Promise<void>,
retainedBytes = 0
): void {
if (this.queuedEvents >= 64 || this.queuedBytes + retainedBytes > 2 * 1024 * 1024) {
this.cancel(subscriptionId, { code: 'rate_limited', retryable: true })
return
}
this.queuedEvents += 1
this.queuedBytes += retainedBytes
record.delivery = record.delivery
.then(async () => {
if (this.isCurrent(subscriptionId, record) && this.canDeliver(subscriptionId, record)) {
@@ -157,6 +176,10 @@ export class MobileWebSubscriptionLedger<
.catch(() => {
this.cancel(subscriptionId, { code: 'unavailable', retryable: true })
})
.finally(() => {
this.queuedEvents -= 1
this.queuedBytes -= retainedBytes
})
}
/** Re-checked at delivery time so a binding revoked while queued cannot publish. */
@@ -0,0 +1,42 @@
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 { executeMobileWebWorkspaceOperation } from './mobile-web-workspace-operations'
type OnceRequest = Extract<
Extract<MobileWebBridgePageMessage, { type: 'request' }>,
{ mode: 'once' }
>
export async function executeWorkspace(
args: MobileWebCapabilityExecutionDependencies,
request: OnceRequest
): Promise<unknown> {
if (request.operation === 'hostCatalog') {
return readMobileWebHostCatalog(args.connectedClient(), request.payload)
}
if (request.operation === 'hostRequest') {
return executeMobileWebHostRequest({
client: args.connectedClient(),
authority: args.workspaceAuthority,
payload: request.payload,
isActive: args.isRequestActive
})
}
if (request.capability !== 'workspace' && request.capability !== 'settings') {
throw new MobileWebBrokerError('unsupported_capability')
}
const result = await executeMobileWebWorkspaceOperation({
capability: request.capability,
operation: request.operation,
payload: request.payload,
client: args.connectedClient(),
authority: args.workspaceAuthority,
snapshots: args.workspaceSnapshots
})
if (request.capability === 'workspace' && request.operation === 'activate') {
args.terminalArtifactAuthority.clear()
}
return result
}
@@ -29,6 +29,7 @@ type StreamRecord = {
Parameters<RpcClient['subscribe']>[3]
>['onTerminalBinaryFrame']
streamIds: Set<number>
serverUnsubscribeMethod?: string
subscriptionId?: string
cancelled: boolean
sent: boolean
@@ -60,7 +61,7 @@ export class MobileRelayRpcStreams {
private readonly streams = new Map<string, StreamRecord>()
private readonly cancelledSubscriptions = new Map<
string,
{ method: string; unsubscribe?: StreamUnsubscribe }
{ method: string; serverUnsubscribeMethod?: string; unsubscribe?: StreamUnsubscribe }
>()
private readonly terminalListeners = new Map<number, (result: unknown) => void>()
private readonly terminalSnapshots = new Map<number, TerminalSnapshotState>()
@@ -75,11 +76,19 @@ export class MobileRelayRpcStreams {
listener: (result: unknown) => void,
subscribeOptions?: Parameters<RpcClient['subscribe']>[3]
): () => void {
if (
subscribeOptions?.serverUnsubscribeMethod &&
this.streams.size + this.cancelledSubscriptions.size >= 128
) {
listener({ type: 'error', message: 'Stream capacity reached' })
return () => {}
}
const id = this.options.nextId()
const stream: StreamRecord = {
method,
params,
listener,
serverUnsubscribeMethod: subscribeOptions?.serverUnsubscribeMethod,
onBinaryFrame: subscribeOptions?.onBinaryFrame,
onTerminalBinaryFrame: subscribeOptions?.onTerminalBinaryFrame,
streamIds: new Set(),
@@ -122,7 +131,11 @@ export class MobileRelayRpcStreams {
this.options.sendFrame({ id: this.options.nextId(), ...cancelled.unsubscribe })
} else if (typeof result.subscriptionId === 'string') {
this.cancelledSubscriptions.delete(response.id)
const unsubscribe = buildReadyStreamUnsubscribe(cancelled.method, result.subscriptionId)
const unsubscribe = buildReadyStreamUnsubscribe(
cancelled.method,
result.subscriptionId,
cancelled.serverUnsubscribeMethod
)
if (unsubscribe) {
this.options.sendFrame({ id: this.options.nextId(), ...unsubscribe })
}
@@ -216,16 +229,29 @@ export class MobileRelayRpcStreams {
}
} else {
const unsubscribe = stream.subscriptionId
? buildReadyStreamUnsubscribe(stream.method, stream.subscriptionId)
? buildReadyStreamUnsubscribe(
stream.method,
stream.subscriptionId,
stream.serverUnsubscribeMethod
)
: null
if (byParams && stream.method === 'session.tabs.subscribe' && !stream.receivedSnapshot) {
// The host registers cleanup only after resolving the initial snapshot.
this.cancelledSubscriptions.set(id, { method: stream.method, unsubscribe: byParams })
} else if (unsubscribe || byParams) {
this.sendUnsubscribe((unsubscribe ?? byParams)!)
} else if (buildServerSubscriptionUnsubscribe(stream.method, 'pending')) {
} else if (
buildServerSubscriptionUnsubscribe(
stream.method,
'pending',
stream.serverUnsubscribeMethod
)
) {
// Keep only the cleanup route while the server assigns its subscription ID.
this.cancelledSubscriptions.set(id, { method: stream.method })
this.cancelledSubscriptions.set(id, {
method: stream.method,
serverUnsubscribeMethod: stream.serverUnsubscribeMethod
})
} else if (stream.subscriptionId) {
this.sendUnsubscribe({
method: stream.method.replace(/\.subscribe$/, '.unsubscribe'),
@@ -25,12 +25,18 @@ function setup(waitForConnected: () => Promise<void> = async () => {}) {
describe('relay server subscription cleanup', () => {
it.each([
['files.watch', 'files.unwatch'],
['accounts.subscribe', 'accounts.unsubscribe']
['accounts.subscribe', 'accounts.unsubscribe'],
['future.feed', 'future.release']
])('releases %s before and after the ready response', async (method, unsubscribeMethod) => {
for (const cancelBeforeReady of [false, true]) {
const { streams, sendFrame, ready } = setup()
const listener = vi.fn()
const cancel = streams.subscribe(method, {}, listener)
const cancel = streams.subscribe(
method,
{},
listener,
method === 'future.feed' ? { serverUnsubscribeMethod: unsubscribeMethod } : undefined
)
await Promise.resolve()
if (cancelBeforeReady) {
cancel()
@@ -0,0 +1,85 @@
import { describe, expect, it, vi } from 'vitest'
import { RpcClientStreamRegistry } from './rpc-client-stream-registry'
function fixture() {
let next = 0
const send = vi.fn().mockReturnValue(true)
const registry = new RpcClientStreamRegistry({
nextId: () => String(++next),
deviceToken: 'private-token',
getState: () => 'connected',
sendEncrypted: send
})
const ready = (subscriptionId: string) =>
registry.handleResponse({
id: '1',
ok: true,
streaming: true,
result: { type: 'ready', subscriptionId },
_meta: { runtimeId: 'host' }
})
return { registry, send, ready }
}
describe('Desktop-advertised stream cleanup', () => {
it.each([true, false])(
'cleans up an arbitrary method when cancellation precedes ready: %s',
(early) => {
const { registry, send, ready } = fixture()
const listener = vi.fn()
const stop = registry.subscribe('future.watch', {}, listener, {
serverUnsubscribeMethod: 'future.release'
})
if (early) {
stop()
}
ready('opaque-lease')
if (!early) {
stop()
}
expect(send).toHaveBeenLastCalledWith(
expect.objectContaining({
method: 'future.release',
params: { subscriptionId: 'opaque-lease' }
})
)
expect(registry.size()).toBe(0)
if (early) {
expect(listener).not.toHaveBeenCalled()
}
}
)
it('does not cancel an old lease while a reconnect is awaiting its new ready frame', () => {
const { registry, send, ready } = fixture()
const stop = registry.subscribe('future.watch', {}, () => {}, {
serverUnsubscribeMethod: 'future.release'
})
ready('old')
registry.markForReplay()
registry.replayAfterAuthentication()
stop()
expect(send.mock.calls.filter(([request]) => request.method === 'future.release')).toHaveLength(
0
)
ready('new')
expect(send).toHaveBeenLastCalledWith(
expect.objectContaining({ params: { subscriptionId: 'new' } })
)
})
it('bounds cancelled streams that have not supplied their cleanup token yet', () => {
const { registry } = fixture()
for (let index = 0; index < 128; index++) {
registry.subscribe('future.watch', {}, () => {}, {
serverUnsubscribeMethod: 'future.release'
})()
}
const listener = vi.fn()
registry.subscribe('future.watch', {}, listener, { serverUnsubscribeMethod: 'future.release' })
expect(registry.size()).toBe(128)
expect(listener).toHaveBeenCalledWith(
expect.objectContaining({ type: 'error', error: { code: 'rate_limited' } })
)
})
})
@@ -1,19 +1,23 @@
export function buildServerSubscriptionUnsubscribe(
method: string,
subscriptionId: string
subscriptionId: string,
cleanupMethod?: string
): { method: string; params: { subscriptionId: string } } | null {
const unsubscribeMethod = {
'browser.screencast': 'browser.screencast.unsubscribe',
'accounts.subscribe': 'accounts.unsubscribe',
'files.watch': 'files.unwatch',
'runtime.clientEvents.subscribe': 'runtime.clientEvents.unsubscribe'
}[method]
const unsubscribeMethod =
cleanupMethod ??
{
'browser.screencast': 'browser.screencast.unsubscribe',
'accounts.subscribe': 'accounts.unsubscribe',
'files.watch': 'files.unwatch',
'runtime.clientEvents.subscribe': 'runtime.clientEvents.unsubscribe'
}[method]
return unsubscribeMethod ? { method: unsubscribeMethod, params: { subscriptionId } } : null
}
export function buildReadyStreamUnsubscribe(
method: string,
subscriptionId: string
subscriptionId: string,
cleanupMethod?: string
): { method: string; params: { subscriptionId: string } } | null {
return buildServerSubscriptionUnsubscribe(method, subscriptionId)
return buildServerSubscriptionUnsubscribe(method, subscriptionId, cleanupMethod)
}
@@ -21,11 +21,13 @@ import type { ConnectionState, RpcResponse, RpcSuccess } from './types'
export type RpcStreamingListener = (result: unknown) => void
export type RpcStreamSubscribeOptions = {
serverUnsubscribeMethod?: string
onBinaryFrame?: (frame: BrowserScreencastFrame) => void
onTerminalBinaryFrame?: (frame: TerminalStreamFrame) => boolean
}
type StreamRequest = {
serverUnsubscribeMethod?: string
method: string
params: unknown
listener: RpcStreamingListener
@@ -56,10 +58,19 @@ export class RpcClientStreamRegistry {
listener: RpcStreamingListener,
subscribeOptions?: RpcStreamSubscribeOptions
): () => void {
if (subscribeOptions?.serverUnsubscribeMethod && this.streams.size >= 128) {
listener({
type: 'error',
error: { code: 'rate_limited' },
message: 'Too many pending streams'
})
return () => {}
}
const id = this.options.nextId()
const stream: StreamRequest = {
method,
params,
serverUnsubscribeMethod: subscribeOptions?.serverUnsubscribeMethod,
listener,
onBinaryFrame: subscribeOptions?.onBinaryFrame,
onTerminalBinaryFrame: subscribeOptions?.onTerminalBinaryFrame
@@ -110,6 +121,7 @@ export class RpcClientStreamRegistry {
this.browserSlot.clearAll()
for (const [id, stream] of this.streams) {
stream.sent = false
stream.subscriptionId = undefined
this.resetTerminalRouting(id)
}
}
@@ -202,6 +214,7 @@ export class RpcClientStreamRegistry {
return
}
if (
stream?.serverUnsubscribeMethod ||
stream?.method === 'runtime.clientEvents.subscribe' ||
stream?.method === 'accounts.subscribe' ||
stream?.method === 'files.watch'
@@ -243,7 +256,12 @@ export class RpcClientStreamRegistry {
if (!stream.subscriptionId) {
return
}
const unsubscribe = buildReadyStreamUnsubscribe(stream.method, stream.subscriptionId)
const unsubscribe = stream.serverUnsubscribeMethod
? {
method: stream.serverUnsubscribeMethod,
params: { subscriptionId: stream.subscriptionId }
}
: buildReadyStreamUnsubscribe(stream.method, stream.subscriptionId)
if (unsubscribe) {
this.sendRpc(unsubscribe.method, unsubscribe.params)
}
+2
View File
@@ -46,6 +46,7 @@ import { AGENT_SESSION_METHODS } from './agent-session'
import { STRUCTURED_AGENT_SESSION_METHODS } from './structured-agent-session'
import { ARTIFACT_METHODS } from './artifacts'
import { MOBILE_WEB_FILE_READ_METHODS } from './mobile-web-file-reads'
import { MOBILE_WEB_FILE_WATCH_METHOD } from './mobile-web-file-watch'
import { MOBILE_WEB_HOST_CATALOG_METHOD } from './mobile-web-host-catalog'
import { MOBILE_WEB_PACKAGE_METHODS } from './mobile-web-package'
import { MOBILE_FILE_WRITE_METHODS } from './mobile-file-write-if-unchanged'
@@ -105,5 +106,6 @@ export const ALL_RPC_METHODS: readonly RpcAnyMethod[] = [
...UPDATER_METHODS,
MOBILE_WEB_HOST_CATALOG_METHOD,
...MOBILE_WEB_FILE_READ_METHODS,
MOBILE_WEB_FILE_WATCH_METHOD,
...MOBILE_WEB_PACKAGE_METHODS
]
@@ -0,0 +1,33 @@
import { afterEach, expect, it, vi } from 'vitest'
import { MOBILE_WEB_FILE_WATCH_METHOD } from './mobile-web-file-watch'
import { FILE_METHODS } from './files'
import { isStreamingMethod, type RpcContext } from '../core'
afterEach(() => vi.restoreAllMocks())
it('preserves stream lifecycle and future fields without publishing the workspace selector', async () => {
const source = FILE_METHODS.find((method) => method.name === 'files.watch')!
if (!isStreamingMethod(source)) {
throw new Error('Expected file watch stream')
}
const events = [
{ type: 'ready', subscriptionId: 'opaque-stream' },
{ type: 'changed', worktree: 'id:private-workspace', events: [], future: { revision: 42 } },
{ type: 'end' }
]
const handler = vi
.spyOn(source, 'handler')
.mockImplementation(async (_params, _context, emit) => {
events.forEach(emit)
})
const context = { connectionId: 'connection', signal: new AbortController().signal } as RpcContext
const params = { worktree: 'id:private-workspace' }
const emit = vi.fn()
await MOBILE_WEB_FILE_WATCH_METHOD.handler(params, context, emit)
expect(handler).toHaveBeenCalledWith(params, context, expect.any(Function))
expect(emit.mock.calls.map(([event]) => event)).toEqual([
events[0],
{ type: 'changed', events: [], future: { revision: 42 } },
events[2]
])
})
@@ -0,0 +1,22 @@
import { defineStreamingMethod, isStreamingMethod } from '../core'
import { FILE_METHODS } from './files'
const source = FILE_METHODS.find((method) => method.name === 'files.watch')
if (!source || !isStreamingMethod(source)) {
throw new Error('Missing file watch stream')
}
const fileWatch = source
export const MOBILE_WEB_FILE_WATCH_METHOD = defineStreamingMethod({
name: 'mobileWeb.files.watch',
params: fileWatch.params,
handler: (params, context, emit) =>
fileWatch.handler(params, context, (event) => {
if (typeof event !== 'object' || event === null || Array.isArray(event)) {
throw new Error('Invalid file watch event')
}
const pageEvent: Record<string, unknown> = { ...event }
delete pageEvent.worktree
emit(pageEvent)
})
})
@@ -21,6 +21,16 @@ const PAGE_METHODS = new Map(
])
)
const fileWatchGrant = {
method: 'mobileWeb.files.watch',
workspaceParam: 'worktree',
mode: 'subscription',
unsubscribeMethod: 'files.unwatch',
maxRequestBytes: 16 * 1024,
maxResponseBytes: 512 * 1024
}
PAGE_METHODS.set(fileWatchGrant.method, fileWatchGrant)
export const MOBILE_WEB_HOST_CATALOG_METHOD = defineMethod({
name: 'mobileWeb.host.catalog',
params: MobileWebHostCatalogPayloadSchema,
@@ -1,3 +1,4 @@
import { subscribeHostSourceControl } from './mobile-web-source-control-host-subscription'
import {
MOBILE_WEB_BRIDGE_PROTOCOL_VERSION,
type MobileWebBridgeMessageContext,
@@ -70,6 +71,7 @@ export class MobileWebBridgeClient {
readonly fileReadChunk!: MobileWebFileRequestClient['readChunk']
readonly fileWrite!: MobileWebFileRequestClient['write']
readonly fileOpen!: MobileWebFileRequestClient['open']
readonly hostSubscribe: MobileWebBridgeSubscriptionClient['subscribeHost']
readonly fileResolveTerminalPath!: MobileWebFileRequestClient['resolveTerminalPath']
readonly fileReadTerminalArtifactChunk!: MobileWebFileRequestClient['readTerminalArtifactChunk']
readonly fileReleaseTerminalArtifact!: MobileWebFileRequestClient['releaseTerminalArtifact']
@@ -208,14 +210,15 @@ export class MobileWebBridgeClient {
this.terminalDeviceInputRequest = terminalRequests.deviceInput.bind(terminalRequests)
Object.assign(this, mobileWebBrowserNavigationClientBindings(this.requests))
this.subscriptions = new MobileWebBridgeSubscriptionClient({
getGrant: (capability) =>
this.grants.get(mobileWebBridgeOperationKey(capability, 'subscribe')),
getGrant: (capability, operation = 'subscribe') =>
this.grants.get(mobileWebBridgeOperationKey(capability, operation)),
postMessage: options.postMessage,
envelope,
createMessageId: (excluded) => this.uniqueMessageId(excluded),
otherPendingCount: () => this.requests.pendingCount(),
requestTimeoutMs: options.requestTimeoutMs
})
this.hostSubscribe = this.subscriptions.subscribeHost.bind(this.subscriptions)
this.account = new MobileWebAccountRequestClient(this.requests, this.subscriptions)
this.agentHistory = new MobileWebAgentHistoryRequestClient(this.requests)
this.speech = new MobileWebSpeechRequestClient(this.requests, this.subscriptions)
@@ -261,7 +264,7 @@ export class MobileWebBridgeClient {
sourceControlSubscribe(
...args: Parameters<MobileWebBridgeSubscriptionClient['subscribeSourceControl']>
): MobileWebBridgeSubscription {
return this.subscriptions.subscribeSourceControl(...args)
return subscribeHostSourceControl(this.requests, this.subscriptions, ...args)
}
browserSubscribe(
@@ -1,3 +1,4 @@
import { hostSubscriptionSetup } from './mobile-web-host-subscription-setup'
import type {
MobileWebBridgeCapability,
MobileWebBridgePageMessage,
@@ -10,16 +11,15 @@ import type {
MobileWebSessionSubscribePayload,
MobileWebWorkspaceChange
} from '../../shared/mobile-web/bridge-operation-contract'
import {
MobileWebTerminalEventSchema,
MobileWebTerminalRequestSchema,
type MobileWebTerminalEvent,
type MobileWebTerminalRequest
import type {
MobileWebTerminalEvent,
MobileWebTerminalRequest
} from '../../shared/mobile-web/terminal-stream-contract'
import { MobileWebBridgeClientError } from './mobile-web-bridge-client-error'
import { deliverMobileWebSubscriptionEvent } from './mobile-web-bridge-subscription-event-delivery'
import {
accountSubscriptionSetup,
terminalSubscriptionSetup,
browserSubscriptionSetup,
nativeChatSubscriptionSetup,
sessionSubscriptionSetup,
@@ -54,7 +54,10 @@ export class MobileWebBridgeSubscriptionClient {
constructor(
private readonly options: {
getGrant: (capability: MobileWebBridgeCapability) => OperationGrant | undefined
getGrant: (
capability: MobileWebBridgeCapability,
operation?: string
) => OperationGrant | undefined
postMessage: (message: MobileWebBridgePageMessage) => boolean
envelope: () => Pick<MobileWebBridgePageMessage, 'version' | 'shellSessionId' | 'buildId'>
createMessageId: (excluded?: string) => string
@@ -95,17 +98,14 @@ export class MobileWebBridgeSubscriptionClient {
onEvent: (event: MobileWebTerminalEvent) => void,
onError: (error: MobileWebBridgeClientError) => void
): MobileWebTerminalBridgeSubscription {
const subscription = this.subscribeWith({
capability: 'terminal',
payload,
payloadSchema: MobileWebTerminalRequestSchema,
eventSchema: MobileWebTerminalEventSchema,
onEvent: (value) => onEvent(value as MobileWebTerminalEvent),
onError
})
const subscription = this.subscribeWith(terminalSubscriptionSetup(payload, onEvent, onError))
return { ...subscription, streamId: subscription.subscriptionId }
}
subscribeHost(...args: Parameters<typeof hostSubscriptionSetup>) {
return this.subscribeWith(hostSubscriptionSetup(...args))
}
subscribeSourceControl(
...args: Parameters<typeof sourceControlSubscriptionSetup>
): MobileWebBridgeSubscription {
@@ -123,7 +123,7 @@ export class MobileWebBridgeSubscriptionClient {
private subscribeWith(
setup: MobileWebBridgeSubscriptionSetup
): MobileWebBridgeSubscription & { subscriptionId: string } {
const grant = this.options.getGrant(setup.capability)
const grant = this.options.getGrant(setup.capability, setup.operation ?? 'subscribe')
const parsedPayload = setup.payloadSchema.safeParse(setup.payload)
const error = mobileWebSubscriptionSetupError({
disposed: this.disposed,
@@ -184,7 +184,7 @@ export class MobileWebBridgeSubscriptionClient {
requestId,
subscriptionId,
capability: setup.capability,
operation: 'subscribe',
operation: setup.operation ?? 'subscribe',
payload: parsedPayload.data
})
if (!posted) {
@@ -1,3 +1,9 @@
import {
MobileWebTerminalRequestSchema,
MobileWebTerminalEventSchema,
type MobileWebTerminalRequest,
type MobileWebTerminalEvent
} from '../../shared/mobile-web/terminal-stream-contract'
import type { z } from 'zod'
import type { MobileWebBridgeCapability } from '../../shared/mobile-web/bridge-contract'
import {
@@ -38,6 +44,7 @@ import {
} from '../../shared/mobile-web/speech-operation-contract'
export type MobileWebBridgeSubscriptionSetup = {
operation?: string
capability: MobileWebBridgeCapability
payload: unknown
payloadSchema: z.ZodType<unknown>
@@ -151,3 +158,18 @@ export function sourceControlSubscriptionSetup(
onError
}
}
export function terminalSubscriptionSetup(
payload: Extract<MobileWebTerminalRequest, { operation: 'subscribe' }>,
onEvent: (event: MobileWebTerminalEvent) => void,
onError: (error: MobileWebBridgeClientError) => void
): MobileWebBridgeSubscriptionSetup {
return {
capability: 'terminal',
payload,
payloadSchema: MobileWebTerminalRequestSchema,
eventSchema: MobileWebTerminalEventSchema,
onEvent: (value) => onEvent(value as MobileWebTerminalEvent),
onError
}
}
@@ -0,0 +1,23 @@
import {
MobileWebHostRequestPayloadSchema,
MobileWebHostResultSchema,
type MobileWebHostRequestPayload
} from '../../shared/mobile-web/host-rpc-contract'
import type { MobileWebBridgeClientError } from './mobile-web-bridge-client-error'
import type { MobileWebBridgeSubscriptionSetup } from './mobile-web-bridge-subscription-setup'
export function hostSubscriptionSetup(
payload: MobileWebHostRequestPayload,
onEvent: (event: unknown) => void,
onError: (error: MobileWebBridgeClientError) => void
): MobileWebBridgeSubscriptionSetup {
return {
capability: 'workspace',
operation: 'hostSubscribe',
payload,
payloadSchema: MobileWebHostRequestPayloadSchema,
eventSchema: MobileWebHostResultSchema,
onEvent,
onError
}
}
@@ -0,0 +1,81 @@
import { MobileWebSourceControlSubscribePayloadSchema } from '../../shared/mobile-web/source-control-operation-contract'
import { MobileWebBridgeClientError } from './mobile-web-bridge-client-error'
import type { MobileWebBridgeSubscriptionClient } from './mobile-web-bridge-subscription-client'
import type { MobileWebOneShotRequestClient } from './mobile-web-one-shot-request-client'
import type { MobileWebBridgeSubscription } from './mobile-web-bridge-subscription'
export function subscribeHostSourceControl(
requests: MobileWebOneShotRequestClient,
subscriptions: MobileWebBridgeSubscriptionClient,
...[payload, onEvent, onError]: Parameters<
MobileWebBridgeSubscriptionClient['subscribeSourceControl']
>
): MobileWebBridgeSubscription {
const legacy = () => subscriptions.subscribeSourceControl(payload, onEvent, onError)
if (
!requests.supports('workspace', 'hostSubscribe') ||
!MobileWebSourceControlSubscribePayloadSchema.safeParse(payload).success
) {
return legacy()
}
let cancelled = false
let ready = false
let current: MobileWebBridgeSubscription = subscriptions.subscribeHost(
{
method: 'mobileWeb.files.watch',
workspaceId: payload.workspaceId,
params: {}
},
(event) => {
if (cancelled || typeof event !== 'object' || event === null) {
return
}
if ('type' in event && event.type === 'changed') {
const events = 'events' in event && Array.isArray(event.events) ? event.events : null
if (!events) {
onError(new MobileWebBridgeClientError('invalid_message', false))
return
}
const overflow =
events.length > 5_000 ||
events.some(
(entry: unknown) =>
typeof entry === 'object' &&
entry !== null &&
'kind' in entry &&
entry.kind === 'overflow'
)
onEvent({ workspaceId: payload.workspaceId, reason: overflow ? 'overflow' : 'changed' })
} else if ('type' in event && (event.type === 'error' || event.type === 'end')) {
onError(new MobileWebBridgeClientError('unavailable', true))
}
},
(error) => {
if (!cancelled && (ready || error.code !== 'unsupported_capability')) {
onError(error)
}
}
)
const settled = current.ready
.catch((error: unknown) => {
if (
!cancelled &&
error instanceof MobileWebBridgeClientError &&
error.code === 'unsupported_capability'
) {
current = legacy()
return current.ready
}
throw error
})
.then(() => {
ready = true
})
return {
ready: settled,
unsubscribe() {
cancelled = true
current.unsubscribe()
}
}
}
@@ -7,6 +7,7 @@ export type MobileWebBridgeOperationKind = 'read' | 'mutation' | 'subscription'
export const MOBILE_WEB_BRIDGE_OPERATIONS = {
workspace: {
hostCatalog: 'read',
hostSubscribe: 'subscription',
hostRequest: 'mutation',
snapshot: 'read',
repositories: 'read',
@@ -21,6 +21,8 @@ export const MobileWebHostRequestPayloadSchema = z
export const MobileWebHostGrantSchema = z.object({
method: MethodSchema,
mode: z.enum(['once', 'subscription']).optional(),
unsubscribeMethod: MethodSchema.optional(),
workspaceParam: z
.string()
.min(1)