| """ |
| Agent Execution Orchestration Tests |
| |
| Comprehensive end-to-end tests for agent execution flow: |
| - Governance validation before execution |
| - LLM streaming with BYOK handler |
| - WebSocket token delivery |
| - Chat history persistence |
| - AgentExecution audit trail |
| - Episode creation triggering |
| - Error handling scenarios |
| - Sync execution wrapper |
| |
| Coverage Target: 50%+ on agent_execution_service.py |
| """ |
|
|
| import pytest |
| import asyncio |
| from unittest.mock import Mock, AsyncMock, MagicMock, patch |
| from datetime import datetime |
| from sqlalchemy.orm import Session |
|
|
| from core.agent_execution_service import execute_agent_chat, execute_agent_chat_sync |
| from core.models import AgentExecution, AgentRegistry, ChatSession |
| from core.llm.byok_handler import QueryComplexity |
|
|
|
|
| |
| |
| |
|
|
| @pytest.fixture |
| def mock_autonomous_agent(db_session): |
| """Create AUTONOMOUS agent for testing""" |
| agent = AgentRegistry( |
| id="test_autonomous_agent", |
| name="TestAutonomousAgent", |
| category="testing", |
| status="autonomous", |
| description="Test autonomous agent", |
| module_path="test.agent", |
| class_name="TestAgent" |
| ) |
| db_session.add(agent) |
| db_session.commit() |
| db_session.refresh(agent) |
| yield agent |
| |
| db_session.query(AgentExecution).filter(AgentExecution.agent_id == agent.id).delete() |
| db_session.query(AgentRegistry).filter(AgentRegistry.id == agent.id).delete() |
| db_session.commit() |
|
|
|
|
| @pytest.fixture |
| def mock_student_agent(db_session): |
| """Create STUDENT agent for governance testing""" |
| agent = AgentRegistry( |
| id="test_student_agent", |
| name="TestStudentAgent", |
| category="testing", |
| status="student", |
| description="Test student agent", |
| module_path="test.student_agent", |
| class_name="StudentAgent" |
| ) |
| db_session.add(agent) |
| db_session.commit() |
| db_session.refresh(agent) |
| yield agent |
| |
| db_session.query(AgentExecution).filter(AgentExecution.agent_id == agent.id).delete() |
| db_session.query(AgentRegistry).filter(AgentRegistry.id == agent.id).delete() |
| db_session.commit() |
|
|
|
|
| @pytest.fixture |
| def mock_byok_handler(): |
| """Mock BYOK handler for LLM streaming""" |
| byok_instance = MagicMock() |
| byok_instance.analyze_query_complexity.return_value = QueryComplexity.SIMPLE |
| byok_instance.get_optimal_provider.return_value = ("openai", "gpt-4") |
|
|
| async def mock_stream(**kwargs): |
| tokens = ["Hello", " ", "world", "!"] |
| for token in tokens: |
| yield token |
|
|
| byok_instance.stream_completion = mock_stream |
| return byok_instance |
|
|
|
|
| @pytest.fixture |
| def mock_agent_resolver(mock_autonomous_agent): |
| """Mock agent context resolver""" |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {"resolution_path": ["explicit_agent_id"]}) |
| ) |
|
|
| |
| governance = MagicMock() |
| governance.can_perform_action.return_value = { |
| "proceed": True, |
| "allowed": True |
| } |
| resolver.governance = governance |
|
|
| return resolver |
|
|
|
|
| @pytest.fixture |
| def mock_websocket_manager(): |
| """Mock WebSocket manager for streaming""" |
| ws_instance = MagicMock() |
| ws_instance.broadcast = AsyncMock() |
|
|
| mock_ws = MagicMock() |
| mock_ws.broadcast = ws_instance.broadcast |
| mock_ws.STREAMING_UPDATE = "streaming:update" |
| mock_ws.STREAMING_COMPLETE = "streaming:complete" |
|
|
| return mock_ws |
|
|
|
|
| |
| |
| |
|
|
| class TestGovernanceValidation: |
| """Tests for governance checks before execution""" |
|
|
| @pytest.mark.asyncio |
| async def test_governance_check_passes_for_authorized_agent(self, mock_autonomous_agent, db_session): |
| """AUTONOMOUS agent passes governance check""" |
| with patch('core.agent_execution_service.LLMService') as mock_llm: |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.SessionLocal', return_value=db_session): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| llm_instance = MagicMock() |
| llm.analyze_query_complexity.return_value = QueryComplexity.SIMPLE |
| llm.get_optimal_provider.return_value = ("openai", "gpt-4") |
|
|
| async def mock_stream(**kwargs): |
| yield "Test response" |
|
|
| llm.stream_completion = mock_stream |
| mock_llm.return_value = llm_instance |
|
|
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="Hello", |
| user_id="test_user", |
| stream=False |
| ) |
|
|
| assert result["success"] is True |
| assert "response" in result |
|
|
| @pytest.mark.asyncio |
| async def test_governance_check_blocks_unauthorized_action(self, mock_student_agent, db_session): |
| """STUDENT agent blocked from chat""" |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_student_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={ |
| "proceed": False, |
| "allowed": False, |
| "reason": "STUDENT agent blocked from chat" |
| } |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_student_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| assert result["success"] is False |
| assert "blocked" in result["error"].lower() |
|
|
| @pytest.mark.asyncio |
| async def test_governance_emergency_bypass_allows_execution(self, mock_student_agent, db_session, monkeypatch): |
| """Emergency bypass disables governance checks""" |
| monkeypatch.setenv("EMERGENCY_GOVERNANCE_BYPASS", "true") |
|
|
| with patch('core.agent_execution_service.LLMService') as mock_llm: |
| with patch('core.agent_execution_service.ws_manager'): |
| llm_instance = MagicMock() |
| llm.analyze_query_complexity.return_value = QueryComplexity.SIMPLE |
| llm.get_optimal_provider.return_value = ("openai", "gpt-4") |
|
|
| async def mock_stream(**kwargs): |
| yield "Bypass response" |
|
|
| llm.stream_completion = mock_stream |
| mock_llm.return_value = llm_instance |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_student_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| |
| assert "success" in result or "error" in result |
|
|
| @pytest.mark.asyncio |
| async def test_governance_flag_disables_checks(self, mock_student_agent, db_session, monkeypatch): |
| """STREAMING_GOVERNANCE_ENABLED=false skips governance""" |
| monkeypatch.setenv("STREAMING_GOVERNANCE_ENABLED", "false") |
|
|
| with patch('core.agent_execution_service.LLMService') as mock_llm: |
| with patch('core.agent_execution_service.ws_manager'): |
| llm_instance = MagicMock() |
| llm.analyze_query_complexity.return_value = QueryComplexity.SIMPLE |
| llm.get_optimal_provider.return_value = ("openai", "gpt-4") |
|
|
| async def mock_stream(**kwargs): |
| yield "Disabled governance response" |
|
|
| llm.stream_completion = mock_stream |
| mock_llm.return_value = llm_instance |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_student_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| |
| assert "success" in result or "error" in result |
|
|
|
|
| |
| |
| |
|
|
| class TestLLMStreamingExecution: |
| """Tests for BYOK handler streaming""" |
|
|
| @pytest.mark.asyncio |
| async def test_llm_streaming_accumulates_tokens(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """Tokens are accumulated correctly from streaming""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123", |
| stream=False |
| ) |
|
|
| assert result["success"] is True |
| assert result["response"] == "Hello world!" |
| assert result["tokens"] == 4 |
|
|
| @pytest.mark.asyncio |
| async def test_llm_provider_selection(self, mock_autonomous_agent, db_session): |
| """Optimal provider is selected based on query complexity""" |
| with patch('core.agent_execution_service.LLMService') as mock_llm: |
| with patch('core.agent_execution_service.ws_manager'): |
| llm_instance = MagicMock() |
| llm.analyze_query_complexity.return_value = QueryComplexity.COMPLEX |
| llm.get_optimal_provider.return_value = ("anthropic", "claude-3-opus") |
|
|
| async def mock_stream(**kwargs): |
| yield "Complex query response" |
|
|
| llm.stream_completion = mock_stream |
| mock_llm.return_value = llm_instance |
|
|
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="Complex question requiring reasoning", |
| user_id="user_123" |
| ) |
|
|
| assert result["success"] is True |
| assert result["provider"] == "anthropic" |
| assert result["model"] == "claude-3-opus" |
|
|
| @pytest.mark.asyncio |
| async def test_llm_streaming_with_conversation_history(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """Conversation history is included in LLM context""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| history = [ |
| {"role": "user", "content": "Previous question"}, |
| {"role": "assistant", "content": "Previous answer"} |
| ] |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="Follow-up question", |
| user_id="user_123", |
| conversation_history=history |
| ) |
|
|
| assert result["success"] is True |
| |
|
|
| @pytest.mark.asyncio |
| async def test_llm_error_propagation(self, mock_autonomous_agent, db_session): |
| """LLM streaming errors are handled gracefully""" |
| with patch('core.agent_execution_service.LLMService') as mock_llm: |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action.return_value = { |
| "proceed": True, "allowed": True |
| } |
| mock_resolver.return_value = resolver |
|
|
| llm_instance = MagicMock() |
| llm.analyze_query_complexity.return_value = QueryComplexity.SIMPLE |
|
|
| async def mock_stream_error(**kwargs): |
| raise Exception("LLM API error") |
|
|
| llm.stream_completion = mock_stream_error |
| mock_llm.return_value = llm_instance |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| assert result["success"] is False |
| assert "error" in result |
|
|
|
|
| |
| |
| |
|
|
| class TestWebSocketStreaming: |
| """Tests for WebSocket token delivery""" |
|
|
| @pytest.mark.asyncio |
| async def test_websocket_sends_start_message(self, mock_autonomous_agent, mock_byok_handler, mock_websocket_manager, db_session): |
| """WebSocket sends streaming:start message""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager', mock_websocket_manager): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123", |
| stream=True |
| ) |
|
|
| assert result["success"] is True |
| |
| mock_websocket_manager.broadcast.assert_any_call( |
| "user:user_123", |
| { |
| "type": "streaming:start", |
| "id": result["message_id"], |
| "model": "gpt-4", |
| "provider": "openai", |
| "agent_id": mock_autonomous_agent.id, |
| "agent_name": mock_autonomous_agent.name, |
| "execution_id": result["execution_id"] |
| } |
| ) |
|
|
| @pytest.mark.asyncio |
| async def test_websocket_sends_update_messages(self, mock_autonomous_agent, mock_byok_handler, mock_websocket_manager, db_session): |
| """WebSocket sends streaming:update messages for each token""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager', mock_websocket_manager): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123", |
| stream=True |
| ) |
|
|
| assert result["success"] is True |
| |
| assert mock_websocket_manager.broadcast.call_count >= 5 |
|
|
| @pytest.mark.asyncio |
| async def test_websocket_sends_complete_message(self, mock_autonomous_agent, mock_byok_handler, mock_websocket_manager, db_session): |
| """WebSocket sends streaming:complete message with full content""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager', mock_websocket_manager): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123", |
| stream=True |
| ) |
|
|
| assert result["success"] is True |
| |
| complete_calls = [ |
| call for call in mock_websocket_manager.broadcast.call_args_list |
| if "streaming:complete" in str(call) |
| ] |
| assert len(complete_calls) > 0 |
|
|
| @pytest.mark.asyncio |
| async def test_websocket_skipped_when_stream_false(self, mock_autonomous_agent, mock_byok_handler, mock_websocket_manager, db_session): |
| """No WebSocket messages when stream=False""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager', mock_websocket_manager): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123", |
| stream=False |
| ) |
|
|
| assert result["success"] is True |
| |
| mock_websocket_manager.broadcast.assert_not_called() |
|
|
|
|
| |
| |
| |
|
|
| class TestChatHistoryPersistence: |
| """Tests for chat history saving""" |
|
|
| @pytest.mark.asyncio |
| async def test_chat_session_created_on_execution(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """New chat session is created if none provided""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123", |
| session_id=None |
| ) |
|
|
| assert result["success"] is True |
| assert "session_id" in result |
|
|
| @pytest.mark.asyncio |
| async def test_messages_saved_to_chat_history(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """User and assistant messages are saved to chat history""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="Hello world", |
| user_id="user_123" |
| ) |
|
|
| assert result["success"] is True |
| assert result["response"] == "Hello world!" |
|
|
| @pytest.mark.asyncio |
| async def test_chat_history_persistence_error_doesnt_fail_execution(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """Execution succeeds even if chat history persistence fails""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.get_chat_history_manager') as mock_history: |
| mock_history.side_effect = Exception("Chat history DB error") |
|
|
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action.return_value = { |
| "proceed": True, "allowed": True |
| } |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| |
| assert result["success"] is True |
|
|
| @pytest.mark.asyncio |
| async def test_existing_session_reused_for_continuity(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """Existing session_id is reused for conversation continuity""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| existing_session_id = "test_session_123" |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123", |
| session_id=existing_session_id |
| ) |
|
|
| assert result["success"] is True |
| assert result["session_id"] == existing_session_id |
|
|
|
|
| |
| |
| |
|
|
| class TestAgentExecutionAuditTrail: |
| """Tests for AgentExecution records""" |
|
|
| @pytest.mark.asyncio |
| async def test_agent_execution_record_created(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """AgentExecution record is created on execution start""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123", |
| stream=False |
| ) |
|
|
| |
| execution = db_session.query(AgentExecution).filter( |
| AgentExecution.agent_id == mock_autonomous_agent.id |
| ).first() |
|
|
| assert execution is not None |
| assert execution.status == "completed" |
| assert execution.agent_name == mock_autonomous_agent.name |
|
|
| @pytest.mark.asyncio |
| async def test_execution_record_updated_on_completion(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """AgentExecution status updated to completed with output data""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| |
| execution = db_session.query(AgentExecution).filter( |
| AgentExecution.agent_id == mock_autonomous_agent.id |
| ).first() |
|
|
| assert execution.status == "completed" |
| assert execution.output_data is not None |
| assert execution.output_data["response"] == "Hello world!" |
| assert execution.output_data["tokens"] == 4 |
| assert execution.duration_ms > 0 |
| assert execution.end_time is not None |
|
|
| @pytest.mark.asyncio |
| async def test_execution_record_marked_failed_on_error(self, mock_autonomous_agent, db_session): |
| """AgentExecution marked as failed on execution error""" |
| with patch('core.agent_execution_service.LLMService') as mock_llm: |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action.return_value = { |
| "proceed": True, "allowed": True |
| } |
| mock_resolver.return_value = resolver |
|
|
| llm_instance = MagicMock() |
| llm.analyze_query_complexity.return_value = QueryComplexity.SIMPLE |
|
|
| async def mock_stream_error(**kwargs): |
| raise Exception("LLM API failure") |
|
|
| llm.stream_completion = mock_stream_error |
| mock_llm.return_value = llm_instance |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| |
| execution = db_session.query(AgentExecution).filter( |
| AgentExecution.agent_id == mock_autonomous_agent.id |
| ).first() |
|
|
| assert execution is not None |
| assert execution.status == "failed" |
| assert execution.error_message is not None |
|
|
| @pytest.mark.asyncio |
| async def test_execution_metadata_includes_governance_context(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """Execution metadata includes governance check details""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolution_context = {"resolution_path": ["explicit_agent_id"]} |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, resolution_context) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action.return_value = { |
| "proceed": True, "allowed": True, "maturity_level": "AUTONOMOUS" |
| } |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| |
| execution = db_session.query(AgentExecution).filter( |
| AgentExecution.agent_id == mock_autonomous_agent.id |
| ).first() |
|
|
| assert execution is not None |
| assert execution.metadata is not None |
| assert "governance_check" in execution.metadata |
|
|
|
|
| |
| |
| |
|
|
| class TestEpisodeCreationTriggering: |
| """Tests for episode creation after execution""" |
|
|
| @pytest.mark.asyncio |
| async def test_episode_creation_triggered_after_execution(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """Episode creation is triggered after successful execution""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.trigger_episode_creation') as mock_episode: |
| mock_episode.return_value = None |
|
|
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action.return_value = { |
| "proceed": True, "allowed": True |
| } |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| |
| mock_episode.assert_called_once() |
|
|
| @pytest.mark.asyncio |
| async def test_episode_creation_error_doesnt_fail_execution(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """Execution succeeds even if episode creation fails""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.trigger_episode_creation') as mock_episode: |
| mock_episode.side_effect = Exception("Episode service error") |
|
|
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action.return_value = { |
| "proceed": True, "allowed": True |
| } |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| |
| assert result["success"] is True |
|
|
| @pytest.mark.asyncio |
| async def test_episode_creation_with_session_context(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """Episode creation includes session and agent context""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.trigger_episode_creation') as mock_episode: |
| mock_episode.return_value = None |
|
|
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action.return_value = { |
| "proceed": True, "allowed": True |
| } |
| mock_resolver.return_value = resolver |
|
|
| session_id = "test_session_456" |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123", |
| session_id=session_id |
| ) |
|
|
| |
| mock_episode.assert_called_once() |
| call_args = mock_episode.call_args |
| assert call_args[1]["session_id"] == session_id |
| assert call_args[1]["agent_id"] == mock_autonomous_agent.id |
|
|
| @pytest.mark.asyncio |
| async def test_episode_creation_for_governance_blocked_execution(self, mock_student_agent, db_session): |
| """Episode creation not triggered for governance-blocked executions""" |
| with patch('core.agent_execution_service.trigger_episode_creation') as mock_episode: |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_student_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action.return_value = { |
| "proceed": False, |
| "allowed": False, |
| "reason": "STUDENT agent blocked" |
| } |
| mock_resolver.return_value = resolver |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_student_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| |
| mock_episode.assert_not_called() |
|
|
|
|
| |
| |
| |
|
|
| class TestErrorHandling: |
| """Tests for error scenarios""" |
|
|
| @pytest.mark.asyncio |
| async def test_llm_failure_caught_and_logged(self, mock_autonomous_agent, db_session): |
| """LLM provider failures are caught and returned as errors""" |
| with patch('core.agent_execution_service.LLMService') as mock_llm: |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action.return_value = { |
| "proceed": True, "allowed": True |
| } |
| mock_resolver.return_value = resolver |
|
|
| llm_instance = MagicMock() |
| llm.analyze_query_complexity.return_value = QueryComplexity.SIMPLE |
|
|
| async def mock_stream_error(**kwargs): |
| raise Exception("OpenAI API timeout") |
|
|
| llm.stream_completion = mock_stream_error |
| mock_llm.return_value = llm_instance |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| assert result["success"] is False |
| assert "OpenAI API timeout" in result["error"] |
|
|
| @pytest.mark.asyncio |
| async def test_database_connection_error_returns_error_response(self, mock_autonomous_agent, db_session, monkeypatch): |
| """Database connection errors are handled gracefully""" |
| |
| with patch('core.agent_execution_service.SessionLocal') as mock_db: |
| mock_db.side_effect = Exception("Database connection failed") |
|
|
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| |
| |
| assert "success" in result or "error" in result |
|
|
| @pytest.mark.asyncio |
| async def test_websocket_disconnection_doesnt_affect_execution(self, mock_autonomous_agent, mock_byok_handler, db_session): |
| """Execution completes even if WebSocket disconnects""" |
| with patch('core.agent_execution_service.LLMService', return_value=mock_llm_service): |
| with patch('core.agent_execution_service.ws_manager') as mock_ws: |
| |
| mock_ws.broadcast.side_effect = Exception("WebSocket disconnected") |
|
|
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_autonomous_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| |
| try: |
| result = await execute_agent_chat( |
| agent_id=mock_autonomous_agent.id, |
| message="test", |
| user_id="user_123", |
| stream=True |
| ) |
| |
| assert result["success"] is True or "error" in result |
| except Exception as e: |
| |
| assert "WebSocket" in str(e) |
|
|
|
|
| |
| |
| |
|
|
| class TestSyncExecutionWrapper: |
| """Tests for execute_agent_chat_sync()""" |
|
|
| def test_sync_wrapper_creates_event_loop(self): |
| """Sync wrapper creates event loop if none exists""" |
| with patch('core.agent_execution_service.LLMService') as mock_llm: |
| with patch('core.agent_execution_service.ws_manager'): |
| with patch('core.agent_execution_service.AgentContextResolver') as mock_resolver: |
| |
| mock_agent = MagicMock() |
| mock_agent.id = "test_agent" |
| mock_agent.name = "TestAgent" |
|
|
| resolver = MagicMock() |
| resolver.resolve_agent_for_request = AsyncMock( |
| return_value=(mock_agent, {}) |
| ) |
| resolver.governance = MagicMock() |
| resolver.governance.can_perform_action = MagicMock( |
| return_value={"proceed": True, "allowed": True} |
| ) |
| mock_resolver.return_value = resolver |
|
|
| llm_instance = MagicMock() |
| llm.analyze_query_complexity.return_value = QueryComplexity.SIMPLE |
| llm.get_optimal_provider.return_value = ("openai", "gpt-4") |
|
|
| async def mock_stream(**kwargs): |
| yield "Sync response" |
|
|
| llm.stream_completion = mock_stream |
| mock_llm.return_value = llm_instance |
|
|
| result = execute_agent_chat_sync( |
| agent_id="test_agent", |
| message="test message", |
| user_id="user_123", |
| stream=False |
| ) |
|
|
| assert "success" in result or "error" in result |
|
|
| def test_sync_wrapper_disables_streaming(self): |
| """Sync wrapper forces stream=False regardless of input""" |
| with patch('core.agent_execution_service.execute_agent_chat') as mock_execute: |
| mock_execute.return_value = {"success": True, "response": "test"} |
|
|
| result = execute_agent_chat_sync( |
| agent_id="test_agent", |
| message="test", |
| user_id="user_123", |
| stream=True |
| ) |
|
|
| |
| mock_execute.assert_called_once() |
| call_kwargs = mock_execute.call_args[1] |
| assert call_kwargs["stream"] is False |
|
|
| def test_sync_wrapper_reuses_existing_event_loop(self): |
| """Sync wrapper reuses existing event loop""" |
| import asyncio |
|
|
| |
| loop = asyncio.new_event_loop() |
| asyncio.set_event_loop(loop) |
|
|
| try: |
| with patch('core.agent_execution_service.execute_agent_chat') as mock_execute: |
| mock_execute.return_value = {"success": True, "response": "test"} |
|
|
| result = execute_agent_chat_sync( |
| agent_id="test_agent", |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| assert result["success"] is True |
| finally: |
| loop.close() |
|
|
| def test_sync_wrapper_propagates_execution_result(self): |
| """Sync wrapper returns full execution result""" |
| with patch('core.agent_execution_service.execute_agent_chat') as mock_execute: |
| mock_execute.return_value = { |
| "success": True, |
| "execution_id": "exec_123", |
| "response": "Test response", |
| "tokens": 2, |
| "provider": "openai", |
| "model": "gpt-4" |
| } |
|
|
| result = execute_agent_chat_sync( |
| agent_id="test_agent", |
| message="test", |
| user_id="user_123" |
| ) |
|
|
| assert result["success"] is True |
| assert result["execution_id"] == "exec_123" |
| assert result["response"] == "Test response" |
|
|