refactor(mobile-web): put the speech ledger on the shared subscription base

The speech ledger was the one shell ledger still hand-rolling its record
map, cancellation and posting because it is push-driven. It maps onto
MobileWebSubscriptionLedger with no base-class change: admit plus a
records.set stands in for open, and enqueue carries the broadcast. The
authority now takes a MobileWebSubscriptionLedgerConfig, which retires
the per-subscription post/closed plumbing, and the broker hands all
three consumers one shared subscriptionPosts() sender.

Speech stays out of replaceClient's closeAll on purpose: its feed is the
shell's device runtime, not the host RPC client, so it survives the swap
and a terminal closure frame would wedge the dictation hook in error
with no resubscribe.

Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb
This commit is contained in:
Jinwoo-H
2026-09-06 15:40:50 -04:00
parent 63cb1ea524
commit 3f069dfc9c
11 changed files with 128 additions and 157 deletions
@@ -4,6 +4,8 @@ import type {
MobileWebBridgeShellMessage
} from '../../../src/shared/mobile-web/bridge-contract'
import { mobileWebBrokerEnvelope } from './mobile-web-broker-envelope'
import { mobileWebSubscriptionClosedPoster } from './mobile-web-subscription-closure'
import type { MobileWebSubscriptionLedgerConfig } from './mobile-web-subscription-ledger'
export class MobileWebBrokerMessageSender {
constructor(
@@ -57,6 +59,15 @@ export class MobileWebBrokerMessageSender {
})
}
/** A ledger only ever needs this channel's liveness, its event frame, and its closure frame. */
subscriptionPosts(): MobileWebSubscriptionLedgerConfig<unknown> {
return {
isActive: this.options.isActive,
postEvent: (subscriptionId, sequence, event) => this.event(subscriptionId, sequence, event),
postClosed: mobileWebSubscriptionClosedPoster(this)
}
}
private async post(message: MobileWebBridgeShellMessage): Promise<void> {
if (this.options.isActive()) {
await this.options.postMessage(message)
@@ -21,7 +21,6 @@ import { MobileWebCapabilityAuthorities } from './mobile-web-capability-authorit
import type { MobileWebCapabilityBrokerOptions } from './mobile-web-capability-broker-options'
import { MobileWebBrokerMessageSender } from './mobile-web-broker-message-sender'
import { MobileWebBrokerReplayGuard } from './mobile-web-broker-replay-guard'
import { mobileWebSubscriptionClosedPoster } from './mobile-web-subscription-closure'
import { rememberMobileWebBrokerRoute } from './mobile-web-broker-route-memory'
import { resolveMobileWebHostNavigationRoute } from './mobile-web-host-navigation-route'
import {
@@ -46,7 +45,7 @@ export class MobileWebCapabilityBroker {
private readonly replay = new MobileWebBrokerReplayGuard()
private readonly subscriptions: MobileWebCapabilitySubscriptions
private readonly terminalStreams: MobileWebTerminalStreams
private readonly speechAuthority = new MobileWebSpeechAuthority()
private readonly speechAuthority: MobileWebSpeechAuthority
private readonly rateLimiter: MobileWebOperationRateLimiter
private readonly commitMessageGeneration = new MobileWebCommitMessageGeneration()
private readonly authorities: MobileWebCapabilityAuthorities
@@ -61,23 +60,22 @@ export class MobileWebCapabilityBroker {
isActive: () => !this.disposed && options.isActive(),
postMessage: options.postMessage
})
const posts = this.messages.subscriptionPosts()
this.subscriptions = new MobileWebCapabilitySubscriptions({
isActive: () => !this.disposed && this.options.isActive(),
messages: this.messages,
...posts,
browserAuthority: this.authorities.browser,
nativeChatAuthority: this.authorities.nativeChat,
workspaceAuthority: this.authorities.workspace
})
this.terminalStreams = new MobileWebTerminalStreams({
isActive: () => !this.disposed && this.options.isActive(),
...posts,
clientId: options.terminalClientId,
now: options.now,
onFlowMetrics: options.onTerminalFlowMetrics,
onResync: options.onTerminalResync,
workspaceAuthority: this.authorities.workspace,
postEvent: this.messages.event.bind(this.messages),
postClosed: mobileWebSubscriptionClosedPoster(this.messages)
workspaceAuthority: this.authorities.workspace
})
this.speechAuthority = new MobileWebSpeechAuthority(posts)
}
async handle(message: MobileWebBridgePageMessage): Promise<void> {
@@ -245,8 +243,6 @@ export class MobileWebCapabilityBroker {
sourceControlSubscriptions: this.subscriptions.sourceControl,
sourceControlBranchCompare: this.authorities.sourceControlBranchCompare,
speechAuthority: this.speechAuthority,
postSpeechEvent: this.messages.event.bind(this.messages),
postSpeechClosed: mobileWebSubscriptionClosedPoster(this.messages),
workspaceSubscriptions: this.subscriptions.workspace,
terminalStreams: this.terminalStreams,
commitMessageGeneration: this.commitMessageGeneration,
@@ -271,9 +271,7 @@ async function subscribeSpeech(args: Deps, request: SubscriptionRequest): Promis
MobileWebSpeechSubscribePayloadSchema.parse(request.payload)
args.speechAuthority.subscribe({
requestId: request.requestId,
subscriptionId: request.subscriptionId,
post: (sequence, event) => args.postSpeechEvent(request.subscriptionId, sequence, event),
closed: (closure) => args.postSpeechClosed(request.subscriptionId, closure)
subscriptionId: request.subscriptionId
})
return null
}
@@ -1,6 +1,4 @@
import type { MobileWebPostSubscriptionClosed } from './mobile-web-subscription-closure'
import type { MobileWebBridgePageMessage } from '../../../src/shared/mobile-web/bridge-contract'
import type { MobileWebSpeechEvent } from '../../../src/shared/mobile-web/speech-operation-contract'
import type { RpcClient } from '../transport/rpc-client'
import type { MobileWebAccountSubscriptions } from './mobile-web-account-subscriptions'
import type { MobileWebAgentHistoryAuthority } from './mobile-web-agent-history-authority'
@@ -43,12 +41,6 @@ export type MobileWebCapabilityExecutionDependencies = {
sourceControlSubscriptions: MobileWebSourceControlSubscriptions
sourceControlBranchCompare: MobileWebSourceControlBranchComparePager
speechAuthority: MobileWebSpeechAuthority
postSpeechEvent: (
subscriptionId: string,
sequence: number,
event: MobileWebSpeechEvent
) => Promise<void>
postSpeechClosed: MobileWebPostSubscriptionClosed
workspaceSubscriptions: MobileWebWorkspaceSubscriptions
terminalStreams: MobileWebTerminalStreams
commitMessageGeneration: MobileWebCommitMessageGeneration
@@ -1,10 +1,11 @@
import { MobileWebAccountSubscriptions } from './mobile-web-account-subscriptions'
import { mobileWebSubscriptionClosedPoster } from './mobile-web-subscription-closure'
import type { MobileWebSubscriptionClosure } from './mobile-web-subscription-closure'
import type { MobileWebSubscriptionLedgerHandle } from './mobile-web-subscription-ledger'
import type {
MobileWebSubscriptionLedgerConfig,
MobileWebSubscriptionLedgerHandle
} from './mobile-web-subscription-ledger'
import type { MobileWebBrowserAuthority } from './mobile-web-browser-authority'
import { MobileWebBrowserStreams } from './mobile-web-browser-streams'
import type { MobileWebBrokerMessageSender } from './mobile-web-broker-message-sender'
import { MobileWebSessionSubscriptions } from './mobile-web-session-subscriptions'
import type { MobileWebNativeChatAuthority } from './mobile-web-native-chat-authority'
import { MobileWebNativeChatSubscriptions } from './mobile-web-native-chat-subscriptions'
@@ -21,17 +22,18 @@ export class MobileWebCapabilitySubscriptions {
readonly workspace: MobileWebWorkspaceSubscriptions
private readonly ledgers: MobileWebSubscriptionLedgerHandle[]
constructor(args: {
isActive: () => boolean
messages: MobileWebBrokerMessageSender
browserAuthority: MobileWebBrowserAuthority
nativeChatAuthority: MobileWebNativeChatAuthority
workspaceAuthority: MobileWebWorkspaceAuthority
}) {
const postEvent = (subscriptionId: string, sequence: number, event: unknown) =>
args.messages.event(subscriptionId, sequence, event)
const postClosed = mobileWebSubscriptionClosedPoster(args.messages)
const shared = { isActive: args.isActive, postEvent, postClosed }
constructor(
args: MobileWebSubscriptionLedgerConfig<unknown> & {
browserAuthority: MobileWebBrowserAuthority
nativeChatAuthority: MobileWebNativeChatAuthority
workspaceAuthority: MobileWebWorkspaceAuthority
}
) {
const shared = {
isActive: args.isActive,
postEvent: args.postEvent,
postClosed: args.postClosed
}
this.account = new MobileWebAccountSubscriptions(shared)
this.browser = new MobileWebBrowserStreams({
...shared,
@@ -7,15 +7,7 @@ import type { MobileWebSpeechRuntime } from './mobile-web-speech-runtime'
describe('MobileWebSpeechAuthority', () => {
it('keeps PCM in the shell and finishes without cancelling a successful transcript', async () => {
const harness = createHarness()
const events: unknown[] = []
harness.authority.subscribe({
requestId: 'request-1',
subscriptionId: 'subscription-1',
post: async (sequence, event) => {
events.push({ sequence, event })
},
closed: vi.fn()
})
harness.authority.subscribe({ requestId: 'request-1', subscriptionId: 'subscription-1' })
harness.sendRequest.mockImplementation(async (method) => {
if (method === 'speech.dictation.finish') {
return success({ text: ` ${'a'.repeat(40_000)} ` })
@@ -40,10 +32,10 @@ describe('MobileWebSpeechAuthority', () => {
harness.sendRequest.mock.calls.filter(([method]) => method === 'speech.dictation.cancel')
).toHaveLength(0)
await vi.waitFor(() =>
expect(events).toEqual([
{ sequence: 0, event: { status: 'recording' } },
{ sequence: 1, event: { status: 'processing' } },
{ sequence: 2, event: { status: 'idle' } }
expect(harness.postEvent.mock.calls).toEqual([
['subscription-1', 0, { status: 'recording' }],
['subscription-1', 1, { status: 'processing' }],
['subscription-1', 2, { status: 'idle' }]
])
)
expect(harness.runtime.releaseKeepAwake).toHaveBeenCalledOnce()
@@ -114,13 +106,7 @@ describe('MobileWebSpeechAuthority', () => {
it('cancels recording when the native audio session is interrupted', async () => {
const harness = createHarness()
const post = vi.fn(async () => {})
harness.authority.subscribe({
requestId: 'request-1',
subscriptionId: 'subscription-1',
post,
closed: vi.fn()
})
harness.authority.subscribe({ requestId: 'request-1', subscriptionId: 'subscription-1' })
harness.sendRequest.mockResolvedValue(success({}))
await harness.authority.start(harness.client)
@@ -133,7 +119,7 @@ describe('MobileWebSpeechAuthority', () => {
)
)
await vi.waitFor(() =>
expect(post).toHaveBeenLastCalledWith(1, {
expect(harness.postEvent).toHaveBeenLastCalledWith('subscription-1', 1, {
status: 'idle',
reason: 'interrupted'
})
@@ -187,11 +173,18 @@ function createHarness() {
} satisfies MobileWebSpeechRuntime
const sendRequest = vi.fn<RpcClient['sendRequest']>()
const client = { sendRequest } as unknown as RpcClient
const postEvent = vi.fn(async () => {})
const postClosed = vi.fn()
return {
authority: new MobileWebSpeechAuthority(async () => runtime),
authority: new MobileWebSpeechAuthority(
{ isActive: () => true, postEvent, postClosed },
async () => runtime
),
runtime,
client,
sendRequest,
postEvent,
postClosed,
get microphone() {
return microphone
},
@@ -1,4 +1,4 @@
import type { MobileWebSubscriptionClosure } from './mobile-web-subscription-closure'
import type { MobileWebSubscriptionLedgerConfig } from './mobile-web-subscription-ledger'
import type {
MobileWebSpeechEvent,
MobileWebSpeechStartResult,
@@ -20,7 +20,7 @@ import { MobileWebSpeechSubscriptions } from './mobile-web-speech-subscriptions'
export class MobileWebSpeechAuthority {
private readonly audio = new MobileWebSpeechAudioForwarder()
private readonly subscriptions = new MobileWebSpeechSubscriptions()
private readonly subscriptions: MobileWebSpeechSubscriptions
private runtime: MobileWebSpeechRuntime | null = null
private runtimePromise: Promise<MobileWebSpeechRuntime> | null = null
private removeMicrophoneListener: (() => void) | null = null
@@ -32,19 +32,14 @@ export class MobileWebSpeechAuthority {
private disposed = false
constructor(
config: MobileWebSubscriptionLedgerConfig<MobileWebSpeechEvent>,
private readonly loadRuntime: () => Promise<MobileWebSpeechRuntime> = async () =>
(await import('./mobile-web-speech-native-runtime')).createMobileWebSpeechNativeRuntime()
) {}
) {
this.subscriptions = new MobileWebSpeechSubscriptions(config)
}
subscribe(args: {
requestId: string
subscriptionId: string
post: (sequence: number, event: MobileWebSpeechEvent) => Promise<void>
closed: (closure: MobileWebSubscriptionClosure) => void
}): void {
if (this.disposed) {
throw new MobileWebBrokerError('invalid_request')
}
subscribe(args: { requestId: string; subscriptionId: string }): void {
this.subscriptions.start(args)
}
@@ -3,21 +3,20 @@ import { MobileWebSpeechSubscriptions } from './mobile-web-speech-subscriptions'
describe('MobileWebSpeechSubscriptions', () => {
it('orders events and drops queued delivery after cancellation', async () => {
const subscriptions = new MobileWebSpeechSubscriptions()
let releaseFirst: (() => void) | undefined
const post = vi.fn((sequence: number) =>
const postEvent = vi.fn((_subscriptionId: string, sequence: number) =>
sequence === 0
? new Promise<void>((resolve) => {
releaseFirst = resolve
})
: Promise.resolve()
)
subscriptions.start({
requestId: 'request-1',
subscriptionId: 'subscription-1',
post,
closed: vi.fn()
const subscriptions = new MobileWebSpeechSubscriptions({
isActive: () => true,
postEvent,
postClosed: vi.fn()
})
subscriptions.start({ requestId: 'request-1', subscriptionId: 'subscription-1' })
subscriptions.post({ status: 'recording' })
subscriptions.post({ status: 'processing' })
@@ -27,24 +26,39 @@ describe('MobileWebSpeechSubscriptions', () => {
await Promise.resolve()
await Promise.resolve()
expect(post).toHaveBeenCalledOnce()
expect(post).toHaveBeenCalledWith(0, { status: 'recording' })
expect(postEvent).toHaveBeenCalledOnce()
expect(postEvent).toHaveBeenCalledWith('subscription-1', 0, { status: 'recording' })
})
it('tells the page when a delivery failure retires the dictation stream', async () => {
const subscriptions = new MobileWebSpeechSubscriptions()
const closed = vi.fn()
subscriptions.start({
requestId: 'request-1',
subscriptionId: 'subscription-1',
post: () => Promise.reject(new Error('post failed')),
closed
const postClosed = vi.fn()
const subscriptions = new MobileWebSpeechSubscriptions({
isActive: () => true,
postEvent: () => Promise.reject(new Error('post failed')),
postClosed
})
subscriptions.start({ requestId: 'request-1', subscriptionId: 'subscription-1' })
subscriptions.post({ status: 'recording' })
await vi.waitFor(() => expect(closed).toHaveBeenCalledOnce())
await vi.waitFor(() => expect(postClosed).toHaveBeenCalledOnce())
expect(closed).toHaveBeenCalledWith({ code: 'unavailable', retryable: true })
expect(postClosed).toHaveBeenCalledWith('subscription-1', {
code: 'unavailable',
retryable: true
})
expect(subscriptions.cancel('subscription-1')).toBeNull()
})
it('refuses a new dictation stream once the ledger is disposed', () => {
const subscriptions = new MobileWebSpeechSubscriptions({
isActive: () => true,
postEvent: async () => {},
postClosed: vi.fn()
})
subscriptions.dispose()
expect(() =>
subscriptions.start({ requestId: 'request-1', subscriptionId: 'subscription-1' })
).toThrow('invalid_request')
})
})
@@ -1,79 +1,35 @@
import type { MobileWebSubscriptionClosure } from './mobile-web-subscription-closure'
import {
MobileWebSubscriptionLedger,
type MobileWebSubscriptionLedgerConfig
} from './mobile-web-subscription-ledger'
import type { MobileWebSpeechEvent } from '../../../src/shared/mobile-web/speech-operation-contract'
import { MobileWebBrokerError } from './mobile-web-broker-error'
type SpeechSubscriber = {
requestId: string
sequence: number
delivery: Promise<void>
post: (sequence: number, event: MobileWebSpeechEvent) => Promise<void>
closed: (closure: MobileWebSubscriptionClosure) => void
}
export class MobileWebSpeechSubscriptions {
private readonly records = new Map<string, SpeechSubscriber>()
/** Push-driven: the shell's own dictation runtime feeds every entry, so there is no host handle to
* open and no `closeAll` on a client swap — the authority publishes `session-replaced` instead. */
export class MobileWebSpeechSubscriptions extends MobileWebSubscriptionLedger<MobileWebSpeechEvent> {
private disposed = false
start(args: {
requestId: string
subscriptionId: string
post: (sequence: number, event: MobileWebSpeechEvent) => Promise<void>
closed: (closure: MobileWebSubscriptionClosure) => void
}): void {
if (this.disposed || this.records.has(args.subscriptionId)) {
constructor(config: MobileWebSubscriptionLedgerConfig<MobileWebSpeechEvent>) {
super({ ...config, operationKey: 'speech.subscribe' })
}
start(args: { requestId: string; subscriptionId: string }): void {
if (this.disposed) {
throw new MobileWebBrokerError('invalid_request')
}
this.records.set(args.subscriptionId, {
requestId: args.requestId,
sequence: 0,
delivery: Promise.resolve(),
post: args.post,
closed: args.closed
})
this.admit(args.subscriptionId)
this.records.set(args.subscriptionId, this.newRecord(args.requestId))
}
post(event: MobileWebSpeechEvent): void {
for (const [subscriptionId, record] of this.records) {
const sequence = record.sequence++
record.delivery = record.delivery
.then(async () => {
if (this.records.get(subscriptionId) !== record) {
return
}
await record.post(sequence, event)
})
.catch(() => {
if (this.records.get(subscriptionId) === record) {
this.records.delete(subscriptionId)
record.closed({ code: 'unavailable', retryable: true })
}
})
this.enqueue(subscriptionId, record, event)
}
}
cancel(subscriptionId: string): string | null {
const record = this.records.get(subscriptionId)
if (!record) {
return null
}
this.records.delete(subscriptionId)
return record.requestId
}
cancelByRequest(requestId: string): void {
for (const [subscriptionId, record] of this.records) {
if (record.requestId === requestId) {
this.records.delete(subscriptionId)
}
}
}
countForOperation(operationKey: string): number {
return operationKey === 'speech.subscribe' ? this.records.size : 0
}
dispose(): void {
override dispose(): void {
this.disposed = true
this.records.clear()
super.dispose()
}
}
@@ -8,6 +8,8 @@ import { MobileWebNativeChatAuthority } from './mobile-web-native-chat-authority
import { MobileWebNativeChatSubscriptions } from './mobile-web-native-chat-subscriptions'
import { MobileWebSessionSubscriptions } from './mobile-web-session-subscriptions'
import { MobileWebSourceControlSubscriptions } from './mobile-web-source-control-subscriptions'
import { MobileWebSpeechSubscriptions } from './mobile-web-speech-subscriptions'
import type { MobileWebSpeechEvent } from '../../../src/shared/mobile-web/speech-operation-contract'
import { MobileWebWorkspaceSubscriptions } from './mobile-web-workspace-subscriptions'
import {
mobileWebHostWorkspaceIdFromHost,
@@ -210,6 +212,18 @@ const LEDGER_CASES: LedgerCase[] = [
})
return host.emit
}
},
{
name: 'speech',
// Push-driven from the shell's dictation runtime, so no host frame can be unusable.
invalidCode: null,
invalid: { status: 'bogus' },
valid: { status: 'recording' },
open: async (posts) => {
const subscriptions = new MobileWebSpeechSubscriptions(posts)
subscriptions.start({ requestId: 'request-1', subscriptionId: SUBSCRIPTION_ID })
return (value) => subscriptions.post(value as MobileWebSpeechEvent)
}
}
]
@@ -75,15 +75,15 @@ describe('subscription ledger teardown', () => {
it('fans closeAll out across every capability ledger', () => {
const messages: MobileWebBridgeShellMessage[] = []
const subscriptions = new MobileWebCapabilitySubscriptions({
const sender = new MobileWebBrokerMessageSender({
context: MOBILE_WEB_BRIDGE_ROUNDTRIP_CONTEXT,
isActive: () => true,
messages: new MobileWebBrokerMessageSender({
context: MOBILE_WEB_BRIDGE_ROUNDTRIP_CONTEXT,
isActive: () => true,
postMessage: (message) => {
messages.push(message)
}
}),
postMessage: (message) => {
messages.push(message)
}
})
const subscriptions = new MobileWebCapabilitySubscriptions({
...sender.subscriptionPosts(),
browserAuthority: new MobileWebBrowserAuthority(randomBytes),
nativeChatAuthority: new MobileWebNativeChatAuthority(randomBytes),
workspaceAuthority: new MobileWebWorkspaceAuthority(randomBytes)