codekingpro/portable-devtools
114k
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 