Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
waiting.py394 linesDownload Raw Back to psycopg
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 
codekingpro/portable-devtools · Team Ai