diff --git a/config/reliability-gates.jsonc b/config/reliability-gates.jsonc index 8d035576165..bb2ea73d9a9 100644 --- a/config/reliability-gates.jsonc +++ b/config/reliability-gates.jsonc @@ -126,20 +126,9 @@ "issue-command read and write", "hosted-review own-store host resolution" ], - "platforms": [ - "macos", - "linux", - "windows" - ], - "providers": [ - "local", - "ssh", - "wsl", - "remote-runtime" - ], - "coveredPlatforms": [ - "macos" - ], + "platforms": ["macos", "linux", "windows"], + "providers": ["local", "ssh", "wsl", "remote-runtime"], + "coveredPlatforms": ["macos"], "coveredProviders": [], "coverageNotes": "Actual hydration and authoritative desktop owner selection with mocked filesystem, Git, inspection and provider boundaries; synthetic local/SSH/runtime rows. No physical provider or native platform certification.", "motivatingLinks": [ @@ -150,9 +139,7 @@ "commands": [ "ORCA_BACKGROUND_LAUNCH=1 pnpm exec vitest run --config config/vitest.config.ts src/main/ipc/hooks/worktree-hook-execution-host.test.ts" ], - "testFiles": [ - "src/main/ipc/hooks/worktree-hook-execution-host.test.ts" - ], + "testFiles": ["src/main/ipc/hooks/worktree-hook-execution-host.test.ts"], "assertionRefs": [ { "file": "src/main/ipc/hooks/worktree-hook-execution-host.test.ts", @@ -201,6 +188,85 @@ ], "demotionRule": "Keep experimental until soak; investigate authority/routing failures without replacing exact-target or no-local-I/O assertions." }, + { + "id": "runtime-files.renderer-document-watch-ownership", + "title": "Filesystem watches stay owned by their renderer document", + "maturity": "experimental", + "protection": "partial", + "owner": "runtime-platform", + "layer": "ipc-contract", + "surfaces": [ + "desktop file explorer watchers", + "desktop editor file watchers", + "direct SSH file watchers" + ], + "platforms": ["macos", "linux", "windows"], + "providers": ["local", "ssh", "wsl"], + "coveredPlatforms": ["macos"], + "coveredProviders": ["local", "ssh", "wsl"], + "coverageNotes": "Actual desktop local/WSL subscription and SSH watcher controller logic under mocked filesystem/provider handles on macOS. Document events use an EventEmitter sender. Live remote/native Windows/Linux/WSL and rendered reload reattachment are unproved; PTY/process liveness and wire formats are unaffected.", + "motivatingLinks": [ + "https://github.com/stablyai/orca/blob/main/src/main/ipc/filesystem-watcher-listener-lifecycle.ts" + ], + "invariant": "A replaced or crashed renderer document retains no file watches or pending ownership and cannot recreate them from an old asynchronous continuation. Shared roots retain live sibling owners; blocked and same-document navigation retain the current document. Shutdown removes registered lifecycle listeners.", + "oracle": "Run repeated reloads/crash, pending local/SSH installation cancellation with late success, aborted-install joiners, reconnect retry snapshots, failed-removal restoration and shutdown/reopen. Require zero obsolete roots/intent/listeners, no stale setup/retry, physical handle disposal and preservation of sibling notifications and same-tick remote handoff.", + "commands": [ + "ORCA_BACKGROUND_LAUNCH=1 pnpm exec vitest run --config config/vitest.config.ts src/main/ipc/filesystem-watcher-document-lifetime.test.ts src/main/ipc/filesystem-watcher-local-unsubscribe.test.ts src/main/ipc/filesystem-watcher-remote-cancellation.test.ts src/main/ipc/filesystem-watcher-native-capacity.test.ts --reporter=dot" + ], + "testFiles": [ + "src/main/ipc/filesystem-watcher-document-lifetime.test.ts", + "src/main/ipc/filesystem-watcher-local-unsubscribe.test.ts", + "src/main/ipc/filesystem-watcher-remote-cancellation.test.ts", + "src/main/ipc/filesystem-watcher-native-capacity.test.ts" + ], + "assertionRefs": [ + { + "file": "src/main/ipc/filesystem-watcher-document-lifetime.test.ts", + "assertions": [ + "retains zero roots and listeners across 15 reloads and a renderer crash", + "does not revive a cancelled %s setup for a joiner whose document reloaded", + "does not re-arm an old provider snapshot into the same WebContents after replacement", + "does not restore a replaced sibling document after an awaited %s removal recovery", + "keeps same-document and blocked navigations alive, and removes lifecycle listeners at shutdown" + ] + } + ], + "evidenceRuns": [ + { + "date": "2026-09-25", + "runner": "local", + "platform": "macos", + "command": "ORCA_BACKGROUND_LAUNCH=1 pnpm exec vitest run --config config/vitest.config.ts src/main/ipc/filesystem-watcher-document-lifetime.test.ts src/main/ipc/filesystem-watcher-local-unsubscribe.test.ts src/main/ipc/filesystem-watcher-remote-cancellation.test.ts src/main/ipc/filesystem-watcher-native-capacity.test.ts --reporter=dot", + "result": "passed", + "durationSeconds": 1.08, + "summary": "45 tests across four focused suites pass. Broader watcher suite passes 146 plus one existing skipped; isolated publication without sibling #22995 passes 145 plus one existing skipped." + } + ], + "runtimeBudget": { + "p95Seconds": 15, + "scope": "Focused deterministic IPC-contract suites; p95 not established." + }, + "flakeHistory": { + "status": "not-started", + "evidence": "Local validation only; no CI soak." + }, + "redGreenEvidence": { + "status": "complete", + "evidence": "Seven permanent regressions fail with HEAD source substituted: cancelled joiner resurrection for local and SSH, restoration into replaced local/SSH documents, stale provider retry, shutdown listener retention, and root retention. All 13 document ownership cases pass after the change." + }, + "performanceBudget": { + "required": true, + "evidence": "15 reloads plus crash: retained local/SSH/desired roots 16→0 each; pending setup aborted; sender lifecycle listeners 1→0. Actual ownership modules with synthetic handles; no native process/heap/latency claim. No new timer, polling, process scan or provider fanout; existing root cleanup runs at document end." + }, + "knownGaps": [ + "No native Windows/Linux/WSL or live SSH execution.", + "No rendered Electron reload-reattachment/paint run or CI soak; broader watcher process isolation is a separate gate." + ], + "promotionCriteria": [ + "Retain lifecycle, count and sibling-delivery assertions and add cross-platform and rendered reload evidence before promotion." + ], + "demotionRule": "Keep experimental until soak evidence; investigate lifecycle failures without suppressing assertions or retrying unexplained failures." + }, { "id": "agent-session.journal-streaming-replay", "title": "Journal replay bounds obsolete revision memory without changing recovery", diff --git a/src/main/ipc/filesystem-watcher-canonical-root-paths.test.ts b/src/main/ipc/filesystem-watcher-canonical-root-paths.test.ts index 45c14d485e0..9aa42e3107e 100644 --- a/src/main/ipc/filesystem-watcher-canonical-root-paths.test.ts +++ b/src/main/ipc/filesystem-watcher-canonical-root-paths.test.ts @@ -1,3 +1,4 @@ +import { createWatcherSender } from './filesystem-watcher-test-sender' /* * macOS FSEvents reports OS-canonical paths (symlinks resolved, on-disk * casing). Before the rewrite at the watcher boundary those events reached the @@ -80,7 +81,7 @@ describe('local filesystem watcher canonical root paths', () => { return { unsubscribe: vi.fn() } as never }) const sendMock = vi.fn() - const sender = { isDestroyed: () => false, send: sendMock, once: vi.fn(), id: 1 } + const sender = createWatcherSender(1, sendMock) await handlers['fs:watchWorktree']({ sender }, { worktreePath }) watcherCallback!(null, events) await vi.waitFor( diff --git a/src/main/ipc/filesystem-watcher-document-lifetime.test.ts b/src/main/ipc/filesystem-watcher-document-lifetime.test.ts new file mode 100644 index 00000000000..9bc42d1ef73 --- /dev/null +++ b/src/main/ipc/filesystem-watcher-document-lifetime.test.ts @@ -0,0 +1,299 @@ +import { EventEmitter } from 'node:events' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { createDebouncedBatch } from './filesystem-watcher-batch-control' + +type WatchArgs = { worktreePath: string; connectionId?: string } +class Sender extends EventEmitter { + constructor(readonly id: number) { + super() + } + isDestroyed = () => false + send = vi.fn() +} +const { handleMock, watchRemote, createLocal, getProvider } = vi.hoisted(() => ({ + handleMock: vi.fn(), + watchRemote: vi.fn(), + createLocal: vi.fn(), + getProvider: vi.fn() +})) +vi.mock('electron', () => ({ ipcMain: { handle: handleMock } })) +vi.mock('node:fs/promises', () => ({ stat: async () => ({ isDirectory: () => true }) })) +vi.mock('./filesystem-watcher-local-events', () => ({ + createLocalWatcher: createLocal, + scheduleLocalBatchFlush: vi.fn() +})) +vi.mock('./parcel-watcher-process', () => ({ disposeWatcherProcess: vi.fn() })) +vi.mock('../providers/ssh-filesystem-dispatch', () => ({ + getSshFilesystemProvider: getProvider, + onSshFilesystemProviderRegistered: () => () => {} +})) +import { + REMOTE_WATCH_RETRY_MS, + watcherLifecycleState as state +} from './filesystem-watcher-lifecycle-state' +import { getLocalWatcherRoot, getRemoteWatcherKey } from './filesystem-watcher-paths' +import { reinstallRemoteWatchersForConnection } from './filesystem-watcher-remote-controller' +import { reinstallRemoteWatchersForConnectionCore } from './filesystem-watcher-remote-provider-rearm' +import { + registerFilesystemWatcherHandlers, + closeAllWatchers, + closeLocalWatcherForWorktreePath, + closeRemoteWatcherForWorktreePath, + restoreLocalWatcherAfterFailedRemoval, + restoreRemoteWatcherAfterFailedRemoval +} from './filesystem-watcher' + +const handlers = new Map Promise>() +function invoke(channel: string, sender: Sender, args: WatchArgs): Promise { + const handler = handlers.get(channel) + if (!handler) { + throw new Error(`Missing ${channel}`) + } + return handler({ sender }, args) +} +const watch = (sender: Sender, args: WatchArgs) => invoke('fs:watchWorktree', sender, args) +const unwatch = (sender: Sender, args: WatchArgs) => invoke('fs:unwatchWorktree', sender, args) +function deferred() { + let resolve!: (value: T) => void + const promise = new Promise((done) => { + resolve = done + }) + return { promise, resolve } +} +function localRoot(unsubscribe = vi.fn(async () => {})) { + return { + subscription: { unsubscribe }, + listeners: new Map(), + batch: createDebouncedBatch(), + rootPath: '/folder' + } +} + +beforeEach(async () => { + await closeAllWatchers() + vi.clearAllMocks() + createLocal.mockImplementation(async () => localRoot()) + watchRemote.mockImplementation(async () => vi.fn()) + getProvider.mockReturnValue({ watch: watchRemote }) + handleMock.mockImplementation((channel, handler) => handlers.set(channel, handler)) + registerFilesystemWatcherHandlers() +}) +afterEach(async () => { + await closeAllWatchers() + vi.useRealTimers() +}) + +describe('filesystem watcher renderer document ownership', () => { + it.each(['did-navigate', 'render-process-gone'])( + 'releases installed local and SSH roots on %s without closing sibling owners', + async (event) => { + const sender = new Sender(1) + const sibling = new Sender(2) + const local = { worktreePath: '/folder' } + const remote = { ...local, connectionId: 'ssh' } + const root = localRoot() + const remoteClose = vi.fn() + createLocal.mockResolvedValue(root) + watchRemote.mockResolvedValue(remoteClose) + await watch(sender, local) + await watch(sender, remote) + await watch(sibling, local) + await watch(sibling, remote) + sender.emit(event) + await Promise.resolve() + expect(root.subscription.unsubscribe).not.toHaveBeenCalled() + expect(remoteClose).not.toHaveBeenCalled() + expect([...root.listeners.keys()]).toEqual([2]) + expect([...state.desiredRemoteWatchers.values()][0].listeners.size).toBe(1) + sibling.emit(event) + await Promise.resolve() + expect(root.subscription.unsubscribe).toHaveBeenCalledTimes(1) + expect(remoteClose).toHaveBeenCalledTimes(1) + expect(state.watchedRoots.size).toBe(0) + expect(state.remoteWatchers.size).toBe(0) + expect(state.desiredRemoteWatchers.size).toBe(0) + expect(sender.eventNames()).toEqual([]) + expect(sibling.eventNames()).toEqual([]) + } + ) + + it('keeps same-document and blocked navigations alive, and removes lifecycle listeners at shutdown', async () => { + const sender = new Sender(1) + for (let cycle = 0; cycle < 3; cycle++) { + await watch(sender, { worktreePath: '/folder' }) + sender.emit('did-start-navigation') + sender.emit('did-navigate-in-page') + expect(state.watchedRoots.size).toBe(1) + for (const event of ['destroyed', 'did-navigate', 'render-process-gone']) { + expect(sender.listenerCount(event)).toBe(1) + } + await closeAllWatchers() + expect(sender.eventNames()).toEqual([]) + } + }) + + it.each(['local', 'ssh'])( + 'does not revive a cancelled %s setup for a joiner whose document reloaded', + async (kind) => { + const sender = new Sender(1) + const joiner = new Sender(2) + const args = { worktreePath: '/folder', ...(kind === 'ssh' ? { connectionId: 'ssh' } : {}) } + const install = deferred & (() => void)>() + const setup = kind === 'ssh' ? watchRemote : createLocal + setup.mockReturnValueOnce(install.promise) + const first = watch(sender, args) + await vi.waitFor(() => expect(setup).toHaveBeenCalledTimes(1)) + await unwatch(sender, args) + await Promise.resolve() + const key = + kind === 'ssh' ? getRemoteWatcherKey('ssh', '/folder') : getLocalWatcherRoot('/folder').key + const token = + kind === 'ssh' + ? state.inFlightRemoteInstalls.get(key) + : state.inFlightLocalInstalls.get(key) + expect(token?.abortController.signal.aborted).toBe(true) + const pendingJoiner = watch(joiner, args) + joiner.emit('did-navigate') + const fresh = watch(joiner, { ...args, worktreePath: '/fresh' }) + const closeLate = vi.fn(async () => {}) + install.resolve(Object.assign(closeLate, localRoot(closeLate))) + await Promise.all([first, pendingJoiner, fresh]) + expect(setup.mock.calls.filter(([path]) => path === '/folder')).toHaveLength(1) + expect(closeLate).toHaveBeenCalledTimes(1) + expect(kind === 'ssh' ? state.remoteWatchers.has(key) : state.watchedRoots.has(key)).toBe( + false + ) + expect(setup).toHaveBeenCalledTimes(2) + } + ) + + it.each(['local', 'ssh'])( + 'aborts pending %s setup on renderer crash and discards late success', + async (kind) => { + const sender = new Sender(1) + const setup = kind === 'ssh' ? watchRemote : createLocal + const install = deferred & (() => void)>() + setup.mockReturnValueOnce(install.promise) + const pending = watch(sender, { + worktreePath: '/folder', + ...(kind === 'ssh' ? { connectionId: 'ssh' } : {}) + }) + await vi.waitFor(() => expect(setup).toHaveBeenCalledTimes(1)) + const token = [ + ...(kind === 'ssh' ? state.inFlightRemoteInstalls : state.inFlightLocalInstalls).values() + ][0] + sender.emit('render-process-gone') + await Promise.resolve() + expect(token.abortController.signal.aborted).toBe(true) + const closeLate = vi.fn(async () => {}) + install.resolve(Object.assign(closeLate, localRoot(closeLate))) + await pending + expect(closeLate).toHaveBeenCalledTimes(1) + expect(state.watchedRoots.size + state.remoteWatchers.size).toBe(0) + } + ) + + it('does not re-arm stale SSH handler or retry snapshots after reload', async () => { + vi.useFakeTimers() + const sender = new Sender(1) + const args = { worktreePath: '/folder', connectionId: 'ssh' } + watchRemote.mockRejectedValueOnce(new Error('temporary unavailable')) + await watch(sender, args) + const retry = deferred<() => void>() + watchRemote.mockReturnValueOnce(retry.promise) + await vi.advanceTimersByTimeAsync(REMOTE_WATCH_RETRY_MS) + expect(watchRemote).toHaveBeenCalledTimes(2) + sender.emit('did-navigate') + // A fresh document can ask for the same root while the old retry is settling. + const fresh = watch(sender, args) + retry.resolve(vi.fn()) + await fresh + expect(state.pendingRemoteWatcherRetries.size).toBe(0) + expect(state.remoteWatcherResyncStates.size).toBe(0) + expect(watchRemote).toHaveBeenCalledTimes(2) + }) + + it('does not re-arm an old provider snapshot into the same WebContents after replacement', async () => { + const sender = new Sender(1) + const args = { worktreePath: '/folder', connectionId: 'ssh' } + await watch(sender, args) + const pending = deferred<'unavailable'>() + const dependencies = { + install: vi.fn(() => pending.promise), + requestResync: vi.fn(), + scheduleRetry: vi.fn(), + scheduleDormant: vi.fn() + } + reinstallRemoteWatchersForConnectionCore('ssh', dependencies) + sender.emit('did-navigate') + await watch(sender, args) + pending.resolve('unavailable') + await pending.promise + await Promise.resolve() + expect(dependencies.scheduleRetry).not.toHaveBeenCalled() + expect(dependencies.requestResync).not.toHaveBeenCalled() + expect(dependencies.scheduleDormant).not.toHaveBeenCalled() + }) + + it.each(['local', 'ssh'])( + 'does not restore a replaced sibling document after an awaited %s removal recovery', + async (kind) => { + const first = new Sender(1) + const sibling = new Sender(2) + const args = { worktreePath: '/folder', ...(kind === 'ssh' ? { connectionId: 'ssh' } : {}) } + await watch(first, args) + await watch(sibling, args) + await (kind === 'ssh' + ? closeRemoteWatcherForWorktreePath('ssh', '/folder') + : closeLocalWatcherForWorktreePath('/folder')) + const install = deferred & (() => void)>() + const setup = kind === 'ssh' ? watchRemote : createLocal + setup.mockReturnValueOnce(install.promise) + const restore = + kind === 'ssh' + ? restoreRemoteWatcherAfterFailedRemoval('ssh', '/folder') + : restoreLocalWatcherAfterFailedRemoval('/folder') + await vi.waitFor(() => expect(setup).toHaveBeenCalledTimes(2)) + sibling.emit('did-navigate') + install.resolve(Object.assign(vi.fn(), localRoot())) + await restore + expect(sibling.send).not.toHaveBeenCalled() + expect(first.send).toHaveBeenCalledTimes(1) + const roots = kind === 'ssh' ? state.remoteWatchers : state.watchedRoots + expect([...roots.values()][0].listeners.size).toBe(1) + } + ) + + it('retains zero roots and listeners across 15 reloads and a renderer crash', async () => { + const sender = new Sender(1) + const closeLocal = vi.fn(async () => {}) + const closeRemote = vi.fn() + createLocal.mockImplementation(async () => localRoot(closeLocal)) + watchRemote.mockResolvedValue(closeRemote) + for (let generation = 0; generation < 16; generation++) { + await watch(sender, { worktreePath: `/folder-${generation}` }) + await watch(sender, { worktreePath: `/folder-${generation}`, connectionId: 'ssh' }) + sender.emit(generation === 15 ? 'render-process-gone' : 'did-navigate') + await Promise.resolve() + expect( + state.watchedRoots.size + state.remoteWatchers.size + state.desiredRemoteWatchers.size + ).toBe(0) + expect(sender.eventNames()).toEqual([]) + } + expect(closeLocal).toHaveBeenCalledTimes(16) + expect(closeRemote).toHaveBeenCalledTimes(16) + }) + + it('clears desired SSH roots while no provider is available and does not re-arm them on reconnect', async () => { + const sender = new Sender(1) + getProvider.mockReturnValue(undefined) + await watch(sender, { worktreePath: '/folder', connectionId: 'ssh' }) + expect(state.pendingRemoteWatcherRetries.size).toBe(1) + sender.emit('render-process-gone') + getProvider.mockReturnValue({ watch: watchRemote }) + reinstallRemoteWatchersForConnection('ssh') + expect(state.pendingRemoteWatcherRetries.size).toBe(0) + expect(state.desiredRemoteWatchers.size).toBe(0) + expect(watchRemote).not.toHaveBeenCalled() + }) +}) diff --git a/src/main/ipc/filesystem-watcher-dormant-rearm.test.ts b/src/main/ipc/filesystem-watcher-dormant-rearm.test.ts index 7d27d4ca780..c4931cd7cb2 100644 --- a/src/main/ipc/filesystem-watcher-dormant-rearm.test.ts +++ b/src/main/ipc/filesystem-watcher-dormant-rearm.test.ts @@ -1,3 +1,4 @@ +import { senderEvents } from './filesystem-watcher-test-sender' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' const { handleMock, getSshFilesystemProviderMock, providerRegistrationListeners } = vi.hoisted( @@ -36,10 +37,11 @@ const DORMANT_FIRST_MS = 60_000 function createSender(id: number): { isDestroyed: () => boolean send: ReturnType + removeListener: ReturnType once: ReturnType id: number } { - return { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id } + return { isDestroyed: () => false, send: vi.fn(), ...senderEvents(), id } } describe('remote filesystem watcher dormant re-arm', () => { diff --git a/src/main/ipc/filesystem-watcher-handlers.ts b/src/main/ipc/filesystem-watcher-handlers.ts index bac74f63589..39f9ef0c184 100644 --- a/src/main/ipc/filesystem-watcher-handlers.ts +++ b/src/main/ipc/filesystem-watcher-handlers.ts @@ -1,10 +1,12 @@ import { ipcMain } from 'electron' +import { isCurrentWatcherSender } from './filesystem-watcher-sender-lifetime' import { onSshFilesystemProviderRegistered } from '../providers/ssh-filesystem-dispatch' import { watcherLifecycleState } from './filesystem-watcher-lifecycle-state' import { getRemoteWatcherKey } from './filesystem-watcher-paths' import { cancelInFlightRemoteInstallIfUnowned, forgetDesiredRemoteWatcher, + registerWatcherSenderCleanup, releaseRemoteWatchListener } from './filesystem-watcher-listener-lifecycle' import { @@ -30,6 +32,7 @@ export function registerFilesystemWatcherHandlers(): void { ipcMain.handle( 'fs:watchWorktree', async (event, args: { worktreePath: string; connectionId?: string }): Promise => { + const senderSignal = registerWatcherSenderCleanup(event.sender) if (args.connectionId) { // Why: a real new watch reopens the subsystem after closeAllWatchers latched it shut (also resets tests between cases). watcherLifecycleState.remoteWatchersClosed = false @@ -42,6 +45,9 @@ export function registerFilesystemWatcherHandlers(): void { args.connectionId, args.worktreePath ) + if (!isCurrentWatcherSender(event.sender, senderSignal)) { + return + } if (result === 'capacity') { // Why straight to the dormant backoff: the cap is full until some other root is released, // which a 1 Hz reinstall cannot bring about — it only adds relay load per refused root. diff --git a/src/main/ipc/filesystem-watcher-large-batch.test.ts b/src/main/ipc/filesystem-watcher-large-batch.test.ts index 3d71581595e..240ca629042 100644 --- a/src/main/ipc/filesystem-watcher-large-batch.test.ts +++ b/src/main/ipc/filesystem-watcher-large-batch.test.ts @@ -1,3 +1,4 @@ +import { createWatcherSender } from './filesystem-watcher-test-sender' import { join, resolve } from 'node:path' import { beforeEach, describe, expect, it, vi } from 'vitest' @@ -68,10 +69,7 @@ describe('local filesystem watcher large batches', () => { return { unsubscribe: vi.fn() } as never }) - await handlers['fs:watchWorktree']( - { sender: { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } }, - { worktreePath } - ) + await handlers['fs:watchWorktree']({ sender: createWatcherSender(1) }, { worktreePath }) const events = Array.from({ length: 200_000 }, (_, index): WatcherEvent => ({ type: 'delete', @@ -92,7 +90,7 @@ describe('local filesystem watcher large batches', () => { return { unsubscribe: vi.fn() } as never }) const worktreePath = resolve('/tmp/repo') - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath }) watcherCallback?.( @@ -124,7 +122,7 @@ describe('local filesystem watcher large batches', () => { }) const worktreePath = resolve('/tmp/repo') const filePath = join(worktreePath, 'a.ts') - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath }) watcherCallback?.(null, [{ type: 'update', path: filePath }]) @@ -152,7 +150,7 @@ describe('local filesystem watcher large batches', () => { return { unsubscribe: vi.fn() } as never }) const worktreePath = resolve('/tmp/repo') - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath }) // Step under the trailing window so only the max wait can force a flush. diff --git a/src/main/ipc/filesystem-watcher-lifecycle-state.ts b/src/main/ipc/filesystem-watcher-lifecycle-state.ts index 7d6a120c0a7..feef5d40373 100644 --- a/src/main/ipc/filesystem-watcher-lifecycle-state.ts +++ b/src/main/ipc/filesystem-watcher-lifecycle-state.ts @@ -65,7 +65,7 @@ export const watcherLifecycleState = { watchedRoots: new Map(), unwatchableRoots: new Set(), // Why: key cleanup by sender WebContents (not per root) to avoid MaxListeners warnings when a workspace has many worktrees open. - senderCleanupRegistered: new Set(), + senderLifetimes: new Map void }>(), pendingTeardowns: new Map>(), // Why: @parcel/watcher unsubscribe does native async work that sender-destroy can start before shutdown, so will-quit must still await it. pendingLocalUnsubscribes: new Set>(), diff --git a/src/main/ipc/filesystem-watcher-listener-lifecycle.ts b/src/main/ipc/filesystem-watcher-listener-lifecycle.ts index 713305069e3..c2dd5f3dc77 100644 --- a/src/main/ipc/filesystem-watcher-listener-lifecycle.ts +++ b/src/main/ipc/filesystem-watcher-listener-lifecycle.ts @@ -8,6 +8,7 @@ import { watcherLifecycleState } from './filesystem-watcher-lifecycle-state' import { cancelLocalBatchFlush } from './filesystem-watcher-batch-control' +import { captureWatcherSenderLifetime } from './filesystem-watcher-sender-lifetime' export function rememberUnwatchableRoot(rootKey: string): void { const { unwatchableRoots } = watcherLifecycleState @@ -172,13 +173,8 @@ export function releaseRemoteWatchListener(key: string, senderId: number): void watcherLifecycleState.remoteWatchers.delete(key) } -export function registerWatcherSenderCleanup(sender: WebContents): void { - if (watcherLifecycleState.senderCleanupRegistered.has(sender.id)) { - return - } - watcherLifecycleState.senderCleanupRegistered.add(sender.id) - sender.once('destroyed', () => { - watcherLifecycleState.senderCleanupRegistered.delete(sender.id) +export function registerWatcherSenderCleanup(sender: WebContents): AbortSignal { + return captureWatcherSenderLifetime(sender, () => { cleanupLocalWatchersForSender(sender.id) cleanupRemoteWatchersForSender(sender.id) }) @@ -215,6 +211,9 @@ function cleanupRemoteWatchersForSender(senderId: number): void { for (const key of Array.from(watcherLifecycleState.desiredRemoteWatchers.keys())) { forgetDesiredRemoteWatcher(key, senderId) } + for (const resync of watcherLifecycleState.remoteWatcherResyncStates.values()) { + resync.listeners.delete(senderId) + } for (const [key, suspended] of watcherLifecycleState.suspendedRemoteWatcherListeners) { suspended.listeners.delete(senderId) if (suspended.listeners.size === 0) { diff --git a/src/main/ipc/filesystem-watcher-local-events.test.ts b/src/main/ipc/filesystem-watcher-local-events.test.ts index b5ec3eab939..343ee2bc3ad 100644 --- a/src/main/ipc/filesystem-watcher-local-events.test.ts +++ b/src/main/ipc/filesystem-watcher-local-events.test.ts @@ -336,7 +336,7 @@ describe('local filesystem watcher flush serialization', () => { // Why real timers: fake-timers' refresh() revives a cleared handle, but Node's is a no-op — the bug only shows on real Timeouts. vi.useRealTimers() statMock.mockResolvedValue({ isDirectory: () => true }) - const listener = { ...sender, id: 7, once: vi.fn() } + const listener = { ...sender, id: 7, removeListener: vi.fn(), once: vi.fn() } try { await subscribeLocalWatcher('/repo', listener as never) watcherCallback?.(null, [{ type: 'delete', path: '/repo/file.ts' }]) diff --git a/src/main/ipc/filesystem-watcher-local-removal.ts b/src/main/ipc/filesystem-watcher-local-removal.ts index a3d4b6ff558..78f015f7b94 100644 --- a/src/main/ipc/filesystem-watcher-local-removal.ts +++ b/src/main/ipc/filesystem-watcher-local-removal.ts @@ -1,3 +1,4 @@ +import { isCurrentWatcherSender } from './filesystem-watcher-sender-lifetime' import type { WebContents } from 'electron' import type { FsChangedPayload } from '../../shared/filesystem-entry-types' import { @@ -10,6 +11,7 @@ import { getLocalWatcherRoot } from './filesystem-watcher-paths' import { watcherLifecycleState } from './filesystem-watcher-lifecycle-state' import { abandonLocalUnsubscribes, + registerWatcherSenderCleanup, clearLocalCapacityRetry, trackDetachedLocalUnsubscribe } from './filesystem-watcher-listener-lifecycle' @@ -124,29 +126,40 @@ export async function restoreLocalWatcherAfterFailedRemoval(worktreePath: string return } watcherLifecycleState.suspendedLocalWatcherListeners.delete(rootKey) - const failures: unknown[] = [] - const failedListeners = new Map() - for (const sender of suspended.listeners.values()) { - if (sender.isDestroyed()) { + const failures: { sender: WebContents; signal: AbortSignal; error: unknown }[] = [] + const owners = Array.from(suspended.listeners.values(), (sender) => ({ + sender, + signal: registerWatcherSenderCleanup(sender) + })) + for (const { sender, signal } of owners) { + if (!isCurrentWatcherSender(sender, signal)) { continue } try { - await subscribeLocalWatcher(suspended.worktreePath, sender) + await subscribeLocalWatcher(suspended.worktreePath, sender, undefined, signal) + if (!isCurrentWatcherSender(sender, signal)) { + continue + } sender.send('fs:changed', { worktreePath: suspended.worktreePath, events: [{ kind: 'overflow', absolutePath: suspended.worktreePath }] } satisfies FsChangedPayload) } catch (error) { - failures.push(error) - failedListeners.set(sender.id, sender) + if (!isCurrentWatcherSender(sender, signal)) { + continue + } + failures.push({ sender, signal, error }) } } - if (failures.length > 0) { + const liveFailures = failures.filter(({ sender, signal }) => + isCurrentWatcherSender(sender, signal) + ) + if (liveFailures.length > 0) { watcherLifecycleState.suspendedLocalWatcherListeners.set(rootKey, { worktreePath: suspended.worktreePath, - listeners: failedListeners + listeners: new Map(liveFailures.map(({ sender }) => [sender.id, sender])) }) - throw failures[0] + throw liveFailures[0].error } } diff --git a/src/main/ipc/filesystem-watcher-local-subscription.ts b/src/main/ipc/filesystem-watcher-local-subscription.ts index e888893dbf4..c667210078a 100644 --- a/src/main/ipc/filesystem-watcher-local-subscription.ts +++ b/src/main/ipc/filesystem-watcher-local-subscription.ts @@ -11,21 +11,25 @@ import { addLocalWatchListener, clearLocalCapacityRetry, rememberUnwatchableRoot, + registerWatcherSenderCleanup, takeLocalCapacityRetryListeners, trackDetachedLocalUnsubscribe } from './filesystem-watcher-listener-lifecycle' import { cancelLocalBatchFlush } from './filesystem-watcher-batch-control' import { scheduleLocalCapacityRetry } from './filesystem-watcher-local-capacity' import { installLocalWatcher } from './filesystem-watcher-local-install' +import { isCurrentWatcherSender } from './filesystem-watcher-sender-lifetime' // ── Subscribe / Unsubscribe ────────────────────────────────────────── export async function subscribeLocalWatcher( worktreePath: string, sender: WebContents, - generation = watcherLifecycleState.localWatcherLifecycleGeneration + generation = watcherLifecycleState.localWatcherLifecycleGeneration, + senderSignal = registerWatcherSenderCleanup(sender) ): Promise { if ( + !isCurrentWatcherSender(sender, senderSignal) || watcherLifecycleState.localWatchersClosed || generation !== watcherLifecycleState.localWatcherLifecycleGeneration ) { @@ -33,7 +37,7 @@ export async function subscribeLocalWatcher( } const finishInstall = beginWatcherInstall(worktreePath) try { - await subscribeWhileRemovalAllowed(worktreePath, sender, generation) + await subscribeWhileRemovalAllowed(worktreePath, sender, generation, senderSignal) } finally { finishInstall() } @@ -42,9 +46,11 @@ export async function subscribeLocalWatcher( async function subscribeWhileRemovalAllowed( worktreePath: string, sender: WebContents, - generation: number + generation: number, + senderSignal: AbortSignal ): Promise { if ( + !isCurrentWatcherSender(sender, senderSignal) || watcherLifecycleState.localWatchersClosed || generation !== watcherLifecycleState.localWatcherLifecycleGeneration ) { @@ -70,6 +76,9 @@ async function subscribeWhileRemovalAllowed( watcherLifecycleState.pendingTeardowns.delete(rootKey) } const capacityRetryListeners = takeLocalCapacityRetryListeners(rootKey) + const retrySignals = new Map( + capacityRetryListeners.map((listener) => [listener.id, registerWatcherSenderCleanup(listener)]) + ) if (root) { for (const listener of capacityRetryListeners) { @@ -91,6 +100,10 @@ async function subscribeWhileRemovalAllowed( } } const result = await pendingInstall + const liveCapacityListeners = capacityRetryListeners.filter((listener) => + isCurrentWatcherSender(listener, retrySignals.get(listener.id)!) + ) + const senderIsCurrent = isCurrentWatcherSender(sender, senderSignal) if ( result === 'cancelled' && !canJoinInstall && @@ -102,33 +115,42 @@ async function subscribeWhileRemovalAllowed( watcherLifecycleState.pendingLocalInstallPromises.delete(rootKey) } const retryListeners = new Map( - capacityRetryListeners.map((listener) => [listener.id, listener]) + liveCapacityListeners.map((listener) => [listener.id, listener]) ) - retryListeners.set(sender.id, sender) + if (senderIsCurrent) { + retryListeners.set(sender.id, sender) + } for (const listener of retryListeners.values()) { if (!listener.isDestroyed()) { - await subscribeWhileRemovalAllowed(worktreePath, listener, generation) + await subscribeWhileRemovalAllowed( + worktreePath, + listener, + generation, + listener === sender ? senderSignal : retrySignals.get(listener.id)! + ) } } return } if (!inFlight) { if (result === 'installed') { - for (const listener of capacityRetryListeners) { + for (const listener of liveCapacityListeners) { addLocalWatchListener(rootKey, listener) } } else if (result === 'capacity') { const retryListeners = new Map( - capacityRetryListeners.map((listener) => [listener.id, listener]) + liveCapacityListeners.map((listener) => [listener.id, listener]) ) - retryListeners.set(sender.id, sender) + if (senderIsCurrent) { + retryListeners.set(sender.id, sender) + } scheduleLocalCapacityRetry(rootKey, worktreePath, retryListeners, subscribeLocalWatcher) } } if ( result === 'installed' && watcherLifecycleState.watchedRoots.has(rootKey) && - !sender.isDestroyed() && + senderIsCurrent && (!inFlight || inFlight.listeners.has(sender.id)) ) { addLocalWatchListener(rootKey, sender) diff --git a/src/main/ipc/filesystem-watcher-local-unsubscribe.test.ts b/src/main/ipc/filesystem-watcher-local-unsubscribe.test.ts index eb161c1ba30..da8b92dbcf8 100644 --- a/src/main/ipc/filesystem-watcher-local-unsubscribe.test.ts +++ b/src/main/ipc/filesystem-watcher-local-unsubscribe.test.ts @@ -1,3 +1,4 @@ +import { createWatcherSender } from './filesystem-watcher-test-sender' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import type * as ParcelWatcherProcess from './parcel-watcher-process' @@ -85,6 +86,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { const sender = { isDestroyed: () => false, send: vi.fn(), + removeListener: vi.fn(), once: vi.fn((event: string, callback: () => void) => { if (event === 'destroyed') { destroyedCallbacks.push(callback) @@ -124,7 +126,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { watcherCallback = callback as typeof watcherCallback return { unsubscribe: unsubscribeMock } as never }) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: '/tmp/repo' }) watcherCallback(new Error('root disappeared'), []) @@ -157,6 +159,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { const sender = { isDestroyed: () => false, send: vi.fn(), + removeListener: vi.fn(), once: vi.fn((event: string, callback: () => void) => { if (event === 'destroyed') { destroyedCallbacks.push(callback) @@ -197,18 +200,8 @@ describe('local filesystem watcher unsubscribe cleanup', () => { subscribeResolvers.push(resolve as (subscription: { unsubscribe: () => void }) => void) }) ) - const senderOne = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } - const senderTwo = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 2 - } + const senderOne = createWatcherSender(1) + const senderTwo = createWatcherSender(2) const watchOne = handlers['fs:watchWorktree']( { sender: senderOne }, @@ -248,12 +241,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { vi.mocked(stat).mockResolvedValue({ isDirectory: () => true } as never) const unsubscribeMock = vi.fn() vi.mocked(subscribeParcelWatcher).mockResolvedValue({ unsubscribe: unsubscribeMock } as never) - const sender = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: '/tmp/repo' }) @@ -274,12 +262,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { vi.mocked(stat).mockResolvedValue({ isDirectory: () => true } as never) const unsubscribeMock = vi.fn() vi.mocked(subscribeParcelWatcher).mockResolvedValue({ unsubscribe: unsubscribeMock } as never) - const sender = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: '/tmp/repo' }) await closeLocalWatcherForWorktreePath('/tmp/repo') @@ -294,12 +277,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { vi.mocked(stat).mockResolvedValue({ isDirectory: () => true } as never) const unsubscribeMock = vi.fn() vi.mocked(subscribeParcelWatcher).mockResolvedValue({ unsubscribe: unsubscribeMock } as never) - const sender = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: watchPath }) await closeLocalWatcherForWorktreePath(closePath) @@ -323,12 +301,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { watcherCallback = callback as typeof watcherCallback return { unsubscribe: unsubscribeMock } as never }) - const sender = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: watchPath }) watcherCallback(null, [{ type: 'update', path: `${watchPath}\\file.txt` }]) @@ -349,7 +322,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { const terminationError = new Error('watcher child did not exit') const unsubscribeMock = vi.fn().mockRejectedValue(terminationError) vi.mocked(subscribeParcelWatcher).mockResolvedValue({ unsubscribe: unsubscribeMock } as never) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: '/tmp/repo' }) @@ -372,6 +345,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { const sender = { isDestroyed: () => false, send: vi.fn(), + removeListener: vi.fn(), once: vi.fn((event: string, callback: () => void) => { if (event === 'destroyed') { destroyedCallbacks.push(callback) @@ -413,6 +387,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { const sender = { isDestroyed: () => false, send: vi.fn(), + removeListener: vi.fn(), once: vi.fn((event: string, callback: () => void) => { if (event === 'destroyed') { destroyedCallbacks.push(callback) @@ -448,7 +423,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { watcherCallback = callback as typeof watcherCallback return { unsubscribe } }) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: '/tmp/repo' }) watcherCallback(terminationError, []) @@ -478,7 +453,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { hooks?.signal?.addEventListener('abort', () => reject(terminationError), { once: true }) }) ) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) const watchPromise = handlers['fs:watchWorktree']( { sender }, { worktreePath: '/tmp/repo' } @@ -499,12 +474,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { vi.mocked(stat).mockResolvedValue({ isDirectory: () => true } as never) const unsubscribeMock = vi.fn() vi.mocked(subscribeParcelWatcher).mockResolvedValue({ unsubscribe: unsubscribeMock } as never) - const sender = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: '/tmp/repo' }) @@ -532,12 +502,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { resolveSubscribe = resolve as typeof resolveSubscribe }) ) - const sender = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1) const watchPromise = handlers['fs:watchWorktree']( { sender }, @@ -581,8 +546,8 @@ describe('local filesystem watcher unsubscribe cleanup', () => { replacementCallback = callback return { unsubscribe: replacementUnsubscribe } }) - const firstSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } - const replacementSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } + const firstSender = createWatcherSender(1) + const replacementSender = createWatcherSender(2) const firstWatch = handlers['fs:watchWorktree']( { sender: firstSender }, @@ -629,9 +594,9 @@ describe('local filesystem watcher unsubscribe cleanup', () => { .mockImplementationOnce(install as never) .mockImplementationOnce(install as never) const lateUnsubscribe = vi.fn() - const firstSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } - const joinerSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } - const reopenSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 3 } + const firstSender = createWatcherSender(1) + const joinerSender = createWatcherSender(2) + const reopenSender = createWatcherSender(3) const first = handlers['fs:watchWorktree']( { sender: firstSender }, { worktreePath: '/tmp/repo' } @@ -675,12 +640,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { resolveSubscribe = resolve as typeof resolveSubscribe }) ) - const sender = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1) const watchPromise = handlers['fs:watchWorktree']( { sender }, @@ -708,12 +668,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { vi.mocked(subscribeParcelWatcher) .mockResolvedValueOnce({ unsubscribe: firstUnsubscribe } as never) .mockResolvedValueOnce({ unsubscribe: replacementUnsubscribe } as never) - const sender = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: '/tmp/repo' }) await closeLocalWatcherForWorktreePath('/tmp/repo') @@ -731,7 +686,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { vi.mocked(stat).mockResolvedValue({ isDirectory: () => true } as never) const firstUnsubscribe = vi.fn() vi.mocked(subscribeParcelWatcher).mockResolvedValue({ unsubscribe: firstUnsubscribe } as never) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: '/tmp/repo' }) await closeLocalWatcherForWorktreePath('/tmp/repo') @@ -751,6 +706,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { const sender = { isDestroyed: () => false, send: vi.fn(), + removeListener: vi.fn(), once: vi.fn((event: string, callback: () => void) => { if (event === 'destroyed') { destroyedCallbacks.push(callback) @@ -779,12 +735,7 @@ describe('local filesystem watcher unsubscribe cleanup', () => { resolveSubscribe = resolve as typeof resolveSubscribe }) ) - const sender = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1) const watchPromise = handlers['fs:watchWorktree']( { sender }, diff --git a/src/main/ipc/filesystem-watcher-native-capacity.test.ts b/src/main/ipc/filesystem-watcher-native-capacity.test.ts index 670b7eb2af3..9d291f7a9ee 100644 --- a/src/main/ipc/filesystem-watcher-native-capacity.test.ts +++ b/src/main/ipc/filesystem-watcher-native-capacity.test.ts @@ -1,3 +1,4 @@ +import { createWatcherSender } from './filesystem-watcher-test-sender' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' const { handleMock, statMock, subscribeViaWatcherProcessMock, disposeWatcherProcessMock } = @@ -71,7 +72,7 @@ describe('native filesystem watcher capacity recovery', () => { subscribeViaWatcherProcessMock .mockRejectedValueOnce(new WatcherChildCapacityError()) .mockResolvedValueOnce({ unsubscribe }) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) const args = { worktreePath: '/tmp/native-capacity-root' } await handlers['fs:watchWorktree']({ sender }, args) @@ -86,7 +87,7 @@ describe('native filesystem watcher capacity recovery', () => { it('cancels the native capacity wait when its renderer unwatches', async () => { const releases = fillWatcherChildCapacity() subscribeViaWatcherProcessMock.mockRejectedValue(new WatcherChildCapacityError()) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) const args = { worktreePath: '/tmp/native-capacity-root' } await handlers['fs:watchWorktree']({ sender }, args) diff --git a/src/main/ipc/filesystem-watcher-real.test.ts b/src/main/ipc/filesystem-watcher-real.test.ts index 7c9556e7ba4..6bd5060fb0e 100644 --- a/src/main/ipc/filesystem-watcher-real.test.ts +++ b/src/main/ipc/filesystem-watcher-real.test.ts @@ -1,3 +1,4 @@ +import { createWatcherSender } from './filesystem-watcher-test-sender' /* * Real (unmocked) @parcel/watcher integration test. * @@ -95,12 +96,7 @@ describe('filesystem-watcher real @parcel/watcher integration', () => { // tmpdir() returns /var, so compare canonical paths instead of aliases. tempDir = await realpath(await mkdtemp(join(tmpdir(), 'orca-fswatch-real-'))) const sendMock = vi.fn() - const sender = { - isDestroyed: () => false, - send: sendMock, - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1, sendMock) // Subscribe resolves only after the native watcher is installed. await handlers['fs:watchWorktree']({ sender }, { worktreePath: tempDir }) @@ -139,7 +135,7 @@ describe('filesystem-watcher real @parcel/watcher integration', () => { await mkdir(join(root.realRoot, 'src'), { recursive: true }) const sendMock = vi.fn() - const sender = { isDestroyed: () => false, send: sendMock, once: vi.fn(), id: 1 } + const sender = createWatcherSender(1, sendMock) await handlers['fs:watchWorktree']({ sender }, { worktreePath: root.aliasRoot }) const expectedPath = join(root.aliasRoot, 'src', 'agent-edit.ts') diff --git a/src/main/ipc/filesystem-watcher-remote-batch.test.ts b/src/main/ipc/filesystem-watcher-remote-batch.test.ts index de1f4fb3d37..fc4566b97c9 100644 --- a/src/main/ipc/filesystem-watcher-remote-batch.test.ts +++ b/src/main/ipc/filesystem-watcher-remote-batch.test.ts @@ -1,3 +1,4 @@ +import { senderEvents } from './filesystem-watcher-test-sender' import { beforeEach, describe, expect, it, vi } from 'vitest' import type { FsChangeEvent } from '../../shared/filesystem-entry-types' @@ -34,7 +35,13 @@ describe('remote filesystem watcher batching', () => { const watchCallbacks: WatchCallback[] = [] function makeSender(overrides: Partial<{ isDestroyed: () => boolean }> = {}) { - return { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1, ...overrides } + return { + isDestroyed: () => false, + send: vi.fn(), + ...senderEvents(), + id: 1, + ...overrides + } } beforeEach(async () => { diff --git a/src/main/ipc/filesystem-watcher-remote-cancellation.test.ts b/src/main/ipc/filesystem-watcher-remote-cancellation.test.ts index 60e25db9040..26140876bb5 100644 --- a/src/main/ipc/filesystem-watcher-remote-cancellation.test.ts +++ b/src/main/ipc/filesystem-watcher-remote-cancellation.test.ts @@ -1,3 +1,4 @@ +import { createWatcherSender } from './filesystem-watcher-test-sender' import { beforeEach, describe, expect, it, vi } from 'vitest' const { handleMock, getSshFilesystemProviderMock } = vi.hoisted(() => ({ @@ -50,8 +51,8 @@ describe('remote filesystem watcher cancellation', () => { ) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) const args = { worktreePath: '/home/me/repo', connectionId: 'conn-1' } - const senderOne = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } - const senderTwo = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } + const senderOne = createWatcherSender(1) + const senderTwo = createWatcherSender(2) const first = handlers['fs:watchWorktree']({ sender: senderOne }, args) as Promise const second = handlers['fs:watchWorktree']({ sender: senderTwo }, args) as Promise @@ -101,8 +102,8 @@ describe('remote filesystem watcher cancellation', () => { }) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) const args = { worktreePath: '/home/me/repo', connectionId: 'conn-1' } - const firstSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } - const secondSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } + const firstSender = createWatcherSender(1) + const secondSender = createWatcherSender(2) const first = handlers['fs:watchWorktree']({ sender: firstSender }, args) as Promise await Promise.resolve() @@ -136,6 +137,7 @@ describe('remote filesystem watcher cancellation', () => { const destroyedSender = { isDestroyed: () => false, send: vi.fn(), + removeListener: vi.fn(), once: vi.fn((event: string, callback: () => void) => { if (event === 'destroyed') { destroyedCallbacks.push(callback) @@ -156,7 +158,7 @@ describe('remote filesystem watcher cancellation', () => { installs.get('/destroyed')?.resolve(vi.fn()) await destroyedWatch - const shutdownSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } + const shutdownSender = createWatcherSender(2) const shutdownArgs = { worktreePath: '/shutdown', connectionId: 'conn-1' } const shutdownWatch = handlers['fs:watchWorktree']( { sender: shutdownSender }, @@ -182,8 +184,8 @@ describe('remote filesystem watcher cancellation', () => { ) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) const args = { worktreePath: '/home/me/repo', connectionId: 'conn-1' } - const senderOne = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } - const senderTwo = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } + const senderOne = createWatcherSender(1) + const senderTwo = createWatcherSender(2) const first = handlers['fs:watchWorktree']({ sender: senderOne }, args) as Promise await Promise.resolve() @@ -216,8 +218,8 @@ describe('remote filesystem watcher cancellation', () => { ) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) const args = { worktreePath: '/home/me/repo', connectionId: 'conn-1' } - const firstSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } - const secondSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } + const firstSender = createWatcherSender(1) + const secondSender = createWatcherSender(2) const first = handlers['fs:watchWorktree']({ sender: firstSender }, args) as Promise await Promise.resolve() @@ -259,9 +261,9 @@ describe('remote filesystem watcher cancellation', () => { const args = { worktreePath: '/home/me/repo', connectionId: 'conn-1' } const reopenArgs = { worktreePath: '/home/me/other', connectionId: 'conn-1' } const lateUnwatch = vi.fn() - const firstSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } - const joinerSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } - const reopenSender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 3 } + const firstSender = createWatcherSender(1) + const joinerSender = createWatcherSender(2) + const reopenSender = createWatcherSender(3) const first = handlers['fs:watchWorktree']({ sender: firstSender }, args) as Promise await Promise.resolve() diff --git a/src/main/ipc/filesystem-watcher-remote-capacity.test.ts b/src/main/ipc/filesystem-watcher-remote-capacity.test.ts index 97a8f82bf33..f67083f845d 100644 --- a/src/main/ipc/filesystem-watcher-remote-capacity.test.ts +++ b/src/main/ipc/filesystem-watcher-remote-capacity.test.ts @@ -1,3 +1,4 @@ +import { createWatcherSender } from './filesystem-watcher-test-sender' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' const { handleMock, getSshFilesystemProviderMock } = vi.hoisted(() => ({ @@ -58,7 +59,7 @@ describe('remote filesystem watcher capacity refusals', () => { throw new Error(WATCH_ROOT_CAPACITY_REFUSAL_MESSAGE) }) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) const args = { worktreePath: '/home/me/repos/one', connectionId: 'conn-capacity' } await handlers['fs:watchWorktree']({ sender }, args) @@ -78,7 +79,7 @@ describe('remote filesystem watcher capacity refusals', () => { throw new Error('Relay channel lost') }) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } + const sender = createWatcherSender(2) const args = { worktreePath: '/home/me/repos/two', connectionId: 'conn-unavailable' } await handlers['fs:watchWorktree']({ sender }, args) diff --git a/src/main/ipc/filesystem-watcher-remote-controller.ts b/src/main/ipc/filesystem-watcher-remote-controller.ts index 1b6953cd292..8a4f4cc2e91 100644 --- a/src/main/ipc/filesystem-watcher-remote-controller.ts +++ b/src/main/ipc/filesystem-watcher-remote-controller.ts @@ -1,6 +1,7 @@ import type { WebContents } from 'electron' import type { RemoteWatcherInstallToken } from './filesystem-watcher-lifecycle-state' import { watcherLifecycleState } from './filesystem-watcher-lifecycle-state' +import { registerWatcherSenderCleanup } from './filesystem-watcher-listener-lifecycle' import { installRemoteWatcherCore, type RemoteWatcherTerminalErrorHandler @@ -14,14 +15,16 @@ export function installRemoteWatcher( sender: WebContents, connectionId: string, worktreePath: string, - generation = watcherLifecycleState.remoteWatcherLifecycleGeneration + generation = watcherLifecycleState.remoteWatcherLifecycleGeneration, + senderSignal = registerWatcherSenderCleanup(sender) ) { return installRemoteWatcherCore( sender, connectionId, worktreePath, handleRemoteWatcherTerminalError, - generation + generation, + senderSignal ) } diff --git a/src/main/ipc/filesystem-watcher-remote-dormant.ts b/src/main/ipc/filesystem-watcher-remote-dormant.ts index 7c168114f51..11783f7bb57 100644 --- a/src/main/ipc/filesystem-watcher-remote-dormant.ts +++ b/src/main/ipc/filesystem-watcher-remote-dormant.ts @@ -1,3 +1,5 @@ +import { isCurrentWatcherSender } from './filesystem-watcher-sender-lifetime' +import { registerWatcherSenderCleanup } from './filesystem-watcher-listener-lifecycle' import { getSshFilesystemProvider } from '../providers/ssh-filesystem-dispatch' import { isWatcherRemovalInProgressError } from './watcher-removal-gate' import type { @@ -74,6 +76,7 @@ async function rearmDormantRemoteWatcher( } const listeners = Array.from(desired.listeners.values()) + const signals = listeners.map(registerWatcherSenderCleanup) let results: Awaited>[] try { results = await Promise.all( @@ -84,6 +87,9 @@ async function rearmDormantRemoteWatcher( // Why: removal owns the key now and either forgets the intent or restores the watch itself. return } + if (!listeners.some((listener, index) => isCurrentWatcherSender(listener, signals[index]))) { + return + } scheduleDormantRemoteWatcherRearmCore( connectionId, worktreePath, @@ -92,10 +98,16 @@ async function rearmDormantRemoteWatcher( ) return } + if (!listeners.some((listener, index) => isCurrentWatcherSender(listener, signals[index]))) { + return + } dependencies.requestResync( key, worktreePath, - listeners.filter((_, index) => results[index] === 'installed') + listeners.filter( + (listener, index) => + results[index] === 'installed' && isCurrentWatcherSender(listener, signals[index]) + ) ) // Why: 'cancelled' means shutdown or the last listener left, so only a refusal stays dormant. if (results.some((result) => result === 'unavailable' || result === 'capacity')) { diff --git a/src/main/ipc/filesystem-watcher-remote-install.ts b/src/main/ipc/filesystem-watcher-remote-install.ts index 3ab5babbd91..b1a6bab61fa 100644 --- a/src/main/ipc/filesystem-watcher-remote-install.ts +++ b/src/main/ipc/filesystem-watcher-remote-install.ts @@ -15,6 +15,7 @@ import type { } from './filesystem-watcher-lifecycle-state' import { watcherLifecycleState } from './filesystem-watcher-lifecycle-state' import { getRemoteWatcherKey } from './filesystem-watcher-paths' +import { isCurrentWatcherSender } from './filesystem-watcher-sender-lifetime' import { addInFlightRemoteInstallListener, addRemoteWatchListener, @@ -34,10 +35,12 @@ export async function installRemoteWatcherCore( connectionId: string, worktreePath: string, onTerminalError: RemoteWatcherTerminalErrorHandler, - generation = watcherLifecycleState.remoteWatcherLifecycleGeneration + generation = watcherLifecycleState.remoteWatcherLifecycleGeneration, + senderSignal = registerWatcherSenderCleanup(sender) ): Promise { // Why: refuse installs racing in after teardown (or a waiter from an earlier lifecycle) so provider.watch() isn't called post-shutdown. if ( + !isCurrentWatcherSender(sender, senderSignal) || watcherLifecycleState.remoteWatchersClosed || generation !== watcherLifecycleState.remoteWatcherLifecycleGeneration ) { @@ -50,7 +53,8 @@ export async function installRemoteWatcherCore( connectionId, worktreePath, onTerminalError, - generation + generation, + senderSignal ) } finally { finishInstall() @@ -62,7 +66,8 @@ async function installRemoteWatcherWhileRemovalAllowed( connectionId: string, worktreePath: string, onTerminalError: RemoteWatcherTerminalErrorHandler, - generation: number + generation: number, + senderSignal: AbortSignal ): Promise { const provider = getSshFilesystemProvider(connectionId) if (!provider || sender.isDestroyed()) { @@ -85,6 +90,9 @@ async function installRemoteWatcherWhileRemovalAllowed( addInFlightRemoteInstallListener(inFlight, sender) } const result = await pendingInstall + if (!isCurrentWatcherSender(sender, senderSignal)) { + return 'cancelled' + } if ( result === 'installed' && watcherLifecycleState.remoteWatchers.has(key) && @@ -108,7 +116,8 @@ async function installRemoteWatcherWhileRemovalAllowed( connectionId, worktreePath, onTerminalError, - generation + generation, + senderSignal ) } return result diff --git a/src/main/ipc/filesystem-watcher-remote-provider-rearm.ts b/src/main/ipc/filesystem-watcher-remote-provider-rearm.ts index 3cc657c4086..5c55c537f07 100644 --- a/src/main/ipc/filesystem-watcher-remote-provider-rearm.ts +++ b/src/main/ipc/filesystem-watcher-remote-provider-rearm.ts @@ -1,3 +1,4 @@ +import { isCurrentWatcherSender } from './filesystem-watcher-sender-lifetime' import type { WebContents } from 'electron' import { isWatcherRemovalInProgressError } from './watcher-removal-gate' import type { @@ -5,7 +6,10 @@ import type { RequestRemoteWatcherResync } from './filesystem-watcher-remote-retry' import { watcherLifecycleState } from './filesystem-watcher-lifecycle-state' -import { clearDormantRemoteWatcher } from './filesystem-watcher-listener-lifecycle' +import { + clearDormantRemoteWatcher, + registerWatcherSenderCleanup +} from './filesystem-watcher-listener-lifecycle' type ScheduleRemoteWatcherRetry = ( sender: WebContents, @@ -79,24 +83,33 @@ export function reinstallRemoteWatchersForConnectionCore( watcherLifecycleState.loggedUnavailableRemoteWatchers.delete(key) const listeners = Array.from(desired.listeners.values()) + const signals = listeners.map(registerWatcherSenderCleanup) + const liveListeners = () => + listeners.filter((listener, index) => isCurrentWatcherSender(listener, signals[index])) void Promise.all( listeners.map((listener) => dependencies.install(listener, desired.connectionId, desired.worktreePath) ) ) .then((results) => { + if (liveListeners().length === 0) { + return + } // Why: events between the transport dropping and this reinstall are gone for good. dependencies.requestResync( key, desired.worktreePath, - listeners.filter((_, index) => results[index] === 'installed') + listeners.filter( + (listener, index) => + results[index] === 'installed' && isCurrentWatcherSender(listener, signals[index]) + ) ) if (results.some((result) => result === 'capacity')) { dependencies.scheduleDormant(desired.connectionId, desired.worktreePath) return } if (results.some((result) => result === 'unavailable')) { - for (const listener of listeners) { + for (const listener of liveListeners()) { dependencies.scheduleRetry( listener, desired.connectionId, @@ -111,7 +124,7 @@ export function reinstallRemoteWatchersForConnectionCore( if (isWatcherRemovalInProgressError(error)) { return } - for (const listener of listeners) { + for (const listener of liveListeners()) { dependencies.scheduleRetry( listener, desired.connectionId, diff --git a/src/main/ipc/filesystem-watcher-remote-rearm.test.ts b/src/main/ipc/filesystem-watcher-remote-rearm.test.ts index 0eecfbd82e9..056b10579a7 100644 --- a/src/main/ipc/filesystem-watcher-remote-rearm.test.ts +++ b/src/main/ipc/filesystem-watcher-remote-rearm.test.ts @@ -1,3 +1,4 @@ +import { createWatcherSender } from './filesystem-watcher-test-sender' import { beforeEach, describe, expect, it, vi } from 'vitest' const { handleMock, getSshFilesystemProviderMock, providerRegistrationListeners } = vi.hoisted( @@ -51,8 +52,8 @@ describe('remote filesystem watcher re-arm', () => { it('still resyncs when a fresh watch beat the failed reinstall to the retry slot', async () => { vi.useFakeTimers() const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) - const senderOne = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } - const senderTwo = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } + const senderOne = createWatcherSender(1) + const senderTwo = createWatcherSender(2) const args = { worktreePath: '/home/me/repo', connectionId: 'conn-1' } getSshFilesystemProviderMock.mockReturnValue({ watch: vi.fn().mockResolvedValue(vi.fn()) }) diff --git a/src/main/ipc/filesystem-watcher-remote-removal.ts b/src/main/ipc/filesystem-watcher-remote-removal.ts index 0e9c741851a..a6c7f2a6636 100644 --- a/src/main/ipc/filesystem-watcher-remote-removal.ts +++ b/src/main/ipc/filesystem-watcher-remote-removal.ts @@ -1,3 +1,4 @@ +import { isCurrentWatcherSender } from './filesystem-watcher-sender-lifetime' import type { WebContents } from 'electron' import type { FsChangedPayload } from '../../shared/filesystem-entry-types' import { getSshFilesystemProvider } from '../providers/ssh-filesystem-dispatch' @@ -5,6 +6,7 @@ import { watcherLifecycleState } from './filesystem-watcher-lifecycle-state' import { getRemoteWatcherKey } from './filesystem-watcher-paths' import { clearDormantRemoteWatcher, + registerWatcherSenderCleanup, clearRemoteWatcherResync } from './filesystem-watcher-listener-lifecycle' import { @@ -71,11 +73,18 @@ export async function restoreRemoteWatcherAfterFailedRemoval( return } watcherLifecycleState.suspendedRemoteWatcherListeners.delete(key) - for (const sender of suspended.listeners.values()) { - if (sender.isDestroyed()) { + const owners = Array.from(suspended.listeners.values(), (sender) => ({ + sender, + signal: registerWatcherSenderCleanup(sender) + })) + for (const { sender, signal } of owners) { + if (!isCurrentWatcherSender(sender, signal)) { + continue + } + const result = await installRemoteWatcher(sender, connectionId, worktreePath, undefined, signal) + if (!isCurrentWatcherSender(sender, signal)) { continue } - const result = await installRemoteWatcher(sender, connectionId, worktreePath) if (result === 'capacity') { scheduleDormantRemoteWatcherRearm(connectionId, worktreePath) } else if (result === 'unavailable') { diff --git a/src/main/ipc/filesystem-watcher-remote-retry.ts b/src/main/ipc/filesystem-watcher-remote-retry.ts index d73f9ed4db4..02c797e9c75 100644 --- a/src/main/ipc/filesystem-watcher-remote-retry.ts +++ b/src/main/ipc/filesystem-watcher-remote-retry.ts @@ -1,3 +1,4 @@ +import { isCurrentWatcherSender } from './filesystem-watcher-sender-lifetime' import type { WebContents } from 'electron' import type { FsChangedPayload } from '../../shared/filesystem-entry-types' import { isWatcherRemovalInProgressError } from './watcher-removal-gate' @@ -8,7 +9,10 @@ import { watcherLifecycleState } from './filesystem-watcher-lifecycle-state' import { getRemoteWatcherKey } from './filesystem-watcher-paths' -import { clearRemoteWatcherResync } from './filesystem-watcher-listener-lifecycle' +import { + clearRemoteWatcherResync, + registerWatcherSenderCleanup +} from './filesystem-watcher-listener-lifecycle' import { isCurrentDesiredRemoteWatcher } from './filesystem-watcher-remote-desired' export type InstallRemoteWatcher = ( @@ -44,6 +48,9 @@ export function scheduleRemoteWatcherRetryCore( resyncOnInstall = false ): void { const key = getRemoteWatcherKey(connectionId, worktreePath) + if (watcherLifecycleState.remoteWatchersClosed || !isCurrentDesiredRemoteWatcher(key, sender)) { + return + } const existingRetry = watcherLifecycleState.pendingRemoteWatcherRetryListeners.get(key) if (existingRetry) { if (!sender.isDestroyed()) { @@ -89,15 +96,24 @@ export function scheduleRemoteWatcherRetryCore( const listeners = Array.from(retry.listeners.values()).filter( (listener) => !listener.isDestroyed() && isCurrentDesiredRemoteWatcher(key, listener) ) + const signals = listeners.map(registerWatcherSenderCleanup) + const liveListeners = () => + listeners.filter((listener, index) => isCurrentWatcherSender(listener, signals[index])) void Promise.all( listeners.map((listener) => dependencies.install(listener, connectionId, worktreePath)) ) .then((results) => { + if (liveListeners().length === 0) { + return + } if (retry.resyncOnInstall) { dependencies.requestResync( key, worktreePath, - listeners.filter((_, index) => results[index] === 'installed') + listeners.filter( + (listener, index) => + results[index] === 'installed' && isCurrentWatcherSender(listener, signals[index]) + ) ) } // Why capacity leaves the fast window: the relay is refusing on a full watch-root cap, and a @@ -108,7 +124,7 @@ export function scheduleRemoteWatcherRetryCore( } // Why: don't re-arm on 'cancelled' (renderer stopped watching) — it would fire a stale overflow when the 60s window expires. if (results.some((result) => result === 'unavailable')) { - for (const listener of listeners) { + for (const listener of liveListeners()) { scheduleRemoteWatcherRetryCore( listener, connectionId, @@ -124,7 +140,7 @@ export function scheduleRemoteWatcherRetryCore( if (isWatcherRemovalInProgressError(error)) { return } - for (const listener of listeners) { + for (const listener of liveListeners()) { scheduleRemoteWatcherRetryCore( listener, connectionId, diff --git a/src/main/ipc/filesystem-watcher-removal-deadline.test.ts b/src/main/ipc/filesystem-watcher-removal-deadline.test.ts index 9eabdb8964e..ca1e731c99a 100644 --- a/src/main/ipc/filesystem-watcher-removal-deadline.test.ts +++ b/src/main/ipc/filesystem-watcher-removal-deadline.test.ts @@ -1,3 +1,4 @@ +import { createWatcherSender } from './filesystem-watcher-test-sender' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import type * as ParcelWatcherProcess from './parcel-watcher-process' @@ -84,12 +85,7 @@ describe('local filesystem watcher removal deadline', () => { // subscribe that ignores the abort signal exercises the deadline rather than the cancel path. // Once, so the wedge cannot leak into the next test. vi.mocked(subscribeViaWatcherProcess).mockImplementationOnce(() => new Promise(() => {})) - const sender = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1) const watchPromise = handlers['fs:watchWorktree']( { sender }, @@ -132,7 +128,7 @@ describe('local filesystem watcher removal deadline', () => { }) ) vi.mocked(subscribeParcelWatcher).mockResolvedValue({ unsubscribe: unsubscribeMock } as never) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: '/tmp/repo' }) @@ -160,12 +156,7 @@ describe('local filesystem watcher removal deadline', () => { try { vi.mocked(stat).mockResolvedValue({ isDirectory: () => true } as never) vi.mocked(subscribeViaWatcherProcess).mockImplementationOnce(() => new Promise(() => {})) - const sender = { - isDestroyed: () => false, - send: vi.fn(), - once: vi.fn(), - id: 1 - } + const sender = createWatcherSender(1) const watchPromise = handlers['fs:watchWorktree']( { sender }, @@ -209,7 +200,7 @@ describe('local filesystem watcher removal deadline', () => { }) ) vi.mocked(subscribeParcelWatcher).mockResolvedValue({ unsubscribe: unsubscribeMock } as never) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']({ sender }, { worktreePath: '/tmp/repo' }) vi.useFakeTimers() diff --git a/src/main/ipc/filesystem-watcher-sender-lifetime.ts b/src/main/ipc/filesystem-watcher-sender-lifetime.ts new file mode 100644 index 00000000000..b72948471b1 --- /dev/null +++ b/src/main/ipc/filesystem-watcher-sender-lifetime.ts @@ -0,0 +1,36 @@ +import type { WebContents } from 'electron' +import { abortWhenRendererGone } from './renderer-lifetime-abort' +import { watcherLifecycleState } from './filesystem-watcher-lifecycle-state' + +export function captureWatcherSenderLifetime( + sender: WebContents, + cleanup: () => void +): AbortSignal { + const existing = watcherLifecycleState.senderLifetimes.get(sender.id) + if (existing) { + return existing.signal + } + if (sender.isDestroyed()) { + return AbortSignal.abort() + } + const lifetime = abortWhenRendererGone(sender) + watcherLifecycleState.senderLifetimes.set(sender.id, lifetime) + lifetime.signal.addEventListener( + 'abort', + () => { + lifetime.dispose() + watcherLifecycleState.senderLifetimes.delete(sender.id) + cleanup() + }, + { once: true } + ) + return lifetime.signal +} + +export function isCurrentWatcherSender(sender: WebContents, signal: AbortSignal): boolean { + return ( + !sender.isDestroyed() && + !signal.aborted && + watcherLifecycleState.senderLifetimes.get(sender.id)?.signal === signal + ) +} diff --git a/src/main/ipc/filesystem-watcher-shutdown.ts b/src/main/ipc/filesystem-watcher-shutdown.ts index 9f9cd2de353..31a57ecab10 100644 --- a/src/main/ipc/filesystem-watcher-shutdown.ts +++ b/src/main/ipc/filesystem-watcher-shutdown.ts @@ -8,7 +8,10 @@ export async function closeAllWatchers(): Promise { // Why: drop the intent with the rest of the state, but keep the provider-registration // subscription — a new fs:watchWorktree reopens the subsystem and still needs the re-arm hook. watcherLifecycleState.desiredRemoteWatchers.clear() - watcherLifecycleState.senderCleanupRegistered.clear() + for (const lifetime of watcherLifecycleState.senderLifetimes.values()) { + lifetime.dispose() + } + watcherLifecycleState.senderLifetimes.clear() watcherLifecycleState.unwatchableRoots.clear() watcherLifecycleState.suspendedLocalWatcherListeners.clear() watcherLifecycleState.suspendedRemoteWatcherListeners.clear() diff --git a/src/main/ipc/filesystem-watcher-terminal-resync.test.ts b/src/main/ipc/filesystem-watcher-terminal-resync.test.ts index 31e87d7a4c9..e8a81d1a7ff 100644 --- a/src/main/ipc/filesystem-watcher-terminal-resync.test.ts +++ b/src/main/ipc/filesystem-watcher-terminal-resync.test.ts @@ -41,6 +41,7 @@ const OVERFLOW_PAYLOAD = { type MockSender = { isDestroyed: () => boolean send: ReturnType + removeListener: ReturnType once: ReturnType id: number destroy: () => void @@ -52,6 +53,7 @@ function createSender(id: number): MockSender { return { isDestroyed: () => destroyed, send: vi.fn(), + removeListener: vi.fn(), once: vi.fn((event: string, handler: () => void) => { if (event === 'destroyed') { destroyedHandlers.push(handler) diff --git a/src/main/ipc/filesystem-watcher-test-sender.ts b/src/main/ipc/filesystem-watcher-test-sender.ts new file mode 100644 index 00000000000..4ff5a5bbcc1 --- /dev/null +++ b/src/main/ipc/filesystem-watcher-test-sender.ts @@ -0,0 +1,18 @@ +import { vi, type Mock } from 'vitest' + +type SenderEvents = { once: Mock; removeListener: Mock } + +export function senderEvents(): SenderEvents { + return { once: vi.fn(), removeListener: vi.fn() } +} + +export function createWatcherSender( + id: number, + send: Mock = vi.fn() +): SenderEvents & { + id: number + isDestroyed: () => boolean + send: Mock +} { + return { id, isDestroyed: () => false, send, ...senderEvents() } +} diff --git a/src/main/ipc/filesystem-watcher-unwatchable-roots.test.ts b/src/main/ipc/filesystem-watcher-unwatchable-roots.test.ts index 8378d55ab0d..78274bbe028 100644 --- a/src/main/ipc/filesystem-watcher-unwatchable-roots.test.ts +++ b/src/main/ipc/filesystem-watcher-unwatchable-roots.test.ts @@ -1,3 +1,4 @@ +import { createWatcherSender } from './filesystem-watcher-test-sender' import { statSync } from 'node:fs' import { beforeEach, describe, expect, it, vi } from 'vitest' @@ -52,7 +53,7 @@ describe('filesystem watcher unwatchable root cache', () => { it('releases the install record when the root is a file', async () => { const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) vi.mocked(stat).mockResolvedValue(statSync(new URL(import.meta.url))) try { await handlers['fs:watchWorktree']({ sender }, { worktreePath: '/tmp/not-directory' }) @@ -66,7 +67,7 @@ describe('filesystem watcher unwatchable root cache', () => { it('evicts oldest failed local roots while suppressing recent retries', async () => { const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) vi.mocked(stat).mockRejectedValue(new Error('missing')) for (let i = 0; i < 257; i += 1) { diff --git a/src/main/ipc/filesystem-watcher.test.ts b/src/main/ipc/filesystem-watcher.test.ts index 2d9181d1155..db6fcf206e5 100644 --- a/src/main/ipc/filesystem-watcher.test.ts +++ b/src/main/ipc/filesystem-watcher.test.ts @@ -1,3 +1,4 @@ +import { createWatcherSender } from './filesystem-watcher-test-sender' import { beforeEach, describe, expect, it, vi } from 'vitest' const { handleMock, getSshFilesystemProviderMock, providerRegistrationListeners } = vi.hoisted( @@ -106,7 +107,7 @@ describe('registerFilesystemWatcherHandlers', () => { vi.mocked(subscribeParcelWatcher).mockResolvedValue({ unsubscribe: vi.fn() } as never) await handlers['fs:watchWorktree']( - { sender: { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } }, + { sender: createWatcherSender(1) }, { worktreePath: 'C:\\repo' } ) @@ -140,7 +141,7 @@ describe('registerFilesystemWatcherHandlers', () => { rootPath: worktreePath } }) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) const args = { worktreePath: '\\\\wsl.localhost\\Ubuntu\\home\\me\\repo' } await expect(handlers['fs:watchWorktree']({ sender }, args)).resolves.toBeUndefined() @@ -162,7 +163,7 @@ describe('registerFilesystemWatcherHandlers', () => { reserveWatcherChild() ) vi.mocked(createWslWatcher).mockRejectedValue(new WatcherChildCapacityError()) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) const args = { worktreePath: '\\\\wsl.localhost\\Ubuntu\\home\\me\\repo' } await handlers['fs:watchWorktree']({ sender }, args) @@ -178,7 +179,7 @@ describe('registerFilesystemWatcherHandlers', () => { it('rejects installs during destructive removal and allows a retry afterward', async () => { vi.mocked(stat).mockResolvedValue({ isDirectory: () => true } as never) vi.mocked(subscribeParcelWatcher).mockResolvedValue({ unsubscribe: vi.fn() } as never) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) const removal = acquireWatcherRemovalGate('/repo') await removal.ready @@ -201,12 +202,12 @@ describe('registerFilesystemWatcherHandlers', () => { await expect( handlers['fs:watchWorktree']( - { sender: { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } }, + { sender: createWatcherSender(1) }, { worktreePath: '/home/me/repo', connectionId: 'conn-1' } ) ).resolves.toBeUndefined() await handlers['fs:watchWorktree']( - { sender: { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } }, + { sender: createWatcherSender(1) }, { worktreePath: '/home/me/repo', connectionId: 'conn-1' } ) @@ -226,7 +227,7 @@ describe('registerFilesystemWatcherHandlers', () => { vi.useFakeTimers() const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) const sendMock = vi.fn() - const sender = { isDestroyed: () => false, send: sendMock, once: vi.fn(), id: 1 } + const sender = createWatcherSender(1, sendMock) const unwatchMock = vi.fn() const watchMock = vi.fn().mockResolvedValue(unwatchMock) getSshFilesystemProviderMock.mockReturnValueOnce(undefined) @@ -261,7 +262,7 @@ describe('registerFilesystemWatcherHandlers', () => { vi.useFakeTimers() const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) const sendMock = vi.fn() - const sender = { isDestroyed: () => false, send: sendMock, once: vi.fn(), id: 1 } + const sender = createWatcherSender(1, sendMock) const unwatchMock = vi.fn() const retryWatchMock = vi.fn().mockResolvedValue(unwatchMock) getSshFilesystemProviderMock @@ -299,7 +300,7 @@ describe('registerFilesystemWatcherHandlers', () => { getSshFilesystemProviderMock.mockReturnValueOnce(undefined) await handlers['fs:watchWorktree']( - { sender: { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } }, + { sender: createWatcherSender(1) }, { worktreePath: '/home/me/repo', connectionId: 'conn-1' } ) @@ -315,8 +316,8 @@ describe('registerFilesystemWatcherHandlers', () => { it('retries all live renderer owners after a terminal relay watch failure', async () => { vi.useFakeTimers() const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) - const senderOne = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } - const senderTwo = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } + const senderOne = createWatcherSender(1) + const senderTwo = createWatcherSender(2) const watchMock = vi.fn().mockResolvedValue(vi.fn()) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) @@ -354,7 +355,7 @@ describe('registerFilesystemWatcherHandlers', () => { }) it('reinstalls an SSH worktree watch when the provider is re-registered after a reconnect', async () => { - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) const staleUnwatch = vi.fn() const watchMock = vi.fn().mockResolvedValue(staleUnwatch) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) @@ -390,7 +391,7 @@ describe('registerFilesystemWatcherHandlers', () => { it('re-arms an SSH watch whose first install found no provider yet', async () => { const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) // A connect slower than the retry window leaves the renderer subscribed with nothing installed. getSshFilesystemProviderMock.mockReturnValue(undefined) @@ -411,7 +412,7 @@ describe('registerFilesystemWatcherHandlers', () => { it('resyncs after a reconnect whose reinstall only succeeded on a retry', async () => { vi.useFakeTimers() const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) getSshFilesystemProviderMock.mockReturnValue({ watch: vi.fn().mockResolvedValue(vi.fn()) }) await handlers['fs:watchWorktree']( @@ -440,7 +441,7 @@ describe('registerFilesystemWatcherHandlers', () => { }) it('does not resurrect an SSH watch the renderer already unwatched', async () => { - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) const watchMock = vi.fn().mockResolvedValue(vi.fn()) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) @@ -462,7 +463,7 @@ describe('registerFilesystemWatcherHandlers', () => { }) it('leaves watches on other connections untouched when one provider re-registers', async () => { - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) const watchMock = vi.fn().mockResolvedValue(vi.fn()) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) @@ -480,8 +481,8 @@ describe('registerFilesystemWatcherHandlers', () => { }) it('reinstalls one shared watch when several senders share a re-registered connection', async () => { - const senderOne = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } - const senderTwo = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } + const senderOne = createWatcherSender(1) + const senderTwo = createWatcherSender(2) const watchMock = vi.fn().mockResolvedValue(vi.fn()) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) @@ -515,8 +516,11 @@ describe('registerFilesystemWatcherHandlers', () => { const sender = { isDestroyed: () => destroyed, send: vi.fn(), - once: vi.fn((_event: string, handler: () => void) => { - destroyHandlers.push(handler) + removeListener: vi.fn(), + once: vi.fn((event: string, handler: () => void) => { + if (event === 'destroyed') { + destroyHandlers.push(handler) + } }), id: 1 } @@ -545,8 +549,8 @@ describe('registerFilesystemWatcherHandlers', () => { it('shares SSH worktree watchers across renderer senders until the last unwatch', async () => { const sendOne = vi.fn() const sendTwo = vi.fn() - const senderOne = { isDestroyed: () => false, send: sendOne, once: vi.fn(), id: 1 } - const senderTwo = { isDestroyed: () => false, send: sendTwo, once: vi.fn(), id: 2 } + const senderOne = createWatcherSender(1, sendOne) + const senderTwo = createWatcherSender(2, sendTwo) const unwatchMock = vi.fn() const watchMock = vi.fn().mockResolvedValue(unwatchMock) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock }) @@ -590,7 +594,7 @@ describe('registerFilesystemWatcherHandlers', () => { const provider = { watch: vi.fn().mockResolvedValue(vi.fn()), closeWatch } getSshFilesystemProviderMock.mockReturnValue(provider) await handlers['fs:watchWorktree']( - { sender: { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } }, + { sender: createWatcherSender(1) }, { worktreePath: '/home/me/repo', connectionId: 'conn-1' } ) @@ -615,7 +619,7 @@ describe('registerFilesystemWatcherHandlers', () => { .mockResolvedValueOnce(replacementUnwatch) const closeWatch = vi.fn().mockResolvedValue(undefined) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock, closeWatch }) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']( { sender }, @@ -642,7 +646,7 @@ describe('registerFilesystemWatcherHandlers', () => { const watchMock = vi.fn().mockResolvedValue(vi.fn()) const closeWatch = vi.fn().mockResolvedValue(undefined) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock, closeWatch }) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) await handlers['fs:watchWorktree']( { sender }, @@ -666,7 +670,7 @@ describe('registerFilesystemWatcherHandlers', () => { const watchMock = vi.fn().mockResolvedValue(firstUnwatch) const closeWatch = vi.fn().mockResolvedValue(undefined) getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock, closeWatch }) - const sender = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } + const sender = createWatcherSender(1) const args = { worktreePath: '/home/me/repo', connectionId: 'conn-1' } await handlers['fs:watchWorktree']({ sender }, args) @@ -685,14 +689,12 @@ describe('registerFilesystemWatcherHandlers', () => { getSshFilesystemProviderMock.mockReturnValue({ watch: watchMock, closeWatch }) const destroyedCallbacks: (() => void)[] = [] const sender = { - isDestroyed: () => false, - send: vi.fn(), + ...createWatcherSender(1), once: vi.fn((event: string, callback: () => void) => { if (event === 'destroyed') { destroyedCallbacks.push(callback) } - }), - id: 1 + }) } const args = { worktreePath: '/home/me/repo', connectionId: 'conn-1' } @@ -709,8 +711,8 @@ describe('registerFilesystemWatcherHandlers', () => { it('preserves remote event routing when acknowledged teardown rejects', async () => { const sendOne = vi.fn() const sendTwo = vi.fn() - const senderOne = { isDestroyed: () => false, send: sendOne, once: vi.fn(), id: 1 } - const senderTwo = { isDestroyed: () => false, send: sendTwo, once: vi.fn(), id: 2 } + const senderOne = createWatcherSender(1, sendOne) + const senderTwo = createWatcherSender(2, sendTwo) const watchMock = vi.fn().mockResolvedValue(vi.fn()) const closeWatch = vi .fn() @@ -747,8 +749,8 @@ describe('registerFilesystemWatcherHandlers', () => { it('dedupes concurrent pending SSH worktree watcher installs', async () => { const sendOne = vi.fn() const sendTwo = vi.fn() - const senderOne = { isDestroyed: () => false, send: sendOne, once: vi.fn(), id: 1 } - const senderTwo = { isDestroyed: () => false, send: sendTwo, once: vi.fn(), id: 2 } + const senderOne = createWatcherSender(1, sendOne) + const senderTwo = createWatcherSender(2, sendTwo) const unwatchMock = vi.fn() let resolveWatch: (unwatch: () => void) => void = () => {} const watchPromise = new Promise<() => void>((resolve) => { @@ -791,8 +793,8 @@ describe('registerFilesystemWatcherHandlers', () => { it('keeps a pending SSH watcher install alive when only one pending sender unwatches', async () => { const sendOne = vi.fn() const sendTwo = vi.fn() - const senderOne = { isDestroyed: () => false, send: sendOne, once: vi.fn(), id: 1 } - const senderTwo = { isDestroyed: () => false, send: sendTwo, once: vi.fn(), id: 2 } + const senderOne = createWatcherSender(1, sendOne) + const senderTwo = createWatcherSender(2, sendTwo) const unwatchMock = vi.fn() let resolveWatch: (unwatch: () => void) => void = () => {} const watchPromise = new Promise<() => void>((resolve) => { @@ -833,14 +835,12 @@ describe('registerFilesystemWatcherHandlers', () => { it('unsubscribes if the sender is destroyed while an SSH watcher is opening', async () => { const destroyedCallbacks: (() => void)[] = [] const sender = { - isDestroyed: () => false, - send: vi.fn(), + ...createWatcherSender(1), once: vi.fn((event: string, callback: () => void) => { if (event === 'destroyed') { destroyedCallbacks.push(callback) } - }), - id: 1 + }) } const unwatchMock = vi.fn() let resolveWatch: (unwatch: () => void) => void = () => {} @@ -871,8 +871,8 @@ describe('registerFilesystemWatcherHandlers', () => { it('revives a pending SSH watcher install when a new sender joins after cancellation', async () => { const args = { worktreePath: '/home/me/repo', connectionId: 'conn-1' } - const senderOne = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 1 } - const senderTwo = { isDestroyed: () => false, send: vi.fn(), once: vi.fn(), id: 2 } + const senderOne = createWatcherSender(1) + const senderTwo = createWatcherSender(2) const unwatchMock = vi.fn() let resolveWatch!: (unwatch: () => void) => void const watchMock = vi.fn().mockReturnValue( @@ -907,14 +907,12 @@ describe('registerFilesystemWatcherHandlers', () => { it('registers one destroyed listener for many SSH worktree watches', async () => { const destroyedCallbacks: (() => void)[] = [] const sender = { - isDestroyed: () => false, - send: vi.fn(), + ...createWatcherSender(99), once: vi.fn((event: string, callback: () => void) => { if (event === 'destroyed') { destroyedCallbacks.push(callback) } - }), - id: 99 + }) } const unwatchMock = vi.fn() const watchMock = vi.fn().mockResolvedValue(unwatchMock) @@ -929,7 +927,7 @@ describe('registerFilesystemWatcherHandlers', () => { // Why: WebContents warns after 10 listeners. The cleanup work still covers // every remote watch by scanning the shared remote watcher registry. - expect(sender.once).toHaveBeenCalledTimes(1) + expect(sender.once).toHaveBeenCalledTimes(3) expect(destroyedCallbacks).toHaveLength(1) destroyedCallbacks[0]()