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