|
| 1 | +<?php |
| 2 | + |
| 3 | +declare(strict_types=1); |
| 4 | + |
| 5 | +namespace Temporal\Tests\Unit\Router; |
| 6 | + |
| 7 | +use DateTimeImmutable; |
| 8 | +use Psr\Log\NullLogger; |
| 9 | +use React\Promise\Deferred; |
| 10 | +use Spiral\Attributes\AnnotationReader; |
| 11 | +use Spiral\Attributes\AttributeReader; |
| 12 | +use Spiral\Attributes\Composite\SelectiveReader; |
| 13 | +use Spiral\Attributes\ReaderInterface; |
| 14 | +use Temporal\DataConverter\DataConverter; |
| 15 | +use Temporal\DataConverter\EncodedValues; |
| 16 | +use Temporal\DataConverter\ValuesInterface; |
| 17 | +use Temporal\Exception\ExceptionInterceptorInterface; |
| 18 | +use Temporal\Interceptor\SimplePipelineProvider; |
| 19 | +use Temporal\Internal\Declaration\Reader\ActivityReader; |
| 20 | +use Temporal\Internal\Marshaller\Mapper\AttributeMapperFactory; |
| 21 | +use Temporal\Internal\Marshaller\Marshaller; |
| 22 | +use Temporal\Internal\Queue\QueueInterface; |
| 23 | +use Temporal\Internal\ServiceContainer; |
| 24 | +use Temporal\Internal\Transport\ClientInterface; |
| 25 | +use Temporal\Internal\Transport\Router\InvokeActivity; |
| 26 | +use Temporal\Tests\Unit\AbstractUnit; |
| 27 | +use Temporal\Worker\Environment\EnvironmentInterface; |
| 28 | +use Temporal\Worker\LoopInterface; |
| 29 | +use Temporal\Worker\Transport\Command\Server\ServerRequest; |
| 30 | +use Temporal\Worker\Transport\Command\Server\TickInfo; |
| 31 | +use Temporal\Worker\Transport\RPCConnectionInterface; |
| 32 | + |
| 33 | +final class InvokeActivityHeartbeatContextTestCase extends AbstractUnit |
| 34 | +{ |
| 35 | + private const NAMESPACE = 'test-namespace'; |
| 36 | + private const TASK_QUEUE = 'test-task-queue'; |
| 37 | + private const ACTIVITY = 'HeartbeatDetailsActivity.ReadSignature'; |
| 38 | + |
| 39 | + public function testHeartbeatDetailsAreDecodedWithTheActivityContext(): void |
| 40 | + { |
| 41 | + $signature = $this->runActivity(); |
| 42 | + |
| 43 | + $this->assertSame( |
| 44 | + \sprintf('act|%s|%s|%s', self::NAMESPACE, self::ACTIVITY, self::TASK_QUEUE), |
| 45 | + $signature, |
| 46 | + ); |
| 47 | + } |
| 48 | + |
| 49 | + private function runActivity(): string |
| 50 | + { |
| 51 | + $converter = new DataConverter(new ContextSignatureConverter()); |
| 52 | + $services = $this->createServices($converter); |
| 53 | + $router = new InvokeActivity($services, $this->createMock(RPCConnectionInterface::class), new SimplePipelineProvider()); |
| 54 | + |
| 55 | + $resolver = new Deferred(); |
| 56 | + $router->handle($this->createRequest($converter), [], $resolver); |
| 57 | + |
| 58 | + $result = null; |
| 59 | + $resolver->promise()->then(static function (ValuesInterface $values) use (&$result): void { |
| 60 | + $result = $values->getValue(0); |
| 61 | + }); |
| 62 | + |
| 63 | + return $result; |
| 64 | + } |
| 65 | + |
| 66 | + private function createRequest(DataConverter $converter): ServerRequest |
| 67 | + { |
| 68 | + $options = [ |
| 69 | + 'name' => self::ACTIVITY, |
| 70 | + 'heartbeatDetails' => 1, |
| 71 | + 'info' => [ |
| 72 | + 'TaskToken' => \base64_encode('token'), |
| 73 | + 'ActivityType' => ['Name' => self::ACTIVITY], |
| 74 | + 'ActivityID' => '1', |
| 75 | + 'WorkflowNamespace' => self::NAMESPACE, |
| 76 | + 'TaskQueue' => self::TASK_QUEUE, |
| 77 | + 'Attempt' => 2, |
| 78 | + ], |
| 79 | + ]; |
| 80 | + |
| 81 | + return new ServerRequest( |
| 82 | + 'InvokeActivity', |
| 83 | + new TickInfo(new DateTimeImmutable()), |
| 84 | + $options, |
| 85 | + EncodedValues::fromValues(['input', 'heartbeat'], $converter), |
| 86 | + ); |
| 87 | + } |
| 88 | + |
| 89 | + private function createServices(DataConverter $converter): ServiceContainer |
| 90 | + { |
| 91 | + $services = new ServiceContainer( |
| 92 | + $this->createMock(LoopInterface::class), |
| 93 | + $this->createMock(EnvironmentInterface::class), |
| 94 | + $this->createMock(ClientInterface::class), |
| 95 | + $this->createMock(ReaderInterface::class), |
| 96 | + $this->createMock(QueueInterface::class), |
| 97 | + new Marshaller(new AttributeMapperFactory(new AttributeReader())), |
| 98 | + $converter, |
| 99 | + $this->createMock(ExceptionInterceptorInterface::class), |
| 100 | + new SimplePipelineProvider(), |
| 101 | + new NullLogger(), |
| 102 | + ); |
| 103 | + |
| 104 | + $reader = new ActivityReader(new SelectiveReader([new AnnotationReader(), new AttributeReader()])); |
| 105 | + foreach ($reader->fromClass(HeartbeatDetailsActivity::class) as $proto) { |
| 106 | + $services->activities->add($proto); |
| 107 | + } |
| 108 | + |
| 109 | + return $services; |
| 110 | + } |
| 111 | +} |
0 commit comments