diff --git a/EvoScientist/cli/channel.py b/EvoScientist/cli/channel.py index 0b713c4..367944b 100644 --- a/EvoScientist/cli/channel.py +++ b/EvoScientist/cli/channel.py @@ -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 diff --git a/EvoScientist/cli/interactive.py b/EvoScientist/cli/interactive.py index b5648f0..e8ad8f0 100644 --- a/EvoScientist/cli/interactive.py +++ b/EvoScientist/cli/interactive.py @@ -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: diff --git a/EvoScientist/stream/display.py b/EvoScientist/stream/display.py index 55bbfe5..e334311 100644 --- a/EvoScientist/stream/display.py +++ b/EvoScientist/stream/display.py @@ -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() # --------------------------------------------------------------------------- diff --git a/tests/test_bus_integration.py b/tests/test_bus_integration.py index a40fb7f..d118c93 100644 --- a/tests/test_bus_integration.py +++ b/tests/test_bus_integration.py @@ -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(