Repository navigation
fix(task-runner-native): settle wait() for a process that already closed - #787
Merged
Merged
Conversation
wait() attached its 'close' listener only when called, so a process that had already closed left it waiting for an event that would never fire again: the promise never settled and the caller hung. Sequential code hides this, because it always waits before spawning the next task. A worker pool does not: it spawns several tasks before awaiting any of them, and a task that fails fast – a bad jar path exits in milliseconds – closes in between. Record what each process leaves behind when it closes, in a WeakMap so a task nobody waits on is forgotten with its process object, and let wait() settle from that record when the process is already gone. stop() keeps returning the output of a killed process: the recorded result only carries output on the failure path, which had already consumed it to build its error.
wait() removed every 'close' listener before attaching its own, which stranded anything else waiting on that event: - a stop() already in flight lost its cleanup listener, so it never resolved and its SIGKILL timer was never cleared; - a second, concurrent wait() deleted the first waiter's listener, so the first promise never settled. Both are the shapes a worker pool produces – several tasks in flight, and an abort that stops them while their waits are outstanding – so they belong with the race this branch already fixes. Remove exactly the listener run() installed, share one result record between concurrent waiters, and cache the output on it: reading consumes it, so waiting twice used to answer the second caller with an empty string.
The class kept a task's state in four places – two Maps keyed by pid, two WeakMaps keyed by the process – and read the output destructively, which is what forced the workarounds this branch accumulated: juggling close listeners, sharing a result record between concurrent waiters, and caching output on it so a second wait() would not answer with an empty string. Keep one record per task instead, created the moment the process exists: its output, and a promise that resolves when it closes. wait() and stop() both read that record and nothing else. The three hangs this branch fixes then stop being cases to handle. Awaiting a promise that already resolved is ordinary, so a closed process needs no lookaside; a promise settles every waiter, so concurrent waits need no shared record; and reading a field is not destructive, so waiting twice answers the same thing. No listener is removed, so a stop() in flight keeps its own. It also drops the two defects left for a follow-up. State is keyed by the process rather than by pid, which the OS recycles, and held weakly, so output is released with the process object instead of accumulating for the runner's lifetime. The synthetic 'error' event that run() emitted on a non-zero exit goes too: it was absorbed by the no-op listener installed beside it and observable to nobody. A spawn failure still arrives as a real 'error', still has a listener so Node does not throw it, and still reaches the caller through wait().
…xplained Follow-up on a review of the record-per-task rework: - stop() returned as soon as the process was gone, but 'exit' fires before 'close' and `killed` is true the moment a signal is *sent*, so it could answer with output that had not been read yet. It now awaits the close in every path, which is what the record already offers. - A process that could not be spawned at all – an unreadable cwd, say – closes with code -2 and no output, so wait() rejected with a message naming neither the cause nor the path. Keep the 'error' the spawn reports and put it in front of the output. - Decode the streams as UTF-8 rather than each chunk on its own, so a multi-byte character split across a chunk boundary survives. - Bound the retained output. A server task lives as long as the pipeline, and every log line it wrote was kept. Keep the last million characters: what a command reports about its work – a failure, or the metadata it prints when it is done – it prints at the end, which is also what sparql-qlever reads back out of it. The test helper now polls for the process to have closed rather than assuming one turn of the event loop is enough, so the tests for waiting after a close cannot silently pass by waiting before it. The already-exited stop() test asserts the output, not merely that something was returned. Drops the nx release lockfile drift that the first commit swept in.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
NativeTaskRunner.wait()could hang forever, in three different ways. All of them need a caller that does not wait for one task before starting the next – which is exactly what the worker pool in #782 (must 1) will do, and why this lands first.The hangs
wait()attached itscloselistener only when it was called, so a process that had already closed left it waiting for an event that would never fire again. The promise never settled: no error, no timeout. Sequential code hides this; a pool does not, and a task that fails fast – a bad jar path exits in milliseconds – closes in between.stop()already guarded for an exited process;wait()did not.stop()in flight, or a concurrentwait().wait()calledremoveAllListeners('close'), which strips every listener, not just the onerun()installed. Astop()in flight lost its cleanup listener, so it never resolved and its SIGKILL timer was never cleared; a second, concurrentwait()deleted the first waiter's listener, so the first promise never settled.wait()answered with an empty string.The fix
The class kept a task's state in four places – two
Maps keyed by pid, twoWeakMaps keyed by the process – and read the output destructively. That is what made each of the above a separate case to handle.It now keeps one record per task, created the moment the process exists: its output, and a promise that resolves when it closes.
wait()andstop()read that record and nothing else.The three hangs stop being cases to handle at all:
stop()in flight keeps its own.Net -79 lines in the class.
Two defects that were going to need a follow-up PR are gone with it: state is keyed by the process rather than by pid – which the OS recycles, so a long-lived runner could serve one task's output to another – and held weakly, so output is released with the process object instead of accumulating for the runner's lifetime (ADR 12).
One behaviour removed
run()emitted a synthetic'error'event on a non-zero exit, which the no-op listener installed two lines below absorbed – observable to nobody. A genuine spawn failure still arrives as a real'error', still has a listener so Node does not throw it, and still reaches the caller throughwait().Tests
Six new cases, four of which timed out rather than failed before the fix – that is what a hang looks like in a suite, and why this went unnoticed: a process that already succeeded; one that already failed; one already stopped; a
stop()in flight alongside await(); two concurrent waiters; and waiting twice. They were written against the previous, more elaborate fix and pass unchanged against this one.