diff --git a/CHANGELOG.md b/CHANGELOG.md index ac2395d2..50d87c8a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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. diff --git a/docs/02-task_types.md b/docs/02-task_types.md index 9701f680..fc40b2c4 100644 --- a/docs/02-task_types.md +++ b/docs/02-task_types.md @@ -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 diff --git a/docs/03-custom_tasks.md b/docs/03-custom_tasks.md index 1020c0cd..f0cacf9d 100644 --- a/docs/03-custom_tasks.md +++ b/docs/03-custom_tasks.md @@ -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 diff --git a/docs/reference/tasks/counter_task.md b/docs/reference/tasks/counter_task.md index 9cc9a5af..79067612 100644 --- a/docs/reference/tasks/counter_task.md +++ b/docs/reference/tasks/counter_task.md @@ -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. diff --git a/src/Model/FlushableTaskInterface.php b/src/Model/FlushableTaskInterface.php index 8f56f5fa..1bb4cadc 100644 --- a/src/Model/FlushableTaskInterface.php +++ b/src/Model/FlushableTaskInterface.php @@ -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 { diff --git a/src/Task/CounterTask.php b/src/Task/CounterTask.php index 9b2020e5..12e2a650 100644 --- a/src/Task/CounterTask.php +++ b/src/Task/CounterTask.php @@ -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; @@ -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 diff --git a/tests/Task/CounterTaskTest.php b/tests/Task/CounterTaskTest.php new file mode 100644 index 00000000..7e71d283 --- /dev/null +++ b/tests/Task/CounterTaskTest.php @@ -0,0 +1,118 @@ +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 + */ + 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(); + } +}