diff --git a/mobile/src/transport/rpc-client-session-tabs-stream.ts b/mobile/src/transport/rpc-client-session-tabs-stream.ts new file mode 100644 index 00000000000..f98a7586921 --- /dev/null +++ b/mobile/src/transport/rpc-client-session-tabs-stream.ts @@ -0,0 +1,34 @@ +import { isSnapshotResult } from './rpc-subscription-result-shapes' + +type SessionTabsStreamState = { + method: string + sent?: boolean + receivedSnapshot?: boolean + cancelled?: boolean +} + +/** Older hosts register a tabs stream only as they emit the first snapshot, so an earlier unsubscribe + * finds nothing there. Newer hosts register on arrival and still send that snapshot, so the hold only delays. */ +function awaitsRegistration(stream: SessionTabsStreamState): boolean { + return ( + stream.method === 'session.tabs.subscribe' && stream.sent === true && !stream.receivedSnapshot + ) +} + +/** Cancels a stream still awaiting its first snapshot; that snapshot then triggers the unsubscribe. */ +export function holdUnsubscribe(stream: SessionTabsStreamState): boolean { + if (!awaitsRegistration(stream)) { + return false + } + stream.cancelled = true + return true +} + +/** Returns true when `result` is the tabs stream's snapshot, i.e. the host has registered it. */ +export function recordSnapshot(stream: SessionTabsStreamState, result: unknown): boolean { + if (stream.method !== 'session.tabs.subscribe' || !isSnapshotResult(result)) { + return false + } + stream.receivedSnapshot = true + return true +} diff --git a/mobile/src/transport/rpc-client-stream-registry.test.ts b/mobile/src/transport/rpc-client-stream-registry.test.ts index 6afcd8a7ef3..ce210333618 100644 --- a/mobile/src/transport/rpc-client-stream-registry.test.ts +++ b/mobile/src/transport/rpc-client-stream-registry.test.ts @@ -79,6 +79,104 @@ describe('RpcClientStreamRegistry', () => { ]) }) + it('unsubscribes a session tabs stream by its own request id', () => { + const { registry, sent } = createRegistry() + const disposeOlder = registry.subscribe( + 'session.tabs.subscribe', + { worktree: 'wt-1' }, + () => {} + ) + registry.subscribe('session.tabs.subscribe', { worktree: 'wt-1' }, () => {}) + const [older, newer] = sent + registry.handleResponse(streamingResponse(older!.id, { type: 'snapshot', tabs: [] })) + + disposeOlder() + + // Without the request id the host sweeps every stream for the worktree, including the newer one. + expect(sent[2]).toMatchObject({ + method: 'session.tabs.unsubscribe', + params: { worktree: 'wt-1', subscriptionId: older!.id } + }) + expect(newer!.id).not.toBe(older!.id) + expect(sent).toHaveLength(3) + }) + + it('holds a session tabs unsubscribe until the host registers the stream', () => { + const { registry, sent } = createRegistry() + const events: unknown[] = [] + const dispose = registry.subscribe('session.tabs.subscribe', { worktree: 'wt-1' }, (event) => + events.push(event) + ) + const subscribe = sent[0]! + + dispose() + // Older hosts register only as they emit the first snapshot, so an earlier unsubscribe finds nothing. + expect(sent).toHaveLength(1) + + registry.handleResponse(streamingResponse(subscribe.id, { type: 'snapshot', tabs: [] })) + + expect(sent[1]).toMatchObject({ + method: 'session.tabs.unsubscribe', + params: { worktree: 'wt-1', subscriptionId: subscribe.id } + }) + // A host that registered on arrival ends the stream; that end must not unsubscribe again. + registry.handleResponse(streamingResponse(subscribe.id, { type: 'end' })) + expect(sent).toHaveLength(2) + expect(events).toEqual([]) + expect(registry.size()).toBe(0) + }) + + it('holds a session tabs unsubscribe again after a reconnect replays the stream', () => { + const { registry, sent } = createRegistry() + const dispose = registry.subscribe('session.tabs.subscribe', { worktree: 'wt-1' }, () => {}) + const subscribe = sent[0]! + registry.handleResponse(streamingResponse(subscribe.id, { type: 'snapshot', tabs: [] })) + + registry.markForReplay() + registry.replayAfterAuthentication() + dispose() + + expect(sent.map((request) => request.method)).toEqual([ + 'session.tabs.subscribe', + 'session.tabs.subscribe' + ]) + registry.handleResponse(streamingResponse(subscribe.id, { type: 'snapshot', tabs: [] })) + expect(sent[2]).toMatchObject({ + method: 'session.tabs.unsubscribe', + params: { worktree: 'wt-1', subscriptionId: subscribe.id } + }) + }) + + it('drops a held session tabs unsubscribe when the subscribe fails or reconnects first', () => { + const { registry, sent } = createRegistry() + const events: unknown[] = [] + const disposeFailed = registry.subscribe( + 'session.tabs.subscribe', + { worktree: 'wt-1' }, + (event) => events.push(event) + ) + const disposeReplayed = registry.subscribe( + 'session.tabs.subscribe', + { worktree: 'wt-2' }, + () => {} + ) + const failed = sent[0]! + disposeFailed() + disposeReplayed() + + registry.handleResponse({ + id: failed.id, + ok: false, + error: { code: 'worktree_not_found', message: 'Worktree not found' } + }) + registry.markForReplay() + registry.replayAfterAuthentication() + + expect(sent).toHaveLength(2) + expect(events).toEqual([]) + expect(registry.size()).toBe(0) + }) + it('ends one transcript stream on dispose and leaves a sibling on the same socket (U-03)', () => { const { registry, sent } = createRegistry() const disposeFirst = registry.subscribe('agentSession.subscribe', { sessionId: 's1' }, () => {}) diff --git a/mobile/src/transport/rpc-client-stream-registry.ts b/mobile/src/transport/rpc-client-stream-registry.ts index ebfdb741997..4062e80b49b 100644 --- a/mobile/src/transport/rpc-client-stream-registry.ts +++ b/mobile/src/transport/rpc-client-stream-registry.ts @@ -14,6 +14,7 @@ import { isTerminalSubscribedResult } from './rpc-subscription-result-shapes' import { RpcClientTerminalStreamRouter } from './rpc-client-terminal-stream-router' +import * as sessionTabsStream from './rpc-client-session-tabs-stream' import type { ConnectionState, RpcResponse, RpcSuccess } from './types' export type RpcStreamingListener = (result: unknown) => void @@ -30,6 +31,7 @@ type StreamRequest = { subscriptionId?: string cancelled?: boolean sent?: boolean + receivedSnapshot?: boolean } type StreamRegistryOptions = { @@ -108,6 +110,7 @@ export class RpcClientStreamRegistry { this.pendingBrowserRequestId = null for (const [id, stream] of this.streams) { stream.sent = false + stream.receivedSnapshot = false // The id named a registration on the closed socket; the replay's ready brings the new one. stream.subscriptionId = undefined this.resetTerminalRouting(id) @@ -169,6 +172,10 @@ export class RpcClientStreamRegistry { return } const result = response.result + if (sessionTabsStream.recordSnapshot(stream, result) && stream.cancelled) { + this.dispose(response.id) + return + } if (isStreamingSubscriptionReadyResult(result)) { stream.subscriptionId = result.subscriptionId if (stream.cancelled) { @@ -215,6 +222,8 @@ export class RpcClientStreamRegistry { // Why: `requestId` names this exact request; hosts that predate it strip it and use the slot. this.sendRpc('terminal.unsubscribe', { ...params, requestId: id }) } + } else if (stream && sessionTabsStream.holdUnsubscribe(stream)) { + return } else { const unsubscribe = buildStreamUnsubscribe(stream?.method, stream?.params, id) if (unsubscribe) { diff --git a/mobile/src/transport/rpc-client.test.ts b/mobile/src/transport/rpc-client.test.ts index e88e6c220e6..ff290f72b84 100644 --- a/mobile/src/transport/rpc-client.test.ts +++ b/mobile/src/transport/rpc-client.test.ts @@ -196,12 +196,21 @@ describe('mobile rpc-client connection timeout', () => { { worktree: 'id:wt-1' }, () => {} ) + const request = sentRequest(socket, 'session.tabs.subscribe') + socket.receive( + `encrypted:${JSON.stringify({ + id: request.id, + ok: true, + streaming: true, + result: { type: 'snapshot', worktree: 'id:wt-1', tabs: [] }, + _meta: { runtimeId: 'r1' } + })}` + ) unsubscribe() - expect( - socket.sent.some((payload) => payload.includes('"method":"session.tabs.unsubscribe"')) - ).toBe(true) - expect(socket.sent.some((payload) => payload.includes('"worktree":"id:wt-1"'))).toBe(true) + expect(sentRequests(socket, 'session.tabs.unsubscribe')).toEqual([ + expect.objectContaining({ params: { worktree: 'id:wt-1', subscriptionId: request.id } }) + ]) client.close() }) diff --git a/mobile/src/transport/rpc-subscription-result-shapes.ts b/mobile/src/transport/rpc-subscription-result-shapes.ts index cf09acdb6ca..4c624658126 100644 --- a/mobile/src/transport/rpc-subscription-result-shapes.ts +++ b/mobile/src/transport/rpc-subscription-result-shapes.ts @@ -19,3 +19,7 @@ export function isStreamingSubscriptionReadyResult( typeof (value as { subscriptionId?: unknown }).subscriptionId === 'string' ) } + +export function isSnapshotResult(value: unknown): value is { type: 'snapshot' } { + return typeof value === 'object' && value !== null && 'type' in value && value.type === 'snapshot' +} diff --git a/src/main/runtime/rpc/methods/session-tabs-unsubscribe.test.ts b/src/main/runtime/rpc/methods/session-tabs-unsubscribe.test.ts index 8490181d0ea..021224bd649 100644 --- a/src/main/runtime/rpc/methods/session-tabs-unsubscribe.test.ts +++ b/src/main/runtime/rpc/methods/session-tabs-unsubscribe.test.ts @@ -1,5 +1,6 @@ import { describe, expect, it, vi } from 'vitest' import type { OrcaRuntimeService } from '../../orca-runtime' +import { RuntimeSubscriptionRegistry } from '../../runtime-subscription-registry' import { RpcDispatcher } from '../dispatcher' import { SESSION_TAB_METHODS } from './session-tabs' @@ -61,11 +62,61 @@ describe('session tab unsubscribe RPC methods', () => { expect(cleanupSubscription).toHaveBeenCalledWith('session.tabs:conn-1:*:sub-all-1') expect(cleanupSubscriptionsByPrefix).not.toHaveBeenCalled() }) + + it('ends only the named stream when a newer one watches the same worktree', async () => { + const { dispatcher, ends } = await subscribeTwiceToOneWorktree() + + await dispatcher.dispatchStreaming( + request('session.tabs.unsubscribe', { worktree: 'id:wt-1', subscriptionId: 'sub-old' }), + vi.fn(), + { connectionId: 'conn-1' } + ) + + await vi.waitFor(() => expect(ends['sub-old']).toBe(1)) + expect(ends['sub-new']).toBe(0) + }) + + it('ends every stream for the worktree when the unsubscribe names no request', async () => { + const { dispatcher, ends } = await subscribeTwiceToOneWorktree() + + await dispatcher.dispatchStreaming( + request('session.tabs.unsubscribe', { worktree: 'id:wt-1' }), + vi.fn(), + { connectionId: 'conn-1' } + ) + + await vi.waitFor(() => expect(ends).toEqual({ 'sub-old': 1, 'sub-new': 1 })) + }) }) +async function subscribeTwiceToOneWorktree() { + const registry = new RuntimeSubscriptionRegistry() + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Partial runtime backed by the real subscription registry; it supplies every member session.tabs subscribe/unsubscribe call. + const runtime = { + ...runtimeWithCleanup(registry.cleanup.bind(registry), registry.cleanupByPrefix.bind(registry)), + registerSubscriptionCleanup: registry.register.bind(registry), + getSubscriptionRegistrationVersion: registry.getRegistrationVersion.bind(registry), + onMobileSessionTabsChanged: () => () => {} + } as unknown as OrcaRuntimeService + const dispatcher = new RpcDispatcher({ runtime, methods: SESSION_TAB_METHODS }) + const ends: Record = { 'sub-old': 0, 'sub-new': 0 } + for (const id of ['sub-old', 'sub-new']) { + await dispatcher.dispatchStreaming( + { ...request('session.tabs.subscribe', { worktree: 'id:wt-1' }), id }, + (message) => { + if (JSON.parse(message).result?.type === 'end') { + ends[id]! += 1 + } + }, + { connectionId: 'conn-1' } + ) + } + return { dispatcher, ends } +} + function runtimeWithCleanup( - cleanupSubscription: ReturnType, - cleanupSubscriptionsByPrefix = vi.fn() + cleanupSubscription: (id: string) => void, + cleanupSubscriptionsByPrefix: (prefix: string, throughVersion?: number) => void = vi.fn() ): OrcaRuntimeService { // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the unsubscribe methods reach only these members; a missing one throws and fails the test. return {