Parcourir la source

fix(subprocess): preserve failed teardown state

pku-xht il y a 1 mois
Parent
commit
07408ff0f1

+ 2 - 2
.agents/notes/implemented/bug-fix/2026-08-11-synchronous-subprocess-exit-cleanup.i18n.yaml

@@ -2,5 +2,5 @@
 # side as of the last confirmed-consistent state. Both languages carry equal authority;
 # after editing either side, bring the other along and re-record with:
 #   pnpm run verify-translation-pairing --write .agents/notes/implemented/bug-fix/2026-08-11-synchronous-subprocess-exit-cleanup.md
-2026-08-11-synchronous-subprocess-exit-cleanup.md: 086bb2f3763af513b07e0c20910ef80116695ca4
-2026-08-11-synchronous-subprocess-exit-cleanup.zh.md: 524fa6b0d4e4f5012c44e108fd0cecc5d7012802
+2026-08-11-synchronous-subprocess-exit-cleanup.md: 1a6664c03ce0ae90b94d210645ac1401e078fae0
+2026-08-11-synchronous-subprocess-exit-cleanup.zh.md: 4c08847b24422d58025b6563785579d84c3ba6a7

+ 2 - 2
.agents/notes/implemented/bug-fix/2026-08-11-synchronous-subprocess-exit-cleanup.md

@@ -12,7 +12,7 @@ The public subprocess seam correctly promises awaited quiescence during normal d
 
 ## Decision
 
-`LocalSubprocessRuntime` installs one synchronous Node `exit` listener in its Cordis effect. The same effect removes the listener only after normal disposal settles. Ordinary and terminal handles remain in the service's existing live sets while asynchronous cleanup is pending, so a shorter outer exit bound still sees and force-terminates them. If awaited disposal reports a cleanup failure, the service invokes the same synchronous final operations before clearing the sets and removing the listener.
+`LocalSubprocessRuntime` installs one synchronous Node `exit` listener in its Cordis effect. The same effect removes the listener only after normal disposal succeeds. Ordinary and terminal handles remain in the service's existing live sets while asynchronous cleanup is pending, so a shorter outer exit bound still sees and force-terminates them. Disposal releases each successfully stopped target individually; failed targets and the listener remain owned so a later host exit can invoke the same synchronous final operations.
 
 The listener uses local-only final operations that are absent from the public `SubprocessHandle` and `SubprocessTerminalHandle` interfaces:
 
@@ -46,6 +46,6 @@ Unit evidence pins synchronous native-owner and fallback delivery, native termin
 
 ## Consequences
 
-Each active local subprocess service contributes one process-global exit listener, removed with the service effect. Fatal exit gives up grace, output draining, and an in-process quiescence proof in exchange for issuing the strongest available local termination before the host disappears. Normal disposal keeps those guarantees and costs unchanged.
+Each active local subprocess service contributes one process-global exit listener. Successful disposal removes it with the service effect; failed disposal retains it with the targets that still require final termination. Fatal exit gives up grace, output draining, and an in-process quiescence proof in exchange for issuing the strongest available local termination before the host disappears. Normal disposal keeps those guarantees and costs unchanged.
 
 The listener cannot cover failures that do not execute JavaScript. Supported Linux terminals signal the scope described by the [containment decision](2026-08-20-subprocess-native-containment.md); fallback terminals still cannot discover a descendant that escaped before the provider observed it.

+ 2 - 2
.agents/notes/implemented/bug-fix/2026-08-11-synchronous-subprocess-exit-cleanup.zh.md

@@ -12,7 +12,7 @@ Status: implemented
 
 ## Decision
 
-`LocalSubprocessRuntime`在自身 Cordis effect中安装一个同步 Node `exit` listener。只有正常 dispose结算后,同一 effect才移除该 listener。异步清理仍在等待时,普通和 terminal handle继续保留在服务已有的存活集合中,因此更短的外层退出上限仍能看到并强制终止它们。等待中的 dispose报告清理失败时,服务会在清空集合并移除 listener前调用同一组同步最终操作。
+`LocalSubprocessRuntime`在自身 Cordis effect中安装一个同步 Node `exit` listener。只有正常 dispose成功后,同一 effect才移除该 listener。异步清理仍在等待时,普通和 terminal handle继续保留在服务已有的存活集合中,因此更短的外层退出上限仍能看到并强制终止它们。dispose会逐个释放已经成功停稳的目标;失败的目标与 listener继续由服务拥有,使后续宿主退出仍能调用同一组同步最终操作。
 
 该 listener使用本地实现私有的最终操作;公共 `SubprocessHandle`和 `SubprocessTerminalHandle`接口不包含这些操作:
 
@@ -46,6 +46,6 @@ Status: implemented
 
 ## Consequences
 
-每个有效的本地 subprocess service都会贡献一个进程全局 exit listener,并随服务 effect移除。致命退出放弃宽限、输出排空与进程内停稳证明,以换取宿主消失前发出本地可用的最强终止操作。正常 dispose的保证与成本保持不变。
+每个有效的本地 subprocess service都会贡献一个进程全局 exit listener。成功的 dispose会随服务 effect移除它;失败的 dispose会让它与仍需最终终止的目标一起保留。致命退出放弃宽限、输出排空与进程内停稳证明,以换取宿主消失前发出本地可用的最强终止操作。正常 dispose的保证与成本保持不变。
 
 listener 无法覆盖不执行 JavaScript 的故障。受支持的 Linux terminal 会向[containment decision](2026-08-20-subprocess-native-containment.zh.md)所述的 scope 发出信号;fallback terminal 仍无法发现 provider 首次观察前已经逃逸的后代。

+ 26 - 14
packages/lsp/lsp-stdio/src/index.ts

@@ -281,23 +281,35 @@ class LocalLspProvider implements LspProvider {
       // synchronous get-or-create so every spawned process remains owned by teardown.
       this.assertActive(querySignal)
       let instance = this.instanceFor(workspaceKey, workspace)
-      try {
-        return await instance.query(request, source, querySignal)
-      } catch (error) {
-        // A selected child can have died while idle or fail during the next write. Queries are
-        // read-only, so replace that transport once and retry transparently.
-        if (!instance.isTransportFailure(error)) throw error
-        await instance.dispose()
-        this.evictIfCurrent(workspaceKey, instance)
-        this.assertActive(querySignal)
-        instance = this.instanceFor(workspaceKey, workspace)
-        return await instance.query(request, source, querySignal)
-      } finally {
-        // Reach quiescence before dropping a dead slot; a replacement must survive this ownership check.
+      let canRetryTransport = true
+      for (;;) {
+        const [queryOutcome] = await Promise.allSettled([
+          instance.query(request, source, querySignal),
+        ])
+        let teardownOutcome: PromiseSettledResult<void> | undefined
         if (instance.dead) {
-          await instance.dispose()
+          ;[teardownOutcome] = await Promise.allSettled([instance.dispose()])
+          // A dead instance is never reusable, even when its final quiescence observation fails.
           this.evictIfCurrent(workspaceKey, instance)
         }
+        if (teardownOutcome?.status === 'rejected') {
+          if (queryOutcome.status === 'rejected') {
+            throw new AggregateError(
+              [queryOutcome.reason, teardownOutcome.reason],
+              'LSP operation and teardown failed',
+            )
+          }
+          throw teardownOutcome.reason
+        }
+        if (queryOutcome.status === 'fulfilled') return queryOutcome.value
+        // A selected child can have died while idle or fail during the next write. Queries are
+        // read-only, so replace that transport once and retry transparently after clean disposal.
+        if (!canRetryTransport || !instance.isTransportFailure(queryOutcome.reason)) {
+          throw queryOutcome.reason
+        }
+        canRetryTransport = false
+        this.assertActive(querySignal)
+        instance = this.instanceFor(workspaceKey, workspace)
       }
     })
   }

+ 12 - 28
packages/lsp/lsp-stdio/src/instance.ts

@@ -97,7 +97,7 @@ export class LspInstance {
     const run = abortable(this.queue, signal)
       .then(() => this.runQuery(request, source, signal))
       .catch(async (error: unknown) => {
-        if (this.isTransportFailure(error)) await this.throwAfterTeardown(error)
+        if (this.isTransportFailure(error)) await this.awaitTeardownAttempt()
         throw error
       })
     // Keep the tail alive regardless of this query's outcome so the next caller still serializes. The
@@ -136,7 +136,7 @@ export class LspInstance {
       await abortable(this.ready, signal)
     } catch (error) {
       if (!this.dead) {
-        await this.throwAfterTeardown(error)
+        await this.awaitTeardownAttempt()
       }
       throw error
     }
@@ -152,8 +152,6 @@ export class LspInstance {
 
     const uri = source.fileUrl
     let opened = false
-    let queryFailed = false
-    let queryFailure: unknown
     try {
       /* v8 ignore next -- guards an abort landing between the ready wait and didOpen; not deterministically reproducible. */
       if (signal?.aborted) throw abortError(signal)
@@ -164,15 +162,12 @@ export class LspInstance {
       } catch (error) {
         // A canceled backpressured write or failed stdin leaves the protocol stream unusable before
         // `opened` can arm the didClose cleanup. Teardown here makes the pool evict the instance.
-        await this.throwAfterTeardown(error)
+        await this.awaitTeardownAttempt()
+        throw error
       }
       opened = true
       const payload = await this.sendRequest(request.operation, uri, request.position, signal)
       return this.normalize(request.operation, payload)
-    } catch (error: unknown) {
-      queryFailed = true
-      queryFailure = error
-      throw error
     } finally {
       // A disposed or closed instance (e.g. an aborted request whose server ignored
       // `$/cancelRequest`) is already tearing down; sending didClose would race that teardown and let
@@ -180,20 +175,10 @@ export class LspInstance {
       if (opened && !this.dead) {
         try {
           await this.connection.notify('textDocument/didClose', { textDocument: { uri } })
-        } catch (closeError: unknown) {
+        } catch (_closeFailure: unknown) {
           // A close-write failure does not replace the settled result/error, but the instance can no
-          // longer be trusted: invalidate it and await bounded process termination. If teardown also
-          // fails, every failure remains visible in operation, close, cleanup order.
-          try {
-            await this.startTeardown()
-          } catch (teardownError: unknown) {
-            throw new AggregateError(
-              queryFailed
-                ? [queryFailure, closeError, teardownError]
-                : [closeError, teardownError],
-              'LSP query cleanup failed',
-            )
-          }
+          // longer be trusted. The provider re-awaits this teardown and owns failure reporting.
+          await this.awaitTeardownAttempt()
         }
       }
     }
@@ -242,7 +227,7 @@ export class LspInstance {
             grace.signal.addEventListener('abort', () => { resolve(false) }, { once: true })
           }),
         ])
-        if (!settled) await this.throwAfterTeardown(error)
+        if (!settled) await this.awaitTeardownAttempt()
       } finally {
         grace[Symbol.dispose]()
       }
@@ -293,14 +278,13 @@ export class LspInstance {
     return this.teardownPromise
   }
 
-  /** Preserve an operation failure when teardown also fails. */
-  private async throwAfterTeardown(error: unknown): Promise<never> {
+  /** Await teardown while leaving its memoized failure for provider-level finalization. */
+  private async awaitTeardownAttempt(): Promise<void> {
     try {
       await this.startTeardown()
-    } catch (teardownError: unknown) {
-      throw new AggregateError([error, teardownError], 'LSP operation and teardown failed')
+    } catch (_teardownFailure: unknown) {
+      // LocalLspProvider re-awaits the same teardown and combines it with the query outcome.
     }
-    throw error
   }
 
   private async tearDown(): Promise<void> {

+ 1 - 110
packages/lsp/lsp-stdio/tests/instance.spec.ts

@@ -1,4 +1,4 @@
-import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
+import { afterEach, beforeEach, describe, expect, it } from 'vitest'
 import { readFileSync } from 'node:fs'
 import { mkdtemp, mkdir, readFile, rm, writeFile, realpath } from 'node:fs/promises'
 import { tmpdir } from 'node:os'
@@ -97,24 +97,6 @@ function scriptInstance(script: string, overrides: Partial<InstanceSpec> = {}):
   return instance
 }
 
-interface TestConnection {
-  readonly closed: Promise<void>
-  waitForProcessTreeExit(signal?: AbortSignal): Promise<boolean>
-}
-
-/** Make the instance's final managed-range observation fail deterministically. */
-function rejectProcessTreeWait(instance: LspInstance, failure: Error): TestConnection {
-  const connection = (instance as unknown as { connection: TestConnection }).connection
-  vi.spyOn(connection, 'waitForProcessTreeExit').mockRejectedValue(failure)
-  return connection
-}
-
-/** A failed teardown cannot be disposed again; await process close and remove it from afterEach. */
-async function releaseFailedInstance(instance: LspInstance, connection: TestConnection): Promise<void> {
-  live = live.filter(candidate => candidate !== instance)
-  await connection.closed
-}
-
 /** An inline server that answers initialize + definition and echoes a location. */
 const RESPONDING_SERVER =
   'let b=Buffer.alloc(0);'
@@ -262,52 +244,6 @@ describe('LspInstance query and abort', () => {
     expect(processAlive(pid)).toBe(false)
   })
 
-  it('preserves a request write failure with a managed-range teardown failure', async () => {
-    const operationFailure = new Error('fixture textDocument/definition failure')
-    const teardownFailure = new Error('managed range observation failed')
-    const instance = makeInstance({}, {
-      shutdownTimeoutMs: 100,
-      killGraceMs: 100,
-    }, failingWriter('textDocument/definition', operationFailure))
-    const connection = rejectProcessTreeWait(instance, teardownFailure)
-    try {
-      const failure = await run(instance, 'goToDefinition').then(
-        () => undefined,
-        (error: unknown) => error,
-      )
-      expect(failure).toMatchObject({
-        errors: [operationFailure, teardownFailure],
-        message: 'LSP operation and teardown failed',
-      })
-    } finally {
-      await releaseFailedInstance(instance, connection)
-    }
-  })
-
-  it('preserves an initialization failure with a managed-range teardown failure', async () => {
-    const teardownFailure = new Error('managed range observation failed')
-    const instance = makeInstance({ LSP_FAKE_ENCODING: 'utf-8' }, {
-      shutdownTimeoutMs: 100,
-      killGraceMs: 100,
-    })
-    const connection = rejectProcessTreeWait(instance, teardownFailure)
-    try {
-      const failure = await run(instance, 'goToDefinition').then(
-        () => undefined,
-        (error: unknown) => error,
-      )
-      expect(failure).toBeInstanceOf(AggregateError)
-      const errors = (failure as AggregateError).errors as unknown[]
-      expect(errors).toHaveLength(2)
-      expect(errors[0]).toBeInstanceOf(Error)
-      expect((errors[0] as Error).message).toContain('unsupported position encoding')
-      expect(errors[1]).toBe(teardownFailure)
-      expect((failure as AggregateError).message).toBe('LSP operation and teardown failed')
-    } finally {
-      await releaseFailedInstance(instance, connection)
-    }
-  })
-
   it('rejects when the server lacks the operation capability', async () => {
     const instance = makeInstance({ LSP_FAKE_CAPS: JSON.stringify({ definitionProvider: false }), LSP_FAKE_DEF: 'null' })
     await expect(run(instance, 'goToDefinition')).rejects.toThrow(/does not support goToDefinition/)
@@ -333,51 +269,6 @@ describe('LspInstance query and abort', () => {
     expect(instance.dead).toBe(true)
   })
 
-  it('reports didClose and teardown failures after a settled result', async () => {
-    const closeFailure = new Error('fixture textDocument/didClose failure')
-    const teardownFailure = new Error('managed range observation failed')
-    const instance = makeInstance({
-      LSP_FAKE_DEF: 'null',
-    }, { shutdownTimeoutMs: 100, killGraceMs: 100 }, failingWriter('textDocument/didClose', closeFailure))
-    const connection = rejectProcessTreeWait(instance, teardownFailure)
-    try {
-      const failure = await run(instance, 'goToDefinition').then(
-        () => undefined,
-        (error: unknown) => error,
-      )
-      expect(failure).toMatchObject({
-        errors: [closeFailure, teardownFailure],
-        message: 'LSP query cleanup failed',
-      })
-    } finally {
-      await releaseFailedInstance(instance, connection)
-    }
-  })
-
-  it('reports query, didClose, and teardown failures in lifecycle order', async () => {
-    const closeFailure = new Error('fixture textDocument/didClose failure')
-    const teardownFailure = new Error('managed range observation failed')
-    const instance = makeInstance({
-      LSP_FAKE_ERROR: '1',
-    }, { shutdownTimeoutMs: 100, killGraceMs: 100 }, failingWriter('textDocument/didClose', closeFailure))
-    const connection = rejectProcessTreeWait(instance, teardownFailure)
-    try {
-      const failure = await run(instance, 'goToDefinition').then(
-        () => undefined,
-        (error: unknown) => error,
-      )
-      expect(failure).toBeInstanceOf(AggregateError)
-      const errors = (failure as AggregateError).errors as unknown[]
-      expect(errors).toHaveLength(3)
-      expect(errors[0]).toBeInstanceOf(Error)
-      expect((errors[0] as Error).message).toContain('server refused')
-      expect(errors[1]).toBe(closeFailure)
-      expect(errors[2]).toBe(teardownFailure)
-      expect((failure as AggregateError).message).toBe('LSP query cleanup failed')
-    } finally {
-      await releaseFailedInstance(instance, connection)
-    }
-  })
 })
 
 describe('LspInstance disposal', () => {

+ 95 - 0
packages/lsp/lsp-stdio/tests/lifecycle.spec.ts

@@ -11,6 +11,7 @@ import LocalSubprocessRuntime from '@deepseek-ai/dsh-subprocess-local'
 import LocalFileSystem from '@deepseek-ai/dsh-fs-local'
 import * as LspLocal from '@deepseek-ai/dsh-lsp-stdio'
 import type { LspLocalServerConfig } from '@deepseek-ai/dsh-lsp-stdio'
+import { LspConnection } from '../src/connection.ts'
 
 const fixtureServer = fileURLToPath(new URL('./fixture-server.ts', import.meta.url))
 
@@ -44,10 +45,12 @@ async function mount(
   fakeEnv: Record<string, string> = {},
   overrides: Partial<LspLocalServerConfig> = {},
   captureProvider?: (provider: LspProvider) => void,
+  configureSubprocess?: (ctx: Context) => void,
 ): Promise<Context> {
   const ctx = new Context()
   await ctx.plugin(Lsp)
   await ctx.plugin(LocalSubprocessRuntime)
+  configureSubprocess?.(ctx)
   await ctx.plugin(LocalFileSystem, { cwd: process.cwd() })
   const register = ctx.lsp.registerProvider.bind(ctx.lsp)
   const registrationSpy = captureProvider === undefined
@@ -165,6 +168,98 @@ describe('lsp-stdio end to end over a fake server', () => {
     await ctx.fiber.dispose()
   })
 
+  it('preserves a query failure with final disposal failure and evicts the instance', async () => {
+    const teardownFailure = new Error('managed range observation failed')
+    let provider: LspProvider | undefined
+    let firstSpawn = true
+    let restoreFirstWait: (() => void) | undefined
+    const ctx = await mount(
+      { LSP_FAKE_ENCODING: 'utf-8', LSP_FAKE_DEF: 'null' },
+      { shutdownTimeoutMs: 100, killGraceMs: 100 },
+      (registered) => { provider = registered },
+      (mounted) => {
+        const spawn = mounted.subprocess.spawn.bind(mounted.subprocess)
+        vi.spyOn(mounted.subprocess, 'spawn').mockImplementation((spec) => {
+          const handle = spawn(spec)
+          if (!firstSpawn) return handle
+          firstSpawn = false
+          const waitForExit = handle.waitForExit.bind(handle)
+          const waitSpy = vi.spyOn(handle, 'waitForExit')
+            .mockImplementation(async (signal) => {
+              await waitForExit(signal)
+              throw teardownFailure
+            })
+          restoreFirstWait = () => { waitSpy.mockRestore() }
+          return handle
+        })
+      },
+    )
+    const failure = await ctx.lsp.query(query('goToDefinition')).then(
+      () => undefined,
+      (error: unknown) => error,
+    )
+    expect(failure).toBeInstanceOf(AggregateError)
+    const errors = (failure as AggregateError).errors as unknown[]
+    expect(errors).toHaveLength(2)
+    expect(errors[0]).toBeInstanceOf(Error)
+    expect((errors[0] as Error).message).toContain('unsupported position encoding')
+    expect(errors[1]).toBe(teardownFailure)
+    expect((failure as AggregateError).message).toBe('LSP operation and teardown failed')
+    restoreFirstWait?.()
+    if (provider === undefined) throw new Error('expected lsp-stdio to register a provider')
+    const instances = (provider as unknown as { readonly instances: ReadonlyMap<string, unknown> }).instances
+    expect(instances.size).toBe(0)
+    await expect(ctx.lsp.query(query('goToDefinition'))).rejects.toThrow(/unsupported position encoding/)
+    expect(instances.size).toBe(0)
+    await ctx.fiber.dispose()
+  })
+
+  it('reports final disposal failure after a settled query and evicts the instance', async () => {
+    const closeFailure = new Error('fixture textDocument/didClose failure')
+    const teardownFailure = new Error('managed range observation failed')
+    const notify = Object.getOwnPropertyDescriptor(LspConnection.prototype, 'notify')?.value as LspConnection['notify']
+    const notifySpy = vi.spyOn(LspConnection.prototype, 'notify').mockImplementation(function (this: LspConnection, method, params) {
+      if (method === 'textDocument/didClose') return Promise.reject(closeFailure)
+      return notify.call(this, method, params)
+    })
+    let provider: LspProvider | undefined
+    let restoreFirstWait: (() => void) | undefined
+    const ctx = await mount(
+      { LSP_FAKE_DEF: 'null' },
+      { shutdownTimeoutMs: 100, killGraceMs: 100 },
+      (registered) => { provider = registered },
+      (mounted) => {
+        const spawn = mounted.subprocess.spawn.bind(mounted.subprocess)
+        let firstSpawn = true
+        vi.spyOn(mounted.subprocess, 'spawn').mockImplementation((spec) => {
+          const handle = spawn(spec)
+          if (!firstSpawn) return handle
+          firstSpawn = false
+          const waitForExit = handle.waitForExit.bind(handle)
+          const waitSpy = vi.spyOn(handle, 'waitForExit').mockImplementation(async (signal) => {
+            await waitForExit(signal)
+            throw teardownFailure
+          })
+          restoreFirstWait = () => { waitSpy.mockRestore() }
+          return handle
+        })
+      },
+    )
+    try {
+      await expect(ctx.lsp.query(query('goToDefinition'))).rejects.toBe(teardownFailure)
+      if (provider === undefined) throw new Error('expected lsp-stdio to register a provider')
+      const instances = (provider as unknown as { readonly instances: ReadonlyMap<string, unknown> }).instances
+      expect(instances.size).toBe(0)
+      restoreFirstWait?.()
+      notifySpy.mockRestore()
+      await expect(ctx.lsp.query(query('goToDefinition'))).resolves.toMatchObject({ kind: 'locations' })
+    } finally {
+      restoreFirstWait?.()
+      notifySpy.mockRestore()
+      await ctx.fiber.dispose()
+    }
+  })
+
   it('rejects a server without transient-open sync (None)', async () => {
     const ctx = await mount({ LSP_FAKE_SYNC: '0', LSP_FAKE_DEF: 'null' })
     await expect(ctx.lsp.query(query('goToDefinition'))).rejects.toThrow(/transient textDocument\/didOpen/)

+ 4 - 3
packages/subprocess/subprocess-local/src/index.ts

@@ -114,9 +114,10 @@ export class LocalSubprocessRuntime extends SubprocessRuntime {
       pending.push(terminal.terminate().then(() => { this.terminals.delete(terminal) }))
     }
     const outcomes = await Promise.allSettled(pending)
-    const failures = outcomes.flatMap<unknown>(outcome => outcome.status === 'rejected'
-      ? [outcome.reason as unknown]
-      : [])
+    const failures: unknown[] = []
+    for (const outcome of outcomes) {
+      if (outcome.status === 'rejected') failures.push(outcome.reason)
+    }
     if (failures.length > 0) this.terminateForHostExit()
     if (failures.length === 1) throw failures[0]
     if (failures.length > 1) throw new AggregateError(failures, 'local subprocess teardown failed')

+ 4 - 1
packages/subprocess/subprocess-local/src/spawn.ts

@@ -507,7 +507,10 @@ export function bindManagedProcess(
     // leaking an unhandled rejection when a caller only invokes terminate().
     void observeRangeExit().catch(() => {})
     kill('SIGTERM')
-    graceTimer = setTimeout(() => { kill('SIGKILL') }, spec.graceMs)
+    graceTimer = setTimeout(() => {
+      graceTimer = undefined
+      kill('SIGKILL')
+    }, spec.graceMs)
   }
 
   const terminateForHostExit = (): void => {

+ 31 - 15
packages/subprocess/subprocess-local/tests/managed-spawn.spec.ts

@@ -180,29 +180,45 @@ describe('managed process binding', () => {
     }
   })
 
-  it('retries after a background range-observation rejection', async () => {
-    const wrapper = spawn(process.execPath, ['-e', 'setInterval(() => {}, 1000)'], {
-      stdio: ['ignore', 'pipe', 'pipe'],
-    })
+  it('retries termination after an expired escalation and range-observation rejection', async () => {
+    vi.useFakeTimers()
     const failure = new Error('range observation failed')
+    const firstObservation = Promise.withResolvers<undefined>()
+    const secondObservation = Promise.withResolvers<undefined>()
     const waitForExit = vi.fn()
-      .mockRejectedValueOnce(failure)
-      .mockResolvedValue(undefined)
-    const handle = bindManagedProcess(spec(), {
-      stdin: wrapper.stdin,
-      stdout: wrapper.stdout,
-      stderr: wrapper.stderr,
-      pid: wrapper.pid,
+      .mockImplementationOnce(() => firstObservation.promise)
+      .mockImplementationOnce(() => secondObservation.promise)
+    const signal = vi.fn()
+    const handle = bindManagedProcess({
+      ...spec(),
+      stdio: { stdin: 'ignore', stdout: 'inherit', stderr: 'inherit' },
+    }, {
+      stdin: null,
+      stdout: null,
+      stderr: null,
+      pid: 4242,
       direct: new Promise(() => {}),
-      owner: { signal: vi.fn(), waitForExit },
+      owner: { signal, waitForExit },
     })
     try {
       handle.terminate()
-      await new Promise(resolve => setImmediate(resolve))
-      await expect(handle.waitForExit()).resolves.toBe(true)
+      const firstWait = handle.waitForExit()
+      expect(signal.mock.calls).toEqual([['SIGTERM']])
+      await vi.advanceTimersByTimeAsync(30)
+      expect(signal.mock.calls).toEqual([['SIGTERM'], ['SIGKILL']])
+      firstObservation.reject(failure)
+      await expect(firstWait).rejects.toBe(failure)
+
+      handle.terminate()
+      const secondWait = handle.waitForExit()
+      expect(signal.mock.calls).toEqual([['SIGTERM'], ['SIGKILL'], ['SIGTERM']])
+      await vi.advanceTimersByTimeAsync(30)
+      expect(signal.mock.calls).toEqual([['SIGTERM'], ['SIGKILL'], ['SIGTERM'], ['SIGKILL']])
+      secondObservation.resolve(undefined)
+      await expect(secondWait).resolves.toBe(true)
       expect(waitForExit).toHaveBeenCalledTimes(2)
     } finally {
-      wrapper.kill('SIGKILL')
+      vi.useRealTimers()
     }
   })
 

+ 5 - 4
packages/subprocess/subprocess-local/tests/process-exit.spec.ts

@@ -15,7 +15,8 @@ interface TreeState { root: number; descendant: number }
 
 const repoRoot = fileURLToPath(new URL('../../../../', import.meta.url))
 const hostScript = fileURLToPath(new URL('./fixtures/process-exit-host.ts', import.meta.url))
-const scenarioTimeoutMs = 30_000
+const scenarioTimeoutMs = process.platform === 'win32' ? 60_000 : 30_000
+const testTimeoutMs = scenarioTimeoutMs + 15_000
 
 function processExists(pid: number): boolean {
   try {
@@ -141,7 +142,7 @@ describe('synchronous cleanup on host exit', () => {
     { trigger: 'direct' as const, expectedCode: 23, diagnostic: undefined },
     { trigger: 'uncaught-exception' as const, expectedCode: 1, diagnostic: 'host-exit-uncaught-exception' },
     { trigger: 'unhandled-rejection' as const, expectedCode: 1, diagnostic: 'host-exit-unhandled-rejection' },
-  ])('removes an ordinary managed tree after $trigger', { timeout: 45_000 }, async ({
+  ])('removes an ordinary managed tree after $trigger', { timeout: testTimeoutMs }, async ({
     trigger,
     expectedCode,
     diagnostic,
@@ -154,7 +155,7 @@ describe('synchronous cleanup on host exit', () => {
 
   it.skipIf(process.platform === 'win32')(
     'removes a terminal root and descendant after direct exit',
-    { timeout: 45_000 },
+    { timeout: testTimeoutMs },
     async () => {
       const { outcome } = await runScenario('terminal', 'direct')
       expect(outcome.exitCode).toBe(23)
@@ -162,7 +163,7 @@ describe('synchronous cleanup on host exit', () => {
     },
   )
 
-  it('preserves normal terminate-and-join disposal and removes the exit listener', { timeout: 45_000 }, async () => {
+  it('preserves normal terminate-and-join disposal and removes the exit listener', { timeout: testTimeoutMs }, async () => {
     const { outcome, disposeCounts } = await runScenario('ordinary', 'dispose')
     expect(outcome.exitCode).toBe(0)
     expect(disposeCounts?.listenersAfterLoad).toBe((disposeCounts?.listenersBefore ?? 0) + 1)