codekingpro/portable-devtools
114k
1"""2Code concerned with waiting in different contexts (blocking, async, etc).3 4These functions are designed to consume the generators returned by the5`generators` module function and to return their final value.6 7"""8 9# Copyright (C) 2020 The Psycopg Team10 11 12import os13import sys14import select15import selectors16from typing import Optional17from asyncio import get_event_loop, wait_for, Event, TimeoutError18from selectors import DefaultSelector19 20from . import errors as e21from .abc import RV, PQGen, PQGenConn, WaitFunc22from ._enums import Wait as Wait, Ready as Ready # re-exported23from ._cmodule import _psycopg24 25WAIT_R = Wait.R26WAIT_W = Wait.W27WAIT_RW = Wait.RW28READY_R = Ready.R29READY_W = Ready.W30READY_RW = Ready.RW31 32 33def wait_selector(gen: PQGen[RV], fileno: int, timeout: Optional[float] = None) -> RV:34 """35 Wait for a generator using the best strategy available.36 37 :param gen: a generator performing database operations and yielding38 `Ready` values when it would block.39 :param fileno: the file descriptor to wait on.40 :param timeout: timeout (in seconds) to check for other interrupt, e.g.41 to allow Ctrl-C.42 :type timeout: float43 :return: whatever `!gen` returns on completion.44 45 Consume `!gen`, scheduling `fileno` for completion when it is reported to46 block. Once ready again send the ready state back to `!gen`.47 """48 try:49 s = next(gen)50 with DefaultSelector() as sel:51 while True:52 sel.register(fileno, s)53 rlist = None54 while not rlist:55 rlist = sel.select(timeout=timeout)56 sel.unregister(fileno)57 # note: this line should require a cast, but mypy doesn't complain58 ready: Ready = rlist[0][1]59 assert s & ready60 s = gen.send(ready)61 62 except StopIteration as ex:63 rv: RV = ex.args[0] if ex.args else None64 return rv65 66 67def wait_conn(gen: PQGenConn[RV], timeout: Optional[float] = None) -> RV:68 """69 Wait for a connection generator using the best strategy available.70 71 :param gen: a generator performing database operations and yielding72 (fd, `Ready`) pairs when it would block.73 :param timeout: timeout (in seconds) to check for other interrupt, e.g.74 to allow Ctrl-C. If zero or None, wait indefinitely.75 :type timeout: float76 :return: whatever `!gen` returns on completion.77 78 Behave like in `wait()`, but take the fileno to wait from the generator79 itself, which might change during processing.80 """81 try:82 fileno, s = next(gen)83 if not timeout:84 timeout = None85 with DefaultSelector() as sel:86 while True:87 sel.register(fileno, s)88 rlist = sel.select(timeout=timeout)89 sel.unregister(fileno)90 if not rlist:91 raise e.ConnectionTimeout("connection timeout expired")92 ready: Ready = rlist[0][1] # type: ignore[assignment]93 fileno, s = gen.send(ready)94 95 except StopIteration as ex:96 rv: RV = ex.args[0] if ex.args else None97 return rv98 99 100async def wait_async(101 gen: PQGen[RV], fileno: int, timeout: Optional[float] = None102) -> RV:103 """104 Coroutine waiting for a generator to complete.105 106 :param gen: a generator performing database operations and yielding107 `Ready` values when it would block.108 :param fileno: the file descriptor to wait on.109 :return: whatever `!gen` returns on completion.110 111 Behave like in `wait()`, but exposing an `asyncio` interface.112 """113 # Use an event to block and restart after the fd state changes.114 # Not sure this is the best implementation but it's a start.115 ev = Event()116 loop = get_event_loop()117 ready: Ready118 s: Wait119 120 def wakeup(state: Ready) -> None:121 nonlocal ready122 ready |= state # type: ignore[assignment]123 ev.set()124 125 try:126 s = next(gen)127 while True:128 reader = s & WAIT_R129 writer = s & WAIT_W130 if not reader and not writer:131 raise e.InternalError(f"bad poll status: {s}")132 ev.clear()133 ready = 0 # type: ignore[assignment]134 if reader:135 loop.add_reader(fileno, wakeup, READY_R)136 if writer:137 loop.add_writer(fileno, wakeup, READY_W)138 try:139 if timeout is None:140 await ev.wait()141 else:142 try:143 await wait_for(ev.wait(), timeout)144 except TimeoutError:145 pass146 finally:147 if reader:148 loop.remove_reader(fileno)149 if writer:150 loop.remove_writer(fileno)151 s = gen.send(ready)152 153 except StopIteration as ex:154 rv: RV = ex.args[0] if ex.args else None155 return rv156 157 158async def wait_conn_async(gen: PQGenConn[RV], timeout: Optional[float] = None) -> RV:159 """160 Coroutine waiting for a connection generator to complete.161 162 :param gen: a generator performing database operations and yielding163 (fd, `Ready`) pairs when it would block.164 :param timeout: timeout (in seconds) to check for other interrupt, e.g.165 to allow Ctrl-C. If zero or None, wait indefinitely.166 :return: whatever `!gen` returns on completion.167 168 Behave like in `wait()`, but take the fileno to wait from the generator169 itself, which might change during processing.170 """171 # Use an event to block and restart after the fd state changes.172 # Not sure this is the best implementation but it's a start.173 ev = Event()174 loop = get_event_loop()175 ready: Ready176 s: Wait177 178 def wakeup(state: Ready) -> None:179 nonlocal ready180 ready = state181 ev.set()182 183 try:184 fileno, s = next(gen)185 if not timeout:186 timeout = None187 while True:188 reader = s & WAIT_R189 writer = s & WAIT_W190 if not reader and not writer:191 raise e.InternalError(f"bad poll status: {s}")192 ev.clear()193 ready = 0 # type: ignore[assignment]194 if reader:195 loop.add_reader(fileno, wakeup, READY_R)196 if writer:197 loop.add_writer(fileno, wakeup, READY_W)198 try:199 await wait_for(ev.wait(), timeout)200 finally:201 if reader:202 loop.remove_reader(fileno)203 if writer:204 loop.remove_writer(fileno)205 fileno, s = gen.send(ready)206 207 except TimeoutError:208 raise e.ConnectionTimeout("connection timeout expired")209 210 except StopIteration as ex:211 rv: RV = ex.args[0] if ex.args else None212 return rv213 214 215# Specialised implementation of wait functions.216 217 218def wait_select(gen: PQGen[RV], fileno: int, timeout: Optional[float] = None) -> RV:219 """220 Wait for a generator using select where supported.221 222 BUG: on Linux, can't select on FD >= 1024. On Windows it's fine.223 """224 try:225 s = next(gen)226 227 empty = ()228 fnlist = (fileno,)229 while True:230 rl, wl, xl = select.select(231 fnlist if s & WAIT_R else empty,232 fnlist if s & WAIT_W else empty,233 fnlist,234 timeout,235 )236 ready = 0237 if rl:238 ready = READY_R239 if wl:240 ready |= READY_W241 if not ready:242 continue243 # assert s & ready244 s = gen.send(ready) # type: ignore245 246 except StopIteration as ex:247 rv: RV = ex.args[0] if ex.args else None248 return rv249 250 251if hasattr(selectors, "EpollSelector"):252 _epoll_evmasks = {253 WAIT_R: select.EPOLLONESHOT | select.EPOLLIN | select.EPOLLERR,254 WAIT_W: select.EPOLLONESHOT | select.EPOLLOUT | select.EPOLLERR,255 WAIT_RW: select.EPOLLONESHOT256 | (select.EPOLLIN | select.EPOLLOUT | select.EPOLLERR),257 }258else:259 _epoll_evmasks = {}260 261 262def wait_epoll(gen: PQGen[RV], fileno: int, timeout: Optional[float] = None) -> RV:263 """264 Wait for a generator using epoll where supported.265 266 Parameters are like for `wait()`. If it is detected that the best selector267 strategy is `epoll` then this function will be used instead of `wait`.268 269 See also: https://linux.die.net/man/2/epoll_ctl270 271 BUG: if the connection FD is closed, `epoll.poll()` hangs. Same for272 EpollSelector. For this reason, wait_poll() is currently preferable.273 To reproduce the bug:274 275 export PSYCOPG_WAIT_FUNC=wait_epoll276 pytest tests/test_concurrency.py::test_concurrent_close277 """278 try:279 s = next(gen)280 281 if timeout is None or timeout < 0:282 timeout = 0283 else:284 timeout = int(timeout * 1000.0)285 286 with select.epoll() as epoll:287 evmask = _epoll_evmasks[s]288 epoll.register(fileno, evmask)289 while True:290 fileevs = None291 while not fileevs:292 fileevs = epoll.poll(timeout)293 ev = fileevs[0][1]294 ready = 0295 if ev & ~select.EPOLLOUT:296 ready = READY_R297 if ev & ~select.EPOLLIN:298 ready |= READY_W299 # assert s & ready300 s = gen.send(ready)301 evmask = _epoll_evmasks[s]302 epoll.modify(fileno, evmask)303 304 except StopIteration as ex:305 rv: RV = ex.args[0] if ex.args else None306 return rv307 308 309if hasattr(selectors, "PollSelector"):310 _poll_evmasks = {311 WAIT_R: select.POLLIN,312 WAIT_W: select.POLLOUT,313 WAIT_RW: select.POLLIN | select.POLLOUT,314 }315else:316 _poll_evmasks = {}317 318 319def wait_poll(gen: PQGen[RV], fileno: int, timeout: Optional[float] = None) -> RV:320 """321 Wait for a generator using poll where supported.322 323 Parameters are like for `wait()`.324 """325 try:326 s = next(gen)327 328 if timeout is None or timeout < 0:329 timeout = 0330 else:331 timeout = int(timeout * 1000.0)332 333 poll = select.poll()334 evmask = _poll_evmasks[s]335 poll.register(fileno, evmask)336 while True:337 fileevs = None338 while not fileevs:339 fileevs = poll.poll(timeout)340 ev = fileevs[0][1]341 ready = 0342 if ev & ~select.POLLOUT:343 ready = READY_R344 if ev & ~select.POLLIN:345 ready |= READY_W346 # assert s & ready347 s = gen.send(ready)348 evmask = _poll_evmasks[s]349 poll.modify(fileno, evmask)350 351 except StopIteration as ex:352 rv: RV = ex.args[0] if ex.args else None353 return rv354 355 356if _psycopg:357 wait_c = _psycopg.wait_c358 359 360# Choose the best wait strategy for the platform.361#362# the selectors objects have a generic interface but come with some overhead,363# so we also offer more finely tuned implementations.364 365wait: WaitFunc366 367# Allow the user to choose a specific function for testing368if "PSYCOPG_WAIT_FUNC" in os.environ:369 fname = os.environ["PSYCOPG_WAIT_FUNC"]370 if not fname.startswith("wait_") or fname not in globals():371 raise ImportError(372 "PSYCOPG_WAIT_FUNC should be the name of an available wait function;"373 f" got {fname!r}"374 )375 wait = globals()[fname]376 377# On Windows, for the moment, avoid using wait_c, because it was reported to378# use excessive CPU (see #645).379# TODO: investigate why.380elif _psycopg and sys.platform != "win32":381 wait = wait_c382 383elif selectors.DefaultSelector is getattr(selectors, "SelectSelector", None):384 # On Windows, SelectSelector should be the default.385 wait = wait_select386 387elif hasattr(selectors, "PollSelector"):388 # On linux, EpollSelector is the default. However, it hangs if the fd is389 # closed while polling.390 wait = wait_poll391 392else:393 wait = wait_selector394 