fix: release file watches when their renderer document ends (#23055)

This commit is contained in:
Neil
2026-09-26 14:29:30 -07:00
committed by GitHub
parent b42de170b1
commit f7abecde01
31 changed files with 715 additions and 239 deletions
+83 -17
View File
@@ -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",
@@ -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(
@@ -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<string, (event: { sender: Sender }, args: WatchArgs) => Promise<void>>()
function invoke(channel: string, sender: Sender, args: WatchArgs): Promise<void> {
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<T>() {
let resolve!: (value: T) => void
const promise = new Promise<T>((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<ReturnType<typeof localRoot> & (() => 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<ReturnType<typeof localRoot> & (() => 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<ReturnType<typeof localRoot> & (() => 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()
})
})
@@ -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<typeof vi.fn>
removeListener: ReturnType<typeof vi.fn>
once: ReturnType<typeof vi.fn>
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', () => {
@@ -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<void> => {
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.
@@ -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.
@@ -65,7 +65,7 @@ export const watcherLifecycleState = {
watchedRoots: new Map<string, WatchedRoot>(),
unwatchableRoots: new Set<string>(),
// Why: key cleanup by sender WebContents (not per root) to avoid MaxListeners warnings when a workspace has many worktrees open.
senderCleanupRegistered: new Set<number>(),
senderLifetimes: new Map<number, { signal: AbortSignal; dispose: () => void }>(),
pendingTeardowns: new Map<string, ReturnType<typeof setTimeout>>(),
// 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<Promise<void>>(),
@@ -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) {
@@ -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' }])
@@ -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<number, WebContents>()
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
}
}
@@ -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<void> {
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<void> {
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)
@@ -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 },
@@ -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)
+3 -7
View File
@@ -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')
@@ -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 () => {
@@ -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<unknown>
const second = handlers['fs:watchWorktree']({ sender: senderTwo }, args) as Promise<unknown>
@@ -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<unknown>
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<unknown>
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<unknown>
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<unknown>
await Promise.resolve()
@@ -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)
@@ -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
)
}
@@ -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<ReturnType<InstallRemoteWatcher>>[]
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')) {
@@ -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<RemoteWatcherInstallResult> {
// 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<RemoteWatcherInstallResult> {
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
@@ -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,
@@ -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()) })
@@ -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') {
@@ -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,
@@ -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()
@@ -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
)
}
+4 -1
View File
@@ -8,7 +8,10 @@ export async function closeAllWatchers(): Promise<void> {
// 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()
@@ -41,6 +41,7 @@ const OVERFLOW_PAYLOAD = {
type MockSender = {
isDestroyed: () => boolean
send: ReturnType<typeof vi.fn>
removeListener: ReturnType<typeof vi.fn>
once: ReturnType<typeof vi.fn>
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)
@@ -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() }
}
@@ -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) {
+45 -47
View File
@@ -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]()