diff --git a/CHANGELOG.md b/CHANGELOG.md index ff166885..aa6404f2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -12,6 +12,25 @@ FIXED - Asynchronous Azure Blob payload uploads and downloads no longer run gzip compression or decompression on the calling event loop, keeping concurrent async operations responsive during large transfers. +- Fixed nested `when_all` and `when_any` tasks leaving their enclosing composite +tasks waiting after completion. Patterns such as +`when_any([cancel, when_all(tasks)])` now complete correctly. + +> [!WARNING] +> **Replay-breaking bug fix:** This correction can change the winner of a nested +> race when replaying an existing orchestration. For example, an inner +> `when_all` could previously finish without completing the enclosing +> `when_any`, allowing a later cancellation to win instead. If the orchestration +> already recorded downstream actions for that cancellation branch, replay with +> the corrected behavior can take a different branch and fail with +> `NonDeterminismError`. Ordinary, non-nested `when_all` and `when_any` usage is +> unchanged; not every nested-composite history is affected. +> +> Before upgrading, allow affected instances to finish on the previous SDK, or +> keep their task hub on the previous SDK and use a separate task hub for new +> instances. Instances already stuck because of this bug may need deliberate +> recovery rather than waiting to drain. Pin or lock the `durabletask` dependency +> while planning the transition, including when using it through a provider. ## v1.10.1 diff --git a/azure-functions-durable/CHANGELOG.md b/azure-functions-durable/CHANGELOG.md index 206c9acc..b97e3da8 100644 --- a/azure-functions-durable/CHANGELOG.md +++ b/azure-functions-durable/CHANGELOG.md @@ -18,6 +18,25 @@ FIXED - Preserved the application's source directory in `context.function_directory` for decorated activities and durable-client functions. +- With a corrected core `durabletask` SDK, nested `when_all` and `when_any` +composites now complete correctly instead of leaving their enclosing composite +waiting. This also applies to the compatibility APIs `context.task_all()` and +`context.task_any()`. + +> [!WARNING] +> **Replay-breaking bug fix:** Existing instances that progressed past a nested +> race under the previous behavior can select a different winner during replay +> and fail with `NonDeterminismError`. Ordinary, non-nested composites are +> unaffected. This warning also applies when upgrading from earlier +> `azure-functions-durable` 2.x prereleases, including release candidates. See +> the [core changelog](../CHANGELOG.md) for the affected scenario and guidance on +> draining or isolating existing instances before upgrading. +> +> Pinning `azure-functions-durable` v2 alone does not pin the core SDK. Its +> open-ended `durabletask[opentelemetry]` dependency allows a fresh dependency +> resolution to select a newer core SDK containing this fix even when the +> provider version is unchanged. Pin or lock `durabletask` as well to control +> when this behavior change is adopted. ## v2.0.0rc1 diff --git a/durabletask-azuremanaged/CHANGELOG.md b/durabletask-azuremanaged/CHANGELOG.md index 895b254a..38823026 100644 --- a/durabletask-azuremanaged/CHANGELOG.md +++ b/durabletask-azuremanaged/CHANGELOG.md @@ -28,6 +28,22 @@ FIXED - With the corresponding core SDK update, asynchronous Azure Blob payload transfers no longer block the event loop during compression or decompression. +- Orchestrations using a corrected core `durabletask` SDK no longer remain +blocked when a nested `when_all` or `when_any` task completes, including patterns +such as `when_any([cancel, when_all(tasks)])`. + +> [!WARNING] +> **Replay-breaking bug fix:** Existing instances that progressed past a nested +> race under the previous behavior can select a different winner during replay +> and fail with `NonDeterminismError`. Ordinary, non-nested composites are +> unaffected. See the [core changelog](../CHANGELOG.md) for the affected scenario +> and guidance on draining or isolating existing instances before upgrading. +> +> Pinning `durabletask.azuremanaged` alone does not pin the core SDK. Its +> open-ended `durabletask` dependency allows a fresh dependency resolution to +> select a newer core SDK containing this fix even when the provider version is +> unchanged. Pin or lock `durabletask` as well to control when this behavior +> change is adopted. ## v1.10.1 diff --git a/durabletask/task.py b/durabletask/task.py index d8235c47..68ae0944 100644 --- a/durabletask/task.py +++ b/durabletask/task.py @@ -643,6 +643,10 @@ def get_exception(self) -> TaskFailedError: raise ValueError('The task has not failed.') return self._exception + def _notify_parent(self) -> None: + if self._parent is not None: + self._parent.on_child_completed(self) + class CompositeTask(Task[T]): """A task that is composed of other tasks.""" @@ -710,6 +714,7 @@ def on_child_completed(self, task: Task[Any]) -> None: # The order of the result MUST match the order of the tasks # provided to the constructor. self._result = [child.get_result() for child in self._tasks] + self._notify_parent() def get_completed_tasks(self) -> int: return self._completed_tasks @@ -727,16 +732,14 @@ def complete(self, result: T): raise ValueError('The task has already completed.') self._result = result self._is_complete = True - if self._parent is not None: - self._parent.on_child_completed(self) + self._notify_parent() def fail(self, message: str, details: Exception | pb.TaskFailureDetails): if self._is_complete: raise ValueError('The task has already completed.') self._exception = TaskFailedError(message, details) self._is_complete = True - if self._parent is not None: - self._parent.on_child_completed(self) + self._notify_parent() class CancellableTask(CompletableTask[T]): @@ -776,8 +779,7 @@ def cancel(self) -> bool: self._is_cancelled = True self._is_complete = True - if self._parent is not None: - self._parent.on_child_completed(self) + self._notify_parent() return True @@ -859,6 +861,7 @@ def on_child_completed(self, task: Task[Any]) -> None: if not self.is_complete: self._is_complete = True self._result = cast(Task[T], task) + self._notify_parent() def when_all(tasks: list[Task[T]]) -> WhenAllTask[T]: diff --git a/tests/durabletask/test_orchestration_executor.py b/tests/durabletask/test_orchestration_executor.py index a7064145..4c213403 100644 --- a/tests/durabletask/test_orchestration_executor.py +++ b/tests/durabletask/test_orchestration_executor.py @@ -2051,6 +2051,74 @@ def test_when_all_handles_pre_completed_children(): assert when_all.is_failed +def test_when_any_completes_when_nested_when_all_succeeds(): + """A nested when_all task must notify its when_any parent.""" + cancel = task.CompletableTask() + first = task.CompletableTask() + second = task.CompletableTask() + all_task = task.when_all([first, second]) + race = task.when_any([cancel, all_task]) + + first.complete("first") + assert not all_task.is_complete + assert not race.is_complete + + second.complete("second") + + assert all_task.is_complete + assert all_task.get_result() == ["first", "second"] + assert race.is_complete + assert race.get_result() is all_task + + +def test_when_any_completes_when_nested_when_all_fails(): + """A failed nested when_all task must notify its when_any parent.""" + cancel = task.CompletableTask() + failed = task.CompletableTask() + completed = task.CompletableTask() + all_task = task.when_all([failed, completed]) + race = task.when_any([cancel, all_task]) + + failed.fail("boom", Exception("boom")) + assert not all_task.is_complete + assert not race.is_complete + + completed.complete("done") + + assert all_task.is_complete + assert all_task.is_failed + assert race.is_complete + assert race.get_result() is all_task + with pytest.raises(task.TaskFailedError, match="boom"): + all_task.get_result() + + +def test_nested_when_any_notifies_parent_once(): + """A nested when_any task must notify its parent only for its winner.""" + first = task.CompletableTask() + second = task.CompletableTask() + sibling = task.CompletableTask() + nested_race = task.when_any([first, second]) + all_task = task.when_all([nested_race, sibling]) + + first.complete("first") + + assert nested_race.is_complete + assert nested_race.get_result() is first + assert all_task.get_completed_tasks() == 1 + assert not all_task.is_complete + + second.complete("second") + + assert all_task.get_completed_tasks() == 1 + assert not all_task.is_complete + + sibling.complete("sibling") + + assert all_task.is_complete + assert all_task.get_result() == [first, "sibling"] + + def test_when_any(): """Tests that a when_any pattern works correctly""" def hello(_, name: str):