MagicStack / MagicStack/asyncpg
Connection.copy_records_to_table hanging when generator passed to records parameter raises exception with long error message
未关闭
还没有人认领这个 Issue。
- 主要语言
- Python
- 星标
- 8.1k
- 派生
- 468
- PR 合并指标
- 30 天内没有已合并 PR
描述
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
#!/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()))
贡献指南
这个仓库没有索引到贡献指南
从这里开始
- 先读完整个 Issue,再读项目的贡献指南。
- 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
- Fork 仓库,在一个分支上完成修改。
- 提交 Pull Request,并在描述里引用这个 Issue 编号。
调研方向
从 Connection.copy_records_to_table 开始,针对 PostgreSQL 运行 repro.py,并在 9995 字节的阈值附近改变 ERROR_BYTES。跟踪生成器异常和 COPY 回滚路径;当异常能够干净地完成且不会挂起时即表示完成,包括消息长度达到或超过该阈值的情况。
由索引模型根据 Issue 内容生成。
评估
- 技术栈
- postgresql, python
- 领域
- databases
- Issue 类型
- 缺陷
- 难度
- 4/5
- 预计耗时
- 3-5 天
- 活跃度
- 活跃
- 描述清晰度
- 基本清楚
- 新手友好度
- 48/100