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 <alpha-eng@stably.ai>
Co-authored-by: m4air <m4air@m4airs-Air.localdomain>
This commit is contained in:
Neil
2026-09-26 15:05:30 -07:00
committed by GitHub
co-authored by OrcaWin m4air
parent 1554f15b6c
commit 3ab183d770
7 changed files with 636 additions and 31 deletions
+84 -19
View File
@@ -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",
@@ -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)
})
})
@@ -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,<title>WebRTC egress probe</title>')
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)
})
`
@@ -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<typeof subscribeRuntimeEnvironmentFromPreload>[0]
type Args = Parameters<typeof subscribeRuntimeEnvironmentFromPreload>[1]
type Callbacks = Parameters<typeof subscribeRuntimeEnvironmentFromPreload>[2]
let listener: Parameters<Ipc['on']>[1] | undefined
const pending: ReturnType<typeof Promise.withResolvers<unknown>>[] = []
const issueRead = Promise.withResolvers<null>()
const ipc: Ipc = {
invoke: vi.fn((channel) => {
if (channel !== 'runtimeEnvironments:subscribe') {
return Promise.resolve()
}
const setup = Promise.withResolvers<unknown>()
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()
})
@@ -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: () => {
@@ -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<typeof vi.fn>
}[] = []
const onEvent = vi.fn()
const manager = createRuntimeClientEventsSync({
getDesiredEnvironmentIds: () => desired,
getSubscriptionKey: (id) => (id === 'A' ? key : id),
subscribe: (environmentId, notify) => {
const pending = Promise.withResolvers<RuntimeClientEventSubscriptionHandle>()
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<void> {
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<typeof Promise.withResolvers<RuntimeClientEventSubscriptionHandle>>
unsubscribe: ReturnType<typeof vi.fn<() => 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<RuntimeClientEventSubscriptionHandle>()
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<RuntimeClientEventSubscriptionHandle>()
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()
})
})
@@ -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<RuntimeClientEventSubscriptionHandle>
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<string, { key: string; unsubscribe: () => void }>()
const pending = new Map<string, { key: string; generation: number }>()
type SubscriptionToken = { key: string; generation: number }
const subscriptions = new Map<
string,
{ key: string; token: SubscriptionToken; unsubscribe: () => void }
>()
const pending = new Map<string, SubscriptionToken>()
const retryTimers = new Map<string, ReturnType<typeof setTimeout>>()
const consecutiveFailures = new Map<string, { key: string; count: number }>()
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) &&