From 641faa96623b185d4ce0b065d81204e5fc039f8d Mon Sep 17 00:00:00 2001 From: Arun Mittal Date: Wed, 10 Jun 2026 15:32:21 -0400 Subject: [PATCH] refactor(proxy): use asyncio.create_task in _iter_with_keepalive MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit asyncio.ensure_future is soft-deprecated since Python 3.10 for the "schedule a coroutine on the running loop" use case. create_task is the modern, idiomatic equivalent and is what every other call site in proxy_server.py uses — the ensure_future call in _iter_with_keepalive was the outlier. Functionally identical here since aiter.__anext__() always returns a coroutine. Test updated to spy on asyncio.create_task instead of ensure_future. --- litellm/proxy/proxy_server.py | 2 +- .../proxy/test_async_data_generator_keepalive.py | 8 ++++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index cfa56944f23..3f3bac4942b 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -7051,7 +7051,7 @@ async def _iter_with_keepalive(aiter, keepalive_seconds: float): try: while True: if pending is None: - pending = asyncio.ensure_future(aiter.__anext__()) + pending = asyncio.create_task(aiter.__anext__()) done, _ = await asyncio.wait({pending}, timeout=keepalive_seconds) if not done: yield _STREAM_KEEPALIVE diff --git a/tests/test_litellm/proxy/test_async_data_generator_keepalive.py b/tests/test_litellm/proxy/test_async_data_generator_keepalive.py index e211af90264..e7ff5b02159 100644 --- a/tests/test_litellm/proxy/test_async_data_generator_keepalive.py +++ b/tests/test_litellm/proxy/test_async_data_generator_keepalive.py @@ -273,10 +273,10 @@ def test_iter_with_keepalive_cancels_pending_task_on_early_close(): # Track the Tasks ``_iter_with_keepalive`` creates so we can assert on # cancellation after the wrapper is closed. created_tasks = [] - original_ensure_future = asyncio.ensure_future + original_create_task = asyncio.create_task - def _spy_ensure_future(coro_or_future, *args, **kwargs): - task = original_ensure_future(coro_or_future, *args, **kwargs) + def _spy_create_task(coro, *args, **kwargs): + task = original_create_task(coro, *args, **kwargs) created_tasks.append(task) return task @@ -287,7 +287,7 @@ def test_iter_with_keepalive_cancels_pending_task_on_early_close(): yield "unreachable" # pragma: no cover async def _run_test(): - with patch.object(asyncio, "ensure_future", _spy_ensure_future): + with patch.object(asyncio, "create_task", _spy_create_task): wrapper = proxy_server_module._iter_with_keepalive( _never_yields().__aiter__(), keepalive_seconds=0.05 )