MagicStack / MagicStack/asyncpg

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

未关闭
#1,354 0 条评论 0 个 reaction 已指派 0 人 在 GitHub 查看

还没有人认领这个 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()))

贡献指南

这个仓库没有索引到贡献指南

从这里开始

  1. 先读完整个 Issue,再读项目的贡献指南。
  2. 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
  3. Fork 仓库,在一个分支上完成修改。
  4. 提交 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

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。