Merge pull request #3486 from pipecat-ai/pk/fix-nova-sonic-reset-conversation
Fix `AWSNovaSonicLLMService.reset_conversation()`
This commit is contained in:
1
changelog/3486.fixed.md
Normal file
1
changelog/3486.fixed.md
Normal file
@@ -0,0 +1 @@
|
|||||||
|
- Fixed `AWSNovaSonicLLMService.reset_conversation()`, which would previously error out. Now it successfully reconnects and "rehydrates" from the context object.
|
||||||
@@ -113,6 +113,14 @@ async def load_conversation(params: FunctionCallParams):
|
|||||||
# "content": f"{AWSNovaSonicLLMService.AWAIT_TRIGGER_ASSISTANT_RESPONSE_INSTRUCTION}",
|
# "content": f"{AWSNovaSonicLLMService.AWAIT_TRIGGER_ASSISTANT_RESPONSE_INSTRUCTION}",
|
||||||
# }
|
# }
|
||||||
# )
|
# )
|
||||||
|
# If the last message isn't from the user, add a message asking for a recap
|
||||||
|
if messages and messages[-1].get("role") != "user":
|
||||||
|
messages.append(
|
||||||
|
{
|
||||||
|
"role": "user",
|
||||||
|
"content": "Can you catch me up on what we were talking about?",
|
||||||
|
}
|
||||||
|
)
|
||||||
params.context.set_messages(messages)
|
params.context.set_messages(messages)
|
||||||
await params.llm.reset_conversation()
|
await params.llm.reset_conversation()
|
||||||
# await params.llm.trigger_assistant_response()
|
# await params.llm.trigger_assistant_response()
|
||||||
|
|||||||
@@ -296,6 +296,7 @@ class AWSNovaSonicLLMService(LLMService):
|
|||||||
self._user_text_buffer = ""
|
self._user_text_buffer = ""
|
||||||
self._assistant_text_buffer = ""
|
self._assistant_text_buffer = ""
|
||||||
self._completed_tool_calls = set()
|
self._completed_tool_calls = set()
|
||||||
|
self._audio_input_started = False
|
||||||
|
|
||||||
file_path = files("pipecat.services.aws.nova_sonic").joinpath("ready.wav")
|
file_path = files("pipecat.services.aws.nova_sonic").joinpath("ready.wav")
|
||||||
with wave.open(file_path.open("rb"), "rb") as wav_file:
|
with wave.open(file_path.open("rb"), "rb") as wav_file:
|
||||||
@@ -532,14 +533,30 @@ class AWSNovaSonicLLMService(LLMService):
|
|||||||
if system_instruction:
|
if system_instruction:
|
||||||
await self._send_text_event(text=system_instruction, role=Role.SYSTEM)
|
await self._send_text_event(text=system_instruction, role=Role.SYSTEM)
|
||||||
|
|
||||||
# Send conversation history
|
# Send conversation history (except for the last message if it's from the
|
||||||
for message in llm_connection_params["messages"]:
|
# user, which we'll send as interactive after starting audio input)
|
||||||
|
messages = llm_connection_params["messages"]
|
||||||
|
last_user_message = None
|
||||||
|
for i, message in enumerate(messages):
|
||||||
# logger.debug(f"Seeding conversation history with message: {message}")
|
# logger.debug(f"Seeding conversation history with message: {message}")
|
||||||
await self._send_text_event(text=message.text, role=message.role)
|
is_last_message = i == len(messages) - 1
|
||||||
|
if is_last_message and message.role == Role.USER:
|
||||||
|
# Save for sending after audio input starts
|
||||||
|
last_user_message = message
|
||||||
|
else:
|
||||||
|
await self._send_text_event(text=message.text, role=message.role)
|
||||||
|
|
||||||
# Start audio input
|
# Start audio input
|
||||||
await self._send_audio_input_start_event()
|
await self._send_audio_input_start_event()
|
||||||
|
|
||||||
|
# Now send the last user message as interactive to trigger bot response
|
||||||
|
if last_user_message:
|
||||||
|
# logger.debug(
|
||||||
|
# f"Sending last user message as interactive to trigger bot response: {last_user_message}")
|
||||||
|
await self._send_text_event(
|
||||||
|
text=last_user_message.text, role=last_user_message.role, interactive=True
|
||||||
|
)
|
||||||
|
|
||||||
# Start receiving events
|
# Start receiving events
|
||||||
self._receive_task = self.create_task(self._receive_task_handler())
|
self._receive_task = self.create_task(self._receive_task_handler())
|
||||||
|
|
||||||
@@ -602,6 +619,7 @@ class AWSNovaSonicLLMService(LLMService):
|
|||||||
self._user_text_buffer = ""
|
self._user_text_buffer = ""
|
||||||
self._assistant_text_buffer = ""
|
self._assistant_text_buffer = ""
|
||||||
self._completed_tool_calls = set()
|
self._completed_tool_calls = set()
|
||||||
|
self._audio_input_started = False
|
||||||
|
|
||||||
logger.info("Finished disconnecting")
|
logger.info("Finished disconnecting")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
@@ -727,8 +745,18 @@ class AWSNovaSonicLLMService(LLMService):
|
|||||||
}}
|
}}
|
||||||
'''
|
'''
|
||||||
await self._send_client_event(audio_content_start)
|
await self._send_client_event(audio_content_start)
|
||||||
|
self._audio_input_started = True
|
||||||
|
|
||||||
async def _send_text_event(self, text: str, role: Role):
|
async def _send_text_event(self, text: str, role: Role, interactive: bool = False):
|
||||||
|
"""Send a text event to the LLM.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
text: The text content to send.
|
||||||
|
role: The role associated with the text (e.g., USER, ASSISTANT, SYSTEM).
|
||||||
|
interactive: Whether the content is interactive. Defaults to False.
|
||||||
|
False: conversation history or system instruction, sent prior to interactive audio
|
||||||
|
True: text input sent during (or at the start of) interactive audio
|
||||||
|
"""
|
||||||
if not self._stream or not self._prompt_name or not text:
|
if not self._stream or not self._prompt_name or not text:
|
||||||
return
|
return
|
||||||
|
|
||||||
@@ -741,7 +769,7 @@ class AWSNovaSonicLLMService(LLMService):
|
|||||||
"promptName": "{self._prompt_name}",
|
"promptName": "{self._prompt_name}",
|
||||||
"contentName": "{content_name}",
|
"contentName": "{content_name}",
|
||||||
"type": "TEXT",
|
"type": "TEXT",
|
||||||
"interactive": true,
|
"interactive": {json.dumps(interactive)},
|
||||||
"role": "{role.value}",
|
"role": "{role.value}",
|
||||||
"textInputConfiguration": {{
|
"textInputConfiguration": {{
|
||||||
"mediaType": "text/plain"
|
"mediaType": "text/plain"
|
||||||
@@ -779,7 +807,7 @@ class AWSNovaSonicLLMService(LLMService):
|
|||||||
await self._send_client_event(text_content_end)
|
await self._send_client_event(text_content_end)
|
||||||
|
|
||||||
async def _send_user_audio_event(self, audio: bytes):
|
async def _send_user_audio_event(self, audio: bytes):
|
||||||
if not self._stream:
|
if not self._stream or not self._audio_input_started:
|
||||||
return
|
return
|
||||||
|
|
||||||
blob = base64.b64encode(audio)
|
blob = base64.b64encode(audio)
|
||||||
|
|||||||
Reference in New Issue
Block a user