MagicStack / MagicStack/asyncpg

Connection.copy_records_to_table hanging when generator passed to records parameter raises exception with long error message

Đang mở
#1,354 0 bình luận 0 reaction 0 người được giao Xem trên GitHub

Chưa có ai nhận issue này.

Ngôn ngữ chính
Python
Star
8.1k
Fork
468
Chỉ số merge pull request
Không có pull request nào được merge trong 30 ngày

Mô tả

When a generator passed to the `records` argument of `Connection.copy_records_to_table` raises an exception that when stringified with `str(exc)` produces a 9996 byte long (or greater) message it makes the COPY hang forever.

Reproduction below

```python
#!/usr/bin/env -S uv run --script
# /// script
# requires-python = "==3.15"
# dependencies = ["asyncpg==0.31.0"]
# ///
"""Repro: asyncpg deadlocks in ROLLBACK when the records generator passed to
copy_records_to_table raises an exception whose str() exceeds PostgreSQL's
10,000-byte frontend message limit.

Run (assumes an existing PostgreSQL; only touches a session-local temp table):

DSN=postgres://user:pass@host:5432/db uv run repro.py

Exit code 1 means the deadlock occurred.
"""

import asyncio
import os
import sys

import asyncpg

DSN = os.environ.get("DSN", "postgres://postgres:postgres@localhost:5432/postgres")
HANG_TIMEOUT = float(os.environ.get("HANG_TIMEOUT", "15"))

# Length in bytes of the error message raised by the record generator.
# > 9995 bytes => PostgreSQL rejects the CopyFail frame ("invalid message
# length") and the deadlock triggers; at or below, the run fails cleanly.
# The threshold is exact: the frame's int32 length field (4 bytes, counted in
# itself) + payload + NUL terminator must not exceed PQ_SMALL_MESSAGE_LIMIT
# (10,000), so 10000 - 4 - 1 = 9995.
ERROR_BYTES = int(os.environ.get("ERROR_BYTES", "9996"))

def records():
for i in range(100):
yield (i,)
raise ValueError("A" * ERROR_BYTES)

async def run_case():
conn = await asyncpg.connect(DSN)
try:
async with conn.transaction():
await conn.execute("CREATE TEMP TABLE repro_t (b int)")
await conn.copy_records_to_table("repro_t", records=records())
except ValueError:
# clean failure, bug not reproduced
pass
finally:
await conn.close()

async def main():
task = asyncio.create_task(run_case())
_, pending = await asyncio.wait([task], timeout=HANG_TIMEOUT)
if pending:
print("Statement timed out, bug reproduced", file=sys.stderr)
return 1
print("Statement executed before timeout, bug not reproduced", file=sys.stderr)
return 0

if __name__ == "__main__":
os._exit(asyncio.run(main()))
```

Hướng dẫn đóng góp

Chưa lập chỉ mục được hướng dẫn đóng góp cho kho mã nguồn này

Bắt đầu từ đâu

  1. Đọc hết issue, rồi đọc hướng dẫn đóng góp của dự án.
  2. Bình luận trên issue rằng bạn sẽ nhận — tránh hai người làm cùng một việc.
  3. Fork repository và làm thay đổi trên một nhánh.
  4. Mở pull request có tham chiếu số hiệu của issue.

Hướng nghiên cứu

Bắt đầu từ Connection.copy_records_to_table và chạy repro.py trên PostgreSQL, thay đổi ERROR_BYTES quanh ngưỡng 9995 byte. Theo dõi đường đi của ngoại lệ generator và quá trình rollback của COPY; được xem là hoàn tất khi ngoại lệ kết thúc sạch sẽ mà không bị treo, kể cả với các thông báo ở ngưỡng hoặc cao hơn ngưỡng.

Do mô hình lập chỉ mục viết ra từ nội dung của issue.

Đánh giá

Công nghệ
postgresql, python
Lĩnh vực
databases
Loại issue
Lỗi
Độ khó
4/5
Thời gian dự kiến
3-5 ngày
Mức độ hoạt động
Sôi nổi
Độ rõ ràng
Khá rõ ràng
Mức phù hợp với người mới
48/100

Nhận issue mới trong hộp thư của bạn

Bản tóm tắt ngắn những issue GitHub phù hợp với người mới.