diff --git a/packages/runtime/src/scr_async.c b/packages/runtime/src/scr_async.c index 131fa1f4..dae4d630 100644 --- a/packages/runtime/src/scr_async.c +++ b/packages/runtime/src/scr_async.c @@ -53,6 +53,7 @@ * their empty implementation makes the absence explicit at link time * without weakening async/timer support. */ bool scr_children_pending(void) { return false; } +bool scr_children_ready(void) { return false; } bool scr_children_reffed_pending(void) { return false; } bool scr_children_failed_pending(void) { return false; } void scr_children_poll(void) {} @@ -2614,6 +2615,11 @@ bool scr_loop_run(ScrPromise *top_level) { /* Pending immediates are always-ready work: no sleep — run due timers * (Node's timers phase precedes check), then the check phase below. */ if (scr_pending_immediates > 0) due = now; + /* Child polling can drain a pipe into the IPC queue after its dispatch + * station has run. Completed sends and disconnects can also become ready + * during that turn. No fd will wake us for this userspace work: return to + * dispatch without sleeping, still allowing due timers and immediates. */ + if (scr_children_ready()) due = now; bool evw = scr_events_watching_fn != NULL && scr_events_watching_fn(); if (io) { if (kids && due > now + SCR_CHILD_POLL_MS) due = now + SCR_CHILD_POLL_MS; diff --git a/packages/runtime/src/scr_child.c b/packages/runtime/src/scr_child.c index 540e1fc9..0cd0ed8b 100644 --- a/packages/runtime/src/scr_child.c +++ b/packages/runtime/src/scr_child.c @@ -5660,6 +5660,20 @@ static bool scr_ipc_pending(void) { return false; } +bool scr_children_ready(void) { + for (ScrIpc *ipc = scr_ipcs; ipc != NULL; ipc = ipc->next) { + if (ipc->disconnect_pending || (ipc->n_pending > 0 && ipc->n_message > 0)) { + return true; + } + bool writer_pending = scr_child_writer_pending(ipc->writer); + if (ipc->n_send > 0 && (ipc->send_error != NULL || !writer_pending)) { + return true; + } + if (ipc->local_closing && !writer_pending) return true; + } + return false; +} + double scr_process_fork_target(double target_count) { if (scr_process_fork_id != -2) return scr_process_fork_id; double target = -1; diff --git a/packages/runtime/src/scr_runtime.h b/packages/runtime/src/scr_runtime.h index e4294b00..66deb40d 100644 --- a/packages/runtime/src/scr_runtime.h +++ b/packages/runtime/src/scr_runtime.h @@ -2806,6 +2806,9 @@ const char *scr_signal_name(int sig); * pending child: non-kqueue platforms, spawn failures awaiting their * first-pass settle, or a child whose exit filter could not be armed. */ bool scr_children_pending(void); +/* Work already queued in userspace: dispatch can progress without another + * pipe/exit notification. Connected channels and blocked writes are not ready. */ +bool scr_children_ready(void); bool scr_children_failed_pending(void); void scr_children_poll(void); bool scr_children_wait(double max_wait_ms); diff --git a/scripts/sandbox-test.mjs b/scripts/sandbox-test.mjs index 8964b488..adeb4df9 100644 --- a/scripts/sandbox-test.mjs +++ b/scripts/sandbox-test.mjs @@ -137,6 +137,7 @@ const hostLaneContractPattern = [ "udp-loopback-pair", "1564-fs-watch.ts", "1470-child-lifecycle.ts", + "2963-child-fork-dispatch/main.ts", "read-all: chunked writes with delays, then EOF", ].join("|"); const hostInvariantContractFiles = [ diff --git a/tests/corpus/2963-child-fork-dispatch/main.ts b/tests/corpus/2963-child-fork-dispatch/main.ts new file mode 100644 index 00000000..e9a9316a --- /dev/null +++ b/tests/corpus/2963-child-fork-dispatch/main.ts @@ -0,0 +1,37 @@ +import { fork } from "node:child_process"; + +const child = fork(new URL("./worker.ts", import.meta.url), [], { + stdio: ["ignore", "ignore", "inherit", "ipc"], +}); + +// Start the deadline after the child is ready so process startup is not part +// of the bound. Queued IPC must not wait for the idle poll timeout (1 second +// per handoff), or for this timer to wake the loop. No periodic wakeups. +let replies = 0; +let sent = 0; +let disconnected = false; +child.once("message", () => { + const deadline = setTimeout(() => { + console.log("IPC stalled"); + }, 500); + child.on("message", (message: { value: number }) => { + replies++; + if (message.value < 3) { + child.send({ value: message.value + 1 }, (error) => { + if (error) throw error; + sent++; + }); + } + }); + child.once("disconnect", () => { + disconnected = true; + clearTimeout(deadline); + }); + child.send({ value: 1 }, (error) => { + if (error) throw error; + sent++; + }); +}); +child.once("exit", (code) => { + console.log("replies", replies, "sent", sent, "disconnected", disconnected, "exit", code); +}); diff --git a/tests/corpus/2963-child-fork-dispatch/worker.ts b/tests/corpus/2963-child-fork-dispatch/worker.ts new file mode 100644 index 00000000..a65966ba --- /dev/null +++ b/tests/corpus/2963-child-fork-dispatch/worker.ts @@ -0,0 +1,7 @@ +process.on("message", (message: { value: number }) => { + process.send?.({ value: message.value }, (error) => { + if (error) throw error; + if (message.value === 3) process.disconnect(); + }); +}); +process.send?.({ value: 0 });