Team Ai
Datasetpublic

codekingpro/portable-devtools

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