mirror of
https://github.com/stablyai/orca.git
synced 2026-10-03 00:02:19 +00:00
* feat(relay): await owned watcher and agent children on shutdown (#16741 T2 P1)
Relay shutdown and the runtime watcher pool now wait for the child processes
they own to physically exit instead of only signalling them:
- WatcherProcessSupervisor tracks launched children and exposes
disposeAndWait; RuntimeWatcherProcessPool retires slots through
RuntimeWatcherDisposalOwners, which retries failed disposals.
- The relay watch registry fences new watches once disposed and joins
in-flight native unsubscribes; FsHandler.disposeWatchers awaits it.
- AgentExecHandler tracks spawned agent children until they close, so a
request timeout is no longer taken as proof the child is gone.
- RelayRuntimeServices.disposeOwnedProcesses awaits both, alongside the
existing stream drains; a failure defers shutdown for the grace retry.
Ported by hunk from #16741 (a68b6f3531). Relay-bundled code avoids
Promise.withResolvers for the Node 18 floor.
* fix(relay): bound agent close on shutdown and send one graceful watcher signal
AgentExecHandler.dispose() now waits at most RELAY_AGENT_CLOSE_DEADLINE_MS for
killed agent trees to close and then rejects with
relay_agent_execution_shutdown_incomplete, so the relay grace lifecycle defers
and retries instead of hanging. A retry re-checks close without signalling a
tree it already killed.
The supervisor's dispose now marks the graceful signal it sends, so an awaited
termination of the same watcher child only escalates to SIGKILL and waits.
* fix(relay): give the watcher teardown join a promise for every root
TeardownTracker.join answers undefined for a root with nothing in flight; the
type-aware await-thenable rule rejects a promise aggregator over those values.
* test(relay): drop type assertions from the shutdown test fixtures
* fix(watcher): take the child process type from the shared child-process boundary
---------
Co-authored-by: m4air <m4air@m4airs-Air.localdomain>
This commit is contained in:
@@ -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<ChildProcess, Promise<void>>()
|
||||
const signalledChildren = new WeakSet<ChildProcess>()
|
||||
|
||||
/** 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<boolean> {
|
||||
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<boolean> {
|
||||
})
|
||||
}
|
||||
|
||||
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<void> {
|
||||
return child.exitCode !== null || child.signalCode !== null
|
||||
? Promise.resolve()
|
||||
: (physicalExitPromises.get(child) ??
|
||||
new Promise<void>((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<number, PendingWatcherUnsubscribe>,
|
||||
onFinished: (exited: boolean) => void
|
||||
): Promise<void> {
|
||||
// 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,
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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()
|
||||
})
|
||||
@@ -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<ChildProcessHandle>()
|
||||
private disposal: Promise<void> | 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<void> {
|
||||
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
|
||||
}
|
||||
}
|
||||
@@ -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, { 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',
|
||||
|
||||
@@ -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.
|
||||
}
|
||||
|
||||
@@ -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<number, PendingWatcherUnsubscribe>()
|
||||
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<void> => 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
|
||||
})
|
||||
)
|
||||
|
||||
@@ -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<FakeWatcherChild> {
|
||||
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
|
||||
})
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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<void>()
|
||||
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<void>()
|
||||
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<void>()
|
||||
let nested: Promise<void> | 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()
|
||||
})
|
||||
@@ -0,0 +1,90 @@
|
||||
import type { RuntimeWatcherPoolSupervisor } from './runtime-watcher-pool-state'
|
||||
|
||||
type Attempt = { pending: Promise<void> | null; failure?: unknown }
|
||||
|
||||
type Completion = { promise: Promise<void>; 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<void>((res, rej) => {
|
||||
resolve = res
|
||||
reject = rej
|
||||
})
|
||||
return { promise, resolve, reject }
|
||||
}
|
||||
|
||||
export class RuntimeWatcherDisposalOwners {
|
||||
private readonly retained = new Map<RuntimeWatcherPoolSupervisor, Attempt>()
|
||||
private shutdown: Promise<void> | null = null
|
||||
|
||||
retire(owner: RuntimeWatcherPoolSupervisor): void {
|
||||
if (!this.retained.has(owner)) {
|
||||
this.start(owner)
|
||||
}
|
||||
}
|
||||
|
||||
disposeAndWait(disposePool: () => void): Promise<void> {
|
||||
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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<void>()
|
||||
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
|
||||
})
|
||||
@@ -1,6 +1,7 @@
|
||||
import type { WatcherProcessSupervisor } from './parcel-watcher-process-supervisor'
|
||||
|
||||
export type RuntimeWatcherPoolSupervisor = Pick<WatcherProcessSupervisor, 'dispose' | 'subscribe'>
|
||||
export type RuntimeWatcherPoolSupervisor = Pick<WatcherProcessSupervisor, 'dispose' | 'subscribe'> &
|
||||
Partial<Pick<WatcherProcessSupervisor, 'disposeAndWait'>>
|
||||
|
||||
export type RuntimeWatcherPoolSlot = {
|
||||
supervisor: RuntimeWatcherPoolSupervisor
|
||||
@@ -10,6 +11,13 @@ export type RuntimeWatcherPoolSlot = {
|
||||
disposed: boolean
|
||||
}
|
||||
|
||||
export function activeWatcherSlots(
|
||||
slots: ReadonlySet<RuntimeWatcherPoolSlot>,
|
||||
isolated: boolean
|
||||
): RuntimeWatcherPoolSlot[] {
|
||||
return [...slots].filter((slot) => slot.isolated === isolated && !slot.retired)
|
||||
}
|
||||
|
||||
export type RuntimeWatcherPoolAssignment = {
|
||||
slot: RuntimeWatcherPoolSlot
|
||||
leases: number
|
||||
|
||||
@@ -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<RuntimeWatcherPoolAssignment>
|
||||
>()
|
||||
private readonly lifecycle = new RuntimeWatcherPoolLifecycle()
|
||||
private disposalOwners = new RuntimeWatcherDisposalOwners()
|
||||
private readonly predecessorBarriers = new RuntimeWatcherPredecessorBarriers()
|
||||
private readonly quarantineQueue: RuntimeWatcherQuarantineQueue<RuntimeWatcherPoolSlot>
|
||||
|
||||
@@ -142,9 +145,12 @@ export class RuntimeWatcherProcessPool {
|
||||
|
||||
resetForTest(): void {
|
||||
this.dispose()
|
||||
this.disposalOwners = new RuntimeWatcherDisposalOwners()
|
||||
this.lifecycle.reset()
|
||||
}
|
||||
|
||||
disposeAndWait = (): Promise<void> => 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<RuntimeWatcherPoolSlot> {
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<typeof ChildProcess>()),
|
||||
spawn: (...args: unknown[]) => spawnMock(...args),
|
||||
execFile: vi.fn()
|
||||
}))
|
||||
|
||||
function fixture() {
|
||||
const methods = new Map<string, MethodHandler>()
|
||||
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<string, unknown> = {}) =>
|
||||
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
|
||||
})
|
||||
@@ -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 {
|
||||
|
||||
@@ -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<string, InFlightExec>()
|
||||
@@ -124,6 +126,10 @@ export class AgentExecHandler {
|
||||
dispatcher.onRequest('agent.cancelExec', (p) => this.cancel(p as CancelParams))
|
||||
}
|
||||
|
||||
dispose(): Promise<void> {
|
||||
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<ExecResult> {
|
||||
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
|
||||
|
||||
@@ -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<void>()
|
||||
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()
|
||||
}
|
||||
})
|
||||
@@ -275,4 +275,8 @@ export class FsHandler {
|
||||
disposeFileStreams(): Promise<void> {
|
||||
return this.streamRegistry.disposeAll()
|
||||
}
|
||||
|
||||
disposeWatchers(): Promise<void> {
|
||||
return this.watchRegistry.disposeAndWait()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,76 @@
|
||||
import { terminateRelaySubprocessTree } from './subprocess-tree-termination'
|
||||
|
||||
type Child = Parameters<typeof terminateRelaySubprocessTree>[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<Child, Promise<void>>()
|
||||
// Why: a retry re-checks close instead of re-killing a tree whose pid may already be reused.
|
||||
private readonly signalled = new WeakSet<Child>()
|
||||
private disposal: Promise<void> | 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<void>((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<void> {
|
||||
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<void>[]): Promise<void> {
|
||||
if (closes.length === 0) {
|
||||
return Promise.resolve()
|
||||
}
|
||||
return new Promise<void>((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()
|
||||
})
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
this.disposed = true
|
||||
return closeRelayWatchesAndWait(
|
||||
this.watches,
|
||||
this.pendingSetups,
|
||||
this.teardownTracker,
|
||||
(state) => this.closeWatch(state)
|
||||
)
|
||||
}
|
||||
|
||||
disposeAndWait = (): Promise<void> =>
|
||||
disposeRelayWatchesAndWait(() => this.closeWatchesAndWait(), this.watcherPool)
|
||||
|
||||
private startInitialWatch(state: RelayWatcherTeardownState): Promise<void> {
|
||||
return startInitialRelayWatch(
|
||||
state,
|
||||
() => this.subscribeState(state),
|
||||
() => emitRelayWatcherOverflow(this.dispatcher, state.rootPath, state.closed),
|
||||
() => this.closeWatch(state)
|
||||
)
|
||||
}
|
||||
|
||||
private subscribeState(state: RelayWatcherTeardownState): Promise<void> {
|
||||
@@ -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<void> {
|
||||
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))
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
})
|
||||
|
||||
@@ -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<typeof createRelayAiVaultService> | 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<void> {
|
||||
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')
|
||||
}
|
||||
|
||||
@@ -5,7 +5,8 @@ import { WatcherProcessSupervisor } from '../main/ipc/parcel-watcher-process-sup
|
||||
export type RelayWatcherProcessPool = Pick<
|
||||
RuntimeWatcherProcessPool,
|
||||
'dispose' | 'forgetRoot' | 'subscribe'
|
||||
>
|
||||
> &
|
||||
Partial<Pick<RuntimeWatcherProcessPool, 'disposeAndWait'>>
|
||||
|
||||
export function getRelayWatcherProcessEntryPath(): string {
|
||||
return join(__dirname, 'relay-watcher.js')
|
||||
|
||||
@@ -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<void>,
|
||||
emitOverflow: () => void,
|
||||
close: () => Promise<void>
|
||||
): Promise<void> {
|
||||
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<void> {
|
||||
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 (
|
||||
|
||||
@@ -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<void>()
|
||||
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<void> }>()
|
||||
// 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<typeof f.pool.subscribe>)
|
||||
const watch = f.registry.watch(f.root)
|
||||
const finished = vi.fn()
|
||||
const shutdown = f.registry.closeWatchesAndWait().then(finished)
|
||||
const close = Promise.withResolvers<void>()
|
||||
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<void>()
|
||||
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<void> }>()
|
||||
// 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<typeof f.pool.subscribe>)
|
||||
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<void>()
|
||||
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<void>()
|
||||
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()
|
||||
})
|
||||
@@ -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<void>,
|
||||
pool: RelayWatcherProcessPool
|
||||
): Promise<void> {
|
||||
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<string, RelayWatcherTeardownState>,
|
||||
pendingSetups: ReadonlyMap<string, RelayWatcherPendingSetup>,
|
||||
tracker: RelayWatcherTeardownTracker,
|
||||
close: (state: RelayWatcherTeardownState) => Promise<void>
|
||||
): Promise<void> {
|
||||
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'
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user