From 6e081c723670f08ac3d74269b4e3f02a3d2ef548 Mon Sep 17 00:00:00 2001 From: woutdenolf Date: Fri, 10 Jul 2026 16:17:32 +0200 Subject: [PATCH] InputMergeActor: fix deadlock --- CHANGELOG.md | 4 ++ src/ewoksppf/bindings.py | 3 +- src/ewoksppf/tests/test_input_merge_actor.py | 38 ++++++++++++ .../test_ppf_actors/pythonActorSelfLoop.py | 3 + src/ewoksppf/tests/test_ppf_workflow26.py | 61 +++++++++++++++++++ 5 files changed, 108 insertions(+), 1 deletion(-) create mode 100644 src/ewoksppf/tests/test_input_merge_actor.py create mode 100644 src/ewoksppf/tests/test_ppf_actors/pythonActorSelfLoop.py create mode 100644 src/ewoksppf/tests/test_ppf_workflow26.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 4e2b8bd..e36b762 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- `InputMergeActor`: possible deadlock for trigger loopback from a downstream node. + ## [3.0.0] - 2026-07-01 ### Added diff --git a/src/ewoksppf/bindings.py b/src/ewoksppf/bindings.py index a6f454d..d43394d 100644 --- a/src/ewoksppf/bindings.py +++ b/src/ewoksppf/bindings.py @@ -242,7 +242,8 @@ def __init__(self, parent=None, name="Input merger", **kw): # after all required triggers arrived self._retained_optional_trigger: Optional[dict] = None - self._lock = threading.Lock() + # Re-entrant in case downstream actor triggers this actor in the same call stack. + self._lock = threading.RLock() def register_input_actor(self, actor: Optional[AbstractActor]): if actor.required: diff --git a/src/ewoksppf/tests/test_input_merge_actor.py b/src/ewoksppf/tests/test_input_merge_actor.py new file mode 100644 index 0000000..5ab950e --- /dev/null +++ b/src/ewoksppf/tests/test_input_merge_actor.py @@ -0,0 +1,38 @@ +import threading + +from pypushflow.AbstractActor import AbstractActor +from pypushflow.ThreadCounter import ThreadCounter + +from ..bindings import InputMergeActor + + +def test_input_merge_actor_reentrant_trigger(): + """ + Deterministic counterpart to test_ppf_workflow26's real + (but timing-dependent) graph reproduction of the same bug. + """ + thread_counter = ThreadCounter() + merger = InputMergeActor(thread_counter=thread_counter, name="merger") + looping_actor = _SelfLoopingActor(merger, thread_counter) + merger.connect(looping_actor) + + thread = threading.Thread(target=merger.trigger, args=({},), daemon=True) + thread.start() + thread.join(timeout=5) + + assert not thread.is_alive(), "InputMergeActor deadlocked on a reentrant trigger" + assert looping_actor.triggered == 2 + assert thread_counter.nthreads == 0 + + +class _SelfLoopingActor(AbstractActor): + def __init__(self, merger: InputMergeActor, thread_counter: ThreadCounter): + super().__init__(thread_counter=thread_counter, name="looping actor") + self.merger = merger + self.triggered = 0 + + def _execute(self, inData: dict, _scope_id=None) -> None: + self.triggered += 1 + if self.triggered == 1: + # call InputMergeActor in the current thread + self.merger.trigger(inData) diff --git a/src/ewoksppf/tests/test_ppf_actors/pythonActorSelfLoop.py b/src/ewoksppf/tests/test_ppf_actors/pythonActorSelfLoop.py new file mode 100644 index 0000000..19e042e --- /dev/null +++ b/src/ewoksppf/tests/test_ppf_actors/pythonActorSelfLoop.py @@ -0,0 +1,3 @@ +def run(index=0, limit=10, **kwargs): + index += 1 + return {"index": index, "has_data": index < limit} diff --git a/src/ewoksppf/tests/test_ppf_workflow26.py b/src/ewoksppf/tests/test_ppf_workflow26.py new file mode 100644 index 0000000..b17171c --- /dev/null +++ b/src/ewoksppf/tests/test_ppf_workflow26.py @@ -0,0 +1,61 @@ +import sys +import threading + +from ewoksppf import execute_graph + + +def workflow26(limit: int): + nodes = [ + { + "id": "loop", + "default_inputs": [ + {"name": "index", "value": 0}, + {"name": "limit", "value": limit}, + ], + "force_start_node": True, + "task_type": "ppfmethod", + "task_identifier": "ewoksppf.tests.test_ppf_actors.pythonActorSelfLoop.run", + }, + ] + links = [ + { + "source": "loop", + "target": "loop", + "conditions": [{"source_output": "has_data", "value": True}], + "map_all_data": True, + }, + ] + graph = {"graph": {"id": "workflow26"}, "links": links, "nodes": nodes} + expected_result = {"_ppfdict": {"index": limit, "limit": limit, "has_data": False}} + return graph, expected_result + + +def test_workflow26(ppf_log_config): + """Workflow that maximizes this race-condition in the workflow execution pool + + .. code-block:: python + + future = self._pool.submit(...) + future.add_done_callback(cb) # if worker already finished, `cb` runs in the current call stack + + See `test_input_merge_actor_reentrant_trigger` for deterministic counterpart. + """ + graph, _ = workflow26(limit=200) + + # Lower GIL switch interval to make it more likely that + # the submitted job finished before calling `future.add_done_callback`. + old_interval = sys.getswitchinterval() + sys.setswitchinterval(1e-6) + try: + thread = threading.Thread( + target=execute_graph, + args=(graph,), + kwargs={"pool_type": "thread"}, + daemon=True, + ) + thread.start() + thread.join(timeout=30) + finally: + sys.setswitchinterval(old_interval) + + assert not thread.is_alive(), "deadlocked InputMergeActor"