From 041e40265b6269dad54ad8f2f3ef23e080f219fb Mon Sep 17 00:00:00 2001 From: Jinwoo-H Date: Sun, 6 Sep 2026 19:08:04 -0400 Subject: [PATCH] feat(mobile): forward native-chat feeds through generic subscriptions Retain future transcript fields, validate host-owned bindings for every event, and use per-stream cleanup tokens. Preserve legacy host/shell fallback and cancellation semantics. Required gates and page export pass. --- .../2026-09-06-long-lived-mobile-shell.md | 22 +++- ...ile-web-host-native-chat-roundtrip.test.ts | 90 ++++++++++++- ...web-host-session-native-chat-operations.ts | 9 +- src/main/runtime/rpc/methods/index.ts | 2 + .../rpc/methods/mobile-web-host-catalog.ts | 16 ++- .../mobile-web-native-chat-stream.test.ts | 118 ++++++++++++++++++ .../methods/mobile-web-native-chat-stream.ts | 85 +++++++++++++ .../src/mobile-web-bridge-client.ts | 9 +- .../mobile-web-host-native-chat-binding.ts | 32 +++++ .../src/mobile-web-host-native-chat-read.ts | 24 ++-- ...obile-web-host-native-chat-subscription.ts | 87 +++++++++++++ .../mobile-web-native-chat-request-client.ts | 17 ++- 12 files changed, 481 insertions(+), 30 deletions(-) create mode 100644 src/main/runtime/rpc/methods/mobile-web-native-chat-stream.test.ts create mode 100644 src/main/runtime/rpc/methods/mobile-web-native-chat-stream.ts create mode 100644 src/mobile-web/src/mobile-web-host-native-chat-binding.ts create mode 100644 src/mobile-web/src/mobile-web-host-native-chat-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 0ec61ac757e..c7fb61fb06b 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 @@ -46,7 +46,7 @@ source; it must not depend on those temporary files to explain remaining work. - [ ] Extend host-advertised metadata without freezing new domain schemas into the shell. Keep catalog queries bounded, not the lifetime method vocabulary. -- [ ] Support page-safe host-owned opaque handles alongside existing workspace +- [x] Support page-safe host-owned opaque handles alongside existing workspace handles; retire authority on document/host/client replacement. - [ ] Preserve intent fingerprints across opaque ID translation. Validate page intent before mapping, recompute host fingerprints afterward, and preserve @@ -73,7 +73,7 @@ folder and SSH workspaces still use their actual execution owner. - [ ] Preserve direct/relay setup, ready, unsubscribe and reconnect behavior. - [x] Reuse the existing subscription ledger with bounded pending event bytes and event count; enforce aggregate subscription ceilings. -- [x] Forward domain event shapes without APK-owned projections (source-control file watch). +- [x] Forward domain event shapes without APK-owned projections (source-control file watch and native-chat transcript feed). - [ ] Migrate native-chat/session/source-control/account feeds. - [ ] Preserve terminal binary capability negotiation, acknowledgements, backpressure and resync; never silently substitute JSON stream semantics. @@ -329,3 +329,21 @@ Review's Back returns to its source route, and Session need not render literal `tabs` text. Updated both readiness checks. One manually assisted Back action was used to inspect this failure; this is not a complete unattended E2E pass. Rerun with the committed fixture before claiming full iOS coverage; Android remains open. + +### Native-chat feed migration + +Read/resource slice committed as `e14164f974f`; harness readiness fixes as +`79f54a8b49e`. Native-chat subscriptions now use the same host-owned opaque +resources and generic stream. The host announces a random cleanup token and +reuses the existing transcript watcher. Every event checks the authoritative +provider binding before publishing; changed bindings close once. Setup failures +release registered watchers. Cancellation during page binding opens no watcher, +and cancellation inside ready delivery cannot publish a subsequent snapshot. +Future transcript/event fields stay intact on the active generic path. + +All required gates pass: mobile 836 files / 5,525 tests; root 321 files / 2,709 +tests. Additional focused cancellation/setup checks pass (mobile 7, root 4). +Logs: `/tmp/orca-ota-e2e/native-chat-stream-gates/`. +Page build: `81e5d5f47e703f4047e8544d5f3812459adf71ed45ec0951a840129797b6e266`. +Next: rerun iOS with the updated fixture, then continue mutations, session/domain +migration and page-owned preferences/routes. No complete platform run claimed yet. diff --git a/mobile/src/mobile-web/mobile-web-host-native-chat-roundtrip.test.ts b/mobile/src/mobile-web/mobile-web-host-native-chat-roundtrip.test.ts index 9f9e8a69989..0cbe7a0ec58 100644 --- a/mobile/src/mobile-web/mobile-web-host-native-chat-roundtrip.test.ts +++ b/mobile/src/mobile-web/mobile-web-host-native-chat-roundtrip.test.ts @@ -59,8 +59,11 @@ function fixture(genericHost = true, genericShell = true) { ok: true, result: { grants: genericHost - ? ['bind', 'read'].map((operation) => ({ + ? ['bind', 'read', 'subscribe'].map((operation) => ({ method: `mobileWeb.nativeChat.${operation}`, + ...(operation === 'subscribe' + ? { mode: 'subscription', unsubscribeMethod: 'nativeChat.unsubscribe' } + : {}), workspaceParam: 'worktree', pageSessionParam: 'pageSession', maxRequestBytes: 16384, @@ -75,12 +78,25 @@ function fixture(genericHost = true, genericShell = true) { } return { ok: true, result: transcript } }) + let emit: (event: unknown) => void = () => {} + const unsubscribe = vi.fn() + const subscribe = vi.fn((_method, _params, listener) => { + emit = listener + return unsubscribe + }) const bridge = createMobileWebBridgeRoundtripFixture({ grants: MOBILE_WEB_PRODUCTION_GRANTS, ...(genericShell ? {} : { shellFeatures: [] }), - rpcClient: { sendRequest } as unknown as RpcClient + rpcClient: { sendRequest, subscribe } as unknown as RpcClient }) - return { ...bridge, sendRequest, transcript } + return { + ...bridge, + sendRequest, + transcript, + subscribe, + unsubscribe, + emit: (event: unknown) => emit(event) + } } describe('native-chat generic read migration', () => { @@ -127,4 +143,72 @@ describe('native-chat generic read migration', () => { } expect(JSON.stringify(f.shellMessages)).not.toContain('private-session') }) + it.each([ + [true, true], + [false, true], + [true, false] + ])('stream host=%s shell=%s', async (host, shell) => { + const f = fixture(host, shell) + const workspaceId = (await f.client.workspaceSnapshot({ limit: 10 })).workspaces[0]!.id + const session = await f.client.sessionSnapshot({ workspaceId }) + const tab = session.tabs[0] + if (tab.type !== 'terminal' || !tab.nativeChatSessionId) { + throw new Error('Missing chat fixture') + } + const onEvent = vi.fn() + const subscription = f.client.nativeChat.subscribeForTab( + tab.id, + { workspaceId, sessionId: tab.nativeChatSessionId, limit: 20 }, + onEvent, + vi.fn() + ) + await subscription.ready + const generic = host && shell + expect(f.subscribe.mock.calls[0][0]).toBe( + generic ? 'mobileWeb.nativeChat.subscribe' : 'nativeChat.subscribe' + ) + const event = { type: 'snapshot', ...f.transcript } + f.emit(event) + await vi.waitFor(() => expect(onEvent).toHaveBeenCalledOnce()) + if (generic) { + expect(onEvent).toHaveBeenCalledWith(event) + expect(f.subscribe.mock.calls[0][1]).toMatchObject({ + pageSession: MOBILE_WEB_BRIDGE_ROUNDTRIP_CONTEXT.shellSessionId, + resourceId: 'opaque-resource' + }) + } + subscription.unsubscribe() + expect(f.unsubscribe).toHaveBeenCalledOnce() + }) + it('does not subscribe after cancellation during host binding', async () => { + const f = fixture() + const workspaceId = (await f.client.workspaceSnapshot({ limit: 10 })).workspaces[0]!.id + const session = await f.client.sessionSnapshot({ workspaceId }) + const tab = session.tabs[0] + if (tab.type !== 'terminal' || !tab.nativeChatSessionId) { + throw new Error('Missing chat fixture') + } + const original = f.sendRequest.getMockImplementation()! + const bound = Promise.withResolvers>>() + f.sendRequest.mockImplementation((...args) => + args[0] === 'mobileWeb.nativeChat.bind' ? bound.promise : original(...args) + ) + const onError = vi.fn() + const subscription = f.client.nativeChat.subscribeForTab( + tab.id, + { workspaceId, sessionId: tab.nativeChatSessionId, limit: 20 }, + vi.fn(), + onError + ) + await vi.waitFor(() => + expect( + f.sendRequest.mock.calls.some(([method]) => method === 'mobileWeb.nativeChat.bind') + ).toBe(true) + ) + subscription.unsubscribe() + bound.resolve({ ok: true, result: { resourceId: 'opaque-resource' } }) + await expect(subscription.ready).rejects.toMatchObject({ code: 'cancelled' }) + expect(f.subscribe).not.toHaveBeenCalled() + expect(onError).not.toHaveBeenCalled() + }) }) diff --git a/mobile/src/session/web-host-session-native-chat-operations.ts b/mobile/src/session/web-host-session-native-chat-operations.ts index d8affb0210b..0405059bb94 100644 --- a/mobile/src/session/web-host-session-native-chat-operations.ts +++ b/mobile/src/session/web-host-session-native-chat-operations.ts @@ -25,11 +25,10 @@ export function webHostSessionNativeChatOperations( return (await client.nativeChat.readability({ workspaceId })).readable }, subscribe(target, limit, onEvent, onError) { - const subscription = client.nativeChatSubscribe( - bridgeTarget(target, { limit }), - onEvent, - onError - ) + const payload = bridgeTarget(target, { limit }) + const subscription = target.terminalId + ? client.nativeChat.subscribeForTab(target.terminalId, payload, onEvent, onError) + : client.nativeChatSubscribe(payload, onEvent, onError) void subscription.ready.catch(() => {}) return subscription.unsubscribe }, diff --git a/src/main/runtime/rpc/methods/index.ts b/src/main/runtime/rpc/methods/index.ts index 7d924fc2400..a2f2ddde28c 100644 --- a/src/main/runtime/rpc/methods/index.ts +++ b/src/main/runtime/rpc/methods/index.ts @@ -1,3 +1,4 @@ +import { MOBILE_WEB_NATIVE_CHAT_STREAM_METHOD } from './mobile-web-native-chat-stream' import { MOBILE_WEB_NATIVE_CHAT_METHODS } from './mobile-web-native-chat' import type { RpcAnyMethod } from '../core' import { STATUS_METHODS } from './status' @@ -109,5 +110,6 @@ export const ALL_RPC_METHODS: readonly RpcAnyMethod[] = [ ...MOBILE_WEB_FILE_READ_METHODS, MOBILE_WEB_FILE_WATCH_METHOD, ...MOBILE_WEB_NATIVE_CHAT_METHODS, + MOBILE_WEB_NATIVE_CHAT_STREAM_METHOD, ...MOBILE_WEB_PACKAGE_METHODS ] 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 dadc3ff31ea..ebb25af160e 100644 --- a/src/main/runtime/rpc/methods/mobile-web-host-catalog.ts +++ b/src/main/runtime/rpc/methods/mobile-web-host-catalog.ts @@ -1,8 +1,11 @@ -import { MobileWebHostCatalogPayloadSchema } from '../../../../shared/mobile-web/host-rpc-contract' +import { + MobileWebHostCatalogPayloadSchema, + type MobileWebHostGrant +} from '../../../../shared/mobile-web/host-rpc-contract' import { defineMethod } from '../core' // Only page-safe results belong here; transport credentials never enter this catalog. -const PAGE_METHODS = new Map( +const PAGE_METHODS = new Map( [ 'git.status', 'git.diff', @@ -24,7 +27,7 @@ const PAGE_METHODS = new Map( ]) ) -const fileWatchGrant = { +const fileWatchGrant: MobileWebHostGrant = { method: 'mobileWeb.files.watch', workspaceParam: 'worktree', mode: 'subscription', @@ -34,6 +37,13 @@ const fileWatchGrant = { } PAGE_METHODS.set(fileWatchGrant.method, fileWatchGrant) +PAGE_METHODS.set('mobileWeb.nativeChat.subscribe', { + ...fileWatchGrant, + method: 'mobileWeb.nativeChat.subscribe', + pageSessionParam: 'pageSession', + unsubscribeMethod: 'nativeChat.unsubscribe' +}) + export const MOBILE_WEB_HOST_CATALOG_METHOD = defineMethod({ name: 'mobileWeb.host.catalog', params: MobileWebHostCatalogPayloadSchema, diff --git a/src/main/runtime/rpc/methods/mobile-web-native-chat-stream.test.ts b/src/main/runtime/rpc/methods/mobile-web-native-chat-stream.test.ts new file mode 100644 index 00000000000..6194b018111 --- /dev/null +++ b/src/main/runtime/rpc/methods/mobile-web-native-chat-stream.test.ts @@ -0,0 +1,118 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest' +import type { RpcContext } from '../core' +const subscribe = vi.hoisted(() => vi.fn()) +vi.mock('./native-chat', async () => { + const { z } = await import('zod') + return { + NATIVE_CHAT_METHODS: [ + { + name: 'nativeChat.subscribe', + stream: true, + params: z.object({}).passthrough(), + handler: subscribe + } + ] + } +}) +import { bindMobileWebNativeChat } from './mobile-web-native-chat-binding' +import { MOBILE_WEB_NATIVE_CHAT_STREAM_METHOD } from './mobile-web-native-chat-stream' + +async function fixture() { + const providerSession = { id: 'provider', transcriptPath: '/private/transcript' } + const current = { worktreeId: 'workspace', agent: 'codex', providerSession } + const runtime = { + listMobileSessionTabs: vi.fn().mockResolvedValue({ + worktree: 'workspace', + tabs: [ + { + id: 'tab', + type: 'terminal', + terminal: 'terminal', + agentStatus: { agentType: 'codex', providerSession } + } + ] + }), + resolveNativeChatTranscriptBinding: vi.fn().mockReturnValue(current), + registerSubscriptionCleanup: vi.fn(), + cleanupSubscription: vi.fn() + } + const context = { runtime, connectionId: 'connection' } as unknown as RpcContext + const scope = { worktree: 'id:workspace', pageSession: 'page' } + const resource = await bindMobileWebNativeChat(context, { ...scope, tabId: 'tab' }) + const params = { ...scope, ...resource, read: { limit: 20, subscriptionId: 'forged' } } + return { runtime, context, params } +} +beforeEach(() => subscribe.mockReset()) +describe('opaque native-chat feed', () => { + it('announces a private cleanup token and forwards future fields', async () => { + const f = await fixture() + const event = { + type: 'snapshot', + messages: [], + hasMore: false, + futureField: { addedByDesktop: true } + } + subscribe.mockImplementationOnce(async (_params, _context, emit) => emit(event)) + const emit = vi.fn() + await MOBILE_WEB_NATIVE_CHAT_STREAM_METHOD.handler(f.params, f.context, emit) + const ready = emit.mock.calls[0][0] + expect(ready).toMatchObject({ type: 'ready', subscriptionId: expect.any(String) }) + expect(ready.subscriptionId).not.toBe('forged') + expect(ready.subscriptionId).not.toContain('provider') + expect(emit.mock.calls[1]).toEqual([event]) + expect(subscribe).toHaveBeenCalledWith( + expect.objectContaining({ + sessionId: 'provider', + terminal: 'terminal', + worktreeId: 'workspace', + subscriptionId: ready.subscriptionId + }), + f.context, + expect.any(Function) + ) + }) + + it('closes once without publishing after its provider binding changes', async () => { + const f = await fixture() + let publish!: (event: unknown) => void + subscribe.mockImplementationOnce(async (_params, _context, emit) => { + publish = emit + }) + const emit = vi.fn() + await MOBILE_WEB_NATIVE_CHAT_STREAM_METHOD.handler(f.params, f.context, emit) + f.runtime.resolveNativeChatTranscriptBinding.mockReturnValue(null) + publish({ type: 'appended', messages: [{ secret: 'wrong-session' }] }) + publish({ type: 'appended', messages: [] }) + expect(emit.mock.calls.map(([event]) => event.type)).toEqual(['ready', 'error']) + expect(JSON.stringify(emit.mock.calls)).not.toContain('wrong-session') + expect(f.runtime.cleanupSubscription).toHaveBeenCalledOnce() + }) + + it('does not publish the initial snapshot after cancellation inside ready delivery', async () => { + const f = await fixture() + const received: unknown[] = [] + let publish!: (event: unknown) => void + subscribe.mockImplementationOnce(async (_params, _context, emit) => { + publish = emit + emit({ type: 'snapshot', messages: [] }) + }) + await MOBILE_WEB_NATIVE_CHAT_STREAM_METHOD.handler(f.params, f.context, (event) => { + received.push(event) + if ((event as { type: string }).type === 'ready') { + publish({ type: 'end' }) + } + }) + expect(received).toEqual([ + { type: 'ready', subscriptionId: expect.any(String) }, + { type: 'end' } + ]) + }) + it('cleans up a watcher when native setup fails after registration', async () => { + const f = await fixture() + subscribe.mockRejectedValueOnce(new Error('Watcher setup failed')) + await expect( + MOBILE_WEB_NATIVE_CHAT_STREAM_METHOD.handler(f.params, f.context, vi.fn()) + ).rejects.toThrow('Watcher setup failed') + expect(f.runtime.cleanupSubscription).toHaveBeenCalledOnce() + }) +}) diff --git a/src/main/runtime/rpc/methods/mobile-web-native-chat-stream.ts b/src/main/runtime/rpc/methods/mobile-web-native-chat-stream.ts new file mode 100644 index 00000000000..9bfe297330d --- /dev/null +++ b/src/main/runtime/rpc/methods/mobile-web-native-chat-stream.ts @@ -0,0 +1,85 @@ +import { randomUUID } from 'node:crypto' +import { z } from 'zod' +import { defineStreamingMethod, isStreamingMethod } from '../core' +import { NATIVE_CHAT_METHODS } from './native-chat' +import { + MobileWebChatScope, + mobileWebNativeChatHostParams, + resolveMobileWebNativeChat +} from './mobile-web-native-chat-binding' + +const source = NATIVE_CHAT_METHODS.find((method) => method.name === 'nativeChat.subscribe') +if (!source || !isStreamingMethod(source)) { + throw new Error('Missing native chat stream') +} +const stream = source + +export const MOBILE_WEB_NATIVE_CHAT_STREAM_METHOD = defineStreamingMethod({ + name: 'mobileWeb.nativeChat.subscribe', + params: MobileWebChatScope.extend({ + resourceId: z.string().min(1).max(160), + read: z.record(z.string(), z.unknown()) + }), + handler: async (params, context, emit) => { + const binding = await resolveMobileWebNativeChat(context, params) + if (context.signal?.aborted) { + return + } + const subscriptionId = randomUUID() + const input = stream.params!.parse({ + ...mobileWebNativeChatHostParams(binding, params.read), + subscriptionId + }) + let ready = false + let closed = false + const announce = () => { + if (!ready) { + ready = true + emit({ type: 'ready', subscriptionId }) + } + } + const cleanup = () => + context.runtime.cleanupSubscription( + `nativeChat:${context.connectionId ?? 'local'}:${subscriptionId}` + ) + try { + await stream.handler(input, context, (event) => { + if (closed) { + return + } + announce() + if (closed) { + return + } + const type = + typeof event === 'object' && event !== null && 'type' in event ? event.type : undefined + if (type === 'end' || type === 'error') { + closed = true + emit(event) + return + } + const current = context.runtime.resolveNativeChatTranscriptBinding(binding.terminal) + if ( + !current || + current.worktreeId !== binding.worktreeId || + current.agent !== binding.agent || + current.providerSession?.id !== binding.sessionId || + current.providerSession.transcriptPath !== binding.transcriptPath + ) { + closed = true + emit({ type: 'error', message: 'Chat session is no longer available' }) + cleanup() + return + } + emit(event) + }) + } catch (error) { + closed = true + cleanup() + throw error + } + if (!closed) { + announce() + } + } +}) diff --git a/src/mobile-web/src/mobile-web-bridge-client.ts b/src/mobile-web/src/mobile-web-bridge-client.ts index 443e18849fb..6c0c98cf838 100644 --- a/src/mobile-web/src/mobile-web-bridge-client.ts +++ b/src/mobile-web/src/mobile-web-bridge-client.ts @@ -204,10 +204,6 @@ export class MobileWebBridgeClient { mobileWebSessionClientBindings(new MobileWebSessionRequestClient(this.requests)) ) this.native = new MobileWebNativeRequestClient(this.requests) - this.nativeChat = new MobileWebNativeChatRequestClient( - this.requests, - this.shellFeatures.has(MOBILE_WEB_SHELL_HOST_PAGE_SESSION_FEATURE) - ) this.markdown = new MobileWebMarkdownRequestClient(this.requests) const terminalRequests = new MobileWebTerminalRequestClient(this.requests) this.terminalRequest = terminalRequests.request.bind(terminalRequests) @@ -222,6 +218,11 @@ export class MobileWebBridgeClient { otherPendingCount: () => this.requests.pendingCount(), requestTimeoutMs: options.requestTimeoutMs }) + this.nativeChat = new MobileWebNativeChatRequestClient( + this.requests, + this.shellFeatures.has(MOBILE_WEB_SHELL_HOST_PAGE_SESSION_FEATURE), + this.subscriptions + ) this.hostSubscribe = this.subscriptions.subscribeHost.bind(this.subscriptions) this.account = new MobileWebAccountRequestClient(this.requests, this.subscriptions) this.agentHistory = new MobileWebAgentHistoryRequestClient(this.requests) diff --git a/src/mobile-web/src/mobile-web-host-native-chat-binding.ts b/src/mobile-web/src/mobile-web-host-native-chat-binding.ts new file mode 100644 index 00000000000..8d35e00de0d --- /dev/null +++ b/src/mobile-web/src/mobile-web-host-native-chat-binding.ts @@ -0,0 +1,32 @@ +import { MobileWebBridgeClientError } from './mobile-web-bridge-client-error' +import { readMobileWebHostMethods, requestMobileWebHost } from './mobile-web-host-request-client' +import type { MobileWebOneShotRequestClient } from './mobile-web-one-shot-request-client' + +export async function bindMobileWebHostNativeChat( + requests: MobileWebOneShotRequestClient, + workspaceId: string, + tabId: string, + method: string +): Promise { + if ( + !requests.supports('workspace', 'hostRequest') || + !requests.supports('workspace', 'hostCatalog') + ) { + return null + } + const methods = ['mobileWeb.nativeChat.bind', method] + const catalog = await readMobileWebHostMethods(requests, methods) + if (!methods.every((method) => catalog.grants.some((grant) => grant.method === method))) { + return null + } + const bound = await requestMobileWebHost(requests, methods[0], workspaceId, { tabId }) + if ( + typeof bound !== 'object' || + bound === null || + !('resourceId' in bound) || + typeof bound.resourceId !== 'string' + ) { + throw new MobileWebBridgeClientError('invalid_message', false) + } + return bound.resourceId +} diff --git a/src/mobile-web/src/mobile-web-host-native-chat-read.ts b/src/mobile-web/src/mobile-web-host-native-chat-read.ts index b91f140af35..0507dbbae40 100644 --- a/src/mobile-web/src/mobile-web-host-native-chat-read.ts +++ b/src/mobile-web/src/mobile-web-host-native-chat-read.ts @@ -1,6 +1,7 @@ +import { bindMobileWebHostNativeChat } from './mobile-web-host-native-chat-binding' import type { MobileWebNativeChatReadResult } from '../../shared/mobile-web/native-chat-operation-contract' import { MobileWebBridgeClientError } from './mobile-web-bridge-client-error' -import { readMobileWebHostMethods, requestMobileWebHost } from './mobile-web-host-request-client' +import { requestMobileWebHost } from './mobile-web-host-request-client' import type { MobileWebOneShotRequestClient } from './mobile-web-one-shot-request-client' type MobileWebHostChatReadResult = MobileWebNativeChatReadResult @@ -17,19 +18,18 @@ export async function readMobileWebHostNativeChat( return legacy() } try { - const methods = ['mobileWeb.nativeChat.bind', 'mobileWeb.nativeChat.read'] - const catalog = await readMobileWebHostMethods(requests, methods) - if (!methods.every((method) => catalog.grants.some((grant) => grant.method === method))) { + const method = 'mobileWeb.nativeChat.read' + const resourceId = await bindMobileWebHostNativeChat( + requests, + target.workspaceId, + target.tabId, + method + ) + if (!resourceId) { return legacy() } - const bound = await requestMobileWebHost(requests, methods[0], target.workspaceId, { - tabId: target.tabId - }) - if (!isRecord(bound) || typeof bound.resourceId !== 'string') { - throw new MobileWebBridgeClientError('invalid_message', false) - } - const result = await requestMobileWebHost(requests, methods[1], target.workspaceId, { - resourceId: bound.resourceId, + const result = await requestMobileWebHost(requests, method, target.workspaceId, { + resourceId, read: { limit: target.limit, ...(target.beforeOffset === undefined ? {} : { beforeOffset: target.beforeOffset }) diff --git a/src/mobile-web/src/mobile-web-host-native-chat-subscription.ts b/src/mobile-web/src/mobile-web-host-native-chat-subscription.ts new file mode 100644 index 00000000000..4529aeb4c4f --- /dev/null +++ b/src/mobile-web/src/mobile-web-host-native-chat-subscription.ts @@ -0,0 +1,87 @@ +import type { MobileWebNativeChatEvent } from '../../shared/mobile-web/native-chat-operation-contract' +import { MobileWebBridgeClientError } from './mobile-web-bridge-client-error' +import type { MobileWebBridgeSubscription } from './mobile-web-bridge-subscription' +import type { MobileWebBridgeSubscriptionClient } from './mobile-web-bridge-subscription-client' +import { bindMobileWebHostNativeChat } from './mobile-web-host-native-chat-binding' +import type { MobileWebOneShotRequestClient } from './mobile-web-one-shot-request-client' + +export function subscribeMobileWebHostNativeChat( + requests: MobileWebOneShotRequestClient, + subscriptions: MobileWebBridgeSubscriptionClient, + tabId: string, + ...[payload, onEvent, onError]: Parameters< + MobileWebBridgeSubscriptionClient['subscribeNativeChat'] + > +): MobileWebBridgeSubscription { + const legacy = () => subscriptions.subscribeNativeChat(payload, onEvent, onError) + if (!requests.supports('workspace', 'hostSubscribe')) { + return legacy() + } + let cancelled = false + let current: MobileWebBridgeSubscription | undefined + const ready = (async () => { + let resourceId: string | null + try { + resourceId = await bindMobileWebHostNativeChat( + requests, + payload.workspaceId, + tabId, + 'mobileWeb.nativeChat.subscribe' + ) + } catch (error) { + if ( + !(error instanceof MobileWebBridgeClientError) || + error.code !== 'unsupported_capability' + ) { + throw error + } + resourceId = null + } + if (cancelled) { + throw new MobileWebBridgeClientError('cancelled', false) + } + current = + resourceId === null + ? legacy() + : subscriptions.subscribeHost( + { + method: 'mobileWeb.nativeChat.subscribe', + workspaceId: payload.workspaceId, + params: { + resourceId, + read: { limit: payload.limit, capabilities: { transcriptPending: 1 } } + } + }, + (event) => { + if (cancelled) { + return + } + if (typeof event !== 'object' || event === null || !('type' in event)) { + onError(new MobileWebBridgeClientError('invalid_message', false)) + return + } + if (event.type !== 'ready') { + onEvent(event as MobileWebNativeChatEvent) + } + }, + onError + ) + await current.ready + })() + void ready.catch((error: unknown) => { + if (!cancelled && !current) { + onError( + error instanceof MobileWebBridgeClientError + ? error + : new MobileWebBridgeClientError('unavailable', true) + ) + } + }) + return { + ready, + unsubscribe() { + cancelled = true + current?.unsubscribe() + } + } +} diff --git a/src/mobile-web/src/mobile-web-native-chat-request-client.ts b/src/mobile-web/src/mobile-web-native-chat-request-client.ts index 7d766dbe126..e1d3c8031a0 100644 --- a/src/mobile-web/src/mobile-web-native-chat-request-client.ts +++ b/src/mobile-web/src/mobile-web-native-chat-request-client.ts @@ -1,3 +1,5 @@ +import { subscribeMobileWebHostNativeChat } from './mobile-web-host-native-chat-subscription' +import type { MobileWebBridgeSubscriptionClient } from './mobile-web-bridge-subscription-client' import { readMobileWebHostNativeChat } from './mobile-web-host-native-chat-read' import { MobileWebNativeChatFileSearchPayloadSchema, @@ -50,9 +52,22 @@ import type { MobileWebBridgeRequestOptions } from './mobile-web-bridge-request- export class MobileWebNativeChatRequestClient { constructor( private readonly requests: MobileWebOneShotRequestClient, - private readonly hostPageSession = false + private readonly hostPageSession = false, + private readonly subscriptions?: MobileWebBridgeSubscriptionClient ) {} + subscribeForTab( + tabId: string, + ...args: Parameters + ) { + if (!this.subscriptions) { + throw new MobileWebBridgeClientError('unsupported_capability', false) + } + return this.hostPageSession + ? subscribeMobileWebHostNativeChat(this.requests, this.subscriptions, tabId, ...args) + : this.subscriptions.subscribeNativeChat(...args) + } + readForTab(payload: MobileWebNativeChatReadPayload, tabId: string) { if (!this.hostPageSession) { return this.read(payload)