From 3ab183d770e884cd90565799c12f8d4eca8b3ce9 Mon Sep 17 00:00:00 2001 From: Neil <4138956+nwparker@users.noreply.github.com> Date: Sat, 26 Sep 2026 15:05:30 -0700 Subject: [PATCH] fix: stop stale runtime event work after renderer cleanup (#23066) * fix: stop stale runtime event work after renderer cleanup * fix(i18n): restore diff note draft catalog entries * test(browser): report the phase of WebRTC probe timeouts * fix: preserve runtime subscription ownership during nested sync * test: isolate hourly identity inputs from release version changes --------- Co-authored-by: OrcaWin Co-authored-by: m4air --- config/reliability-gates.jsonc | 103 ++++-- config/scripts/hourly-build-version.test.mjs | 10 + ...owser-route-webrtc-egress.electron.test.ts | 25 +- ...untime-client-ipc-bridge-ownership.test.ts | 179 +++++++++++ .../ipc-events/runtime-client-ipc-bridge.ts | 5 +- ...ntime-client-events-sync-ownership.test.ts | 297 ++++++++++++++++++ .../src/hooks/runtime-client-events-sync.ts | 48 ++- 7 files changed, 636 insertions(+), 31 deletions(-) create mode 100644 src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge-ownership.test.ts create mode 100644 src/renderer/src/hooks/runtime-client-events-sync-ownership.test.ts diff --git a/config/reliability-gates.jsonc b/config/reliability-gates.jsonc index bb2ea73d9a9..63b3e99f45a 100644 --- a/config/reliability-gates.jsonc +++ b/config/reliability-gates.jsonc @@ -22,25 +22,10 @@ "preview settings updates", "preview box fit and grid claim" ], - "platforms": [ - "macos", - "linux", - "windows" - ], - "providers": [ - "local", - "daemon", - "ssh", - "wsl", - "remote-runtime" - ], - "coveredPlatforms": [ - "macos" - ], - "coveredProviders": [ - "local", - "remote-runtime" - ], + "platforms": ["macos", "linux", "windows"], + "providers": ["local", "daemon", "ssh", "wsl", "remote-runtime"], + "coveredPlatforms": ["macos"], + "coveredProviders": ["local", "remote-runtime"], "coverageNotes": "Mounted production component and real settings action with cloned IPC replies, inert xterm/input adapters, and simulated metrics. Local/remote-qualified renderer identities covered; no physical provider, native geometry, focus, IME or rendered app claim. Subsequent hidden/offscreen Electron component proof uses real xterm and synthetic snapshots, preserving the rendered terminal across unrelated updates and measuring equivalent fallback scale/fit requests after a font-size change.", "motivatingLinks": [ "https://github.com/stablyai/orca/blob/main/src/renderer/src/components/dashboard-popout/AgentTerminalPreview.tsx", @@ -17734,6 +17719,86 @@ ], "demotionRule": "Demote or quarantine if the gate flakes once without a product bug or harness bug filed to the owner." }, + { + "id": "runtime-events.renderer-subscription-ownership", + "title": "Runtime events stay owned by their renderer subscription", + "maturity": "experimental", + "protection": "partial", + "owner": "runtime-platform", + "layer": "ipc-contract", + "surfaces": [ + "renderer runtime events", + "remote runtime reconnect replay", + "renderer Linear refresh" + ], + "platforms": ["macos", "linux", "windows"], + "providers": ["remote-runtime"], + "coveredPlatforms": ["macos"], + "coveredProviders": ["remote-runtime"], + "coverageNotes": "Actual renderer manager, bridge, store actions, adapter and preload dispatcher with synthetic IPC boundaries on macOS. No rendered/native lifecycle, live peer or real provider network claim.", + "motivatingLinks": [ + "https://github.com/stablyai/orca/blob/main/src/renderer/src/hooks/runtime-client-events-sync.ts" + ], + "invariant": "Only the current subscription attempt can cause renderer event work or reconnect recovery. Current early frames and other hosts survive; canceled late setup handles dispose.", + "oracle": "Hold setup, clean up the bridge, deliver100 events and a replay, and require zero provider read dispatches/cache publications/recovery/re-subscriptions. Require initial/live delivery, same-ID replacement, rekey, stop during synchronous setup, sibling host survival and late disposal.", + "commands": [ + "ORCA_BACKGROUND_LAUNCH=1 pnpm exec vitest run --config config/vitest.config.ts src/renderer/src/hooks/runtime-client-events-sync-ownership.test.ts src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge-ownership.test.ts src/renderer/src/hooks/runtime-client-events-sync.test.ts src/renderer/src/hooks/ipc-events/runtime-reconnect-host-status.test.ts src/renderer/src/runtime/runtime-client-events.test.ts src/preload/runtime-environment-subscriptions.test.ts" + ], + "testFiles": [ + "src/renderer/src/hooks/runtime-client-events-sync-ownership.test.ts", + "src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge-ownership.test.ts", + "src/renderer/src/hooks/runtime-client-events-sync.test.ts", + "src/renderer/src/hooks/ipc-events/runtime-reconnect-host-status.test.ts", + "src/renderer/src/runtime/runtime-client-events.test.ts", + "src/preload/runtime-environment-subscriptions.test.ts" + ], + "assertionRefs": [ + { + "file": "src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge-ownership.test.ts", + "assertions": [ + "does no Linear read dispatch or cache publication after cleanup while setup is pending", + "does no replay recovery or resubscription after cleanup", + "accepts a fresh bridge early frame while rejecting the previous bridge frame" + ] + } + ], + "evidenceRuns": [ + { + "date": "2026-09-25", + "runner": "local", + "platform": "macos", + "command": "ORCA_BACKGROUND_LAUNCH=1 pnpm exec vitest run --config config/vitest.config.ts src/renderer/src/hooks/runtime-client-events-sync-ownership.test.ts src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge-ownership.test.ts src/renderer/src/hooks/runtime-client-events-sync.test.ts src/renderer/src/hooks/ipc-events/runtime-reconnect-host-status.test.ts src/renderer/src/runtime/runtime-client-events.test.ts src/preload/runtime-environment-subscriptions.test.ts", + "result": "passed", + "durationSeconds": 2.12, + "summary": "41 tests /6 suites pass;13 new regressions have12 failures /1 pass against original production source." + } + ], + "runtimeBudget": { + "p95Seconds": 15, + "scope": "Focused deterministic renderer/preload suites; p95 not established." + }, + "flakeHistory": { + "status": "not-started", + "evidence": "Local validation only; no CI soak." + }, + "redGreenEvidence": { + "status": "complete", + "evidence": "12/13 new tests fail original source; candidate41 tests pass.100 stale dispatches/publications become0 each." + }, + "performanceBudget": { + "required": true, + "evidence": "No new timers, scans or IPC. Per-event O(1) owner guard adds paired median0.026ms/10k pending events and0.062ms/10k settled events. Each benchmark delivers2.2million events with no per-event desired/key reads. Synthetic timing, no rendered latency claim." + }, + "knownGaps": [ + "No native Windows/Linux/WSL, live SSH/remote or rendered app lifecycle test.", + "Already-admitted asynchronous handlers are not canceled; physical listener cleanup still waits for setup settlement.", + "No CI soak or mounted render/network/heap measurement." + ], + "promotionCriteria": [ + "Retain actual composed dispatch/publication and late-disposal oracles; add CI soak and cross-platform lifecycle evidence before promotion." + ], + "demotionRule": "Keep experimental until soak; investigate failures without suppressing ownership assertions or unexplained retries." + }, { "id": "runtime-files.watcher-process-isolation", "title": "Runtime and SSH relay watcher faults stay process-isolated without disrupting host services", diff --git a/config/scripts/hourly-build-version.test.mjs b/config/scripts/hourly-build-version.test.mjs index 861f8941cf1..dfb88c37332 100644 --- a/config/scripts/hourly-build-version.test.mjs +++ b/config/scripts/hourly-build-version.test.mjs @@ -144,4 +144,14 @@ describe('getHourlyBuildIdentity', () => { expect(identity.version).toBe('1.4.203-hourly.202609142000') expect(identity.buildNumber).toBe(5) }) + + it('keeps a newer package version as the hourly base floor', () => { + const identity = getHourlyBuildIdentity(new Date('2026-09-14T20:00:00Z'), { + packageVersion: '1.4.214', + publishedVersions: ['v1.4.202', 'v1.4.203-hourly.202609140417'], + releaseNames: ['1.4.203 • 04 • Sep 13, 9:17PM • 2ce252f'] + }) + expect(identity.version).toBe('1.4.214-hourly.202609142000') + expect(identity.buildNumber).toBe(1) + }) }) diff --git a/src/main/browser/browser-route-webrtc-egress.electron.test.ts b/src/main/browser/browser-route-webrtc-egress.electron.test.ts index 013a852e810..329f500cdd8 100644 --- a/src/main/browser/browser-route-webrtc-egress.electron.test.ts +++ b/src/main/browser/browser-route-webrtc-egress.electron.test.ts @@ -36,6 +36,12 @@ const dgram = require('node:dgram') const net = require('node:net') const os = require('node:os') const { writeFileSync } = require('node:fs') +let phase = 'app-ready' + +function enterPhase(next) { + phase = next + console.error('[webrtc-egress] ' + phase) +} function bind(socket, host) { return new Promise((resolve, reject) => { @@ -71,19 +77,24 @@ async function probe() { const tcp = net.createServer((socket) => socket.destroy()) const packets = [] udp.on('message', (message) => packets.push(message.length)) + enterPhase('bind-listeners') const [udpAddress, tcpAddress] = await Promise.all([ bind(udp, '0.0.0.0'), listen(tcp, '127.0.0.1') ]) const partition = 'persist:webrtc-egress-${protectedGuest}-' + Date.now() const routeSession = session.fromPartition(partition, { cache: false }) + enterPhase('configure-proxy') await routeSession.setProxy({ mode: 'fixed_servers', proxyRules: 'socks5://127.0.0.1:' + tcpAddress.port, proxyBypassRules: '<-loopback>' }) + enterPhase('close-connections') await routeSession.closeAllConnections() + enterPhase('resolve-proxy') const resolvedProxy = await routeSession.resolveProxy('https://example.invalid/') + enterPhase('create-window') const window = new BrowserWindow({ show: false, webPreferences: { partition, sandbox: true, nodeIntegration: false, contextIsolation: true } @@ -92,6 +103,7 @@ async function probe() { window.webContents.setWebRTCIPHandlingPolicy('disable_non_proxied_udp') } const policy = window.webContents.getWebRTCIPHandlingPolicy() + enterPhase('load-page') await window.loadURL('data:text/html,WebRTC egress probe') const target = viewerAddress() const script = \` @@ -107,8 +119,11 @@ async function probe() { peer.close() })() \` + enterPhase('renderer-webrtc') await window.webContents.executeJavaScript(script) + enterPhase('drain-packets') await new Promise((resolve) => setTimeout(resolve, 500)) + enterPhase('cleanup') window.destroy() udp.close() tcp.close() @@ -116,16 +131,22 @@ async function probe() { } async function run() { - const timeout = setTimeout(() => app.exit(2), 20000) + const timeout = setTimeout(() => { + writeFileSync(${JSON.stringify(resultPath)}, JSON.stringify({ + error: 'WebRTC egress probe timed out', phase, protectedGuest: ${protectedGuest} + })) + app.exit(2) + }, 20000) await app.whenReady() const result = await probe() + enterPhase('write-result') writeFileSync(${JSON.stringify(resultPath)}, JSON.stringify(result)) clearTimeout(timeout) app.quit() } run().catch((error) => { - writeFileSync(${JSON.stringify(resultPath)}, JSON.stringify({ error: String(error?.stack || error) })) + writeFileSync(${JSON.stringify(resultPath)}, JSON.stringify({ error: String(error?.stack || error), phase, protectedGuest: ${protectedGuest} })) app.exit(1) }) ` diff --git a/src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge-ownership.test.ts b/src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge-ownership.test.ts new file mode 100644 index 00000000000..cd1b231da18 --- /dev/null +++ b/src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge-ownership.test.ts @@ -0,0 +1,179 @@ +import { afterEach, beforeEach, expect, it, vi } from 'vitest' +import { useAppStore } from '@/store' +import { subscribeRuntimeEnvironmentFromPreload } from '../../../../preload/runtime-environment-subscriptions' +import { tagRuntimeSubscriptionReplayResponse } from '../../../../shared/runtime-subscription-replay' +import { createCompatibleRuntimeStatusResponse } from '@/runtime/runtime-compatibility-test-fixture' +import { registerRuntimeClientIpcBridge } from './runtime-client-ipc-bridge' + +const initialState = useAppStore.getState() + +beforeEach(() => vi.useFakeTimers()) +afterEach(() => { + useAppStore.setState(initialState, true) + vi.useRealTimers() + vi.unstubAllGlobals() +}) + +function createHarness() { + type Ipc = Parameters[0] + type Args = Parameters[1] + type Callbacks = Parameters[2] + let listener: Parameters[1] | undefined + const pending: ReturnType>[] = [] + const issueRead = Promise.withResolvers() + const ipc: Ipc = { + invoke: vi.fn((channel) => { + if (channel !== 'runtimeEnvironments:subscribe') { + return Promise.resolve() + } + const setup = Promise.withResolvers() + pending.push(setup) + return setup.promise + }), + send: vi.fn(), + on: (_channel, callback) => { + listener = callback + }, + removeListener: vi.fn() + } + const getIssue = vi.fn(() => issueRead.promise) + const refreshStatus = vi.fn(async () => true) + let nextId = 0 + vi.stubGlobal('window', { + api: { + runtimeEnvironments: { + subscribe: (args: Args, callbacks: Callbacks) => + subscribeRuntimeEnvironmentFromPreload(ipc, args, callbacks, () => `sub-${nextId++}`), + call: vi.fn(async () => ({ id: 'r', ok: true, result: [] })) + }, + linear: { getIssue } + } + }) + const status = createCompatibleRuntimeStatusResponse() + if (!status.ok) { + throw new Error('expected valid runtime status fixture') + } + useAppStore.setState({ + settings: null, + runtimeEnvironments: [ + { + id: 'host-a', + name: 'Host A', + createdAt: 1, + updatedAt: 1, + lastUsedAt: null, + runtimeId: null, + endpoints: [ + { id: 'ws-a', kind: 'websocket', label: 'WS', endpoint: 'ws://example.invalid' } + ], + preferredEndpointId: 'ws-a' + } + ], + runtimeStatusByEnvironmentId: new Map([['host-a', { status: status.result, checkedAt: 1 }]]), + linearIssueCache: {}, + linearSearchCache: {}, + linearListCache: {}, + linearProjectIssueCache: {}, + linearCustomViewIssueCache: {}, + checkLinearConnection: vi.fn(async () => {}), + refreshRuntimeEnvironmentStatus: refreshStatus + }) + const starts: (() => void)[] = [] + const start = (): (() => void) => { + const unsubs: (() => void)[] = [] + const unsubscribeStore = registerRuntimeClientIpcBridge(unsubs, { + worktreeChangeRefreshQueue: { enqueue: vi.fn(), dispose: vi.fn() }, + activateNotifiedWorktree: vi.fn(async () => {}) + }) + const stop = (): void => { + unsubscribeStore() + unsubs.forEach((unsubscribe) => unsubscribe()) + } + starts.push(stop) + return stop + } + return { + start, + getIssue, + refreshStatus, + pending, + ipc, + emit: (index: number, replay = false) => { + const response = { + id: 'r', + ok: true as const, + _meta: { runtimeId: 'remote-runtime' }, + result: { + type: 'linearLinkedIssueUpdated', + identifier: 'ISSUE-1', + workspaceId: 'workspace-a' + } + } + listener?.(null, { + subscriptionId: `sub-${index}`, + type: 'response', + response: replay ? tagRuntimeSubscriptionReplayResponse(response) : response + }) + }, + finish: async () => { + starts.forEach((stop) => stop()) + pending.forEach((setup, index) => + setup.resolve({ subscriptionId: `sub-${index}`, requestId: 'r' }) + ) + issueRead.resolve(null) + for (let index = 0; index < 30; index += 1) { + await Promise.resolve() + } + } + } +} + +it('does no Linear read dispatch or cache publication after cleanup while setup is pending', async () => { + const h = createHarness() + let publications = 0 + const stopCounting = useAppStore.subscribe((state, previous) => { + if (state.linearIssueCache !== previous.linearIssueCache) { + publications += 1 + } + }) + try { + h.start()() + for (let index = 0; index < 100; index += 1) { + h.emit(0) + } + expect(h.getIssue).not.toHaveBeenCalled() + expect(publications).toBe(0) + } finally { + stopCounting() + await h.finish() + } + expect(h.ipc.removeListener).toHaveBeenCalledOnce() +}) + +it('does no replay recovery or resubscription after cleanup', async () => { + const h = createHarness() + try { + h.start()() + h.emit(0, true) + expect(h.refreshStatus).not.toHaveBeenCalled() + expect(h.pending).toHaveLength(1) + expect(h.getIssue).not.toHaveBeenCalled() + } finally { + await h.finish() + } +}) + +it('accepts a fresh bridge early frame while rejecting the previous bridge frame', async () => { + const h = createHarness() + try { + h.start()() + h.start() + h.emit(0) + h.emit(1) + expect(h.getIssue).toHaveBeenCalledOnce() + expect(h.getIssue).toHaveBeenCalledWith({ id: 'ISSUE-1', workspaceId: 'workspace-a' }) + } finally { + await h.finish() + } + expect(h.ipc.removeListener).toHaveBeenCalledOnce() +}) diff --git a/src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge.ts b/src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge.ts index a5db87bbce1..ae014d70230 100644 --- a/src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge.ts +++ b/src/renderer/src/hooks/ipc-events/runtime-client-ipc-bridge.ts @@ -137,7 +137,7 @@ export function registerRuntimeClientIpcBridge( const runtimeClientEventsSync = createRuntimeClientEventsSync({ getDesiredEnvironmentIds: () => getRuntimeClientEventEnvironmentIds(useAppStore.getState()), getSubscriptionKey: (environmentId) => buildRuntimeClientEventEnvironmentKey([environmentId]), - subscribe: (environmentId, onEvent, onError) => { + subscribe: (environmentId, onEvent, onError, isCurrent) => { const sshGeneration = getEnvironmentSshStateGeneration(environmentId) const runtimeGeneration = getRuntimeEnvironmentConnectionGeneration(environmentId) const runtimeRevision = getRuntimeEnvironmentRevision(environmentId) @@ -154,6 +154,9 @@ export function registerRuntimeClientIpcBridge( }, onError, () => { + if (!isCurrent()) { + return + } invalidateRuntimeClientEventReplay({ getSshStateReference: () => useAppStore.getState().sshStateByEnvironment, refreshRuntimeStatus: () => { diff --git a/src/renderer/src/hooks/runtime-client-events-sync-ownership.test.ts b/src/renderer/src/hooks/runtime-client-events-sync-ownership.test.ts new file mode 100644 index 00000000000..8383c863781 --- /dev/null +++ b/src/renderer/src/hooks/runtime-client-events-sync-ownership.test.ts @@ -0,0 +1,297 @@ +import { describe, expect, it, vi } from 'vitest' +import type { RuntimeClientEvent } from '../../../shared/runtime-client-events' +import { + createRuntimeClientEventsSync, + type RuntimeClientEventSubscriptionHandle +} from './runtime-client-events-sync' + +function makeHarness() { + let desired = ['A', 'B'] + let key = 'A:1' + const records: { + environmentId: string + emit: () => void + resolve: () => void + reject: () => void + unsubscribe: ReturnType + }[] = [] + const onEvent = vi.fn() + const manager = createRuntimeClientEventsSync({ + getDesiredEnvironmentIds: () => desired, + getSubscriptionKey: (id) => (id === 'A' ? key : id), + subscribe: (environmentId, notify) => { + const pending = Promise.withResolvers() + const unsubscribe = vi.fn() + const event: RuntimeClientEvent = { type: 'reposChanged' } + records.push({ + environmentId, + emit: () => notify(event), + resolve: () => pending.resolve({ unsubscribe }), + reject: () => pending.reject(new Error('setup failed')), + unsubscribe + }) + notify(event) + return pending.promise + }, + onEvent + }) + return { + manager, + onEvent, + records, + setDesired: (next: string[]) => { + desired = next + }, + rekey: () => { + key = 'A:2' + } + } +} + +async function settle(): Promise { + for (let index = 0; index < 4; index += 1) { + await Promise.resolve() + } +} + +describe('runtime event subscription ownership', () => { + it.each([{ desired: ['A'] }, { desired: ['A', 'B'] }])( + 'keeps the nested desired set after an initial frame replaces $desired', + async ({ desired }) => { + const h = makeHarness() + h.setDesired(desired) + h.onEvent.mockImplementation((id: string) => { + if (id === 'A') { + h.setDesired(['C']) + h.manager.sync() + } + }) + h.manager.sync() + expect(h.records.map(({ environmentId }) => environmentId)).toEqual(['A', 'C']) + h.onEvent.mockClear() + h.records.forEach((record) => record.emit()) + expect(h.onEvent.mock.calls.map(([id]) => id)).toEqual(['C']) + h.records.forEach((record) => record.resolve()) + await settle() + expect(h.records[0].unsubscribe).toHaveBeenCalledOnce() + expect(h.records[1].unsubscribe).not.toHaveBeenCalled() + h.manager.sync() + expect(h.records).toHaveLength(2) + h.manager.stop() + expect(h.records[1].unsubscribe).toHaveBeenCalledOnce() + } + ) + + it('keeps a retained owner when unsubscribe replaces an outer empty desired set', async () => { + const h = makeHarness() + h.manager.sync() + h.records.forEach((record) => record.resolve()) + await settle() + h.records[0].unsubscribe.mockImplementation(() => { + h.setDesired(['B', 'C']) + h.manager.sync() + }) + h.setDesired([]) + h.manager.sync() + expect(h.records[0].unsubscribe).toHaveBeenCalledOnce() + expect(h.records[1].unsubscribe).not.toHaveBeenCalled() + h.onEvent.mockClear() + h.records.forEach((record) => record.emit()) + expect(h.onEvent.mock.calls.map(([id]) => id)).toEqual(['B', 'C']) + h.records[2].resolve() + await settle() + expect(h.records[2].unsubscribe).not.toHaveBeenCalled() + h.manager.stop() + h.records.forEach((record) => expect(record.unsubscribe).toHaveBeenCalledOnce()) + }) + + it('detaches stopped subscriptions before unsubscribe starts a new owner', async () => { + const h = makeHarness() + h.manager.sync() + h.records.forEach((record) => record.resolve()) + await settle() + h.records[0].unsubscribe.mockImplementation(() => { + h.setDesired(['C']) + h.manager.sync() + }) + h.manager.stop() + expect(h.records[0].unsubscribe).toHaveBeenCalledOnce() + expect(h.records[1].unsubscribe).toHaveBeenCalledOnce() + h.onEvent.mockClear() + h.records.forEach((record) => record.emit()) + expect(h.onEvent.mock.calls.map(([id]) => id)).toEqual(['C']) + h.records[2].resolve() + await settle() + expect(h.records[2].unsubscribe).not.toHaveBeenCalled() + h.manager.stop() + expect(h.records[2].unsubscribe).toHaveBeenCalledOnce() + }) + + it('preserves a new owner started synchronously by the last initial frame', async () => { + let desired = ['A'] + const records: { + emit: () => void + setup: ReturnType> + unsubscribe: ReturnType void>> + }[] = [] + const onEvent = vi.fn((id: string) => { + if (id === 'A') { + manager.stop() + desired = ['B'] + manager.sync() + } + }) + const manager = createRuntimeClientEventsSync({ + getDesiredEnvironmentIds: () => desired, + subscribe: (_id, notify) => { + const setup = Promise.withResolvers() + const unsubscribe = vi.fn() + const emit = (): void => notify({ type: 'reposChanged' }) + records.push({ emit, setup, unsubscribe }) + emit() + return setup.promise + }, + onEvent + }) + manager.sync() + onEvent.mockClear() + records[0].emit() + records[1].emit() + expect(onEvent.mock.calls.map(([id]) => id)).toEqual(['B']) + records.forEach(({ setup, unsubscribe }) => setup.resolve({ unsubscribe })) + await settle() + expect(records[0].unsubscribe).toHaveBeenCalledOnce() + expect(records[1].unsubscribe).not.toHaveBeenCalled() + manager.stop() + }) + + it('stops starting further subscriptions when an initial frame synchronously stops the owner', async () => { + const pending = Promise.withResolvers() + const unsubscribe = vi.fn() + const subscribe = vi.fn((_id: string, notify: (event: RuntimeClientEvent) => void) => { + notify({ type: 'reposChanged' }) + return pending.promise + }) + const manager = createRuntimeClientEventsSync({ + getDesiredEnvironmentIds: () => ['A', 'B'], + subscribe, + onEvent: () => manager.stop() + }) + manager.sync() + expect(subscribe).toHaveBeenCalledOnce() + pending.resolve({ unsubscribe }) + await settle() + expect(unsubscribe).toHaveBeenCalledOnce() + }) + + it('admits synchronous initial frames, pending frames and settled live frames', async () => { + const h = makeHarness() + h.manager.sync() + expect(h.onEvent).toHaveBeenCalledTimes(2) + h.records[0].emit() + h.records[0].resolve() + await settle() + h.records[0].emit() + expect(h.onEvent).toHaveBeenCalledTimes(4) + h.manager.stop() + }) + + it('drops pending and settled callbacks after stop and releases the late handle', async () => { + const h = makeHarness() + h.manager.sync() + h.records[1].resolve() + await settle() + h.manager.stop() + h.onEvent.mockClear() + h.records.forEach((record) => record.emit()) + expect(h.onEvent).not.toHaveBeenCalled() + h.records[0].resolve() + await settle() + h.records.forEach((record) => expect(record.unsubscribe).toHaveBeenCalledOnce()) + }) + + it.each([false, true])( + 'drops a replaced owner while another host stays live (settled=%s)', + async (settled) => { + const h = makeHarness() + h.manager.sync() + h.records[1].resolve() + if (settled) { + h.records[0].resolve() + } + await settle() + h.setDesired(['B']) + h.manager.sync() + h.setDesired(['A', 'B']) + h.manager.sync() + h.onEvent.mockClear() + h.records[0].emit() + h.records[1].emit() + h.records[2].emit() + expect(h.onEvent.mock.calls.map(([id]) => id)).toEqual(['B', 'A']) + h.records[0].resolve() + h.records[2].resolve() + await settle() + expect(h.records[0].unsubscribe).toHaveBeenCalledOnce() + expect(h.records[2].unsubscribe).not.toHaveBeenCalled() + h.manager.stop() + } + ) + + it('rejects callbacks from the old key and admits its replacement before setup completes', async () => { + const h = makeHarness() + h.manager.sync() + h.rekey() + h.manager.sync() + h.onEvent.mockClear() + h.records[0].emit() + h.records[2].emit() + expect(h.onEvent).toHaveBeenCalledTimes(1) + h.records[0].resolve() + h.records[2].resolve() + await settle() + expect(h.records[0].unsubscribe).toHaveBeenCalledOnce() + h.manager.stop() + }) + + it('permits restarting after stop without reviving the old callback', () => { + const h = makeHarness() + h.manager.sync() + h.manager.stop() + h.manager.sync() + h.onEvent.mockClear() + h.records[0].emit() + h.records[2].emit() + expect(h.onEvent).toHaveBeenCalledTimes(1) + h.manager.stop() + }) + + it('drops callbacks from a rejected setup while retry is waiting', async () => { + const h = makeHarness() + const warning = vi.spyOn(console, 'warn').mockImplementation(() => {}) + try { + h.manager.sync() + h.records[0].reject() + await settle() + h.onEvent.mockClear() + h.records[0].emit() + expect(h.onEvent).not.toHaveBeenCalled() + } finally { + h.manager.stop() + warning.mockRestore() + } + }) + + it('revokes an owner before its unsubscribe callback can reenter delivery', async () => { + const h = makeHarness() + h.manager.sync() + h.records[0].resolve() + await settle() + h.records[0].unsubscribe.mockImplementation(h.records[0].emit) + h.onEvent.mockClear() + h.setDesired(['B']) + h.manager.sync() + expect(h.onEvent).not.toHaveBeenCalled() + h.manager.stop() + }) +}) diff --git a/src/renderer/src/hooks/runtime-client-events-sync.ts b/src/renderer/src/hooks/runtime-client-events-sync.ts index 6388c35e3ac..ae44fe78907 100644 --- a/src/renderer/src/hooks/runtime-client-events-sync.ts +++ b/src/renderer/src/hooks/runtime-client-events-sync.ts @@ -13,7 +13,9 @@ export type RuntimeClientEventsSyncDeps = { subscribe: ( environmentId: string, onEvent: (event: RuntimeClientEvent) => void, - onError: (error: unknown) => void + onError: (error: unknown) => void, + /** Use for subscription-side recovery that runs before event delivery. */ + isCurrent: () => boolean ) => Promise onEvent: (environmentId: string, event: RuntimeClientEvent) => void /** Base retry delay; doubles per consecutive failure up to retryMaxDelayMs. */ @@ -49,8 +51,12 @@ export type RuntimeClientEventsSync = { export function createRuntimeClientEventsSync( deps: RuntimeClientEventsSyncDeps ): RuntimeClientEventsSync { - const subscriptions = new Map void }>() - const pending = new Map() + type SubscriptionToken = { key: string; generation: number } + const subscriptions = new Map< + string, + { key: string; token: SubscriptionToken; unsubscribe: () => void } + >() + const pending = new Map() const retryTimers = new Map>() const consecutiveFailures = new Map() const retryDelayMs = deps.retryDelayMs ?? 1_000 @@ -58,6 +64,7 @@ export function createRuntimeClientEventsSync( const random = deps.random ?? Math.random const getSubscriptionKey = deps.getSubscriptionKey ?? ((environmentId: string) => environmentId) let generation = 0 + let syncInvocation = 0 const clearRetryTimer = (environmentId: string): void => { const retryTimer = retryTimers.get(environmentId) @@ -109,9 +116,7 @@ export function createRuntimeClientEventsSync( const stop = (): void => { generation += 1 - for (const subscription of subscriptions.values()) { - subscription.unsubscribe() - } + const stoppedSubscriptions = [...subscriptions.values()] subscriptions.clear() pending.clear() for (const retryTimer of retryTimers.values()) { @@ -119,9 +124,14 @@ export function createRuntimeClientEventsSync( } retryTimers.clear() consecutiveFailures.clear() + for (const subscription of stoppedSubscriptions) { + subscription.unsubscribe() + } } const sync = (): void => { + const syncGeneration = generation + const currentSyncInvocation = ++syncInvocation const desiredIds = new Set(deps.getDesiredEnvironmentIds()) for (const environmentId of retryTimers.keys()) { if (desiredIds.has(environmentId)) { @@ -136,14 +146,20 @@ export function createRuntimeClientEventsSync( } for (const [environmentId, subscription] of subscriptions) { + if (syncGeneration !== generation || currentSyncInvocation !== syncInvocation) { + return + } if (desiredIds.has(environmentId) && subscription.key === getSubscriptionKey(environmentId)) { continue } - subscription.unsubscribe() subscriptions.delete(environmentId) + subscription.unsubscribe() } for (const environmentId of desiredIds) { + if (syncGeneration !== generation || currentSyncInvocation !== syncInvocation) { + return + } const subscriptionKey = getSubscriptionKey(environmentId) const pendingSubscription = pending.get(environmentId) if (pendingSubscription && pendingSubscription.key !== subscriptionKey) { @@ -159,13 +175,23 @@ export function createRuntimeClientEventsSync( const subscribeGeneration = generation const pendingSubscriptionToken = { key: subscriptionKey, generation: subscribeGeneration } pending.set(environmentId, pendingSubscriptionToken) + const isCurrent = (): boolean => + subscribeGeneration === generation && + (pending.get(environmentId) === pendingSubscriptionToken || + subscriptions.get(environmentId)?.token === pendingSubscriptionToken) void deps .subscribe( environmentId, - (event) => deps.onEvent(environmentId, event), + (event) => { + // Initial frames arrive before setup settles; only the owning attempt may deliver. + if (isCurrent()) { + deps.onEvent(environmentId, event) + } + }, (error) => { console.warn('[runtime-client-events] subscription error:', error) - } + }, + isCurrent ) .then((subscription) => { const isCurrentPending = pending.get(environmentId) === pendingSubscriptionToken @@ -192,6 +218,7 @@ export function createRuntimeClientEventsSync( consecutiveFailures.delete(environmentId) subscriptions.set(environmentId, { key: subscriptionKey, + token: pendingSubscriptionToken, unsubscribe: subscription.unsubscribe }) }) @@ -224,6 +251,9 @@ export function createRuntimeClientEventsSync( }) } + if (syncGeneration !== generation || currentSyncInvocation !== syncInvocation) { + return + } for (const [environmentId, pendingSubscription] of pending) { if ( desiredIds.has(environmentId) &&