Skip to content

Commit 1216b5a

Browse files
authored
fix: don't lose worker output on teardown, deflake timing-sensitive tests (#10842)
1 parent d487659 commit 1216b5a

9 files changed

Lines changed: 124 additions & 30 deletions

File tree

‎packages/vitest/src/node/pools/workers/forksWorker.ts‎

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ import type { Writable } from 'node:stream'
33
import type { PoolOptions, PoolWorker, WorkerRequest } from '../types'
44
import { fork } from 'node:child_process'
55
import { resolve } from 'node:path'
6+
import { streamFlushed } from './utils'
67

78
const SIGKILL_TIMEOUT = 500 /** jest does 500ms by default, let's follow it */
89

@@ -49,14 +50,16 @@ export class ForksPoolWorker implements PoolWorker {
4950
serialization: 'advanced',
5051
})
5152

53+
// `end: false`: the logger streams are shared by every worker, so one
54+
// ending worker stream must not end them for everyone else
5255
if (this._fork.stdout) {
5356
this.stdout.setMaxListeners(1 + this.stdout.getMaxListeners())
54-
this._fork.stdout.pipe(this.stdout)
57+
this._fork.stdout.pipe(this.stdout, { end: false })
5558
}
5659

5760
if (this._fork.stderr) {
5861
this.stderr.setMaxListeners(1 + this.stderr.getMaxListeners())
59-
this._fork.stderr.pipe(this.stderr)
62+
this._fork.stderr.pipe(this.stderr, { end: false })
6063
}
6164
}
6265

@@ -87,12 +90,14 @@ export class ForksPoolWorker implements PoolWorker {
8790
clearTimeout(sigkillTimeout)
8891

8992
if (fork.stdout) {
90-
fork.stdout?.unpipe(this.stdout)
93+
await streamFlushed(fork.stdout)
94+
fork.stdout.unpipe(this.stdout)
9195
this.stdout.setMaxListeners(this.stdout.getMaxListeners() - 1)
9296
}
9397

9498
if (fork.stderr) {
95-
fork.stderr?.unpipe(this.stderr)
99+
await streamFlushed(fork.stderr)
100+
fork.stderr.unpipe(this.stderr)
96101
this.stderr.setMaxListeners(this.stderr.getMaxListeners() - 1)
97102
}
98103

‎packages/vitest/src/node/pools/workers/threadsWorker.ts‎

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import type { Writable } from 'node:stream'
22
import type { PoolOptions, PoolWorker, WorkerRequest } from '../types'
33
import { resolve } from 'node:path'
44
import { Worker } from 'node:worker_threads'
5+
import { streamFlushed } from './utils'
56

67
/** @experimental */
78
export class ThreadsPoolWorker implements PoolWorker {
@@ -46,20 +47,31 @@ export class ThreadsPoolWorker implements PoolWorker {
4647
stderr: true,
4748
})
4849

50+
// `end: false`: the logger streams are shared by every worker, so one
51+
// ending worker stream must not end them for everyone else
4952
this.stdout.setMaxListeners(1 + this.stdout.getMaxListeners())
50-
this._thread.stdout.pipe(this.stdout)
53+
this._thread.stdout.pipe(this.stdout, { end: false })
5154

5255
this.stderr.setMaxListeners(1 + this.stderr.getMaxListeners())
53-
this._thread.stderr.pipe(this.stderr)
56+
this._thread.stderr.pipe(this.stderr, { end: false })
5457
}
5558

5659
async stop(): Promise<void> {
57-
await this.thread.terminate()
58-
59-
this._thread?.stdout?.unpipe(this.stdout)
60+
const thread = this.thread
61+
// `terminate()` makes node drain the stdio still queued on the worker's
62+
// message port into these readables; keep the pipes attached until the
63+
// streams end so late output still reaches the logger streams
64+
const flushed = Promise.all([
65+
streamFlushed(thread.stdout),
66+
streamFlushed(thread.stderr),
67+
])
68+
await thread.terminate()
69+
await flushed
70+
71+
thread.stdout.unpipe(this.stdout)
6072
this.stdout.setMaxListeners(this.stdout.getMaxListeners() - 1)
6173

62-
this._thread?.stderr?.unpipe(this.stderr)
74+
thread.stderr.unpipe(this.stderr)
6375
this.stderr.setMaxListeners(this.stderr.getMaxListeners() - 1)
6476

6577
this._thread = undefined
Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
import type { Readable } from 'node:stream'
2+
3+
// After a worker dies, its remaining stdio is drained into the parent-side
4+
// readables asynchronously. Waiting for `end`/`close` before unpiping ensures
5+
// the tail of the output still reaches the shared logger streams.
6+
export function streamFlushed(stream: Readable): Promise<unknown> {
7+
if (stream.readableEnded || stream.destroyed) {
8+
return Promise.resolve()
9+
}
10+
return new Promise((resolve) => {
11+
stream.once('end', resolve)
12+
stream.once('close', resolve)
13+
})
14+
}

‎packages/vitest/src/runtime/workers/init.ts‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,27 @@ const __vitest_worker_response__ = true
4949
const memoryUsage = process.memoryUsage.bind(process)
5050
let reportMemory = false
5151

52+
// In worker threads stdio is proxied to the parent over a MessagePort with a
53+
// backpressure protocol: a chunk stays buffered inside the worker until the
54+
// parent acks the previous one. The pool starts `runner.stop()` as soon as it
55+
// receives `testfileFinished`, and `thread.terminate()` halts the worker before
56+
// buffered chunks are ever posted, losing output. An empty write's callback
57+
// only fires after every previously buffered chunk has been acked, so awaiting
58+
// it before signaling completion guarantees the output reached the parent.
59+
// A cheap no-op for forks, where stdio goes through OS pipes.
60+
function flushStdio(): Promise<unknown> {
61+
const flush = (stream: NodeJS.WriteStream) =>
62+
new Promise((resolve) => {
63+
try {
64+
stream.write('', () => resolve(undefined))
65+
}
66+
catch {
67+
resolve(undefined)
68+
}
69+
})
70+
return Promise.all([flush(process.stdout), flush(process.stderr)])
71+
}
72+
5273
let traces!: Traces
5374

5475
/** @experimental */
@@ -168,6 +189,8 @@ export function init(worker: Options): void {
168189
)
169190
const error = await runPromise
170191

192+
await flushStdio()
193+
171194
send({
172195
type: 'testfileFinished',
173196
__vitest_worker_response__,
@@ -226,6 +249,8 @@ export function init(worker: Options): void {
226249
)
227250
const error = await runPromise
228251

252+
await flushStdio()
253+
229254
send({
230255
type: 'testfileFinished',
231256
__vitest_worker_response__,
@@ -275,11 +300,15 @@ export function init(worker: Options): void {
275300

276301
persistCompileCache()
277302

303+
await flushStdio()
304+
278305
send({ type: 'stopped', error, __vitest_worker_response__ })
279306
}
280307
catch (error) {
281308
persistCompileCache()
282309

310+
await flushStdio()
311+
283312
send({ type: 'stopped', error: serializeError(error), __vitest_worker_response__ })
284313
}
285314

‎test/browser/fixtures/server-url/vitest.config.ts‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,10 @@ export default defineConfig({
2323
!!process.env.TEST_HTTPS && basicSsl(),
2424
],
2525
test: {
26-
api: process.env.TEST_HTTPS ? 51122 : 51133,
26+
// below the OS ephemeral port range (32768+ on Linux): a kernel-assigned
27+
// outbound socket holding the fixed port would make Vite silently bind
28+
// port+1 and fail the exact-port assertions
29+
api: process.env.TEST_HTTPS ? 31122 : 31133,
2730
browser: {
2831
enabled: true,
2932
provider: configuredProvider,

‎test/browser/specs/server-url.test.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ test('server-url http', async () => {
1313
const url = ctx?.projects[0].vite.resolvedUrls?.local[0]
1414
expect(stderr).toBe('')
1515
expect.assert(url)
16-
expect(new URL(url).port).toBe('51133')
16+
expect(new URL(url).port).toBe('31133')
1717
})
1818

1919
test('server-url https', async () => {
@@ -25,6 +25,6 @@ test('server-url https', async () => {
2525
expect(stderr).toBe('')
2626
const url = ctx?.projects[0].vite.resolvedUrls?.local[0]
2727
expect.assert(url)
28-
expect(new URL(url).port).toBe('51122')
28+
expect(new URL(url).port).toBe('31122')
2929
expect(stdout).toReportSummaryTestFiles({ passed: instances.length })
3030
})

‎test/e2e/test/concurrent.test.ts‎

Lines changed: 28 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,14 @@ import { runInlineTests } from '../../test-utils'
1010
// <- *
1111
// <------
1212

13-
const deadlockSource = `
13+
// In the deadlocking variant "c" resolves the deadlock only after "b" reported
14+
// its timeout: whether a deadlocked test is reported as timed out depends on
15+
// its own elapsed time the moment the deadlock resolves, and "b" starts its
16+
// clock a few event-loop turns after "a", so an unconditional resolve can
17+
// release "b" while it is still within its own budget. The passing variant
18+
// must not gate: nothing fails there, so the gate would never open.
19+
function deadlockSource(gateOnTimeout: boolean) {
20+
return `
1421
import { describe, expect, test } from 'vitest'
1522
import { createDefer } from '@vitest/utils/helpers'
1623
@@ -20,14 +27,18 @@ describe.concurrent('wrapper', () => {
2027
createDefer<void>(),
2128
createDefer<void>(),
2229
]
30+
const bTimedOut = createDefer<void>()
2331
2432
test('a', async () => {
2533
expect(1).toBe(1)
2634
defers[0].resolve()
2735
await defers[2]
2836
})
2937
30-
test('b', async () => {
38+
test('b', async ({ onTestFailed }) => {
39+
onTestFailed(() => {
40+
bTimedOut.resolve()
41+
})
3142
expect(1).toBe(1)
3243
await defers[0]
3344
defers[1].resolve()
@@ -37,14 +48,16 @@ describe.concurrent('wrapper', () => {
3748
test('c', async () => {
3849
expect(1).toBe(1)
3950
await defers[1]
51+
${gateOnTimeout ? 'await bTimedOut' : ''}
4052
defers[2].resolve()
4153
})
4254
})
4355
`
56+
}
4457

4558
test('deadlocks with insufficient maxConcurrency', async () => {
4659
const { errorTree } = await runInlineTests({
47-
'basic.test.ts': deadlockSource,
60+
'basic.test.ts': deadlockSource(true),
4861
}, {
4962
maxConcurrency: 2,
5063
testTimeout: 500,
@@ -74,7 +87,7 @@ test('deadlocks with insufficient maxConcurrency', async () => {
7487

7588
test('passes when maxConcurrency is high enough', async () => {
7689
const { stderr, errorTree } = await runInlineTests({
77-
'basic.test.ts': deadlockSource,
90+
'basic.test.ts': deadlockSource(false),
7891
}, {
7992
maxConcurrency: 3,
8093
})
@@ -93,7 +106,8 @@ test('passes when maxConcurrency is high enough', async () => {
93106
`)
94107
})
95108

96-
const suiteDeadlockSource = `
109+
function suiteDeadlockSource(gateOnTimeout: boolean) {
110+
return `
97111
import { describe, expect, test } from 'vitest'
98112
import { createDefer } from '@vitest/utils/helpers'
99113
@@ -103,6 +117,7 @@ describe.concurrent('wrapper', () => {
103117
createDefer<void>(),
104118
createDefer<void>(),
105119
]
120+
const bTimedOut = createDefer<void>()
106121
107122
describe('1st suite', () => {
108123
test('a', async () => {
@@ -111,7 +126,10 @@ describe.concurrent('wrapper', () => {
111126
await defers[2]
112127
})
113128
114-
test('b', async () => {
129+
test('b', async ({ onTestFailed }) => {
130+
onTestFailed(() => {
131+
bTimedOut.resolve()
132+
})
115133
expect(1).toBe(1)
116134
await defers[0]
117135
defers[1].resolve()
@@ -123,15 +141,17 @@ describe.concurrent('wrapper', () => {
123141
test('c', async () => {
124142
expect(1).toBe(1)
125143
await defers[1]
144+
${gateOnTimeout ? 'await bTimedOut' : ''}
126145
defers[2].resolve()
127146
})
128147
})
129148
})
130149
`
150+
}
131151

132152
test('suite deadlocks with insufficient maxConcurrency', async () => {
133153
const { errorTree } = await runInlineTests({
134-
'basic.test.ts': suiteDeadlockSource,
154+
'basic.test.ts': suiteDeadlockSource(true),
135155
}, {
136156
maxConcurrency: 2,
137157
testTimeout: 500,
@@ -162,7 +182,7 @@ test('suite deadlocks with insufficient maxConcurrency', async () => {
162182

163183
test('suite passes when maxConcurrency is high enough', async () => {
164184
const { stderr, errorTree } = await runInlineTests({
165-
'basic.test.ts': suiteDeadlockSource,
185+
'basic.test.ts': suiteDeadlockSource(false),
166186
}, {
167187
maxConcurrency: 3,
168188
})

‎test/e2e/test/reporters/import-durations.test.ts‎

Lines changed: 9 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -64,19 +64,21 @@ describe('import durations', () => {
6464
}, 40000)
6565

6666
it('should handle tests with no imports gracefully', async () => {
67-
const { exitCode, ctx } = await runVitest({
67+
const { ctx } = await runVitest({
6868
root,
6969
include: ['**/ok.test.ts'],
7070
experimental: { importDurations: { limit: 10 } },
7171
})
7272

73-
expect(exitCode).toBe(0)
74-
7573
const capturedFiles = ctx!.state.getFiles()
7674

7775
expect(capturedFiles).toHaveLength(1)
7876

7977
const file = capturedFiles[0]
78+
// assert on the run's own state, not `exitCode`: `process.exitCode` is
79+
// process-global, and with `isolate: false` a stray unhandled rejection
80+
// leaked by an earlier test file in this worker flips it to 1
81+
expect(file.result?.state).toBe('pass')
8082
expect(file.importDurations).toBeDefined()
8183
expect(file.importDurations?.[file.filepath].totalTime).toBeGreaterThanOrEqual(0)
8284
expect(file.importDurations?.[file.filepath].selfTime).toBeGreaterThanOrEqual(0)
@@ -172,7 +174,7 @@ describe('import durations', () => {
172174

173175
it('should fail when failOnDanger is enabled and threshold exceeded', async () => {
174176
// With default danger threshold (500ms), should NOT fail (imports are ~75ms)
175-
const { exitCode: exitCodeDefault, stderr: stderrDefault } = await runVitest({
177+
const { ctx: ctxDefault, stderr: stderrDefault } = await runVitest({
176178
root,
177179
include: ['**/import-durations.test.ts'],
178180
experimental: {
@@ -182,7 +184,9 @@ describe('import durations', () => {
182184
},
183185
})
184186

185-
expect(exitCodeDefault).toBe(0)
187+
// see "should handle tests with no imports gracefully" for why `exitCode`
188+
// is not asserted here
189+
expect(ctxDefault!.state.getFiles()[0]?.result?.state).toBe('pass')
186190
expect(stderrDefault).not.toContain('exceeded the danger threshold')
187191

188192
// With lower danger threshold (50ms), should fail (imports are ~75ms > 50ms)

0 commit comments

Comments
 (0)