| 107 | |
| 108 | |
| 109 | class ConcurrentExecutorGenResults(_ConcurrentExecutor): |
| 110 | |
| 111 | def _put_result(self, result, idx, success): |
| 112 | with self._condition: |
| 113 | heappush(self._results_queue, (idx, ExecutionResult(success, result))) |
| 114 | self._execute_next() |
| 115 | self._condition.notify() |
| 116 | |
| 117 | def _results(self): |
| 118 | with self._condition: |
| 119 | while self._current < self._exec_count: |
| 120 | while not self._results_queue or self._results_queue[0][0] != self._current: |
| 121 | self._condition.wait() |
| 122 | while self._results_queue and self._results_queue[0][0] == self._current: |
| 123 | _, res = heappop(self._results_queue) |
| 124 | try: |
| 125 | self._condition.release() |
| 126 | if self._fail_fast and not res[0]: |
| 127 | raise res[1] |
| 128 | yield res |
| 129 | finally: |
| 130 | self._condition.acquire() |
| 131 | self._current += 1 |
| 132 | |
| 133 | |
| 134 | class ConcurrentExecutorListResults(_ConcurrentExecutor): |
no outgoing calls
no test coverage detected