Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
copy.py920 linesDownload Raw Back to psycopg
1"""2psycopg copy support3"""4 5# Copyright (C) 2020 The Psycopg Team6 7import re8import queue9import struct10import asyncio11import threading12from abc import ABC, abstractmethod13from types import TracebackType14from typing import Any, AsyncIterator, Dict, Generic, Iterator, List, Match, IO15from typing import Optional, Sequence, Tuple, Type, TypeVar, Union, TYPE_CHECKING16 17from . import pq18from . import adapt19from . import errors as e20from .abc import Buffer, ConnectionType, PQGen, Transformer21from ._compat import create_task22from .pq.misc import connection_summary23from ._cmodule import _psycopg24from ._encodings import pgconn_encoding25from .generators import copy_from, copy_to, copy_end26 27if TYPE_CHECKING:28    from .cursor import BaseCursor, Cursor29    from .cursor_async import AsyncCursor30    from .connection import Connection  # noqa: F40131    from .connection_async import AsyncConnection  # noqa: F40132 33PY_TEXT = adapt.PyFormat.TEXT34PY_BINARY = adapt.PyFormat.BINARY35 36TEXT = pq.Format.TEXT37BINARY = pq.Format.BINARY38 39COPY_IN = pq.ExecStatus.COPY_IN40COPY_OUT = pq.ExecStatus.COPY_OUT41 42ACTIVE = pq.TransactionStatus.ACTIVE43 44# Size of data to accumulate before sending it down the network. We fill a45# buffer this size field by field, and when it passes the threshold size46# we ship it, so it may end up being bigger than this.47BUFFER_SIZE = 32 * 102448 49# Maximum data size we want to queue to send to the libpq copy. Sending a50# buffer too big to be handled can cause an infinite loop in the libpq51# (#255) so we want to split it in more digestable chunks.52MAX_BUFFER_SIZE = 4 * BUFFER_SIZE53# Note: making this buffer too large, e.g.54# MAX_BUFFER_SIZE = 1024 * 102455# makes operations *way* slower! Probably triggering some quadraticity56# in the libpq memory management and data sending.57 58# Max size of the write queue of buffers. More than that copy will block59# Each buffer should be around BUFFER_SIZE size.60QUEUE_SIZE = 102461 62 63class BaseCopy(Generic[ConnectionType]):64    """65    Base implementation for the copy user interface.66 67    Two subclasses expose real methods with the sync/async differences.68 69    The difference between the text and binary format is managed by two70    different `Formatter` subclasses.71 72    Writing (the I/O part) is implemented in the subclasses by a `Writer` or73    `AsyncWriter` instance. Normally writing implies sending copy data to a74    database, but a different writer might be chosen, e.g. to stream data into75    a file for later use.76    """77 78    _Self = TypeVar("_Self", bound="BaseCopy[Any]")79 80    formatter: "Formatter"81 82    def __init__(83        self,84        cursor: "BaseCursor[ConnectionType, Any]",85        *,86        binary: Optional[bool] = None,87    ):88        self.cursor = cursor89        self.connection = cursor.connection90        self._pgconn = self.connection.pgconn91 92        result = cursor.pgresult93        if result:94            self._direction = result.status95            if self._direction != COPY_IN and self._direction != COPY_OUT:96                raise e.ProgrammingError(97                    "the cursor should have performed a COPY operation;"98                    f" its status is {pq.ExecStatus(self._direction).name} instead"99                )100        else:101            self._direction = COPY_IN102 103        if binary is None:104            binary = bool(result and result.binary_tuples)105 106        tx: Transformer = getattr(cursor, "_tx", None) or adapt.Transformer(cursor)107        if binary:108            self.formatter = BinaryFormatter(tx)109        else:110            self.formatter = TextFormatter(tx, encoding=pgconn_encoding(self._pgconn))111 112        self._finished = False113 114    def __repr__(self) -> str:115        cls = f"{self.__class__.__module__}.{self.__class__.__qualname__}"116        info = connection_summary(self._pgconn)117        return f"<{cls} {info} at 0x{id(self):x}>"118 119    def _enter(self) -> None:120        if self._finished:121            raise TypeError("copy blocks can be used only once")122 123    def set_types(self, types: Sequence[Union[int, str]]) -> None:124        """125        Set the types expected in a COPY operation.126 127        The types must be specified as a sequence of oid or PostgreSQL type128        names (e.g. ``int4``, ``timestamptz[]``).129 130        This operation overcomes the lack of metadata returned by PostgreSQL131        when a COPY operation begins:132 133        - On :sql:`COPY TO`, `!set_types()` allows to specify what types the134          operation returns. If `!set_types()` is not used, the data will be135          returned as unparsed strings or bytes instead of Python objects.136 137        - On :sql:`COPY FROM`, `!set_types()` allows to choose what type the138          database expects. This is especially useful in binary copy, because139          PostgreSQL will apply no cast rule.140 141        """142        registry = self.cursor.adapters.types143        oids = [t if isinstance(t, int) else registry.get_oid(t) for t in types]144 145        if self._direction == COPY_IN:146            self.formatter.transformer.set_dumper_types(oids, self.formatter.format)147        else:148            self.formatter.transformer.set_loader_types(oids, self.formatter.format)149 150    # High level copy protocol generators (state change of the Copy object)151 152    def _read_gen(self) -> PQGen[Buffer]:153        if self._finished:154            return memoryview(b"")155 156        res = yield from copy_from(self._pgconn)157        if isinstance(res, memoryview):158            return res159 160        # res is the final PGresult161        self._finished = True162 163        # This result is a COMMAND_OK which has info about the number of rows164        # returned, but not about the columns, which is instead an information165        # that was received on the COPY_OUT result at the beginning of COPY.166        # So, don't replace the results in the cursor, just update the rowcount.167        nrows = res.command_tuples168        self.cursor._rowcount = nrows if nrows is not None else -1169        return memoryview(b"")170 171    def _read_row_gen(self) -> PQGen[Optional[Tuple[Any, ...]]]:172        data = yield from self._read_gen()173        if not data:174            return None175 176        row = self.formatter.parse_row(data)177        if row is None:178            # Get the final result to finish the copy operation179            yield from self._read_gen()180            self._finished = True181            return None182 183        return row184 185    def _end_copy_out_gen(self, exc: Optional[BaseException]) -> PQGen[None]:186        if not exc:187            return188 189        if self._pgconn.transaction_status != ACTIVE:190            # The server has already finished to send copy data. The connection191            # is already in a good state.192            return193 194        # Throw a cancel to the server, then consume the rest of the copy data195        # (which might or might not have been already transferred entirely to196        # the client, so we won't necessary see the exception associated with197        # canceling).198        self.connection.cancel()199        try:200            while (yield from self._read_gen()):201                pass202        except e.QueryCanceled:203            pass204 205 206class Copy(BaseCopy["Connection[Any]"]):207    """Manage a :sql:`COPY` operation.208 209    :param cursor: the cursor where the operation is performed.210    :param binary: if `!True`, write binary format.211    :param writer: the object to write to destination. If not specified, write212        to the `!cursor` connection.213 214    Choosing `!binary` is not necessary if the cursor has executed a215    :sql:`COPY` operation, because the operation result describes the format216    too. The parameter is useful when a `!Copy` object is created manually and217    no operation is performed on the cursor, such as when using ``writer=``\\218    `~psycopg.copy.FileWriter`.219 220    """221 222    __module__ = "psycopg"223 224    writer: "Writer"225 226    def __init__(227        self,228        cursor: "Cursor[Any]",229        *,230        binary: Optional[bool] = None,231        writer: Optional["Writer"] = None,232    ):233        super().__init__(cursor, binary=binary)234        if not writer:235            writer = LibpqWriter(cursor)236 237        self.writer = writer238        self._write = writer.write239 240    def __enter__(self: BaseCopy._Self) -> BaseCopy._Self:241        self._enter()242        return self243 244    def __exit__(245        self,246        exc_type: Optional[Type[BaseException]],247        exc_val: Optional[BaseException],248        exc_tb: Optional[TracebackType],249    ) -> None:250        self.finish(exc_val)251 252    # End user sync interface253 254    def __iter__(self) -> Iterator[Buffer]:255        """Implement block-by-block iteration on :sql:`COPY TO`."""256        while True:257            data = self.read()258            if not data:259                break260            yield data261 262    def read(self) -> Buffer:263        """264        Read an unparsed row after a :sql:`COPY TO` operation.265 266        Return an empty string when the data is finished.267        """268        return self.connection.wait(self._read_gen())269 270    def rows(self) -> Iterator[Tuple[Any, ...]]:271        """272        Iterate on the result of a :sql:`COPY TO` operation record by record.273 274        Note that the records returned will be tuples of unparsed strings or275        bytes, unless data types are specified using `set_types()`.276        """277        while True:278            record = self.read_row()279            if record is None:280                break281            yield record282 283    def read_row(self) -> Optional[Tuple[Any, ...]]:284        """285        Read a parsed row of data from a table after a :sql:`COPY TO` operation.286 287        Return `!None` when the data is finished.288 289        Note that the records returned will be tuples of unparsed strings or290        bytes, unless data types are specified using `set_types()`.291        """292        return self.connection.wait(self._read_row_gen())293 294    def write(self, buffer: Union[Buffer, str]) -> None:295        """296        Write a block of data to a table after a :sql:`COPY FROM` operation.297 298        If the :sql:`COPY` is in binary format `!buffer` must be `!bytes`. In299        text mode it can be either `!bytes` or `!str`.300        """301        data = self.formatter.write(buffer)302        if data:303            self._write(data)304 305    def write_row(self, row: Sequence[Any]) -> None:306        """Write a record to a table after a :sql:`COPY FROM` operation."""307        data = self.formatter.write_row(row)308        if data:309            self._write(data)310 311    def finish(self, exc: Optional[BaseException]) -> None:312        """Terminate the copy operation and free the resources allocated.313 314        You shouldn't need to call this function yourself: it is usually called315        by exit. It is available if, despite what is documented, you end up316        using the `Copy` object outside a block.317        """318        if self._direction == COPY_IN:319            data = self.formatter.end()320            if data:321                self._write(data)322            self.writer.finish(exc)323            self._finished = True324        else:325            self.connection.wait(self._end_copy_out_gen(exc))326 327 328class Writer(ABC):329    """330    A class to write copy data somewhere.331    """332 333    @abstractmethod334    def write(self, data: Buffer) -> None:335        """336        Write some data to destination.337        """338        ...339 340    def finish(self, exc: Optional[BaseException] = None) -> None:341        """342        Called when write operations are finished.343 344        If operations finished with an error, it will be passed to ``exc``.345        """346        pass347 348 349class LibpqWriter(Writer):350    """351    A `Writer` to write copy data to a Postgres database.352    """353 354    def __init__(self, cursor: "Cursor[Any]"):355        self.cursor = cursor356        self.connection = cursor.connection357        self._pgconn = self.connection.pgconn358 359    def write(self, data: Buffer) -> None:360        if len(data) <= MAX_BUFFER_SIZE:361            # Most used path: we don't need to split the buffer in smaller362            # bits, so don't make a copy.363            self.connection.wait(copy_to(self._pgconn, data))364        else:365            # Copy a buffer too large in chunks to avoid causing a memory366            # error in the libpq, which may cause an infinite loop (#255).367            for i in range(0, len(data), MAX_BUFFER_SIZE):368                self.connection.wait(369                    copy_to(self._pgconn, data[i : i + MAX_BUFFER_SIZE])370                )371 372    def finish(self, exc: Optional[BaseException] = None) -> None:373        bmsg: Optional[bytes]374        if exc:375            msg = f"error from Python: {type(exc).__qualname__} - {exc}"376            bmsg = msg.encode(pgconn_encoding(self._pgconn), "replace")377        else:378            bmsg = None379 380        try:381            res = self.connection.wait(copy_end(self._pgconn, bmsg))382        # The QueryCanceled is expected if we sent an exception message to383        # pgconn.put_copy_end(). The Python exception that generated that384        # cancelling is more important, so don't clobber it.385        except e.QueryCanceled:386            if not bmsg:387                raise388        else:389            self.cursor._results = [res]390 391 392class QueuedLibpqWriter(LibpqWriter):393    """394    A writer using a buffer to queue data to write to a Postgres database.395 396    `write()` returns immediately, so that the main thread can be CPU-bound397    formatting messages, while a worker thread can be IO-bound waiting to write398    on the connection.399    """400 401    def __init__(self, cursor: "Cursor[Any]"):402        super().__init__(cursor)403 404        self._queue: queue.Queue[Buffer] = queue.Queue(maxsize=QUEUE_SIZE)405        self._worker: Optional[threading.Thread] = None406        self._worker_error: Optional[BaseException] = None407 408    def worker(self) -> None:409        """Push data to the server when available from the copy queue.410 411        Terminate reading when the queue receives a false-y value, or in case412        of error.413 414        The function is designed to be run in a separate thread.415        """416        try:417            while True:418                data = self._queue.get(block=True, timeout=24 * 60 * 60)419                if not data:420                    break421                self.connection.wait(copy_to(self._pgconn, data))422        except BaseException as ex:423            # Propagate the error to the main thread.424            self._worker_error = ex425 426    def write(self, data: Buffer) -> None:427        if not self._worker:428            # warning: reference loop, broken by _write_end429            self._worker = threading.Thread(target=self.worker)430            self._worker.daemon = True431            self._worker.start()432 433        # If the worker thread raies an exception, re-raise it to the caller.434        if self._worker_error:435            raise self._worker_error436 437        if len(data) <= MAX_BUFFER_SIZE:438            # Most used path: we don't need to split the buffer in smaller439            # bits, so don't make a copy.440            self._queue.put(data)441        else:442            # Copy a buffer too large in chunks to avoid causing a memory443            # error in the libpq, which may cause an infinite loop (#255).444            for i in range(0, len(data), MAX_BUFFER_SIZE):445                self._queue.put(data[i : i + MAX_BUFFER_SIZE])446 447    def finish(self, exc: Optional[BaseException] = None) -> None:448        self._queue.put(b"")449 450        if self._worker:451            self._worker.join()452            self._worker = None  # break the loop453 454        # Check if the worker thread raised any exception before terminating.455        if self._worker_error:456            raise self._worker_error457 458        super().finish(exc)459 460 461class FileWriter(Writer):462    """463    A `Writer` to write copy data to a file-like object.464 465    :param file: the file where to write copy data. It must be open for writing466        in binary mode.467    """468 469    def __init__(self, file: IO[bytes]):470        self.file = file471 472    def write(self, data: Buffer) -> None:473        self.file.write(data)474 475 476class AsyncCopy(BaseCopy["AsyncConnection[Any]"]):477    """Manage an asynchronous :sql:`COPY` operation."""478 479    __module__ = "psycopg"480 481    writer: "AsyncWriter"482 483    def __init__(484        self,485        cursor: "AsyncCursor[Any]",486        *,487        binary: Optional[bool] = None,488        writer: Optional["AsyncWriter"] = None,489    ):490        super().__init__(cursor, binary=binary)491 492        if not writer:493            writer = AsyncLibpqWriter(cursor)494 495        self.writer = writer496        self._write = writer.write497 498    async def __aenter__(self: BaseCopy._Self) -> BaseCopy._Self:499        self._enter()500        return self501 502    async def __aexit__(503        self,504        exc_type: Optional[Type[BaseException]],505        exc_val: Optional[BaseException],506        exc_tb: Optional[TracebackType],507    ) -> None:508        await self.finish(exc_val)509 510    async def __aiter__(self) -> AsyncIterator[Buffer]:511        while True:512            data = await self.read()513            if not data:514                break515            yield data516 517    async def read(self) -> Buffer:518        return await self.connection.wait(self._read_gen())519 520    async def rows(self) -> AsyncIterator[Tuple[Any, ...]]:521        while True:522            record = await self.read_row()523            if record is None:524                break525            yield record526 527    async def read_row(self) -> Optional[Tuple[Any, ...]]:528        return await self.connection.wait(self._read_row_gen())529 530    async def write(self, buffer: Union[Buffer, str]) -> None:531        data = self.formatter.write(buffer)532        if data:533            await self._write(data)534 535    async def write_row(self, row: Sequence[Any]) -> None:536        data = self.formatter.write_row(row)537        if data:538            await self._write(data)539 540    async def finish(self, exc: Optional[BaseException]) -> None:541        if self._direction == COPY_IN:542            data = self.formatter.end()543            if data:544                await self._write(data)545            await self.writer.finish(exc)546            self._finished = True547        else:548            await self.connection.wait(self._end_copy_out_gen(exc))549 550 551class AsyncWriter(ABC):552    """553    A class to write copy data somewhere (for async connections).554    """555 556    @abstractmethod557    async def write(self, data: Buffer) -> None:558        ...559 560    async def finish(self, exc: Optional[BaseException] = None) -> None:561        pass562 563 564class AsyncLibpqWriter(AsyncWriter):565    """566    An `AsyncWriter` to write copy data to a Postgres database.567    """568 569    def __init__(self, cursor: "AsyncCursor[Any]"):570        self.cursor = cursor571        self.connection = cursor.connection572        self._pgconn = self.connection.pgconn573 574    async def write(self, data: Buffer) -> None:575        if len(data) <= MAX_BUFFER_SIZE:576            # Most used path: we don't need to split the buffer in smaller577            # bits, so don't make a copy.578            await self.connection.wait(copy_to(self._pgconn, data))579        else:580            # Copy a buffer too large in chunks to avoid causing a memory581            # error in the libpq, which may cause an infinite loop (#255).582            for i in range(0, len(data), MAX_BUFFER_SIZE):583                await self.connection.wait(584                    copy_to(self._pgconn, data[i : i + MAX_BUFFER_SIZE])585                )586 587    async def finish(self, exc: Optional[BaseException] = None) -> None:588        bmsg: Optional[bytes]589        if exc:590            msg = f"error from Python: {type(exc).__qualname__} - {exc}"591            bmsg = msg.encode(pgconn_encoding(self._pgconn), "replace")592        else:593            bmsg = None594 595        try:596            res = await self.connection.wait(copy_end(self._pgconn, bmsg))597        # The QueryCanceled is expected if we sent an exception message to598        # pgconn.put_copy_end(). The Python exception that generated that599        # cancelling is more important, so don't clobber it.600        except e.QueryCanceled:601            if not bmsg:602                raise603        else:604            self.cursor._results = [res]605 606 607class AsyncQueuedLibpqWriter(AsyncLibpqWriter):608    """609    An `AsyncWriter` using a buffer to queue data to write.610 611    `write()` returns immediately, so that the main thread can be CPU-bound612    formatting messages, while a worker thread can be IO-bound waiting to write613    on the connection.614    """615 616    def __init__(self, cursor: "AsyncCursor[Any]"):617        super().__init__(cursor)618 619        self._queue: asyncio.Queue[Buffer] = asyncio.Queue(maxsize=QUEUE_SIZE)620        self._worker: Optional[asyncio.Future[None]] = None621 622    async def worker(self) -> None:623        """Push data to the server when available from the copy queue.624 625        Terminate reading when the queue receives a false-y value.626 627        The function is designed to be run in a separate task.628        """629        while True:630            data = await self._queue.get()631            if not data:632                break633            await self.connection.wait(copy_to(self._pgconn, data))634 635    async def write(self, data: Buffer) -> None:636        if not self._worker:637            self._worker = create_task(self.worker())638 639        if len(data) <= MAX_BUFFER_SIZE:640            # Most used path: we don't need to split the buffer in smaller641            # bits, so don't make a copy.642            await self._queue.put(data)643        else:644            # Copy a buffer too large in chunks to avoid causing a memory645            # error in the libpq, which may cause an infinite loop (#255).646            for i in range(0, len(data), MAX_BUFFER_SIZE):647                await self._queue.put(data[i : i + MAX_BUFFER_SIZE])648 649    async def finish(self, exc: Optional[BaseException] = None) -> None:650        await self._queue.put(b"")651 652        if self._worker:653            await asyncio.gather(self._worker)654            self._worker = None  # break reference loops if any655 656        await super().finish(exc)657 658 659class Formatter(ABC):660    """661    A class which understand a copy format (text, binary).662    """663 664    format: pq.Format665 666    def __init__(self, transformer: Transformer):667        self.transformer = transformer668        self._write_buffer = bytearray()669        self._row_mode = False  # true if the user is using write_row()670 671    @abstractmethod672    def parse_row(self, data: Buffer) -> Optional[Tuple[Any, ...]]:673        ...674 675    @abstractmethod676    def write(self, buffer: Union[Buffer, str]) -> Buffer:677        ...678 679    @abstractmethod680    def write_row(self, row: Sequence[Any]) -> Buffer:681        ...682 683    @abstractmethod684    def end(self) -> Buffer:685        ...686 687 688class TextFormatter(Formatter):689    format = TEXT690 691    def __init__(self, transformer: Transformer, encoding: str = "utf-8"):692        super().__init__(transformer)693        self._encoding = encoding694 695    def parse_row(self, data: Buffer) -> Optional[Tuple[Any, ...]]:696        if data:697            return parse_row_text(data, self.transformer)698        else:699            return None700 701    def write(self, buffer: Union[Buffer, str]) -> Buffer:702        data = self._ensure_bytes(buffer)703        self._signature_sent = True704        return data705 706    def write_row(self, row: Sequence[Any]) -> Buffer:707        # Note down that we are writing in row mode: it means we will have708        # to take care of the end-of-copy marker too709        self._row_mode = True710 711        format_row_text(row, self.transformer, self._write_buffer)712        if len(self._write_buffer) > BUFFER_SIZE:713            buffer, self._write_buffer = self._write_buffer, bytearray()714            return buffer715        else:716            return b""717 718    def end(self) -> Buffer:719        buffer, self._write_buffer = self._write_buffer, bytearray()720        return buffer721 722    def _ensure_bytes(self, data: Union[Buffer, str]) -> Buffer:723        if isinstance(data, str):724            return data.encode(self._encoding)725        else:726            # Assume, for simplicity, that the user is not passing stupid727            # things to the write function. If that's the case, things728            # will fail downstream.729            return data730 731 732class BinaryFormatter(Formatter):733    format = BINARY734 735    def __init__(self, transformer: Transformer):736        super().__init__(transformer)737        self._signature_sent = False738 739    def parse_row(self, data: Buffer) -> Optional[Tuple[Any, ...]]:740        if not self._signature_sent:741            if data[: len(_binary_signature)] != _binary_signature:742                raise e.DataError(743                    "binary copy doesn't start with the expected signature"744                )745            self._signature_sent = True746            data = data[len(_binary_signature) :]747 748        elif data == _binary_trailer:749            return None750 751        return parse_row_binary(data, self.transformer)752 753    def write(self, buffer: Union[Buffer, str]) -> Buffer:754        data = self._ensure_bytes(buffer)755        self._signature_sent = True756        return data757 758    def write_row(self, row: Sequence[Any]) -> Buffer:759        # Note down that we are writing in row mode: it means we will have760        # to take care of the end-of-copy marker too761        self._row_mode = True762 763        if not self._signature_sent:764            self._write_buffer += _binary_signature765            self._signature_sent = True766 767        format_row_binary(row, self.transformer, self._write_buffer)768        if len(self._write_buffer) > BUFFER_SIZE:769            buffer, self._write_buffer = self._write_buffer, bytearray()770            return buffer771        else:772            return b""773 774    def end(self) -> Buffer:775        # If we have sent no data we need to send the signature776        # and the trailer777        if not self._signature_sent:778            self._write_buffer += _binary_signature779            self._write_buffer += _binary_trailer780 781        elif self._row_mode:782            # if we have sent data already, we have sent the signature783            # too (either with the first row, or we assume that in784            # block mode the signature is included).785            # Write the trailer only if we are sending rows (with the786            # assumption that who is copying binary data is sending the787            # whole format).788            self._write_buffer += _binary_trailer789 790        buffer, self._write_buffer = self._write_buffer, bytearray()791        return buffer792 793    def _ensure_bytes(self, data: Union[Buffer, str]) -> Buffer:794        if isinstance(data, str):795            raise TypeError("cannot copy str data in binary mode: use bytes instead")796        else:797            # Assume, for simplicity, that the user is not passing stupid798            # things to the write function. If that's the case, things799            # will fail downstream.800            return data801 802 803def _format_row_text(804    row: Sequence[Any], tx: Transformer, out: Optional[bytearray] = None805) -> bytearray:806    """Convert a row of objects to the data to send for copy."""807    if out is None:808        out = bytearray()809 810    if not row:811        out += b"\n"812        return out813 814    for item in row:815        if item is not None:816            dumper = tx.get_dumper(item, PY_TEXT)817            b = dumper.dump(item)818            out += _dump_re.sub(_dump_sub, b)819        else:820            out += rb"\N"821        out += b"\t"822 823    out[-1:] = b"\n"824    return out825 826 827def _format_row_binary(828    row: Sequence[Any], tx: Transformer, out: Optional[bytearray] = None829) -> bytearray:830    """Convert a row of objects to the data to send for binary copy."""831    if out is None:832        out = bytearray()833 834    out += _pack_int2(len(row))835    adapted = tx.dump_sequence(row, [PY_BINARY] * len(row))836    for b in adapted:837        if b is not None:838            out += _pack_int4(len(b))839            out += b840        else:841            out += _binary_null842 843    return out844 845 846def _parse_row_text(data: Buffer, tx: Transformer) -> Tuple[Any, ...]:847    if not isinstance(data, bytes):848        data = bytes(data)849    fields = data.split(b"\t")850    fields[-1] = fields[-1][:-1]  # drop \n851    row = [None if f == b"\\N" else _load_re.sub(_load_sub, f) for f in fields]852    return tx.load_sequence(row)853 854 855def _parse_row_binary(data: Buffer, tx: Transformer) -> Tuple[Any, ...]:856    row: List[Optional[Buffer]] = []857    nfields = _unpack_int2(data, 0)[0]858    pos = 2859    for i in range(nfields):860        length = _unpack_int4(data, pos)[0]861        pos += 4862        if length >= 0:863            row.append(data[pos : pos + length])864            pos += length865        else:866            row.append(None)867 868    return tx.load_sequence(row)869 870 871_pack_int2 = struct.Struct("!h").pack872_pack_int4 = struct.Struct("!i").pack873_unpack_int2 = struct.Struct("!h").unpack_from874_unpack_int4 = struct.Struct("!i").unpack_from875 876_binary_signature = (877    b"PGCOPY\n\xff\r\n\0"  # Signature878    b"\x00\x00\x00\x00"  # flags879    b"\x00\x00\x00\x00"  # extra length880)881_binary_trailer = b"\xff\xff"882_binary_null = b"\xff\xff\xff\xff"883 884_dump_re = re.compile(b"[\b\t\n\v\f\r\\\\]")885_dump_repl = {886    b"\b": b"\\b",887    b"\t": b"\\t",888    b"\n": b"\\n",889    b"\v": b"\\v",890    b"\f": b"\\f",891    b"\r": b"\\r",892    b"\\": b"\\\\",893}894 895 896def _dump_sub(m: Match[bytes], __map: Dict[bytes, bytes] = _dump_repl) -> bytes:897    return __map[m.group(0)]898 899 900_load_re = re.compile(b"\\\\[btnvfr\\\\]")901_load_repl = {v: k for k, v in _dump_repl.items()}902 903 904def _load_sub(m: Match[bytes], __map: Dict[bytes, bytes] = _load_repl) -> bytes:905    return __map[m.group(0)]906 907 908# Override functions with fast versions if available909if _psycopg:910    format_row_text = _psycopg.format_row_text911    format_row_binary = _psycopg.format_row_binary912    parse_row_text = _psycopg.parse_row_text913    parse_row_binary = _psycopg.parse_row_binary914 915else:916    format_row_text = _format_row_text917    format_row_binary = _format_row_binary918    parse_row_text = _parse_row_text919    parse_row_binary = _parse_row_binary920 
codekingpro/portable-devtools · Team Ai