diff --git a/UPGRADE.md b/UPGRADE.md index 45631dad07..f674212f96 100644 --- a/UPGRADE.md +++ b/UPGRADE.md @@ -56,14 +56,19 @@ Platform ``` Accordingly, `Bridge\MiniMax\MiniMaxResultConverter` no longer takes an HTTP client, API key, - endpoint or clock, only an optional `MiniMaxJobClient` that creates the job handles; polling moved to the new `Bridge\MiniMax\MiniMaxJobClient`. Code building the + endpoint or clock - polling moved to the new `Bridge\MiniMax\MiniMaxJobClient`. Code building the bridge through `Bridge\MiniMax\Factory` is unaffected. ```diff -$converter = new MiniMaxResultConverter($httpClient, $apiKey, $endpoint, $clock); - +$converter = new MiniMaxResultConverter($jobClient); + +$converter = new MiniMaxResultConverter(); ``` + Its only argument left is the name of the provider it belongs to, which it stamps onto the + handles of the jobs it starts so they can be resolved later. It defaults to `minimax` and only + needs to be passed when the provider was registered under a different name - which + `Factory::createProvider()` does on its own. + Store ----- diff --git a/docs/components/platform.rst b/docs/components/platform.rst index 6f99093add..bd5c9fe694 100644 --- a/docs/components/platform.rst +++ b/docs/components/platform.rst @@ -1127,10 +1127,11 @@ The handle holds no connection and no client, only what is needed to ask the pro again, so it can be stored and picked up somewhere else entirely:: // in the process that started the job - $repository->save($handle->getId(), $handle->toString()); + $jobId = $handle->getId(); + $repository->save($jobId, $handle->toString()); // in a worker, possibly much later - $handle = JobHandle::fromString($repository->load($id)); + $handle = JobHandle::fromString($repository->load($jobId)); if ($jobClient->getStatus($handle)->is(JobStateCase::SUCCEEDED)) { $result = $jobClient->getResult($handle); @@ -1175,27 +1176,35 @@ be given ten minutes in a worker and five seconds inside a web request. Say so p $result = $runner->wait($jobClient, $handle, maxDuration: 5); A budget passed to the runner's constructor applies to every job it waits for and sits between the -two: it overrules what a job asks for, and a single call overrules it in turn. +two: it overrules what a job asks for, and a single call overrules it in turn. Whichever wins, it is +spent as wall clock rather than as a number of polls - asking the provider takes time too, so five +seconds mean five seconds and not five requests that may each take one. + +Before polling, the runner asks the client whether the handle is one it can resolve +(:method:`Symfony\\AI\\Platform\\Job\\JobClientInterface::supports`) - what that means is the +bridge's own judgement. A handle the client turns down raises an ``InvalidArgumentException`` rather +than a request that fails halfway. In a Symfony application a runner using the application clock is available as ``ai.platform.job_runner`` and autowired through :class:`Symfony\\AI\\Platform\\Job\\JobRunner`. It carries no budget of its own, so the same shared service serves a job finishing in seconds and one running for minutes. Each job-capable platform also registers its client as -``ai.platform.job_client.``, autowired by argument name:: +``ai.platform.job_client.``, autowired by the platform name as argument name - so the argument +of a MiniMax job client has to be called ``$minimax``:: public function __construct( private JobRunner $jobRunner, - private JobClientInterface $minimaxJobClient, + private JobClientInterface $minimax, ) { } public function __invoke(JobHandle $handle): void { // trust the job - $this->jobRunner->wait($this->minimaxJobClient, $handle); + $this->jobRunner->wait($this->minimax, $handle); // or bound it to what a request can afford - $this->jobRunner->wait($this->minimaxJobClient, $handle, maxDuration: 5); + $this->jobRunner->wait($this->minimax, $handle, maxDuration: 5); } An application holding handles of several providers picks the client by the name the handle carries, @@ -1210,7 +1219,10 @@ from a locator over the ``ai.platform.job_client`` tag:: public function __invoke(JobHandle $handle): void { - $this->jobRunner->wait($this->jobClients->get($handle->getProvider()), $handle); + // A handle only names a provider when the bridge that created it stated one. + $provider = $handle->getProvider() ?? throw new \InvalidArgumentException('The job handle does not name a provider.'); + + $this->jobRunner->wait($this->jobClients->get($provider), $handle); } The runner throws a :class:`Symfony\\AI\\Platform\\Exception\\JobFailedException` when the provider diff --git a/examples/minimax/.gitignore b/examples/minimax/.gitignore index bf01b35b98..ca59d0d582 100644 --- a/examples/minimax/.gitignore +++ b/examples/minimax/.gitignore @@ -1,2 +1,4 @@ +minimax-speech.mp3 +minimax-video-job.json minimax-video.mp4 text-to-image.jpg diff --git a/examples/minimax/text-to-speech-async.php b/examples/minimax/text-to-speech-async.php index 30be300873..c68cfbb301 100644 --- a/examples/minimax/text-to-speech-async.php +++ b/examples/minimax/text-to-speech-async.php @@ -15,10 +15,10 @@ require_once dirname(__DIR__).'/bootstrap.php'; -$provider = Factory::createProvider(env('MINI_MAX_API_KEY'), http_client()); +$platform = Factory::createPlatform(env('MINI_MAX_API_KEY'), http_client()); // The async endpoint enqueues a task, so the invocation hands back a job handle instead of audio. -$handle = $provider->invoke('speech-2.6-hd', new Text('The real danger is not that computers start thinking like people, but that people start thinking like computers.'), [ +$handle = $platform->invoke('speech-2.6-hd', new Text('The real danger is not that computers start thinking like people, but that people start thinking like computers.'), [ 'async' => true, 'voice_setting' => [ 'voice_id' => 'English_expressive_narrator', diff --git a/examples/minimax/text-to-video.php b/examples/minimax/text-to-video.php index 11393acbb2..8bc541df67 100644 --- a/examples/minimax/text-to-video.php +++ b/examples/minimax/text-to-video.php @@ -15,11 +15,11 @@ require_once dirname(__DIR__).'/bootstrap.php'; -$provider = Factory::createProvider(env('MINI_MAX_API_KEY'), http_client()); +$platform = Factory::createPlatform(env('MINI_MAX_API_KEY'), http_client()); // Video generation is asynchronous: MiniMax accepts the request and answers with a task, so the // invocation returns a handle instead of a video. -$handle = $provider->invoke('MiniMax-Hailuo-02', new Text('A cat playing the piano on a stage, cinematic lighting'), [ +$handle = $platform->invoke('MiniMax-Hailuo-02', new Text('A cat playing the piano on a stage, cinematic lighting'), [ 'duration' => 6, 'resolution' => '768P', ])->asJob(); diff --git a/examples/minimax/video-job-resume.php b/examples/minimax/video-job-resume.php index f0e5a5c181..d4eada117a 100644 --- a/examples/minimax/video-job-resume.php +++ b/examples/minimax/video-job-resume.php @@ -29,8 +29,8 @@ $storage = __DIR__.'/minimax-video-job.json'; if (!is_file($storage)) { - $provider = Factory::createProvider(env('MINI_MAX_API_KEY'), http_client()); - $handle = $provider->invoke('MiniMax-Hailuo-02', new Text('A cat playing the piano on a stage, cinematic lighting'), [ + $platform = Factory::createPlatform(env('MINI_MAX_API_KEY'), http_client()); + $handle = $platform->invoke('MiniMax-Hailuo-02', new Text('A cat playing the piano on a stage, cinematic lighting'), [ 'duration' => 6, 'resolution' => '768P', ])->asJob(); diff --git a/src/ai-bundle/CHANGELOG.md b/src/ai-bundle/CHANGELOG.md index f6bf5a04a6..0047ad3e91 100644 --- a/src/ai-bundle/CHANGELOG.md +++ b/src/ai-bundle/CHANGELOG.md @@ -1,5 +1,6 @@ CHANGELOG -== +========= + 0.14 ---- diff --git a/src/ai-bundle/src/AiBundle.php b/src/ai-bundle/src/AiBundle.php index 1cae817175..fb972bd98f 100644 --- a/src/ai-bundle/src/AiBundle.php +++ b/src/ai-bundle/src/AiBundle.php @@ -1088,7 +1088,6 @@ private function processPlatformConfig(string $type, array $platform, ContainerB $platform['api_key'], new Reference($platform['http_client'], ContainerInterface::NULL_ON_INVALID_REFERENCE), $platform['endpoint'], - 'minimax', ]) ->addTag('ai.platform.job_client', ['key' => 'minimax'])); $container->registerAliasForArgument($jobClientId, JobClientInterface::class, 'minimax'); diff --git a/src/platform/CHANGELOG.md b/src/platform/CHANGELOG.md index 7e8cfdeeec..c65792be32 100644 --- a/src/platform/CHANGELOG.md +++ b/src/platform/CHANGELOG.md @@ -1,5 +1,6 @@ CHANGELOG -== +========= + 0.14 ---- diff --git a/src/platform/src/Bridge/MiniMax/CHANGELOG.md b/src/platform/src/Bridge/MiniMax/CHANGELOG.md index beb6b7f8b8..a91edf3ad3 100644 --- a/src/platform/src/Bridge/MiniMax/CHANGELOG.md +++ b/src/platform/src/Bridge/MiniMax/CHANGELOG.md @@ -5,7 +5,7 @@ CHANGELOG ---- * Add model information to token usage extraction - * [BC BREAK] Stop polling asynchronous tasks inside `MiniMaxResultConverter`. Video generation and asynchronous speech synthesis now return a `Result\JobResult` carrying a serializable job handle, resolved through the new `MiniMaxJobClient`, built by `Factory::createJobClient()` and creating the handles — see the platform `UPGRADE` notes. `MiniMaxResultConverter` no longer takes an HTTP client, API key, endpoint or clock, only an optional `MiniMaxJobClient` + * [BC BREAK] Stop polling asynchronous tasks inside `MiniMaxResultConverter`. Video generation and asynchronous speech synthesis now return a `Result\JobResult` carrying a serializable job handle, resolved through the new `MiniMaxJobClient`, built by `Factory::createJobClient()` — see the platform `UPGRADE` notes. `MiniMaxResultConverter` no longer takes an HTTP client, API key, endpoint or clock, only an optional provider name, which it stamps onto the handles it creates 0.11 ---- diff --git a/src/platform/src/Bridge/MiniMax/Factory.php b/src/platform/src/Bridge/MiniMax/Factory.php index 986aec87cc..36c1fbd6ed 100644 --- a/src/platform/src/Bridge/MiniMax/Factory.php +++ b/src/platform/src/Bridge/MiniMax/Factory.php @@ -27,25 +27,26 @@ */ final class Factory { + private const DEFAULT_ENDPOINT = 'https://api.minimax.io/v1'; + /** * @param non-empty-string $name */ public static function createProvider( #[\SensitiveParameter] string $apiKey, ?HttpClientInterface $httpClient = null, - string $endpoint = 'https://api.minimax.io/v1', + string $endpoint = self::DEFAULT_ENDPOINT, ModelCatalogInterface $modelCatalog = new ModelCatalog(), ?Contract $contract = null, ?EventDispatcherInterface $eventDispatcher = null, string $name = 'minimax', ): ProviderInterface { $httpClient = $httpClient instanceof EventSourceHttpClient ? $httpClient : new EventSourceHttpClient($httpClient); - $jobClient = self::createJobClient($apiKey, $httpClient, $endpoint, $name); return new Provider( $name, [new MiniMaxClient($httpClient, $apiKey, $endpoint)], - [new MiniMaxResultConverter($jobClient)], + [new MiniMaxResultConverter($name)], $modelCatalog, $contract ?? MiniMaxContract::create(), $eventDispatcher, @@ -55,16 +56,13 @@ public static function createProvider( /** * The client resolving the jobs this bridge hands out - typically in a worker picking up a * stored handle, without a provider or platform at hand. - * - * @param string $name the provider name stated on the handles this client creates */ public static function createJobClient( #[\SensitiveParameter] string $apiKey, ?HttpClientInterface $httpClient = null, - string $endpoint = 'https://api.minimax.io/v1', - string $name = 'minimax', + string $endpoint = self::DEFAULT_ENDPOINT, ): MiniMaxJobClient { - return new MiniMaxJobClient($httpClient ?? new EventSourceHttpClient(), $apiKey, $endpoint, $name); + return new MiniMaxJobClient($httpClient ?? new EventSourceHttpClient(), $apiKey, $endpoint); } /** @@ -73,7 +71,7 @@ public static function createJobClient( public static function createPlatform( #[\SensitiveParameter] string $apiKey, ?HttpClientInterface $httpClient = null, - string $endpoint = 'https://api.minimax.io/v1', + string $endpoint = self::DEFAULT_ENDPOINT, ModelCatalogInterface $modelCatalog = new ModelCatalog(), ?Contract $contract = null, ?EventDispatcherInterface $eventDispatcher = null, diff --git a/src/platform/src/Bridge/MiniMax/MiniMaxJobClient.php b/src/platform/src/Bridge/MiniMax/MiniMaxJobClient.php index f13d266d26..2eb4ff3f8b 100644 --- a/src/platform/src/Bridge/MiniMax/MiniMaxJobClient.php +++ b/src/platform/src/Bridge/MiniMax/MiniMaxJobClient.php @@ -28,8 +28,7 @@ * MiniMax answers such a request with a `task_id`, exposes the task under an endpoint-specific query * path, and delivers the payload as a file that has to be looked up and downloaded separately. Both * the query path and the expected MIME type are carried in the {@see JobHandle}, put there by - * {@see MiniMaxResultConverter} which knows the endpoint the task came from. It creates the handle - * through this client, which names the provider it serves. + * {@see MiniMaxResultConverter} which knows the endpoint the task came from. * * @author Johannes Wachter */ @@ -55,19 +54,9 @@ public function __construct( private readonly HttpClientInterface $httpClient, #[\SensitiveParameter] private readonly string $apiKey, private readonly string $endpoint = 'https://api.minimax.io/v1', - private readonly string $provider = 'minimax', ) { } - /** - * @param array $data - * @param int $maxDuration how long the task may reasonably take, in seconds - */ - public function createHandle(string $taskId, array $data, int $maxDuration): JobHandle - { - return new JobHandle($taskId, $data, $this->provider, $maxDuration); - } - public function supports(JobHandle $handle): bool { return \is_string($handle->get('query_path')); @@ -75,25 +64,16 @@ public function supports(JobHandle $handle): bool public function getStatus(JobHandle $handle): JobStatus { - $data = $this->query($handle); - - $raw = (string) ($data['status'] ?? ''); - $case = self::STATES[strtolower($raw)] ?? JobStateCase::UNKNOWN; - - $error = $data['base_resp']['status_msg'] ?? null; - - return new JobStatus($case, $raw, \is_string($error) && '' !== $error ? $error : null); + return $this->toStatus($this->query($handle)); } public function getResult(JobHandle $handle): ResultInterface { $data = $this->query($handle); + $status = $this->toStatus($data); - $raw = (string) ($data['status'] ?? ''); - $case = self::STATES[strtolower($raw)] ?? JobStateCase::UNKNOWN; - - if (JobStateCase::SUCCEEDED !== $case) { - throw new JobFailedException(new JobStatus($case, $raw), \sprintf('The MiniMax task "%s" is not ready to be fetched, its status is "%s".', $handle->getId(), $raw)); + if (!$status->is(JobStateCase::SUCCEEDED)) { + throw new JobFailedException($status, \sprintf('The MiniMax task "%s" is not ready to be fetched, its status is "%s".', $handle->getId(), $status->getRaw())); } // The file identifier can already be known from the submit response; the query response wins @@ -116,6 +96,33 @@ public function getResult(JobHandle $handle): ResultInterface return new BinaryResult($payload, \is_string($mimeType) ? $mimeType : null); } + /** + * `status_msg` is a failure message only when `status_code` is not 0 - a healthy one says "success". + * + * @param array $data + */ + private function toStatus(array $data): JobStatus + { + $raw = (string) ($data['status'] ?? ''); + $case = self::STATES[strtolower($raw)] ?? JobStateCase::UNKNOWN; + + $statusCode = $data['base_resp']['status_code'] ?? 0; + + if (0 === $statusCode) { + return new JobStatus($case, $raw); + } + + $message = $data['base_resp']['status_msg'] ?? null; + + // No state at all next to an error code is a task MiniMax will not talk about. + if ('' === $raw) { + $case = JobStateCase::FAILED; + $raw = 'status code '.(\is_scalar($statusCode) ? (string) $statusCode : 'unknown'); + } + + return new JobStatus($case, $raw, \is_string($message) && '' !== $message ? $message : null); + } + /** * @return array */ diff --git a/src/platform/src/Bridge/MiniMax/MiniMaxResultConverter.php b/src/platform/src/Bridge/MiniMax/MiniMaxResultConverter.php index e147b10a6b..6abe222a22 100644 --- a/src/platform/src/Bridge/MiniMax/MiniMaxResultConverter.php +++ b/src/platform/src/Bridge/MiniMax/MiniMaxResultConverter.php @@ -14,6 +14,7 @@ use Symfony\AI\Platform\Exception\IncompleteStreamException; use Symfony\AI\Platform\Exception\RuntimeException; use Symfony\AI\Platform\FinishReason\FinishReasonAwareTrait; +use Symfony\AI\Platform\Job\JobHandle; use Symfony\AI\Platform\Model; use Symfony\AI\Platform\Result\BinaryResult; use Symfony\AI\Platform\Result\ChoiceResult; @@ -28,7 +29,6 @@ use Symfony\AI\Platform\Result\TextResult; use Symfony\AI\Platform\ResultConverterInterface; use Symfony\AI\Platform\TokenUsage\TokenUsageExtractorInterface; -use Symfony\Component\HttpClient\EventSourceHttpClient; /** * @author Guillaume Loulier @@ -47,16 +47,12 @@ final class MiniMaxResultConverter implements ResultConverterInterface private const VIDEO_MAX_DURATION = 600; - private readonly MiniMaxJobClient $jobClient; - /** - * @param MiniMaxJobClient|null $jobClient creates the handles of the jobs this converter starts, so - * they name the provider the client serves + * @param string $provider the name stamped onto the handles of the jobs this converter starts */ - public function __construct(?MiniMaxJobClient $jobClient = null) - { - // Only used to create handles, never to send a request. - $this->jobClient = $jobClient ?? new MiniMaxJobClient(new EventSourceHttpClient(), ''); + public function __construct( + private readonly string $provider = 'minimax', + ) { } public function supports(Model $model): bool @@ -206,8 +202,8 @@ private function decodeHexAudio(array $data): string /** * MiniMax answered with a task identifier instead of a payload, so the invocation produces a * reference to that task rather than a result. Resolving it - polling, and downloading the file - * it produces - is the job of {@see MiniMaxJobClient}, which therefore creates the handle; the - * handle carries what that client needs to know about the endpoint the task came from. + * it produces - is the job of {@see MiniMaxJobClient}; the handle carries what that client needs + * to know about the endpoint the task came from. * * @param array $data * @param int $maxDuration how long this endpoint may reasonably take, in seconds @@ -218,11 +214,11 @@ private function startJob(array $data, string $queryPath, string $mimeType, int { $taskId = $data['task_id'] ?? throw new RuntimeException('The MiniMax response does not contain a task identifier.'); - return new JobResult($this->jobClient->createHandle((string) $taskId, [ + return new JobResult(new JobHandle((string) $taskId, [ 'query_path' => $queryPath, 'mime_type' => $mimeType, 'archive_member' => $archiveMember, 'file_id' => $data['file_id'] ?? null, - ], $maxDuration)); + ], $this->provider, $maxDuration)); } } diff --git a/src/platform/src/Bridge/MiniMax/Tests/FactoryTest.php b/src/platform/src/Bridge/MiniMax/Tests/FactoryTest.php index c1e7d094f0..9726e109ce 100644 --- a/src/platform/src/Bridge/MiniMax/Tests/FactoryTest.php +++ b/src/platform/src/Bridge/MiniMax/Tests/FactoryTest.php @@ -32,9 +32,4 @@ public function testTheJobHandleCarriesTheNameTheProviderWasCreatedWith() $this->assertSame('789', $handle->getId()); $this->assertSame('minimax-eu', $handle->getProvider()); } - - public function testTheJobClientCreatesHandlesForTheGivenName() - { - $this->assertSame('minimax-eu', Factory::createJobClient('key', name: 'minimax-eu')->createHandle('789', [], 600)->getProvider()); - } } diff --git a/src/platform/src/Bridge/MiniMax/Tests/MiniMaxJobClientTest.php b/src/platform/src/Bridge/MiniMax/Tests/MiniMaxJobClientTest.php index 0fc583fdbd..a5796fe319 100644 --- a/src/platform/src/Bridge/MiniMax/Tests/MiniMaxJobClientTest.php +++ b/src/platform/src/Bridge/MiniMax/Tests/MiniMaxJobClientTest.php @@ -31,19 +31,6 @@ */ final class MiniMaxJobClientTest extends TestCase { - public function testItCreatesHandlesForTheProviderItServes() - { - $handle = (new MiniMaxJobClient(new MockHttpClient(), 'key', provider: 'minimax-eu')) - ->createHandle('123', ['query_path' => 'query/video_generation'], 600); - - $this->assertSame('123', $handle->getId()); - $this->assertSame('minimax-eu', $handle->getProvider()); - $this->assertSame('query/video_generation', $handle->get('query_path')); - $this->assertSame(600, $handle->getMaxDuration()); - - $this->assertSame('minimax', (new MiniMaxJobClient(new MockHttpClient(), 'key'))->createHandle('123', [], 600)->getProvider()); - } - public function testItOnlySupportsHandlesCarryingAQueryPath() { $jobClient = new MiniMaxJobClient(new MockHttpClient(), 'key'); @@ -91,6 +78,60 @@ public function testItExposesTheProvidersFailureMessage() $this->assertSame('invalid params', $status->getError()); } + public function testItReportsNoErrorWhileTheJobIsHealthy() + { + $httpClient = new MockHttpClient(new JsonMockResponse([ + 'status' => 'Processing', + 'base_resp' => ['status_code' => 0, 'status_msg' => 'success'], + ])); + + $status = (new MiniMaxJobClient($httpClient, 'key'))->getStatus($this->handle()); + + $this->assertTrue($status->is(JobStateCase::RUNNING)); + $this->assertNull($status->getError()); + } + + public function testARejectedQueryEndsTheJobInsteadOfLookingLikeAnUnknownState() + { + $httpClient = new MockHttpClient(new JsonMockResponse([ + 'base_resp' => ['status_code' => 2013, 'status_msg' => 'invalid params'], + ])); + + $status = (new MiniMaxJobClient($httpClient, 'key'))->getStatus($this->handle()); + + $this->assertTrue($status->is(JobStateCase::FAILED)); + $this->assertTrue($status->isTerminal()); + $this->assertSame('status code 2013', $status->getRaw()); + $this->assertSame('invalid params', $status->getError()); + } + + public function testAStateThisBridgeDoesNotKnowStaysNonTerminal() + { + $httpClient = new MockHttpClient(new JsonMockResponse([ + 'status' => 'Rescheduled', + 'base_resp' => ['status_code' => 2013, 'status_msg' => 'invalid params'], + ])); + + $status = (new MiniMaxJobClient($httpClient, 'key'))->getStatus($this->handle()); + + $this->assertTrue($status->is(JobStateCase::UNKNOWN)); + $this->assertFalse($status->isTerminal()); + $this->assertSame('Rescheduled', $status->getRaw()); + $this->assertSame('invalid params', $status->getError()); + } + + public function testARejectedQueryFailsTheRunnerRightAwayInsteadOfSpendingTheBudget() + { + $httpClient = new MockHttpClient(new JsonMockResponse([ + 'base_resp' => ['status_code' => 1008, 'status_msg' => 'insufficient balance'], + ])); + + $this->expectException(JobFailedException::class); + $this->expectExceptionMessage('insufficient balance'); + + (new JobRunner(new MockClock()))->wait(new MiniMaxJobClient($httpClient, 'key'), $this->handle()); + } + public function testItLooksUpTheFileAndDownloadsIt() { $httpClient = new MockHttpClient([ diff --git a/src/platform/src/Job/JobClientInterface.php b/src/platform/src/Job/JobClientInterface.php index 9b1af0a725..f0cbd5b024 100644 --- a/src/platform/src/Job/JobClientInterface.php +++ b/src/platform/src/Job/JobClientInterface.php @@ -17,13 +17,16 @@ /** * Resolves asynchronous jobs previously started through `Platform::invoke()`. * - * Implementations perform exactly one request per call and never sleep: how often a job is polled, - * and for how long, is the caller's decision - see {@see JobRunner} for the blocking variant. + * Implementations never sleep: how often a job is polled, and for how long, is the caller's + * decision - see {@see JobRunner} for the blocking variant. * * @author Johannes Wachter */ interface JobClientInterface { + /** + * Whether this client can resolve the handle; a {@see JobRunner} refuses one it declines. + */ public function supports(JobHandle $handle): bool; /** @@ -34,7 +37,7 @@ public function supports(JobHandle $handle): bool; public function getStatus(JobHandle $handle): JobStatus; /** - * Fetches the finished job's result. + * Fetches the finished job's result, in as many requests as the provider needs to hand it over. * * Only meaningful once {@see getStatus()} reported {@see JobStateCase::SUCCEEDED}; implementations * throw when the job did not finish successfully. diff --git a/src/platform/src/Job/JobRunner.php b/src/platform/src/Job/JobRunner.php index 6a2f26aa1c..8b479f37e8 100644 --- a/src/platform/src/Job/JobRunner.php +++ b/src/platform/src/Job/JobRunner.php @@ -49,6 +49,11 @@ final class JobRunner */ private const DEFAULT_MAX_DURATION = 120; + /** + * Below the clock's own resolution, what is left of the budget is rounding. + */ + private const CLOCK_RESOLUTION = 0.000001; + /** * @param float $pollInterval seconds to wait between two polls * @param int|null $maxDuration seconds to wait before giving up, for every job this runner @@ -72,17 +77,22 @@ public function __construct( * @param int|null $maxDuration seconds to wait for this job, overruling both the runner's own * budget and what the job asks for * - * @throws JobFailedException when the job reached a terminal state without a result - * @throws JobTimeoutException when the job was still running after the last poll + * @throws InvalidArgumentException when the job client cannot resolve this handle + * @throws JobFailedException when the job reached a terminal state without a result + * @throws JobTimeoutException when the job was still running after the last poll */ public function wait(JobClientInterface $jobClient, JobHandle $handle, ?int $maxDuration = null): DeferredResult { self::assertDuration($maxDuration); + if (!$jobClient->supports($handle)) { + throw new InvalidArgumentException(\sprintf('The job "%s" of provider "%s" cannot be resolved by "%s".', $handle->getId(), $handle->getProvider() ?? 'unknown', $jobClient::class)); + } + $budget = $maxDuration ?? $this->maxDuration ?? $handle->getMaxDuration() ?? self::DEFAULT_MAX_DURATION; - $maxPolls = $this->maxPollsFor($budget); + $deadline = $this->now() + $budget; - for ($poll = 1; $poll <= $maxPolls; ++$poll) { + while (true) { $status = $jobClient->getStatus($handle); if ($status->is(JobStateCase::SUCCEEDED)) { @@ -95,21 +105,20 @@ public function wait(JobClientInterface $jobClient, JobHandle $handle, ?int $max throw new JobFailedException($status, \sprintf('The job "%s" ended as "%s".%s', $handle->getId(), $status->getRaw(), null !== $status->getError() ? ' '.$status->getError() : '')); } - // Not after the last poll: the caller would only wait for a status nobody reads. - if ($poll < $maxPolls) { - $this->clock->sleep($this->pollInterval); + // Sleeping past the deadline would only wait for a status nobody reads. + if ($this->now() + $this->pollInterval + self::CLOCK_RESOLUTION >= $deadline) { + break; } + + $this->clock->sleep($this->pollInterval); } throw new JobTimeoutException($handle, \sprintf('The job "%s" did not finish within %d second(s). It may still be running - keep the handle and wait for it again later, or allow more time via the "maxDuration" argument.', $handle->getId(), $budget)); } - /** - * Turns a duration into a number of polls at the configured interval. - */ - private function maxPollsFor(int $maxDuration): int + private function now(): float { - return max(1, (int) ceil($maxDuration / $this->pollInterval)); + return (float) $this->clock->now()->format('U.u'); } private static function assertDuration(?int $maxDuration): void diff --git a/src/platform/tests/Fixtures/Job/ScriptedJobClient.php b/src/platform/tests/Fixtures/Job/ScriptedJobClient.php index fdc51504b7..1066d06c26 100644 --- a/src/platform/tests/Fixtures/Job/ScriptedJobClient.php +++ b/src/platform/tests/Fixtures/Job/ScriptedJobClient.php @@ -30,6 +30,15 @@ final class ScriptedJobClient implements JobClientInterface public int $resultCalls = 0; + public bool $supports = true; + + /** + * Lets a test spend time inside a poll. + * + * @var (\Closure(): void)|null + */ + public ?\Closure $onStatus = null; + /** * @var list */ @@ -42,13 +51,17 @@ public function __construct(JobStatus ...$statuses) public function supports(JobHandle $handle): bool { - return true; + return $this->supports; } public function getStatus(JobHandle $handle): JobStatus { ++$this->statusCalls; + if (null !== $this->onStatus) { + ($this->onStatus)(); + } + return array_shift($this->statuses) ?? throw new LogicException('The runner polled more often than the test scripted.'); } diff --git a/src/platform/tests/Job/JobRunnerTest.php b/src/platform/tests/Job/JobRunnerTest.php index 996628c9de..3f412ffe0b 100644 --- a/src/platform/tests/Job/JobRunnerTest.php +++ b/src/platform/tests/Job/JobRunnerTest.php @@ -212,6 +212,36 @@ public function testItGivesUpAfterTheLastPollAndHandsTheHandleBack() $this->assertSame(0, $jobClient->resultCalls); } + public function testTheBudgetCountsTheTimeTheProviderTakesToAnswer() + { + $clock = new MockClock('2026-01-01 00:00:00'); + + // Every poll takes two seconds at the provider, on top of the one second between polls. + $jobClient = new ScriptedJobClient(...array_fill(0, 10, new JobStatus(JobStateCase::RUNNING, 'Processing'))); + $jobClient->onStatus = static fn () => $clock->sleep(2.0); + + try { + (new JobRunner($clock, 1.0))->wait($jobClient, new JobHandle('task-1'), maxDuration: 9); + $this->fail(\sprintf('Expected a "%s".', JobTimeoutException::class)); + } catch (JobTimeoutException) { + } + + // Three polls at three seconds each, not nine at one second between them. + $this->assertSame(3, $jobClient->statusCalls); + $this->assertSame('2026-01-01 00:00:08', $clock->now()->format('Y-m-d H:i:s')); + } + + public function testItRefusesAHandleTheJobClientCannotResolve() + { + $jobClient = $this->jobClient(); + $jobClient->supports = false; + + $this->expectException(InvalidArgumentException::class); + $this->expectExceptionMessage('cannot be resolved by'); + + (new JobRunner(new MockClock()))->wait($jobClient, new JobHandle('task-1', [], 'other-provider')); + } + public function testItRejectsANonsensicalPollInterval() { $this->expectException(InvalidArgumentException::class);