Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 17 additions & 1 deletion gateway/platforms/weixin.py
Original file line number Diff line number Diff line change
Expand Up @@ -2005,7 +2005,23 @@ async def send_weixin_direct(

live_adapter = _LIVE_ADAPTERS.get(resolved_token)
send_session = getattr(live_adapter, '_send_session', None)
if live_adapter is not None and send_session is not None and not send_session.closed:
can_reuse = (
live_adapter is not None
and send_session is not None
and not send_session.closed
)
# When called from a tool handler via _run_async (gateway context),
# the coroutine runs in a fresh thread with its own event loop.
# aiohttp sessions cannot be used across loops, so detect the
# mismatch and fall through to the standalone path.
if can_reuse and live_adapter._poll_task is not None:
try:
task_loop = live_adapter._poll_task.get_loop()
except AttributeError:
task_loop = None # Python <3.9 — fall through to safe path
if task_loop is not asyncio.get_running_loop():
can_reuse = False
if can_reuse:
last_result: Optional[SendResult] = None
cleaned = live_adapter.format_message(message)
if cleaned:
Expand Down
74 changes: 73 additions & 1 deletion tests/gateway/test_weixin.py
Original file line number Diff line number Diff line change
Expand Up @@ -757,4 +757,76 @@ def test_send_file_sets_voice_metadata_for_silk_payload(
assert voice_item.get("playtime", 0) == 0
assert voice_item["encode_type"] == 6
assert voice_item["sample_rate"] == 24000
assert voice_item["bits_per_sample"] == 16


class TestWeixinDirectSendCrossLoop:
"""Verify send_weixin_direct falls back to a fresh session when the
live adapter's event loop differs from the caller's (e.g. when
_run_async spawns a new thread + loop for tool handlers)."""

@patch.object(weixin, "_LIVE_ADAPTERS", {})
@patch.object(weixin, "aiohttp")
@patch("gateway.platforms.weixin._make_ssl_connector")
def test_falls_back_to_new_session_when_no_live_adapter(self, _mock_connector, mock_aiohttp):
"""No live adapter → must create a fresh session (Path 2)."""
mock_session = AsyncMock()
mock_session.closed = False
mock_aiohttp.ClientSession.return_value.__aenter__ = AsyncMock(return_value=mock_session)
mock_aiohttp.ClientSession.return_value.__aexit__ = AsyncMock(return_value=False)

# WeixinAdapter.send is async — mock it
with patch.object(WeixinAdapter, "send", new_callable=AsyncMock) as send_mock:
send_mock.return_value = SendResult(success=True, message_id="msg-1")
result = asyncio.run(
weixin.send_weixin_direct(
extra={"account_id": "acct1", "base_url": "https://ilink.example.com", "cdn_base_url": "https://cdn.example.com"},
token="tok",
chat_id="wxid_user",
message="hello",
)
)
assert result["success"] is True

@patch.object(weixin, "aiohttp")
@patch("gateway.platforms.weixin._make_ssl_connector")
def test_falls_back_to_new_session_on_loop_mismatch(self, _mock_connector, mock_aiohttp):
"""Live adapter exists but in a different event loop → must create
a fresh session (Path 2)."""
# Create a fake poll_task in a different event loop
async def _noop():
pass

other_loop = asyncio.new_event_loop()
try:
other_task = other_loop.create_task(_noop())
finally:
other_loop.close()

fake_session = AsyncMock()
fake_session.closed = False
fake_adapter = AsyncMock()
fake_adapter._send_session = fake_session
fake_adapter._poll_task = other_task

weixin._LIVE_ADAPTERS["tok"] = fake_adapter
try:
mock_session = AsyncMock()
mock_session.closed = False
mock_aiohttp.ClientSession.return_value.__aenter__ = AsyncMock(return_value=mock_session)
mock_aiohttp.ClientSession.return_value.__aexit__ = AsyncMock(return_value=False)

with patch.object(WeixinAdapter, "send", new_callable=AsyncMock) as send_mock:
send_mock.return_value = SendResult(success=True, message_id="msg-2")
result = asyncio.run(
weixin.send_weixin_direct(
extra={"account_id": "acct1", "base_url": "https://ilink.example.com", "cdn_base_url": "https://cdn.example.com"},
token="tok",
chat_id="wxid_user",
message="hello",
)
)
assert result["success"] is True
# send should NOT have been called on the fake adapter
fake_adapter.send.assert_not_awaited()
finally:
weixin._LIVE_ADAPTERS.pop("tok", None)