mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 08:03:20 +00:00
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.
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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<RpcClient['subscribe']>((_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<Awaited<ReturnType<RpcClient['sendRequest']>>>()
|
||||
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()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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
|
||||
},
|
||||
|
||||
@@ -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
|
||||
]
|
||||
|
||||
@@ -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<string, MobileWebHostGrant>(
|
||||
[
|
||||
'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,
|
||||
|
||||
@@ -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()
|
||||
})
|
||||
})
|
||||
@@ -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()
|
||||
}
|
||||
}
|
||||
})
|
||||
@@ -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)
|
||||
|
||||
@@ -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<string | null> {
|
||||
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
|
||||
}
|
||||
@@ -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 })
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<MobileWebBridgeSubscriptionClient['subscribeNativeChat']>
|
||||
) {
|
||||
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)
|
||||
|
||||
Reference in New Issue
Block a user