"""The async layer used a thread pool it shared with the whole process. ``asyncio.to_thread`` runs on the event loop's default executor. That executor belongs to the process, not to us, and it is sized from the CPU count — ``min(32, cpu_count + 4)``, which is six threads on a two-CPU container. Every blocking call in ``informix_db.aio`` went through it. Two consequences, both silent. **Cancelled calls hold their threads.** ``asyncio.to_thread`` cannot interrupt a worker, so a cancelled await leaves the thread running the wire call until the read timeout expires. Cancellation is ordinary in a web app — a client disconnect cancels the request task — so a handful of them pins every thread in the shared pool. Unrelated ``to_thread`` work anywhere else in the process then stops dead, and so does the driver. **Pool concurrency was capped by an unrelated number.** A pool with ``max_size=20`` on a two-CPU box ran six queries at a time, and nothing said so. Each connection now owns one thread. That is the right size rather than a compromise: the sync connection serializes every wire operation on its own lock, so a second thread could do nothing but wait for the first. It also avoids a deadlock that a shared pool-sized executor invites — with N threads and N connections, N tasks blocked in ``acquire`` occupy every thread while the connection they wait for is held by a task that now needs a thread to finish and release it. """ from __future__ import annotations import asyncio import threading import time from concurrent.futures import ThreadPoolExecutor import pytest from informix_db import aio from tests.conftest import ConnParams pytestmark = pytest.mark.integration def _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": 10.0, "read_timeout": 20.0, "autocommit": True, } @pytest.mark.asyncio async def test_queries_run_with_the_default_executor_fully_held( conn_params: ConnParams, ) -> None: """The starvation case, with real queries. Four unrelated jobs occupy every thread of a four-worker default executor; the driver must not be waiting behind them.""" loop = asyncio.get_running_loop() previous = loop._default_executor loop.set_default_executor(ThreadPoolExecutor(max_workers=4)) conns = [await aio.connect(**_kw(conn_params)) for _ in range(4)] cursors = [await c.cursor() for c in conns] gate = threading.Event() hogs = [ asyncio.create_task(asyncio.to_thread(gate.wait, 30)) for _ in range(4) ] await asyncio.sleep(0.3) async def query(cursor, i: int): await cursor.execute(f"SELECT FIRST 1 tabid + {i} FROM systables") return await cursor.fetchone() try: rows = await asyncio.wait_for( asyncio.gather(*(query(c, i) for i, c in enumerate(cursors))), timeout=10.0, ) assert len(rows) == 4 assert all(r is not None for r in rows) finally: gate.set() await asyncio.gather(*hogs, return_exceptions=True) for c in conns: await c.close() # set_default_executor rejects None, which is what the loop # starts with before anything has used to_thread. loop.set_default_executor(previous or ThreadPoolExecutor()) @pytest.mark.asyncio async def test_each_connection_gets_its_own_thread( conn_params: ConnParams, ) -> None: a = await aio.connect(**_kw(conn_params)) b = await aio.connect(**_kw(conn_params)) try: assert a._executor is not b._executor assert a._executor._max_workers == 1, ( "more than one thread per connection cannot help — the wire " "lock serializes them anyway" ) finally: await a.close() await b.close() @pytest.mark.asyncio async def test_cursor_runs_on_its_connections_thread( conn_params: ConnParams, ) -> None: """A cursor must not fall back to the default executor, or half the work goes back to being shared.""" conn = await aio.connect(**_kw(conn_params)) try: cur = await conn.cursor() await cur.execute("SELECT FIRST 1 tabid FROM systables") names: list[str] = [] await conn._run(lambda: names.append(threading.current_thread().name)) assert names[0].startswith("informix-conn") assert cur._run.__self__ is conn finally: await conn.close() @pytest.mark.asyncio async def test_pool_reuses_one_thread_per_connection( conn_params: ConnParams, ) -> None: """The executor lives on the sync connection, so a connection handed out, returned, and handed out again keeps the same thread instead of spawning one per acquire.""" pool = await aio.create_pool(**_kw(conn_params), min_size=1, max_size=2) try: seen = [] for _ in range(4): async with pool.connection() as conn: cur = await conn.cursor() await cur.execute("SELECT FIRST 1 tabid FROM systables") assert await cur.fetchone() is not None seen.append(id(conn._executor)) assert len(set(seen)) <= 2, ( f"expected at most one executor per pooled connection, saw " f"{len(set(seen))}" ) finally: await pool.close() @pytest.mark.asyncio async def test_pool_concurrency_is_not_capped_by_the_default_executor( conn_params: ConnParams, ) -> None: """Six concurrent pooled queries against a two-worker default executor. Under the old arrangement at most two could run.""" loop = asyncio.get_running_loop() previous = loop._default_executor pool = await aio.create_pool(**_kw(conn_params), min_size=6, max_size=6) loop.set_default_executor(ThreadPoolExecutor(max_workers=2)) try: async def one(i: int): async with pool.connection(timeout=15.0) as conn: cur = await conn.cursor() # Long enough that serialized execution would be obvious. await cur.execute( "SELECT FIRST 200 a.tabid FROM systables a, systables b" ) return len(await cur.fetchall()) started = time.monotonic() results = await asyncio.wait_for( asyncio.gather(*(one(i) for i in range(6))), timeout=25.0 ) assert all(r > 0 for r in results) assert len(results) == 6 elapsed = time.monotonic() - started assert elapsed < 20.0, f"queries appear serialized ({elapsed:.1f}s)" finally: await pool.close() # set_default_executor rejects None, which is what the loop # starts with before anything has used to_thread. loop.set_default_executor(previous or ThreadPoolExecutor()) @pytest.mark.asyncio async def test_closing_a_connection_stops_its_thread( conn_params: ConnParams, ) -> None: """One thread per connection is only affordable if the thread goes away with the connection.""" before = threading.active_count() conns = [await aio.connect(**_kw(conn_params)) for _ in range(5)] for c in conns: cur = await c.cursor() await cur.execute("SELECT FIRST 1 tabid FROM systables") for c in conns: await c.close() del conns # Threads exit asynchronously after shutdown; give them a moment. for _ in range(50): if threading.active_count() <= before + 1: break await asyncio.sleep(0.1) assert threading.active_count() <= before + 1, ( f"connection threads outlived their connections " f"({before} -> {threading.active_count()})" ) @pytest.mark.asyncio async def test_row_factory_reaches_the_async_paths( conn_params: ConnParams, ) -> None: """The async layer wraps the sync cursor, so row_factory should flow through untouched. Each of the four async fetch routes is a separate call site, and async iteration goes through __anext__ rather than __next__.""" import informix_db conn = await aio.connect(row_factory=informix_db.Row, **_kw(conn_params)) try: cur = await conn.cursor() sql = "SELECT FIRST 2 tabid, tabname FROM systables ORDER BY tabid" await cur.execute(sql) one = await cur.fetchone() assert one[0] == one["tabid"] == one.tabid await cur.execute(sql) assert [r.tabid for r in await cur.fetchall()] == [1, 2] await cur.execute(sql) assert all(r.tabname for r in await cur.fetchmany(2)) await cur.execute(sql) assert [r.tabid async for r in cur] == [1, 2] finally: await conn.close() @pytest.mark.asyncio async def test_async_pool_forwards_the_row_factory( conn_params: ConnParams, ) -> None: import informix_db pool = await aio.create_pool( row_factory=informix_db.Row, min_size=1, max_size=2, **_kw(conn_params) ) try: async with pool.connection() as conn: cur = await conn.cursor() await cur.execute("SELECT FIRST 1 tabid FROM systables") assert (await cur.fetchone()).tabid == 1 finally: await pool.close()