Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@
APIConnectOptions,
NotGivenOr,
)
from livekit.agents.utils import is_given
from livekit.agents.utils import aio, is_given
from livekit.agents.voice.io import TimedString

from .log import logger
Expand Down Expand Up @@ -964,6 +964,7 @@ def __init__(
self._idle_connection_timeout = idle_connection_timeout
self._pool: _ConnectionPool | None = None
self._pool_lock = asyncio.Lock()
self._prewarm_task: asyncio.Task[None] | None = None
self._streams = weakref.WeakSet[SynthesizeStream]()
self._sentence_tokenizer = (
tokenizer
Expand Down Expand Up @@ -1087,7 +1088,7 @@ def _ensure_session(self) -> aiohttp.ClientSession:
return self._session

def prewarm(self) -> None:
asyncio.create_task(self._prewarm_impl())
self._prewarm_task = asyncio.create_task(self._prewarm_impl())

async def _prewarm_impl(self) -> None:
# Just ensure the pool is created - first acquire will establish a connection
Expand All @@ -1109,6 +1110,13 @@ def stream(
return stream

async def aclose(self) -> None:
# Cancel first: _prewarm_impl calls _get_pool, which builds a new pool when
# self._pool is None. A prewarm still in flight here would otherwise recreate
# the pool after it was closed, leaving a live connection nothing owns.
if self._prewarm_task is not None:
await aio.cancel_and_wait(self._prewarm_task)
self._prewarm_task = None

for stream in list(self._streams):
await stream.aclose()

Expand Down