codekingpro/portable-devtools
114k
1"""2psycopg cursor objects3"""4 5# Copyright (C) 2020 The Psycopg Team6 7from functools import partial8from types import TracebackType9from typing import Any, Generic, Iterable, Iterator, List10from typing import Optional, NoReturn, Sequence, Tuple, Type, TypeVar11from typing import overload, TYPE_CHECKING12from warnings import warn13from contextlib import contextmanager14 15from . import pq16from . import adapt17from . import errors as e18from .abc import ConnectionType, Query, Params, PQGen19from .copy import Copy, Writer as CopyWriter20from .rows import Row, RowMaker, RowFactory21from ._column import Column22from .pq.misc import connection_summary23from ._queries import PostgresQuery, PostgresClientQuery24from ._pipeline import Pipeline25from ._encodings import pgconn_encoding26from ._preparing import Prepare27from .generators import execute, fetch, send28 29if TYPE_CHECKING:30 from .abc import Transformer31 from .pq.abc import PGconn, PGresult32 from .connection import Connection33 34TEXT = pq.Format.TEXT35BINARY = pq.Format.BINARY36 37EMPTY_QUERY = pq.ExecStatus.EMPTY_QUERY38COMMAND_OK = pq.ExecStatus.COMMAND_OK39TUPLES_OK = pq.ExecStatus.TUPLES_OK40COPY_OUT = pq.ExecStatus.COPY_OUT41COPY_IN = pq.ExecStatus.COPY_IN42COPY_BOTH = pq.ExecStatus.COPY_BOTH43FATAL_ERROR = pq.ExecStatus.FATAL_ERROR44SINGLE_TUPLE = pq.ExecStatus.SINGLE_TUPLE45PIPELINE_ABORTED = pq.ExecStatus.PIPELINE_ABORTED46 47ACTIVE = pq.TransactionStatus.ACTIVE48 49 50class BaseCursor(Generic[ConnectionType, Row]):51 __slots__ = """52 _conn format _adapters arraysize _closed _results pgresult _pos53 _iresult _rowcount _query _tx _last_query _row_factory _make_row54 _pgconn _execmany_returning55 __weakref__56 """.split()57 58 ExecStatus = pq.ExecStatus59 60 _tx: "Transformer"61 _make_row: RowMaker[Row]62 _pgconn: "PGconn"63 64 def __init__(self, connection: ConnectionType):65 self._conn = connection66 self.format = TEXT67 self._pgconn = connection.pgconn68 self._adapters = adapt.AdaptersMap(connection.adapters)69 self.arraysize = 170 self._closed = False71 self._last_query: Optional[Query] = None72 self._reset()73 74 def _reset(self, reset_query: bool = True) -> None:75 self._results: List["PGresult"] = []76 self.pgresult: Optional["PGresult"] = None77 self._pos = 078 self._iresult = 079 self._rowcount = -180 self._query: Optional[PostgresQuery]81 # None if executemany() not executing, True/False according to returning state82 self._execmany_returning: Optional[bool] = None83 if reset_query:84 self._query = None85 86 def __repr__(self) -> str:87 cls = f"{self.__class__.__module__}.{self.__class__.__qualname__}"88 info = connection_summary(self._pgconn)89 if self._closed:90 status = "closed"91 elif self.pgresult:92 status = pq.ExecStatus(self.pgresult.status).name93 else:94 status = "no result"95 return f"<{cls} [{status}] {info} at 0x{id(self):x}>"96 97 @property98 def connection(self) -> ConnectionType:99 """The connection this cursor is using."""100 return self._conn101 102 @property103 def adapters(self) -> adapt.AdaptersMap:104 return self._adapters105 106 @property107 def closed(self) -> bool:108 """`True` if the cursor is closed."""109 return self._closed110 111 @property112 def description(self) -> Optional[List[Column]]:113 """114 A list of `Column` objects describing the current resultset.115 116 `!None` if the current resultset didn't return tuples.117 """118 res = self.pgresult119 120 # We return columns if we have nfields, but also if we don't but121 # the query said we got tuples (mostly to handle the super useful122 # query "SELECT ;"123 if res and (124 res.nfields or res.status == TUPLES_OK or res.status == SINGLE_TUPLE125 ):126 return [Column(self, i) for i in range(res.nfields)]127 else:128 return None129 130 @property131 def rowcount(self) -> int:132 """Number of records affected by the precedent operation."""133 return self._rowcount134 135 @property136 def rownumber(self) -> Optional[int]:137 """Index of the next row to fetch in the current result.138 139 `!None` if there is no result to fetch.140 """141 tuples = self.pgresult and self.pgresult.status == TUPLES_OK142 return self._pos if tuples else None143 144 def setinputsizes(self, sizes: Sequence[Any]) -> None:145 # no-op146 pass147 148 def setoutputsize(self, size: Any, column: Optional[int] = None) -> None:149 # no-op150 pass151 152 def nextset(self) -> Optional[bool]:153 """154 Move to the result set of the next query executed through `executemany()`155 or to the next result set if `execute()` returned more than one.156 157 Return `!True` if a new result is available, which will be the one158 methods `!fetch*()` will operate on.159 """160 # Raise a warning if people is calling nextset() in pipeline mode161 # after a sequence of execute() in pipeline mode. Pipeline accumulating162 # execute() results in the cursor is an unintended difference w.r.t.163 # non-pipeline mode.164 if self._execmany_returning is None and self._conn._pipeline:165 warn(166 "using nextset() in pipeline mode for several execute() is"167 " deprecated and will be dropped in 3.2; please use different"168 " cursors to receive more than one result",169 DeprecationWarning,170 )171 172 if self._iresult < len(self._results) - 1:173 self._select_current_result(self._iresult + 1)174 return True175 else:176 return None177 178 @property179 def statusmessage(self) -> Optional[str]:180 """181 The command status tag from the last SQL command executed.182 183 `!None` if the cursor doesn't have a result available.184 """185 msg = self.pgresult.command_status if self.pgresult else None186 return msg.decode() if msg else None187 188 def _make_row_maker(self) -> RowMaker[Row]:189 raise NotImplementedError190 191 #192 # Generators for the high level operations on the cursor193 #194 # Like for sync/async connections, these are implemented as generators195 # so that different concurrency strategies (threads,asyncio) can use their196 # own way of waiting (or better, `connection.wait()`).197 #198 199 def _execute_gen(200 self,201 query: Query,202 params: Optional[Params] = None,203 *,204 prepare: Optional[bool] = None,205 binary: Optional[bool] = None,206 ) -> PQGen[None]:207 """Generator implementing `Cursor.execute()`."""208 yield from self._start_query(query)209 pgq = self._convert_query(query, params)210 results = yield from self._maybe_prepare_gen(211 pgq, prepare=prepare, binary=binary212 )213 if self._conn._pipeline:214 yield from self._conn._pipeline._communicate_gen()215 else:216 assert results is not None217 self._check_results(results)218 self._results = results219 self._select_current_result(0)220 221 self._last_query = query222 223 for cmd in self._conn._prepared.get_maintenance_commands():224 yield from self._conn._exec_command(cmd)225 226 def _executemany_gen_pipeline(227 self, query: Query, params_seq: Iterable[Params], returning: bool228 ) -> PQGen[None]:229 """230 Generator implementing `Cursor.executemany()` with pipelines available.231 """232 pipeline = self._conn._pipeline233 assert pipeline234 235 yield from self._start_query(query)236 if not returning:237 self._rowcount = 0238 239 assert self._execmany_returning is None240 self._execmany_returning = returning241 242 first = True243 for params in params_seq:244 if first:245 pgq = self._convert_query(query, params)246 self._query = pgq247 first = False248 else:249 pgq.dump(params)250 251 yield from self._maybe_prepare_gen(pgq, prepare=True)252 yield from pipeline._communicate_gen()253 254 self._last_query = query255 256 if returning:257 yield from pipeline._fetch_gen(flush=True)258 259 for cmd in self._conn._prepared.get_maintenance_commands():260 yield from self._conn._exec_command(cmd)261 262 def _executemany_gen_no_pipeline(263 self, query: Query, params_seq: Iterable[Params], returning: bool264 ) -> PQGen[None]:265 """266 Generator implementing `Cursor.executemany()` with pipelines not available.267 """268 yield from self._start_query(query)269 if not returning:270 self._rowcount = 0271 first = True272 for params in params_seq:273 if first:274 pgq = self._convert_query(query, params)275 self._query = pgq276 first = False277 else:278 pgq.dump(params)279 280 results = yield from self._maybe_prepare_gen(pgq, prepare=True)281 assert results is not None282 self._check_results(results)283 if returning:284 self._results.extend(results)285 else:286 # In non-returning case, set rowcount to the cumulated number287 # of rows of executed queries.288 for res in results:289 self._rowcount += res.command_tuples or 0290 291 if self._results:292 self._select_current_result(0)293 294 self._last_query = query295 296 for cmd in self._conn._prepared.get_maintenance_commands():297 yield from self._conn._exec_command(cmd)298 299 def _maybe_prepare_gen(300 self,301 pgq: PostgresQuery,302 *,303 prepare: Optional[bool] = None,304 binary: Optional[bool] = None,305 ) -> PQGen[Optional[List["PGresult"]]]:306 # Check if the query is prepared or needs preparing307 prep, name = self._get_prepared(pgq, prepare)308 if prep is Prepare.NO:309 # The query must be executed without preparing310 self._execute_send(pgq, binary=binary)311 else:312 # If the query is not already prepared, prepare it.313 if prep is Prepare.SHOULD:314 self._send_prepare(name, pgq)315 if not self._conn._pipeline:316 (result,) = yield from execute(self._pgconn)317 if result.status == FATAL_ERROR:318 raise e.error_from_result(result, encoding=self._encoding)319 # Then execute it.320 self._send_query_prepared(name, pgq, binary=binary)321 322 # Update the prepare state of the query.323 # If an operation requires to flush our prepared statements cache,324 # it will be added to the maintenance commands to execute later.325 key = self._conn._prepared.maybe_add_to_cache(pgq, prep, name)326 327 if self._conn._pipeline:328 queued = None329 if key is not None:330 queued = (key, prep, name)331 self._conn._pipeline.result_queue.append((self, queued))332 return None333 334 # run the query335 results = yield from execute(self._pgconn)336 337 if key is not None:338 self._conn._prepared.validate(key, prep, name, results)339 340 return results341 342 def _get_prepared(343 self, pgq: PostgresQuery, prepare: Optional[bool] = None344 ) -> Tuple[Prepare, bytes]:345 return self._conn._prepared.get(pgq, prepare)346 347 def _stream_send_gen(348 self,349 query: Query,350 params: Optional[Params] = None,351 *,352 binary: Optional[bool] = None,353 ) -> PQGen[None]:354 """Generator to send the query for `Cursor.stream()`."""355 yield from self._start_query(query)356 pgq = self._convert_query(query, params)357 self._execute_send(pgq, binary=binary, force_extended=True)358 self._pgconn.set_single_row_mode()359 self._last_query = query360 yield from send(self._pgconn)361 362 def _stream_fetchone_gen(self, first: bool) -> PQGen[Optional["PGresult"]]:363 res = yield from fetch(self._pgconn)364 if res is None:365 return None366 367 status = res.status368 if status == SINGLE_TUPLE:369 self.pgresult = res370 self._tx.set_pgresult(res, set_loaders=first)371 if first:372 self._make_row = self._make_row_maker()373 return res374 375 elif status == TUPLES_OK or status == COMMAND_OK:376 # End of single row results377 while res:378 res = yield from fetch(self._pgconn)379 if status != TUPLES_OK:380 raise e.ProgrammingError(381 "the operation in stream() didn't produce a result"382 )383 return None384 385 else:386 # Errors, unexpected values387 return self._raise_for_result(res)388 389 def _start_query(self, query: Optional[Query] = None) -> PQGen[None]:390 """Generator to start the processing of a query.391 392 It is implemented as generator because it may send additional queries,393 such as `begin`.394 """395 if self.closed:396 raise e.InterfaceError("the cursor is closed")397 398 self._reset()399 if not self._last_query or (self._last_query is not query):400 self._last_query = None401 self._tx = adapt.Transformer(self)402 yield from self._conn._start_query()403 404 def _start_copy_gen(405 self, statement: Query, params: Optional[Params] = None406 ) -> PQGen[None]:407 """Generator implementing sending a command for `Cursor.copy()."""408 409 # The connection gets in an unrecoverable state if we attempt COPY in410 # pipeline mode. Forbid it explicitly.411 if self._conn._pipeline:412 raise e.NotSupportedError("COPY cannot be used in pipeline mode")413 414 yield from self._start_query()415 416 # Merge the params client-side417 if params:418 pgq = PostgresClientQuery(self._tx)419 pgq.convert(statement, params)420 statement = pgq.query421 422 query = self._convert_query(statement)423 424 self._execute_send(query, binary=False)425 results = yield from execute(self._pgconn)426 if len(results) != 1:427 raise e.ProgrammingError("COPY cannot be mixed with other operations")428 429 self._check_copy_result(results[0])430 self._results = results431 self._select_current_result(0)432 433 def _execute_send(434 self,435 query: PostgresQuery,436 *,437 force_extended: bool = False,438 binary: Optional[bool] = None,439 ) -> None:440 """441 Implement part of execute() before waiting common to sync and async.442 443 This is not a generator, but a normal non-blocking function.444 """445 if binary is None:446 fmt = self.format447 else:448 fmt = BINARY if binary else TEXT449 450 self._query = query451 452 if self._conn._pipeline:453 # In pipeline mode always use PQsendQueryParams - see #314454 # Multiple statements in the same query are not allowed anyway.455 self._conn._pipeline.command_queue.append(456 partial(457 self._pgconn.send_query_params,458 query.query,459 query.params,460 param_formats=query.formats,461 param_types=query.types,462 result_format=fmt,463 )464 )465 elif force_extended or query.params or fmt == BINARY:466 self._pgconn.send_query_params(467 query.query,468 query.params,469 param_formats=query.formats,470 param_types=query.types,471 result_format=fmt,472 )473 else:474 # If we can, let's use simple query protocol,475 # as it can execute more than one statement in a single query.476 self._pgconn.send_query(query.query)477 478 def _convert_query(479 self, query: Query, params: Optional[Params] = None480 ) -> PostgresQuery:481 pgq = PostgresQuery(self._tx)482 pgq.convert(query, params)483 return pgq484 485 def _check_results(self, results: List["PGresult"]) -> None:486 """487 Verify that the results of a query are valid.488 489 Verify that the query returned at least one result and that they all490 represent a valid result from the database.491 """492 if not results:493 raise e.InternalError("got no result from the query")494 495 for res in results:496 status = res.status497 if status != TUPLES_OK and status != COMMAND_OK and status != EMPTY_QUERY:498 self._raise_for_result(res)499 500 def _raise_for_result(self, result: "PGresult") -> NoReturn:501 """502 Raise an appropriate error message for an unexpected database result503 """504 status = result.status505 if status == FATAL_ERROR:506 raise e.error_from_result(result, encoding=self._encoding)507 elif status == PIPELINE_ABORTED:508 raise e.PipelineAborted("pipeline aborted")509 elif status == COPY_IN or status == COPY_OUT or status == COPY_BOTH:510 raise e.ProgrammingError(511 "COPY cannot be used with this method; use copy() instead"512 )513 else:514 raise e.InternalError(515 "unexpected result status from query:" f" {pq.ExecStatus(status).name}"516 )517 518 def _select_current_result(519 self, i: int, format: Optional[pq.Format] = None520 ) -> None:521 """522 Select one of the results in the cursor as the active one.523 """524 self._iresult = i525 res = self.pgresult = self._results[i]526 527 # Note: the only reason to override format is to correctly set528 # binary loaders on server-side cursors, because send_describe_portal529 # only returns a text result.530 self._tx.set_pgresult(res, format=format)531 532 self._pos = 0533 534 if res.status == TUPLES_OK:535 self._rowcount = self.pgresult.ntuples536 537 # COPY_OUT has never info about nrows. We need such result for the538 # columns in order to return a `description`, but not overwrite the539 # cursor rowcount (which was set by the Copy object).540 elif res.status != COPY_OUT:541 nrows = self.pgresult.command_tuples542 self._rowcount = nrows if nrows is not None else -1543 544 self._make_row = self._make_row_maker()545 546 def _set_results_from_pipeline(self, results: List["PGresult"]) -> None:547 self._check_results(results)548 first_batch = not self._results549 550 if self._execmany_returning is None:551 # Received from execute()552 self._results.extend(results)553 if first_batch:554 self._select_current_result(0)555 556 else:557 # Received from executemany()558 if self._execmany_returning:559 self._results.extend(results)560 if first_batch:561 self._select_current_result(0)562 else:563 # In non-returning case, set rowcount to the cumulated number of564 # rows of executed queries.565 for res in results:566 self._rowcount += res.command_tuples or 0567 568 def _send_prepare(self, name: bytes, query: PostgresQuery) -> None:569 if self._conn._pipeline:570 self._conn._pipeline.command_queue.append(571 partial(572 self._pgconn.send_prepare,573 name,574 query.query,575 param_types=query.types,576 )577 )578 self._conn._pipeline.result_queue.append(None)579 else:580 self._pgconn.send_prepare(name, query.query, param_types=query.types)581 582 def _send_query_prepared(583 self, name: bytes, pgq: PostgresQuery, *, binary: Optional[bool] = None584 ) -> None:585 if binary is None:586 fmt = self.format587 else:588 fmt = BINARY if binary else TEXT589 590 if self._conn._pipeline:591 self._conn._pipeline.command_queue.append(592 partial(593 self._pgconn.send_query_prepared,594 name,595 pgq.params,596 param_formats=pgq.formats,597 result_format=fmt,598 )599 )600 else:601 self._pgconn.send_query_prepared(602 name, pgq.params, param_formats=pgq.formats, result_format=fmt603 )604 605 def _check_result_for_fetch(self) -> None:606 if self.closed:607 raise e.InterfaceError("the cursor is closed")608 res = self.pgresult609 if not res:610 raise e.ProgrammingError("no result available")611 612 status = res.status613 if status == TUPLES_OK:614 return615 elif status == FATAL_ERROR:616 raise e.error_from_result(res, encoding=self._encoding)617 elif status == PIPELINE_ABORTED:618 raise e.PipelineAborted("pipeline aborted")619 else:620 raise e.ProgrammingError("the last operation didn't produce a result")621 622 def _check_copy_result(self, result: "PGresult") -> None:623 """624 Check that the value returned in a copy() operation is a legit COPY.625 """626 status = result.status627 if status == COPY_IN or status == COPY_OUT:628 return629 elif status == FATAL_ERROR:630 raise e.error_from_result(result, encoding=self._encoding)631 else:632 raise e.ProgrammingError(633 "copy() should be used only with COPY ... TO STDOUT or COPY ..."634 f" FROM STDIN statements, got {pq.ExecStatus(status).name}"635 )636 637 def _scroll(self, value: int, mode: str) -> None:638 self._check_result_for_fetch()639 assert self.pgresult640 if mode == "relative":641 newpos = self._pos + value642 elif mode == "absolute":643 newpos = value644 else:645 raise ValueError(f"bad mode: {mode}. It should be 'relative' or 'absolute'")646 if not 0 <= newpos < self.pgresult.ntuples:647 raise IndexError("position out of bound")648 self._pos = newpos649 650 def _close(self) -> None:651 """Non-blocking part of closing. Common to sync/async."""652 # Don't reset the query because it may be useful to investigate after653 # an error.654 self._reset(reset_query=False)655 self._closed = True656 657 @property658 def _encoding(self) -> str:659 return pgconn_encoding(self._pgconn)660 661 662class Cursor(BaseCursor["Connection[Any]", Row]):663 __module__ = "psycopg"664 __slots__ = ()665 _Self = TypeVar("_Self", bound="Cursor[Any]")666 667 @overload668 def __init__(self: "Cursor[Row]", connection: "Connection[Row]"):669 ...670 671 @overload672 def __init__(673 self: "Cursor[Row]",674 connection: "Connection[Any]",675 *,676 row_factory: RowFactory[Row],677 ):678 ...679 680 def __init__(681 self,682 connection: "Connection[Any]",683 *,684 row_factory: Optional[RowFactory[Row]] = None,685 ):686 super().__init__(connection)687 self._row_factory = row_factory or connection.row_factory688 689 def __enter__(self: _Self) -> _Self:690 return self691 692 def __exit__(693 self,694 exc_type: Optional[Type[BaseException]],695 exc_val: Optional[BaseException],696 exc_tb: Optional[TracebackType],697 ) -> None:698 self.close()699 700 def close(self) -> None:701 """702 Close the current cursor and free associated resources.703 """704 self._close()705 706 @property707 def row_factory(self) -> RowFactory[Row]:708 """Writable attribute to control how result rows are formed."""709 return self._row_factory710 711 @row_factory.setter712 def row_factory(self, row_factory: RowFactory[Row]) -> None:713 self._row_factory = row_factory714 if self.pgresult:715 self._make_row = row_factory(self)716 717 def _make_row_maker(self) -> RowMaker[Row]:718 return self._row_factory(self)719 720 def execute(721 self: _Self,722 query: Query,723 params: Optional[Params] = None,724 *,725 prepare: Optional[bool] = None,726 binary: Optional[bool] = None,727 ) -> _Self:728 """729 Execute a query or command to the database.730 """731 try:732 with self._conn.lock:733 self._conn.wait(734 self._execute_gen(query, params, prepare=prepare, binary=binary)735 )736 except e._NO_TRACEBACK as ex:737 raise ex.with_traceback(None)738 return self739 740 def executemany(741 self,742 query: Query,743 params_seq: Iterable[Params],744 *,745 returning: bool = False,746 ) -> None:747 """748 Execute the same command with a sequence of input data.749 """750 try:751 if Pipeline.is_supported():752 # If there is already a pipeline, ride it, in order to avoid753 # sending unnecessary Sync.754 with self._conn.lock:755 p = self._conn._pipeline756 if p:757 self._conn.wait(758 self._executemany_gen_pipeline(query, params_seq, returning)759 )760 # Otherwise, make a new one761 if not p:762 with self._conn.pipeline(), self._conn.lock:763 self._conn.wait(764 self._executemany_gen_pipeline(query, params_seq, returning)765 )766 else:767 with self._conn.lock:768 self._conn.wait(769 self._executemany_gen_no_pipeline(query, params_seq, returning)770 )771 except e._NO_TRACEBACK as ex:772 raise ex.with_traceback(None)773 774 def stream(775 self,776 query: Query,777 params: Optional[Params] = None,778 *,779 binary: Optional[bool] = None,780 ) -> Iterator[Row]:781 """782 Iterate row-by-row on a result from the database.783 """784 if self._pgconn.pipeline_status:785 raise e.ProgrammingError("stream() cannot be used in pipeline mode")786 787 with self._conn.lock:788 try:789 self._conn.wait(self._stream_send_gen(query, params, binary=binary))790 first = True791 while self._conn.wait(self._stream_fetchone_gen(first)):792 # We know that, if we got a result, it has a single row.793 rec: Row = self._tx.load_row(0, self._make_row) # type: ignore794 yield rec795 first = False796 797 except e._NO_TRACEBACK as ex:798 raise ex.with_traceback(None)799 800 finally:801 if self._pgconn.transaction_status == ACTIVE:802 # Try to cancel the query, then consume the results803 # already received.804 self._conn.cancel()805 try:806 while self._conn.wait(self._stream_fetchone_gen(first=False)):807 pass808 except Exception:809 pass810 811 # Try to get out of ACTIVE state. Just do a single attempt, which812 # should work to recover from an error or query cancelled.813 try:814 self._conn.wait(self._stream_fetchone_gen(first=False))815 except Exception:816 pass817 818 def fetchone(self) -> Optional[Row]:819 """820 Return the next record from the current recordset.821 822 Return `!None` the recordset is finished.823 824 :rtype: Optional[Row], with Row defined by `row_factory`825 """826 self._fetch_pipeline()827 self._check_result_for_fetch()828 record = self._tx.load_row(self._pos, self._make_row)829 if record is not None:830 self._pos += 1831 return record832 833 def fetchmany(self, size: int = 0) -> List[Row]:834 """835 Return the next `!size` records from the current recordset.836 837 `!size` default to `!self.arraysize` if not specified.838 839 :rtype: Sequence[Row], with Row defined by `row_factory`840 """841 self._fetch_pipeline()842 self._check_result_for_fetch()843 assert self.pgresult844 845 if not size:846 size = self.arraysize847 records = self._tx.load_rows(848 self._pos,849 min(self._pos + size, self.pgresult.ntuples),850 self._make_row,851 )852 self._pos += len(records)853 return records854 855 def fetchall(self) -> List[Row]:856 """857 Return all the remaining records from the current recordset.858 859 :rtype: Sequence[Row], with Row defined by `row_factory`860 """861 self._fetch_pipeline()862 self._check_result_for_fetch()863 assert self.pgresult864 records = self._tx.load_rows(self._pos, self.pgresult.ntuples, self._make_row)865 self._pos = self.pgresult.ntuples866 return records867 868 def __iter__(self) -> Iterator[Row]:869 self._fetch_pipeline()870 self._check_result_for_fetch()871 872 def load(pos: int) -> Optional[Row]:873 return self._tx.load_row(pos, self._make_row)874 875 while True:876 row = load(self._pos)877 if row is None:878 break879 self._pos += 1880 yield row881 882 def scroll(self, value: int, mode: str = "relative") -> None:883 """884 Move the cursor in the result set to a new position according to mode.885 886 If `!mode` is ``'relative'`` (default), `!value` is taken as offset to887 the current position in the result set; if set to ``'absolute'``,888 `!value` states an absolute target position.889 890 Raise `!IndexError` in case a scroll operation would leave the result891 set. In this case the position will not change.892 """893 self._fetch_pipeline()894 self._scroll(value, mode)895 896 @contextmanager897 def copy(898 self,899 statement: Query,900 params: Optional[Params] = None,901 *,902 writer: Optional[CopyWriter] = None,903 ) -> Iterator[Copy]:904 """905 Initiate a :sql:`COPY` operation and return an object to manage it.906 907 :rtype: Copy908 """909 try:910 with self._conn.lock:911 self._conn.wait(self._start_copy_gen(statement, params))912 913 with Copy(self, writer=writer) as copy:914 yield copy915 except e._NO_TRACEBACK as ex:916 raise ex.with_traceback(None)917 918 # If a fresher result has been set on the cursor by the Copy object,919 # read its properties (especially rowcount).920 self._select_current_result(0)921 922 def _fetch_pipeline(self) -> None:923 if (924 self._execmany_returning is not False925 and not self.pgresult926 and self._conn._pipeline927 ):928 with self._conn.lock:929 self._conn.wait(self._conn._pipeline._fetch_gen(flush=True))930 