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
33 changes: 10 additions & 23 deletions gateway/platforms/weixin.py
Original file line number Diff line number Diff line change
Expand Up @@ -1983,29 +1983,16 @@ 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:
last_result: Optional[SendResult] = None
cleaned = live_adapter.format_message(message)
if cleaned:
last_result = await live_adapter.send(chat_id, cleaned)
if not last_result.success:
return {"error": f"Weixin send failed: {last_result.error}"}

for media_path, _is_voice in media_files or []:
ext = Path(media_path).suffix.lower()
if ext in {".jpg", ".jpeg", ".png", ".gif", ".webp", ".bmp"}:
last_result = await live_adapter.send_image_file(chat_id, media_path)
else:
last_result = await live_adapter.send_document(chat_id, media_path)
if not last_result.success:
return {"error": f"Weixin media send failed: {last_result.error}"}

return {
"success": True,
"platform": "weixin",
"chat_id": chat_id,
"message_id": last_result.message_id if last_result else None,
"context_token_used": bool(context_token),
}
# The live adapter's aiohttp sessions are bound to the gateway event loop.
# Reusing them directly from cron/send_message helper paths can run under a
# different event loop (or even a different thread), which trips aiohttp
# timeout/task checks with errors like:
# "Timeout context manager should be used inside a task"
#
# For one-shot delivery helpers, prefer an isolated session instead of
# cross-loop adapter reuse. This keeps cron/send_message delivery safe even
# while the gateway is running.
live_adapter = None

async with aiohttp.ClientSession(trust_env=True, connector=_make_ssl_connector()) as session:
adapter = WeixinAdapter(
Expand Down
36 changes: 35 additions & 1 deletion tests/gateway/test_weixin.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@
from gateway.config import GatewayConfig, HomeChannel, Platform, _apply_env_overrides
from gateway.platforms.base import SendResult
from gateway.platforms import weixin
from gateway.platforms.weixin import ContextTokenStore, WeixinAdapter
from gateway.platforms.weixin import ContextTokenStore, WeixinAdapter, send_weixin_direct
from tools.send_message_tool import _parse_target_ref, _send_to_platform


Expand Down Expand Up @@ -308,6 +308,40 @@ def test_send_to_platform_routes_weixin_media_to_native_helper(self, send_weixin
media_files=[("/tmp/demo.png", False)],
)

def test_send_weixin_direct_avoids_cross_loop_live_adapter_reuse(self):
class _DummyClientSession:
def __init__(self, *args, **kwargs):
self.closed = False

async def __aenter__(self):
return self

async def __aexit__(self, exc_type, exc, tb):
self.closed = True
return False

live_adapter = AsyncMock()
live_adapter._send_session = type("_LiveSession", (), {"closed": False})()
live_adapter.send.side_effect = AssertionError("should not reuse live adapter across loops")

with patch.dict(weixin._LIVE_ADAPTERS, {"test-token": live_adapter}, clear=True), \
patch.object(weixin.ContextTokenStore, "restore", return_value=None), \
patch.object(weixin.ContextTokenStore, "get", return_value=None), \
patch("gateway.platforms.weixin.aiohttp.ClientSession", _DummyClientSession), \
patch.object(WeixinAdapter, "send", new=AsyncMock(return_value=SendResult(success=True, message_id="msg-direct"))):
result = asyncio.run(
send_weixin_direct(
extra={"account_id": "test-account"},
token="test-token",
chat_id="wxid_test123",
message="hello from cron",
)
)

assert result["success"] is True
assert result["message_id"] == "msg-direct"
live_adapter.send.assert_not_awaited()


class TestWeixinChunkDelivery:
def _connected_adapter(self) -> WeixinAdapter:
Expand Down