rtvi: add RTVIProcessor.handle_message()
This commit is contained in:
@@ -343,6 +343,9 @@ class RTVIProcessor(FrameProcessor):
|
|||||||
self._client_ready = True
|
self._client_ready = True
|
||||||
await self._maybe_send_bot_ready()
|
await self._maybe_send_bot_ready()
|
||||||
|
|
||||||
|
async def handle_message(self, message: RTVIMessage):
|
||||||
|
await self._message_queue.put(message)
|
||||||
|
|
||||||
async def handle_function_call(
|
async def handle_function_call(
|
||||||
self,
|
self,
|
||||||
function_name: str,
|
function_name: str,
|
||||||
@@ -492,20 +495,21 @@ class RTVIProcessor(FrameProcessor):
|
|||||||
async def _message_task_handler(self):
|
async def _message_task_handler(self):
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
frame = await self._message_queue.get()
|
message = await self._message_queue.get()
|
||||||
await self._handle_message(frame)
|
await self._handle_message(message)
|
||||||
self._message_queue.task_done()
|
self._message_queue.task_done()
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
break
|
break
|
||||||
|
|
||||||
async def _handle_message(self, frame: TransportMessageFrame):
|
async def _handle_transport_message(self, frame: TransportMessageFrame):
|
||||||
try:
|
try:
|
||||||
message = RTVIMessage.model_validate(frame.message)
|
message = RTVIMessage.model_validate(frame.message)
|
||||||
|
await self._message_queue.put(message)
|
||||||
except ValidationError as e:
|
except ValidationError as e:
|
||||||
await self.send_error(f"Invalid incoming message: {e}")
|
await self.send_error(f"Invalid RTVI transport message: {e}")
|
||||||
logger.warning(f"Invalid incoming message: {e}")
|
logger.warning(f"Invalid RTVI transport message: {e}")
|
||||||
return
|
|
||||||
|
|
||||||
|
async def _handle_message(self, message: RTVIMessage):
|
||||||
try:
|
try:
|
||||||
match message.type:
|
match message.type:
|
||||||
case "client-ready":
|
case "client-ready":
|
||||||
@@ -531,8 +535,8 @@ class RTVIProcessor(FrameProcessor):
|
|||||||
await self._send_error_response(message.id, f"Unsupported type {message.type}")
|
await self._send_error_response(message.id, f"Unsupported type {message.type}")
|
||||||
|
|
||||||
except ValidationError as e:
|
except ValidationError as e:
|
||||||
await self._send_error_response(message.id, f"Invalid incoming message: {e}")
|
await self._send_error_response(message.id, f"Invalid message: {e}")
|
||||||
logger.warning(f"Invalid incoming message: {e}")
|
logger.warning(f"Invalid message: {e}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self._send_error_response(message.id, f"Exception processing message: {e}")
|
await self._send_error_response(message.id, f"Exception processing message: {e}")
|
||||||
logger.warning(f"Exception processing message: {e}")
|
logger.warning(f"Exception processing message: {e}")
|
||||||
|
|||||||
Reference in New Issue
Block a user