File size: 5,974 Bytes
6172a47 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 | import asyncio
from unittest.mock import MagicMock
import pytest
from messaging.handler import ClaudeMessageHandler
from messaging.trees.data import MessageState
@pytest.fixture
def handler_integration(mock_platform, mock_cli_manager, mock_session_store):
# Use real TreeQueueManager
handler = ClaudeMessageHandler(mock_platform, mock_cli_manager, mock_session_store)
return handler
async def mock_async_gen(events):
for e in events:
yield e
@pytest.mark.asyncio
async def test_full_conversation_flow_single_user(
handler_integration, mock_platform, mock_cli_manager, incoming_message_factory
):
# 1. First message
msg1 = incoming_message_factory(text="message 1", message_id="m1")
mock_platform.queue_send_message.return_value = "s1"
# Mock CLI session for m1
mock_session1 = MagicMock()
mock_session1.start_task.return_value = mock_async_gen(
[
{"type": "session_info", "session_id": "sess1"},
{
"type": "assistant",
"message": {"content": [{"type": "text", "text": "Reply 1"}]},
},
{"type": "exit", "code": 0, "stderr": None},
]
)
mock_cli_manager.get_or_create_session.return_value = (
mock_session1,
"pending_1",
True,
)
await handler_integration.handle_message(msg1)
# Wait for processing
tree = handler_integration.tree_queue.get_tree_for_node("m1")
for _ in range(10):
if tree.get_node("m1").state.value == MessageState.COMPLETED.value:
break
await asyncio.sleep(0.01)
assert tree.get_node("m1").state.value == MessageState.COMPLETED.value
assert tree.get_node("m1").session_id == "sess1"
mock_session1.start_task.assert_called_with(
"message 1", session_id=None, fork_session=False
)
# 2. Reply to m1
msg2 = incoming_message_factory(
text="message 2", message_id="m2", reply_to_message_id="m1"
)
mock_platform.queue_send_message.return_value = "s2"
# Mock CLI session for m2
mock_session2 = MagicMock()
mock_session2.start_task.return_value = mock_async_gen(
[
{"type": "session_info", "session_id": "sess2"},
{
"type": "assistant",
"message": {"content": [{"type": "text", "text": "Reply 2"}]},
},
{"type": "exit", "code": 0, "stderr": None},
]
)
mock_cli_manager.get_or_create_session.reset_mock()
mock_cli_manager.get_or_create_session.return_value = (
mock_session2,
"pending_2",
True,
)
await handler_integration.handle_message(msg2)
# Wait for processing
for _ in range(10):
if tree.get_node("m2").state.value == MessageState.COMPLETED.value:
break
await asyncio.sleep(0.01)
assert tree.get_node("m2").state.value == MessageState.COMPLETED.value
assert tree.get_node("m2").parent_id == "m1"
mock_cli_manager.get_or_create_session.assert_called_with(session_id=None)
mock_session2.start_task.assert_called_with(
"message 2", session_id="sess1", fork_session=True
)
@pytest.mark.asyncio
async def test_error_propagation_chain(
handler_integration, mock_platform, mock_cli_manager, incoming_message_factory
):
msg1 = incoming_message_factory(text="m1", message_id="m1")
mock_platform.queue_send_message.return_value = "s1"
mock_session1 = MagicMock()
mock_session1.start_task.return_value = mock_async_gen(
[{"type": "error", "error": {"message": "failed"}}]
)
mock_cli_manager.get_or_create_session.return_value = (
mock_session1,
"sess1",
False,
)
await handler_integration.handle_message(msg1)
tree = handler_integration.tree_queue.get_tree_for_node("m1")
msg2 = incoming_message_factory(
text="m2", message_id="m2", reply_to_message_id="m1"
)
await handler_integration.handle_message(msg2)
# Wait for m1 to fail
for _ in range(20):
if tree.get_node("m1").state.value == MessageState.ERROR.value:
break
await asyncio.sleep(0.01)
# Give a tiny bit of time for propagation and skipping in processor
await asyncio.sleep(0.05)
assert tree.get_node("m1").state.value == MessageState.ERROR.value
assert tree.get_node("m2").state.value == MessageState.ERROR.value
assert "Parent failed" in tree.get_node("m2").error_message
@pytest.mark.asyncio
async def test_concurrent_replies_to_different_trees(
handler_integration, mock_platform, mock_cli_manager, incoming_message_factory
):
msg1 = incoming_message_factory(text="t1", message_id="t1")
msg2 = incoming_message_factory(text="t2", message_id="t2")
mock_session1 = MagicMock()
mock_session1.start_task.return_value = mock_async_gen(
[{"type": "exit", "code": 0}]
)
mock_session2 = MagicMock()
mock_session2.start_task.return_value = mock_async_gen(
[{"type": "exit", "code": 0}]
)
mock_cli_manager.get_or_create_session.side_effect = [
(mock_session1, "s1", False),
(mock_session2, "s2", False),
]
await handler_integration.handle_message(msg1)
await handler_integration.handle_message(msg2)
# Wait for both
for _ in range(20):
node1 = handler_integration.tree_queue.get_node("t1")
node2 = handler_integration.tree_queue.get_node("t2")
if (
node1
and node2
and node1.state.value == MessageState.COMPLETED.value
and node2.state.value == MessageState.COMPLETED.value
):
break
await asyncio.sleep(0.01)
assert (
handler_integration.tree_queue.get_node("t1").state.value
== MessageState.COMPLETED.value
)
assert (
handler_integration.tree_queue.get_node("t2").state.value
== MessageState.COMPLETED.value
)
|