Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions phpstan.neon
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
26 changes: 26 additions & 0 deletions tests/integration/Temporal/Fixtures/HeartbeatActivities.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
<?php

declare(strict_types=1);

namespace integration\Temporal\Fixtures;

use Gplanchat\Durable\Attribute\AsActivityMethod;

/**
* The activities that heartbeat, hosted by a framework's own activity worker rather than by the
* suite's `worker.php` (#518).
*/
interface HeartbeatActivities
{
/** Heartbeats once a second for `$seconds`, then returns how many beats it sent. */
#[AsActivityMethod('heartbeat.for')]
public function heartbeatFor(int $seconds): int;

/**
* Heartbeats once a second until a heartbeat answers that the server asked for cancellation,
* then writes it to `$marker`: the test reads the file, since a cancelled activity's result
* reaches no one.
*/
#[AsActivityMethod('heartbeat.until_cancelled')]
public function heartbeatUntilCancelled(string $marker): string;
}
42 changes: 42 additions & 0 deletions tests/integration/Temporal/Fixtures/HeartbeatingActivities.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
<?php

declare(strict_types=1);

namespace integration\Temporal\Fixtures;

use Gplanchat\Durable\Port\ActivityHeartbeatSenderInterface;

/**
* Takes its sender the way an application's activity does: injected by the host's container (#510).
*/
final class HeartbeatingActivities implements HeartbeatActivities
{
private const MAX_SECONDS = 60;

public function __construct(private readonly ActivityHeartbeatSenderInterface $heartbeat) {}

public function heartbeatFor(int $seconds): int
{
for ($beat = 1; $beat <= $seconds; ++$beat) {
sleep(1);
$this->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';
}
}
29 changes: 29 additions & 0 deletions tests/integration/Temporal/Fixtures/IntegrationWorkflows.php
Original file line number Diff line number Diff line change
Expand Up @@ -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. */
Expand Down
87 changes: 87 additions & 0 deletions tests/integration/Temporal/HeartbeatOnAHostTestCase.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
<?php

declare(strict_types=1);

namespace integration\Temporal;

use Temporal\Api\Common\V1\WorkflowExecution;
use Temporal\Api\Enums\V1\WorkflowExecutionStatus;
use Temporal\Api\Workflowservice\V1\DescribeWorkflowExecutionRequest;

/**
* The heartbeat scenarios of #510 against a real server, with the activities hosted by a
* framework's activity worker rather than by this suite's own (#518). The workflows run on the
* suite's workflow worker; only the activity side changes host.
*
* The workers are killed in the parent's tearDown(), which runs whether the test passed or not.
*/
abstract class HeartbeatOnAHostTestCase extends TemporalServerTestCase
{
private const MARKER_TIMEOUT_SECONDS = 30.0;

/** @var list<string> */
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<string, mixed> $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()));
}
}
66 changes: 66 additions & 0 deletions tests/integration/Temporal/Hosts/laravel-activity.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
<?php

declare(strict_types=1);

/*
* The Temporal activity worker as a Laravel application gets it: `DurableServiceProvider` on a bare
* container with `backend: temporal`, the activities resolved from that container, so the sender
* they inject is the provider's (#510, #518). The loop is `durable:temporal-worker --role=activity`
* with `--max-time`: illuminate/console is not a root dependency, so the command itself cannot run.
*
* DURABLE_HEARTBEAT=noop hands the activities the no-op sender instead, to show the heartbeat test
* red. Only this script reads it: nothing under src/ knows it.
*
* Required by worker.php, whose arguments it reads: <address> <namespace> <taskQueue> <role> [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 <address> <namespace> <taskQueue> <role> [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");
}
}
54 changes: 54 additions & 0 deletions tests/integration/Temporal/Hosts/magento-activity.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
<?php

declare(strict_types=1);

/*
* The Temporal activity worker as the Magento module gets it: `RuntimeFactory::activityWorker()`,
* with the shared sender `di.xml` hands the activities (#510, #518). The factory is plain PHP, so
* it runs here without a Mage-OS install; what Magento adds around it is the ObjectManager and the
* `bin/magento` command that turns this loop.
*
* DURABLE_HEARTBEAT=noop hands the activities the no-op sender instead, to show the heartbeat test
* red. Only this script reads it: nothing under src/ knows it.
*
* Required by worker.php, whose arguments it reads: <address> <namespace> <taskQueue> <role> [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 <address> <namespace> <taskQueue> <role> [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");
}
}
17 changes: 17 additions & 0 deletions tests/integration/Temporal/LaravelHostedHeartbeatTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
<?php

declare(strict_types=1);

namespace integration\Temporal;

/**
* The heartbeat scenarios with the activities hosted by the Laravel integration's activity worker
* (see Hosts/laravel-activity.php).
*/
final class LaravelHostedHeartbeatTest extends HeartbeatOnAHostTestCase
{
protected function activityWorkerRole(): string
{
return 'laravel-activity';
}
}
17 changes: 17 additions & 0 deletions tests/integration/Temporal/MagentoHostedHeartbeatTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
<?php

declare(strict_types=1);

namespace integration\Temporal;

/**
* The heartbeat scenarios with the activities hosted by the Magento integration's activity worker
* (see Hosts/magento-activity.php).
*/
final class MagentoHostedHeartbeatTest extends HeartbeatOnAHostTestCase
{
protected function activityWorkerRole(): string
{
return 'magento-activity';
}
}
10 changes: 9 additions & 1 deletion tests/integration/Temporal/TemporalServerTestCase.php
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,15 @@ protected function setUp(): void
$this->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
Expand Down
8 changes: 8 additions & 0 deletions tests/integration/Temporal/worker.php
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Loading