Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions .changeset/childprocess-process-group-wait.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
---
"@effect/platform-node-shared": patch
"effect": patch
---

Wait for Node child process groups to exit during scoped release and `kill`.

After signalling a process group, both operations now wait for its leader and descendants. Without `forceKillAfter`, the wait is limited to one second and never escalates. With `forceKillAfter`, the group receives `SIGKILL` at the deadline, followed by a final wait of up to one second. Native timers keep escalation working under a `TestClock`, and cleanup no longer depends on stdio closing.

`exitCode` and `isRunning` remain tied to the leader's exit, and a leader that already exited successfully still leaves its group untouched. Process group checks count zombies, so cleanup may wait for the full bound under a non-reaping PID 1.
6 changes: 4 additions & 2 deletions packages/effect/src/unstable/process/ChildProcess.ts
Original file line number Diff line number Diff line change
Expand Up @@ -247,8 +247,10 @@ export interface KillOptions {
/**
* The duration of time to wait after the child process has been terminated
* before forcefully killing the child process by sending it the `"SIGKILL"`
* signal. Defaults to `undefined`, which means that no timeout will be
* enforced by default.
* signal. Defaults to `undefined`, so `"SIGKILL"` is never sent.
*
* The spawner decides whether to terminate descendants and how long to wait.
* See its platform module documentation, such as `NodeChildProcessSpawner`.
*/
readonly forceKillAfter?: Duration.Input | undefined
}
Expand Down
122 changes: 79 additions & 43 deletions packages/platform/node-shared/src/NodeChildProcessSpawner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,17 @@
* process at a time, wiring the selected source stream (`stdout`, `stderr`,
* `all`, or `fdN`) to the destination `stdin` or `fdN`.
*
* Scoped release and `kill` wait for the signalled process group. Without
* `forceKillAfter`, the wait is limited to one second and never escalates.
* With it, cleanup may take the configured duration plus a final one-second
* wait after `SIGKILL`. Zombie descendants can consume either full bound. On
* Windows, `taskkill` terminates the tree and only the leader's exit is awaited.
*
* @since 4.0.0
*/
import type * as Arr from "effect/Array"
import * as Deferred from "effect/Deferred"
import * as Duration from "effect/Duration"
import * as Effect from "effect/Effect"
import * as Exit from "effect/Exit"
import * as FileSystem from "effect/FileSystem"
Expand Down Expand Up @@ -66,6 +73,25 @@ const toPlatformError = (
type ExitCodeWithSignal = readonly [code: number | null, signal: NodeJS.Signals | null]
type ExitSignal = Deferred.Deferred<ExitCodeWithSignal>

const processGroupGraceMillis = 1_000
const processGroupPollIntervalMillis = 10

/** The leader has not exited, or a member of its process group still exists. */
const isProcessAlive = (childProcess: NodeChildProcess.ChildProcess, exitSignal: ExitSignal): boolean => {
if (!Deferred.isDoneUnsafe(exitSignal)) {
return true
}
if (globalThis.process.platform === "win32") {
return false
}
try {
globalThis.process.kill(-childProcess.pid!, 0)
return true
} catch {
return false
}
}

const taskkill = (
childProcess: NodeChildProcess.ChildProcess,
onExit: (error: NodeChildProcess.ExecException | null) => void = () => {}
Expand Down Expand Up @@ -413,26 +439,56 @@ const make = Effect.gen(function*() {
return Effect.void
})

const withTimeout = (
/** Waits for the leader and its process group, bounded by native time. */
const awaitProcessExit = (
childProcess: NodeChildProcess.ChildProcess,
exitSignal: ExitSignal,
timeoutMillis: number
): Effect.Effect<void> =>
Effect.callback<void>((resume) => {
const deadline = Date.now() + timeoutMillis
let timer: NodeJS.Timeout | undefined
const stop = () => {
clearTimeout(timer)
childProcess.removeListener("exit", poll)
}
const poll = () => {
clearTimeout(timer)
if (Date.now() >= deadline || !isProcessAlive(childProcess, exitSignal)) {
stop()
resume(Effect.void)
return
}
timer = setTimeout(poll, processGroupPollIntervalMillis)
}
// spawn's exit listener completes exitSignal before this listener runs
childProcess.on("exit", poll)
poll()
return Effect.sync(stop)
})

const terminateProcessGroup = Effect.fnUntraced(function*(
command: ChildProcess.StandardCommand,
childProcess: NodeChildProcess.ChildProcess,
exitSignal: ExitSignal,
options: ChildProcess.KillOptions | undefined
) =>
<A, E, R>(
kill: (
command: ChildProcess.StandardCommand,
childProcess: NodeChildProcess.ChildProcess,
signal: NodeJS.Signals
) => Effect.Effect<A, E, R>
) => {
const killSignal = options?.killSignal ?? "SIGTERM"
return Predicate.isUndefined(options?.forceKillAfter)
? kill(command, childProcess, killSignal)
: Effect.timeoutOrElse(kill(command, childProcess, killSignal), {
duration: options.forceKillAfter,
orElse: () => kill(command, childProcess, "SIGKILL")
})
}
) {
const signalGroup = (signal: NodeJS.Signals) =>
killProcessGroup(command, childProcess, signal).pipe(
Effect.catch(() => killProcess(command, childProcess, signal))
)
yield* signalGroup(options?.killSignal ?? "SIGTERM")
if (Predicate.isUndefined(options?.forceKillAfter)) {
yield* awaitProcessExit(childProcess, exitSignal, processGroupGraceMillis)
} else {
yield* awaitProcessExit(childProcess, exitSignal, Duration.toMillis(options.forceKillAfter))
if (isProcessAlive(childProcess, exitSignal)) {
yield* signalGroup("SIGKILL")
yield* awaitProcessExit(childProcess, exitSignal, processGroupGraceMillis)
}
}
yield* Deferred.await(exitSignal)
})

/**
* Get the appropriate source stream from a process handle based on the
Expand Down Expand Up @@ -485,28 +541,17 @@ const make = Effect.gen(function*() {
spawn(cmd, buildSpawnOptions(cmd.options, { cwd, env, stdio }, process.platform)),
Effect.fnUntraced(function*([childProcess, exitSignal]) {
const exited = yield* Deferred.isDone(exitSignal)
const killWithTimeout = withTimeout(childProcess, cmd, cmd.options)
if (exited) {
// Process already exited, check if children need cleanup
const [code] = yield* Deferred.await(exitSignal)
if (code !== 0 && Predicate.isNotNull(code)) {
// Non-zero exit code ,attempt to clean up process group
return yield* Effect.ignore(killWithTimeout(killProcessGroup))
yield* Effect.ignore(killProcessGroup(cmd, childProcess, cmd.options.killSignal ?? "SIGTERM"))
}
return yield* Effect.void
return
}
if (!isReferenced) {
return yield* Effect.void
return
}
// Process is still running, kill it
return yield* killWithTimeout((command, childProcess, signal) =>
killProcessGroup(command, childProcess, signal).pipe(
Effect.catch(() => killProcess(command, childProcess, signal)),
Effect.andThen(Deferred.await(exitSignal))
)
).pipe(
Effect.ignore
)
yield* Effect.ignore(terminateProcessGroup(cmd, childProcess, exitSignal, cmd.options))
})
)

Expand Down Expand Up @@ -543,17 +588,8 @@ const make = Effect.gen(function*() {
const error = new globalThis.Error(`Process interrupted due to receipt of signal: '${signal}'`)
return Effect.fail(toPlatformError("exitCode", error, cmd))
})
const kill = (options?: ChildProcess.KillOptions | undefined) => {
const killWithTimeout = withTimeout(childProcess, cmd, options)
return killWithTimeout((command, childProcess, signal) =>
killProcessGroup(command, childProcess, signal).pipe(
Effect.catch(() => killProcess(command, childProcess, signal)),
Effect.andThen(Deferred.await(exitSignal))
)
).pipe(
Effect.asVoid
)
}
const kill = (options?: ChildProcess.KillOptions | undefined) =>
terminateProcessGroup(cmd, childProcess, exitSignal, options)

return makeHandle({
pid,
Expand Down
133 changes: 133 additions & 0 deletions packages/platform/node-shared/test/NodeChildProcessSpawner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,15 @@ import * as NodeFileSystem from "@effect/platform-node-shared/NodeFileSystem"
import * as NodePath from "@effect/platform-node-shared/NodePath"
import { assert, describe, it } from "@effect/vitest"
import * as ChildProcessSpawnerTest from "effect-test/unstable/process/ChildProcessSpawnerTest"
import * as Deferred from "effect/Deferred"
import * as Effect from "effect/Effect"
import * as Exit from "effect/Exit"
import * as FileSystem from "effect/FileSystem"
import * as Layer from "effect/Layer"
import * as Scope from "effect/Scope"
import * as Stream from "effect/Stream"
import * as ChildProcess from "effect/unstable/process/ChildProcess"
import { join } from "node:path"

const NodeServices = NodeChildProcessSpawner.layer.pipe(
Layer.provideMerge(Layer.mergeAll(
Expand Down Expand Up @@ -90,3 +95,131 @@ it.live("kills every process in a pipeline", () =>
assert.strictEqual(rootFinalSize, rootSizeAfterKill)
assert.strictEqual(childFinalSize, childSizeAfterKill)
}).pipe(Effect.scoped, Effect.provide(NodeServices)))

const processGroupFixture = join(__dirname, "fixtures", "process-group.ts")

// Use native timers under TestClock.
const liveSleep = (millis: number) =>
Effect.callback<void>((resume) => {
const timer = setTimeout(() => resume(Effect.void), millis)
return Effect.sync(() => clearTimeout(timer))
})

const liveTimeout = (millis: number) => <A, E, R>(effect: Effect.Effect<A, E, R>) =>
Effect.raceFirst(effect, liveSleep(millis).pipe(Effect.andThen(Effect.die(new Error("timed out")))))

const startProcessGroup = (mode: "exit-on-signal" | "ignore-signal", options?: ChildProcess.CommandOptions) =>
Effect.gen(function*() {
const fs = yield* FileSystem.FileSystem
const directory = yield* fs.makeTempDirectoryScoped()
const marker = `${directory}/marker`
const scope = yield* Scope.make()
const handle = yield* Scope.provide(scope)(ChildProcess.make(
process.execPath,
[processGroupFixture, "leader", mode, marker],
{ stdin: "ignore", ...options }
))
const ready = yield* Deferred.make<number>()
yield* handle.stdout.pipe(
Stream.decodeText,
Stream.splitLines,
Stream.runForEach((line) =>
line.startsWith("READY ") ? Deferred.succeed(ready, Number(line.slice("READY ".length))) : Effect.void
),
Effect.forkScoped
)
const descendantPid = yield* Deferred.await(ready).pipe(liveTimeout(5_000))
return { handle, descendantPid, marker, scope }
})

const killDescendant = (pid: number) =>
Effect.sync(() => {
try {
process.kill(pid, "SIGKILL")
} catch {}
})

const assertHeartbeatStopped = (marker: string) =>
Effect.gen(function*() {
const fs = yield* FileSystem.FileSystem
const sizeAfterKill = (yield* fs.stat(marker)).size
yield* liveSleep(100)
const finalSize = (yield* fs.stat(marker)).size
assert.strictEqual(finalSize, sizeAfterKill)
})

const timed = <A, E, R>(effect: Effect.Effect<A, E, R>) =>
Effect.gen(function*() {
const start = Date.now()
yield* effect
return Date.now() - start
})

describe.skipIf(process.platform === "win32")("process group cleanup", () => {
it.live("scope release waits for descendants that outlive the leader", () =>
Effect.gen(function*() {
const fs = yield* FileSystem.FileSystem
const { handle, marker, scope } = yield* startProcessGroup("exit-on-signal")

yield* Scope.close(scope, Exit.void)

assert.isFalse(yield* handle.isRunning)
assert.isTrue(yield* fs.exists(marker))
}).pipe(Effect.scoped, Effect.provide(NodeServices)))

it.live("scope release force kills descendants that ignore the kill signal", () =>
Effect.gen(function*() {
const { marker, scope } = yield* startProcessGroup("ignore-signal", { forceKillAfter: "200 millis" })

yield* Scope.close(scope, Exit.void)

yield* assertHeartbeatStopped(marker)
}).pipe(Effect.scoped, Effect.provide(NodeServices)))

it.live("kill force kills descendants that ignore the kill signal", () =>
Effect.gen(function*() {
const { handle, marker, scope } = yield* startProcessGroup("ignore-signal")

yield* handle.kill({ forceKillAfter: "200 millis" })

assert.isFalse(yield* handle.isRunning)
yield* assertHeartbeatStopped(marker)
yield* Scope.close(scope, Exit.void)
}).pipe(Effect.scoped, Effect.provide(NodeServices)))

it.effect("forceKillAfter escalation does not depend on the Effect clock", () =>
Effect.gen(function*() {
const { marker, scope } = yield* startProcessGroup("ignore-signal", { forceKillAfter: "200 millis" })

const releaseMillis = yield* timed(Scope.close(scope, Exit.void))

assert.isBelow(releaseMillis, 2_000)
yield* assertHeartbeatStopped(marker)
}).pipe(Effect.scoped, Effect.provide(NodeServices)))

it.live("scope release returns when a descendant holds the inherited pipe without forceKillAfter", () =>
Effect.gen(function*() {
const { descendantPid, handle, scope } = yield* startProcessGroup("ignore-signal")

yield* Effect.gen(function*() {
const releaseMillis = yield* timed(Scope.close(scope, Exit.void))

assert.isFalse(yield* handle.isRunning)
assert.isAtLeast(releaseMillis, 1_000)
assert.isBelow(releaseMillis, 3_000)
}).pipe(Effect.ensuring(killDescendant(descendantPid)))
}).pipe(Effect.scoped, Effect.provide(NodeServices)))
})

it.live("scope release returns when stdout is unread and backpressured", () =>
Effect.gen(function*() {
const releaseMillis = yield* timed(Effect.scoped(Effect.gen(function*() {
yield* ChildProcess.make(process.execPath, [
"-e",
"process.stdout.write(\"x\".repeat(1024 * 1024)); setInterval(() => {}, 1000)"
], { stdin: "ignore" })
yield* Effect.sleep("100 millis")
})))

assert.isBelow(releaseMillis, 2_000)
}).pipe(Effect.provide(NodeServices)))
29 changes: 29 additions & 0 deletions packages/platform/node-shared/test/fixtures/process-group.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
import { spawn } from "node:child_process"
import { appendFileSync, writeFileSync } from "node:fs"

// The leader and descendant share a process group and inherited stdio. The
// leader keeps the default SIGTERM behavior. The descendant either exits 200ms
// after SIGTERM or ignores it while writing heartbeats to the marker file.
const [role, mode, marker] = process.argv.slice(2)

if (role === "leader") {
spawn(process.execPath, [process.argv[1], "descendant", mode, marker], {
stdio: ["ignore", "inherit", "inherit"]
})
setInterval(() => {}, 1_000)
} else {
if (mode === "exit-on-signal") {
process.on("SIGTERM", () => {
setTimeout(() => {
writeFileSync(marker, "exited")
process.exit(0)
}, 200)
})
} else {
process.on("SIGTERM", () => {})
appendFileSync(marker, "x")
setInterval(() => appendFileSync(marker, "x"), 10)
}
setTimeout(() => process.exit(1), 5_000)
process.stdout.write(`READY ${process.pid}\n`)
}
Loading