Patch loop to make it reentrant.
(loop)
| 140 | |
| 141 | |
| 142 | def _patch_loop(loop): |
| 143 | """Patch loop to make it reentrant.""" |
| 144 | |
| 145 | def run_forever(self): |
| 146 | with manage_run(self), manage_asyncgens(self): |
| 147 | while True: |
| 148 | self._run_once() |
| 149 | if self._stopping: |
| 150 | break |
| 151 | self._stopping = False |
| 152 | |
| 153 | def run_until_complete(self, future): |
| 154 | with manage_run(self): |
| 155 | f = asyncio.ensure_future(future, loop=self) |
| 156 | if f is not future: |
| 157 | f._log_destroy_pending = False |
| 158 | while not f.done(): |
| 159 | self._run_once() |
| 160 | if self._stopping: |
| 161 | break |
| 162 | if not f.done(): |
| 163 | raise RuntimeError("Loop stopped before Future completed!") |
| 164 | return f.result() |
| 165 | |
| 166 | def _run_once(self): |
| 167 | """Simplified re-implementation of asyncio's _run_once.""" |
| 168 | ready = self._ready |
| 169 | scheduled = self._scheduled |
| 170 | while scheduled and scheduled[0]._cancelled: |
| 171 | heappop(scheduled) |
| 172 | timeout = ( |
| 173 | 0 |
| 174 | if ready or self._stopping |
| 175 | else min(max(scheduled[0]._when - self.time(), 0), 86400) |
| 176 | if scheduled |
| 177 | else None |
| 178 | ) |
| 179 | event_list = self._selector.select(timeout) |
| 180 | self._process_events(event_list) |
| 181 | end_time = self.time() + self._clock_resolution |
| 182 | while scheduled and scheduled[0]._when < end_time: |
| 183 | handle = heappop(scheduled) |
| 184 | ready.append(handle) |
| 185 | for _ in range(len(ready)): |
| 186 | if not ready: |
| 187 | break |
| 188 | handle = ready.popleft() |
| 189 | if not handle._cancelled: |
| 190 | if sys.version_info < (3, 14, 0): |
| 191 | curr_task = curr_tasks.pop(self, None) |
| 192 | else: |
| 193 | try: |
| 194 | curr_task = asyncio.tasks._swap_current_task( |
| 195 | self, None |
| 196 | ) |
| 197 | except KeyError: |
| 198 | curr_task = None |
| 199 | try: |
no test coverage detected
searching dependent graphs…