Skip to content
Open
Show file tree
Hide file tree
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 @@ -162,7 +162,7 @@ async def _create_session(self, livekit_url: str, livekit_token: str, room_name:
if is_given(self._max_duration):
body["maxDuration"] = self._max_duration

for attempt in range(self._conn_options.max_retry):
for attempt in range(self._conn_options.max_retry + 1):

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔴 Malformed responses create duplicate Runway sessions

With max_retry=1, a malformed 200 response makes _create_session retry after Runway creates the first session. Only the second session ID is saved, leaving the first session active without a cancellation path.

Learn more

A successful HTTP POST can create a realtime session even when its response body cannot be decoded. In _create_session, response.json() runs inside the broad exception handler. A decode error enters the same retry path as a failed request, so another POST creates a second session. Cleanup via _cancel_runway_realtime_session knows only the last successfully decoded ID and cannot cancel the earlier session.

Example: Set max_retry=1. Runway creates session A but returns a malformed 200 body. The retry creates session B and saves B's ID. On disconnection, cleanup cancels B while A remains active.

Recommended fix: Distinguish invalid responses to successful POSTs from retryable connection or status errors. Fail without re-POSTing when JSON decoding fails or a successful response lacks a valid session ID; preserve the underlying error so the caller can diagnose the incomplete creation.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch — confirmed against the code: response.json() is inside the try, so a decode error on a 2xx falls into the generic except Exception and the loop POSTs again, while only the last decoded id is kept in _realtime_session_id for _cancel_runway_realtime_session. A 2xx whose payload carries no string id is a second, quieter variant: it returns with no session id at all.

The reachability is partly pre-existing — before this change the loop made max_retry requests, so a decode error already re-POSTed whenever max_retry >= 2; this PR shifts that to max_retry >= 1 and adds one attempt everywhere, so I won't pretend the amplification is unrelated to the diff.

I'd rather keep this PR scoped to #7604 (send the initial request when retries are disabled) and treat "don't re-POST when the request was accepted but the payload is unusable" as its own change, since it alters the retryable/non-retryable classification rather than the request count — that's also how the Tavus side handled it in #7494. Happy to open that follow-up (plus the same classification for bey/bithuman/avatario/liveavatar, which still carry the off-by-one) if that shape is welcome here.

try:
async with self._ensure_http_session().post(
f"{self._api_url}/v1/realtime_sessions",
Expand Down Expand Up @@ -198,7 +198,7 @@ async def _create_session(self, livekit_url: str, livekit_token: str, room_name:
else:
logger.exception("failed to call Runway avatar API")

if attempt < self._conn_options.max_retry - 1:
if attempt < self._conn_options.max_retry:
await asyncio.sleep(self._conn_options.retry_interval)

raise APIConnectionError("Failed to start Runway Avatar Session after all retries")
Expand Down
109 changes: 109 additions & 0 deletions tests/test_plugin_runway_avatar.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,109 @@
"""Tests for the Runway avatar plugin: the retry budget of the session request.

`APIConnectOptions.max_retry` counts retries, not requests, so `max_retry=0` still
owes one initial POST — the contract the core avatar gateway, Anam, Tavus and
Synthesia keep (#7604).
"""

from __future__ import annotations

from typing import Any

import aiohttp
import pytest

from livekit.agents import APIConnectionError, APIConnectOptions, APIStatusError
from livekit.plugins.runway import AvatarSession

pytestmark = pytest.mark.unit

_ROOM = "room-name"


class _Response:
def __init__(self, status: int) -> None:
self.status = status
self.ok = status < 400

async def text(self) -> str:
return f"status {self.status}"

async def json(self) -> dict[str, Any]:
return {"id": "realtime-session-1"}

async def __aenter__(self) -> _Response:
return self

async def __aexit__(self, *exc: object) -> None:
return None


class _ScriptedSession:
"""Answers each POST with the next step in `script`; an exception instance in the
script is raised instead, and the last step repeats."""

def __init__(self, script: list[int | BaseException]) -> None:
self._script = list(script)
self.posts = 0

def post(self, url: str, **kwargs: Any) -> _Response:
self.posts += 1
step = self._script.pop(0) if len(self._script) > 1 else self._script[0]
if isinstance(step, BaseException):
raise step
return _Response(step)


def _avatar(session: _ScriptedSession, max_retry: int) -> AvatarSession:
avatar = AvatarSession(
preset_id="preset-id",
api_key="test-key",
conn_options=APIConnectOptions(max_retry=max_retry, retry_interval=0.0),
)
avatar._http_session = session # type: ignore[assignment]
avatar._local_participant_identity = "local-agent"
return avatar


async def test_zero_retries_still_sends_the_initial_request():
"""`range(max_retry)` with max_retry=0 looped zero times, so startup failed
without ever asking the API — retries were disabled along with the request."""
session = _ScriptedSession([200])

await _avatar(session, max_retry=0)._create_session("wss://livekit.test", "token", _ROOM)

assert session.posts == 1


@pytest.mark.parametrize("max_retry", [0, 1, 3])
async def test_a_persistent_connection_error_is_retried_max_retry_times(
max_retry: int,
):
"""One initial request plus `max_retry` retries, and the last error is kept."""
session = _ScriptedSession([aiohttp.ClientConnectionError("connection refused")])

with pytest.raises(APIConnectionError):
await _avatar(session, max_retry=max_retry)._create_session(
"wss://livekit.test", "token", _ROOM
)

assert session.posts == max_retry + 1


async def test_a_retryable_error_followed_by_success_starts_the_session():
session = _ScriptedSession([aiohttp.ClientConnectionError("connection refused"), 200])

await _avatar(session, max_retry=1)._create_session("wss://livekit.test", "token", _ROOM)

assert session.posts == 2


async def test_a_client_error_is_not_retried():
"""A bad key or an unknown preset fails the same way every time."""
session = _ScriptedSession([401])

with pytest.raises(APIStatusError) as exc_info:
await _avatar(session, max_retry=3)._create_session("wss://livekit.test", "token", _ROOM)

assert session.posts == 1
assert exc_info.value.status_code == 401
Loading