informix-db/tests/test_recovery_and_concurrency.py
Ryan Malloy 4af1c90c6f Failed statements bricked the connection; cancelled acquires starved the
pool (2026.08.31.2)

Two more found by fuzzing rather than by users, outside the type system.
Both are ordinary-path bugs a happy-path test cannot reach.

A failed statement was never released. Successful DML sent PREPARE ->
EXECUTE -> RELEASE; failing DML sent PREPARE -> EXECUTE and stopped. The
statement stayed allocated server-side, collided with the next PREPARE,
and every subsequent call on that connection returned a nonsense error
(-255 "Not in transaction" under autocommit, -285 otherwise) whose
offset pointed back at the FAILED sql rather than the new statement.

In practice: one duplicate-key violation bricked the connection. An
INSERT tripping a unique constraint is about the most routine error an
application can hit — every insert-if-not-exists pattern makes them —
and afterwards nothing on that connection worked.

The docstring on the parameterised path already described this hazard
and guarded the parameter-ENCODING failure; it never covered EXECUTE
itself failing. Fixed in all three paths (DML, parameterised DML, and
SELECT — whose fetch loop had the same gap, and whose GC-time finalizer
only covers scrollable cursors).

Cancelling a pool acquire leaked the connection permanently.
asyncio.to_thread cannot interrupt its worker, so a task cancelled while
blocked in pool.acquire() left the worker to finish and hand back a
connection nobody owned — checked out, never returned, one slot gone per
occurrence until the pool was dead.

Not theoretical for anything serving HTTP: a client disconnect cancels
the request task, and under load those cancellations land precisely
while waiting for a connection. The pool dies one slot at a time, only
under load, and PoolTimeoutError points nowhere near the cause.
acquire() now shields the inner future and returns whatever the worker
produced if the caller went away; add_done_callback fires immediately on
an already-resolved future, so the "worker finished just before the
cancel" race is the same code path rather than a second one.

The harness covered three things the type fuzzer cannot see: fetch
batching (NFETCH is a 4096-BYTE budget, so batch edges move with row
width — counts 0..1025 across three schema widths, every fetch style
required to agree), error recovery (eight error classes, repeated, each
followed by a known-good query on the same cursor and a fresh one), and
concurrency (threads sharing one connection, pooled threads that fail
before releasing, transaction isolation across checkouts, async
cancellation). Every worker reads back a value only it supplied, so a
crossed wire fails on data even when nothing raises.

Two harness false positives worth recording: Informix silently truncates
over-long strings (a literal SQL insert does the same, so the driver
matches the server), and it rejects a bare `?` in a projection — a cast
is required.

356/356 integration on 15, 14.10 and 12.10 (was 326).
2026-08-31 17:13:43 -06:00

460 lines
17 KiB
Python

"""Regression tests for error recovery, cursor lifecycle, and concurrency.
The type-matrix suite covers what a value looks like on the wire. This
covers everything around it: what happens after a statement fails, what
happens to abandoned cursors, and what happens when more than one thread
or task is involved.
Two bugs here were found by fuzzing, and both are the kind that never
show up in a single-threaded happy-path test:
* **A failed statement was never released.** Successful DML sent
PREPARE → EXECUTE → RELEASE; failing DML sent PREPARE → EXECUTE and
stopped. The leaked statement collided with the next PREPARE, and
every subsequent call on that connection returned a nonsense error
whose offset pointed back at the *failed* SQL. A duplicate-key
violation — about the most ordinary error an application can hit —
killed the connection outright.
* **Cancelling a pool acquire leaked the connection.** ``asyncio.to_thread``
cannot interrupt its worker, so a cancelled waiter left the worker to
finish and hand back a connection nobody owned. Under HTTP load, where
client disconnects cancel request tasks precisely while they wait for
a connection, the pool dies one slot at a time.
The tests that follow assert on *data* wherever they can, not merely the
absence of an exception: a worker reads back a value only it wrote, so a
crossed wire fails even when nothing raises.
"""
from __future__ import annotations
import asyncio
import contextlib
import threading
import pytest
import informix_db
from informix_db import aio
from tests.conftest import ConnParams
pytestmark = pytest.mark.integration
def _connect(conn_params: ConnParams, **kw) -> informix_db.Connection:
return informix_db.connect(
host=conn_params.host,
port=conn_params.port,
user=conn_params.user,
password=conn_params.password,
database=conn_params.database,
server=conn_params.server,
connect_timeout=10.0,
# Bounded: a leaked statement used to manifest as a hang.
read_timeout=25.0,
**kw,
)
def _pool_kw(conn_params: ConnParams) -> dict:
return {
"host": conn_params.host,
"port": conn_params.port,
"user": conn_params.user,
"password": conn_params.password,
"database": conn_params.database,
"server": conn_params.server,
"connect_timeout": 15.0,
"read_timeout": 30.0,
}
# ---------------------------------------------------------------------------
# Error recovery — a failed statement must not poison the connection
# ---------------------------------------------------------------------------
FAILING_SQL = [
pytest.param("SELECT FROM WHERE", id="syntax"),
pytest.param("SELECT * FROM no_such_table_xyz", id="no_such_table"),
pytest.param("SELECT no_such_col FROM systables", id="no_such_column"),
pytest.param("SELECT no_such_function_xyz(1) FROM systables", id="bad_func"),
pytest.param("SELECT 1/0 FROM systables", id="div_zero"),
pytest.param("SELECT 'notanumber'::INT FROM systables", id="bad_cast"),
]
@pytest.mark.parametrize("sql", FAILING_SQL)
def test_connection_survives_failed_statement(
conn_params: ConnParams, sql: str
) -> None:
with _connect(conn_params, autocommit=True) as conn:
cur = conn.cursor()
for _ in range(4): # repeat: a leak accumulates
with pytest.raises(informix_db.Error):
cur.execute(sql)
cur.fetchall()
cur.execute("SELECT FIRST 1 tabid FROM systables ORDER BY tabid")
assert cur.fetchone() is not None
@pytest.mark.parametrize("autocommit", [True, False])
def test_duplicate_key_does_not_kill_the_connection(
conn_params: ConnParams, autocommit: bool
) -> None:
"""The bug that started this file. A unique-constraint violation left
the prepared statement allocated server-side; the next statement then
failed with a nonsense error (-255 "Not in transaction" in autocommit,
-285 otherwise) and the connection never recovered."""
with _connect(conn_params, autocommit=autocommit) as conn:
cur = conn.cursor()
cur.execute("CREATE TEMP TABLE t_dup (k INT PRIMARY KEY)")
cur.execute("INSERT INTO t_dup VALUES (1)")
if not autocommit:
conn.commit()
for i in range(5):
with pytest.raises(informix_db.IntegrityError):
cur.execute("INSERT INTO t_dup VALUES (1)")
if not autocommit:
conn.rollback()
cur.execute("SELECT k FROM t_dup")
assert cur.fetchall() == [(1,)], f"broken after {i + 1} violations"
def test_failed_dml_with_params_releases_statement(
conn_params: ConnParams,
) -> None:
with _connect(conn_params, autocommit=True) as conn:
cur = conn.cursor()
cur.execute("CREATE TEMP TABLE t_nn (s VARCHAR(8) NOT NULL, k INT)")
for _ in range(4):
with pytest.raises(informix_db.Error):
cur.execute("INSERT INTO t_nn VALUES (?, ?)", (None, 1))
cur.execute("SELECT COUNT(*) FROM t_nn")
assert cur.fetchone() == (0,)
def test_new_cursor_works_after_error(conn_params: ConnParams) -> None:
with _connect(conn_params, autocommit=True) as conn:
cur = conn.cursor()
with pytest.raises(informix_db.Error):
cur.execute("SELECT * FROM no_such_table_xyz")
other = conn.cursor()
other.execute("SELECT FIRST 1 tabid FROM systables")
assert other.fetchone() is not None
other.close()
# ---------------------------------------------------------------------------
# Fetch batching — NFETCH is a byte budget, so batch edges move with width
# ---------------------------------------------------------------------------
@pytest.mark.parametrize("n", [0, 1, 2, 3, 63, 64, 65, 199, 200, 201, 512, 1025])
def test_row_count_exact_across_fetch_styles(
conn_params: ConnParams, n: int
) -> None:
"""Every fetch style must agree, at counts that straddle plausible
server batch boundaries. A batching off-by-one shows up as a wrong
count, a duplicate, or a dropped row."""
with _connect(conn_params, autocommit=True) as conn:
cur = conn.cursor()
cur.execute(
"CREATE TEMP TABLE t_batch "
"(k INT, a VARCHAR(32), b INT, c VARCHAR(32), d INT)"
)
if n:
cur.executemany(
"INSERT INTO t_batch VALUES (?, ?, ?, ?, ?)",
[(i, f"a{i}", i * 2, f"c{i}", i * 3) for i in range(n)],
)
want = [(i,) for i in range(n)]
sql = "SELECT k FROM t_batch ORDER BY k"
cur.execute(sql)
assert cur.fetchall() == want
cur.execute(sql)
got = []
while (row := cur.fetchone()) is not None:
got.append(row)
assert got == want
for size in (1, 7, 100):
cur.execute(sql)
got = []
while chunk := cur.fetchmany(size):
got.extend(chunk)
assert got == want, f"fetchmany({size}) disagreed at n={n}"
cur.execute(sql)
assert list(cur) == want
# ---------------------------------------------------------------------------
# Cursor lifecycle
# ---------------------------------------------------------------------------
def test_abandoned_partial_fetches_do_not_leak(
conn_params: ConnParams,
) -> None:
with _connect(conn_params, autocommit=True) as conn:
cur = conn.cursor()
cur.execute("CREATE TEMP TABLE t_aband (k INT)")
cur.executemany(
"INSERT INTO t_aband VALUES (?)", [(i,) for i in range(200)]
)
for i in range(25):
c = conn.cursor()
c.execute("SELECT k FROM t_aband ORDER BY k")
c.fetchone()
if i % 2:
c.close()
else:
del c # rely on the finalizer
cur.execute("SELECT COUNT(*) FROM t_aband")
assert cur.fetchone() == (200,)
def test_re_execute_mid_fetch(conn_params: ConnParams) -> None:
with _connect(conn_params, autocommit=True) as conn:
cur = conn.cursor()
cur.execute("CREATE TEMP TABLE t_reex (k INT)")
cur.executemany(
"INSERT INTO t_reex VALUES (?)", [(i,) for i in range(100)]
)
for _ in range(10):
cur.execute("SELECT k FROM t_reex ORDER BY k")
cur.fetchone() # abandon mid-fetch
cur.execute("SELECT k FROM t_reex ORDER BY k")
assert len(cur.fetchall()) == 100
def test_interleaved_cursors_on_one_connection(
conn_params: ConnParams,
) -> None:
with _connect(conn_params, autocommit=True) as conn:
cur = conn.cursor()
cur.execute("CREATE TEMP TABLE t_inter (k INT)")
cur.executemany(
"INSERT INTO t_inter VALUES (?)", [(i,) for i in range(50)]
)
a, b = conn.cursor(), conn.cursor()
a.execute("SELECT k FROM t_inter ORDER BY k")
b.execute("SELECT k FROM t_inter ORDER BY k DESC")
assert [a.fetchone() for _ in range(5)] == [(i,) for i in range(5)]
assert [b.fetchone() for _ in range(5)] == [
(i,) for i in range(49, 44, -1)
]
a.close()
b.close()
# ---------------------------------------------------------------------------
# Concurrency
# ---------------------------------------------------------------------------
def test_threads_sharing_one_connection_do_not_interleave(
conn_params: ConnParams,
) -> None:
"""One connection is one socket. Each thread reads back a value only
it supplied, so crossed wires fail even if nothing raises."""
failures: list[str] = []
lock = threading.Lock()
with _connect(conn_params, autocommit=True) as conn:
def worker(tid: int) -> None:
try:
for r in range(10):
cur = conn.cursor()
token = tid * 1000 + r
cur.execute(
"SELECT FIRST 1 ?::INT, ?::VARCHAR(16) FROM systables",
(token, f"t{tid}"),
)
row = cur.fetchone()
if row != (token, f"t{tid}"):
with lock:
failures.append(f"thread {tid}: got {row!r}")
return
cur.close()
except Exception as exc:
with lock:
failures.append(f"thread {tid}: {type(exc).__name__}: {exc}")
threads = [threading.Thread(target=worker, args=(i,)) for i in range(6)]
for t in threads:
t.start()
for t in threads:
t.join()
assert not failures, failures
def test_pool_hands_out_clean_connections_under_load(
conn_params: ConnParams,
) -> None:
"""Borrowers that fail must still return a usable connection."""
failures: list[str] = []
lock = threading.Lock()
pool = informix_db.create_pool(
**_pool_kw(conn_params), min_size=2, max_size=4,
acquire_timeout=30.0, autocommit=True,
)
try:
def worker(tid: int) -> None:
try:
for r in range(8):
with pool.connection() as conn:
cur = conn.cursor()
if r % 3 == 0:
with contextlib.suppress(informix_db.Error):
cur.execute("SELECT * FROM no_such_tbl_xyz")
token = tid * 1000 + r
cur.execute(
"SELECT FIRST 1 ?::INT FROM systables", (token,)
)
got = cur.fetchone()
if got != (token,):
with lock:
failures.append(
f"thread {tid}: got {got!r}, want {token}"
)
return
cur.close()
except Exception as exc:
with lock:
failures.append(f"thread {tid}: {type(exc).__name__}: {exc}")
threads = [threading.Thread(target=worker, args=(i,)) for i in range(6)]
for t in threads:
t.start()
for t in threads:
t.join()
finally:
with contextlib.suppress(Exception):
pool.close()
assert not failures, failures
def test_pool_does_not_leak_transactions_between_borrowers(
conn_params: ConnParams,
) -> None:
"""max_size=1 so both borrowers get the same underlying connection."""
pool = informix_db.create_pool(
**_pool_kw(conn_params), min_size=1, max_size=1,
acquire_timeout=30.0, autocommit=False,
)
try:
with pool.connection() as conn:
cur = conn.cursor()
with contextlib.suppress(Exception):
cur.execute("DROP TABLE t_pool_txn")
cur.execute("CREATE TABLE t_pool_txn (k INT)")
conn.commit()
try:
for i in range(5):
with pool.connection() as conn: # dirties, never commits
conn.cursor().execute(
"INSERT INTO t_pool_txn VALUES (?)", (i,)
)
with pool.connection() as conn: # must not see it
cur = conn.cursor()
cur.execute("SELECT COUNT(*) FROM t_pool_txn")
assert cur.fetchone() == (0,), "transaction leaked"
conn.rollback()
finally:
with pool.connection() as conn:
cur = conn.cursor()
with contextlib.suppress(Exception):
cur.execute("DROP TABLE t_pool_txn")
conn.commit()
finally:
with contextlib.suppress(Exception):
pool.close()
# ---------------------------------------------------------------------------
# Async
# ---------------------------------------------------------------------------
async def test_async_pool_survives_cancelled_acquires(
conn_params: ConnParams,
) -> None:
"""Cancelling a task blocked on acquire used to destroy that pool slot
permanently: ``asyncio.to_thread`` cannot interrupt its worker, so the
worker finished and handed back a connection nobody owned.
Two holders keep both connections checked out; three waiters block on
acquire and are cancelled there. Afterwards the pool must still serve.
"""
pool = await aio.create_pool(
**_pool_kw(conn_params), min_size=1, max_size=2, autocommit=True
)
try:
release = asyncio.Event()
async def holder(i: int) -> None:
async with pool.connection() as conn:
cur = await conn.cursor()
await cur.execute(
"SELECT FIRST 1 ?::INT FROM systables", (i,)
)
await cur.fetchone()
await release.wait()
async def waiter(i: int) -> tuple | None:
async with pool.connection() as conn:
cur = await conn.cursor()
await cur.execute(
"SELECT FIRST 1 ?::INT FROM systables", (i,)
)
return await cur.fetchone()
for cycle in range(2):
release.clear()
holders = [asyncio.create_task(holder(i)) for i in range(2)]
await asyncio.sleep(0.5) # both connections held
waiters = [asyncio.create_task(waiter(100 + i)) for i in range(3)]
await asyncio.sleep(0.5) # all blocked on acquire
for t in waiters:
t.cancel()
await asyncio.gather(*waiters, return_exceptions=True)
release.set()
await asyncio.gather(*holders, return_exceptions=True)
await asyncio.sleep(0.3) # orphan returns land
got = await asyncio.wait_for(waiter(200 + cycle), timeout=15)
assert got == (200 + cycle,), (
f"pool starved after cycle {cycle} — cancelled acquires leaked"
)
finally:
with contextlib.suppress(Exception):
await pool.close()
async def test_async_pool_concurrent_tasks_get_their_own_data(
conn_params: ConnParams,
) -> None:
pool = await aio.create_pool(
**_pool_kw(conn_params), min_size=2, max_size=4, autocommit=True
)
try:
async def worker(tid: int) -> None:
for r in range(6):
async with pool.connection() as conn:
cur = await conn.cursor()
token = tid * 1000 + r
await cur.execute(
"SELECT FIRST 1 ?::INT FROM systables", (token,)
)
assert await cur.fetchone() == (token,)
await asyncio.gather(*(worker(i) for i in range(6)))
finally:
with contextlib.suppress(Exception):
await pool.close()