Skip to content
Merged
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
5 changes: 5 additions & 0 deletions .changeset/breezy-rockets-camp.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"livekit-agents": patch
---

pipelineagent: fix speech_committed never called
21 changes: 14 additions & 7 deletions livekit-agents/livekit/agents/pipeline/pipeline_agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -624,8 +624,8 @@ async def _synthesize_answer_task(
playing_speech = self._playing_speech
if playing_speech is not None and playing_speech.initialized:
if (
not playing_speech.user_question or playing_speech.user_commited
) and not playing_speech.speech_commited:
not playing_speech.user_question or playing_speech.user_committed
) and not playing_speech.speech_committed:
# the speech is playing but not committed yet, add it to the chat context for this new reply synthesis
copied_ctx.messages.append(
ChatMessage.create(
Expand Down Expand Up @@ -681,7 +681,7 @@ def _commit_user_question_if_needed() -> None:
if (
not user_question
or synthesis_handle.interrupted
or speech_handle.user_commited
or speech_handle.user_committed
):
return

Expand All @@ -707,7 +707,7 @@ def _commit_user_question_if_needed() -> None:
self.emit("user_speech_committed", user_msg)

self._transcribed_text = self._transcribed_text[len(user_question) :]
speech_handle.mark_user_commited()
speech_handle.mark_user_committed()

# wait for the play_handle to finish and check every 1s if the user question should be committed
_commit_user_question_if_needed()
Expand Down Expand Up @@ -737,11 +737,18 @@ def _commit_user_question_if_needed() -> None:
if is_using_tools and not interrupted:
assert isinstance(speech_handle.source, LLMStream)
assert (
not user_question or speech_handle.user_commited
not user_question or speech_handle.user_committed
), "user speech should have been committed before using tools"

llm_stream = speech_handle.source

if collected_text:
msg = ChatMessage.create(text=collected_text, role="assistant")
self._chat_ctx.messages.append(msg)

speech_handle.mark_speech_committed()
self.emit("agent_speech_committed", msg)

# execute functions
call_ctx = AgentCallContext(self, llm_stream)
tk = _CallContextVar.set(call_ctx)
Expand Down Expand Up @@ -822,7 +829,7 @@ def _commit_user_question_if_needed() -> None:
_CallContextVar.reset(tk)

if speech_handle.add_to_chat_ctx and (
not user_question or speech_handle.user_commited
not user_question or speech_handle.user_committed
):
self._chat_ctx.messages.extend(extra_tools_messages)

Expand All @@ -832,7 +839,7 @@ def _commit_user_question_if_needed() -> None:
msg = ChatMessage.create(text=collected_text, role="assistant")
self._chat_ctx.messages.append(msg)

speech_handle.mark_speech_commited()
speech_handle.mark_speech_committed()

if interrupted:
self.emit("agent_speech_interrupted", msg)
Expand Down
20 changes: 10 additions & 10 deletions livekit-agents/livekit/agents/pipeline/speech_handle.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,11 +25,11 @@ def __init__(
# is_reply is True when the speech is answering to a user question
self._is_reply = is_reply
self._user_question = user_question
self._user_commited = False
self._user_committed = False

self._init_fut: asyncio.Future[None] = asyncio.Future()
self._initialized = False
self._speech_commited = False # speech committed (interrupted or not)
self._speech_committed = False # speech committed (interrupted or not)

# source and synthesis_handle are None until the speech is initialized
self._source: str | LLMStream | AsyncIterable[str] | None = None
Expand Down Expand Up @@ -81,19 +81,19 @@ def initialize(
self._initialized = True
self._init_fut.set_result(None)

def mark_user_commited(self) -> None:
self._user_commited = True
def mark_user_committed(self) -> None:
self._user_committed = True

def mark_speech_commited(self) -> None:
self._speech_commited = True
def mark_speech_committed(self) -> None:
self._speech_committed = True

@property
def user_commited(self) -> bool:
return self._user_commited
def user_committed(self) -> bool:
return self._user_committed

@property
def speech_commited(self) -> bool:
return self._speech_commited
def speech_committed(self) -> bool:
return self._speech_committed

@property
def id(self) -> str:
Expand Down