update
This commit is contained in:
+18
-11
@@ -334,12 +334,15 @@ def _create_channel_handler():
|
||||
metadata=msg.metadata,
|
||||
)
|
||||
|
||||
# Wait for main thread to process and set response
|
||||
response = await asyncio.to_thread(
|
||||
_ChannelState.get_response, msg_id, 300 # 5 minute timeout
|
||||
)
|
||||
# Wait indefinitely for main thread to process and set response
|
||||
# (no timeout - let the agent work as long as needed)
|
||||
await asyncio.to_thread(event.wait)
|
||||
|
||||
return response or "No response"
|
||||
# Get the response
|
||||
with _ChannelState._response_lock:
|
||||
response = _ChannelState.pending_responses.pop(msg_id, {}).get("response", "")
|
||||
|
||||
return response if response else "(empty response)"
|
||||
|
||||
return handler
|
||||
|
||||
@@ -466,13 +469,17 @@ def cmd_interactive(
|
||||
console.print(f"[bold blue]>[/bold blue] {msg.content} [dim]{source_tag}[/dim]")
|
||||
console.print()
|
||||
|
||||
# Use SAME _run_streaming as CLI input — full Live experience
|
||||
response_text = _run_streaming(
|
||||
state["agent"], msg.content, state["thread_id"], show_thinking, interactive=True
|
||||
)
|
||||
try:
|
||||
# Use SAME _run_streaming as CLI input — full Live experience
|
||||
response_text = _run_streaming(
|
||||
state["agent"], msg.content, state["thread_id"], show_thinking, interactive=True
|
||||
)
|
||||
|
||||
# Set response for channel handler to retrieve
|
||||
_ChannelState.set_response(msg.msg_id, response_text or "")
|
||||
# Set response for channel handler to retrieve
|
||||
_ChannelState.set_response(msg.msg_id, response_text or "")
|
||||
except Exception as e:
|
||||
console.print(f"[red]Channel processing error: {e}[/red]")
|
||||
_ChannelState.set_response(msg.msg_id, f"Error: {e}")
|
||||
|
||||
_print_separator()
|
||||
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
"""Tests for event loop management in streaming display."""
|
||||
|
||||
import asyncio
|
||||
import pytest
|
||||
from unittest.mock import Mock, patch
|
||||
|
||||
from EvoScientist.stream.display import _get_event_loop, _create_event_loop
|
||||
@@ -199,30 +198,3 @@ class TestEventLoopThreadSafety:
|
||||
|
||||
# Cleanup
|
||||
loop.close()
|
||||
|
||||
@pytest.mark.skipif(
|
||||
True, # Skip by default as threading tests can be flaky
|
||||
reason="Thread test can be flaky in CI"
|
||||
)
|
||||
def test_worker_thread_creates_loop(self):
|
||||
"""Worker thread should be able to create its own loop."""
|
||||
import threading
|
||||
|
||||
result = {"loop": None, "error": None}
|
||||
|
||||
def thread_func():
|
||||
try:
|
||||
# In a worker thread, get_event_loop() may raise RuntimeError
|
||||
# Our code should handle this by creating a new loop
|
||||
result["loop"] = _get_event_loop()
|
||||
except Exception as e:
|
||||
result["error"] = e
|
||||
|
||||
thread = threading.Thread(target=thread_func)
|
||||
thread.start()
|
||||
thread.join()
|
||||
|
||||
# Should either get a loop or handle the error gracefully
|
||||
assert result["error"] is None or isinstance(result["error"], RuntimeError)
|
||||
if result["loop"]:
|
||||
assert not result["loop"].is_closed()
|
||||
|
||||
@@ -1,13 +1,5 @@
|
||||
"""Smoke tests verifying package structure is intact."""
|
||||
|
||||
import os
|
||||
import pytest
|
||||
|
||||
needs_api_key = pytest.mark.skipif(
|
||||
not os.getenv("ANTHROPIC_API_KEY"),
|
||||
reason="ANTHROPIC_API_KEY not set",
|
||||
)
|
||||
|
||||
|
||||
def test_import_stream_utils():
|
||||
from EvoScientist.stream.utils import (
|
||||
@@ -42,10 +34,3 @@ def test_import_tools():
|
||||
from EvoScientist.tools import think_tool
|
||||
assert think_tool is not None
|
||||
assert hasattr(think_tool, "invoke")
|
||||
|
||||
|
||||
@needs_api_key
|
||||
def test_import_package_exports():
|
||||
from EvoScientist import EvoScientist_agent
|
||||
# Just verify they are importable; don't call them without API key
|
||||
assert EvoScientist_agent is not None
|
||||
|
||||
Reference in New Issue
Block a user