%%% @doc End-to-end test suite for Python asyncio TCP integration. -module(py_async_e2e_SUITE). -include_lib("common_test/include/ct.hrl"). -export([ all/0, init_per_suite/1, end_per_suite/1, init_per_testcase/2, end_per_testcase/2 ]). -export([ test_asyncio_sleep/1, test_asyncio_gather/1, test_asyncio_tcp_echo/1, test_asyncio_concurrent_tcp/1 ]). all() -> [ test_asyncio_sleep, test_asyncio_gather, test_asyncio_tcp_echo, test_asyncio_concurrent_tcp ]. init_per_suite(Config) -> {ok, _} = application:ensure_all_started(erlang_python), %% Ensure contexts are running {ok, _} = py:start_contexts(), %% Install Erlang event loop policy for asyncio.run() Ctx = py:context(1), ok = py:exec(Ctx, <<"import erlang; erlang.install()">>), Config. end_per_suite(_Config) -> ok = application:stop(erlang_python), ok. init_per_testcase(_TestCase, Config) -> Config. end_per_testcase(_TestCase, _Config) -> ok. %% ============================================================================ %% Test Cases - All imports included in the code block %% ============================================================================ %% Test basic asyncio.sleep timing test_asyncio_sleep(_Config) -> Ctx = py:context(1), ok = py:exec(Ctx, <<" import asyncio import time async def timed_sleep(): start = time.monotonic() await asyncio.sleep(0.05) return time.monotonic() - start elapsed = asyncio.run(timed_sleep()) assert elapsed >= 0.04, f'Expected >= 0.04s, got {elapsed:.3f}s' ">>), ok. %% Test asyncio.gather runs tasks concurrently test_asyncio_gather(_Config) -> Ctx = py:context(1), ok = py:exec(Ctx, <<" import asyncio import time async def task(val): await asyncio.sleep(0.05) return val async def main(): start = time.monotonic() r = await asyncio.gather(task(1), task(2), task(3)) elapsed = time.monotonic() - start assert list(r) == [1, 2, 3], f'Expected [1, 2, 3], got {list(r)}' # Allow more time on CI (0.3s instead of 0.15s) assert elapsed < 0.3, f'Expected < 0.3s, got {elapsed:.3f}s' asyncio.run(main()) ">>), ok. %% Test TCP echo server/client using asyncio test_asyncio_tcp_echo(_Config) -> Ctx = py:context(1), ok = py:exec(Ctx, <<" import asyncio async def handler(r, w): data = await r.read(100) w.write(data) await w.drain() w.close() await w.wait_closed() async def test(): srv = await asyncio.start_server(handler, '127.0.0.1', 0) port = srv.sockets[0].getsockname()[1] r, w = await asyncio.open_connection('127.0.0.1', port) w.write(b'hello') await w.drain() resp = await r.read(100) w.close() await w.wait_closed() srv.close() await srv.wait_closed() assert resp == b'hello', f'Expected b\"hello\", got {resp}' asyncio.run(test()) ">>), ok. %% Test multiple concurrent TCP connections test_asyncio_concurrent_tcp(_Config) -> Ctx = py:context(1), ok = py:exec(Ctx, <<" import asyncio async def handler(r, w): data = await r.read(100) w.write(b're:' + data) await w.drain() w.close() await w.wait_closed() async def req(port, msg): r, w = await asyncio.open_connection('127.0.0.1', port) w.write(msg) await w.drain() resp = await r.read(100) w.close() await w.wait_closed() return resp async def test(): srv = await asyncio.start_server(handler, '127.0.0.1', 0) port = srv.sockets[0].getsockname()[1] results = await asyncio.gather( req(port, b'1'), req(port, b'2'), req(port, b'3') ) srv.close() await srv.wait_closed() assert set(results) == {b're:1', b're:2', b're:3'}, f'Expected {{b\"re:1\", b\"re:2\", b\"re:3\"}}, got {set(results)}' asyncio.run(test()) ">>), ok.