MagicStack / MagicStack/asyncpg

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

Aperta
#1,354 0 commenti 0 reazioni 0 assegnatari Vedi su GitHub

Nessuno ha ancora preso questa issue.

Lingua principale
Python
Stelle
8.1k
Fork
468
Metriche di merge delle PR
Nessuna PR unita negli ultimi 30g

Descrizione

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()))
```

Guida per i contributori

Nessuna guida per i contributori indicizzata per questo repository

Come iniziare

  1. Leggi tutta la issue e poi la guida ai contributi del progetto.
  2. Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
  3. Fai un fork del repository e lavora su un branch.
  4. Apri una pull request che faccia riferimento al numero della issue.

Direzione di ricerca

Inizia da Connection.copy_records_to_table ed esegui repro.py su PostgreSQL, variando ERROR_BYTES intorno alla soglia di 9995 byte. Segui il percorso dell’eccezione del generatore e del rollback di COPY; il lavoro è completato quando l’eccezione termina correttamente senza bloccarsi, anche per i messaggi pari o superiori alla soglia.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Valutazione

Stack tecnologico
postgresql, python
Ambito
databases
Tipo di issue
Bug
Difficoltà
4/5
Tempo stimato
3-5 giorni
Stato di attività
Attiva
Chiarezza
Abbastanza chiara
Idoneità per principianti
48/100

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.