From a0eee2fc21aeea36169717e98cc58b024f4e8b9d Mon Sep 17 00:00:00 2001 From: Jinwoo-H Date: Sun, 6 Sep 2026 18:45:47 -0400 Subject: [PATCH] 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. --- .../2026-09-06-long-lived-mobile-shell.md | 23 ++- .../mobile-web-capability-broker.ts | 27 ++-- ...ile-web-capability-dispatch-census.test.ts | 2 +- .../mobile-web-capability-execution-arms.ts | 42 ++--- ...e-web-capability-execution-dependencies.ts | 2 + ...obile-web-capability-schema-corpus.test.ts | 5 +- .../mobile-web-capability-subscriptions.ts | 7 + .../mobile-web/mobile-web-host-requests.ts | 25 ++- .../mobile-web-host-stream-roundtrip.test.ts | 112 +++++++++++++ .../mobile-web-host-subscriptions.test.ts | 148 ++++++++++++++++++ .../mobile-web-host-subscriptions.ts | 104 ++++++++++++ ...eb-mutation-reauthorization-census.test.ts | 3 +- .../mobile-web-production-grants.ts | 1 + .../mobile-web-request-accounting.ts | 48 +++++- .../mobile-web-subscription-capacity.test.ts | 48 ++++++ .../mobile-web-subscription-ledger.ts | 37 ++++- .../mobile-web-workspace-capability.ts | 42 +++++ .../src/transport/mobile-relay-rpc-streams.ts | 36 ++++- .../mobile-relay-subscription-cleanup.test.ts | 10 +- .../rpc-client-generic-stream-cleanup.test.ts | 85 ++++++++++ .../rpc-client-server-subscription.ts | 22 +-- .../transport/rpc-client-stream-registry.ts | 20 ++- src/main/runtime/rpc/methods/index.ts | 2 + .../rpc/methods/mobile-web-file-watch.test.ts | 33 ++++ .../rpc/methods/mobile-web-file-watch.ts | 22 +++ .../rpc/methods/mobile-web-host-catalog.ts | 10 ++ .../src/mobile-web-bridge-client.ts | 9 +- .../mobile-web-bridge-subscription-client.ts | 32 ++-- .../mobile-web-bridge-subscription-setup.ts | 22 +++ .../src/mobile-web-host-subscription-setup.ts | 23 +++ ...le-web-source-control-host-subscription.ts | 81 ++++++++++ .../mobile-web/bridge-operation-registry.ts | 1 + src/shared/mobile-web/host-rpc-contract.ts | 2 + 33 files changed, 988 insertions(+), 98 deletions(-) create mode 100644 mobile/src/mobile-web/mobile-web-host-stream-roundtrip.test.ts create mode 100644 mobile/src/mobile-web/mobile-web-host-subscriptions.test.ts create mode 100644 mobile/src/mobile-web/mobile-web-host-subscriptions.ts create mode 100644 mobile/src/mobile-web/mobile-web-subscription-capacity.test.ts create mode 100644 mobile/src/mobile-web/mobile-web-workspace-capability.ts create mode 100644 mobile/src/transport/rpc-client-generic-stream-cleanup.test.ts create mode 100644 src/main/runtime/rpc/methods/mobile-web-file-watch.test.ts create mode 100644 src/main/runtime/rpc/methods/mobile-web-file-watch.ts create mode 100644 src/mobile-web/src/mobile-web-host-subscription-setup.ts create mode 100644 src/mobile-web/src/mobile-web-source-control-host-subscription.ts diff --git a/docs/reference/plans/2026-09-06-long-lived-mobile-shell.md b/docs/reference/plans/2026-09-06-long-lived-mobile-shell.md index 1b88ce06a63..92584ec7112 100644 --- a/docs/reference/plans/2026-09-06-long-lived-mobile-shell.md +++ b/docs/reference/plans/2026-09-06-long-lived-mobile-shell.md @@ -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. diff --git a/mobile/src/mobile-web/mobile-web-capability-broker.ts b/mobile/src/mobile-web/mobile-web-capability-broker.ts index 96269f4854f..03d24500125 100644 --- a/mobile/src/mobile-web/mobile-web-capability-broker.ts +++ b/mobile/src/mobile-web/mobile-web-capability-broker.ts @@ -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, diff --git a/mobile/src/mobile-web/mobile-web-capability-dispatch-census.test.ts b/mobile/src/mobile-web/mobile-web-capability-dispatch-census.test.ts index b01f18d17e1..2f9b96beae9 100644 --- a/mobile/src/mobile-web/mobile-web-capability-dispatch-census.test.ts +++ b/mobile/src/mobile-web/mobile-web-capability-dispatch-census.test.ts @@ -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', () => { diff --git a/mobile/src/mobile-web/mobile-web-capability-execution-arms.ts b/mobile/src/mobile-web/mobile-web-capability-execution-arms.ts index 82a92e719a9..a19ee171bfc 100644 --- a/mobile/src/mobile-web/mobile-web-capability-execution-arms.ts +++ b/mobile/src/mobile-web/mobile-web-capability-execution-arms.ts @@ -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 type OnceRequest = Extract @@ -70,35 +69,6 @@ async function executeBrowser(args: Deps, request: OnceRequest): Promise { - 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 { 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 { + 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({ diff --git a/mobile/src/mobile-web/mobile-web-capability-execution-dependencies.ts b/mobile/src/mobile-web/mobile-web-capability-execution-dependencies.ts index 7f1a50fdc38..ab9b50f64d3 100644 --- a/mobile/src/mobile-web/mobile-web-capability-execution-dependencies.ts +++ b/mobile/src/mobile-web/mobile-web-capability-execution-dependencies.ts @@ -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 diff --git a/mobile/src/mobile-web/mobile-web-capability-schema-corpus.test.ts b/mobile/src/mobile-web/mobile-web-capability-schema-corpus.test.ts index 2795c3d6ef7..54682d663ce 100644 --- a/mobile/src/mobile-web/mobile-web-capability-schema-corpus.test.ts +++ b/mobile/src/mobile-web/mobile-web-capability-schema-corpus.test.ts @@ -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)}` } + : {}) }) } diff --git a/mobile/src/mobile-web/mobile-web-capability-subscriptions.ts b/mobile/src/mobile-web/mobile-web-capability-subscriptions.ts index 476cab8942e..b312345b8bc 100644 --- a/mobile/src/mobile-web/mobile-web-capability-subscriptions.ts +++ b/mobile/src/mobile-web/mobile-web-capability-subscriptions.ts @@ -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, diff --git a/mobile/src/mobile-web/mobile-web-host-requests.ts b/mobile/src/mobile-web/mobile-web-host-requests.ts index 16385305824..9b40d1b2631 100644 --- a/mobile/src/mobile-web/mobile-web-host-requests.ts +++ b/mobile/src/mobile-web/mobile-web-host-requests.ts @@ -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 { +} + +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 { + 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) diff --git a/mobile/src/mobile-web/mobile-web-host-stream-roundtrip.test.ts b/mobile/src/mobile-web/mobile-web-host-stream-roundtrip.test.ts new file mode 100644 index 00000000000..05b598bc159 --- /dev/null +++ b/mobile/src/mobile-web/mobile-web-host-stream-roundtrip.test.ts @@ -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((_method, _params, listener) => { + emit = listener + return unsubscribe + }) + const sendRequest = vi.fn().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() + }) +}) diff --git a/mobile/src/mobile-web/mobile-web-host-subscriptions.test.ts b/mobile/src/mobile-web/mobile-web-host-subscriptions.test.ts new file mode 100644 index 00000000000..36a8a9026ab --- /dev/null +++ b/mobile/src/mobile-web/mobile-web-host-subscriptions.test.ts @@ -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((_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((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() + 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() + }) +}) diff --git a/mobile/src/mobile-web/mobile-web-host-subscriptions.ts b/mobile/src/mobile-web/mobile-web-host-subscriptions.ts new file mode 100644 index 00000000000..693f92900ac --- /dev/null +++ b/mobile/src/mobile-web/mobile-web-host-subscriptions.ts @@ -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 & { + workspaceAuthority: MobileWebWorkspaceAuthority + } + ) { + super({ ...config, operationKey: 'workspace.hostSubscribe' }) + } + + async start( + args: Omit & { + requestId: string + subscriptionId: string + } + ): Promise { + 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) + } +} diff --git a/mobile/src/mobile-web/mobile-web-mutation-reauthorization-census.test.ts b/mobile/src/mobile-web/mobile-web-mutation-reauthorization-census.test.ts index f8674db222c..5dc8290071a 100644 --- a/mobile/src/mobile-web/mobile-web-mutation-reauthorization-census.test.ts +++ b/mobile/src/mobile-web/mobile-web-mutation-reauthorization-census.test.ts @@ -20,6 +20,7 @@ const REAUTHORIZATION_SITES: Record = { '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 { function dispatchModules(sources: Map, 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'] diff --git a/mobile/src/mobile-web/mobile-web-production-grants.ts b/mobile/src/mobile-web/mobile-web-production-grants.ts index cc2c597cf90..77a87ae1f38 100644 --- a/mobile/src/mobile-web/mobile-web-production-grants.ts +++ b/mobile/src/mobile-web/mobile-web-production-grants.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), diff --git a/mobile/src/mobile-web/mobile-web-request-accounting.ts b/mobile/src/mobile-web/mobile-web-request-accounting.ts index 4fa264f1a72..993b7b02cfd 100644 --- a/mobile/src/mobile-web/mobile-web-request-accounting.ts +++ b/mobile/src/mobile-web/mobile-web-request-accounting.ts @@ -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 + 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 + ) +} diff --git a/mobile/src/mobile-web/mobile-web-subscription-capacity.test.ts b/mobile/src/mobile-web/mobile-web-subscription-capacity.test.ts new file mode 100644 index 00000000000..f66b4dbffe5 --- /dev/null +++ b/mobile/src/mobile-web/mobile-web-subscription-capacity.test.ts @@ -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) + }) +}) diff --git a/mobile/src/mobile-web/mobile-web-subscription-ledger.ts b/mobile/src/mobile-web/mobile-web-subscription-ledger.ts index a3a088396b8..5d33d5f05e6 100644 --- a/mobile/src/mobile-web/mobile-web-subscription-ledger.ts +++ b/mobile/src/mobile-web/mobile-web-subscription-ledger.ts @@ -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() + private queuedEvents = 0 + private queuedBytes = 0 constructor(protected readonly options: MobileWebSubscriptionLedgerOptions) {} @@ -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 { + protected enqueueTask( + subscriptionId: string, + record: TRecord, + task: () => Promise, + 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. */ diff --git a/mobile/src/mobile-web/mobile-web-workspace-capability.ts b/mobile/src/mobile-web/mobile-web-workspace-capability.ts new file mode 100644 index 00000000000..6a12706fad7 --- /dev/null +++ b/mobile/src/mobile-web/mobile-web-workspace-capability.ts @@ -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, + { mode: 'once' } +> + +export async function executeWorkspace( + args: MobileWebCapabilityExecutionDependencies, + request: OnceRequest +): Promise { + 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 +} diff --git a/mobile/src/transport/mobile-relay-rpc-streams.ts b/mobile/src/transport/mobile-relay-rpc-streams.ts index 4652bd8079b..6188fa18ab8 100644 --- a/mobile/src/transport/mobile-relay-rpc-streams.ts +++ b/mobile/src/transport/mobile-relay-rpc-streams.ts @@ -29,6 +29,7 @@ type StreamRecord = { Parameters[3] >['onTerminalBinaryFrame'] streamIds: Set + serverUnsubscribeMethod?: string subscriptionId?: string cancelled: boolean sent: boolean @@ -60,7 +61,7 @@ export class MobileRelayRpcStreams { private readonly streams = new Map() private readonly cancelledSubscriptions = new Map< string, - { method: string; unsubscribe?: StreamUnsubscribe } + { method: string; serverUnsubscribeMethod?: string; unsubscribe?: StreamUnsubscribe } >() private readonly terminalListeners = new Map void>() private readonly terminalSnapshots = new Map() @@ -75,11 +76,19 @@ export class MobileRelayRpcStreams { listener: (result: unknown) => void, subscribeOptions?: Parameters[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'), diff --git a/mobile/src/transport/mobile-relay-subscription-cleanup.test.ts b/mobile/src/transport/mobile-relay-subscription-cleanup.test.ts index 85799eeb11a..75b7fcc0217 100644 --- a/mobile/src/transport/mobile-relay-subscription-cleanup.test.ts +++ b/mobile/src/transport/mobile-relay-subscription-cleanup.test.ts @@ -25,12 +25,18 @@ function setup(waitForConnected: () => Promise = 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() diff --git a/mobile/src/transport/rpc-client-generic-stream-cleanup.test.ts b/mobile/src/transport/rpc-client-generic-stream-cleanup.test.ts new file mode 100644 index 00000000000..64165973974 --- /dev/null +++ b/mobile/src/transport/rpc-client-generic-stream-cleanup.test.ts @@ -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' } }) + ) + }) +}) diff --git a/mobile/src/transport/rpc-client-server-subscription.ts b/mobile/src/transport/rpc-client-server-subscription.ts index dd2c430c7a0..7d2fb65ac0b 100644 --- a/mobile/src/transport/rpc-client-server-subscription.ts +++ b/mobile/src/transport/rpc-client-server-subscription.ts @@ -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) } diff --git a/mobile/src/transport/rpc-client-stream-registry.ts b/mobile/src/transport/rpc-client-stream-registry.ts index ad5476b7ce4..2b80bc2eddc 100644 --- a/mobile/src/transport/rpc-client-stream-registry.ts +++ b/mobile/src/transport/rpc-client-stream-registry.ts @@ -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) } diff --git a/src/main/runtime/rpc/methods/index.ts b/src/main/runtime/rpc/methods/index.ts index 97db5066f31..07c74cba3b2 100644 --- a/src/main/runtime/rpc/methods/index.ts +++ b/src/main/runtime/rpc/methods/index.ts @@ -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 ] diff --git a/src/main/runtime/rpc/methods/mobile-web-file-watch.test.ts b/src/main/runtime/rpc/methods/mobile-web-file-watch.test.ts new file mode 100644 index 00000000000..3efd5fd3dc8 --- /dev/null +++ b/src/main/runtime/rpc/methods/mobile-web-file-watch.test.ts @@ -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] + ]) +}) diff --git a/src/main/runtime/rpc/methods/mobile-web-file-watch.ts b/src/main/runtime/rpc/methods/mobile-web-file-watch.ts new file mode 100644 index 00000000000..93de3259196 --- /dev/null +++ b/src/main/runtime/rpc/methods/mobile-web-file-watch.ts @@ -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 = { ...event } + delete pageEvent.worktree + emit(pageEvent) + }) +}) diff --git a/src/main/runtime/rpc/methods/mobile-web-host-catalog.ts b/src/main/runtime/rpc/methods/mobile-web-host-catalog.ts index 77a59d690bc..c4a37ebe38d 100644 --- a/src/main/runtime/rpc/methods/mobile-web-host-catalog.ts +++ b/src/main/runtime/rpc/methods/mobile-web-host-catalog.ts @@ -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, diff --git a/src/mobile-web/src/mobile-web-bridge-client.ts b/src/mobile-web/src/mobile-web-bridge-client.ts index 7484346b912..e7d604e18b5 100644 --- a/src/mobile-web/src/mobile-web-bridge-client.ts +++ b/src/mobile-web/src/mobile-web-bridge-client.ts @@ -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 ): MobileWebBridgeSubscription { - return this.subscriptions.subscribeSourceControl(...args) + return subscribeHostSourceControl(this.requests, this.subscriptions, ...args) } browserSubscribe( diff --git a/src/mobile-web/src/mobile-web-bridge-subscription-client.ts b/src/mobile-web/src/mobile-web-bridge-subscription-client.ts index f35e734fb58..342de4f7fcd 100644 --- a/src/mobile-web/src/mobile-web-bridge-subscription-client.ts +++ b/src/mobile-web/src/mobile-web-bridge-subscription-client.ts @@ -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 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) { + return this.subscribeWith(hostSubscriptionSetup(...args)) + } + subscribeSourceControl( ...args: Parameters ): 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) { diff --git a/src/mobile-web/src/mobile-web-bridge-subscription-setup.ts b/src/mobile-web/src/mobile-web-bridge-subscription-setup.ts index 9f839fc618b..095aab41df7 100644 --- a/src/mobile-web/src/mobile-web-bridge-subscription-setup.ts +++ b/src/mobile-web/src/mobile-web-bridge-subscription-setup.ts @@ -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 @@ -151,3 +158,18 @@ export function sourceControlSubscriptionSetup( onError } } + +export function terminalSubscriptionSetup( + payload: Extract, + onEvent: (event: MobileWebTerminalEvent) => void, + onError: (error: MobileWebBridgeClientError) => void +): MobileWebBridgeSubscriptionSetup { + return { + capability: 'terminal', + payload, + payloadSchema: MobileWebTerminalRequestSchema, + eventSchema: MobileWebTerminalEventSchema, + onEvent: (value) => onEvent(value as MobileWebTerminalEvent), + onError + } +} diff --git a/src/mobile-web/src/mobile-web-host-subscription-setup.ts b/src/mobile-web/src/mobile-web-host-subscription-setup.ts new file mode 100644 index 00000000000..31890c6cf78 --- /dev/null +++ b/src/mobile-web/src/mobile-web-host-subscription-setup.ts @@ -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 + } +} diff --git a/src/mobile-web/src/mobile-web-source-control-host-subscription.ts b/src/mobile-web/src/mobile-web-source-control-host-subscription.ts new file mode 100644 index 00000000000..cc3851e96b8 --- /dev/null +++ b/src/mobile-web/src/mobile-web-source-control-host-subscription.ts @@ -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() + } + } +} diff --git a/src/shared/mobile-web/bridge-operation-registry.ts b/src/shared/mobile-web/bridge-operation-registry.ts index 448199438ba..5e06049326a 100644 --- a/src/shared/mobile-web/bridge-operation-registry.ts +++ b/src/shared/mobile-web/bridge-operation-registry.ts @@ -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', diff --git a/src/shared/mobile-web/host-rpc-contract.ts b/src/shared/mobile-web/host-rpc-contract.ts index 06568c50971..322b158ccce 100644 --- a/src/shared/mobile-web/host-rpc-contract.ts +++ b/src/shared/mobile-web/host-rpc-contract.ts @@ -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)