diff --git a/src/Parallel/ValueObject/ParallelProcess.php b/src/Parallel/ValueObject/ParallelProcess.php index d9a105071ff..131fd4f03fc 100644 --- a/src/Parallel/ValueObject/ParallelProcess.php +++ b/src/Parallel/ValueObject/ParallelProcess.php @@ -19,6 +19,7 @@ /** * Inspired at @see https://raw.githubusercontent.com/phpstan/phpstan-src/master/src/Parallel/Process.php + * @see \Rector\Tests\Parallel\ValueObject\ParallelProcessTest */ final class ParallelProcess { @@ -63,7 +64,12 @@ public function start(callable $onData, callable $onError, callable $onExit): vo } $this->stdErr = $tmp; - $this->process = new Process($this->command, null, null, [ + + // on Unix, the command runs in a wrapping shell; exec replaces the shell with the worker, + // so terminating the process stops the worker itself, not only the shell + $command = DIRECTORY_SEPARATOR === '\\' ? $this->command : 'exec ' . $this->command; + + $this->process = new Process($command, null, null, [ 2 => $this->stdErr, // todo is it fine to not have 0 and 1 FD? ]); @@ -98,6 +104,9 @@ public function request(array $data): void $this->cancelTimer(); $this->encoder->write($data); $this->timer = $this->loop->addTimer($this->timetoutInSeconds, function (): void { + // a worker that does not answer in time cannot be asked to stop either + $this->process->terminate(); + $onError = $this->onError; $errorMessage = sprintf('Child process timed out after %d seconds', $this->timetoutInSeconds); @@ -107,7 +116,6 @@ public function request(array $data): void public function quit(): void { - $this->cancelTimer(); if (! $this->process->isRunning()) { return; } @@ -117,10 +125,14 @@ public function quit(): void } // the process can be quit before its connection is bound, e.g. on quitAll() after an error; - // in that case the encoder was never set - if (isset($this->encoder)) { - $this->encoder->end(); + // such a worker cannot be asked to stop + if (! isset($this->encoder)) { + $this->process->terminate(); + return; } + + // a busy worker keeps its timeout, so it is still terminated when it never finishes its job + $this->encoder->end(); } public function bindConnection(Decoder $decoder, Encoder $encoder): void diff --git a/tests/Parallel/ValueObject/ParallelProcessTest.php b/tests/Parallel/ValueObject/ParallelProcessTest.php new file mode 100644 index 00000000000..f8c296c4857 --- /dev/null +++ b/tests/Parallel/ValueObject/ParallelProcessTest.php @@ -0,0 +1,162 @@ +streamSelectLoop = new StreamSelectLoop(); + $this->pidFile = (string) tempnam(sys_get_temp_dir(), 'rector_parallel_process_test'); + } + + protected function tearDown(): void + { + unlink($this->pidFile); + + if (is_resource($this->workerSocket)) { + fclose($this->workerSocket); + } + } + + public function testQuitTerminatesWorkerThatIsNotConnected(): void + { + $parallelProcess = $this->startWorker(); + + $parallelProcess->quit(); + + $this->assertWorkerIsTerminated(); + $this->assertSame([], $this->errorMessages); + } + + public function testTimeoutTerminatesWorker(): void + { + $parallelProcess = $this->startWorker(); + $this->bindConnection($parallelProcess); + + $parallelProcess->request([]); + + $this->assertWorkerIsTerminated(); + $this->assertSame(['Child process timed out after 1 seconds'], $this->errorMessages); + } + + public function testQuitKeepsTimeoutOfBusyWorker(): void + { + $parallelProcess = $this->startWorker(); + $this->bindConnection($parallelProcess); + $parallelProcess->request([]); + + $parallelProcess->quit(); + + $this->assertWorkerIsTerminated(); + $this->assertSame(['Child process timed out after 1 seconds'], $this->errorMessages); + } + + /** + * Starts a worker that neither connects nor answers, like one that is still booting or stuck in an endless loop + */ + private function startWorker(): ParallelProcess + { + $command = sprintf( + '%s -r %s -- %s', + escapeshellarg(PHP_BINARY), + escapeshellarg('file_put_contents($argv[1], getmypid()); sleep(10);'), + escapeshellarg($this->pidFile) + ); + + $parallelProcess = new ParallelProcess($command, $this->streamSelectLoop, 1); + $parallelProcess->start( + static function (): void { + }, + function (Throwable $throwable): void { + $this->errorMessages[] = $throwable->getMessage(); + }, + function (): void { + $this->hasExited = true; + $this->streamSelectLoop->stop(); + } + ); + + $this->runLoopUntil(fn (): bool => $this->readWorkerPid() > 0); + + return $parallelProcess; + } + + private function bindConnection(ParallelProcess $parallelProcess): void + { + $sockets = stream_socket_pair(STREAM_PF_UNIX, STREAM_SOCK_STREAM, STREAM_IPPROTO_IP); + $this->assertIsArray($sockets); + + // the other end stays open but silent, like a connected worker that never answers + [$this->workerSocket, $mainSocket] = $sockets; + + $duplexResourceStream = new DuplexResourceStream($mainSocket, $this->streamSelectLoop); + $parallelProcess->bindConnection(new Decoder($duplexResourceStream, true), new Encoder($duplexResourceStream)); + } + + private function assertWorkerIsTerminated(): void + { + $this->runLoopUntil(fn (): bool => $this->hasExited); + + $this->assertTrue($this->hasExited); + + // the worker itself must be gone, not only a shell wrapping it + $this->assertFalse(posix_kill($this->readWorkerPid(), 0)); + } + + /** + * @param callable(): bool $condition + */ + private function runLoopUntil(callable $condition): void + { + $periodicTimer = $this->streamSelectLoop->addPeriodicTimer(0.05, function () use ($condition): void { + if ($condition()) { + $this->streamSelectLoop->stop(); + } + }); + + $timeoutTimer = $this->streamSelectLoop->addTimer(5, function (): void { + $this->streamSelectLoop->stop(); + }); + + $this->streamSelectLoop->run(); + + $this->streamSelectLoop->cancelTimer($periodicTimer); + $this->streamSelectLoop->cancelTimer($timeoutTimer); + } + + private function readWorkerPid(): int + { + clearstatcache(); + + return (int) file_get_contents($this->pidFile); + } +}