feat(channel): improve error handling and logging in channel management
feat(interactive): refactor message sending to channel for better error handling fix(display): remove unnecessary print statement in final results display test(bus): update tests to match changes in _bus_inbound_consumer signature
This commit is contained in:
+14
-11
@@ -12,6 +12,7 @@ for the main thread to set a response via ``_set_channel_response()``.
|
||||
import asyncio
|
||||
import logging
|
||||
import queue
|
||||
import time
|
||||
import threading
|
||||
import uuid
|
||||
from dataclasses import dataclass
|
||||
@@ -37,7 +38,7 @@ class ChannelMessage:
|
||||
content: str
|
||||
sender: str
|
||||
channel_type: str
|
||||
metadata: Any = None
|
||||
metadata: dict | None = None
|
||||
# Filled by the bus consumer so the main thread can send callbacks
|
||||
channel_ref: Any = None # Channel instance (for thinking / todo / file)
|
||||
bus_ref: Any = None # MessageBus (for publishing outbound)
|
||||
@@ -116,8 +117,8 @@ def _channels_stop(channel_type: str | None = None) -> None:
|
||||
_manager.stop_all(), _bus_loop,
|
||||
)
|
||||
future.result(timeout=10)
|
||||
except Exception:
|
||||
pass
|
||||
except Exception as e:
|
||||
_channel_logger.debug(f"Error stopping channels: {e}")
|
||||
if _manager:
|
||||
_manager.bus.stop()
|
||||
if _bus_thread:
|
||||
@@ -136,8 +137,8 @@ def _channels_stop(channel_type: str | None = None) -> None:
|
||||
_manager.remove_channel(channel_type), _bus_loop,
|
||||
)
|
||||
future.result(timeout=5)
|
||||
except Exception:
|
||||
pass
|
||||
except Exception as e:
|
||||
_channel_logger.debug(f"Error removing channel {channel_type}: {e}")
|
||||
|
||||
if _manager and not _manager.running_channels():
|
||||
_cli_agent = None
|
||||
@@ -170,7 +171,7 @@ def _start_channels_bus_mode(config, agent, thread_id: str, show_thinking: bool
|
||||
|
||||
async def _run():
|
||||
consumer = asyncio.create_task(
|
||||
_bus_inbound_consumer(mgr.bus, mgr, show_thinking)
|
||||
_bus_inbound_consumer(mgr.bus, mgr)
|
||||
)
|
||||
try:
|
||||
await mgr.start_all()
|
||||
@@ -193,7 +194,6 @@ def _start_channels_bus_mode(config, agent, thread_id: str, show_thinking: bool
|
||||
thread.start()
|
||||
|
||||
# Wait briefly for the loop to start
|
||||
import time
|
||||
for _ in range(20):
|
||||
if _bus_loop is not None:
|
||||
break
|
||||
@@ -218,9 +218,7 @@ def _add_channel_to_running_bus(channel_type: str, config) -> None:
|
||||
future.result(timeout=10)
|
||||
|
||||
|
||||
async def _bus_inbound_consumer(
|
||||
bus, manager, show_thinking: bool = True,
|
||||
) -> None:
|
||||
async def _bus_inbound_consumer(bus, manager) -> None:
|
||||
"""Consume inbound messages from bus and bridge to the main CLI thread.
|
||||
|
||||
This does NOT invoke the agent. It enqueues a ``ChannelMessage`` on
|
||||
@@ -262,7 +260,12 @@ async def _bus_inbound_consumer(
|
||||
event = _enqueue_channel_message(cm)
|
||||
|
||||
# Wait (non-blocking for asyncio) until main thread sets response
|
||||
await asyncio.to_thread(event.wait)
|
||||
_RESPONSE_TIMEOUT = 600 # 10 minutes max per message
|
||||
replied = await asyncio.to_thread(event.wait, _RESPONSE_TIMEOUT)
|
||||
if not replied:
|
||||
_channel_logger.warning(
|
||||
f"[bus] Response timeout ({_RESPONSE_TIMEOUT}s) for {cm.msg_id}"
|
||||
)
|
||||
response = _pop_channel_response(cm.msg_id) or "No response"
|
||||
|
||||
# Publish the response back through the bus → channel
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
"""Interactive CLI mode and single-shot execution."""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
import queue
|
||||
import sys
|
||||
@@ -45,6 +46,8 @@ import EvoScientist.cli.channel as _ch_mod
|
||||
from .mcp_ui import _cmd_mcp
|
||||
from .skills_cmd import _cmd_list_skills, _cmd_install_skill, _cmd_uninstall_skill
|
||||
|
||||
_channel_logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# Banner
|
||||
@@ -465,62 +468,39 @@ def cmd_interactive(
|
||||
rx.append("]", style="dim")
|
||||
console.print(rx)
|
||||
_print_separator()
|
||||
console.print()
|
||||
|
||||
def _send_to_channel(coro, label: str, timeout: int = 15) -> None:
|
||||
"""Schedule an async channel send on the bus loop."""
|
||||
loop = _ch_mod._bus_loop
|
||||
if not loop:
|
||||
return
|
||||
try:
|
||||
asyncio.run_coroutine_threadsafe(coro, loop).result(timeout=timeout)
|
||||
except Exception as e:
|
||||
_channel_logger.debug(f"{label} send failed: {e}")
|
||||
|
||||
def _send_thinking_to_channel(thinking: str) -> None:
|
||||
"""Send thinking text to the channel (sync callback)."""
|
||||
ch = msg.channel_ref
|
||||
loop = _ch_mod._bus_loop
|
||||
if ch and loop and ch.send_thinking:
|
||||
try:
|
||||
future = asyncio.run_coroutine_threadsafe(
|
||||
ch.send_thinking_message(
|
||||
sender=msg.chat_id,
|
||||
thinking=thinking,
|
||||
metadata=msg.metadata,
|
||||
),
|
||||
loop,
|
||||
)
|
||||
future.result(timeout=15)
|
||||
except Exception:
|
||||
pass
|
||||
if ch and ch.send_thinking:
|
||||
_send_to_channel(
|
||||
ch.send_thinking_message(sender=msg.chat_id, thinking=thinking, metadata=msg.metadata),
|
||||
"Thinking",
|
||||
)
|
||||
|
||||
def _send_todo_to_channel(items: list[dict]) -> None:
|
||||
"""Send todo list to the channel (sync callback)."""
|
||||
from ..channels.consumer import _format_todo_list
|
||||
ch = msg.channel_ref
|
||||
loop = _ch_mod._bus_loop
|
||||
if ch and loop:
|
||||
try:
|
||||
future = asyncio.run_coroutine_threadsafe(
|
||||
ch.send_todo_message(
|
||||
sender=msg.chat_id,
|
||||
content=_format_todo_list(items),
|
||||
metadata=msg.metadata,
|
||||
),
|
||||
loop,
|
||||
)
|
||||
future.result(timeout=15)
|
||||
except Exception:
|
||||
pass
|
||||
if msg.channel_ref:
|
||||
_send_to_channel(
|
||||
msg.channel_ref.send_todo_message(sender=msg.chat_id, content=_format_todo_list(items), metadata=msg.metadata),
|
||||
"Todo",
|
||||
)
|
||||
|
||||
def _send_media_to_channel(file_path: str) -> None:
|
||||
"""Send media file back through the channel (sync callback)."""
|
||||
ch = msg.channel_ref
|
||||
loop = _ch_mod._bus_loop
|
||||
if ch and loop:
|
||||
try:
|
||||
future = asyncio.run_coroutine_threadsafe(
|
||||
ch.send_media(
|
||||
recipient=msg.chat_id,
|
||||
file_path=file_path,
|
||||
metadata=msg.metadata,
|
||||
),
|
||||
loop,
|
||||
)
|
||||
future.result(timeout=30)
|
||||
except Exception as e:
|
||||
console.print(f"[dim]Media send failed: {e}[/dim]")
|
||||
if msg.channel_ref:
|
||||
_send_to_channel(
|
||||
msg.channel_ref.send_media(recipient=msg.chat_id, file_path=file_path, metadata=msg.metadata),
|
||||
"Media", timeout=30,
|
||||
)
|
||||
|
||||
meta = _build_metadata(state["workspace_dir"], model)
|
||||
try:
|
||||
@@ -545,7 +525,7 @@ def cmd_interactive(
|
||||
_print_separator()
|
||||
|
||||
# Redraw the ❯ prompt on a new line after separator
|
||||
sys.stdout.write("\n\033[34;1m\u276f\033[0m ")
|
||||
sys.stdout.write("\033[34;1m\u276f\033[0m ")
|
||||
sys.stdout.flush()
|
||||
|
||||
async def _check_channel_queue() -> None:
|
||||
@@ -644,7 +624,7 @@ def cmd_interactive(
|
||||
continue
|
||||
|
||||
if user_input.lower().startswith("/mcp"):
|
||||
_cmd_mcp(user_input[4:])
|
||||
_cmd_mcp(user_input[len("/mcp"):])
|
||||
continue
|
||||
|
||||
if user_input.lower().startswith("/channel"):
|
||||
@@ -663,6 +643,7 @@ def cmd_interactive(
|
||||
state["agent"], user_input, state["thread_id"],
|
||||
show_thinking, interactive=True, metadata=meta,
|
||||
)
|
||||
console.print()
|
||||
_print_separator()
|
||||
|
||||
except KeyboardInterrupt:
|
||||
|
||||
@@ -547,7 +547,6 @@ def display_final_results(
|
||||
clean_response = clean_response.rstrip().removesuffix("...").rstrip()
|
||||
console.print()
|
||||
console.print(Markdown(clean_response or state.response_text))
|
||||
console.print()
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@@ -88,7 +88,7 @@ class TestBusInboundConsumer:
|
||||
manager.register(ch)
|
||||
|
||||
consumer = asyncio.create_task(
|
||||
_bus_inbound_consumer(bus, manager, False)
|
||||
_bus_inbound_consumer(bus, manager)
|
||||
)
|
||||
|
||||
await bus.publish_inbound(InboundMessage(
|
||||
@@ -143,7 +143,7 @@ class TestBusInboundConsumer:
|
||||
manager.register(ch)
|
||||
|
||||
consumer = asyncio.create_task(
|
||||
_bus_inbound_consumer(bus, manager, False)
|
||||
_bus_inbound_consumer(bus, manager)
|
||||
)
|
||||
|
||||
await bus.publish_inbound(InboundMessage(
|
||||
@@ -191,7 +191,7 @@ class TestBusInboundConsumer:
|
||||
manager.register(ch)
|
||||
|
||||
consumer = asyncio.create_task(
|
||||
_bus_inbound_consumer(bus, manager, False)
|
||||
_bus_inbound_consumer(bus, manager)
|
||||
)
|
||||
|
||||
await bus.publish_inbound(InboundMessage(
|
||||
@@ -238,7 +238,7 @@ class TestBusInboundConsumer:
|
||||
manager.register(ch)
|
||||
|
||||
consumer = asyncio.create_task(
|
||||
_bus_inbound_consumer(bus, manager, False)
|
||||
_bus_inbound_consumer(bus, manager)
|
||||
)
|
||||
|
||||
await bus.publish_inbound(InboundMessage(
|
||||
|
||||
Reference in New Issue
Block a user