diff --git a/src/main/ipc/parcel-watcher-child-termination.ts b/src/main/ipc/parcel-watcher-child-termination.ts index fc498410438..e7067dc32d6 100644 --- a/src/main/ipc/parcel-watcher-child-termination.ts +++ b/src/main/ipc/parcel-watcher-child-termination.ts @@ -10,6 +10,13 @@ export const WATCHER_PROCESS_HARD_KILL_DELAY_MS = 5_000 export const WATCHER_PROCESS_EXIT_DEADLINE_MS = RUNTIME_FILE_WATCH_EXIT_DEADLINE_MS const physicalExitPromises = new WeakMap>() +const signalledChildren = new WeakSet() + +/** Sends the graceful signal once; a later awaited termination only escalates and waits. */ +export function signalWatcherChild(child: ChildProcess): void { + signalledChildren.add(child) + child.kill() +} export function registerWatcherChildPhysicalExit(child: ChildProcess): () => void { let resolveExit: () => void = () => undefined @@ -86,6 +93,10 @@ export function terminateWatcherChild(child: ChildProcess): Promise { hardKillTimer.unref?.() const exitDeadlineTimer = setTimeout(() => finish(false), WATCHER_PROCESS_EXIT_DEADLINE_MS) exitDeadlineTimer.unref?.() + if (signalledChildren.has(child)) { + return + } + signalledChildren.add(child) try { child.kill() } catch { @@ -94,11 +105,10 @@ export function terminateWatcherChild(child: ChildProcess): Promise { }) } -export function createWatcherChildTerminationFailure(child: ChildProcess): WatcherProcessFailure { - const physicalExit = - child.exitCode !== null || child.signalCode !== null - ? Promise.resolve() - : (physicalExitPromises.get(child) ?? +export function watcherChildPhysicalExit(child: ChildProcess): Promise { + return child.exitCode !== null || child.signalCode !== null + ? Promise.resolve() + : (physicalExitPromises.get(child) ?? new Promise((resolve) => { const finish = (): void => { child.removeListener('exit', finish) @@ -108,11 +118,14 @@ export function createWatcherChildTerminationFailure(child: ChildProcess): Watch child.once('exit', finish) child.once('close', finish) })) +} + +export function createWatcherChildTerminationFailure(child: ChildProcess): WatcherProcessFailure { return new WatcherProcessFailure( 'file watcher process did not exit after termination deadline', 'supervisor', 'process_unavailable', - physicalExit + watcherChildPhysicalExit(child) ) } @@ -127,10 +140,12 @@ export async function terminateIdleWatcherChild( pendingUnsubscribes: Map, onFinished: (exited: boolean) => void ): Promise { + // Windows directory handles require physical exit, not merely an accepted signal. try { await requireWatcherChildTermination(child) onFinished(true) } catch (error) { + // Idle children retain capacity but cannot double-watch; the owner may remain reusable. onFinished(false) resolvePendingWatcherUnsubscribes( pendingUnsubscribes, diff --git a/src/main/ipc/parcel-watcher-entry-path.ts b/src/main/ipc/parcel-watcher-entry-path.ts index 229dcdc65e2..4633739fe15 100644 --- a/src/main/ipc/parcel-watcher-entry-path.ts +++ b/src/main/ipc/parcel-watcher-entry-path.ts @@ -4,6 +4,14 @@ import { join } from 'node:path' type ElectronAppPath = { getAppPath(): string; isPackaged(): boolean } +export function watcherProcessEntryExists(entryPath: string): boolean { + if (existsSync(entryPath)) { + return true + } + console.error(`[parcel-watcher-process] entry not found at ${entryPath}; refusing fail-open`) + return false +} + // Why the port and not require('electron'): this module is reachable from plain-Node // fork entries, where the literal text require("electron") fails the build guard even // inside a try/catch. hasAppEnvironment() gives the same "no app root here" answer. diff --git a/src/main/ipc/parcel-watcher-owned-children.test.ts b/src/main/ipc/parcel-watcher-owned-children.test.ts new file mode 100644 index 00000000000..cbb369cf3e2 --- /dev/null +++ b/src/main/ipc/parcel-watcher-owned-children.test.ts @@ -0,0 +1,120 @@ +import { EventEmitter } from 'node:events' +import type { ChildProcess } from 'node:child_process' +import { afterEach, expect, it, vi } from 'vitest' +import { WatcherOwnedChildren } from './parcel-watcher-owned-children' +import { + registerWatcherChildPhysicalExit, + signalWatcherChild, + WATCHER_PROCESS_EXIT_DEADLINE_MS, + WATCHER_PROCESS_HARD_KILL_DELAY_MS +} from './parcel-watcher-child-termination' + +afterEach(() => vi.useRealTimers()) + +function child() { + const exitState: { exitCode: number | null; signalCode: NodeJS.Signals | null } = { + exitCode: null, + signalCode: null + } + const events = Object.assign(new EventEmitter(), exitState, { kill: vi.fn(() => true) }) + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: termination reads only the exit fields, kill and events stubbed here. + const process = events as unknown as ChildProcess + const physicalExit = registerWatcherChildPhysicalExit(process) + events.on('exit', physicalExit) + events.on('close', physicalExit) + events.on('error', () => {}) + return { process, events, close: () => events.emit('close') } +} + +it('joins disposal and does not confuse disconnect/error with physical exit', async () => { + const owner = new WatcherOwnedChildren() + const c = child() + owner.track(c.process) + const logical = vi.fn() + const disposed = vi.fn() + const first = owner.disposeAndWait(logical) + expect(owner.disposeAndWait(logical)).toBe(first) + const result = first.then(disposed) + c.events.emit('disconnect') + c.events.emit('error', new Error('spawn or IPC failure')) + await Promise.resolve() + expect(disposed).not.toHaveBeenCalled() + expect(logical).toHaveBeenCalledOnce() + c.close() + await result + expect(disposed).toHaveBeenCalledOnce() +}) + +it('includes children already signaled by synchronous or retired-owner disposal', async () => { + const owner = new WatcherOwnedChildren() + const old = child() + const replacement = child() + owner.track(old.process) + old.process.kill() + owner.track(replacement.process) + const disposed = vi.fn() + const result = owner.disposeAndWait(() => {}).then(disposed) + replacement.close() + await Promise.resolve() + expect(disposed).not.toHaveBeenCalled() + old.close() + await result +}) + +it('retains a child after termination deadline failure and supports explicit retry', async () => { + vi.useFakeTimers() + const owner = new WatcherOwnedChildren() + const c = child() + owner.track(c.process) + const result = owner.disposeAndWait(() => {}).catch((error: unknown) => error) + await vi.advanceTimersByTimeAsync(WATCHER_PROCESS_EXIT_DEADLINE_MS) + expect(await result).toMatchObject({ message: 'watcher_owned_children_shutdown_incomplete' }) + const completed = vi.fn() + const retry = owner.disposeAndWait(() => {}).then(completed) + await Promise.resolve() + expect(completed).not.toHaveBeenCalled() + c.close() + await retry +}) + +it('waits for other children even when logical disposal and one kill fail', async () => { + const owner = new WatcherOwnedChildren() + const failed = child() + const pending = child() + owner.track(failed.process) + owner.track(pending.process) + failed.events.kill.mockImplementationOnce(() => { + throw new Error('kill failed') + }) + const logicalFailure = new Error('logical cleanup failed') + const finished = vi.fn() + const result = owner + .disposeAndWait(() => { + throw logicalFailure + }) + .catch((error: unknown) => { + finished() + return error + }) + await Promise.resolve() + expect(finished).not.toHaveBeenCalled() + pending.close() + expect(await result).toMatchObject({ errors: [logicalFailure, expect.any(Error)] }) + const retry = owner.disposeAndWait(() => {}) + failed.close() + await retry +}) + +it('skips a second graceful signal after the supervisor sent one but still escalates', async () => { + vi.useFakeTimers() + const owner = new WatcherOwnedChildren() + const c = child() + owner.track(c.process) + const disposal = owner.disposeAndWait(() => signalWatcherChild(c.process)) + expect(c.events.kill).toHaveBeenCalledTimes(1) + expect(c.events.kill).toHaveBeenCalledWith() + await vi.advanceTimersByTimeAsync(WATCHER_PROCESS_HARD_KILL_DELAY_MS) + expect(c.events.kill).toHaveBeenLastCalledWith('SIGKILL') + c.events.emit('exit', null, 'SIGKILL') + await expect(disposal).resolves.toBeUndefined() +}) diff --git a/src/main/ipc/parcel-watcher-owned-children.ts b/src/main/ipc/parcel-watcher-owned-children.ts new file mode 100644 index 00000000000..7ac044ac906 --- /dev/null +++ b/src/main/ipc/parcel-watcher-owned-children.ts @@ -0,0 +1,47 @@ +import type { ChildProcessHandle } from '../../shared/child-process/process-spec' +import { + watcherChildPhysicalExit, + requireWatcherChildTermination +} from './parcel-watcher-child-termination' + +export class WatcherOwnedChildren { + private readonly children = new Set() + private disposal: Promise | null = null + + track(child: ChildProcessHandle): ChildProcessHandle { + this.children.add(child) + // Reuse launch-owned physical-exit evidence, including close without exit after spawn failure. + void watcherChildPhysicalExit(child).then(() => { + this.children.delete(child) + }) + return child + } + + disposeAndWait(disposeLogicalOwner: () => void): Promise { + if (this.disposal) { + return this.disposal + } + const failures: unknown[] = [] + try { + disposeLogicalOwner() + } catch (error) { + failures.push(error) + } + const cleanup = Promise.allSettled([...this.children].map(requireWatcherChildTermination)) + .then((results) => { + for (const result of results) { + if (result.status === 'rejected') { + failures.push(result.reason) + } + } + if (failures.length > 0) { + throw new AggregateError(failures, 'watcher_owned_children_shutdown_incomplete') + } + }) + .finally(() => { + this.disposal = null + }) + this.disposal = cleanup + return cleanup + } +} diff --git a/src/main/ipc/parcel-watcher-process-entry.test.ts b/src/main/ipc/parcel-watcher-process-entry.test.ts index 6ebdb5385dd..513775acb29 100644 --- a/src/main/ipc/parcel-watcher-process-entry.test.ts +++ b/src/main/ipc/parcel-watcher-process-entry.test.ts @@ -86,8 +86,13 @@ describe('parcel watcher process canary', () => { await vi.advanceTimersByTimeAsync(0) expect(sendMock).toHaveBeenCalledWith( - expect.objectContaining({ op: 'subscribe-failed', id: 7 }) + expect.objectContaining({ op: 'subscribe-failed', id: 7 }), + expect.any(Function) ) + const failedSend = sendMock.mock.calls.find(([message]) => message.op === 'subscribe-failed') + expect(() => + failedSend![1](Object.assign(new Error('host disconnected'), { code: 'EPIPE' })) + ).not.toThrow() expect(watchMock).not.toHaveBeenCalledWith('/repo/.git', expect.anything(), expect.anything()) }) @@ -313,10 +318,13 @@ describe('parcel watcher process canary', () => { finishActiveCrawl?.({ unsubscribe: vi.fn().mockResolvedValue(undefined) }) await vi.advanceTimersByTimeAsync(0) - expect(sendMock).toHaveBeenCalledWith({ op: 'unsubscribed', id: 2 }) + expect(sendMock).toHaveBeenCalledWith({ op: 'unsubscribed', id: 2 }, expect.any(Function)) expect(subscribeMock).toHaveBeenCalledTimes(2) - expect(sendMock).not.toHaveBeenCalledWith({ op: 'subscribe-started', id: 2 }) - expect(sendMock).not.toHaveBeenCalledWith({ op: 'subscribed', id: 2 }) + expect(sendMock).not.toHaveBeenCalledWith( + { op: 'subscribe-started', id: 2 }, + expect.any(Function) + ) + expect(sendMock).not.toHaveBeenCalledWith({ op: 'subscribed', id: 2 }, expect.any(Function)) }) it('unsubscribes a late cancel after the crawl already finished', async () => { @@ -332,13 +340,16 @@ describe('parcel watcher process canary', () => { process.emit('message', { op: 'subscribe', id: 1, dir: '/finished', opts: {} }) await vi.advanceTimersByTimeAsync(0) - expect(sendMock).toHaveBeenCalledWith({ op: 'subscribed', id: 1 }) + expect(sendMock).toHaveBeenCalledWith({ op: 'subscribed', id: 1 }, expect.any(Function)) process.emit('message', { op: 'cancel-subscribe', id: 1 }) await vi.advanceTimersByTimeAsync(0) expect(unsubscribe).toHaveBeenCalledTimes(1) - expect(sendMock).toHaveBeenCalledWith({ op: 'unsubscribed', id: 1 }) - expect(sendMock).not.toHaveBeenCalledWith({ op: 'cancel-requires-restart', id: 1 }) + expect(sendMock).toHaveBeenCalledWith({ op: 'unsubscribed', id: 1 }, expect.any(Function)) + expect(sendMock).not.toHaveBeenCalledWith( + { op: 'cancel-requires-restart', id: 1 }, + expect.any(Function) + ) }) it('reports native unsubscribe rejection without acknowledging handle release', async () => { @@ -356,12 +367,15 @@ describe('parcel watcher process canary', () => { process.emit('message', { op: 'unsubscribe', id: 1 }) await vi.advanceTimersByTimeAsync(0) - expect(sendMock).toHaveBeenCalledWith({ - op: 'unsubscribe-failed', - id: 1, - message: 'native handle still active' - }) - expect(sendMock).not.toHaveBeenCalledWith({ op: 'unsubscribed', id: 1 }) + expect(sendMock).toHaveBeenCalledWith( + { + op: 'unsubscribe-failed', + id: 1, + message: 'native handle still active' + }, + expect.any(Function) + ) + expect(sendMock).not.toHaveBeenCalledWith({ op: 'unsubscribed', id: 1 }, expect.any(Function)) }) it('asks the host to restart when an active crawl is cancelled', async () => { @@ -380,10 +394,13 @@ describe('parcel watcher process canary', () => { await vi.advanceTimersByTimeAsync(0) process.emit('message', { op: 'cancel-subscribe', id: 1 }) - expect(sendMock).toHaveBeenCalledWith({ op: 'cancel-requires-restart', id: 1 }) + expect(sendMock).toHaveBeenCalledWith( + { op: 'cancel-requires-restart', id: 1 }, + expect.any(Function) + ) finishCrawl?.({ unsubscribe: vi.fn().mockResolvedValue(undefined) }) await vi.advanceTimersByTimeAsync(0) - expect(sendMock).not.toHaveBeenCalledWith({ op: 'subscribed', id: 1 }) + expect(sendMock).not.toHaveBeenCalledWith({ op: 'subscribed', id: 1 }, expect.any(Function)) }) it('still restarts after consecutive missed events once every subscription is live', async () => { @@ -481,11 +498,14 @@ describe('parcel watcher process canary', () => { callback?.(null, [{ type: 'update', path: '/repo/after-overflow.txt' }]) await vi.advanceTimersByTimeAsync(0) - expect(sendMock).toHaveBeenCalledWith({ - op: 'watch-error', - id: 1, - message: 'Events were dropped by the FSEvents client. File system must be re-scanned.' - }) + expect(sendMock).toHaveBeenCalledWith( + { + op: 'watch-error', + id: 1, + message: 'Events were dropped by the FSEvents client. File system must be re-scanned.' + }, + expect.any(Function) + ) expect(sendMock).toHaveBeenCalledWith( { op: 'events', diff --git a/src/main/ipc/parcel-watcher-process-entry.ts b/src/main/ipc/parcel-watcher-process-entry.ts index 0de05010f47..d2328821ade 100644 --- a/src/main/ipc/parcel-watcher-process-entry.ts +++ b/src/main/ipc/parcel-watcher-process-entry.ts @@ -118,7 +118,7 @@ async function startCanary(getStableActivityRevision: () => number | null): Prom function main(): void { const send = (message: WatcherToHostMessage): void => { try { - process.send?.(message) + process.send?.(message, () => undefined) } catch { // Host is gone; the disconnect handler below exits this process. } diff --git a/src/main/ipc/parcel-watcher-process-supervisor.ts b/src/main/ipc/parcel-watcher-process-supervisor.ts index 3f0a534c198..926ff036466 100644 --- a/src/main/ipc/parcel-watcher-process-supervisor.ts +++ b/src/main/ipc/parcel-watcher-process-supervisor.ts @@ -1,8 +1,7 @@ import type { ChildProcess } from 'node:child_process' -import { existsSync } from 'node:fs' import { restartCancelledWatcherChild } from './parcel-watcher-cancellation-restart' import { WatcherCancellationTracker } from './parcel-watcher-cancellation-tracker' -import { getWatcherProcessEntryPath } from './parcel-watcher-entry-path' +import { getWatcherProcessEntryPath, watcherProcessEntryExists } from './parcel-watcher-entry-path' import { removeWatcherCanaryDirectory } from './parcel-watcher-canary-directory' import * as termination from './parcel-watcher-child-termination' import { launchWatcherChild } from './parcel-watcher-child-launch' @@ -39,6 +38,7 @@ import { } from './parcel-watcher-supervisor-subscribe' import { disposeWatcherSupervisor } from './parcel-watcher-supervisor-disposal' import { handleWatcherSupervisorMessage } from './parcel-watcher-supervisor-message' +import { WatcherOwnedChildren } from './parcel-watcher-owned-children' export class WatcherProcessSupervisor { private child: ChildProcess | null = null @@ -52,6 +52,7 @@ export class WatcherProcessSupervisor { private readonly pendingUnsubscribes = new Map() private readonly cancelledSubscribes = new WatcherCancellationTracker() private readonly capacityWait = new WatcherSupervisorCapacityWait() + private ownedChildren = new WatcherOwnedChildren() constructor(private readonly options: WatcherProcessSupervisorOptions = {}) {} @@ -107,6 +108,7 @@ export class WatcherProcessSupervisor { resetForTest(): void { this.dispose() + this.ownedChildren = new WatcherOwnedChildren() this.shutdownRequested = false this.terminatingChild = null this.terminationQueue.resetForTest() @@ -114,6 +116,8 @@ export class WatcherProcessSupervisor { resetWatcherChildRegistryForTest() } + disposeAndWait = (): Promise => this.ownedChildren.disposeAndWait(() => this.dispose()) + private ensureWatcherProcess( entryPath = this.options.entryPath ?? getWatcherProcessEntryPath() ): ChildProcess | null { @@ -126,8 +130,7 @@ export class WatcherProcessSupervisor { if (this.crashFuse.isOpen()) { return null } - if (!existsSync(entryPath)) { - console.error(`[parcel-watcher-process] entry not found at ${entryPath}; refusing fail-open`) + if (!watcherProcessEntryExists(entryPath)) { return null } const launched = launchWatcherChild( @@ -145,7 +148,7 @@ export class WatcherProcessSupervisor { return null } this.canaryDir = launched.canaryDir - this.child = launched.child + this.child = this.ownedChildren.track(launched.child) return launched.child } @@ -251,14 +254,9 @@ export class WatcherProcessSupervisor { } this.child = null this.terminatingChild = proc - // Why: destructive Windows cleanup must await exit to release directory handles. this.canaryDir = removeWatcherCanaryDirectory(this.canaryDir) return this.terminationQueue.track( termination.terminateIdleWatcherChild(proc, this.pendingUnsubscribes, () => { - // Why: an idle child owns zero records, so a missed exit deadline has no - // double-watch hazard; the child keeps its capacity reservation until - // physical exit, and poisoning this supervisor would permanently end - // local watching — the shared singleton has no retire-and-replace path. this.terminatingChild = null }) ) diff --git a/src/main/ipc/parcel-watcher-supervisor-disposal.test.ts b/src/main/ipc/parcel-watcher-supervisor-disposal.test.ts new file mode 100644 index 00000000000..d3aaa15770c --- /dev/null +++ b/src/main/ipc/parcel-watcher-supervisor-disposal.test.ts @@ -0,0 +1,67 @@ +import { afterEach, beforeEach, expect, it, vi } from 'vitest' +import { acknowledgeWatcherSubscribe, FakeWatcherChild } from './parcel-watcher-process-test-child' +import { resetWatcherChildRegistryForTest } from './parcel-watcher-child-registry' + +const { forkMock } = vi.hoisted(() => ({ forkMock: vi.fn() })) +vi.mock('node:child_process', () => ({ fork: forkMock })) +vi.mock('node:fs', () => ({ + existsSync: vi.fn(() => true), + mkdtempSync: vi.fn(() => '/tmp/orca-watcher-disposal-test'), + rmSync: vi.fn() +})) +import { WatcherProcessSupervisor } from './parcel-watcher-process-supervisor' + +const children: FakeWatcherChild[] = [] +beforeEach(() => { + resetWatcherChildRegistryForTest() + forkMock.mockImplementation(() => { + const child = new FakeWatcherChild() + children.push(child) + return child + }) +}) +afterEach(() => { + for (const child of children.splice(0)) { + child.emit('close') + } + vi.clearAllMocks() + resetWatcherChildRegistryForTest() +}) + +async function subscribe(supervisor: WatcherProcessSupervisor): Promise { + const pending = supervisor.subscribe('/repo', vi.fn(), {}) + const child = children.at(-1)! + acknowledgeWatcherSubscribe(child) + await pending + return child +} + +it('awaited supervisor disposal retains a child after prior synchronous disposal', async () => { + const supervisor = new WatcherProcessSupervisor({ useInProcessVitestFallback: false }) + const child = await subscribe(supervisor) + supervisor.dispose() + const finished = vi.fn() + const shutdown = supervisor.disposeAndWait().then(finished) + await new Promise((resolve) => setImmediate(resolve)) + expect(finished).not.toHaveBeenCalled() + await expect(supervisor.subscribe('/another', vi.fn(), {})).rejects.toThrow() + expect(forkMock).toHaveBeenCalledOnce() + child.emit('close') + await shutdown + expect(finished).toHaveBeenCalledOnce() +}) + +it('test reset isolates old physical-exit callbacks from the new lifetime owner', async () => { + const supervisor = new WatcherProcessSupervisor({ useInProcessVitestFallback: false }) + const old = await subscribe(supervisor) + const oldShutdown = supervisor.disposeAndWait() + supervisor.resetForTest() + const current = await subscribe(supervisor) + const finished = vi.fn() + const shutdown = supervisor.disposeAndWait().then(finished) + old.emit('close') + await oldShutdown + expect(finished).not.toHaveBeenCalled() + current.emit('close') + await shutdown +}) diff --git a/src/main/ipc/parcel-watcher-supervisor-disposal.ts b/src/main/ipc/parcel-watcher-supervisor-disposal.ts index 42dc918e05b..07bc73ef3ee 100644 --- a/src/main/ipc/parcel-watcher-supervisor-disposal.ts +++ b/src/main/ipc/parcel-watcher-supervisor-disposal.ts @@ -9,6 +9,7 @@ import { resetPendingSubscribeAttempt, takePendingSubscribe } from './parcel-watcher-pending-subscribe' +import { signalWatcherChild } from './parcel-watcher-child-termination' import { watcherHostFailure } from './parcel-watcher-process-failure' import type { WatcherProcessSubscriptionRecord } from './parcel-watcher-process-subscription' @@ -27,6 +28,8 @@ export function disposeWatcherSupervisor( resolvePendingWatcherUnsubscribes(pendingUnsubscribes) cancelledSubscribes.completeAll() records.clear() - child?.kill() + if (child) { + signalWatcherChild(child) + } return removeWatcherCanaryDirectory(canaryDir) } diff --git a/src/main/ipc/runtime-watcher-disposal-owners.test.ts b/src/main/ipc/runtime-watcher-disposal-owners.test.ts new file mode 100644 index 00000000000..a889db510a4 --- /dev/null +++ b/src/main/ipc/runtime-watcher-disposal-owners.test.ts @@ -0,0 +1,103 @@ +import { expect, it, vi } from 'vitest' +import { RuntimeWatcherDisposalOwners } from './runtime-watcher-disposal-owners' + +function owner() { + return { + dispose: vi.fn(), + subscribe: vi.fn(), + disposeAndWait: vi.fn(async () => {}) + } +} + +it('retains retired cleanup and joins concurrent shutdown callers', async () => { + const owners = new RuntimeWatcherDisposalOwners() + const retired = owner() + const pending = Promise.withResolvers() + retired.disposeAndWait.mockReturnValue(pending.promise) + owners.retire(retired) + const shutdown = owners.disposeAndWait(() => {}) + expect(owners.disposeAndWait(() => {})).toBe(shutdown) + const finished = vi.fn() + const result = shutdown.then(finished) + await Promise.resolve() + expect(finished).not.toHaveBeenCalled() + expect(retired.disposeAndWait).toHaveBeenCalledOnce() + pending.resolve() + await result + await owners.disposeAndWait(() => {}) + expect(retired.disposeAndWait).toHaveBeenCalledOnce() +}) + +it('does not retry a newly failed attempt until a later explicit shutdown call', async () => { + const owners = new RuntimeWatcherDisposalOwners() + const failed = owner() + const sibling = owner() + const pending = Promise.withResolvers() + const failure = new Error('child still live') + failed.disposeAndWait.mockRejectedValueOnce(failure) + sibling.disposeAndWait.mockReturnValue(pending.promise) + const finished = vi.fn() + const shutdown = owners.disposeAndWait(() => { + owners.retire(failed) + owners.retire(sibling) + }) + const result = shutdown.catch((error: unknown) => { + finished() + return error + }) + await new Promise((resolve) => setImmediate(resolve)) + expect(finished).not.toHaveBeenCalled() + expect(owners.disposeAndWait(() => {})).toBe(shutdown) + expect(failed.disposeAndWait).toHaveBeenCalledOnce() + pending.resolve() + expect(await result).toMatchObject({ errors: [failure] }) + await owners.disposeAndWait(() => {}) + expect(failed.disposeAndWait).toHaveBeenCalledTimes(2) + expect(sibling.disposeAndWait).toHaveBeenCalledOnce() +}) + +it('retains synchronous failures and attempts sibling owners', async () => { + const owners = new RuntimeWatcherDisposalOwners() + const failed = owner() + const sibling = owner() + failed.disposeAndWait.mockImplementationOnce(() => { + throw new Error('sync failure') + }) + await expect( + owners.disposeAndWait(() => { + owners.retire(failed) + owners.retire(sibling) + }) + ).rejects.toThrow('watcher_pool_shutdown_incomplete') + expect(sibling.disposeAndWait).toHaveBeenCalledOnce() + await owners.disposeAndWait(() => {}) + expect(failed.disposeAndWait).toHaveBeenCalledTimes(2) +}) + +it('publishes the retained attempt before a synchronous disposal callback reenters', async () => { + const owners = new RuntimeWatcherDisposalOwners() + const child = owner() + const pending = Promise.withResolvers() + let nested: Promise | undefined + child.disposeAndWait.mockImplementation(() => { + owners.retire(child) + nested = owners.disposeAndWait(() => {}) + return pending.promise + }) + const shutdown = owners.disposeAndWait(() => owners.retire(child)) + expect(nested).toBe(shutdown) + expect(child.disposeAndWait).toHaveBeenCalledOnce() + pending.resolve() + await shutdown +}) + +it('refuses to acknowledge a supervisor without a physical-disposal API', async () => { + const owners = new RuntimeWatcherDisposalOwners() + const legacy = { dispose: vi.fn(), subscribe: vi.fn() } + await expect(owners.disposeAndWait(() => owners.retire(legacy))).rejects.toMatchObject({ + errors: [ + expect.objectContaining({ message: 'watcher_supervisor_awaited_disposal_unavailable' }) + ] + }) + expect(legacy.dispose).toHaveBeenCalledOnce() +}) diff --git a/src/main/ipc/runtime-watcher-disposal-owners.ts b/src/main/ipc/runtime-watcher-disposal-owners.ts new file mode 100644 index 00000000000..e88c741844f --- /dev/null +++ b/src/main/ipc/runtime-watcher-disposal-owners.ts @@ -0,0 +1,90 @@ +import type { RuntimeWatcherPoolSupervisor } from './runtime-watcher-pool-state' + +type Attempt = { pending: Promise | null; failure?: unknown } + +type Completion = { promise: Promise; resolve: () => void; reject: (error: unknown) => void } + +// Why not Promise.withResolvers: the relay bundles this pool and still targets Node 18 hosts. +function createCompletion(): Completion { + let resolve!: () => void + let reject!: (error: unknown) => void + const promise = new Promise((res, rej) => { + resolve = res + reject = rej + }) + return { promise, resolve, reject } +} + +export class RuntimeWatcherDisposalOwners { + private readonly retained = new Map() + private shutdown: Promise | null = null + + retire(owner: RuntimeWatcherPoolSupervisor): void { + if (!this.retained.has(owner)) { + this.start(owner) + } + } + + disposeAndWait(disposePool: () => void): Promise { + if (this.shutdown) { + return this.shutdown + } + const retry = [...this.retained].filter(([, attempt]) => !attempt.pending) + const completion = createCompletion() + this.shutdown = completion.promise + const failures: unknown[] = [] + try { + disposePool() + } catch (error) { + failures.push(error) + } + for (const [owner, attempt] of retry) { + if (this.retained.get(owner) === attempt) { + this.start(owner) + } + } + const pending = [...this.retained.values()].map( + (attempt) => attempt.pending ?? Promise.reject(attempt.failure) + ) + void Promise.allSettled(pending).then((results) => { + for (const result of results) { + if (result.status === 'rejected') { + failures.push(result.reason) + } + } + this.shutdown = null + if (failures.length > 0) { + completion.reject(new AggregateError(failures, 'watcher_pool_shutdown_incomplete')) + } else { + completion.resolve() + } + }) + return completion.promise + } + + private start(owner: RuntimeWatcherPoolSupervisor): void { + const completion = createCompletion() + const attempt: Attempt = { pending: completion.promise } + this.retained.set(owner, attempt) + void completion.promise.catch(() => {}) + const failed = (error: unknown): void => { + attempt.pending = null + attempt.failure = error + completion.reject(error) + } + try { + if (!owner.disposeAndWait) { + owner.dispose() + throw new Error('watcher_supervisor_awaited_disposal_unavailable') + } + void owner.disposeAndWait().then(() => { + if (this.retained.get(owner) === attempt) { + this.retained.delete(owner) + } + completion.resolve() + }, failed) + } catch (error) { + failed(error) + } + } +} diff --git a/src/main/ipc/runtime-watcher-pool-shutdown.test.ts b/src/main/ipc/runtime-watcher-pool-shutdown.test.ts new file mode 100644 index 00000000000..1f37f3540c5 --- /dev/null +++ b/src/main/ipc/runtime-watcher-pool-shutdown.test.ts @@ -0,0 +1,79 @@ +import { expect, it, vi } from 'vitest' +import { RuntimeWatcherProcessPool } from './runtime-watcher-process-pool' +import type { WatcherProcessHooks } from './parcel-watcher-process-subscription' +import { WatcherProcessFailure } from './parcel-watcher-process-failure' + +function supervisor() { + const pending = Promise.withResolvers() + let hooks: WatcherProcessHooks | undefined + return { + pending, + dispose: vi.fn(), + disposeAndWait: vi.fn(() => pending.promise), + subscribe: vi.fn( + async ( + _dir: string, + _callback: unknown, + _options: unknown, + options?: WatcherProcessHooks + ) => { + hooks = options + return { unsubscribe: vi.fn(async () => {}) } + } + ), + fail: () => + hooks?.onTerminalError?.( + new WatcherProcessFailure('failed', 'supervisor', 'process_unavailable') + ) + } +} + +it('awaits retired and replacement supervisors after logical retirement removed the old slot', async () => { + const old = supervisor() + const current = supervisor() + const create = vi.fn().mockReturnValueOnce(old).mockReturnValueOnce(current) + const pool = new RuntimeWatcherProcessPool({ createSupervisor: create }) + await pool.subscribe('/first', vi.fn(), {}) + old.fail() + await new Promise((resolve) => setImmediate(resolve)) + expect(old.disposeAndWait).toHaveBeenCalledOnce() + await pool.subscribe('/second', vi.fn(), {}) + const finished = vi.fn() + const shutdown = pool.disposeAndWait() + expect(pool.disposeAndWait()).toBe(shutdown) + const result = shutdown.then(finished) + current.pending.resolve() + await new Promise((resolve) => setImmediate(resolve)) + expect(finished).not.toHaveBeenCalled() + await expect(pool.subscribe('/third', vi.fn(), {})).rejects.toThrow('disposed') + expect(create).toHaveBeenCalledTimes(2) + old.pending.resolve() + await result +}) + +it('retains failure after synchronous pool disposal and retries it explicitly', async () => { + const child = supervisor() + const failure = new Error('child still running') + child.disposeAndWait.mockRejectedValueOnce(failure) + const pool = new RuntimeWatcherProcessPool({ createSupervisor: () => child }) + await pool.subscribe('/first', vi.fn(), {}) + const shutdown = pool.disposeAndWait() + await expect(shutdown).rejects.toMatchObject({ errors: [failure] }) + const retry = pool.disposeAndWait() + expect(child.disposeAndWait).toHaveBeenCalledTimes(2) + child.pending.resolve() + await retry +}) + +it('does not duplicate termination when a queued retirement races with synchronous disposal', async () => { + const child = supervisor() + const pool = new RuntimeWatcherProcessPool({ createSupervisor: () => child }) + await pool.subscribe('/first', vi.fn(), {}) + child.fail() + pool.dispose() + const shutdown = pool.disposeAndWait() + await new Promise((resolve) => setImmediate(resolve)) + expect(child.disposeAndWait).toHaveBeenCalledOnce() + child.pending.resolve() + await shutdown +}) diff --git a/src/main/ipc/runtime-watcher-pool-state.ts b/src/main/ipc/runtime-watcher-pool-state.ts index e6c016b9844..a99c9b93895 100644 --- a/src/main/ipc/runtime-watcher-pool-state.ts +++ b/src/main/ipc/runtime-watcher-pool-state.ts @@ -1,6 +1,7 @@ import type { WatcherProcessSupervisor } from './parcel-watcher-process-supervisor' -export type RuntimeWatcherPoolSupervisor = Pick +export type RuntimeWatcherPoolSupervisor = Pick & + Partial> export type RuntimeWatcherPoolSlot = { supervisor: RuntimeWatcherPoolSupervisor @@ -10,6 +11,13 @@ export type RuntimeWatcherPoolSlot = { disposed: boolean } +export function activeWatcherSlots( + slots: ReadonlySet, + isolated: boolean +): RuntimeWatcherPoolSlot[] { + return [...slots].filter((slot) => slot.isolated === isolated && !slot.retired) +} + export type RuntimeWatcherPoolAssignment = { slot: RuntimeWatcherPoolSlot leases: number diff --git a/src/main/ipc/runtime-watcher-process-pool.ts b/src/main/ipc/runtime-watcher-process-pool.ts index 4bafc38208d..852fd5abdf4 100644 --- a/src/main/ipc/runtime-watcher-process-pool.ts +++ b/src/main/ipc/runtime-watcher-process-pool.ts @@ -7,15 +7,17 @@ import type { WatcherProcessSubscription } from './parcel-watcher-process-subscription' import { RuntimeWatcherPendingAssignment } from './runtime-watcher-pending-assignment' +import { RuntimeWatcherDisposalOwners } from './runtime-watcher-disposal-owners' import { RuntimeWatcherPoolLifecycle } from './runtime-watcher-pool-lifecycle' import { RuntimeWatcherPredecessorBarriers } from './runtime-watcher-predecessor-barriers' import { RuntimeWatcherQuarantineQueue } from './runtime-watcher-quarantine-queue' import { handleRuntimeWatcherSubscriptionFailure } from './runtime-watcher-subscription-failure' -import type { - RuntimeWatcherPoolAssignment, - RuntimeWatcherPoolSlot, - RuntimeWatcherPoolSupervisor, - RuntimeWatcherProcessPoolOptions +import { + activeWatcherSlots, + type RuntimeWatcherPoolAssignment, + type RuntimeWatcherPoolSlot, + type RuntimeWatcherPoolSupervisor, + type RuntimeWatcherProcessPoolOptions } from './runtime-watcher-pool-state' export type { RuntimeWatcherProcessPoolOptions } from './runtime-watcher-pool-state' @@ -39,6 +41,7 @@ export class RuntimeWatcherProcessPool { RuntimeWatcherPendingAssignment >() private readonly lifecycle = new RuntimeWatcherPoolLifecycle() + private disposalOwners = new RuntimeWatcherDisposalOwners() private readonly predecessorBarriers = new RuntimeWatcherPredecessorBarriers() private readonly quarantineQueue: RuntimeWatcherQuarantineQueue @@ -142,9 +145,12 @@ export class RuntimeWatcherProcessPool { resetForTest(): void { this.dispose() + this.disposalOwners = new RuntimeWatcherDisposalOwners() this.lifecycle.reset() } + disposeAndWait = (): Promise => this.disposalOwners.disposeAndWait(() => this.dispose()) + forgetRoot(dir: string): void { // Physical subscriptions release their assignment through unsubscribe or // terminal callbacks. This clears fault history after setup gives up. @@ -221,9 +227,7 @@ export class RuntimeWatcherProcessPool { } private sharedSlot(): RuntimeWatcherPoolSlot { - const sharedSlots = Array.from(this.activeSlots).filter( - (slot) => !slot.isolated && !slot.retired - ) + const sharedSlots = activeWatcherSlots(this.activeSlots, false) if (sharedSlots.length < this.maxSharedSupervisors) { return this.createSlot(false) } @@ -233,9 +237,7 @@ export class RuntimeWatcherProcessPool { } private quarantineSlot(dir: string): RuntimeWatcherPoolSlot | Promise { - const quarantineSlots = Array.from(this.activeSlots).filter( - (slot) => slot.isolated && !slot.retired - ) + const quarantineSlots = activeWatcherSlots(this.activeSlots, true) if (quarantineSlots.length < this.maxQuarantineSupervisors) { return this.createSlot(true) } @@ -273,7 +275,11 @@ export class RuntimeWatcherProcessPool { slot.roots.clear() // Why: failAllSubscriptions is still iterating callbacks; defer disposal // so every logical root receives the supervisor failure first. + const owners = this.disposalOwners queueMicrotask(() => { + if (this.disposalOwners !== owners) { + return + } this.disposeSlot(slot) this.drainQuarantineWaiters() }) @@ -308,8 +314,7 @@ export class RuntimeWatcherProcessPool { private drainQuarantineWaiters(): void { while ( this.quarantineQueue.length > 0 && - Array.from(this.activeSlots).filter((slot) => slot.isolated && !slot.retired).length < - this.maxQuarantineSupervisors + activeWatcherSlots(this.activeSlots, true).length < this.maxQuarantineSupervisors ) { this.quarantineQueue.grantNext(this.createSlot(true)) } @@ -320,7 +325,7 @@ export class RuntimeWatcherProcessPool { return } slot.disposed = true - slot.supervisor.dispose() + this.disposalOwners.retire(slot.supervisor) this.allSlots.delete(slot) } } diff --git a/src/relay/agent-exec-disposal.test.ts b/src/relay/agent-exec-disposal.test.ts new file mode 100644 index 00000000000..6e770ff5b9e --- /dev/null +++ b/src/relay/agent-exec-disposal.test.ts @@ -0,0 +1,148 @@ +import { execFile } from 'node:child_process' +import type * as ChildProcess from 'node:child_process' +import { afterEach, beforeEach, expect, it, vi } from 'vitest' +import { AgentExecHandler } from './agent-exec-handler' +import { RELAY_AGENT_CLOSE_DEADLINE_MS } from './relay-agent-process-lifetime' +import { createFakeChild, requestContext } from './agent-exec-handler-test-harness' +import type { MethodHandler, RelayDispatcher } from './dispatcher' + +// Why an untyped mock: the fake child stubs only what AgentExecHandler reads from a ChildProcess. +const { spawnMock } = vi.hoisted(() => ({ spawnMock: vi.fn() })) +vi.mock('child_process', async (importOriginal) => ({ + ...(await importOriginal()), + spawn: (...args: unknown[]) => spawnMock(...args), + execFile: vi.fn() +})) + +function fixture() { + const methods = new Map() + const handler = new AgentExecHandler( + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the handler only registers request methods. + { + onRequest: (name: string, method: MethodHandler) => methods.set(name, method) + } as unknown as RelayDispatcher + ) + const exec = (params: Record = {}) => + methods.get('agent.execNonInteractive')!( + { binary: 'agent', timeoutMs: 1000, ...params }, + requestContext() + ) + return { handler, exec } +} + +beforeEach(() => { + vi.useFakeTimers() + spawnMock.mockReset() + vi.mocked(execFile).mockReset() +}) +afterEach(() => { + vi.useRealTimers() + vi.restoreAllMocks() +}) + +it('waits for every child close, including commands without a cwd and replaced lanes', async () => { + const f = fixture() + const children = [createFakeChild(), createFakeChild(), createFakeChild()] + for (const child of children) { + spawnMock.mockReturnValueOnce(child) + } + const requests = [f.exec(), f.exec({ cwd: '/repo' }), f.exec({ cwd: '/repo' })] + let finished = false + const disposal = f.handler.dispose().then(() => { + finished = true + }) + await Promise.resolve() + expect(finished).toBe(false) + if (process.platform === 'win32') { + expect(execFile).toHaveBeenCalledWith( + 'taskkill', + ['/pid', '12345', '/T', '/F'], + expect.any(Function) + ) + } else { + for (const child of children) { + expect(child.kill).toHaveBeenCalledWith('SIGKILL') + } + } + children[0].emit('close', null) + children[2].emit('close', null) + await Promise.resolve() + expect(finished).toBe(false) + children[1].emit('close', null) + await disposal + expect(finished).toBe(true) + await Promise.all(requests) +}) + +it('refuses execution once disposal begins and after it completes', async () => { + const f = fixture() + const child = createFakeChild() + spawnMock.mockReturnValue(child) + const request = f.exec() + const disposal = f.handler.dispose() + await expect(f.exec()).rejects.toThrow() + expect(spawnMock).toHaveBeenCalledTimes(1) + child.emit('close', null) + await Promise.all([request, disposal]) + await expect(f.exec()).rejects.toThrow() + expect(spawnMock).toHaveBeenCalledTimes(1) +}) + +it('retains timed-out children until close rather than treating RPC settlement as exit', async () => { + const f = fixture() + const child = createFakeChild() + spawnMock.mockReturnValue(child) + const request = f.exec({ cwd: '/repo' }) + await vi.advanceTimersByTimeAsync(1000) + await expect(request).resolves.toMatchObject({ timedOut: true }) + let finished = false + const disposal = f.handler.dispose().then(() => { + finished = true + }) + await Promise.resolve() + expect(finished).toBe(false) + child.emit('error', new Error('late child error')) + await Promise.resolve() + expect(finished).toBe(false) + child.emit('close', null) + await disposal + expect(finished).toBe(true) + expect(child.listenerCount('error')).toBe(0) + expect(child.listenerCount('close')).toBe(0) +}) + +it('resolves disposal with no children or only already-closed children', async () => { + await fixture().handler.dispose() + const f = fixture() + const child = createFakeChild() + spawnMock.mockReturnValue(child) + const request = f.exec() + child.emit('close', 0) + await request + await f.handler.dispose() + expect(child.kill).not.toHaveBeenCalled() +}) + +it('rejects within the close deadline and re-checks without re-killing on retry', async () => { + const f = fixture() + const child = createFakeChild() + spawnMock.mockReturnValue(child) + const request = f.exec({ timeoutMs: 60_000 }) + const signals = () => + process.platform === 'win32' + ? vi.mocked(execFile).mock.calls.length + : child.kill.mock.calls.length + const first = f.handler.dispose() + const firstOutcome = first.then( + () => 'resolved', + (error: Error) => error.message + ) + expect(signals()).toBe(1) + await vi.advanceTimersByTimeAsync(RELAY_AGENT_CLOSE_DEADLINE_MS) + await expect(firstOutcome).resolves.toBe('relay_agent_execution_shutdown_incomplete') + const retry = f.handler.dispose() + expect(signals()).toBe(1) + child.emit('close', null) + await expect(retry).resolves.toBeUndefined() + await request +}) diff --git a/src/relay/agent-exec-handler.test.ts b/src/relay/agent-exec-handler.test.ts index 87a2ef0403b..1cb2820b1ec 100644 --- a/src/relay/agent-exec-handler.test.ts +++ b/src/relay/agent-exec-handler.test.ts @@ -523,6 +523,10 @@ describe('AgentExecHandler', () => { } expect(child.stdout.listenerCount('data')).toBe(0) expect(child.stderr.listenerCount('data')).toBe(0) + // Shutdown keeps tracking the child until it physically closes. + expect(child.listenerCount('error')).toBe(1) + expect(child.listenerCount('close')).toBe(1) + child.emit('close', null) expect(child.listenerCount('error')).toBe(0) expect(child.listenerCount('close')).toBe(0) } finally { diff --git a/src/relay/agent-exec-handler.ts b/src/relay/agent-exec-handler.ts index 152e8bd72fb..4314dbbc6a7 100644 --- a/src/relay/agent-exec-handler.ts +++ b/src/relay/agent-exec-handler.ts @@ -6,6 +6,7 @@ import type { RelayDispatcher, RequestContext } from './dispatcher' import { applyTerminalGitCredentialPromptGuard } from '../shared/terminal-git-credential-guard' import { mergeGitConfigEnvProtocol } from '../shared/git-credential-prompt-env' import { terminateRelaySubprocessTree } from './subprocess-tree-termination' +import { RelayAgentProcessLifetime } from './relay-agent-process-lifetime' const DEFAULT_TIMEOUT_MS = 60_000 const MAX_TIMEOUT_MS = 5 * 60 * 1000 @@ -109,6 +110,7 @@ type ExecResult = { * and a clean exit code instead of an interactive session. */ export class AgentExecHandler { + private readonly processLifetime = new RelayAgentProcessLifetime() // Why: commit-message and PR-field generation can run together for one cwd; // operation lanes let cancel target only the user-visible job that stopped. private inFlightByLane = new Map() @@ -124,6 +126,10 @@ export class AgentExecHandler { dispatcher.onRequest('agent.cancelExec', (p) => this.cancel(p as CancelParams)) } + dispose(): Promise { + return this.processLifetime.dispose() + } + private async cancel(params: CancelParams): Promise<{ canceled: boolean }> { const cwd = typeof params.cwd === 'string' ? params.cwd : '' const entry = this.inFlightByLane.get(this.laneKey(cwd, params.operation)) @@ -135,6 +141,7 @@ export class AgentExecHandler { } private async exec(params: ExecParams, context?: RequestContext): Promise { + this.processLifetime.assertAdmission() const binary = typeof params.binary === 'string' ? params.binary : '' if (!binary) { throw new Error('agent.execNonInteractive: binary is required') @@ -188,6 +195,7 @@ export class AgentExecHandler { return } + this.processLifetime.track(child) let stdout = '' let stderr = '' let stdoutBytes = 0 diff --git a/src/relay/fs-file-stream-shutdown.test.ts b/src/relay/fs-file-stream-shutdown.test.ts new file mode 100644 index 00000000000..c26fb73b2c4 --- /dev/null +++ b/src/relay/fs-file-stream-shutdown.test.ts @@ -0,0 +1,114 @@ +import { afterEach, beforeEach, expect, it, vi } from 'vitest' +import { mkdtemp, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import type { RelayDispatcher, RequestContext } from './dispatcher' +import { readRelayFileStreamMetadata } from './fs-handler-file-read' +import { RelayStreamRegistry } from './fs-stream-registry' +import { STREAM_CHUNK_SIZE } from './protocol' + +let directory: string +let filePath: string +beforeEach(async () => { + directory = await mkdtemp(join(tmpdir(), 'orca-file-stream-shutdown-')) + filePath = join(directory, 'sample.png') + await writeFile(filePath, Buffer.alloc(12, 42)) +}) +afterEach(async () => { + vi.restoreAllMocks() + await rm(directory, { recursive: true, force: true }) +}) + +function fixture() { + const registry = new RelayStreamRegistry() + const bulk = vi.fn(async () => {}) + const terminal = vi.fn((_id, _method, _params, settled) => { + settled({ ok: true }) + return true + }) + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the file stream reads only these two dispatcher methods. + const dispatcher = { notifyBulk: bulk, tryNotifyClient: terminal } as unknown as RelayDispatcher + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the stream reads only clientId and isStale from its context. + const context = { clientId: 1, isStale: () => false } as RequestContext + const start = (paceWithAcks = false) => + readRelayFileStreamMetadata(filePath, dispatcher, registry, context, { + clientId: 1, + paceWithAcks + }) + return { registry, bulk, terminal, start } +} + +it('waits for a scheduled pump and permanently fences later metadata requests', async () => { + const f = fixture() + const scheduled: (() => void)[] = [] + // A real, already-cleared handle satisfies the return type without scheduling anything. + const handle = setImmediate(() => {}) + clearImmediate(handle) + const schedule = vi.spyOn(globalThis, 'setImmediate').mockImplementation((callback) => { + scheduled.push(callback) + return handle + }) + await f.start() + schedule.mockRestore() + const finished = vi.fn() + const drain = f.registry.disposeAll().then(finished) + await new Promise((resolve) => setImmediate(resolve)) + expect(finished).not.toHaveBeenCalled() + await expect(f.start()).rejects.toThrow('relay_file_stream_shutdown_fenced') + scheduled.forEach((run) => run()) + await drain + expect(f.bulk).not.toHaveBeenCalled() + expect(f.terminal).not.toHaveBeenCalled() +}) + +it.each(['resolve', 'reject'] as const)( + 'waits for an in-flight final chunk to %s without a terminal frame', + async (outcome) => { + const f = fixture() + const pending = Promise.withResolvers() + f.bulk.mockReturnValue(pending.promise) + await f.start() + await vi.waitFor(() => expect(f.bulk).toHaveBeenCalledOnce()) + const finished = vi.fn() + const drain = f.registry.disposeAll().then(finished) + await new Promise((resolve) => setImmediate(resolve)) + expect(finished).not.toHaveBeenCalled() + if (outcome === 'resolve') { + pending.resolve() + } else { + pending.reject(new Error('transport closed')) + } + await drain + expect(f.terminal).not.toHaveBeenCalled() + } +) + +it('wakes an ACK-parked producer during shutdown', async () => { + const f = fixture() + await writeFile(filePath, Buffer.alloc(STREAM_CHUNK_SIZE * 6, 42)) + await f.start(true) + await vi.waitFor(() => expect(f.bulk).toHaveBeenCalledTimes(4)) + await f.registry.disposeAll() + expect(f.bulk).toHaveBeenCalledTimes(4) + expect(f.terminal).not.toHaveBeenCalled() +}) + +it('retains a binary-probe handle when its close fails and retries shutdown cleanup', async () => { + const f = fixture() + filePath = join(directory, 'sample.txt') + await writeFile(filePath, 'text content') + const failure = new Error('probe close failed') + const release = f.registry.releaseUnregisteredHandle.bind(f.registry) + vi.spyOn(f.registry, 'releaseUnregisteredHandle').mockImplementationOnce((handle) => { + vi.spyOn(handle, 'close').mockRejectedValueOnce(failure) + return release(handle) + }) + try { + await expect(f.start()).rejects.toBe(failure) + expect(f.registry.size()).toBe(1) + await f.registry.disposeAll() + expect(f.registry.size()).toBe(0) + } finally { + await f.registry.disposeAll() + } +}) diff --git a/src/relay/fs-handler.ts b/src/relay/fs-handler.ts index 39b22984766..87ef3444d83 100644 --- a/src/relay/fs-handler.ts +++ b/src/relay/fs-handler.ts @@ -275,4 +275,8 @@ export class FsHandler { disposeFileStreams(): Promise { return this.streamRegistry.disposeAll() } + + disposeWatchers(): Promise { + return this.watchRegistry.disposeAndWait() + } } diff --git a/src/relay/relay-agent-process-lifetime.ts b/src/relay/relay-agent-process-lifetime.ts new file mode 100644 index 00000000000..9e290d3e952 --- /dev/null +++ b/src/relay/relay-agent-process-lifetime.ts @@ -0,0 +1,76 @@ +import { terminateRelaySubprocessTree } from './subprocess-tree-termination' + +type Child = Parameters[0] + +// Why 10s: the same physical-exit deadline the relay's watcher children get after a kill. +export const RELAY_AGENT_CLOSE_DEADLINE_MS = 10_000 + +/** Request timeout/cancellation is not evidence that its host child has closed. */ +export class RelayAgentProcessLifetime { + private fenced = false + private readonly children = new Map>() + // Why: a retry re-checks close instead of re-killing a tree whose pid may already be reused. + private readonly signalled = new WeakSet() + private disposal: Promise | null = null + + constructor(private readonly closeDeadlineMs = RELAY_AGENT_CLOSE_DEADLINE_MS) {} + + assertAdmission(): void { + if (this.fenced) { + throw new Error('relay_agent_execution_shutdown_fenced') + } + } + + track(child: Child): void { + this.assertAdmission() + // Why: no Promise.withResolvers — the relay bundle still targets Node 18 hosts. + let resolveClosed!: () => void + const closed = new Promise((resolve) => { + resolveClosed = resolve + }) + this.children.set(child, closed) + // A late error after request timeout must not remove physical-close tracking or crash the host. + const onError = () => {} + child.on('error', onError) + child.once('close', () => { + child.off('error', onError) + this.children.delete(child) + resolveClosed() + }) + } + + dispose(): Promise { + this.fenced = true + if (this.disposal) { + return this.disposal + } + const pending = [...this.children.entries()] + for (const [child] of pending) { + if (!this.signalled.has(child)) { + this.signalled.add(child) + terminateRelaySubprocessTree(child) + } + } + const disposal = this.waitForClose(pending.map(([, closed]) => closed)).finally(() => { + this.disposal = null + }) + this.disposal = disposal + return disposal + } + + private waitForClose(closes: Promise[]): Promise { + if (closes.length === 0) { + return Promise.resolve() + } + return new Promise((resolve, reject) => { + const timer = setTimeout(() => { + reject(new Error('relay_agent_execution_shutdown_incomplete')) + }, this.closeDeadlineMs) + timer.unref?.() + void Promise.all(closes).then(() => { + clearTimeout(timer) + resolve() + }) + }) + } +} diff --git a/src/relay/relay-filesystem-watch-registry.ts b/src/relay/relay-filesystem-watch-registry.ts index 14f5d5b2ee8..e0e538b0736 100644 --- a/src/relay/relay-filesystem-watch-registry.ts +++ b/src/relay/relay-filesystem-watch-registry.ts @@ -1,6 +1,5 @@ import type { RelayDispatcher, RequestContext } from './dispatcher' import { MAX_BATCHED_WATCHER_EVENTS } from '../main/ipc/filesystem-watcher-event-batch' -import { isWatcherProcessFailure } from '../main/ipc/parcel-watcher-process-failure' import { WATCHER_IGNORE_DIRS, buildParcelWatcherIgnoreOptions @@ -10,7 +9,7 @@ import { type RelayWatcherProcessPool } from './relay-watcher-process-pool' import { emitRelayWatcherEvents, emitRelayWatcherOverflow } from './relay-watcher-event-emitter' -import { awaitRelayWatcherSetup, shouldRetryInitialRelayWatch } from './relay-watcher-setup-wait' +import { awaitRelayWatcherSetupForClient, startInitialRelayWatch } from './relay-watcher-setup-wait' import { RelayWatcherTeardownTracker, type RelayWatcherTeardownState @@ -27,6 +26,7 @@ import { releaseStaleRelayWatches } from './relay-watcher-stale-client-release' import { PromiseSettlementWaiters } from '../shared/promise-settlement-waiters' import { joinRelayWatcherPendingSetup } from './relay-watcher-pending-setup-join' import { createRelayWatcherState } from './relay-watcher-state' +import { closeRelayWatchesAndWait, disposeRelayWatchesAndWait } from './relay-watcher-shutdown' const RELAY_WATCH_OPTIONS = buildParcelWatcherIgnoreOptions(WATCHER_IGNORE_DIRS) @@ -40,6 +40,7 @@ export class RelayFilesystemWatchRegistry { ) private readonly teardownTracker: RelayWatcherTeardownTracker private readonly removalFence: RelayWatcherRemovalFence + private disposed = false constructor( private readonly dispatcher: RelayDispatcher, @@ -58,6 +59,9 @@ export class RelayFilesystemWatchRegistry { } watch(rootPath: string, context?: RequestContext, watchId?: number): Promise { + if (this.disposed) { + return Promise.reject(new Error('relay_watcher_shutdown_fenced')) + } const rootKey = normalizeRuntimePathForComparison(rootPath) if (this.removalFence.isActive(rootKey)) { return Promise.reject(new Error('Remote worktree deletion already in progress')) @@ -71,7 +75,9 @@ export class RelayFilesystemWatchRegistry { if (watchId !== undefined) { existing.clientWatchIds.set(clientId, watchId) } - return this.awaitSetupForClient(existing, clientId, context) + return awaitRelayWatcherSetupForClient(existing, context, () => + this.releaseWatchClient(existing, clientId) + ) } void this.closeWatch(existing).catch(() => {}) } @@ -104,6 +110,9 @@ export class RelayFilesystemWatchRegistry { if (capacityRelease) { await capacityRelease } + if (this.disposed) { + throw new Error('relay_watcher_shutdown_fenced') + } const clientId = context?.clientId ?? 0 const isStale = context?.isStale ?? (() => false) const existing = this.watches.get(rootKey) @@ -112,7 +121,9 @@ export class RelayFilesystemWatchRegistry { if (watchId !== undefined) { existing.clientWatchIds.set(clientId, watchId) } - await this.awaitSetupForClient(existing, clientId, context) + await awaitRelayWatcherSetupForClient(existing, context, () => + this.releaseWatchClient(existing, clientId) + ) return } @@ -121,7 +132,9 @@ export class RelayFilesystemWatchRegistry { const state = createRelayWatcherState(rootKey, rootPath, clientId, isStale, watchId) this.watches.set(rootKey, state) state.setupWaiters = new PromiseSettlementWaiters(this.startInitialWatch(state)) - await this.awaitSetupForClient(state, clientId, context) + await awaitRelayWatcherSetupForClient(state, context, () => + this.releaseWatchClient(state, clientId) + ) } unwatch(rootPath: string, context?: RequestContext): void { @@ -183,27 +196,31 @@ export class RelayFilesystemWatchRegistry { } dispose(): void { + this.disposed = true this.watches.forEach((state) => void this.closeWatch(state).catch(() => {})) this.watcherPool.dispose() } - private async startInitialWatch(state: RelayWatcherTeardownState): Promise { - try { - await this.subscribeState(state) - } catch (firstError) { - if (!state.closed && shouldRetryInitialRelayWatch(firstError)) { - try { - await this.subscribeState(state) - emitRelayWatcherOverflow(this.dispatcher, state.rootPath, state.closed) - return - } catch (quarantineError) { - void this.closeWatch(state).catch(() => {}) - throw quarantineError - } - } - void this.closeWatch(state).catch(() => {}) - throw firstError - } + closeWatchesAndWait(): Promise { + this.disposed = true + return closeRelayWatchesAndWait( + this.watches, + this.pendingSetups, + this.teardownTracker, + (state) => this.closeWatch(state) + ) + } + + disposeAndWait = (): Promise => + disposeRelayWatchesAndWait(() => this.closeWatchesAndWait(), this.watcherPool) + + private startInitialWatch(state: RelayWatcherTeardownState): Promise { + return startInitialRelayWatch( + state, + () => this.subscribeState(state), + () => emitRelayWatcherOverflow(this.dispatcher, state.rootPath, state.closed), + () => this.closeWatch(state) + ) } private subscribeState(state: RelayWatcherTeardownState): Promise { @@ -245,7 +262,13 @@ export class RelayFilesystemWatchRegistry { state.generation !== generation || this.watches.get(state.rootKey) !== state ) { + if (state.closed) { + state.subscription = subscription + } await subscription.unsubscribe() + if (state.subscription === subscription) { + state.subscription = null + } return } state.subscription = subscription @@ -276,30 +299,6 @@ export class RelayFilesystemWatchRegistry { }) } - private async awaitSetupForClient( - state: RelayWatcherTeardownState, - clientId: number, - context?: RequestContext - ): Promise { - try { - await awaitRelayWatcherSetup(state.setupWaiters, context?.signal) - } catch (error) { - this.releaseWatchClient(state, clientId) - const expectedAbort = - (error instanceof Error && error.name === 'AbortError') || - (isWatcherProcessFailure(error) && error.code === 'subscribe_aborted') - if (expectedAbort) { - return - } - const message = error instanceof Error ? error.message : String(error) - process.stderr.write(`[relay] File watcher not available for ${state.rootPath}: ${message}\n`) - throw error - } - if (context?.isStale()) { - this.releaseWatchClient(state, clientId) - } - } - private releaseClientWatches(clientId: number): void { this.watches.forEach((state) => this.releaseWatchClient(state, clientId)) } diff --git a/src/relay/relay-runtime-owned-process-shutdown.test.ts b/src/relay/relay-runtime-owned-process-shutdown.test.ts index 48644137c73..15bf7c7742f 100644 --- a/src/relay/relay-runtime-owned-process-shutdown.test.ts +++ b/src/relay/relay-runtime-owned-process-shutdown.test.ts @@ -4,21 +4,24 @@ import { RelayRuntimeServices } from './relay-runtime-services' afterEach(() => vi.restoreAllMocks()) function fixture() { + const agents = vi.fn(async () => {}) const skill = vi.fn(async () => {}) const vault = vi.fn(async () => {}) const responses = vi.fn(async () => {}) const fileStreams = vi.fn(async () => {}) + const watchers = vi.fn(async () => {}) // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: disposeOwnedProcesses only reads the owners stubbed here. const runtime = Object.assign(Object.create(RelayRuntimeServices.prototype), { + agentExecHandler: { dispose: agents }, skillInstallHandler: { dispose: skill }, aiVaultService: { dispose: vault }, responseStreams: { disposeAllAndWait: responses }, - fsHandler: { disposeFileStreams: fileStreams } + fsHandler: { disposeFileStreams: fileStreams, disposeWatchers: watchers } }) as RelayRuntimeServices - return { runtime, skill, vault, responses, fileStreams } + return { runtime, agents, skill, vault, responses, fileStreams, watchers } } -it.each(['responses', 'fileStreams'] as const)( +it.each(['agents', 'responses', 'fileStreams', 'watchers'] as const)( 'does not acknowledge cleanup before %s have settled', async (owner) => { const f = fixture() @@ -70,3 +73,19 @@ it('handles hosts without a vault service', async () => { await expect(f.runtime.disposeOwnedProcesses()).resolves.toBeUndefined() expect(f.vault).not.toHaveBeenCalled() }) + +it('rejects a failed agent or watcher shutdown after cleaning every other owner', async () => { + const f = fixture() + const agentError = new Error('agent cleanup failed') + const watcherError = new Error('watcher cleanup failed') + f.agents.mockRejectedValueOnce(agentError) + f.watchers.mockRejectedValueOnce(watcherError) + await expect(f.runtime.disposeOwnedProcesses()).rejects.toMatchObject({ + message: 'relay_owned_process_shutdown_incomplete', + errors: [agentError, watcherError] + }) + expect(f.skill).toHaveBeenCalledOnce() + expect(f.vault).toHaveBeenCalledOnce() + expect(f.fileStreams).toHaveBeenCalledOnce() + await expect(f.runtime.disposeOwnedProcesses()).resolves.toBeUndefined() +}) diff --git a/src/relay/relay-runtime-services.ts b/src/relay/relay-runtime-services.ts index d4d179449ba..e52fa6afddf 100644 --- a/src/relay/relay-runtime-services.ts +++ b/src/relay/relay-runtime-services.ts @@ -32,6 +32,7 @@ export class RelayRuntimeServices { readonly fsHandler: FsHandler readonly gitHandler: GitHandler readonly skillInstallHandler: SkillInstallHandler + readonly agentExecHandler: AgentExecHandler private readonly responseStreams: GitResponseStreamRegistry private readonly aiVaultService: ReturnType | null private readonly sessionSearch: { dispose(): void } | null @@ -79,7 +80,7 @@ export class RelayRuntimeServices { this.skillInstallHandler = new SkillInstallHandler(dispatcher) const externalAutomationsHandler = new ExternalAutomationsHandler(dispatcher) const portScanHandler = new PortScanHandler(dispatcher) - const agentExecHandler = new AgentExecHandler(dispatcher) + this.agentExecHandler = new AgentExecHandler(dispatcher) const workspaceSessionHandler = new WorkspaceSessionHandler(dispatcher) const relayPlatform = parseUnameToRelayPlatform(process.platform, process.arch) const hostPlatform = relayPlatform ? getRemoteHostPlatform(relayPlatform) : undefined @@ -105,7 +106,7 @@ export class RelayRuntimeServices { this.skillInstallHandler, externalAutomationsHandler, portScanHandler, - agentExecHandler, + this.agentExecHandler, workspaceSessionHandler, new AiVaultHandler(dispatcher, { hostPlatform, @@ -125,12 +126,18 @@ export class RelayRuntimeServices { // pumps and file descriptors behind it are only proven gone by these registry drains. async disposeOwnedProcesses(): Promise { const failures: unknown[] = [] + const agents = this.agentExecHandler.dispose().catch((error: unknown) => { + failures.push(error) + }) const responses = this.responseStreams.disposeAllAndWait().catch((error: unknown) => { failures.push(error) }) const fileStreams = this.fsHandler.disposeFileStreams().catch((error: unknown) => { failures.push(error) }) + const watchers = this.fsHandler.disposeWatchers().catch((error: unknown) => { + failures.push(error) + }) await this.skillInstallHandler.dispose().catch((error) => { relayLogLine( `[relay] Skill upload cleanup failed: ${error instanceof Error ? error.message : String(error)}` @@ -141,10 +148,12 @@ export class RelayRuntimeServices { `[relay] AI Vault sidecar shutdown failed: ${error instanceof Error ? error.message : String(error)}` ) }) + await agents await responses await fileStreams - // Why: an unclosed fd defers shutdown so the next attempt retries it; skill/AI Vault - // cleanup stays log-and-continue until the T2 lifecycle port. + await watchers + // Why: an unclosed fd, watcher child or agent child defers shutdown so the next attempt retries + // it; skill/AI Vault cleanup stays log-and-continue until a later T2 slice. if (failures.length > 0) { throw new AggregateError(failures, 'relay_owned_process_shutdown_incomplete') } diff --git a/src/relay/relay-watcher-process-pool.ts b/src/relay/relay-watcher-process-pool.ts index e0093b9a166..0ff1fe1a59e 100644 --- a/src/relay/relay-watcher-process-pool.ts +++ b/src/relay/relay-watcher-process-pool.ts @@ -5,7 +5,8 @@ import { WatcherProcessSupervisor } from '../main/ipc/parcel-watcher-process-sup export type RelayWatcherProcessPool = Pick< RuntimeWatcherProcessPool, 'dispose' | 'forgetRoot' | 'subscribe' -> +> & + Partial> export function getRelayWatcherProcessEntryPath(): string { return join(__dirname, 'relay-watcher.js') diff --git a/src/relay/relay-watcher-setup-wait.ts b/src/relay/relay-watcher-setup-wait.ts index a4c7b348b02..f09dfb7470b 100644 --- a/src/relay/relay-watcher-setup-wait.ts +++ b/src/relay/relay-watcher-setup-wait.ts @@ -1,5 +1,55 @@ import { isWatcherProcessFailure } from '../main/ipc/parcel-watcher-process-failure' import type { PromiseSettlementWaiters } from '../shared/promise-settlement-waiters' +import type { RequestContext } from './dispatcher' +import type { RelayWatcherTeardownState } from './relay-watcher-teardown-tracker' + +export async function startInitialRelayWatch( + state: RelayWatcherTeardownState, + subscribe: () => Promise, + emitOverflow: () => void, + close: () => Promise +): Promise { + try { + await subscribe() + } catch (firstError) { + if (!state.closed && shouldRetryInitialRelayWatch(firstError)) { + try { + await subscribe() + emitOverflow() + return + } catch (quarantineError) { + void close().catch(() => {}) + throw quarantineError + } + } + void close().catch(() => {}) + throw firstError + } +} + +export async function awaitRelayWatcherSetupForClient( + state: RelayWatcherTeardownState, + context: RequestContext | undefined, + releaseClient: () => void +): Promise { + try { + await awaitRelayWatcherSetup(state.setupWaiters, context?.signal) + } catch (error) { + releaseClient() + const expectedAbort = + (error instanceof Error && error.name === 'AbortError') || + (isWatcherProcessFailure(error) && error.code === 'subscribe_aborted') + if (expectedAbort) { + return + } + const message = error instanceof Error ? error.message : String(error) + process.stderr.write(`[relay] File watcher not available for ${state.rootPath}: ${message}\n`) + throw error + } + if (context?.isStale()) { + releaseClient() + } +} export function shouldRetryInitialRelayWatch(error: unknown): boolean { return ( diff --git a/src/relay/relay-watcher-shutdown.test.ts b/src/relay/relay-watcher-shutdown.test.ts new file mode 100644 index 00000000000..9cf371969b2 --- /dev/null +++ b/src/relay/relay-watcher-shutdown.test.ts @@ -0,0 +1,143 @@ +import { expect, it, vi } from 'vitest' +import { join } from 'node:path' +import { tmpdir } from 'node:os' +import type { RelayDispatcher } from './dispatcher' +import { RelayFilesystemWatchRegistry } from './relay-filesystem-watch-registry' + +function fixture() { + const unsubscribe = vi.fn(async () => {}) + const pool = { + subscribe: vi.fn(async () => ({ unsubscribe })), + forgetRoot: vi.fn(), + dispose: vi.fn(), + disposeAndWait: vi.fn(async () => {}) + } + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: shutdown paths never reach the dispatcher. + const registry = new RelayFilesystemWatchRegistry({} as RelayDispatcher, pool) + const root = join(tmpdir(), 'orca-watcher-shutdown-fixture') + return { registry, pool, unsubscribe, root } +} + +it('fences new watches synchronously and joins an in-flight native unsubscribe', async () => { + const f = fixture() + await f.registry.watch(f.root) + const pending = Promise.withResolvers() + f.unsubscribe.mockReturnValue(pending.promise) + const finished = vi.fn() + const shutdown = f.registry.closeWatchesAndWait().then(finished) + const duplicate = f.registry.closeWatchesAndWait() + await expect(f.registry.watch(f.root)).rejects.toThrow('relay_watcher_shutdown_fenced') + await new Promise((resolve) => setImmediate(resolve)) + expect(finished).not.toHaveBeenCalled() + expect(f.unsubscribe).toHaveBeenCalledOnce() + expect(f.pool.dispose).not.toHaveBeenCalled() + pending.resolve() + await Promise.all([shutdown, duplicate]) + expect(f.pool.forgetRoot).toHaveBeenCalledWith(f.root) +}) + +it('retains failed teardown and retries without admitting replacement watches', async () => { + const f = fixture() + await f.registry.watch(f.root) + const failure = new Error('native close failed') + f.unsubscribe.mockRejectedValueOnce(failure) + await expect(f.registry.closeWatchesAndWait()).rejects.toMatchObject({ + message: 'relay_watcher_shutdown_incomplete', + errors: [failure] + }) + expect(f.pool.forgetRoot).not.toHaveBeenCalled() + await expect(f.registry.watch(f.root)).rejects.toThrow('relay_watcher_shutdown_fenced') + await f.registry.closeWatchesAndWait() + expect(f.unsubscribe).toHaveBeenCalledTimes(2) + expect(f.pool.forgetRoot).toHaveBeenCalledOnce() +}) + +it('waits for a late subscription and its unsubscribe after setup has begun', async () => { + const f = fixture() + const setup = Promise.withResolvers<{ unsubscribe: () => Promise }>() + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the registry reads only unsubscribe from a subscription. + f.pool.subscribe.mockReturnValue(setup.promise as ReturnType) + const watch = f.registry.watch(f.root) + const finished = vi.fn() + const shutdown = f.registry.closeWatchesAndWait().then(finished) + const close = Promise.withResolvers() + f.unsubscribe.mockReturnValue(close.promise) + setup.resolve({ unsubscribe: f.unsubscribe }) + await vi.waitFor(() => expect(f.unsubscribe).toHaveBeenCalledOnce()) + expect(finished).not.toHaveBeenCalled() + close.resolve() + await Promise.all([watch, shutdown]) + expect(f.pool.subscribe).toHaveBeenCalledOnce() +}) + +it('joins teardown already started by client unwatch', async () => { + const f = fixture() + await f.registry.watch(f.root) + const close = Promise.withResolvers() + f.unsubscribe.mockReturnValue(close.promise) + f.registry.unwatch(f.root) + const finished = vi.fn() + const shutdown = f.registry.closeWatchesAndWait().then(finished) + await new Promise((resolve) => setImmediate(resolve)) + expect(finished).not.toHaveBeenCalled() + close.resolve() + await shutdown + expect(f.unsubscribe).toHaveBeenCalledOnce() +}) + +it('retains a late subscription whose cleanup fails instead of treating it as failed setup', async () => { + const f = fixture() + const setup = Promise.withResolvers<{ unsubscribe: () => Promise }>() + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the registry reads only unsubscribe from a subscription. + f.pool.subscribe.mockReturnValue(setup.promise as ReturnType) + const failure = new Error('late native close failed') + f.unsubscribe.mockRejectedValueOnce(failure) + const watch = f.registry.watch(f.root).catch((error: unknown) => error) + const shutdown = f.registry.closeWatchesAndWait().catch((error: unknown) => error) + setup.resolve({ unsubscribe: f.unsubscribe }) + expect(await watch).toBe(failure) + expect(await shutdown).toMatchObject({ + message: 'relay_watcher_shutdown_incomplete', + errors: [failure] + }) + expect(f.pool.forgetRoot).not.toHaveBeenCalled() + await f.registry.closeWatchesAndWait() + expect(f.unsubscribe).toHaveBeenCalledTimes(2) + expect(f.pool.forgetRoot).toHaveBeenCalledOnce() +}) + +it('does not publish a replacement watch when its predecessor closes after shutdown fencing', async () => { + const f = fixture() + await f.registry.watch(f.root) + const close = Promise.withResolvers() + f.unsubscribe.mockReturnValue(close.promise) + f.registry.unwatch(f.root) + const replacement = f.registry.watch(f.root).catch((error: unknown) => error) + const shutdown = f.registry.closeWatchesAndWait() + close.resolve() + expect(await replacement).toMatchObject({ message: 'relay_watcher_shutdown_fenced' }) + await shutdown + expect(f.pool.subscribe).toHaveBeenCalledOnce() +}) + +it('full watcher disposal waits for physical pool exit after native unsubscribe completes', async () => { + const f = fixture() + await f.registry.watch(f.root) + const pending = Promise.withResolvers() + f.pool.disposeAndWait.mockReturnValue(pending.promise) + const finished = vi.fn() + const shutdown = f.registry.disposeAndWait().then(finished) + await vi.waitFor(() => expect(f.pool.forgetRoot).toHaveBeenCalledOnce()) + expect(finished).not.toHaveBeenCalled() + pending.resolve() + await shutdown +}) + +it('does not acknowledge a failed pool exit after native subscriptions close', async () => { + const f = fixture() + await f.registry.watch(f.root) + const error = new Error('pool exit unproven') + f.pool.disposeAndWait.mockRejectedValueOnce(error) + await expect(f.registry.disposeAndWait()).rejects.toMatchObject({ errors: [error] }) + await f.registry.disposeAndWait() +}) diff --git a/src/relay/relay-watcher-shutdown.ts b/src/relay/relay-watcher-shutdown.ts new file mode 100644 index 00000000000..bd0963adc52 --- /dev/null +++ b/src/relay/relay-watcher-shutdown.ts @@ -0,0 +1,57 @@ +import type { RelayWatcherPendingSetup } from './relay-watcher-setup-tracking' +import type { RelayWatcherProcessPool } from './relay-watcher-process-pool' +import type { + RelayWatcherTeardownState, + RelayWatcherTeardownTracker +} from './relay-watcher-teardown-tracker' + +export async function disposeRelayWatchesAndWait( + closeWatches: () => Promise, + pool: RelayWatcherProcessPool +): Promise { + const close = closeWatches() + const children = Promise.resolve().then(() => { + if (!pool.disposeAndWait) { + throw new Error('relay_watcher_pool_awaited_disposal_unavailable') + } + return pool.disposeAndWait() + }) + const results = await Promise.allSettled([close, children]) + const failures = results.filter((result) => result.status === 'rejected') + if (failures.length > 0) { + throw new AggregateError( + failures.map((failure) => failure.reason), + 'relay_watcher_shutdown_incomplete' + ) + } +} + +export async function closeRelayWatchesAndWait( + watches: ReadonlyMap, + pendingSetups: ReadonlyMap, + tracker: RelayWatcherTeardownTracker, + close: (state: RelayWatcherTeardownState) => Promise +): Promise { + const states = new Set(watches.values()) + for (const root of tracker.rootPaths()) { + const failed = tracker.failedState(root) + if (failed) { + states.add(failed) + } + } + const closures = Promise.allSettled([...states].map(close)) + // Setup refusal is expected after fencing; retained teardown failures remain authoritative. + await Promise.allSettled([...pendingSetups.values()].map((setup) => setup.promise)) + const results = await closures + // Why the fallback: join answers undefined for a root with nothing left in flight. + const remaining = await Promise.allSettled( + tracker.rootPaths().map((root) => tracker.join(root) ?? Promise.resolve()) + ) + const failures = [...results, ...remaining].filter((result) => result.status === 'rejected') + if (failures.length > 0) { + throw new AggregateError( + [...new Set(failures.map((failure) => failure.reason))], + 'relay_watcher_shutdown_incomplete' + ) + } +} diff --git a/src/relay/relay-watcher-teardown-tracker.ts b/src/relay/relay-watcher-teardown-tracker.ts index 785e4dd16ea..d9feff34b9f 100644 --- a/src/relay/relay-watcher-teardown-tracker.ts +++ b/src/relay/relay-watcher-teardown-tracker.ts @@ -47,7 +47,7 @@ export class RelayWatcherTeardownTracker { }, (error) => { const physicalExit = isWatcherProcessFailure(error) ? error.physicalExit : undefined - if (!subscription && !physicalExit) { + if (!subscription && !state.subscription && !physicalExit) { this.failed.delete(state.rootKey) this.forgetRoot(state.rootPath) return