Skip to content

Commit 8cc66fe

Browse files
committed
fix: fail fast when benchmark workers exit
Reject worker readiness when startup exits and terminate sibling workers so benchmark failures cannot hang. Validate recovery event records against concrete schemas before assertions. See #2
1 parent 3dd55a7 commit 8cc66fe

5 files changed

Lines changed: 143 additions & 39 deletions

File tree

‎benchmarks/processes.ts‎

Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,59 @@
1+
import type { ChildProcess } from "node:child_process"
2+
3+
export function waitForWorkerReady(worker: ChildProcess): Promise<void> {
4+
const priorExit = workerExitError(worker, "before ready")
5+
if (priorExit) return Promise.reject(priorExit)
6+
7+
return new Promise((resolvePromise, reject) => {
8+
const cleanup = () => {
9+
worker.off("error", onError)
10+
worker.off("exit", onExit)
11+
worker.off("message", onMessage)
12+
}
13+
const onError = (error: Error) => {
14+
cleanup()
15+
reject(error)
16+
}
17+
const onExit = (code: number | null, signal: NodeJS.Signals | null) => {
18+
cleanup()
19+
reject(workerExitError({ exitCode: code, signalCode: signal }, "before ready"))
20+
}
21+
const onMessage = (message: string) => {
22+
if (message !== "ready") return
23+
cleanup()
24+
resolvePromise()
25+
}
26+
worker.once("error", onError)
27+
worker.once("exit", onExit)
28+
worker.on("message", onMessage)
29+
})
30+
}
31+
32+
export function waitForWorkerExit(worker: ChildProcess): Promise<void> {
33+
const priorExit = workerExitError(worker)
34+
if (priorExit) {
35+
return worker.exitCode === 0 ? Promise.resolve() : Promise.reject(priorExit)
36+
}
37+
38+
return new Promise((resolvePromise, reject) => {
39+
worker.once("error", reject)
40+
worker.once("exit", (code, signal) => {
41+
if (code === 0) {
42+
resolvePromise()
43+
return
44+
}
45+
reject(workerExitError({ exitCode: code, signalCode: signal }))
46+
})
47+
})
48+
}
49+
50+
function workerExitError(
51+
worker: Pick<ChildProcess, "exitCode" | "signalCode">,
52+
phase?: string,
53+
): Error | undefined {
54+
if (worker.exitCode === null && worker.signalCode === null) return undefined
55+
const suffix = phase ? ` ${phase}` : ""
56+
return new Error(
57+
`benchmark worker exited${suffix} with code ${worker.exitCode} and signal ${worker.signalCode}`,
58+
)
59+
}

‎benchmarks/run.ts‎

Lines changed: 13 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import { fileURLToPath } from "node:url"
66
import { fork, type ChildProcess } from "node:child_process"
77
import { performance } from "node:perf_hooks"
88
import type { MessageReference } from "solid-objects"
9+
import { waitForWorkerExit, waitForWorkerReady } from "./processes.ts"
910
import { BenchmarkCounter, benchmarkRuntime, type BenchmarkDatabase } from "./shared.ts"
1011

1112
type Shape = "warm-hot" | "warm-many" | "cold-many"
@@ -82,7 +83,7 @@ try {
8283
} finally {
8384
shutdown.abort()
8485
await running
85-
const workerExits = workers.map(waitForExit)
86+
const workerExits = workers.map(waitForWorkerExit)
8687
for (const worker of workers) worker.send("stop")
8788
await Promise.all(workerExits)
8889
await runtime.testing.reset()
@@ -232,18 +233,17 @@ async function spawnWorkers(options: {
232233
{ stdio: ["ignore", "ignore", "inherit", "ipc"] },
233234
),
234235
)
235-
await Promise.all(
236-
workers.map(
237-
(worker) =>
238-
new Promise<void>((resolvePromise, reject) => {
239-
worker.once("error", reject)
240-
worker.on("message", (message) => {
241-
if (message === "ready") resolvePromise()
242-
})
243-
}),
244-
),
245-
)
246-
return workers
236+
try {
237+
await Promise.all(workers.map(waitForWorkerReady))
238+
return workers
239+
} catch (error) {
240+
const workerExits = workers.map(waitForWorkerExit)
241+
for (const worker of workers) {
242+
if (worker.exitCode === null && worker.signalCode === null) worker.kill()
243+
}
244+
await Promise.allSettled(workerExits)
245+
throw error
246+
}
247247
}
248248

249249
async function readDatabaseVersion(runtime: ReturnType<typeof benchmarkRuntime>): Promise<string> {
@@ -263,24 +263,6 @@ async function readDatabaseVersion(runtime: ReturnType<typeof benchmarkRuntime>)
263263
})
264264
}
265265

266-
function waitForExit(child: ChildProcess): Promise<void> {
267-
if (child.exitCode !== null) {
268-
return child.exitCode === 0
269-
? Promise.resolve()
270-
: Promise.reject(new Error(`benchmark worker exited ${child.exitCode}`))
271-
}
272-
if (child.signalCode !== null) {
273-
return Promise.reject(new Error(`benchmark worker exited with ${child.signalCode}`))
274-
}
275-
return new Promise((resolvePromise, reject) => {
276-
child.once("error", reject)
277-
child.once("exit", (code) => {
278-
if (code === 0) resolvePromise()
279-
else reject(new Error(`benchmark worker exited ${code}`))
280-
})
281-
})
282-
}
283-
284266
function percentile(values: number[], percent: number): number {
285267
const ordered = [...values].sort((left, right) => left - right)
286268
const index = Math.max(0, Math.ceil((percent / 100) * ordered.length) - 1)

‎examples/failure-recovery/demo.ts‎

Lines changed: 53 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,18 @@ interface WorkerMessage {
1515
processed?: number
1616
}
1717

18+
interface SerializationEvent {
19+
event: "start" | "finish"
20+
messageId: string
21+
at: number
22+
}
23+
24+
interface ExternalEffectEvent {
25+
messageId: string
26+
attempt: number
27+
processId: number
28+
}
29+
1830
const directory = await mkdtemp(join(tmpdir(), "solid-objects-recovery-"))
1931
const databasePath = join(directory, "state.sqlite3")
2032
const runtime = createRuntime({
@@ -58,7 +70,10 @@ async function proveSerialization(): Promise<{ finalState: number; overlap: fals
5870
const workers = [spawnWorker(), spawnWorker()]
5971
await Promise.all(workers.map(({ finished }) => finished))
6072
await Promise.all(messages.map((message) => message.result()))
61-
const events = await jsonLines(join(controlDirectory, "serialization.jsonl"))
73+
const events = await readJsonLines(
74+
join(controlDirectory, "serialization.jsonl"),
75+
parseSerializationEvent,
76+
)
6277
assert.equal(events.length, 4)
6378
const starts = events.filter((event) => event.event === "start")
6479
const finishes = events.filter((event) => event.event === "finish")
@@ -121,7 +136,10 @@ async function recoveryResult(options: {
121136
const stored = await runtime.repository.findMessage(options.message.id)
122137
const attempts = Number(stored?.attempt_count)
123138
const snapshot = await options.reference.snapshot()
124-
const effects = await jsonLines(join(options.controlDirectory, "external-effects.jsonl"))
139+
const effects = await readJsonLines(
140+
join(options.controlDirectory, "external-effects.jsonl"),
141+
parseExternalEffectEvent,
142+
)
125143
assert.equal(attempts, 2)
126144
assert.equal(snapshot.count, 1)
127145
assert.equal(effects.length, 2)
@@ -175,12 +193,39 @@ function spawnWorker(): {
175193
}
176194
}
177195

178-
async function jsonLines(path: string): Promise<Array<Record<string, unknown>>> {
179-
return (await readFile(path, "utf8"))
180-
.trim()
181-
.split("\n")
182-
.filter(Boolean)
183-
.map((line) => JSON.parse(line) as Record<string, unknown>)
196+
async function readJsonLines<Value>(
197+
path: string,
198+
parse: (line: string) => Value,
199+
): Promise<Value[]> {
200+
return (await readFile(path, "utf8")).trim().split("\n").filter(Boolean).map(parse)
201+
}
202+
203+
function parseSerializationEvent(line: string): SerializationEvent {
204+
const event = JSON.parse(line) as Partial<SerializationEvent>
205+
if (
206+
(event.event !== "start" && event.event !== "finish") ||
207+
typeof event.messageId !== "string" ||
208+
typeof event.at !== "number"
209+
) {
210+
throw new TypeError("invalid serialization event")
211+
}
212+
return { event: event.event, messageId: event.messageId, at: event.at }
213+
}
214+
215+
function parseExternalEffectEvent(line: string): ExternalEffectEvent {
216+
const event = JSON.parse(line) as Partial<ExternalEffectEvent>
217+
if (
218+
typeof event.messageId !== "string" ||
219+
typeof event.attempt !== "number" ||
220+
typeof event.processId !== "number"
221+
) {
222+
throw new TypeError("invalid external effect event")
223+
}
224+
return {
225+
messageId: event.messageId,
226+
attempt: event.attempt,
227+
processId: event.processId,
228+
}
184229
}
185230

186231
async function wait(milliseconds: number): Promise<void> {

‎test/benchmark-processes.test.ts‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,17 @@
1+
import { fork } from "node:child_process"
2+
import { fileURLToPath } from "node:url"
3+
import { describe, expect, it } from "vitest"
4+
import { waitForWorkerReady } from "../benchmarks/processes.js"
5+
6+
describe("benchmark worker processes", () => {
7+
it("rejects when a worker exits before signaling readiness", async () => {
8+
const worker = fork(
9+
fileURLToPath(new URL("./fixtures/benchmark-worker-exit.mjs", import.meta.url)),
10+
{ stdio: ["ignore", "ignore", "ignore", "ipc"] },
11+
)
12+
13+
await expect(waitForWorkerReady(worker)).rejects.toThrow(
14+
"benchmark worker exited before ready with code 7 and signal null",
15+
)
16+
})
17+
})
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
process.exitCode = 7

0 commit comments

Comments
 (0)