Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
generators.py334 linesDownload Raw Back to psycopg
1"""2Generators implementing communication protocols with the libpq3 4Certain operations (connection, querying) are an interleave of libpq calls and5waiting for the socket to be ready. This module contains the code to execute6the operations, yielding a polling state whenever there is to wait. The7functions in the `waiting` module are the ones who wait more or less8cooperatively for the socket to be ready and make these generators continue.9 10All these generators yield pairs (fileno, `Wait`) whenever an operation would11block. The generator can be restarted sending the appropriate `Ready` state12when the file descriptor is ready.13 14"""15 16# Copyright (C) 2020 The Psycopg Team17 18import logging19from typing import List, Optional, Union20 21from . import pq22from . import errors as e23from .abc import Buffer, PipelineCommand, PQGen, PQGenConn24from .pq.abc import PGconn, PGresult25from .waiting import Wait, Ready26from ._compat import Deque27from ._cmodule import _psycopg28from ._encodings import pgconn_encoding, conninfo_encoding29 30OK = pq.ConnStatus.OK31BAD = pq.ConnStatus.BAD32 33POLL_OK = pq.PollingStatus.OK34POLL_READING = pq.PollingStatus.READING35POLL_WRITING = pq.PollingStatus.WRITING36POLL_FAILED = pq.PollingStatus.FAILED37 38COMMAND_OK = pq.ExecStatus.COMMAND_OK39COPY_OUT = pq.ExecStatus.COPY_OUT40COPY_IN = pq.ExecStatus.COPY_IN41COPY_BOTH = pq.ExecStatus.COPY_BOTH42PIPELINE_SYNC = pq.ExecStatus.PIPELINE_SYNC43 44WAIT_R = Wait.R45WAIT_W = Wait.W46WAIT_RW = Wait.RW47READY_R = Ready.R48READY_W = Ready.W49READY_RW = Ready.RW50 51logger = logging.getLogger(__name__)52 53 54def _connect(conninfo: str) -> PQGenConn[PGconn]:55    """56    Generator to create a database connection without blocking.57 58    """59    conn = pq.PGconn.connect_start(conninfo.encode())60    while True:61        if conn.status == BAD:62            encoding = conninfo_encoding(conninfo)63            raise e.OperationalError(64                f"connection is bad: {pq.error_message(conn, encoding=encoding)}",65                pgconn=conn,66            )67 68        status = conn.connect_poll()69        if status == POLL_OK:70            break71        elif status == POLL_READING:72            yield conn.socket, WAIT_R73        elif status == POLL_WRITING:74            yield conn.socket, WAIT_W75        elif status == POLL_FAILED:76            encoding = conninfo_encoding(conninfo)77            raise e.OperationalError(78                f"connection failed: {pq.error_message(conn, encoding=encoding)}",79                pgconn=e.finish_pgconn(conn),80            )81        else:82            raise e.InternalError(83                f"unexpected poll status: {status}", pgconn=e.finish_pgconn(conn)84            )85 86    conn.nonblocking = 187    return conn88 89 90def _execute(pgconn: PGconn) -> PQGen[List[PGresult]]:91    """92    Generator sending a query and returning results without blocking.93 94    The query must have already been sent using `pgconn.send_query()` or95    similar. Flush the query and then return the result using nonblocking96    functions.97 98    Return the list of results returned by the database (whether success99    or error).100    """101    yield from _send(pgconn)102    rv = yield from _fetch_many(pgconn)103    return rv104 105 106def _send(pgconn: PGconn) -> PQGen[None]:107    """108    Generator to send a query to the server without blocking.109 110    The query must have already been sent using `pgconn.send_query()` or111    similar. Flush the query and then return the result using nonblocking112    functions.113 114    After this generator has finished you may want to cycle using `fetch()`115    to retrieve the results available.116    """117    while True:118        f = pgconn.flush()119        if f == 0:120            break121 122        ready = yield WAIT_RW123        if ready & READY_R:124            # This call may read notifies: they will be saved in the125            # PGconn buffer and passed to Python later, in `fetch()`.126            pgconn.consume_input()127 128 129def _fetch_many(pgconn: PGconn) -> PQGen[List[PGresult]]:130    """131    Generator retrieving results from the database without blocking.132 133    The query must have already been sent to the server, so pgconn.flush() has134    already returned 0.135 136    Return the list of results returned by the database (whether success137    or error).138    """139    results: List[PGresult] = []140    while True:141        res = yield from _fetch(pgconn)142        if not res:143            break144 145        results.append(res)146        status = res.status147        if status == COPY_IN or status == COPY_OUT or status == COPY_BOTH:148            # After entering copy mode the libpq will create a phony result149            # for every request so let's break the endless loop.150            break151 152        if status == PIPELINE_SYNC:153            # PIPELINE_SYNC is not followed by a NULL, but we return it alone154            # similarly to other result sets.155            assert len(results) == 1, results156            break157 158    return results159 160 161def _fetch(pgconn: PGconn) -> PQGen[Optional[PGresult]]:162    """163    Generator retrieving a single result from the database without blocking.164 165    The query must have already been sent to the server, so pgconn.flush() has166    already returned 0.167 168    Return a result from the database (whether success or error).169    """170    if pgconn.is_busy():171        yield WAIT_R172        while True:173            pgconn.consume_input()174            if not pgconn.is_busy():175                break176            yield WAIT_R177 178    _consume_notifies(pgconn)179 180    return pgconn.get_result()181 182 183def _pipeline_communicate(184    pgconn: PGconn, commands: Deque[PipelineCommand]185) -> PQGen[List[List[PGresult]]]:186    """Generator to send queries from a connection in pipeline mode while also187    receiving results.188 189    Return a list results, including single PIPELINE_SYNC elements.190    """191    results = []192 193    while True:194        ready = yield WAIT_RW195 196        if ready & READY_R:197            pgconn.consume_input()198            _consume_notifies(pgconn)199 200            res: List[PGresult] = []201            while not pgconn.is_busy():202                r = pgconn.get_result()203                if r is None:204                    if not res:205                        break206                    results.append(res)207                    res = []208                else:209                    status = r.status210                    if status == PIPELINE_SYNC:211                        assert not res212                        results.append([r])213                    elif status == COPY_IN or status == COPY_OUT or status == COPY_BOTH:214                        # This shouldn't happen, but insisting hard enough, it will.215                        # For instance, in test_executemany_badquery(), with the COPY216                        # statement and the AsyncClientCursor, which disables217                        # prepared statements).218                        # Bail out from the resulting infinite loop.219                        raise e.NotSupportedError(220                            "COPY cannot be used in pipeline mode"221                        )222                    else:223                        res.append(r)224 225        if ready & READY_W:226            pgconn.flush()227            if not commands:228                break229            commands.popleft()()230 231    return results232 233 234def _consume_notifies(pgconn: PGconn) -> None:235    # Consume notifies236    while True:237        n = pgconn.notifies()238        if not n:239            break240        if pgconn.notify_handler:241            pgconn.notify_handler(n)242 243 244def notifies(pgconn: PGconn) -> PQGen[List[pq.PGnotify]]:245    yield WAIT_R246    pgconn.consume_input()247 248    ns = []249    while True:250        n = pgconn.notifies()251        if n:252            ns.append(n)253        else:254            break255 256    return ns257 258 259def copy_from(pgconn: PGconn) -> PQGen[Union[memoryview, PGresult]]:260    while True:261        nbytes, data = pgconn.get_copy_data(1)262        if nbytes != 0:263            break264 265        # would block266        yield WAIT_R267        pgconn.consume_input()268 269    if nbytes > 0:270        # some data271        return data272 273    # Retrieve the final result of copy274    results = yield from _fetch_many(pgconn)275    if len(results) > 1:276        # TODO: too brutal? Copy worked.277        raise e.ProgrammingError("you cannot mix COPY with other operations")278    result = results[0]279    if result.status != COMMAND_OK:280        encoding = pgconn_encoding(pgconn)281        raise e.error_from_result(result, encoding=encoding)282 283    return result284 285 286def copy_to(pgconn: PGconn, buffer: Buffer) -> PQGen[None]:287    # Retry enqueuing data until successful.288    #289    # WARNING! This can cause an infinite loop if the buffer is too large. (see290    # ticket #255). We avoid it in the Copy object by splitting a large buffer291    # into smaller ones. We prefer to do it there instead of here in order to292    # do it upstream the queue decoupling the writer task from the producer one.293    while pgconn.put_copy_data(buffer) == 0:294        yield WAIT_W295 296 297def copy_end(pgconn: PGconn, error: Optional[bytes]) -> PQGen[PGresult]:298    # Retry enqueuing end copy message until successful299    while pgconn.put_copy_end(error) == 0:300        yield WAIT_W301 302    # Repeat until it the message is flushed to the server303    while True:304        yield WAIT_W305        f = pgconn.flush()306        if f == 0:307            break308 309    # Retrieve the final result of copy310    (result,) = yield from _fetch_many(pgconn)311    if result.status != COMMAND_OK:312        encoding = pgconn_encoding(pgconn)313        raise e.error_from_result(result, encoding=encoding)314 315    return result316 317 318# Override functions with fast versions if available319if _psycopg:320    connect = _psycopg.connect321    execute = _psycopg.execute322    send = _psycopg.send323    fetch_many = _psycopg.fetch_many324    fetch = _psycopg.fetch325    pipeline_communicate = _psycopg.pipeline_communicate326 327else:328    connect = _connect329    execute = _execute330    send = _send331    fetch_many = _fetch_many332    fetch = _fetch333    pipeline_communicate = _pipeline_communicate334 
codekingpro/portable-devtools · Team Ai