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 @@ -390,6 +390,7 @@ def __init__(
self._session = http_session
self._reconnect_event = asyncio.Event()
self._speaking = False # Track if we're currently in a speech segment
self._last_partial_text = ""
self._audio_duration_collector = PeriodicCollector(
callback=self._on_audio_duration_report,
duration=5.0,
Expand Down Expand Up @@ -538,6 +539,7 @@ async def recv_task(ws: aiohttp.ClientWebSocketResponse) -> None:
while True:
try:
ws = await self._connect_ws()
self._last_partial_text = ""
if self._opts.previous_text:
# Must be the first input_audio_chunk on the connection.
await ws.send_str(
Expand Down Expand Up @@ -676,7 +678,9 @@ def _process_stream_event(self, data: dict) -> None:
if message_type == "partial_transcript":
logger.debug("Received message type partial_transcript: %s", data)

if text:
if text and text != self._last_partial_text:
self._last_partial_text = text

# Send START_OF_SPEECH if we're not already speaking
if not self._speaking:
self._event_ch.send_nowait(
Expand All @@ -697,6 +701,8 @@ def _process_stream_event(self, data: dict) -> None:
):
# Final committed transcripts - these are sent to the LLM/TTS layer in LiveKit agents
# and trigger agent responses (unlike partial transcripts which are UI-only)
self._last_partial_text = ""

if text:
# Send START_OF_SPEECH if we're not already speaking
if not self._speaking:
Expand Down
55 changes: 55 additions & 0 deletions tests/test_plugin_elevenlabs_stt.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ def _new_stream(*, server_vad=NOT_GIVEN) -> elevenlabs_stt.SpeechStream:
stream._language = None
stream._event_ch = _EventSink()
stream._speaking = False
stream._last_partial_text = ""
stream._start_time_offset = 0.0
return stream

Expand All @@ -66,6 +67,60 @@ def _committed_transcript(text: str) -> dict:
}


def _partial_transcript(text: str) -> dict:
return {"message_type": "partial_transcript", "text": text, "words": []}


def _interim_texts(stream: elevenlabs_stt.SpeechStream) -> list[str]:
return [
event.alternatives[0].text
for event in stream._event_ch.events
if event.type == stt.SpeechEventType.INTERIM_TRANSCRIPT
]


def test_advancing_partial_transcripts_are_forwarded() -> None:
stream = _new_stream(server_vad={"vad_silence_threshold_secs": 0.5})

stream._process_stream_event(_partial_transcript("yeah"))
stream._process_stream_event(_partial_transcript("yeah please"))

assert _interim_texts(stream) == ["yeah", "yeah please"]


def test_re_sent_partial_transcript_is_dropped() -> None:
stream = _new_stream(server_vad={"vad_silence_threshold_secs": 0.5})

for _ in range(5):
stream._process_stream_event(_partial_transcript("yeah please"))

assert _interim_texts(stream) == ["yeah please"]
assert [event.type for event in stream._event_ch.events] == [
stt.SpeechEventType.START_OF_SPEECH,
stt.SpeechEventType.INTERIM_TRANSCRIPT,
]


def test_the_same_words_after_a_commit_are_forwarded_again() -> None:
stream = _new_stream(server_vad={"vad_silence_threshold_secs": 0.5})

stream._process_stream_event(_partial_transcript("right"))
stream._process_stream_event(_committed_transcript("right"))
stream._process_stream_event(_partial_transcript("right"))

assert _interim_texts(stream) == ["right", "right"]


def test_the_same_words_after_an_empty_commit_are_forwarded_again() -> None:
stream = _new_stream(server_vad=None)

stream._process_stream_event(_partial_transcript("right"))
stream._process_stream_event(_committed_transcript(""))
stream._process_stream_event(_partial_transcript("right"))

assert _interim_texts(stream) == ["right", "right"]


def test_server_vad_commit_emits_end_of_speech() -> None:
stream = _new_stream(server_vad={"vad_silence_threshold_secs": 0.5})

Expand Down