Revert "fix: recover paired structured session mirror on host swap"

This reverts commit 81bfca0007.
This commit is contained in:
Merge Sim
2026-09-04 01:09:28 -07:00
parent 81bfca0007
commit d54189408f
3 changed files with 37 additions and 104 deletions
@@ -34,7 +34,6 @@ const PRIMARY_GROUP = 'primary-group'
const SECONDARY_GROUP = 'secondary-group'
afterEach(() => {
delete (window as unknown as { __ORCA_WEB_CLIENT__?: boolean }).__ORCA_WEB_CLIENT__
resetLocalStructuredSessionVersionForTests()
resetWebSessionFocusIntentForTests()
resetWebSessionTabsSnapshotFreshnessForTests()
@@ -234,58 +233,6 @@ describe('local structured session tab projection', () => {
}
})
it('refreshes host capabilities and reconnects when the paired client closes a subscription', async () => {
vi.useFakeTimers()
const priorApi = window.api
;(window as unknown as { __ORCA_WEB_CLIENT__?: boolean }).__ORCA_WEB_CLIENT__ = true
const onCloseCallbacks: (() => void)[] = []
const subscribe = vi.fn(
async (
_args: unknown,
callbackOrCallbacks: ((response: unknown) => void) | { onClose?: () => void },
lifecycle?: { onClose?: () => void }
) => {
const callbacks =
typeof callbackOrCallbacks === 'function' ? lifecycle : callbackOrCallbacks
onCloseCallbacks.push(callbacks?.onClose ?? (() => undefined))
return { unsubscribe: vi.fn(), sendBinary: vi.fn() }
}
)
const getStatus = vi
.fn()
.mockResolvedValue({ capabilities: [STRUCTURED_AGENT_SESSION_RUNTIME_CAPABILITY] })
Object.defineProperty(window, 'api', {
configurable: true,
value: {
runtime: {
getStatus,
call: vi.fn().mockResolvedValue({ ok: true, result: { snapshots: [] } })
},
runtimeEnvironments: {
subscribe
}
}
})
try {
await startLocalStructuredSessionTabsSync({
isDisposed: () => false,
setUnsubscribe: () => undefined
})
expect(subscribe).toHaveBeenCalledOnce()
// Replacing the active paired host closes the child subscription without an RPC end frame.
onCloseCallbacks[0]?.()
await vi.advanceTimersByTimeAsync(249)
expect(subscribe).toHaveBeenCalledOnce()
await vi.advanceTimersByTimeAsync(1)
await vi.waitFor(() => expect(subscribe).toHaveBeenCalledTimes(2))
expect(getStatus.mock.calls.length).toBeGreaterThanOrEqual(2)
} finally {
Object.defineProperty(window, 'api', { configurable: true, value: priorApi })
}
})
it('ignores an in-flight inventory response after toggle-off clears the mirror', async () => {
let resolveInventory: ((response: unknown) => void) | undefined
const pendingInventory = new Promise((resolve) => {
@@ -1,6 +1,5 @@
import { useEffect } from 'react'
import { useAppStore } from '../store'
import { isWebClientLocation } from '../lib/web-client-location'
import { clearLocalStructuredSessionTabs } from './local-structured-session-tabs-sync/snapshot-apply'
import { startLocalStructuredSessionTabsSync } from './local-structured-session-tabs-sync/subscription'
@@ -24,11 +23,6 @@ export function useLocalStructuredSessionTabsSync(): void {
(state) => state.workspaceSessionReady && state.terminalStartupRestorationReady
)
const enabled = useAppStore((state) => state.settings?.experimentalStructuredNativeChat === true)
// The web preload swaps its sole active runtime environment in place when a new host is paired.
// Restart the host-owned mirror so status/capability reads cannot bleed across that boundary.
const webEnvironmentId = useAppStore((state) =>
isWebClientLocation() ? (state.runtimeEnvironments[0]?.id ?? null) : null
)
useEffect(() => {
if (!ready) {
return
@@ -48,7 +42,6 @@ export function useLocalStructuredSessionTabsSync(): void {
return () => {
disposed = true
unsubscribe()
clearLocalStructuredSessionTabs()
}
}, [enabled, ready, webEnvironmentId])
}, [enabled, ready])
}
@@ -1,8 +1,6 @@
import { STRUCTURED_AGENT_SESSION_RUNTIME_CAPABILITY } from '../../../../shared/protocol-version'
import type { RuntimeMobileSessionTabsResult } from '../../../../shared/runtime-types'
import type { RuntimeRpcResponse } from '../../../../shared/runtime-rpc-envelope'
import { refreshLocalRuntimeCapabilities } from '../local-runtime-capabilities'
import { isWebClientLocation } from '../../lib/web-client-location'
import {
isCurrentLocalStructuredSessionGeneration,
localStructuredSessionGeneration
@@ -41,17 +39,6 @@ export async function startLocalStructuredSessionTabsSync(args: {
let reconnectTimer: ReturnType<typeof setTimeout> | null = null
let reconnectAttempt = 0
let activeHandle: { unsubscribe: () => void } | null = null
const handleSubscriptionTermination = (): void => {
if (!isCurrent()) {
return
}
subscriptionGeneration += 1
if (activeHandle !== null) {
activeHandle.unsubscribe()
activeHandle = null
}
scheduleSubscribeRetry()
}
const scheduleSubscribeRetry = (): void => {
if (!isCurrent() || reconnectTimer !== null) {
return
@@ -60,8 +47,7 @@ export async function startLocalStructuredSessionTabsSync(args: {
reconnectAttempt += 1
reconnectTimer = setTimeout(() => {
reconnectTimer = null
void refreshLocalRuntimeCapabilities()
.then(() => refreshLocalStructuredSessionTabs(syncGeneration))
void refreshLocalStructuredSessionTabs(syncGeneration)
.catch((error) => console.warn('[structured-session-tabs] resync failed', error))
.finally(() => {
if (isCurrent()) {
@@ -79,35 +65,42 @@ export async function startLocalStructuredSessionTabsSync(args: {
}
const generation = ++subscriptionGeneration
let handle: { unsubscribe: () => void } | null = null
const onResponse = (response: RuntimeRpcResponse<unknown>): void => {
if (!isCurrent() || generation !== subscriptionGeneration) {
return
handle = await window.api.runtime.subscribe(
{ method: 'session.tabs.subscribeAll', params: {} },
(response) => {
if (!isCurrent() || generation !== subscriptionGeneration) {
return
}
if (!response.ok) {
// A streaming RPC can terminate with an error response before its
// handle resolves; fence that generation and retry the subscription.
subscriptionGeneration += 1
handle?.unsubscribe()
if (activeHandle === handle) {
activeHandle = null
}
scheduleSubscribeRetry()
return
}
const event = response.result as SessionTabsEvent
if (event.type === 'snapshots') {
applyStructuredSessionTabSnapshots(event.snapshots)
} else if (event.type === 'snapshot' || event.type === 'updated') {
applyStructuredSessionTabSnapshots([event])
} else if (event.type === 'end' && generation === subscriptionGeneration) {
// Reattach with one refresh so a runtime-restart boundary cannot strand stale tabs.
subscriptionGeneration += 1
handle?.unsubscribe()
if (activeHandle === handle) {
activeHandle = null
}
if (reconnectTimer !== null) {
clearTimeout(reconnectTimer)
}
scheduleSubscribeRetry()
}
}
if (!response.ok) {
// A streaming RPC can terminate with an error response before its
// handle resolves; fence that generation and retry the subscription.
handleSubscriptionTermination()
return
}
const event = response.result as SessionTabsEvent
if (event.type === 'snapshots') {
applyStructuredSessionTabSnapshots(event.snapshots)
} else if (event.type === 'snapshot' || event.type === 'updated') {
applyStructuredSessionTabSnapshots([event])
} else if (event.type === 'end' && generation === subscriptionGeneration) {
// Reattach with one refresh so a runtime-restart boundary cannot strand stale tabs.
handleSubscriptionTermination()
}
}
handle = isWebClientLocation()
? await window.api.runtimeEnvironments.subscribe(
{ selector: 'active', method: 'session.tabs.subscribeAll', params: {} },
{ onResponse, onClose: handleSubscriptionTermination }
)
: await window.api.runtime.subscribe(
{ method: 'session.tabs.subscribeAll', params: {} },
onResponse
)
)
if (!isCurrent() || generation !== subscriptionGeneration) {
handle.unsubscribe()
} else {