diff --git a/phpstan.neon b/phpstan.neon index 718611fbd..aef873d37 100644 --- a/phpstan.neon +++ b/phpstan.neon @@ -64,6 +64,7 @@ parameters: paths: - tests/unit/DurableLaravel/* - tests/unit/Bridge/Illuminate/* + - tests/integration/Temporal/Hosts/* # A test asserts at runtime what its types promise: an assertion PHPStan already proves is # still the check the test makes, and what keeps it from being a test without assertion. - identifier: staticMethod.alreadyNarrowedType diff --git a/tests/integration/Temporal/Fixtures/HeartbeatActivities.php b/tests/integration/Temporal/Fixtures/HeartbeatActivities.php new file mode 100644 index 000000000..0c2e4f22b --- /dev/null +++ b/tests/integration/Temporal/Fixtures/HeartbeatActivities.php @@ -0,0 +1,26 @@ +heartbeat->sendHeartbeat(['beat' => $beat]); + } + + return $seconds; + } + + public function heartbeatUntilCancelled(string $marker): string + { + for ($beat = 1; $beat <= self::MAX_SECONDS; ++$beat) { + sleep(1); + if ($this->heartbeat->sendHeartbeat(['beat' => $beat])) { + file_put_contents($marker, 'cancel-requested'); + + return 'cancelled'; + } + } + file_put_contents($marker, 'never-cancelled'); + + return 'never cancelled'; + } +} diff --git a/tests/integration/Temporal/Fixtures/IntegrationWorkflows.php b/tests/integration/Temporal/Fixtures/IntegrationWorkflows.php index 779518323..dc1b4d18d 100644 --- a/tests/integration/Temporal/Fixtures/IntegrationWorkflows.php +++ b/tests/integration/Temporal/Fixtures/IntegrationWorkflows.php @@ -175,6 +175,35 @@ public static function registerWorkflows(WorkflowRegistry $registry): void $env->childWorkflowStub(DoublerWorkflow::class)->run((int) ($input['value'] ?? 0)), )]; }); + + // One attempt: an activity whose heartbeats never arrive fails on its heartbeat timeout + // instead of being retried until the test gives up (#518). + $registry->registerFactory('HeartbeatsThroughItsTimeout', static fn(array $input) => static fn(WorkflowEnvironment $env): array => [ + 'beats' => $env->await($env->activityStub(HeartbeatActivities::class, self::heartbeatOptions())->heartbeatFor((int) ($input['seconds'] ?? 30))), + ]); + + // The deadline cancels the activity it outlives: the server records the request and hands + // it to the activity's next heartbeat (#518). The run then stays open for a few heartbeats: + // once it closes, a heartbeat hears "not found", not "cancel requested". + $registry->registerFactory('CancelsItsHeartbeatingActivity', static fn(array $input) => static function (WorkflowEnvironment $env) use ($input): array { + try { + $env->await($env->activityStub(HeartbeatActivities::class, self::heartbeatOptions())->heartbeatUntilCancelled((string) ($input['marker'] ?? '')), Duration::seconds(3)); + } catch (DeadlineExceededException) { + $env->sleep(Duration::seconds(10)); + + return ['deadline' => true]; + } + + return ['deadline' => false]; + }); + } + + private static function heartbeatOptions(): ActivityOptions + { + return new ActivityOptions( + RetryLimit::once(), + timeouts: ActivityTimeouts::attempt(Duration::seconds(60))->withHeartbeat(Duration::seconds(5)), + ); } /** Exposed for {@see DoublerWorkflow}, which lives outside this class. */ diff --git a/tests/integration/Temporal/HeartbeatOnAHostTestCase.php b/tests/integration/Temporal/HeartbeatOnAHostTestCase.php new file mode 100644 index 000000000..9c7d5a08f --- /dev/null +++ b/tests/integration/Temporal/HeartbeatOnAHostTestCase.php @@ -0,0 +1,87 @@ + */ + private array $markers = []; + + protected function tearDown(): void + { + foreach ($this->markers as $marker) { + @unlink($marker); + } + parent::tearDown(); + } + + /** + * Thirty seconds of heartbeats, one a second, under a five-second heartbeat timeout and a single + * attempt: the activity completes only if its heartbeats reach the server through the sender its + * host injected. With the no-op sender, the server times it out after five seconds and the run + * never completes (until #544 it did not even fail: the bounded wait below reports it). + */ + public function testAHeartbeatingActivityOutlivesItsHeartbeatTimeout(): void + { + self::assertSame(['beats' => 30], $this->runWithin('HeartbeatsThroughItsTimeout', ['seconds' => 30], 90.0)); + } + + /** + * The workflow's deadline cancels the activity it outlives; the server hands the request to the + * activity's next heartbeat, and sendHeartbeat() returns true. + */ + public function testACancellationRequestedByTheServerReachesSendHeartbeat(): void + { + $marker = $this->markers[] = sys_get_temp_dir() . '/durable-heartbeat-' . bin2hex(random_bytes(6)); + + self::assertSame(['deadline' => true], $this->runWithin('CancelsItsHeartbeatingActivity', ['marker' => $marker], 60.0)); + + $deadline = microtime(true) + self::MARKER_TIMEOUT_SECONDS; + while (!is_file($marker) && microtime(true) < $deadline) { + usleep(250_000); + } + self::assertFileExists($marker, \sprintf('No heartbeat of the %s-hosted activity reported the cancellation within %.0f s.', $this->activityWorkerRole(), self::MARKER_TIMEOUT_SECONDS)); + self::assertSame('cancel-requested', file_get_contents($marker)); + } + + /** + * Starts the workflow and returns its result, failing after `$seconds` of wall-clock time: the + * status is read with a describe, which does not block, and the result only once the run closed. + * + * @param array $input + */ + private function runWithin(string $workflowType, array $input, float $seconds): mixed + { + $executionId = $this->startWorkflow($workflowType, $input); + $request = new DescribeWorkflowExecutionRequest([ + 'namespace' => $this->connection->namespace->name(), + 'execution' => new WorkflowExecution(['workflow_id' => $this->workflowId($executionId)]), + ]); + + $deadline = microtime(true) + $seconds; + do { + $status = $this->client->DescribeWorkflowExecution($request, [], ['timeout' => 5_000_000])->getWorkflowExecutionInfo()?->getStatus(); + if (WorkflowExecutionStatus::WORKFLOW_EXECUTION_STATUS_RUNNING !== $status) { + return $this->workflowClient()->pollForCompletion($executionId, 0, 1); + } + usleep(250_000); + } while (microtime(true) < $deadline); + + self::fail(\sprintf('The %s run did not close within %.0f s; its activities are hosted by the %s worker.', $workflowType, $seconds, $this->activityWorkerRole())); + } +} diff --git a/tests/integration/Temporal/Hosts/laravel-activity.php b/tests/integration/Temporal/Hosts/laravel-activity.php new file mode 100644 index 000000000..d0239065b --- /dev/null +++ b/tests/integration/Temporal/Hosts/laravel-activity.php @@ -0,0 +1,66 @@ + [transport]. + */ + +use Gplanchat\Bridge\Temporal\Worker\TemporalActivityWorker; +use Gplanchat\Durable\Activity\NullActivityHeartbeatSender; +use Gplanchat\Durable\Activity\PayloadToContractMethodInvoker; +use Gplanchat\Durable\Laravel\DurableServiceProvider; +use Gplanchat\Durable\Port\ActivityHeartbeatSenderInterface; +use Gplanchat\Durable\RegistryActivityExecutor; +use Illuminate\Container\Container; +use integration\Temporal\Fixtures\HeartbeatActivities; +use integration\Temporal\Fixtures\HeartbeatingActivities; + +$arguments = $_SERVER['argv'] ?? null; +if (!\is_array($arguments) || !isset($arguments[1], $arguments[2], $arguments[3])) { + fwrite(\STDERR, "Usage: worker.php
[transport]\n"); + + exit(1); +} +[$address, $namespace, $taskQueue] = [$arguments[1], $arguments[2], $arguments[3]]; +$transport = $arguments[5] ?? 'auto'; + +$dsn = \sprintf( + 'temporal://%s?namespace=%s&journal_task_queue=%3$s&workflow_task_queue=%3$s&activity_task_queue=%3$s&transport=%4$s', + $address, + rawurlencode($namespace), + rawurlencode($taskQueue), + rawurlencode($transport), +); + +$app = new Container(); +$app->instance('config', new ArrayObject(['durable' => ['backend' => 'temporal', 'temporal' => ['dsn' => $dsn]]], ArrayObject::ARRAY_AS_PROPS)); +(new DurableServiceProvider($app))->register(); + +if ('noop' === getenv('DURABLE_HEARTBEAT')) { + $app->instance(ActivityHeartbeatSenderInterface::class, new NullActivityHeartbeatSender()); +} +$activities = $app->make(HeartbeatingActivities::class); +$executor = $app->make(RegistryActivityExecutor::class); +foreach ((new ReflectionClass(HeartbeatActivities::class))->getMethods() as $method) { + $name = $method->getAttributes(Gplanchat\Durable\Attribute\AsActivityMethod::class)[0]->newInstance()->name; + $executor->register($name, new PayloadToContractMethodInvoker($activities, HeartbeatActivities::class, $method->getName())); +} + +$worker = $app->make(TemporalActivityWorker::class); +$deadline = microtime(true) + (float) (getenv('DURABLE_HOST_WORKER_MAX_TIME') ?: 180); +while (microtime(true) < $deadline) { + try { + $worker->pollOnce(); + } catch (Throwable $e) { + fwrite(STDERR, 'laravel activity worker: ' . $e::class . ': ' . $e->getMessage() . "\n"); + } +} diff --git a/tests/integration/Temporal/Hosts/magento-activity.php b/tests/integration/Temporal/Hosts/magento-activity.php new file mode 100644 index 000000000..0177415bb --- /dev/null +++ b/tests/integration/Temporal/Hosts/magento-activity.php @@ -0,0 +1,54 @@ + [transport]. + */ + +use Gplanchat\Durable\Activity\NullActivityHeartbeatSender; +use Gplanchat\DurableModule\Runtime\RuntimeFactory; +use Gplanchat\DurableModule\Runtime\SharedActivityHeartbeatSender; +use integration\Temporal\Fixtures\HeartbeatingActivities; + +$arguments = $_SERVER['argv'] ?? null; +if (!\is_array($arguments) || !isset($arguments[1], $arguments[2], $arguments[3])) { + fwrite(\STDERR, "Usage: worker.php
[transport]\n"); + + exit(1); +} +[$address, $namespace, $taskQueue] = [$arguments[1], $arguments[2], $arguments[3]]; +$transport = $arguments[5] ?? 'auto'; + +$dsn = \sprintf( + 'temporal://%s?namespace=%s&journal_task_queue=%3$s&workflow_task_queue=%3$s&activity_task_queue=%3$s&transport=%4$s', + $address, + rawurlencode($namespace), + rawurlencode($taskQueue), + rawurlencode($transport), +); + +$shared = new SharedActivityHeartbeatSender(); +$factory = new RuntimeFactory( + activityHandlers: [new HeartbeatingActivities('noop' === getenv('DURABLE_HEARTBEAT') ? new NullActivityHeartbeatSender() : $shared)], + temporalDsn: $dsn, + heartbeat: $shared, +); + +$worker = $factory->activityWorker(); +$deadline = microtime(true) + (float) (getenv('DURABLE_HOST_WORKER_MAX_TIME') ?: 180); +while (microtime(true) < $deadline) { + try { + $worker->pollOnce(); + } catch (Throwable $e) { + fwrite(STDERR, 'magento activity worker: ' . $e::class . ': ' . $e->getMessage() . "\n"); + } +} diff --git a/tests/integration/Temporal/LaravelHostedHeartbeatTest.php b/tests/integration/Temporal/LaravelHostedHeartbeatTest.php new file mode 100644 index 000000000..ca64f0f7f --- /dev/null +++ b/tests/integration/Temporal/LaravelHostedHeartbeatTest.php @@ -0,0 +1,17 @@ +client = WorkflowServiceClientFactory::create($this->connection); $this->spawnWorker('workflow'); - $this->spawnWorker('activity'); + $this->spawnWorker($this->activityWorkerRole()); + } + + /** + * Who hosts the activities: this suite's own worker by default, or a framework's (#518). + */ + protected function activityWorkerRole(): string + { + return 'activity'; } public static function transportFromEnv(): string diff --git a/tests/integration/Temporal/worker.php b/tests/integration/Temporal/worker.php index d58da85d4..90fb0c1c5 100644 --- a/tests/integration/Temporal/worker.php +++ b/tests/integration/Temporal/worker.php @@ -42,6 +42,14 @@ ); $client = WorkflowServiceClientFactory::create($connection); +// An activity worker a framework builds, not this script: the host's container wires the sender +// its activities inject (#518). +if (\in_array($role, ['laravel-activity', 'magento-activity'], true)) { + require __DIR__ . '/Hosts/' . $role . '.php'; + + exit(0); +} + if ('workflow' === $role) { $registry = new WorkflowRegistry(); IntegrationWorkflows::registerWorkflows($registry);