From 2ce7760db273878a7b3bfc40cfa69cf0654609d7 Mon Sep 17 00:00:00 2001 From: Dmitriy Derepko Date: Tue, 29 Sep 2026 18:03:42 +0400 Subject: [PATCH 1/2] feat(lambda): run a Temporal worker as an AWS Lambda custom runtime Adds bin/lambda, which supervises RoadRunner for the duration of each invocation and stops it inside the shutdown buffer. The environment and each invocation are readonly DTOs, and unit tests cover the config validation, the Runtime API, the /proc parsing and the invocation loop. --- bin/lambda | 20 ++ composer.json | 3 + src/Lambda/Clock.php | 20 ++ src/Lambda/Config.php | 58 +++++ src/Lambda/Environment.php | 75 ++++++ .../Exception/ConfigurationException.php | 14 ++ src/Lambda/Http/Client.php | 81 +++++++ src/Lambda/Http/Response.php | 34 +++ src/Lambda/Http/TransportException.php | 14 ++ src/Lambda/Invocation.php | 24 ++ src/Lambda/Process/Handle.php | 89 +++++++ src/Lambda/Process/Signal.php | 19 ++ src/Lambda/Process/SignalTrap.php | 35 +++ src/Lambda/Process/Tree.php | 97 ++++++++ src/Lambda/RoadRunner/ConfigFile.php | 39 ++++ src/Lambda/RoadRunner/Process.php | 135 +++++++++++ src/Lambda/Runtime.php | 181 +++++++++++++++ src/Lambda/RuntimeApi.php | 97 ++++++++ tests/Unit/Lambda/ConfigTestCase.php | 85 +++++++ tests/Unit/Lambda/EnvironmentTestCase.php | 120 ++++++++++ tests/Unit/Lambda/Process/TreeTestCase.php | 36 +++ tests/Unit/Lambda/RecordingLogger.php | 18 ++ .../Lambda/RoadRunner/ConfigFileTestCase.php | 78 +++++++ tests/Unit/Lambda/RuntimeApiStub.php | 183 +++++++++++++++ tests/Unit/Lambda/RuntimeApiTestCase.php | 123 ++++++++++ tests/Unit/Lambda/RuntimeTestCase.php | 217 ++++++++++++++++++ 26 files changed, 1895 insertions(+) create mode 100755 bin/lambda create mode 100644 src/Lambda/Clock.php create mode 100644 src/Lambda/Config.php create mode 100644 src/Lambda/Environment.php create mode 100644 src/Lambda/Exception/ConfigurationException.php create mode 100644 src/Lambda/Http/Client.php create mode 100644 src/Lambda/Http/Response.php create mode 100644 src/Lambda/Http/TransportException.php create mode 100644 src/Lambda/Invocation.php create mode 100644 src/Lambda/Process/Handle.php create mode 100644 src/Lambda/Process/Signal.php create mode 100644 src/Lambda/Process/SignalTrap.php create mode 100644 src/Lambda/Process/Tree.php create mode 100644 src/Lambda/RoadRunner/ConfigFile.php create mode 100644 src/Lambda/RoadRunner/Process.php create mode 100644 src/Lambda/Runtime.php create mode 100644 src/Lambda/RuntimeApi.php create mode 100644 tests/Unit/Lambda/ConfigTestCase.php create mode 100644 tests/Unit/Lambda/EnvironmentTestCase.php create mode 100644 tests/Unit/Lambda/Process/TreeTestCase.php create mode 100644 tests/Unit/Lambda/RecordingLogger.php create mode 100644 tests/Unit/Lambda/RoadRunner/ConfigFileTestCase.php create mode 100644 tests/Unit/Lambda/RuntimeApiStub.php create mode 100644 tests/Unit/Lambda/RuntimeApiTestCase.php create mode 100644 tests/Unit/Lambda/RuntimeTestCase.php diff --git a/bin/lambda b/bin/lambda new file mode 100755 index 000000000..d79ec4ac3 --- /dev/null +++ b/bin/lambda @@ -0,0 +1,20 @@ +#!/usr/bin/env php +validate(); + } + + private function validate(): void + { + if ($this->gracefulTimeoutMs < 1_000) { + throw new ConfigurationException( + "TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS must be at least 1000, got {$this->gracefulTimeoutMs}", + ); + } + + $minimumBuffer = $this->minimumShutdownBufferMs(); + if ($this->shutdownBufferMs < $minimumBuffer) { + throw new ConfigurationException(\sprintf( + 'TEMPORAL_LAMBDA_SHUTDOWN_BUFFER_MS (%dms) is too small: it must reserve ' + . 'TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS (%dms) plus the SIGKILL escalation, the Runtime ' + . 'API response and the poll granularity, so at least %dms', + $this->shutdownBufferMs, + $this->gracefulTimeoutMs, + $minimumBuffer, + )); + } + } + + private function minimumShutdownBufferMs(): int + { + return $this->gracefulTimeoutMs + + Process::SIGKILL_SLACK_MS + + RuntimeApi::RESPONSE_RESERVE_MS + + 2 * \intdiv(Process::POLL_INTERVAL_US, 1000); + } +} diff --git a/src/Lambda/Environment.php b/src/Lambda/Environment.php new file mode 100644 index 000000000..cfd3721ae --- /dev/null +++ b/src/Lambda/Environment.php @@ -0,0 +1,75 @@ +handle($url, $timeoutMs); + \curl_setopt($handle, \CURLOPT_HEADERFUNCTION, static function (\CurlHandle $handle, string $line) use (&$headers): int { + $parts = \explode(':', $line, 2); + if (\count($parts) === 2) { + $headers[\strtolower(\trim($parts[0]))] = \trim($parts[1]); + } + + return \strlen($line); + }); + + return $this->send($handle, $headers); + } + + public function postJson(string $url, string $body, int $timeoutMs): Response + { + $handle = $this->handle($url, $timeoutMs); + \curl_setopt($handle, \CURLOPT_POST, true); + \curl_setopt($handle, \CURLOPT_POSTFIELDS, $body); + \curl_setopt($handle, \CURLOPT_HTTPHEADER, ['Content-Type: application/json']); + + $headers = []; + + return $this->send($handle, $headers); + } + + private function handle(string $url, int $timeoutMs): \CurlHandle + { + $handle = \curl_init($url); + if ($handle === false) { + throw new TransportException('Unable to initialise a curl handle'); + } + + \curl_setopt($handle, \CURLOPT_RETURNTRANSFER, true); + \curl_setopt($handle, \CURLOPT_FAILONERROR, false); + \curl_setopt($handle, \CURLOPT_CONNECTTIMEOUT, self::CONNECT_TIMEOUT_SECONDS); + \curl_setopt($handle, \CURLOPT_TIMEOUT_MS, $timeoutMs); + + return $handle; + } + + /** + * @param array $headers + */ + private function send(\CurlHandle $handle, array &$headers): Response + { + $body = \curl_exec($handle); + if ($body === false) { + throw new TransportException(\curl_error($handle)); + } + + return new Response( + status: (int) \curl_getinfo($handle, \CURLINFO_RESPONSE_CODE), + body: (string) $body, + headers: $headers, + ); + } +} diff --git a/src/Lambda/Http/Response.php b/src/Lambda/Http/Response.php new file mode 100644 index 000000000..23b13d143 --- /dev/null +++ b/src/Lambda/Http/Response.php @@ -0,0 +1,34 @@ + $headers Header names are lower-cased + */ + public function __construct( + public readonly int $status, + public readonly string $body, + public readonly array $headers = [], + ) {} + + public function isSuccessful(): bool + { + return $this->status >= 200 && $this->status <= 299; + } + + public function header(string $name): ?string + { + return $this->headers[\strtolower($name)] ?? null; + } +} diff --git a/src/Lambda/Http/TransportException.php b/src/Lambda/Http/TransportException.php new file mode 100644 index 000000000..f94ab0386 --- /dev/null +++ b/src/Lambda/Http/TransportException.php @@ -0,0 +1,14 @@ + $command + */ + public static function start(array $command, string $workingDirectory): self + { + $process = \proc_open( + $command, + [0 => ['pipe', 'r'], 1 => \STDERR, 2 => \STDERR], + $pipes, + $workingDirectory, + ); + + if (!\is_resource($process)) { + throw new \RuntimeException('proc_open() failed for: ' . \implode(' ', $command)); + } + + $handle = new self(); + $handle->process = $process; + $handle->pid = \proc_get_status($process)['pid']; + + \fclose($pipes[0]); + + return $handle; + } + + public function pid(): int + { + return $this->pid; + } + + public function exitCode(): ?int + { + if ($this->exitCode !== null) { + return $this->exitCode; + } + + if ($this->process === null) { + return null; + } + + $status = \proc_get_status($this->process); + if ($status['running']) { + return null; + } + + return $this->exitCode = $status['exitcode']; + } + + public function signal(int $signal): bool + { + if ($this->process === null) { + return false; + } + + return \proc_terminate($this->process, $signal); + } + + public function reap(): void + { + if ($this->process === null) { + return; + } + + $process = $this->process; + $this->process = null; + \proc_close($process); + } +} diff --git a/src/Lambda/Process/Signal.php b/src/Lambda/Process/Signal.php new file mode 100644 index 000000000..c24374462 --- /dev/null +++ b/src/Lambda/Process/Signal.php @@ -0,0 +1,19 @@ + $signals + * @param callable(int): void $handler + * @return bool Whether the signals could be trapped at all + */ + public static function trap(array $signals, callable $handler): bool + { + if (!\function_exists('pcntl_async_signals') || !\function_exists('pcntl_signal')) { + return false; + } + + \pcntl_async_signals(true); + + foreach ($signals as $signal) { + \pcntl_signal($signal, $handler); + } + + return true; + } +} diff --git a/src/Lambda/Process/Tree.php b/src/Lambda/Process/Tree.php new file mode 100644 index 000000000..05a14dd0b --- /dev/null +++ b/src/Lambda/Process/Tree.php @@ -0,0 +1,97 @@ + + */ + public static function childrenOf(int $pid): array + { + $children = []; + + $statFiles = \glob('/proc/[0-9]*/stat'); + + foreach ($statFiles === false ? [] : $statFiles as $statFile) { + $stat = @\file_get_contents($statFile); + if ($stat === false) { + continue; + } + + if (self::parentPid($stat) === $pid) { + $children[] = (int) \basename(\dirname($statFile)); + } + } + + return $children; + } + + public static function parentPid(string $stat): ?int + { + $closingParen = \strrpos($stat, ')'); + if ($closingParen === false) { + return null; + } + + $fields = \explode(' ', \substr($stat, $closingParen + 2)); + $parent = $fields[1] ?? ''; + + return \ctype_digit($parent) ? (int) $parent : null; + } + + /** + * @param list $pids + * @return int Number of processes the signal reached + */ + public static function killAll(array $pids): int + { + if (!\function_exists('posix_kill')) { + return 0; + } + + $killed = 0; + foreach ($pids as $pid) { + if (\posix_kill($pid, Signal::SIGKILL)) { + ++$killed; + } + } + + return $killed; + } + + public static function reapReparented(): void + { + if (!\function_exists('pcntl_waitpid')) { + return; + } + + $deadline = Clock::nowMs() + self::REAP_TIMEOUT_MS; + + while (Clock::nowMs() < $deadline) { + $reaped = \pcntl_waitpid(-1, $status, \WNOHANG); + + if ($reaped === -1) { + return; + } + + if ($reaped === 0) { + \usleep(self::REAP_POLL_INTERVAL_US); + } + } + } +} diff --git a/src/Lambda/RoadRunner/ConfigFile.php b/src/Lambda/RoadRunner/ConfigFile.php new file mode 100644 index 000000000..26eba7f7b --- /dev/null +++ b/src/Lambda/RoadRunner/ConfigFile.php @@ -0,0 +1,39 @@ +roadRunnerConfigTemplate); + if ($template === false) { + throw new ConfigurationException( + "RoadRunner config is not readable: {$config->roadRunnerConfigTemplate}", + ); + } + + $rendered = \rtrim($template, "\n") + . "\n\nendure:\n grace_period: {$config->gracefulTimeoutMs}ms\n"; + + if (\file_put_contents(self::PATH, $rendered) === false) { + throw new ConfigurationException('Unable to write ' . self::PATH); + } + + return self::PATH; + } +} diff --git a/src/Lambda/RoadRunner/Process.php b/src/Lambda/RoadRunner/Process.php new file mode 100644 index 000000000..a8c350c66 --- /dev/null +++ b/src/Lambda/RoadRunner/Process.php @@ -0,0 +1,135 @@ +handle = Handle::start([ + $this->config->roadRunnerBinary, + 'serve', + '-w', $this->config->taskRoot, + '-c', ConfigFile::PATH, + ], $this->config->taskRoot); + } + + public function exitCode(): ?int + { + return $this->handle?->exitCode(); + } + + public function stop(): void + { + if ($this->handle === null) { + return; + } + + $alreadyExited = $this->handle->exitCode(); + if ($alreadyExited !== null) { + $this->release(); + $this->logger->error( + "roadrunner had already exited with code {$alreadyExited} when the shutdown started", + ['requestId' => $this->requestId], + ); + + return; + } + + $startedAt = Clock::nowMs(); + $this->stopDeadlineMs ??= $startedAt + $this->config->gracefulTimeoutMs + self::SIGKILL_SLACK_MS; + $this->handle->signal(Signal::SIGTERM); + + while (Clock::nowMs() < $this->stopDeadlineMs) { + $exitCode = $this->handle->exitCode(); + if ($exitCode !== null) { + $this->release(); + $this->logExit($exitCode, Clock::nowMs() - $startedAt); + + return; + } + + \usleep(self::POLL_INTERVAL_US); + } + + $this->kill($startedAt); + } + + private function kill(int $startedAt): void + { + $handle = $this->handle; + if ($handle === null) { + return; + } + + $orphans = Tree::childrenOf($handle->pid()); + + if (!$handle->signal(Signal::SIGKILL)) { + $this->logger->error('SIGKILL to roadrunner failed', ['requestId' => $this->requestId]); + } + + $this->release(); + + $killed = Tree::killAll($orphans); + Tree::reapReparented(); + + $this->logger->error(\sprintf( + 'roadrunner did not stop within %dms, SIGKILLed it and %d worker process(es); ' + . 'in-flight tasks were dropped and will be retried by Temporal', + Clock::nowMs() - $startedAt, + $killed, + ), ['requestId' => $this->requestId]); + } + + private function logExit(int $exitCode, int $elapsedMs): void + { + if ($exitCode === 0) { + $this->logger->info("roadrunner drained and exited in {$elapsedMs}ms", ['requestId' => $this->requestId]); + + return; + } + + $this->logger->error(\sprintf( + 'roadrunner exited with code %d after %dms: the drain did not finish within ' + . 'TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS (%dms); in-flight tasks were dropped and ' + . 'will be retried by Temporal', + $exitCode, + $elapsedMs, + $this->config->gracefulTimeoutMs, + ), ['requestId' => $this->requestId]); + } + + private function release(): void + { + $this->handle?->reap(); + $this->handle = null; + } +} diff --git a/src/Lambda/Runtime.php b/src/Lambda/Runtime.php new file mode 100644 index 000000000..61caec116 --- /dev/null +++ b/src/Lambda/Runtime.php @@ -0,0 +1,181 @@ +runtimeApi), + $logger ?? new StderrLogger(), + ); + } + + public static function main(): never + { + $logger = new StderrLogger(); + + try { + self::create(logger: $logger)->run(); + } catch (ConfigurationException $e) { + $logger->error('configuration error: ' . $e->getMessage()); + self::reportInitError($e); + + exit(1); + } catch (\Throwable $e) { + $logger->error('runtime loop crashed: ' . $e->getMessage()); + + exit(1); + } + + exit(0); + } + + public function run(): void + { + ConfigFile::render($this->config); + $this->trapSignals(); + \register_shutdown_function(fn() => $this->roadRunner?->stop()); + + $this->log('info', \sprintf( + 'runtime ready: buffer=%dms graceful=%dms config=%s', + $this->config->shutdownBufferMs, + $this->config->gracefulTimeoutMs, + $this->config->roadRunnerConfigTemplate, + )); + + while (true) { + $invocation = $this->api->nextInvocation(); + $this->requestId = $invocation->requestId; + + $error = $this->invoke($invocation->deadlineMs); + if ($error !== null) { + $this->log('error', 'invocation failed: ' . $error->getMessage()); + } + + $this->acknowledge($invocation->requestId, $error); + $this->requestId = '-'; + } + } + + private static function reportInitError(\Throwable $error): void + { + $runtimeApi = Environment::runtimeApi(); + if ($runtimeApi === null) { + return; + } + + (new RuntimeApi($runtimeApi))->reportInitError($error); + } + + private function acknowledge(string $requestId, ?\Throwable $error): void + { + try { + if ($error === null) { + $this->api->respond($requestId); + } else { + $this->api->reportInvocationError($requestId, $error); + } + } catch (\Throwable $e) { + $this->log('error', 'acknowledgement failed: ' . $e->getMessage()); + } + } + + private function invoke(int $deadlineMs): ?\Throwable + { + $startedAt = Clock::nowMs(); + $stopAtMs = $deadlineMs - $this->config->shutdownBufferMs; + $this->log('info', \sprintf('invocation started, remaining=%dms', $deadlineMs - $startedAt)); + + try { + $runMs = $stopAtMs - Clock::nowMs(); + if ($runMs < self::MINIMUM_RUN_MS) { + throw new \RuntimeException(\sprintf( + 'Insufficient invocation time: %dms left to poll after reserving a %dms shutdown buffer', + $runMs, + $this->config->shutdownBufferMs, + )); + } + + $this->roadRunner = new Process($this->config, $this->logger, $this->requestId); + $this->roadRunner->start(); + $this->log('info', "roadrunner started, polling for {$runMs}ms"); + + $this->pollUntil($stopAtMs); + $this->log('info', 'shutdown window reached, stopping roadrunner'); + + return null; + } catch (\Throwable $e) { + return $e; + } finally { + $this->roadRunner?->stop(); + $this->roadRunner = null; + $this->log('info', \sprintf( + 'invocation finished in %dms, remaining=%dms', + Clock::nowMs() - $startedAt, + $deadlineMs - Clock::nowMs(), + )); + } + } + + private function log(string $level, string $message): void + { + $this->logger->log($level, $message, ['requestId' => $this->requestId]); + } + + private function pollUntil(int $stopAtMs): void + { + while (Clock::nowMs() < $stopAtMs) { + $exitCode = $this->roadRunner?->exitCode(); + if ($exitCode !== null) { + throw new \RuntimeException("RoadRunner exited before the shutdown window with code {$exitCode}"); + } + + \usleep(Process::POLL_INTERVAL_US); + } + } + + private function trapSignals(): void + { + SignalTrap::trap([Signal::SIGTERM, Signal::SIGINT], function (int $received): void { + $this->log('error', "received signal {$received}, stopping roadrunner"); + + exit(128 + $received); + }); + } +} diff --git a/src/Lambda/RuntimeApi.php b/src/Lambda/RuntimeApi.php new file mode 100644 index 000000000..e9d6fba79 --- /dev/null +++ b/src/Lambda/RuntimeApi.php @@ -0,0 +1,97 @@ +client = $client ?? new Client(); + } + + public function nextInvocation(): Invocation + { + $response = $this->client->get($this->url('runtime/invocation/next'), Client::NO_TIMEOUT); + $this->assertSuccessful($response, 'fetch the next invocation'); + + $requestId = $response->header('Lambda-Runtime-Aws-Request-Id') ?? ''; + if ($requestId === '') { + throw new \RuntimeException('Runtime API response is missing Lambda-Runtime-Aws-Request-Id'); + } + + $deadlineMs = (int) ($response->header('Lambda-Runtime-Deadline-Ms') ?? 0); + if ($deadlineMs <= 0) { + throw new \RuntimeException('Runtime API response is missing Lambda-Runtime-Deadline-Ms'); + } + + return new Invocation($requestId, $deadlineMs); + } + + public function respond(string $requestId): void + { + $this->post("runtime/invocation/{$requestId}/response", 'null', 'send the invocation response'); + } + + public function reportInvocationError(string $requestId, \Throwable $error): void + { + $this->post( + "runtime/invocation/{$requestId}/error", + $this->encodeError($error), + 'report the invocation error', + ); + } + + public function reportInitError(\Throwable $error): void + { + $this->post('runtime/init/error', $this->encodeError($error), 'report the init error'); + } + + private function post(string $path, string $body, string $what): void + { + $response = $this->client->postJson($this->url($path), $body, self::RESPONSE_RESERVE_MS); + + $this->assertSuccessful($response, $what); + } + + private function assertSuccessful(Response $response, string $what): void + { + if (!$response->isSuccessful()) { + throw new \RuntimeException( + \rtrim("Failed to {$what}: HTTP {$response->status} {$response->body}"), + ); + } + } + + private function url(string $path): string + { + return "http://{$this->host}/" . self::VERSION . "/{$path}"; + } + + private function encodeError(\Throwable $error): string + { + return (string) \json_encode([ + 'errorType' => \str_replace('\\', '.', $error::class), + 'errorMessage' => $error->getMessage(), + 'stackTrace' => \explode("\n", $error->getTraceAsString()), + ]); + } +} diff --git a/tests/Unit/Lambda/ConfigTestCase.php b/tests/Unit/Lambda/ConfigTestCase.php new file mode 100644 index 000000000..bd6f290ee --- /dev/null +++ b/tests/Unit/Lambda/ConfigTestCase.php @@ -0,0 +1,85 @@ +expectException(ConfigurationException::class); + $this->expectExceptionMessage(\sprintf( + 'TEMPORAL_LAMBDA_SHUTDOWN_BUFFER_MS (%dms) is too small: it must reserve ' + . 'TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS (%dms) plus the SIGKILL escalation, the Runtime ' + . 'API response and the poll granularity, so at least %dms', + $minimum - 1, + self::MINIMUM_GRACEFUL_MS, + $minimum, + )); + + self::config(shutdownBufferMs: $minimum - 1, gracefulTimeoutMs: self::MINIMUM_GRACEFUL_MS); + } + + public function testBufferExactlyAtTheMinimumIsAccepted(): void + { + $minimum = self::minimumBuffer(self::MINIMUM_GRACEFUL_MS); + + $config = self::config(shutdownBufferMs: $minimum, gracefulTimeoutMs: self::MINIMUM_GRACEFUL_MS); + + self::assertSame($minimum, $config->shutdownBufferMs); + } + + public function testABiggerGracefulTimeoutRaisesTheMinimumBuffer(): void + { + $accepted = self::minimumBuffer(5_000); + + $this->expectException(ConfigurationException::class); + + self::config(shutdownBufferMs: $accepted - 1, gracefulTimeoutMs: 5_000); + } + + public function testGracefulTimeoutBelowOneSecondIsRejected(): void + { + $this->expectException(ConfigurationException::class); + $this->expectExceptionMessage( + 'TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS must be at least 1000, got 999', + ); + + self::config(shutdownBufferMs: 30_000, gracefulTimeoutMs: 999); + } + + private static function minimumBuffer(int $gracefulTimeoutMs): int + { + return $gracefulTimeoutMs + + Process::SIGKILL_SLACK_MS + + RuntimeApi::RESPONSE_RESERVE_MS + + 2 * \intdiv(Process::POLL_INTERVAL_US, 1000); + } + + private static function config(int $shutdownBufferMs, int $gracefulTimeoutMs): Config + { + return new Config( + runtimeApi: '127.0.0.1:9001', + taskRoot: '/var/task', + roadRunnerBinary: 'rr', + roadRunnerConfigTemplate: '/var/task/.rr.yaml', + shutdownBufferMs: $shutdownBufferMs, + gracefulTimeoutMs: $gracefulTimeoutMs, + ); + } +} diff --git a/tests/Unit/Lambda/EnvironmentTestCase.php b/tests/Unit/Lambda/EnvironmentTestCase.php new file mode 100644 index 000000000..90071e2bf --- /dev/null +++ b/tests/Unit/Lambda/EnvironmentTestCase.php @@ -0,0 +1,120 @@ +runtimeApi); + self::assertSame('/opt/task', $config->taskRoot); + self::assertSame('/opt/rr', $config->roadRunnerBinary); + self::assertSame('/opt/task/custom.yaml', $config->roadRunnerConfigTemplate); + self::assertSame(9_000, $config->shutdownBufferMs); + self::assertSame(6_000, $config->gracefulTimeoutMs); + } + + public function testDefaultsApplyWhenOnlyTheRuntimeApiIsSet(): void + { + \putenv('AWS_LAMBDA_RUNTIME_API=127.0.0.1:9001'); + + $config = Environment::capture(); + + self::assertSame('/var/task', $config->taskRoot); + self::assertSame('rr', $config->roadRunnerBinary); + self::assertSame('/var/task/.rr.yaml', $config->roadRunnerConfigTemplate); + self::assertSame(7_000, $config->shutdownBufferMs); + self::assertSame(5_000, $config->gracefulTimeoutMs); + } + + public function testTheConfigTemplateFollowsTheTaskRoot(): void + { + \putenv('AWS_LAMBDA_RUNTIME_API=127.0.0.1:9001'); + \putenv('LAMBDA_TASK_ROOT=/opt/task'); + + self::assertSame('/opt/task/.rr.yaml', Environment::capture()->roadRunnerConfigTemplate); + } + + public function testAMissingRuntimeApiIsRejected(): void + { + $this->expectException(ConfigurationException::class); + $this->expectExceptionMessage( + 'AWS_LAMBDA_RUNTIME_API is not set: this script must run as a Lambda runtime', + ); + + Environment::capture(); + } + + public function testAnEmptyVariableIsTreatedAsUnset(): void + { + \putenv('AWS_LAMBDA_RUNTIME_API='); + + $this->expectException(ConfigurationException::class); + + Environment::capture(); + } + + public function testANonNumericTimeoutIsRejected(): void + { + \putenv('AWS_LAMBDA_RUNTIME_API=127.0.0.1:9001'); + \putenv('TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS=5s'); + + $this->expectException(ConfigurationException::class); + $this->expectExceptionMessage( + 'TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS must be a positive integer, got "5s"', + ); + + Environment::capture(); + } + + protected function setUp(): void + { + parent::setUp(); + + $this->clearEnvironment(); + } + + protected function tearDown(): void + { + $this->clearEnvironment(); + + parent::tearDown(); + } + + private function clearEnvironment(): void + { + foreach (self::KEYS as $key) { + \putenv($key); + } + } +} diff --git a/tests/Unit/Lambda/Process/TreeTestCase.php b/tests/Unit/Lambda/Process/TreeTestCase.php new file mode 100644 index 000000000..33247a504 --- /dev/null +++ b/tests/Unit/Lambda/Process/TreeTestCase.php @@ -0,0 +1,36 @@ + */ + public array $messages = []; + + public function log($level, \Stringable|string $message, array $context = []): void + { + $this->messages[] = (string) $message; + } +} diff --git a/tests/Unit/Lambda/RoadRunner/ConfigFileTestCase.php b/tests/Unit/Lambda/RoadRunner/ConfigFileTestCase.php new file mode 100644 index 000000000..332a48ea9 --- /dev/null +++ b/tests/Unit/Lambda/RoadRunner/ConfigFileTestCase.php @@ -0,0 +1,78 @@ +template, "rpc:\n listen: tcp://127.0.0.1:6001\n"); + + $path = ConfigFile::render($this->config(gracefulTimeoutMs: 4_500)); + + self::assertSame( + "rpc:\n listen: tcp://127.0.0.1:6001\n\nendure:\n grace_period: 4500ms\n", + \file_get_contents($path), + ); + } + + public function testTrailingBlankLinesDoNotProduceAStrayGap(): void + { + \file_put_contents($this->template, "rpc:\n listen: tcp://127.0.0.1:6001\n\n\n\n"); + + $path = ConfigFile::render($this->config(gracefulTimeoutMs: 5_000)); + + self::assertSame( + "rpc:\n listen: tcp://127.0.0.1:6001\n\nendure:\n grace_period: 5000ms\n", + \file_get_contents($path), + ); + } + + public function testRenderOverwritesAPreviousRender(): void + { + \file_put_contents($this->template, "rpc:\n listen: tcp://127.0.0.1:6001\n"); + + ConfigFile::render($this->config(gracefulTimeoutMs: 1_000)); + $path = ConfigFile::render($this->config(gracefulTimeoutMs: 2_000)); + + self::assertStringNotContainsString('1000ms', (string) \file_get_contents($path)); + self::assertStringContainsString('2000ms', (string) \file_get_contents($path)); + } + + protected function setUp(): void + { + parent::setUp(); + + $this->template = \tempnam(\sys_get_temp_dir(), 'rr-template-'); + } + + protected function tearDown(): void + { + @\unlink($this->template); + @\unlink(ConfigFile::PATH); + + parent::tearDown(); + } + + private function config(int $gracefulTimeoutMs): Config + { + return new Config( + runtimeApi: '127.0.0.1:9001', + taskRoot: \sys_get_temp_dir(), + roadRunnerBinary: 'rr', + roadRunnerConfigTemplate: $this->template, + shutdownBufferMs: 30_000, + gracefulTimeoutMs: $gracefulTimeoutMs, + ); + } +} diff --git a/tests/Unit/Lambda/RuntimeApiStub.php b/tests/Unit/Lambda/RuntimeApiStub.php new file mode 100644 index 000000000..0f0ee081b --- /dev/null +++ b/tests/Unit/Lambda/RuntimeApiStub.php @@ -0,0 +1,183 @@ +invocations = $invocations; + $this->writeSettings(); + } + + public function setPollWindowMs(int $pollWindowMs): void + { + $this->pollWindowMs = $pollWindowMs; + $this->writeSettings(); + } + + public function omitHeaders(): void + { + $this->sendHeaders = false; + $this->writeSettings(); + } + + public function failAcknowledgementsWith(int $status): void + { + $this->acknowledgementStatus = $status; + $this->writeSettings(); + } + + public function host(): string + { + return "127.0.0.1:{$this->port}"; + } + + /** + * @return list + */ + public function acknowledgements(): array + { + return \array_column($this->acknowledgementLog(), 'path'); + } + + public function lastBody(): string + { + $log = $this->acknowledgementLog(); + + return $log === [] ? '' : (string) \end($log)['body']; + } + + public function start(): void + { + $this->port = self::freePort(); + $this->writeSettings(); + \file_put_contents($this->router(), self::ROUTER); + + $this->process = new Process( + [\PHP_BINARY, '-S', $this->host(), $this->router()], + env: ['STUB_STATE' => $this->stateDir], + timeout: null, + ); + $this->process->start(); + + $deadline = \microtime(true) + self::START_TIMEOUT_SECONDS; + while (\microtime(true) < $deadline) { + $connection = @\fsockopen('127.0.0.1', $this->port, $code, $message, 0.2); + if ($connection !== false) { + \fclose($connection); + + return; + } + + \usleep(50_000); + } + + throw new \RuntimeException('The Runtime API stub did not start'); + } + + public function stop(): void + { + $this->process?->stop(timeout: 5); + $this->process = null; + } + + private static function freePort(): int + { + $socket = \stream_socket_server('tcp://127.0.0.1:0', $code, $message); + if ($socket === false) { + throw new \RuntimeException("Unable to reserve a port: {$message}"); + } + + $name = (string) \stream_socket_get_name($socket, false); + \fclose($socket); + + return (int) \substr($name, \strrpos($name, ':') + 1); + } + + /** + * @return list + */ + private function acknowledgementLog(): array + { + $log = @\file_get_contents($this->stateDir . '/acknowledgements.json'); + if ($log === false) { + return []; + } + + return \array_map( + static fn(string $line): array => (array) \json_decode($line, true), + \array_filter(\explode("\n", \trim($log))), + ); + } + + private function writeSettings(): void + { + \file_put_contents($this->stateDir . '/settings.json', (string) \json_encode([ + 'invocations' => $this->invocations, + 'pollWindowMs' => $this->pollWindowMs, + 'sendHeaders' => $this->sendHeaders, + 'acknowledgementStatus' => $this->acknowledgementStatus, + ])); + } + + private function router(): string + { + return $this->stateDir . '/runtime-api.php'; + } + + private const ROUTER = <<<'PHP' + = (int) $settings['invocations']) { + \http_response_code(500); + + return; + } + + if ($settings['sendHeaders'] === true) { + \header('Lambda-Runtime-Aws-Request-Id: request-' . $served); + \header('Lambda-Runtime-Deadline-Ms: ' . ( + (int) (\microtime(true) * 1000) + (int) $settings['pollWindowMs'] + )); + } + + return; + } + + \file_put_contents( + $state . '/acknowledgements.json', + \json_encode(['path' => $path, 'body' => (string) \file_get_contents('php://input')]) . "\n", + \FILE_APPEND, + ); + + \http_response_code((int) $settings['acknowledgementStatus']); + PHP; +} diff --git a/tests/Unit/Lambda/RuntimeApiTestCase.php b/tests/Unit/Lambda/RuntimeApiTestCase.php new file mode 100644 index 000000000..9a77df207 --- /dev/null +++ b/tests/Unit/Lambda/RuntimeApiTestCase.php @@ -0,0 +1,123 @@ +api->nextInvocation(); + + self::assertSame('request-0', $invocation->requestId); + self::assertGreaterThanOrEqual($before, $invocation->deadlineMs); + } + + public function testMissingRequestIdHeaderIsRejected(): void + { + $this->stub->omitHeaders(); + + $this->expectExceptionMessage('Runtime API response is missing Lambda-Runtime-Aws-Request-Id'); + + $this->api->nextInvocation(); + } + + public function testAnExhaustedRuntimeApiReportsTheHttpStatus(): void + { + $this->api->nextInvocation(); + + $this->expectExceptionMessage('Failed to fetch the next invocation: HTTP 500'); + + $this->api->nextInvocation(); + } + + public function testResponseIsPostedForTheGivenRequestId(): void + { + $this->api->respond('request-7'); + + self::assertSame( + ['/2018-06-01/runtime/invocation/request-7/response'], + $this->stub->acknowledgements(), + ); + } + + public function testInitErrorCarriesTheExceptionType(): void + { + $this->api->reportInitError(new \LogicException('template is unreadable')); + + self::assertSame(['/2018-06-01/runtime/init/error'], $this->stub->acknowledgements()); + self::assertSame( + [ + 'errorType' => 'LogicException', + 'errorMessage' => 'template is unreadable', + ], + \array_intersect_key( + (array) \json_decode($this->stub->lastBody(), true), + ['errorType' => null, 'errorMessage' => null], + ), + ); + } + + public function testANamespacedExceptionTypeUsesDotsAsSeparators(): void + { + $this->api->reportInitError(new \Temporal\Lambda\Exception\ConfigurationException('nope')); + + self::assertSame( + 'Temporal.Lambda.Exception.ConfigurationException', + ((array) \json_decode($this->stub->lastBody(), true))['errorType'], + ); + } + + public function testAFailedAcknowledgementReportsTheHttpStatus(): void + { + $this->stub->failAcknowledgementsWith(503); + + $this->expectExceptionMessage('Failed to send the invocation response: HTTP 503'); + + $this->api->respond('request-0'); + } + + protected function setUp(): void + { + parent::setUp(); + + $this->stateDir = \sys_get_temp_dir() . '/temporal-lambda-api-' . \bin2hex(\random_bytes(6)); + \mkdir($this->stateDir); + + $this->stub = new RuntimeApiStub($this->stateDir, 5_000); + $this->stub->start(); + $this->api = new RuntimeApi($this->stub->host()); + } + + protected function tearDown(): void + { + $this->stub->stop(); + + foreach ((array) \glob($this->stateDir . '/{,.}[!.,..]*', \GLOB_BRACE) as $file) { + @\unlink((string) $file); + } + + \rmdir($this->stateDir); + + parent::tearDown(); + } +} diff --git a/tests/Unit/Lambda/RuntimeTestCase.php b/tests/Unit/Lambda/RuntimeTestCase.php new file mode 100644 index 000000000..f1bf8adc5 --- /dev/null +++ b/tests/Unit/Lambda/RuntimeTestCase.php @@ -0,0 +1,217 @@ +fakeRoadRunner('\\sleep(30);'); + + $this->runUntilTheStubStops($binary); + + self::assertSame( + ['/2018-06-01/runtime/invocation/request-0/response'], + $this->api->acknowledgements(), + ); + } + + public function testRoadRunnerIsStoppedBeforeTheDeadline(): void + { + $binary = $this->fakeRoadRunner($this->recordStart('runs') . ' \\sleep(30);'); + + $elapsed = $this->runUntilTheStubStops($binary); + + self::assertSame('x', \file_get_contents($this->workDir . '/runs')); + self::assertLessThan(self::POLL_WINDOW_MS + self::BUFFER_MS, $elapsed); + self::assertSame([], $this->survivingProcesses($binary)); + } + + public function testRoadRunnerCrashIsReportedAsAnInvocationError(): void + { + $binary = $this->fakeRoadRunner('exit(3);'); + + $this->runUntilTheStubStops($binary); + + self::assertSame( + ['/2018-06-01/runtime/invocation/request-0/error'], + $this->api->acknowledgements(), + ); + self::assertStringContainsString( + 'RoadRunner exited before the shutdown window with code 3', + $this->api->lastBody(), + ); + self::assertContains( + 'roadrunner had already exited with code 3 when the shutdown started', + $this->logger->messages, + ); + } + + public function testInvocationWithoutEnoughTimeIsReportedWithoutStartingRoadRunner(): void + { + $binary = $this->fakeRoadRunner($this->recordStart('started') . ' \\sleep(30);'); + $this->api->setPollWindowMs(0); + + $this->runUntilTheStubStops($binary); + + self::assertSame( + ['/2018-06-01/runtime/invocation/request-0/error'], + $this->api->acknowledgements(), + ); + self::assertStringContainsString('Insufficient invocation time', $this->api->lastBody()); + self::assertFileDoesNotExist($this->workDir . '/started'); + } + + public function testARoadRunnerThatIgnoresSigtermIsKilled(): void + { + $binary = $this->fakeRoadRunner( + '\\pcntl_async_signals(true); \\pcntl_signal(\\SIGTERM, static fn() => null); while (true) { \\sleep(1); }', + ); + + $elapsed = $this->runUntilTheStubStops($binary); + + self::assertSame([], $this->survivingProcesses($binary)); + self::assertLessThan(self::POLL_WINDOW_MS + self::BUFFER_MS, $elapsed); + self::assertStringContainsString( + 'SIGKILLed it', + \implode("\n", $this->logger->messages), + ); + } + + public function testEveryInvocationGetsItsOwnRoadRunner(): void + { + $binary = $this->fakeRoadRunner($this->recordStart('runs') . ' \\sleep(30);'); + $this->api->setInvocations(2); + + $this->runUntilTheStubStops($binary); + + self::assertSame( + [ + '/2018-06-01/runtime/invocation/request-0/response', + '/2018-06-01/runtime/invocation/request-1/response', + ], + $this->api->acknowledgements(), + ); + self::assertSame('xx', \file_get_contents($this->workDir . '/runs')); + } + + protected function setUp(): void + { + parent::setUp(); + + $this->workDir = \sys_get_temp_dir() . '/temporal-lambda-' . \bin2hex(\random_bytes(6)); + \mkdir($this->workDir); + \file_put_contents($this->workDir . '/.rr.yaml', "rpc:\n listen: tcp://127.0.0.1:6001\n"); + + $this->logger = new RecordingLogger(); + $this->api = new RuntimeApiStub($this->workDir, self::POLL_WINDOW_MS + self::BUFFER_MS); + $this->api->start(); + } + + protected function tearDown(): void + { + $this->api->stop(); + + foreach ((array) \glob($this->workDir . '/{,.}[!.,..]*', \GLOB_BRACE) as $file) { + @\unlink((string) $file); + } + + \rmdir($this->workDir); + @\unlink(ConfigFile::PATH); + + parent::tearDown(); + } + + private function runUntilTheStubStops(string $binary): int + { + $config = new Config( + runtimeApi: $this->api->host(), + taskRoot: $this->workDir, + roadRunnerBinary: $binary, + roadRunnerConfigTemplate: $this->workDir . '/.rr.yaml', + shutdownBufferMs: self::BUFFER_MS, + gracefulTimeoutMs: self::GRACEFUL_MS, + ); + + $runtime = new Runtime($config, new RuntimeApi($config->runtimeApi), $this->logger); + + $startedAt = Clock::nowMs(); + try { + $runtime->run(); + self::fail('The runtime loop returned instead of failing on the exhausted stub'); + } catch (\RuntimeException $e) { + self::assertStringContainsString('fetch the next invocation', $e->getMessage()); + } + + return Clock::nowMs() - $startedAt; + } + + /** + * @return list + */ + private function survivingProcesses(string $binary): array + { + $output = []; + \exec('ps -eo pid=,command=', $output); + + $survivors = []; + foreach ($output as $line) { + if (\str_contains($line, $binary)) { + $survivors[] = (int) \trim($line); + } + } + + return $survivors; + } + + private function recordStart(string $file): string + { + return \sprintf( + '\\file_put_contents(%s, "x", \\FILE_APPEND);', + \var_export($this->workDir . '/' . $file, true), + ); + } + + private function fakeRoadRunner(string $body): string + { + $path = $this->workDir . '/rr'; + \file_put_contents($path, "#!/usr/bin/env php\n Date: Fri, 2 Oct 2026 18:42:59 +0400 Subject: [PATCH 2/2] refactor(lambda): drop the PHP runtime in favour of the RoadRunner plugin The invocation loop now lives in roadrunner-temporal: it serves the Lambda Runtime API and cycles only the Temporal workers, so the PHP pools survive every invocation and the SDK needs no Lambda-specific code. --- bin/lambda | 20 -- composer.json | 3 - src/Lambda/Clock.php | 20 -- src/Lambda/Config.php | 58 ----- src/Lambda/Environment.php | 75 ------ .../Exception/ConfigurationException.php | 14 -- src/Lambda/Http/Client.php | 81 ------- src/Lambda/Http/Response.php | 34 --- src/Lambda/Http/TransportException.php | 14 -- src/Lambda/Invocation.php | 24 -- src/Lambda/Process/Handle.php | 89 ------- src/Lambda/Process/Signal.php | 19 -- src/Lambda/Process/SignalTrap.php | 35 --- src/Lambda/Process/Tree.php | 97 -------- src/Lambda/RoadRunner/ConfigFile.php | 39 ---- src/Lambda/RoadRunner/Process.php | 135 ----------- src/Lambda/Runtime.php | 181 --------------- src/Lambda/RuntimeApi.php | 97 -------- tests/Unit/Lambda/ConfigTestCase.php | 85 ------- tests/Unit/Lambda/EnvironmentTestCase.php | 120 ---------- tests/Unit/Lambda/Process/TreeTestCase.php | 36 --- tests/Unit/Lambda/RecordingLogger.php | 18 -- .../Lambda/RoadRunner/ConfigFileTestCase.php | 78 ------- tests/Unit/Lambda/RuntimeApiStub.php | 183 --------------- tests/Unit/Lambda/RuntimeApiTestCase.php | 123 ---------- tests/Unit/Lambda/RuntimeTestCase.php | 217 ------------------ 26 files changed, 1895 deletions(-) delete mode 100755 bin/lambda delete mode 100644 src/Lambda/Clock.php delete mode 100644 src/Lambda/Config.php delete mode 100644 src/Lambda/Environment.php delete mode 100644 src/Lambda/Exception/ConfigurationException.php delete mode 100644 src/Lambda/Http/Client.php delete mode 100644 src/Lambda/Http/Response.php delete mode 100644 src/Lambda/Http/TransportException.php delete mode 100644 src/Lambda/Invocation.php delete mode 100644 src/Lambda/Process/Handle.php delete mode 100644 src/Lambda/Process/Signal.php delete mode 100644 src/Lambda/Process/SignalTrap.php delete mode 100644 src/Lambda/Process/Tree.php delete mode 100644 src/Lambda/RoadRunner/ConfigFile.php delete mode 100644 src/Lambda/RoadRunner/Process.php delete mode 100644 src/Lambda/Runtime.php delete mode 100644 src/Lambda/RuntimeApi.php delete mode 100644 tests/Unit/Lambda/ConfigTestCase.php delete mode 100644 tests/Unit/Lambda/EnvironmentTestCase.php delete mode 100644 tests/Unit/Lambda/Process/TreeTestCase.php delete mode 100644 tests/Unit/Lambda/RecordingLogger.php delete mode 100644 tests/Unit/Lambda/RoadRunner/ConfigFileTestCase.php delete mode 100644 tests/Unit/Lambda/RuntimeApiStub.php delete mode 100644 tests/Unit/Lambda/RuntimeApiTestCase.php delete mode 100644 tests/Unit/Lambda/RuntimeTestCase.php diff --git a/bin/lambda b/bin/lambda deleted file mode 100755 index d79ec4ac3..000000000 --- a/bin/lambda +++ /dev/null @@ -1,20 +0,0 @@ -#!/usr/bin/env php -validate(); - } - - private function validate(): void - { - if ($this->gracefulTimeoutMs < 1_000) { - throw new ConfigurationException( - "TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS must be at least 1000, got {$this->gracefulTimeoutMs}", - ); - } - - $minimumBuffer = $this->minimumShutdownBufferMs(); - if ($this->shutdownBufferMs < $minimumBuffer) { - throw new ConfigurationException(\sprintf( - 'TEMPORAL_LAMBDA_SHUTDOWN_BUFFER_MS (%dms) is too small: it must reserve ' - . 'TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS (%dms) plus the SIGKILL escalation, the Runtime ' - . 'API response and the poll granularity, so at least %dms', - $this->shutdownBufferMs, - $this->gracefulTimeoutMs, - $minimumBuffer, - )); - } - } - - private function minimumShutdownBufferMs(): int - { - return $this->gracefulTimeoutMs - + Process::SIGKILL_SLACK_MS - + RuntimeApi::RESPONSE_RESERVE_MS - + 2 * \intdiv(Process::POLL_INTERVAL_US, 1000); - } -} diff --git a/src/Lambda/Environment.php b/src/Lambda/Environment.php deleted file mode 100644 index cfd3721ae..000000000 --- a/src/Lambda/Environment.php +++ /dev/null @@ -1,75 +0,0 @@ -handle($url, $timeoutMs); - \curl_setopt($handle, \CURLOPT_HEADERFUNCTION, static function (\CurlHandle $handle, string $line) use (&$headers): int { - $parts = \explode(':', $line, 2); - if (\count($parts) === 2) { - $headers[\strtolower(\trim($parts[0]))] = \trim($parts[1]); - } - - return \strlen($line); - }); - - return $this->send($handle, $headers); - } - - public function postJson(string $url, string $body, int $timeoutMs): Response - { - $handle = $this->handle($url, $timeoutMs); - \curl_setopt($handle, \CURLOPT_POST, true); - \curl_setopt($handle, \CURLOPT_POSTFIELDS, $body); - \curl_setopt($handle, \CURLOPT_HTTPHEADER, ['Content-Type: application/json']); - - $headers = []; - - return $this->send($handle, $headers); - } - - private function handle(string $url, int $timeoutMs): \CurlHandle - { - $handle = \curl_init($url); - if ($handle === false) { - throw new TransportException('Unable to initialise a curl handle'); - } - - \curl_setopt($handle, \CURLOPT_RETURNTRANSFER, true); - \curl_setopt($handle, \CURLOPT_FAILONERROR, false); - \curl_setopt($handle, \CURLOPT_CONNECTTIMEOUT, self::CONNECT_TIMEOUT_SECONDS); - \curl_setopt($handle, \CURLOPT_TIMEOUT_MS, $timeoutMs); - - return $handle; - } - - /** - * @param array $headers - */ - private function send(\CurlHandle $handle, array &$headers): Response - { - $body = \curl_exec($handle); - if ($body === false) { - throw new TransportException(\curl_error($handle)); - } - - return new Response( - status: (int) \curl_getinfo($handle, \CURLINFO_RESPONSE_CODE), - body: (string) $body, - headers: $headers, - ); - } -} diff --git a/src/Lambda/Http/Response.php b/src/Lambda/Http/Response.php deleted file mode 100644 index 23b13d143..000000000 --- a/src/Lambda/Http/Response.php +++ /dev/null @@ -1,34 +0,0 @@ - $headers Header names are lower-cased - */ - public function __construct( - public readonly int $status, - public readonly string $body, - public readonly array $headers = [], - ) {} - - public function isSuccessful(): bool - { - return $this->status >= 200 && $this->status <= 299; - } - - public function header(string $name): ?string - { - return $this->headers[\strtolower($name)] ?? null; - } -} diff --git a/src/Lambda/Http/TransportException.php b/src/Lambda/Http/TransportException.php deleted file mode 100644 index f94ab0386..000000000 --- a/src/Lambda/Http/TransportException.php +++ /dev/null @@ -1,14 +0,0 @@ - $command - */ - public static function start(array $command, string $workingDirectory): self - { - $process = \proc_open( - $command, - [0 => ['pipe', 'r'], 1 => \STDERR, 2 => \STDERR], - $pipes, - $workingDirectory, - ); - - if (!\is_resource($process)) { - throw new \RuntimeException('proc_open() failed for: ' . \implode(' ', $command)); - } - - $handle = new self(); - $handle->process = $process; - $handle->pid = \proc_get_status($process)['pid']; - - \fclose($pipes[0]); - - return $handle; - } - - public function pid(): int - { - return $this->pid; - } - - public function exitCode(): ?int - { - if ($this->exitCode !== null) { - return $this->exitCode; - } - - if ($this->process === null) { - return null; - } - - $status = \proc_get_status($this->process); - if ($status['running']) { - return null; - } - - return $this->exitCode = $status['exitcode']; - } - - public function signal(int $signal): bool - { - if ($this->process === null) { - return false; - } - - return \proc_terminate($this->process, $signal); - } - - public function reap(): void - { - if ($this->process === null) { - return; - } - - $process = $this->process; - $this->process = null; - \proc_close($process); - } -} diff --git a/src/Lambda/Process/Signal.php b/src/Lambda/Process/Signal.php deleted file mode 100644 index c24374462..000000000 --- a/src/Lambda/Process/Signal.php +++ /dev/null @@ -1,19 +0,0 @@ - $signals - * @param callable(int): void $handler - * @return bool Whether the signals could be trapped at all - */ - public static function trap(array $signals, callable $handler): bool - { - if (!\function_exists('pcntl_async_signals') || !\function_exists('pcntl_signal')) { - return false; - } - - \pcntl_async_signals(true); - - foreach ($signals as $signal) { - \pcntl_signal($signal, $handler); - } - - return true; - } -} diff --git a/src/Lambda/Process/Tree.php b/src/Lambda/Process/Tree.php deleted file mode 100644 index 05a14dd0b..000000000 --- a/src/Lambda/Process/Tree.php +++ /dev/null @@ -1,97 +0,0 @@ - - */ - public static function childrenOf(int $pid): array - { - $children = []; - - $statFiles = \glob('/proc/[0-9]*/stat'); - - foreach ($statFiles === false ? [] : $statFiles as $statFile) { - $stat = @\file_get_contents($statFile); - if ($stat === false) { - continue; - } - - if (self::parentPid($stat) === $pid) { - $children[] = (int) \basename(\dirname($statFile)); - } - } - - return $children; - } - - public static function parentPid(string $stat): ?int - { - $closingParen = \strrpos($stat, ')'); - if ($closingParen === false) { - return null; - } - - $fields = \explode(' ', \substr($stat, $closingParen + 2)); - $parent = $fields[1] ?? ''; - - return \ctype_digit($parent) ? (int) $parent : null; - } - - /** - * @param list $pids - * @return int Number of processes the signal reached - */ - public static function killAll(array $pids): int - { - if (!\function_exists('posix_kill')) { - return 0; - } - - $killed = 0; - foreach ($pids as $pid) { - if (\posix_kill($pid, Signal::SIGKILL)) { - ++$killed; - } - } - - return $killed; - } - - public static function reapReparented(): void - { - if (!\function_exists('pcntl_waitpid')) { - return; - } - - $deadline = Clock::nowMs() + self::REAP_TIMEOUT_MS; - - while (Clock::nowMs() < $deadline) { - $reaped = \pcntl_waitpid(-1, $status, \WNOHANG); - - if ($reaped === -1) { - return; - } - - if ($reaped === 0) { - \usleep(self::REAP_POLL_INTERVAL_US); - } - } - } -} diff --git a/src/Lambda/RoadRunner/ConfigFile.php b/src/Lambda/RoadRunner/ConfigFile.php deleted file mode 100644 index 26eba7f7b..000000000 --- a/src/Lambda/RoadRunner/ConfigFile.php +++ /dev/null @@ -1,39 +0,0 @@ -roadRunnerConfigTemplate); - if ($template === false) { - throw new ConfigurationException( - "RoadRunner config is not readable: {$config->roadRunnerConfigTemplate}", - ); - } - - $rendered = \rtrim($template, "\n") - . "\n\nendure:\n grace_period: {$config->gracefulTimeoutMs}ms\n"; - - if (\file_put_contents(self::PATH, $rendered) === false) { - throw new ConfigurationException('Unable to write ' . self::PATH); - } - - return self::PATH; - } -} diff --git a/src/Lambda/RoadRunner/Process.php b/src/Lambda/RoadRunner/Process.php deleted file mode 100644 index a8c350c66..000000000 --- a/src/Lambda/RoadRunner/Process.php +++ /dev/null @@ -1,135 +0,0 @@ -handle = Handle::start([ - $this->config->roadRunnerBinary, - 'serve', - '-w', $this->config->taskRoot, - '-c', ConfigFile::PATH, - ], $this->config->taskRoot); - } - - public function exitCode(): ?int - { - return $this->handle?->exitCode(); - } - - public function stop(): void - { - if ($this->handle === null) { - return; - } - - $alreadyExited = $this->handle->exitCode(); - if ($alreadyExited !== null) { - $this->release(); - $this->logger->error( - "roadrunner had already exited with code {$alreadyExited} when the shutdown started", - ['requestId' => $this->requestId], - ); - - return; - } - - $startedAt = Clock::nowMs(); - $this->stopDeadlineMs ??= $startedAt + $this->config->gracefulTimeoutMs + self::SIGKILL_SLACK_MS; - $this->handle->signal(Signal::SIGTERM); - - while (Clock::nowMs() < $this->stopDeadlineMs) { - $exitCode = $this->handle->exitCode(); - if ($exitCode !== null) { - $this->release(); - $this->logExit($exitCode, Clock::nowMs() - $startedAt); - - return; - } - - \usleep(self::POLL_INTERVAL_US); - } - - $this->kill($startedAt); - } - - private function kill(int $startedAt): void - { - $handle = $this->handle; - if ($handle === null) { - return; - } - - $orphans = Tree::childrenOf($handle->pid()); - - if (!$handle->signal(Signal::SIGKILL)) { - $this->logger->error('SIGKILL to roadrunner failed', ['requestId' => $this->requestId]); - } - - $this->release(); - - $killed = Tree::killAll($orphans); - Tree::reapReparented(); - - $this->logger->error(\sprintf( - 'roadrunner did not stop within %dms, SIGKILLed it and %d worker process(es); ' - . 'in-flight tasks were dropped and will be retried by Temporal', - Clock::nowMs() - $startedAt, - $killed, - ), ['requestId' => $this->requestId]); - } - - private function logExit(int $exitCode, int $elapsedMs): void - { - if ($exitCode === 0) { - $this->logger->info("roadrunner drained and exited in {$elapsedMs}ms", ['requestId' => $this->requestId]); - - return; - } - - $this->logger->error(\sprintf( - 'roadrunner exited with code %d after %dms: the drain did not finish within ' - . 'TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS (%dms); in-flight tasks were dropped and ' - . 'will be retried by Temporal', - $exitCode, - $elapsedMs, - $this->config->gracefulTimeoutMs, - ), ['requestId' => $this->requestId]); - } - - private function release(): void - { - $this->handle?->reap(); - $this->handle = null; - } -} diff --git a/src/Lambda/Runtime.php b/src/Lambda/Runtime.php deleted file mode 100644 index 61caec116..000000000 --- a/src/Lambda/Runtime.php +++ /dev/null @@ -1,181 +0,0 @@ -runtimeApi), - $logger ?? new StderrLogger(), - ); - } - - public static function main(): never - { - $logger = new StderrLogger(); - - try { - self::create(logger: $logger)->run(); - } catch (ConfigurationException $e) { - $logger->error('configuration error: ' . $e->getMessage()); - self::reportInitError($e); - - exit(1); - } catch (\Throwable $e) { - $logger->error('runtime loop crashed: ' . $e->getMessage()); - - exit(1); - } - - exit(0); - } - - public function run(): void - { - ConfigFile::render($this->config); - $this->trapSignals(); - \register_shutdown_function(fn() => $this->roadRunner?->stop()); - - $this->log('info', \sprintf( - 'runtime ready: buffer=%dms graceful=%dms config=%s', - $this->config->shutdownBufferMs, - $this->config->gracefulTimeoutMs, - $this->config->roadRunnerConfigTemplate, - )); - - while (true) { - $invocation = $this->api->nextInvocation(); - $this->requestId = $invocation->requestId; - - $error = $this->invoke($invocation->deadlineMs); - if ($error !== null) { - $this->log('error', 'invocation failed: ' . $error->getMessage()); - } - - $this->acknowledge($invocation->requestId, $error); - $this->requestId = '-'; - } - } - - private static function reportInitError(\Throwable $error): void - { - $runtimeApi = Environment::runtimeApi(); - if ($runtimeApi === null) { - return; - } - - (new RuntimeApi($runtimeApi))->reportInitError($error); - } - - private function acknowledge(string $requestId, ?\Throwable $error): void - { - try { - if ($error === null) { - $this->api->respond($requestId); - } else { - $this->api->reportInvocationError($requestId, $error); - } - } catch (\Throwable $e) { - $this->log('error', 'acknowledgement failed: ' . $e->getMessage()); - } - } - - private function invoke(int $deadlineMs): ?\Throwable - { - $startedAt = Clock::nowMs(); - $stopAtMs = $deadlineMs - $this->config->shutdownBufferMs; - $this->log('info', \sprintf('invocation started, remaining=%dms', $deadlineMs - $startedAt)); - - try { - $runMs = $stopAtMs - Clock::nowMs(); - if ($runMs < self::MINIMUM_RUN_MS) { - throw new \RuntimeException(\sprintf( - 'Insufficient invocation time: %dms left to poll after reserving a %dms shutdown buffer', - $runMs, - $this->config->shutdownBufferMs, - )); - } - - $this->roadRunner = new Process($this->config, $this->logger, $this->requestId); - $this->roadRunner->start(); - $this->log('info', "roadrunner started, polling for {$runMs}ms"); - - $this->pollUntil($stopAtMs); - $this->log('info', 'shutdown window reached, stopping roadrunner'); - - return null; - } catch (\Throwable $e) { - return $e; - } finally { - $this->roadRunner?->stop(); - $this->roadRunner = null; - $this->log('info', \sprintf( - 'invocation finished in %dms, remaining=%dms', - Clock::nowMs() - $startedAt, - $deadlineMs - Clock::nowMs(), - )); - } - } - - private function log(string $level, string $message): void - { - $this->logger->log($level, $message, ['requestId' => $this->requestId]); - } - - private function pollUntil(int $stopAtMs): void - { - while (Clock::nowMs() < $stopAtMs) { - $exitCode = $this->roadRunner?->exitCode(); - if ($exitCode !== null) { - throw new \RuntimeException("RoadRunner exited before the shutdown window with code {$exitCode}"); - } - - \usleep(Process::POLL_INTERVAL_US); - } - } - - private function trapSignals(): void - { - SignalTrap::trap([Signal::SIGTERM, Signal::SIGINT], function (int $received): void { - $this->log('error', "received signal {$received}, stopping roadrunner"); - - exit(128 + $received); - }); - } -} diff --git a/src/Lambda/RuntimeApi.php b/src/Lambda/RuntimeApi.php deleted file mode 100644 index e9d6fba79..000000000 --- a/src/Lambda/RuntimeApi.php +++ /dev/null @@ -1,97 +0,0 @@ -client = $client ?? new Client(); - } - - public function nextInvocation(): Invocation - { - $response = $this->client->get($this->url('runtime/invocation/next'), Client::NO_TIMEOUT); - $this->assertSuccessful($response, 'fetch the next invocation'); - - $requestId = $response->header('Lambda-Runtime-Aws-Request-Id') ?? ''; - if ($requestId === '') { - throw new \RuntimeException('Runtime API response is missing Lambda-Runtime-Aws-Request-Id'); - } - - $deadlineMs = (int) ($response->header('Lambda-Runtime-Deadline-Ms') ?? 0); - if ($deadlineMs <= 0) { - throw new \RuntimeException('Runtime API response is missing Lambda-Runtime-Deadline-Ms'); - } - - return new Invocation($requestId, $deadlineMs); - } - - public function respond(string $requestId): void - { - $this->post("runtime/invocation/{$requestId}/response", 'null', 'send the invocation response'); - } - - public function reportInvocationError(string $requestId, \Throwable $error): void - { - $this->post( - "runtime/invocation/{$requestId}/error", - $this->encodeError($error), - 'report the invocation error', - ); - } - - public function reportInitError(\Throwable $error): void - { - $this->post('runtime/init/error', $this->encodeError($error), 'report the init error'); - } - - private function post(string $path, string $body, string $what): void - { - $response = $this->client->postJson($this->url($path), $body, self::RESPONSE_RESERVE_MS); - - $this->assertSuccessful($response, $what); - } - - private function assertSuccessful(Response $response, string $what): void - { - if (!$response->isSuccessful()) { - throw new \RuntimeException( - \rtrim("Failed to {$what}: HTTP {$response->status} {$response->body}"), - ); - } - } - - private function url(string $path): string - { - return "http://{$this->host}/" . self::VERSION . "/{$path}"; - } - - private function encodeError(\Throwable $error): string - { - return (string) \json_encode([ - 'errorType' => \str_replace('\\', '.', $error::class), - 'errorMessage' => $error->getMessage(), - 'stackTrace' => \explode("\n", $error->getTraceAsString()), - ]); - } -} diff --git a/tests/Unit/Lambda/ConfigTestCase.php b/tests/Unit/Lambda/ConfigTestCase.php deleted file mode 100644 index bd6f290ee..000000000 --- a/tests/Unit/Lambda/ConfigTestCase.php +++ /dev/null @@ -1,85 +0,0 @@ -expectException(ConfigurationException::class); - $this->expectExceptionMessage(\sprintf( - 'TEMPORAL_LAMBDA_SHUTDOWN_BUFFER_MS (%dms) is too small: it must reserve ' - . 'TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS (%dms) plus the SIGKILL escalation, the Runtime ' - . 'API response and the poll granularity, so at least %dms', - $minimum - 1, - self::MINIMUM_GRACEFUL_MS, - $minimum, - )); - - self::config(shutdownBufferMs: $minimum - 1, gracefulTimeoutMs: self::MINIMUM_GRACEFUL_MS); - } - - public function testBufferExactlyAtTheMinimumIsAccepted(): void - { - $minimum = self::minimumBuffer(self::MINIMUM_GRACEFUL_MS); - - $config = self::config(shutdownBufferMs: $minimum, gracefulTimeoutMs: self::MINIMUM_GRACEFUL_MS); - - self::assertSame($minimum, $config->shutdownBufferMs); - } - - public function testABiggerGracefulTimeoutRaisesTheMinimumBuffer(): void - { - $accepted = self::minimumBuffer(5_000); - - $this->expectException(ConfigurationException::class); - - self::config(shutdownBufferMs: $accepted - 1, gracefulTimeoutMs: 5_000); - } - - public function testGracefulTimeoutBelowOneSecondIsRejected(): void - { - $this->expectException(ConfigurationException::class); - $this->expectExceptionMessage( - 'TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS must be at least 1000, got 999', - ); - - self::config(shutdownBufferMs: 30_000, gracefulTimeoutMs: 999); - } - - private static function minimumBuffer(int $gracefulTimeoutMs): int - { - return $gracefulTimeoutMs - + Process::SIGKILL_SLACK_MS - + RuntimeApi::RESPONSE_RESERVE_MS - + 2 * \intdiv(Process::POLL_INTERVAL_US, 1000); - } - - private static function config(int $shutdownBufferMs, int $gracefulTimeoutMs): Config - { - return new Config( - runtimeApi: '127.0.0.1:9001', - taskRoot: '/var/task', - roadRunnerBinary: 'rr', - roadRunnerConfigTemplate: '/var/task/.rr.yaml', - shutdownBufferMs: $shutdownBufferMs, - gracefulTimeoutMs: $gracefulTimeoutMs, - ); - } -} diff --git a/tests/Unit/Lambda/EnvironmentTestCase.php b/tests/Unit/Lambda/EnvironmentTestCase.php deleted file mode 100644 index 90071e2bf..000000000 --- a/tests/Unit/Lambda/EnvironmentTestCase.php +++ /dev/null @@ -1,120 +0,0 @@ -runtimeApi); - self::assertSame('/opt/task', $config->taskRoot); - self::assertSame('/opt/rr', $config->roadRunnerBinary); - self::assertSame('/opt/task/custom.yaml', $config->roadRunnerConfigTemplate); - self::assertSame(9_000, $config->shutdownBufferMs); - self::assertSame(6_000, $config->gracefulTimeoutMs); - } - - public function testDefaultsApplyWhenOnlyTheRuntimeApiIsSet(): void - { - \putenv('AWS_LAMBDA_RUNTIME_API=127.0.0.1:9001'); - - $config = Environment::capture(); - - self::assertSame('/var/task', $config->taskRoot); - self::assertSame('rr', $config->roadRunnerBinary); - self::assertSame('/var/task/.rr.yaml', $config->roadRunnerConfigTemplate); - self::assertSame(7_000, $config->shutdownBufferMs); - self::assertSame(5_000, $config->gracefulTimeoutMs); - } - - public function testTheConfigTemplateFollowsTheTaskRoot(): void - { - \putenv('AWS_LAMBDA_RUNTIME_API=127.0.0.1:9001'); - \putenv('LAMBDA_TASK_ROOT=/opt/task'); - - self::assertSame('/opt/task/.rr.yaml', Environment::capture()->roadRunnerConfigTemplate); - } - - public function testAMissingRuntimeApiIsRejected(): void - { - $this->expectException(ConfigurationException::class); - $this->expectExceptionMessage( - 'AWS_LAMBDA_RUNTIME_API is not set: this script must run as a Lambda runtime', - ); - - Environment::capture(); - } - - public function testAnEmptyVariableIsTreatedAsUnset(): void - { - \putenv('AWS_LAMBDA_RUNTIME_API='); - - $this->expectException(ConfigurationException::class); - - Environment::capture(); - } - - public function testANonNumericTimeoutIsRejected(): void - { - \putenv('AWS_LAMBDA_RUNTIME_API=127.0.0.1:9001'); - \putenv('TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS=5s'); - - $this->expectException(ConfigurationException::class); - $this->expectExceptionMessage( - 'TEMPORAL_LAMBDA_GRACEFUL_TIMEOUT_MS must be a positive integer, got "5s"', - ); - - Environment::capture(); - } - - protected function setUp(): void - { - parent::setUp(); - - $this->clearEnvironment(); - } - - protected function tearDown(): void - { - $this->clearEnvironment(); - - parent::tearDown(); - } - - private function clearEnvironment(): void - { - foreach (self::KEYS as $key) { - \putenv($key); - } - } -} diff --git a/tests/Unit/Lambda/Process/TreeTestCase.php b/tests/Unit/Lambda/Process/TreeTestCase.php deleted file mode 100644 index 33247a504..000000000 --- a/tests/Unit/Lambda/Process/TreeTestCase.php +++ /dev/null @@ -1,36 +0,0 @@ - */ - public array $messages = []; - - public function log($level, \Stringable|string $message, array $context = []): void - { - $this->messages[] = (string) $message; - } -} diff --git a/tests/Unit/Lambda/RoadRunner/ConfigFileTestCase.php b/tests/Unit/Lambda/RoadRunner/ConfigFileTestCase.php deleted file mode 100644 index 332a48ea9..000000000 --- a/tests/Unit/Lambda/RoadRunner/ConfigFileTestCase.php +++ /dev/null @@ -1,78 +0,0 @@ -template, "rpc:\n listen: tcp://127.0.0.1:6001\n"); - - $path = ConfigFile::render($this->config(gracefulTimeoutMs: 4_500)); - - self::assertSame( - "rpc:\n listen: tcp://127.0.0.1:6001\n\nendure:\n grace_period: 4500ms\n", - \file_get_contents($path), - ); - } - - public function testTrailingBlankLinesDoNotProduceAStrayGap(): void - { - \file_put_contents($this->template, "rpc:\n listen: tcp://127.0.0.1:6001\n\n\n\n"); - - $path = ConfigFile::render($this->config(gracefulTimeoutMs: 5_000)); - - self::assertSame( - "rpc:\n listen: tcp://127.0.0.1:6001\n\nendure:\n grace_period: 5000ms\n", - \file_get_contents($path), - ); - } - - public function testRenderOverwritesAPreviousRender(): void - { - \file_put_contents($this->template, "rpc:\n listen: tcp://127.0.0.1:6001\n"); - - ConfigFile::render($this->config(gracefulTimeoutMs: 1_000)); - $path = ConfigFile::render($this->config(gracefulTimeoutMs: 2_000)); - - self::assertStringNotContainsString('1000ms', (string) \file_get_contents($path)); - self::assertStringContainsString('2000ms', (string) \file_get_contents($path)); - } - - protected function setUp(): void - { - parent::setUp(); - - $this->template = \tempnam(\sys_get_temp_dir(), 'rr-template-'); - } - - protected function tearDown(): void - { - @\unlink($this->template); - @\unlink(ConfigFile::PATH); - - parent::tearDown(); - } - - private function config(int $gracefulTimeoutMs): Config - { - return new Config( - runtimeApi: '127.0.0.1:9001', - taskRoot: \sys_get_temp_dir(), - roadRunnerBinary: 'rr', - roadRunnerConfigTemplate: $this->template, - shutdownBufferMs: 30_000, - gracefulTimeoutMs: $gracefulTimeoutMs, - ); - } -} diff --git a/tests/Unit/Lambda/RuntimeApiStub.php b/tests/Unit/Lambda/RuntimeApiStub.php deleted file mode 100644 index 0f0ee081b..000000000 --- a/tests/Unit/Lambda/RuntimeApiStub.php +++ /dev/null @@ -1,183 +0,0 @@ -invocations = $invocations; - $this->writeSettings(); - } - - public function setPollWindowMs(int $pollWindowMs): void - { - $this->pollWindowMs = $pollWindowMs; - $this->writeSettings(); - } - - public function omitHeaders(): void - { - $this->sendHeaders = false; - $this->writeSettings(); - } - - public function failAcknowledgementsWith(int $status): void - { - $this->acknowledgementStatus = $status; - $this->writeSettings(); - } - - public function host(): string - { - return "127.0.0.1:{$this->port}"; - } - - /** - * @return list - */ - public function acknowledgements(): array - { - return \array_column($this->acknowledgementLog(), 'path'); - } - - public function lastBody(): string - { - $log = $this->acknowledgementLog(); - - return $log === [] ? '' : (string) \end($log)['body']; - } - - public function start(): void - { - $this->port = self::freePort(); - $this->writeSettings(); - \file_put_contents($this->router(), self::ROUTER); - - $this->process = new Process( - [\PHP_BINARY, '-S', $this->host(), $this->router()], - env: ['STUB_STATE' => $this->stateDir], - timeout: null, - ); - $this->process->start(); - - $deadline = \microtime(true) + self::START_TIMEOUT_SECONDS; - while (\microtime(true) < $deadline) { - $connection = @\fsockopen('127.0.0.1', $this->port, $code, $message, 0.2); - if ($connection !== false) { - \fclose($connection); - - return; - } - - \usleep(50_000); - } - - throw new \RuntimeException('The Runtime API stub did not start'); - } - - public function stop(): void - { - $this->process?->stop(timeout: 5); - $this->process = null; - } - - private static function freePort(): int - { - $socket = \stream_socket_server('tcp://127.0.0.1:0', $code, $message); - if ($socket === false) { - throw new \RuntimeException("Unable to reserve a port: {$message}"); - } - - $name = (string) \stream_socket_get_name($socket, false); - \fclose($socket); - - return (int) \substr($name, \strrpos($name, ':') + 1); - } - - /** - * @return list - */ - private function acknowledgementLog(): array - { - $log = @\file_get_contents($this->stateDir . '/acknowledgements.json'); - if ($log === false) { - return []; - } - - return \array_map( - static fn(string $line): array => (array) \json_decode($line, true), - \array_filter(\explode("\n", \trim($log))), - ); - } - - private function writeSettings(): void - { - \file_put_contents($this->stateDir . '/settings.json', (string) \json_encode([ - 'invocations' => $this->invocations, - 'pollWindowMs' => $this->pollWindowMs, - 'sendHeaders' => $this->sendHeaders, - 'acknowledgementStatus' => $this->acknowledgementStatus, - ])); - } - - private function router(): string - { - return $this->stateDir . '/runtime-api.php'; - } - - private const ROUTER = <<<'PHP' - = (int) $settings['invocations']) { - \http_response_code(500); - - return; - } - - if ($settings['sendHeaders'] === true) { - \header('Lambda-Runtime-Aws-Request-Id: request-' . $served); - \header('Lambda-Runtime-Deadline-Ms: ' . ( - (int) (\microtime(true) * 1000) + (int) $settings['pollWindowMs'] - )); - } - - return; - } - - \file_put_contents( - $state . '/acknowledgements.json', - \json_encode(['path' => $path, 'body' => (string) \file_get_contents('php://input')]) . "\n", - \FILE_APPEND, - ); - - \http_response_code((int) $settings['acknowledgementStatus']); - PHP; -} diff --git a/tests/Unit/Lambda/RuntimeApiTestCase.php b/tests/Unit/Lambda/RuntimeApiTestCase.php deleted file mode 100644 index 9a77df207..000000000 --- a/tests/Unit/Lambda/RuntimeApiTestCase.php +++ /dev/null @@ -1,123 +0,0 @@ -api->nextInvocation(); - - self::assertSame('request-0', $invocation->requestId); - self::assertGreaterThanOrEqual($before, $invocation->deadlineMs); - } - - public function testMissingRequestIdHeaderIsRejected(): void - { - $this->stub->omitHeaders(); - - $this->expectExceptionMessage('Runtime API response is missing Lambda-Runtime-Aws-Request-Id'); - - $this->api->nextInvocation(); - } - - public function testAnExhaustedRuntimeApiReportsTheHttpStatus(): void - { - $this->api->nextInvocation(); - - $this->expectExceptionMessage('Failed to fetch the next invocation: HTTP 500'); - - $this->api->nextInvocation(); - } - - public function testResponseIsPostedForTheGivenRequestId(): void - { - $this->api->respond('request-7'); - - self::assertSame( - ['/2018-06-01/runtime/invocation/request-7/response'], - $this->stub->acknowledgements(), - ); - } - - public function testInitErrorCarriesTheExceptionType(): void - { - $this->api->reportInitError(new \LogicException('template is unreadable')); - - self::assertSame(['/2018-06-01/runtime/init/error'], $this->stub->acknowledgements()); - self::assertSame( - [ - 'errorType' => 'LogicException', - 'errorMessage' => 'template is unreadable', - ], - \array_intersect_key( - (array) \json_decode($this->stub->lastBody(), true), - ['errorType' => null, 'errorMessage' => null], - ), - ); - } - - public function testANamespacedExceptionTypeUsesDotsAsSeparators(): void - { - $this->api->reportInitError(new \Temporal\Lambda\Exception\ConfigurationException('nope')); - - self::assertSame( - 'Temporal.Lambda.Exception.ConfigurationException', - ((array) \json_decode($this->stub->lastBody(), true))['errorType'], - ); - } - - public function testAFailedAcknowledgementReportsTheHttpStatus(): void - { - $this->stub->failAcknowledgementsWith(503); - - $this->expectExceptionMessage('Failed to send the invocation response: HTTP 503'); - - $this->api->respond('request-0'); - } - - protected function setUp(): void - { - parent::setUp(); - - $this->stateDir = \sys_get_temp_dir() . '/temporal-lambda-api-' . \bin2hex(\random_bytes(6)); - \mkdir($this->stateDir); - - $this->stub = new RuntimeApiStub($this->stateDir, 5_000); - $this->stub->start(); - $this->api = new RuntimeApi($this->stub->host()); - } - - protected function tearDown(): void - { - $this->stub->stop(); - - foreach ((array) \glob($this->stateDir . '/{,.}[!.,..]*', \GLOB_BRACE) as $file) { - @\unlink((string) $file); - } - - \rmdir($this->stateDir); - - parent::tearDown(); - } -} diff --git a/tests/Unit/Lambda/RuntimeTestCase.php b/tests/Unit/Lambda/RuntimeTestCase.php deleted file mode 100644 index f1bf8adc5..000000000 --- a/tests/Unit/Lambda/RuntimeTestCase.php +++ /dev/null @@ -1,217 +0,0 @@ -fakeRoadRunner('\\sleep(30);'); - - $this->runUntilTheStubStops($binary); - - self::assertSame( - ['/2018-06-01/runtime/invocation/request-0/response'], - $this->api->acknowledgements(), - ); - } - - public function testRoadRunnerIsStoppedBeforeTheDeadline(): void - { - $binary = $this->fakeRoadRunner($this->recordStart('runs') . ' \\sleep(30);'); - - $elapsed = $this->runUntilTheStubStops($binary); - - self::assertSame('x', \file_get_contents($this->workDir . '/runs')); - self::assertLessThan(self::POLL_WINDOW_MS + self::BUFFER_MS, $elapsed); - self::assertSame([], $this->survivingProcesses($binary)); - } - - public function testRoadRunnerCrashIsReportedAsAnInvocationError(): void - { - $binary = $this->fakeRoadRunner('exit(3);'); - - $this->runUntilTheStubStops($binary); - - self::assertSame( - ['/2018-06-01/runtime/invocation/request-0/error'], - $this->api->acknowledgements(), - ); - self::assertStringContainsString( - 'RoadRunner exited before the shutdown window with code 3', - $this->api->lastBody(), - ); - self::assertContains( - 'roadrunner had already exited with code 3 when the shutdown started', - $this->logger->messages, - ); - } - - public function testInvocationWithoutEnoughTimeIsReportedWithoutStartingRoadRunner(): void - { - $binary = $this->fakeRoadRunner($this->recordStart('started') . ' \\sleep(30);'); - $this->api->setPollWindowMs(0); - - $this->runUntilTheStubStops($binary); - - self::assertSame( - ['/2018-06-01/runtime/invocation/request-0/error'], - $this->api->acknowledgements(), - ); - self::assertStringContainsString('Insufficient invocation time', $this->api->lastBody()); - self::assertFileDoesNotExist($this->workDir . '/started'); - } - - public function testARoadRunnerThatIgnoresSigtermIsKilled(): void - { - $binary = $this->fakeRoadRunner( - '\\pcntl_async_signals(true); \\pcntl_signal(\\SIGTERM, static fn() => null); while (true) { \\sleep(1); }', - ); - - $elapsed = $this->runUntilTheStubStops($binary); - - self::assertSame([], $this->survivingProcesses($binary)); - self::assertLessThan(self::POLL_WINDOW_MS + self::BUFFER_MS, $elapsed); - self::assertStringContainsString( - 'SIGKILLed it', - \implode("\n", $this->logger->messages), - ); - } - - public function testEveryInvocationGetsItsOwnRoadRunner(): void - { - $binary = $this->fakeRoadRunner($this->recordStart('runs') . ' \\sleep(30);'); - $this->api->setInvocations(2); - - $this->runUntilTheStubStops($binary); - - self::assertSame( - [ - '/2018-06-01/runtime/invocation/request-0/response', - '/2018-06-01/runtime/invocation/request-1/response', - ], - $this->api->acknowledgements(), - ); - self::assertSame('xx', \file_get_contents($this->workDir . '/runs')); - } - - protected function setUp(): void - { - parent::setUp(); - - $this->workDir = \sys_get_temp_dir() . '/temporal-lambda-' . \bin2hex(\random_bytes(6)); - \mkdir($this->workDir); - \file_put_contents($this->workDir . '/.rr.yaml', "rpc:\n listen: tcp://127.0.0.1:6001\n"); - - $this->logger = new RecordingLogger(); - $this->api = new RuntimeApiStub($this->workDir, self::POLL_WINDOW_MS + self::BUFFER_MS); - $this->api->start(); - } - - protected function tearDown(): void - { - $this->api->stop(); - - foreach ((array) \glob($this->workDir . '/{,.}[!.,..]*', \GLOB_BRACE) as $file) { - @\unlink((string) $file); - } - - \rmdir($this->workDir); - @\unlink(ConfigFile::PATH); - - parent::tearDown(); - } - - private function runUntilTheStubStops(string $binary): int - { - $config = new Config( - runtimeApi: $this->api->host(), - taskRoot: $this->workDir, - roadRunnerBinary: $binary, - roadRunnerConfigTemplate: $this->workDir . '/.rr.yaml', - shutdownBufferMs: self::BUFFER_MS, - gracefulTimeoutMs: self::GRACEFUL_MS, - ); - - $runtime = new Runtime($config, new RuntimeApi($config->runtimeApi), $this->logger); - - $startedAt = Clock::nowMs(); - try { - $runtime->run(); - self::fail('The runtime loop returned instead of failing on the exhausted stub'); - } catch (\RuntimeException $e) { - self::assertStringContainsString('fetch the next invocation', $e->getMessage()); - } - - return Clock::nowMs() - $startedAt; - } - - /** - * @return list - */ - private function survivingProcesses(string $binary): array - { - $output = []; - \exec('ps -eo pid=,command=', $output); - - $survivors = []; - foreach ($output as $line) { - if (\str_contains($line, $binary)) { - $survivors[] = (int) \trim($line); - } - } - - return $survivors; - } - - private function recordStart(string $file): string - { - return \sprintf( - '\\file_put_contents(%s, "x", \\FILE_APPEND);', - \var_export($this->workDir . '/' . $file, true), - ); - } - - private function fakeRoadRunner(string $body): string - { - $path = $this->workDir . '/rr'; - \file_put_contents($path, "#!/usr/bin/env php\n