From 319cf1ea5c17dbdaa8eab3ac3a987d019ef4b183 Mon Sep 17 00:00:00 2001 From: Henri Payno Date: Tue, 23 Jun 2026 11:09:48 +0200 Subject: [PATCH 1/9] ewoksorange.gui.owwidgets.threaded: move the task execution handling into a specific `MultiThreadedTaskExecutor` class. --- src/ewoksorange/bindings/taskexecuter.py | 1 + src/ewoksorange/gui/concurrency/threaded.py | 136 ++++++++++++++++++++ src/ewoksorange/gui/owwidgets/threaded.py | 115 +++++------------ src/ewoksorange/tests/test_task_executor.py | 43 +++++++ 4 files changed, 214 insertions(+), 81 deletions(-) diff --git a/src/ewoksorange/bindings/taskexecuter.py b/src/ewoksorange/bindings/taskexecuter.py index aa985e3e..4a36114c 100644 --- a/src/ewoksorange/bindings/taskexecuter.py +++ b/src/ewoksorange/bindings/taskexecuter.py @@ -1,6 +1,7 @@ import warnings from ..gui.concurrency.base import TaskExecutor # noqa F401 +from ..gui.concurrency.threaded import MultiThreadedTaskExecutor # noqa F401 from ..gui.concurrency.threaded import ThreadedTaskExecutor # noqa F401 warnings.warn( diff --git a/src/ewoksorange/gui/concurrency/threaded.py b/src/ewoksorange/gui/concurrency/threaded.py index 31b3bbf8..686c358f 100644 --- a/src/ewoksorange/gui/concurrency/threaded.py +++ b/src/ewoksorange/gui/concurrency/threaded.py @@ -1,6 +1,12 @@ +from dataclasses import dataclass +from typing import Callable +from typing import Dict +from typing import Iterable from typing import Optional +from AnyQt.QtCore import QObject from AnyQt.QtCore import QThread +from AnyQt.QtCore import pyqtSignal as Signal from .base import TaskExecutor @@ -29,3 +35,133 @@ def cancel_running_task(self): """ if self.current_task is not None: self.current_task.cancel() + + +@dataclass +class _TaskExecutorState: + task_executor: ThreadedTaskExecutor + callbacks: Iterable[Callable[[ThreadedTaskExecutor], None]] + started: bool = False + + +class MultiThreadedTaskExecutor(QObject): + """Create and execute each Ewoks task in its own dedicated thread.""" + + sigComputationStarted = Signal() + """Signal emitted when a computation is started""" + + sigComputationEnded = Signal() + """Signal emitted when a computation is ended""" + + def __init__(self, ewokstaskclass): + super().__init__() + self.__ewokstaskclass = ewokstaskclass + self.__task_executors: Dict[int, _TaskExecutorState] = dict() + self.__last_output_variables = dict() + self.__last_task_succeeded = None + self.__last_task_done = None + self.__last_task_exception = None + + def create_task( + self, + _callbacks: Iterable[Callable[[ThreadedTaskExecutor], None]] = tuple(), + log_missing_inputs: bool = False, + **kwargs, + ) -> None: + """Create the next task to be executed in a dedicated thread.""" + task_executor = ThreadedTaskExecutor(ewokstaskclass=self.__ewokstaskclass) + task_executor.create_task(log_missing_inputs=log_missing_inputs, **kwargs) + self.__add_task_executor(task_executor, _callbacks) + + def execute_task(self) -> None: + """Execute the task created by :meth:`create_task`.""" + state = self.__get_pending_state() + if state is None: + return + + task_executor = state.task_executor + state.started = True + + if task_executor.has_task: + task_executor.finished.connect(self.__process_ended) + self.sigComputationStarted.emit() + task_executor.start() + else: + task_executor.finished.emit() + + def __process_ended(self): + self.__process_ended_direct(self.sender()) + + def __process_ended_direct(self, task_executor: ThreadedTaskExecutor): + state = self.__task_executors.get(id(task_executor)) + if state is None: + return + + self.__last_output_variables = task_executor.output_variables + self.__last_task_succeeded = task_executor.succeeded + self.__last_task_done = task_executor.done + self.__last_task_exception = task_executor.exception + + try: + for callback in state.callbacks: + callback(task_executor) + self.sigComputationEnded.emit() + finally: + self.__remove_task_executor(task_executor) + + def stop(self, timeout: Optional[float] = None, wait: bool = False) -> None: + """Stop all tracked task threads.""" + for state in list(self.__task_executors.values()): + task_executor = state.task_executor + if task_executor.receivers(task_executor.finished) > 0: + task_executor.finished.disconnect(self.__process_ended) + task_executor.stop(timeout=timeout, wait=wait) + self.__task_executors.clear() + + def cancel_running_tasks(self, wait=True): + """Request cancellation of all running tasks.""" + for state in list(self.__task_executors.values()): + if not state.started: + continue + task_executor = state.task_executor + task_executor.cancel_running_task() + task_executor.stop(wait=wait) + + def __add_task_executor( + self, + task_executor: ThreadedTaskExecutor, + callbacks: Iterable[Callable[[ThreadedTaskExecutor], None]], + ) -> None: + self.__task_executors[id(task_executor)] = _TaskExecutorState( + task_executor=task_executor, + callbacks=tuple(callbacks), + ) + + def __remove_task_executor(self, task_executor: ThreadedTaskExecutor) -> None: + if task_executor is None: + return + if task_executor.receivers(task_executor.finished) > 0: + task_executor.finished.disconnect(self.__process_ended) + self.__task_executors.pop(id(task_executor), None) + + def __get_pending_state(self) -> Optional[_TaskExecutorState]: + for state in self.__task_executors.values(): + if not state.started: + return state + return None + + @property + def task_succeeded(self) -> Optional[bool]: + return self.__last_task_succeeded + + @property + def task_done(self) -> Optional[bool]: + return self.__last_task_done + + @property + def task_exception(self) -> Optional[Exception]: + return self.__last_task_exception + + @property + def output_variables(self) -> dict: + return self.__last_output_variables diff --git a/src/ewoksorange/gui/owwidgets/threaded.py b/src/ewoksorange/gui/owwidgets/threaded.py index e45493a2..82e42789 100644 --- a/src/ewoksorange/gui/owwidgets/threaded.py +++ b/src/ewoksorange/gui/owwidgets/threaded.py @@ -6,11 +6,10 @@ import logging from contextlib import contextmanager -from typing import Dict from typing import Optional -from typing import Tuple from ..concurrency.queued import TaskExecutorQueue +from ..concurrency.threaded import MultiThreadedTaskExecutor from ..concurrency.threaded import ThreadedTaskExecutor from ..qt_utils.progress import QProgress from .base import OWEwoksBaseWidget @@ -183,113 +182,67 @@ class OWEwoksWidgetOneThreadPerRun(_OWEwoksThreadedBaseWidget, **ow_build_opts): def __init__(self, *args, **kwargs): """ - Initialize per-run executor storage. + Initialize the multi-thread task executor. """ super().__init__(*args, **kwargs) - self.__task_executors: Dict[int, Tuple[ThreadedTaskExecutor, bool]] = dict() - self.__last_output_variables = dict() - self.__last_task_succeeded = None - self.__last_task_done = None - self.__last_task_exception = None + self.__task_executor = MultiThreadedTaskExecutor( + ewokstaskclass=self.ewokstaskclass + ) def _execute_ewoks_task(self, propagate: bool, log_missing_inputs: bool) -> None: """ - Create a fresh ThreadedTaskExecutor, register it, and start it if it has work. + Submit a task to the multi-thread executor. :param propagate: Whether to propagate outputs after execution. :param log_missing_inputs: Whether to log missing input warnings. """ - task_executor = ThreadedTaskExecutor(ewokstaskclass=self.ewokstaskclass) - task_executor.create_task( - log_missing_inputs=log_missing_inputs, **self._get_task_arguments() - ) - with self.__init_task_executor(task_executor, propagate): - if task_executor.has_task: - with self._ewoks_task_start_context(): - task_executor.start() - else: - task_executor.finished.emit() - - @contextmanager - def __init_task_executor(self, task_executor, propagate: bool): - """ - Register a task executor and connect its finished callback for safe cleanup. - - :param task_executor: The ThreadedTaskExecutor instance. - :param propagate: Propagate flag to store with the executor. - """ - task_executor.finished.connect(self._ewoks_task_finished_callback) - self.__add_task_executor(task_executor, propagate) - try: - yield - except Exception: - task_executor.finished.disconnect(self._ewoks_task_finished_callback) - self.__remove_task_executor(task_executor) - raise - - def __disconnect_all_task_executors(self): - """Disconnect all connected finished signals from tracked executors.""" - for task_executor, _ in self.__task_executors.values(): - if task_executor.receivers(task_executor.finished) > 0: - task_executor.finished.disconnect(self._ewoks_task_finished_callback) + with self._ewoks_task_start_context(): + self.__task_executor.create_task( + _callbacks=( + lambda task_executor: self._ewoks_task_finished_callback( + task_executor, propagate + ), + ), + log_missing_inputs=log_missing_inputs, + **self._get_task_arguments(), + ) + self.__task_executor.execute_task() - def _ewoks_task_finished_callback(self): + def _ewoks_task_finished_callback( + self, task_executor: ThreadedTaskExecutor, propagate: bool + ): """ Slot invoked when a per-run executor finishes; stores its outputs and optionally propagates. """ with self._ewoks_task_finished_context(): - task_executor = None - try: - task_executor = self.sender() - self.__last_output_variables = task_executor.output_variables - self.__last_task_succeeded = task_executor.succeeded - self.__last_task_done = task_executor.done - self.__last_task_exception = task_executor.exception - self.__post_task_exception = None - propagate = self.__is_task_executor_propagated(task_executor) - if propagate: - self.propagate_downstream(succeeded=task_executor.succeeded) - finally: - self.__remove_task_executor(task_executor) + self.__post_task_exception = None + if propagate: + self.propagate_downstream(succeeded=task_executor.succeeded) def _cleanup_task_executor(self): - """Disconnect and quit all tracked executors on widget cleanup.""" - self.__disconnect_all_task_executors() - for task_executor, _ in self.__task_executors.values(): - task_executor.quit() - self.__task_executors.clear() - - def __add_task_executor(self, task_executor, propagate: bool): - """Internal: register a new task executor with its propagate flag.""" - self.__task_executors[id(task_executor)] = task_executor, propagate - - def __remove_task_executor(self, task_executor: ThreadedTaskExecutor): - """Internal: unregister a task executor and disconnect its signals.""" - if task_executor is None: - return - if task_executor.receivers(task_executor.finished) > 0: - task_executor.finished.disconnect(self._ewoks_task_finished_callback) - self.__task_executors.pop(id(task_executor), None) - - def __is_task_executor_propagated(self, task_executor) -> bool: - """Return whether the given executor was registered to propagate.""" - return self.__task_executors.get(id(task_executor), (None, False))[1] + """Stop all tracked per-run executors on widget cleanup.""" + self.__task_executor.stop() + self.__task_executor = None @property def task_succeeded(self) -> Optional[bool]: - return self.__last_task_succeeded + return self.__task_executor.task_succeeded @property def task_done(self) -> Optional[bool]: - return self.__last_task_done + return self.__task_executor.task_done @property def task_exception(self) -> Optional[Exception]: - return self.__last_task_exception + return self.__task_executor.task_exception def _get_task_outputs(self) -> dict: """Return the last finished task's outputs.""" - return self.__last_output_variables + return self.__task_executor.output_variables + + def cancel_running_task(self): + """Request cancellation of all running per-run tasks.""" + self.__task_executor.cancel_running_tasks() class OWEwoksWidgetWithTaskStack(_OWEwoksThreadedBaseWidget, **ow_build_opts): diff --git a/src/ewoksorange/tests/test_task_executor.py b/src/ewoksorange/tests/test_task_executor.py index 11a02a51..9a6b34be 100644 --- a/src/ewoksorange/tests/test_task_executor.py +++ b/src/ewoksorange/tests/test_task_executor.py @@ -7,6 +7,7 @@ from ..gui.concurrency.base import TaskExecutor from ..gui.concurrency.queued import TaskExecutorQueue +from ..gui.concurrency.threaded import MultiThreadedTaskExecutor from ..gui.concurrency.threaded import ThreadedTaskExecutor from ..gui.qt_utils.app import QtEvent @@ -51,6 +52,48 @@ def finished_callback(): executor.finished.disconnect(finished_callback) +def test_multi_threaded_task_executor(qtapp): + class MyObject(QObject): + def __init__(self, expected_results): + self.results = None + self.expected_results = expected_results + self.finished = QtEvent() + + def finished_callback(self, task_executor): + self.results = { + k: v.value for k, v in task_executor.output_variables.items() + } + assert self.results == self.expected_results + self.finished.set() + + executor = MultiThreadedTaskExecutor(ewokstaskclass=SumTask) + objects = [ + MyObject({"result": 3}), + MyObject({"result": 7}), + MyObject({"result": 11}), + ] + inputs = [ + {"a": 1, "b": 2}, + {"a": 3, "b": 4}, + {"a": 5, "b": 6}, + ] + + for obj, input_values in zip(objects, inputs): + executor.create_task( + inputs=input_values, + _callbacks=(obj.finished_callback,), + ) + executor.execute_task() + + for obj in objects: + assert obj.finished.wait(timeout=3) + assert obj.results == obj.expected_results + + expected_results = [obj.expected_results for obj in objects] + results = {k: v.value for k, v in executor.output_variables.items()} + assert results in expected_results + + def test_threaded_task_executor_queue(qtapp): class MyObject(QObject): def __init__(self): From b1a89c72cddfc2355e2073a08f9a78df20d27ef8 Mon Sep 17 00:00:00 2001 From: Henri Payno Date: Tue, 23 Jun 2026 16:30:25 +0200 Subject: [PATCH 2/9] MultiThreadedTaskExecutor: make sure '__process_ended_direct' is called. --- src/ewoksorange/gui/concurrency/threaded.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/ewoksorange/gui/concurrency/threaded.py b/src/ewoksorange/gui/concurrency/threaded.py index 686c358f..d7a77b08 100644 --- a/src/ewoksorange/gui/concurrency/threaded.py +++ b/src/ewoksorange/gui/concurrency/threaded.py @@ -87,7 +87,7 @@ def execute_task(self) -> None: self.sigComputationStarted.emit() task_executor.start() else: - task_executor.finished.emit() + self.__process_ended_direct(task_executor) def __process_ended(self): self.__process_ended_direct(self.sender()) From 7b91506eb5acd5a916d78acca62281ce4d68f5b6 Mon Sep 17 00:00:00 2001 From: Henri Payno Date: Wed, 24 Jun 2026 13:38:45 +0200 Subject: [PATCH 3/9] Simplify `MultiThreadedTaskExecutor`: task creation is forced when executing the task. --- src/ewoksorange/gui/concurrency/threaded.py | 83 +++++++++++---------- src/ewoksorange/gui/owwidgets/threaded.py | 3 +- src/ewoksorange/tests/test_task_executor.py | 2 +- 3 files changed, 46 insertions(+), 42 deletions(-) diff --git a/src/ewoksorange/gui/concurrency/threaded.py b/src/ewoksorange/gui/concurrency/threaded.py index d7a77b08..a22b3956 100644 --- a/src/ewoksorange/gui/concurrency/threaded.py +++ b/src/ewoksorange/gui/concurrency/threaded.py @@ -1,7 +1,7 @@ from dataclasses import dataclass from typing import Callable -from typing import Dict from typing import Iterable +from typing import List from typing import Optional from AnyQt.QtCore import QObject @@ -39,8 +39,10 @@ def cancel_running_task(self): @dataclass class _TaskExecutorState: - task_executor: ThreadedTaskExecutor callbacks: Iterable[Callable[[ThreadedTaskExecutor], None]] + task_kwargs: dict + log_missing_inputs: bool = False + task_executor: Optional[ThreadedTaskExecutor] = None started: bool = False @@ -56,30 +58,33 @@ class MultiThreadedTaskExecutor(QObject): def __init__(self, ewokstaskclass): super().__init__() self.__ewokstaskclass = ewokstaskclass - self.__task_executors: Dict[int, _TaskExecutorState] = dict() + self.__task_executors: List[_TaskExecutorState] = [] self.__last_output_variables = dict() self.__last_task_succeeded = None self.__last_task_done = None self.__last_task_exception = None - def create_task( + def execute_task( self, _callbacks: Iterable[Callable[[ThreadedTaskExecutor], None]] = tuple(), log_missing_inputs: bool = False, **kwargs, - ) -> None: - """Create the next task to be executed in a dedicated thread.""" + ) -> Optional[ThreadedTaskExecutor]: + """Execute a prepared task, or directly create and execute a new one.""" task_executor = ThreadedTaskExecutor(ewokstaskclass=self.__ewokstaskclass) - task_executor.create_task(log_missing_inputs=log_missing_inputs, **kwargs) - self.__add_task_executor(task_executor, _callbacks) + task_executor.create_task( + log_missing_inputs=log_missing_inputs, + **kwargs, + ) - def execute_task(self) -> None: - """Execute the task created by :meth:`create_task`.""" - state = self.__get_pending_state() - if state is None: - return + state = _TaskExecutorState( + callbacks=tuple(_callbacks), + task_kwargs=kwargs, + log_missing_inputs=log_missing_inputs, + task_executor=task_executor, + ) + self.__add_task_executor(state) - task_executor = state.task_executor state.started = True if task_executor.has_task: @@ -89,11 +94,20 @@ def execute_task(self) -> None: else: self.__process_ended_direct(task_executor) + return task_executor + def __process_ended(self): self.__process_ended_direct(self.sender()) def __process_ended_direct(self, task_executor: ThreadedTaskExecutor): - state = self.__task_executors.get(id(task_executor)) + state = next( + ( + state + for state in self.__task_executors + if state.task_executor is task_executor + ), + None, + ) if state is None: return @@ -111,44 +125,35 @@ def __process_ended_direct(self, task_executor: ThreadedTaskExecutor): def stop(self, timeout: Optional[float] = None, wait: bool = False) -> None: """Stop all tracked task threads.""" - for state in list(self.__task_executors.values()): + for state in list(self.__task_executors): task_executor = state.task_executor - if task_executor.receivers(task_executor.finished) > 0: - task_executor.finished.disconnect(self.__process_ended) - task_executor.stop(timeout=timeout, wait=wait) + if task_executor is not None: + if task_executor.receivers(task_executor.finished) > 0: + task_executor.finished.disconnect(self.__process_ended) + task_executor.stop(timeout=timeout, wait=wait) self.__task_executors.clear() def cancel_running_tasks(self, wait=True): """Request cancellation of all running tasks.""" - for state in list(self.__task_executors.values()): - if not state.started: + for state in list(self.__task_executors): + if not state.started or state.task_executor is None: continue task_executor = state.task_executor task_executor.cancel_running_task() task_executor.stop(wait=wait) - def __add_task_executor( - self, - task_executor: ThreadedTaskExecutor, - callbacks: Iterable[Callable[[ThreadedTaskExecutor], None]], - ) -> None: - self.__task_executors[id(task_executor)] = _TaskExecutorState( - task_executor=task_executor, - callbacks=tuple(callbacks), - ) + def __add_task_executor(self, state: _TaskExecutorState) -> None: + self.__task_executors.append(state) def __remove_task_executor(self, task_executor: ThreadedTaskExecutor) -> None: if task_executor is None: return - if task_executor.receivers(task_executor.finished) > 0: - task_executor.finished.disconnect(self.__process_ended) - self.__task_executors.pop(id(task_executor), None) - - def __get_pending_state(self) -> Optional[_TaskExecutorState]: - for state in self.__task_executors.values(): - if not state.started: - return state - return None + for state in list(self.__task_executors): + if state.task_executor is task_executor: + if task_executor.receivers(task_executor.finished) > 0: + task_executor.finished.disconnect(self.__process_ended) + self.__task_executors.remove(state) + break @property def task_succeeded(self) -> Optional[bool]: diff --git a/src/ewoksorange/gui/owwidgets/threaded.py b/src/ewoksorange/gui/owwidgets/threaded.py index 82e42789..5e2a575b 100644 --- a/src/ewoksorange/gui/owwidgets/threaded.py +++ b/src/ewoksorange/gui/owwidgets/threaded.py @@ -197,7 +197,7 @@ def _execute_ewoks_task(self, propagate: bool, log_missing_inputs: bool) -> None :param log_missing_inputs: Whether to log missing input warnings. """ with self._ewoks_task_start_context(): - self.__task_executor.create_task( + self.__task_executor.execute_task( _callbacks=( lambda task_executor: self._ewoks_task_finished_callback( task_executor, propagate @@ -206,7 +206,6 @@ def _execute_ewoks_task(self, propagate: bool, log_missing_inputs: bool) -> None log_missing_inputs=log_missing_inputs, **self._get_task_arguments(), ) - self.__task_executor.execute_task() def _ewoks_task_finished_callback( self, task_executor: ThreadedTaskExecutor, propagate: bool diff --git a/src/ewoksorange/tests/test_task_executor.py b/src/ewoksorange/tests/test_task_executor.py index 9a6b34be..4e8880e0 100644 --- a/src/ewoksorange/tests/test_task_executor.py +++ b/src/ewoksorange/tests/test_task_executor.py @@ -79,7 +79,7 @@ def finished_callback(self, task_executor): ] for obj, input_values in zip(objects, inputs): - executor.create_task( + executor.execute_task( inputs=input_values, _callbacks=(obj.finished_callback,), ) From 9013b142dbce696bdb64eb952165942809625619 Mon Sep 17 00:00:00 2001 From: Henri Payno Date: Wed, 24 Jun 2026 13:42:14 +0200 Subject: [PATCH 4/9] ThreadedTaskExecutor: remove useless 'started' field --- src/ewoksorange/gui/concurrency/threaded.py | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/src/ewoksorange/gui/concurrency/threaded.py b/src/ewoksorange/gui/concurrency/threaded.py index a22b3956..4d5bce68 100644 --- a/src/ewoksorange/gui/concurrency/threaded.py +++ b/src/ewoksorange/gui/concurrency/threaded.py @@ -43,7 +43,6 @@ class _TaskExecutorState: task_kwargs: dict log_missing_inputs: bool = False task_executor: Optional[ThreadedTaskExecutor] = None - started: bool = False class MultiThreadedTaskExecutor(QObject): @@ -85,8 +84,6 @@ def execute_task( ) self.__add_task_executor(state) - state.started = True - if task_executor.has_task: task_executor.finished.connect(self.__process_ended) self.sigComputationStarted.emit() @@ -136,7 +133,7 @@ def stop(self, timeout: Optional[float] = None, wait: bool = False) -> None: def cancel_running_tasks(self, wait=True): """Request cancellation of all running tasks.""" for state in list(self.__task_executors): - if not state.started or state.task_executor is None: + if state.task_executor is None or not state.task_executor.isRunning(): continue task_executor = state.task_executor task_executor.cancel_running_task() From 7304e70a11e5cfc45801c5083b4657df03277097 Mon Sep 17 00:00:00 2001 From: Henri Payno Date: Wed, 24 Jun 2026 13:45:50 +0200 Subject: [PATCH 5/9] MultiThreadedTaskExecutor: add `_getState` function for conveniance. --- src/ewoksorange/gui/concurrency/threaded.py | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/src/ewoksorange/gui/concurrency/threaded.py b/src/ewoksorange/gui/concurrency/threaded.py index 4d5bce68..3e9c9273 100644 --- a/src/ewoksorange/gui/concurrency/threaded.py +++ b/src/ewoksorange/gui/concurrency/threaded.py @@ -96,8 +96,8 @@ def execute_task( def __process_ended(self): self.__process_ended_direct(self.sender()) - def __process_ended_direct(self, task_executor: ThreadedTaskExecutor): - state = next( + def _getState(self, task_executor: ThreadedTaskExecutor) -> _TaskExecutorState: + return next( ( state for state in self.__task_executors @@ -105,6 +105,9 @@ def __process_ended_direct(self, task_executor: ThreadedTaskExecutor): ), None, ) + + def __process_ended_direct(self, task_executor: ThreadedTaskExecutor): + state = self._getState() if state is None: return From 7208fa8722318c47bacb180420d2febf7aee61aa Mon Sep 17 00:00:00 2001 From: Henri Payno Date: Wed, 24 Jun 2026 13:47:03 +0200 Subject: [PATCH 6/9] MultiThreadedTaskExecutor: move emission of 'sigComputationEnded' in the 'finally' section. --- src/ewoksorange/gui/concurrency/threaded.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/ewoksorange/gui/concurrency/threaded.py b/src/ewoksorange/gui/concurrency/threaded.py index 3e9c9273..bd757053 100644 --- a/src/ewoksorange/gui/concurrency/threaded.py +++ b/src/ewoksorange/gui/concurrency/threaded.py @@ -119,8 +119,8 @@ def __process_ended_direct(self, task_executor: ThreadedTaskExecutor): try: for callback in state.callbacks: callback(task_executor) - self.sigComputationEnded.emit() finally: + self.sigComputationEnded.emit() self.__remove_task_executor(task_executor) def stop(self, timeout: Optional[float] = None, wait: bool = False) -> None: From 0b297e5ef882faebb6c1225748100d10ce87c97f Mon Sep 17 00:00:00 2001 From: Henri Payno Date: Wed, 24 Jun 2026 13:48:24 +0200 Subject: [PATCH 7/9] MultiThreadedTaskExecutor: improve 'stop' function --- src/ewoksorange/gui/concurrency/threaded.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/src/ewoksorange/gui/concurrency/threaded.py b/src/ewoksorange/gui/concurrency/threaded.py index bd757053..e15cce82 100644 --- a/src/ewoksorange/gui/concurrency/threaded.py +++ b/src/ewoksorange/gui/concurrency/threaded.py @@ -127,10 +127,11 @@ def stop(self, timeout: Optional[float] = None, wait: bool = False) -> None: """Stop all tracked task threads.""" for state in list(self.__task_executors): task_executor = state.task_executor - if task_executor is not None: - if task_executor.receivers(task_executor.finished) > 0: - task_executor.finished.disconnect(self.__process_ended) - task_executor.stop(timeout=timeout, wait=wait) + if task_executor is None: + continue + if task_executor.receivers(task_executor.finished) > 0: + task_executor.finished.disconnect(self.__process_ended) + task_executor.stop(timeout=timeout, wait=wait) self.__task_executors.clear() def cancel_running_tasks(self, wait=True): From 0764e0c45be92b7a9344e148d71ade68d1a95b16 Mon Sep 17 00:00:00 2001 From: Henri Payno Date: Wed, 24 Jun 2026 14:06:44 +0200 Subject: [PATCH 8/9] MultiThreadedTaskExecutor: fix missing parameter to '_getState' --- src/ewoksorange/gui/concurrency/threaded.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/ewoksorange/gui/concurrency/threaded.py b/src/ewoksorange/gui/concurrency/threaded.py index e15cce82..c840f182 100644 --- a/src/ewoksorange/gui/concurrency/threaded.py +++ b/src/ewoksorange/gui/concurrency/threaded.py @@ -107,7 +107,7 @@ def _getState(self, task_executor: ThreadedTaskExecutor) -> _TaskExecutorState: ) def __process_ended_direct(self, task_executor: ThreadedTaskExecutor): - state = self._getState() + state = self._getState(task_executor=task_executor) if state is None: return From 8886be6ec8387f0f134d9e2dbe79d219e6016328 Mon Sep 17 00:00:00 2001 From: payno Date: Mon, 29 Jun 2026 08:38:29 +0200 Subject: [PATCH 9/9] Update src/ewoksorange/gui/concurrency/threaded.py MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Loïc Huder <42204205+loichuder@users.noreply.github.com> --- src/ewoksorange/gui/concurrency/threaded.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/ewoksorange/gui/concurrency/threaded.py b/src/ewoksorange/gui/concurrency/threaded.py index c840f182..0d547f3e 100644 --- a/src/ewoksorange/gui/concurrency/threaded.py +++ b/src/ewoksorange/gui/concurrency/threaded.py @@ -40,7 +40,7 @@ def cancel_running_task(self): @dataclass class _TaskExecutorState: callbacks: Iterable[Callable[[ThreadedTaskExecutor], None]] - task_kwargs: dict + task_kwargs: Dict[str, Any] log_missing_inputs: bool = False task_executor: Optional[ThreadedTaskExecutor] = None