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 CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ Latest
* [#201](https://github.com/cleverage/process-bundle/issues/201) Fix minor defects: error messages of TransformerTrait, RulesTransformer and ExpressionLanguageMapTransformer, useless `setRequired()` in ImplodeTransformer and SprintfTransformer, stray namespace in TrimTransformer, wrong or missing docblocks. Add tests.
* [#220](https://github.com/cleverage/process-bundle/issues/220) Fix the stop error strategy: throw a `ProcessFailedException` (with the original exception as `previous`) instead of a `FatalError`, so that the command exits with a non-zero code when a process fails and `ProcessLauncherTask` detects failed sub-processes. Update documentation, add tests.
* [#221](https://github.com/cleverage/process-bundle/issues/221) Fix CsvSplitterTask: each produced file contains `max_lines` data lines (instead of `max_lines - 2`), no infinite loop with `max_lines` <= 2 (`max_lines` must now be an integer greater than 0), no header-only file at the end. Update documentation, add tests.
* [#223](https://github.com/cleverage/process-bundle/issues/223) Fix CounterTask: the final count is outputted once, as `flush()` may be called several times. Document that `flush()` implementations must be idempotent. Update documentation, add tests.

## Deprecated
* [#189](https://github.com/cleverage/process-bundle/issues/189) EventDispatcherTask: when `event_name` is set, listening to `CleverAge\ProcessBundle\Event\EventDispatcherTaskEvent` is deprecated (the event is still dispatched under its class name, with an `E_USER_DEPRECATED` error, if it has listeners). Listen to the configured `event_name` instead: the BC layer will be removed in v6.0.
Expand Down
3 changes: 1 addition & 2 deletions docs/02-task_types.md
Original file line number Diff line number Diff line change
Expand Up @@ -102,8 +102,7 @@ Examples:
- [SimpleBatchTask](reference/tasks/simple_batch_task.md) groups inputs by batches of `batch_count` elements: each
full batch is outputted during `execute`, and the last incomplete batch during `flush`.
- [CounterTask](reference/tasks/counter_task.md) outputs the count every `flush_every` items, and the current count on
`flush` (unless it is a multiple of `flush_every`): as `flush` may be called several times, the same count can be
outputted more than once.
`flush` (unless it is a multiple of `flush_every`, or it was already outputted by a previous `flush`).

## Initializable tasks

Expand Down
4 changes: 3 additions & 1 deletion docs/03-custom_tasks.md
Original file line number Diff line number Diff line change
Expand Up @@ -173,7 +173,9 @@ memory footprint.
[notice about blocking tasks](02-task_types.md#blocking-tasks) about memory usage.

Tasks should not be both Iterable and Blocking. If you need to buffer data and output it by chunks, look at
`CleverAge\ProcessBundle\Model\FlushableTaskInterface` (see [flushable tasks](02-task_types.md#flushable-tasks)).
`CleverAge\ProcessBundle\Model\FlushableTaskInterface` (see [flushable tasks](02-task_types.md#flushable-tasks)). As
`flush` may be called several times on the same task, it must be idempotent: once the buffer is flushed, a new call
must skip the state (`ProcessState::setSkipped(true)`) instead of outputting the same data again.

## Transformers

Expand Down
3 changes: 1 addition & 2 deletions docs/reference/tasks/counter_task.md
Original file line number Diff line number Diff line change
Expand Up @@ -53,5 +53,4 @@ Notes

* `flush()` can be called several times during a process (see
[Advanced workflow](../../04-advanced_workflow.md)), e.g. once at the end of the upstream iteration and once when
the counter itself is resolved. Each call outputs the current count again (unless it is a multiple of `flush_every`),
so the final count may be sent more than once to the next tasks.
the counter itself is resolved. The final count is only outputted by the first call; the next ones are skipped.
4 changes: 4 additions & 0 deletions src/Model/FlushableTaskInterface.php
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,10 @@

/**
* When iterations are over, this allows task that have some inner buffer to flush it to the output.
*
* flush() may be called several times on the same task during a process (once for each resolved ancestor, and each time
* an upstream iterable task finishes its iterations): implementations must be idempotent, and skip the state
* (ProcessState::setSkipped(true)) when there is nothing new to output.
*/
interface FlushableTaskInterface extends TaskInterface
{
Expand Down
14 changes: 11 additions & 3 deletions src/Task/CounterTask.php
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,11 @@ class CounterTask extends AbstractConfigurableTask implements FlushableTaskInter
{
protected int $counter = 0;

/**
* Count already outputted by a flush, as flush() may be called several times.
*/
protected ?int $flushedCounter = null;

public function execute(ProcessState $state): void
{
++$this->counter;
Expand All @@ -43,11 +48,14 @@ public function execute(ProcessState $state): void
public function flush(ProcessState $state): void
{
$modulo = $this->getOption($state, 'flush_every');
if (0 === $this->counter % $modulo) {
if (0 === $this->counter % $modulo || $this->counter === $this->flushedCounter) {
$state->setSkipped(true);
} else {
$state->setOutput($this->counter);

return;
}

$this->flushedCounter = $this->counter;
$state->setOutput($this->counter);
}

protected function configureOptions(OptionsResolver $resolver): void
Expand Down
118 changes: 118 additions & 0 deletions tests/Task/CounterTaskTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,118 @@
<?php

declare(strict_types=1);

/*
* This file is part of the CleverAge/ProcessBundle package.
*
* Copyright (c) Clever-Age
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

namespace CleverAge\ProcessBundle\Tests\Task;

use CleverAge\ProcessBundle\Configuration\ProcessConfiguration;
use CleverAge\ProcessBundle\Configuration\TaskConfiguration;
use CleverAge\ProcessBundle\Context\ContextualOptionResolver;
use CleverAge\ProcessBundle\Model\AbstractConfigurableTask;
use CleverAge\ProcessBundle\Model\ProcessHistory;
use CleverAge\ProcessBundle\Model\ProcessState;
use CleverAge\ProcessBundle\Task\CounterTask;
use PHPUnit\Framework\TestCase;

#[\PHPUnit\Framework\Attributes\CoversClass(CounterTask::class)]
#[\PHPUnit\Framework\Attributes\UsesClass(AbstractConfigurableTask::class)]
#[\PHPUnit\Framework\Attributes\UsesClass(ProcessConfiguration::class)]
#[\PHPUnit\Framework\Attributes\UsesClass(TaskConfiguration::class)]
#[\PHPUnit\Framework\Attributes\UsesClass(ContextualOptionResolver::class)]
#[\PHPUnit\Framework\Attributes\UsesClass(ProcessHistory::class)]
#[\PHPUnit\Framework\Attributes\UsesClass(ProcessState::class)]
class CounterTaskTest extends TestCase
{
public function testOutputsTheCountEveryFlushEvery(): void
{
[$task, $state] = $this->createTask(2);

self::assertSame([null, 2, null, 4, null], $this->executeAll($task, $state, 5));
}

public function testFlushOutputsTheFinalCountOnlyOnce(): void
{
[$task, $state] = $this->createTask(2);
$this->executeAll($task, $state, 5);

self::assertSame(5, $this->flush($task, $state));
self::assertNull($this->flush($task, $state));
self::assertNull($this->flush($task, $state));
}

public function testFlushSkipsWhenTheCountWasAlreadyOutputted(): void
{
[$task, $state] = $this->createTask(2);
$this->executeAll($task, $state, 4);

self::assertNull($this->flush($task, $state));
}

public function testFlushSkipsWithoutExecution(): void
{
[$task, $state] = $this->createTask(2);

self::assertNull($this->flush($task, $state));
}

public function testFlushOutputsANewCountAfterNewExecutions(): void
{
[$task, $state] = $this->createTask(3);
$this->executeAll($task, $state, 2);
self::assertSame(2, $this->flush($task, $state));

$this->executeAll($task, $state, 2);
self::assertSame(4, $this->flush($task, $state));
self::assertNull($this->flush($task, $state));
}

/**
* @return array{CounterTask, ProcessState}
*/
private function createTask(int $flushEvery): array
{
$processConfiguration = new ProcessConfiguration('test', []);
$state = new ProcessState($processConfiguration, new ProcessHistory($processConfiguration));
$state->setContextualOptionResolver(new ContextualOptionResolver());
$state->setContext([]);
$state->setTaskConfiguration(new TaskConfiguration('counter', CounterTask::class, ['flush_every' => $flushEvery]));
$task = new CounterTask();
$task->initialize($state);

return [$task, $state];
}

/**
* Execute the task $count times, collecting outputs (null when skipped).
*
* @return list<mixed>
*/
private function executeAll(CounterTask $task, ProcessState $state, int $count): array
{
$outputs = [];
for ($i = 0; $i < $count; ++$i) {
$state->reset(false);
$state->setInput($i);
$task->execute($state);
$outputs[] = $state->isSkipped() ? null : $state->getOutput();
}

return $outputs;
}

private function flush(CounterTask $task, ProcessState $state): mixed
{
$state->reset(true);
$task->flush($state);

return $state->isSkipped() ? null : $state->getOutput();
}
}
Loading