Просмотр исходного кода

fix(subprocess): close native runner lifecycle gaps

pku-xht 1 месяц назад
Родитель
Сommit
81374783d3

+ 8 - 2
packages/subprocess/subprocess-local/src/managed-owner.ts

@@ -21,12 +21,18 @@ export interface ManagedProcessLaunch {
 }
 
 /**
- * Observe wrapper close from the moment it is spawned.
+ * Observe wrapper close from the moment it is spawned and contain its error
+ * event while the runner-result path converts launch failures into rejection.
  * @param child - direct child or native wrapper.
  * @returns promise settled by the ChildProcess close event.
  */
 export function observeChildClose(child: ChildProcess): Promise<void> {
-  return new Promise((resolve) => { child.once('close', () => { resolve() }) })
+  return new Promise((resolve) => {
+    child.once('error', () => {
+      // runnerDirectResult reports the wrapper failure through the handle.
+    })
+    child.once('close', () => { resolve() })
+  })
 }
 
 /**

+ 15 - 7
packages/subprocess/subprocess-local/src/runner-launch.ts

@@ -1,7 +1,8 @@
 /** Parent-side launch and direct-result transport for native runners. */
 
 import type { ChildProcess, StdioOptions } from 'node:child_process'
-import { existsSync, readFileSync } from 'node:fs'
+import { readFileSync } from 'node:fs'
+import { extname } from 'node:path'
 import { fileURLToPath } from 'node:url'
 import { setTimeout as sleepMs } from 'node:timers/promises'
 import type { SubprocessOutcome, SubprocessSpawnSpec } from '@deepseek-ai/dsh-subprocess'
@@ -21,15 +22,18 @@ const RUNNER_EVENT_POLL_MS = 100
 const PACKAGED_RUNNER_ARG = '--dsh-internal-subprocess-runner'
 
 /**
- * Resolve the built runner in production or its source entry in repository execution.
+ * Resolve the runner entry from the current module's source or built plane.
+ * @param moduleUrl - executing module URL; defaults to this module.
  * @returns Node executable and runner argv prefix.
  */
-export function spawnRunnerInvocation(): string[] {
+export function spawnRunnerInvocation(moduleUrl = import.meta.url): string[] {
   if ('pkg' in process) return [process.execPath, PACKAGED_RUNNER_ARG]
+  if (extname(fileURLToPath(moduleUrl)) === '.ts') {
+    const sourceEntry = fileURLToPath(import.meta.resolve('@deepseek-ai/dsh-subprocess-local/src/spawn-runner.ts'))
+    return [process.execPath, '--import', 'tsx/esm', sourceEntry]
+  }
   const builtEntry = fileURLToPath(import.meta.resolve('@deepseek-ai/dsh-subprocess-local/spawn-runner'))
-  if (existsSync(builtEntry)) return [process.execPath, builtEntry]
-  const sourceEntry = fileURLToPath(import.meta.resolve('@deepseek-ai/dsh-subprocess-local/src/spawn-runner.ts'))
-  return [process.execPath, '--import', 'tsx/esm', sourceEntry]
+  return [process.execPath, builtEntry]
 }
 
 /**
@@ -114,6 +118,10 @@ async function waitForDirectResult(
   let wrapperClosed = false
   void closed.then(() => { wrapperClosed = true })
   for (;;) {
+    // A read started before close may return a stale snapshot after close has
+    // become visible. Only a read started after close can prove no terminal
+    // event was written before the runner exited.
+    const closedBeforeRead = wrapperClosed
     const events = await readRunnerEventsAsync(files.eventsPath)
     for (const event of events.slice(seen)) {
       if (event.type === 'exit') return { exitCode: event.exitCode, signal: event.signal }
@@ -121,7 +129,7 @@ async function waitForDirectResult(
     }
     seen = Math.max(seen, events.length, initial.length)
     // oxlint-disable-next-line typescript/no-unnecessary-condition -- child close mutates this flag asynchronously.
-    if (wrapperClosed) {
+    if (closedBeforeRead) {
       const known = missingResult?.()
       if (known !== undefined) return known
       throw new Error('native subprocess runner exited without a direct-command result')

+ 60 - 26
packages/subprocess/subprocess-local/tests/spawn-runner.spec.ts

@@ -1,4 +1,4 @@
-import { spawnSync } from 'node:child_process'
+import { spawn, spawnSync } from 'node:child_process'
 import type { ChildProcess } from 'node:child_process'
 import { existsSync, statSync, writeFileSync } from 'node:fs'
 import { join } from 'node:path'
@@ -12,6 +12,7 @@ import {
   runnerStdio,
   spawnRunnerInvocation,
 } from '../src/runner-launch.ts'
+import { observeChildClose } from '../src/managed-owner.ts'
 import {
   appendRunnerEvent,
   cleanupRunnerFiles,
@@ -59,31 +60,13 @@ function runRunner(invocation: string[], requestPath: string, eventsPath: string
 }
 
 describe('spawn runner transport', () => {
-  it('selects built and source runner entries according to artifact availability', async () => {
-    vi.resetModules()
-    vi.doMock('node:fs', async importOriginal => ({
-      ...await importOriginal<typeof import('node:fs')>(),
-      existsSync: () => true,
-    }))
-    try {
-      const built = await import('../src/runner-launch.ts')
-      expect(built.spawnRunnerInvocation()).toEqual([process.execPath, builtEntry])
-    } finally {
-      vi.doUnmock('node:fs')
-      vi.resetModules()
-    }
-
-    vi.doMock('node:fs', async importOriginal => ({
-      ...await importOriginal<typeof import('node:fs')>(),
-      existsSync: () => false,
-    }))
-    try {
-      const source = await import('../src/runner-launch.ts')
-      expect(source.spawnRunnerInvocation()).toEqual(sourceInvocation)
-    } finally {
-      vi.doUnmock('node:fs')
-      vi.resetModules()
-    }
+  it('selects the runner entry from the current execution plane', () => {
+    expect(spawnRunnerInvocation(new URL('../lib/index.js', import.meta.url).href)).toEqual([
+      process.execPath,
+      builtEntry,
+    ])
+    expect(spawnRunnerInvocation(new URL('../src/runner-launch.ts', import.meta.url).href)).toEqual(sourceInvocation)
+    expect(spawnRunnerInvocation()).toEqual(sourceInvocation)
   })
 
   it('re-enters a packaged executable through its private runner dispatch', () => {
@@ -309,6 +292,57 @@ describe('spawn runner transport', () => {
     }
   })
 
+  it('requires an event snapshot started after wrapper close before reporting a missing result', async () => {
+    const staleRead = Promise.withResolvers<Awaited<ReturnType<typeof readRunnerEventsAsync>>>()
+    let readCount = 0
+    vi.resetModules()
+    vi.doMock('../src/runner-protocol.ts', async (importOriginal) => {
+      const actual = await importOriginal<typeof import('../src/runner-protocol.ts')>()
+      return {
+        ...actual,
+        readRunnerEventsAsync: vi.fn(async (eventsPath: string) => {
+          readCount += 1
+          if (readCount === 1) return staleRead.promise
+          return actual.readRunnerEvents(eventsPath)
+        }),
+      }
+    })
+    const files = createRunnerFiles({ argv: ['node'], cwd: '.', env: {} })
+    try {
+      appendRunnerEvent(files.eventsPath, { type: 'started', pid: 456 })
+      const closed = Promise.withResolvers<undefined>()
+      const isolated = await import('../src/runner-launch.ts')
+      const result = isolated.runnerDirectResult(fakeChild(123), files, closed.promise)
+      expect(readCount).toBe(1)
+      closed.resolve()
+      await Promise.resolve()
+      appendRunnerEvent(files.eventsPath, { type: 'exit', exitCode: 0, signal: null })
+      staleRead.resolve([{ type: 'started', pid: 456 }])
+      await expect(result.direct).resolves.toEqual({ exitCode: 0, signal: null })
+      expect(readCount).toBe(2)
+    } finally {
+      cleanupRunnerFiles(files)
+      vi.doUnmock('../src/runner-protocol.ts')
+      vi.resetModules()
+    }
+  })
+
+  it('contains wrapper spawn errors while publishing the runner startup rejection', async () => {
+    const files = createRunnerFiles({ argv: ['node'], cwd: '.', env: {} })
+    try {
+      const child = spawn(`missing-dsh-native-runner-${String(process.pid)}-${String(Date.now())}`, [], {
+        stdio: 'ignore',
+      })
+      const closed = observeChildClose(child)
+      const result = runnerDirectResult(child, files, closed)
+      expect(result.pid).toBe(-1)
+      await expect(result.direct).rejects.toThrow('runner failed to start')
+      await expect(closed).resolves.toBeUndefined()
+    } finally {
+      cleanupRunnerFiles(files)
+    }
+  })
+
   it('reports runner startup failure and handshake timeout without leaking request files', async () => {
     const missingChild = createRunnerFiles({ argv: ['node'], cwd: '.', env: {} })
     const missingResult = runnerDirectResult(fakeChild(undefined), missingChild, new Promise<void>(() => {}))