Compare commits

..

4 Commits

Author SHA1 Message Date
fdc310 1ed107c9d5 feat(dynamic-form): enhance select handling with empty option support and UUID filtering 2026-07-06 13:36:58 +08:00
fdc310 0cce418956 feat(itchat): improve login session handling and error reporting 2026-07-02 14:40:13 +08:00
fdc310 d4e8ccd161 feat: enhance itchat adapter with runtime status tracking and UI integration 2026-07-02 13:56:24 +08:00
fdc310 78fb40a28a feat: add itchat-uos WeChat adapter with QR code login
- Add itchat-uos adapter supporting personal WeChat via QR code login
- Implement message/event converters for text, image, voice, sharing types
- Bridge sync itchat callbacks to async LangBot pipeline via asyncio
- Add QR login API endpoints with session management
- Add frontend QR code login dialog integration
- Fix plugin connector handler attribute check
- Use fresh Core instance per login to avoid singleton state pollution
2026-07-01 18:03:26 +08:00
55 changed files with 2023 additions and 4620 deletions
+3 -2
View File
@@ -1,6 +1,6 @@
[project]
name = "langbot"
version = "4.10.5"
version = "4.10.4"
description = "Production-grade platform for building agentic IM bots"
readme = "README.md"
license-files = ["LICENSE"]
@@ -22,6 +22,7 @@ dependencies = [
"discord-py>=2.5.2",
"pynacl>=1.5.0", # Required for Discord voice support
"gewechat-client>=0.1.5",
"itchat-uos>=1.5.0.dev",
"lark-oapi>=1.5.5",
"mcp>=1.25.0",
"nakuru-project-idk>=0.0.2.1",
@@ -70,7 +71,7 @@ dependencies = [
"chromadb>=1.0.0,<2.0.0",
"qdrant-client (>=1.15.1,<2.0.0)",
"pyseekdb==1.1.0.post3",
"langbot-plugin==0.4.10",
"langbot-plugin==0.4.6",
"asyncpg>=0.30.0",
"line-bot-sdk>=3.19.0",
"matrix-nio>=0.25.2",
@@ -138,39 +138,6 @@ class MonitoringRouterGroup(group.RouterGroup):
}
)
@self.route('/tool-calls', methods=['GET'], auth_type=group.AuthType.USER_TOKEN)
async def get_tool_calls() -> str:
"""Get tool call records"""
bot_ids = quart.request.args.getlist('botId')
pipeline_ids = quart.request.args.getlist('pipelineId')
session_ids = quart.request.args.getlist('sessionId')
start_time_str = quart.request.args.get('startTime')
end_time_str = quart.request.args.get('endTime')
limit = int(quart.request.args.get('limit', 100))
offset = int(quart.request.args.get('offset', 0))
start_time = parse_iso_datetime(start_time_str)
end_time = parse_iso_datetime(end_time_str)
tool_calls, total = await self.ap.monitoring_service.get_tool_calls(
bot_ids=bot_ids if bot_ids else None,
pipeline_ids=pipeline_ids if pipeline_ids else None,
session_ids=session_ids if session_ids else None,
start_time=start_time,
end_time=end_time,
limit=limit,
offset=offset,
)
return self.success(
data={
'tool_calls': tool_calls,
'total': total,
'limit': limit,
'offset': offset,
}
)
@self.route('/embedding-calls', methods=['GET'], auth_type=group.AuthType.USER_TOKEN)
async def get_embedding_calls() -> str:
"""Get embedding call records"""
@@ -317,16 +284,6 @@ class MonitoringRouterGroup(group.RouterGroup):
offset=0,
)
# Get tool calls
tool_calls, tool_calls_total = await self.ap.monitoring_service.get_tool_calls(
bot_ids=bot_ids if bot_ids else None,
pipeline_ids=pipeline_ids if pipeline_ids else None,
start_time=start_time,
end_time=end_time,
limit=limit,
offset=0,
)
# Get sessions
sessions, sessions_total = await self.ap.monitoring_service.get_sessions(
bot_ids=bot_ids if bot_ids else None,
@@ -361,14 +318,12 @@ class MonitoringRouterGroup(group.RouterGroup):
'overview': overview,
'messages': messages,
'llmCalls': llm_calls,
'toolCalls': tool_calls,
'embeddingCalls': embedding_calls,
'sessions': sessions,
'errors': errors,
'totalCount': {
'messages': messages_total,
'llmCalls': llm_calls_total,
'toolCalls': tool_calls_total,
'embeddingCalls': embedding_calls_total,
'sessions': sessions_total,
'errors': errors_total,
@@ -1,6 +1,7 @@
import quart
import mimetypes
import asyncio
import os
from ... import group
from langbot.pkg.utils import importutil
@@ -650,3 +651,224 @@ class AdaptersRouterGroup(group.RouterGroup):
if session and session.get('task') and not session['task'].done():
session['task'].cancel()
return self.success(data={})
# -----------------------------------------------------------------------
# Itchat WeChat QR Code Login
# -----------------------------------------------------------------------
_itchat_login_sessions: dict = {}
_ITCHAT_SESSION_TTL = 600 # 10 minutes (allows multiple QR regenerations)
def _cleanup_expired_itchat_sessions():
import time
now = time.time()
expired = [
sid for sid, s in _itchat_login_sessions.items() if now - s.get('created_at', 0) > _ITCHAT_SESSION_TTL
]
for sid in expired:
session = _itchat_login_sessions.pop(sid, None)
if session:
core = session.get('core')
if core:
try:
core.alive = False
core.isLogging = False
except Exception:
pass
@self.route('/itchat/login', methods=['POST'])
async def _() -> str:
"""Start itchat WeChat QR code login. Returns session_id + QR code data URL."""
import uuid
import time
import base64
import threading
_cleanup_expired_itchat_sessions()
session_id = str(uuid.uuid4())
loop = asyncio.get_running_loop()
status_dir = os.path.join('data', 'itchat')
os.makedirs(status_dir, exist_ok=True)
qr_path = os.path.join(status_dir, f'{session_id}-QR.png')
session = {
'status': 'pending',
'qr_data_url': None,
'expire_at': None,
'nickname': None,
'error': None,
'created_at': time.time(),
'thread': None,
'logged_in': threading.Event(),
'core': None,
}
_itchat_login_sessions[session_id] = session
def _run_itchat_login():
try:
from itchat.core import Core
from itchat.content import TEXT as _TEXT
from langbot.pkg.platform.sources.itchat import ItchatAdapter
for f in (qr_path,):
try:
os.remove(f)
except OSError:
pass
_core = Core()
session['core'] = _core
def on_login():
try:
_core.get_friends(update=True)
user_info = _core.loginInfo.get('User', {})
nick = ItchatAdapter._get_obj_value(user_info, 'NickName', 'unknown')
wxid = ItchatAdapter._get_obj_value(user_info, 'UserName')
except Exception:
nick = 'unknown'
wxid = ''
session['nickname'] = nick
session['wxid'] = wxid
print(f'[itchat-login] Login success: {nick}', flush=True)
# Dump login status so the adapter can hot-reload it
try:
if not wxid:
raise ValueError('Unable to detect WeChat wxid after login')
account_status_path = ItchatAdapter.login_status_path_for_account(wxid)
_core.dump_login_status(account_status_path)
session['login_status_path'] = account_status_path
session['status'] = 'success'
print(f'[itchat-login] Session saved to {account_status_path}', flush=True)
except Exception as e:
session['status'] = 'error'
session['error'] = str(e)
print(f'[itchat-login] Failed to save session: {e}', flush=True)
finally:
session['logged_in'].set()
# Stop the message loop - we only needed the session for QR login
_core.alive = False
def on_qr(**kwargs):
qr_bytes = kwargs.get('qrcode', b'')
status = kwargs.get('status', '')
print(f'[itchat-login] QR callback: status={status}, bytes={len(qr_bytes)}', flush=True)
if status == '200':
return
# Only update QR image on new QR generation (status='0')
# or when status changes to '408' (timeout, QR may refresh)
if qr_bytes and status == '0':
b64 = base64.b64encode(qr_bytes).decode('utf-8')
def _update():
session['qr_data_url'] = f'data:image/png;base64,{b64}'
session['expire_at'] = time.time() + 120
session['status'] = 'waiting'
loop.call_soon_threadsafe(_update)
# Register a dummy text handler
@_core.msg_register([_TEXT])
def _dummy(msg):
pass
print('[itchat-login] Step 3: Calling auto_login...', flush=True)
_core.auto_login(
hotReload=False,
loginCallback=on_login,
qrCallback=on_qr,
)
print('[itchat-login] Step 4: auto_login returned, starting run...', flush=True)
_core.run(blockThread=True)
print('[itchat-login] Step 5: run() returned', flush=True)
except SystemExit as e:
print(f'[itchat-login] SystemExit: {e}', flush=True)
session['status'] = 'error'
session['error'] = f'itchat exited: {e}'
session['logged_in'].set()
except Exception as e:
import traceback
print(f'[itchat-login] Exception: {traceback.format_exc()}', flush=True)
session['status'] = 'error'
session['error'] = str(e)
session['logged_in'].set()
t = threading.Thread(target=_run_itchat_login, daemon=True)
t.start()
session['thread'] = t
# Wait for QR code to be ready (max 15 seconds)
for _ in range(30):
if session['qr_data_url'] or session['error'] or session['status'] == 'success':
break
await asyncio.sleep(0.5)
if session['error']:
return self.http_status(502, -1, session['error'])
if session['status'] == 'success':
return self.success(
data={
'session_id': session_id,
'status': 'success',
'nickname': session['nickname'],
'wxid': session.get('wxid', ''),
}
)
if not session['qr_data_url']:
session['status'] = 'error'
session['error'] = 'Timeout waiting for QR code'
return self.http_status(504, -1, 'Timeout waiting for QR code')
return self.success(
data={
'session_id': session_id,
'qr_data_url': session['qr_data_url'],
'expire_at': session['expire_at'],
}
)
@self.route('/itchat/login/status/<session_id>', methods=['GET'])
async def _(session_id: str) -> str:
"""Poll itchat login status."""
session = _itchat_login_sessions.get(session_id)
if not session:
return self.http_status(404, -1, 'Session not found')
data = {
'status': session['status'],
'qr_data_url': session['qr_data_url'],
'expire_at': session['expire_at'],
}
if session['status'] == 'success':
data['nickname'] = session.get('nickname', '')
data['wxid'] = session.get('wxid', '')
_itchat_login_sessions.pop(session_id, None)
elif session['status'] == 'error':
data['error'] = session['error']
_itchat_login_sessions.pop(session_id, None)
return self.success(data=data)
@self.route('/itchat/login/<session_id>', methods=['DELETE'])
async def _(session_id: str) -> str:
"""Cancel and clean up an itchat login session."""
session = _itchat_login_sessions.pop(session_id, None)
if session:
core = session.get('core')
if core:
try:
core.alive = False
core.isLogging = False
except Exception:
pass
thread = session.get('thread')
if thread and thread.is_alive():
# Thread is daemon, will die with the process
pass
return self.success(data={})
@@ -29,11 +29,11 @@ class MCPRouterGroup(group.RouterGroup):
traceback.print_exc()
return self.http_status(500, -1, f'Failed to create MCP server: {str(e)}')
@self.route(
'/servers/<path:server_name>', methods=['GET', 'PUT', 'DELETE'], auth_type=group.AuthType.USER_TOKEN
)
@self.route('/servers/<server_name>', methods=['GET', 'PUT', 'DELETE'], auth_type=group.AuthType.USER_TOKEN)
async def _(server_name: str) -> str:
"""获取、更新或删除MCP服务器配置"""
from urllib.parse import unquote
server_name = unquote(server_name)
server_data = await self.ap.mcp_service.get_mcp_server_by_name(server_name)
@@ -58,15 +58,17 @@ class MCPRouterGroup(group.RouterGroup):
except Exception as e:
return self.http_status(500, -1, f'Failed to delete MCP server: {str(e)}')
@self.route('/servers/<path:server_name>/test', methods=['POST'], auth_type=group.AuthType.USER_TOKEN)
@self.route('/servers/<server_name>/test', methods=['POST'], auth_type=group.AuthType.USER_TOKEN)
async def _(server_name: str) -> str:
"""测试MCP服务器连接"""
from urllib.parse import unquote
server_name = unquote(server_name)
server_data = await quart.request.json
task_id = await self.ap.mcp_service.test_mcp_server(server_name=server_name, server_data=server_data)
return self.success(data={'task_id': task_id})
@self.route('/servers/<path:server_name>/resources', methods=['GET'], auth_type=group.AuthType.USER_TOKEN)
@self.route('/servers/<server_name>/resources', methods=['GET'], auth_type=group.AuthType.USER_TOKEN)
async def _(server_name: str) -> str:
"""Get resources from an MCP server"""
server_name = unquote(server_name)
@@ -84,9 +86,7 @@ class MCPRouterGroup(group.RouterGroup):
except Exception as e:
return self.http_status(500, -1, f'Failed to get resources: {str(e)}')
@self.route(
'/servers/<path:server_name>/resource-templates', methods=['GET'], auth_type=group.AuthType.USER_TOKEN
)
@self.route('/servers/<server_name>/resource-templates', methods=['GET'], auth_type=group.AuthType.USER_TOKEN)
async def _(server_name: str) -> str:
"""Get resource templates from an MCP server"""
server_name = unquote(server_name)
@@ -96,20 +96,7 @@ class MCPRouterGroup(group.RouterGroup):
except Exception as e:
return self.http_status(500, -1, f'Failed to get resource templates: {str(e)}')
@self.route('/servers/<path:server_name>/logs', methods=['GET'], auth_type=group.AuthType.USER_TOKEN)
async def _(server_name: str) -> str:
"""Get logs from an MCP server"""
server_name = unquote(server_name)
try:
limit = int(quart.request.args.get('limit', 200))
except (TypeError, ValueError):
limit = 200
limit = min(limit, 500)
level = quart.request.args.get('level') or None
logs = await self.ap.mcp_service.get_mcp_server_logs(server_name, limit=limit, level=level)
return self.success(data={'logs': logs})
@self.route('/servers/<path:server_name>/resources/read', methods=['POST'], auth_type=group.AuthType.USER_TOKEN)
@self.route('/servers/<server_name>/resources/read', methods=['POST'], auth_type=group.AuthType.USER_TOKEN)
async def _(server_name: str) -> str:
"""Read a resource from an MCP server"""
server_name = unquote(server_name)
+2
View File
@@ -57,6 +57,8 @@ class BotService:
runtime_bot = await self.ap.platform_mgr.get_bot_by_uuid(bot_uuid)
if runtime_bot is not None:
adapter_runtime_values['bot_account_id'] = runtime_bot.adapter.bot_account_id
if hasattr(runtime_bot.adapter, 'get_runtime_status'):
adapter_runtime_values['runtime_status'] = runtime_bot.adapter.get_runtime_status()
# Webhook URL for unified webhook adapters (independent of bot running state)
if persistence_bot['adapter'] in [
@@ -243,7 +243,6 @@ class MaintenanceService:
tables = {
'messages': persistence_monitoring.MonitoringMessage.id,
'llm_calls': persistence_monitoring.MonitoringLLMCall.id,
'tool_calls': persistence_monitoring.MonitoringToolCall.id,
'embedding_calls': persistence_monitoring.MonitoringEmbeddingCall.id,
'errors': persistence_monitoring.MonitoringError.id,
'sessions': persistence_monitoring.MonitoringSession.session_id,
+2 -41
View File
@@ -48,17 +48,6 @@ class MCPService:
if total_extensions >= max_extensions:
raise ValueError(f'Maximum number of extensions ({max_extensions}) reached')
server_name = str(server_data.get('name') or '').strip()
if not server_name:
raise ValueError('MCP server name is required')
server_data['name'] = server_name
existing_result = await self.ap.persistence_mgr.execute_async(
sqlalchemy.select(persistence_mcp.MCPServer).where(persistence_mcp.MCPServer.name == server_name)
)
if existing_result.first() is not None:
raise ValueError(f'MCP server already exists: {server_name}')
server_data['uuid'] = str(uuid.uuid4())
await self.ap.persistence_mgr.execute_async(sqlalchemy.insert(persistence_mcp.MCPServer).values(server_data))
@@ -188,22 +177,10 @@ class MCPService:
persisted_session = runtime_mcp_session
async def _refresh_and_report() -> None:
# Testing a persisted server should REUSE its live shared-session
# process, not rebuild it. Try a lightweight refresh (a real
# list_tools probe over the existing connection) first; only fall
# back to a full start() when the session has no live connection
# to probe (never connected, or the process is actually gone).
needs_start = persisted_session.status == MCPSessionStatus.ERROR or persisted_session.session is None
if needs_start:
if persisted_session.status == MCPSessionStatus.ERROR:
await persisted_session.start()
else:
try:
await persisted_session.refresh()
except Exception:
# The live connection was stale/dropped: reconnect once
# (reusing the live managed process where possible) and
# re-probe, instead of reporting a false failure.
await persisted_session.start()
await persisted_session.refresh()
# Surface the discovered tools so the config page can render them
# even for an already-hosted server.
ctx.metadata['runtime_info'] = persisted_session.get_runtime_info_dict()
@@ -244,19 +221,3 @@ class MCPService:
context=ctx,
)
return wrapper.id
async def get_mcp_server_logs(self, server_name: str, limit: int = 200, level: str | None = None) -> list[dict]:
"""Get recent log lines captured from the MCP server's stderr."""
session = self.ap.tool_mgr.mcp_tool_loader.get_session(server_name)
if not session:
return []
# Get logs from the session's buffer
logs = list(session._log_buffer)
# Filter by level if specified
if level:
logs = [log for log in logs if log.get('level') == level]
# Return the most recent 'limit' logs
return logs[-limit:]
@@ -2,7 +2,6 @@ from __future__ import annotations
import uuid
import datetime
import json
import sqlalchemy
from ....core import app
@@ -51,12 +50,6 @@ class MonitoringService:
persistence_monitoring.MonitoringLLMCall.timestamp,
persistence_monitoring.MonitoringLLMCall.id,
),
(
'monitoring_tool_calls',
persistence_monitoring.MonitoringToolCall,
persistence_monitoring.MonitoringToolCall.timestamp,
persistence_monitoring.MonitoringToolCall.id,
),
(
'monitoring_embedding_calls',
persistence_monitoring.MonitoringEmbeddingCall,
@@ -138,68 +131,6 @@ class MonitoringService:
await autocommit_conn.execute(sqlalchemy.text('PRAGMA wal_checkpoint(TRUNCATE)'))
await autocommit_conn.execute(sqlalchemy.text('VACUUM'))
def _serialize_tool_payload(self, payload: object, max_length: int = 20000) -> str | None:
"""Serialize tool arguments/results for monitoring storage."""
if payload is None:
return None
if isinstance(payload, str):
text = payload
else:
try:
text = json.dumps(payload, ensure_ascii=False, default=str)
except Exception:
text = str(payload)
if len(text) <= max_length:
return text
return f'{text[:max_length]}... [truncated {len(text) - max_length} chars]'
async def _get_message_for_tool_context(
self,
message_id: str | None = None,
session_id: str | None = None,
):
if message_id:
result = await self.ap.persistence_mgr.execute_async(
sqlalchemy.select(persistence_monitoring.MonitoringMessage).where(
persistence_monitoring.MonitoringMessage.id == message_id
)
)
row = result.first()
if row:
return row[0]
if not session_id:
return None
user_query = (
sqlalchemy.select(persistence_monitoring.MonitoringMessage)
.where(
sqlalchemy.and_(
persistence_monitoring.MonitoringMessage.session_id == session_id,
persistence_monitoring.MonitoringMessage.role == 'user',
)
)
.order_by(persistence_monitoring.MonitoringMessage.timestamp.desc())
.limit(1)
)
result = await self.ap.persistence_mgr.execute_async(user_query)
row = result.first()
if row:
return row[0]
any_query = (
sqlalchemy.select(persistence_monitoring.MonitoringMessage)
.where(persistence_monitoring.MonitoringMessage.session_id == session_id)
.order_by(persistence_monitoring.MonitoringMessage.timestamp.desc())
.limit(1)
)
result = await self.ap.persistence_mgr.execute_async(any_query)
row = result.first()
return row[0] if row else None
# ========== Recording Methods ==========
async def record_message(
@@ -289,57 +220,6 @@ class MonitoringService:
return call_id
async def record_tool_call(
self,
tool_name: str,
tool_source: str,
duration: int,
status: str = 'success',
bot_id: str | None = None,
bot_name: str | None = None,
pipeline_id: str | None = None,
pipeline_name: str | None = None,
session_id: str | None = None,
message_id: str | None = None,
arguments: object | None = None,
result: object | None = None,
error_message: str | None = None,
) -> str:
"""Record a tool call."""
context_message = await self._get_message_for_tool_context(message_id=message_id, session_id=session_id)
if context_message:
bot_id = bot_id or context_message.bot_id
bot_name = bot_name or context_message.bot_name
pipeline_id = pipeline_id or context_message.pipeline_id
pipeline_name = pipeline_name or context_message.pipeline_name
session_id = session_id or context_message.session_id
message_id = message_id or context_message.id
call_id = str(uuid.uuid4())
call_data = {
'id': call_id,
'timestamp': datetime.datetime.now(datetime.timezone.utc).replace(tzinfo=None),
'tool_name': tool_name,
'tool_source': tool_source,
'duration': max(0, duration),
'status': status,
'bot_id': bot_id or 'unknown',
'bot_name': bot_name or 'Unknown',
'pipeline_id': pipeline_id or 'unknown',
'pipeline_name': pipeline_name or 'Unknown',
'session_id': session_id,
'message_id': message_id,
'arguments': self._serialize_tool_payload(arguments),
'result': self._serialize_tool_payload(result),
'error_message': self._serialize_tool_payload(error_message),
}
await self.ap.persistence_mgr.execute_async(
sqlalchemy.insert(persistence_monitoring.MonitoringToolCall).values(call_data)
)
return call_id
async def record_embedding_call(
self,
model_name: str,
@@ -869,58 +749,6 @@ class MonitoringService:
total,
)
async def get_tool_calls(
self,
bot_ids: list[str] | None = None,
pipeline_ids: list[str] | None = None,
session_ids: list[str] | None = None,
start_time: datetime.datetime | None = None,
end_time: datetime.datetime | None = None,
limit: int = 100,
offset: int = 0,
) -> tuple[list[dict], int]:
"""Get tool calls with filters"""
conditions = []
if bot_ids:
conditions.append(persistence_monitoring.MonitoringToolCall.bot_id.in_(bot_ids))
if pipeline_ids:
conditions.append(persistence_monitoring.MonitoringToolCall.pipeline_id.in_(pipeline_ids))
if session_ids:
conditions.append(persistence_monitoring.MonitoringToolCall.session_id.in_(session_ids))
if start_time:
conditions.append(persistence_monitoring.MonitoringToolCall.timestamp >= start_time)
if end_time:
conditions.append(persistence_monitoring.MonitoringToolCall.timestamp <= end_time)
count_query = sqlalchemy.select(sqlalchemy.func.count(persistence_monitoring.MonitoringToolCall.id))
if conditions:
count_query = count_query.where(sqlalchemy.and_(*conditions))
count_result = await self.ap.persistence_mgr.execute_async(count_query)
total = count_result.scalar() or 0
query = sqlalchemy.select(persistence_monitoring.MonitoringToolCall).order_by(
persistence_monitoring.MonitoringToolCall.timestamp.desc()
)
if conditions:
query = query.where(sqlalchemy.and_(*conditions))
query = query.limit(limit).offset(offset)
result = await self.ap.persistence_mgr.execute_async(query)
tool_calls_rows = result.all()
return (
[
self.ap.persistence_mgr.serialize_model(
persistence_monitoring.MonitoringToolCall, row[0] if isinstance(row, tuple) else row
)
for row in tool_calls_rows
],
total,
)
async def get_embedding_calls(
self,
start_time: datetime.datetime | None = None,
@@ -1143,34 +971,6 @@ class MonitoringService:
else:
error_llm_calls += 1
# Get tool calls for this session
tool_query = (
sqlalchemy.select(persistence_monitoring.MonitoringToolCall)
.where(persistence_monitoring.MonitoringToolCall.session_id == session_id)
.order_by(persistence_monitoring.MonitoringToolCall.timestamp.asc())
)
tool_result = await self.ap.persistence_mgr.execute_async(tool_query)
tool_rows = tool_result.all()
tool_calls = [
self.ap.persistence_mgr.serialize_model(
persistence_monitoring.MonitoringToolCall, row[0] if isinstance(row, tuple) else row
)
for row in tool_rows
]
total_tool_calls = len(tool_rows)
success_tool_calls = 0
error_tool_calls = 0
total_tool_duration = 0
for row in tool_rows:
tool_call = row[0] if isinstance(row, tuple) else row
total_tool_duration += tool_call.duration
if tool_call.status == 'success':
success_tool_calls += 1
else:
error_tool_calls += 1
# Get errors for this session
error_query = (
sqlalchemy.select(persistence_monitoring.MonitoringError)
@@ -1214,14 +1014,6 @@ class MonitoringService:
'total_tokens': total_tokens,
'average_duration_ms': int(total_duration / total_llm_calls) if total_llm_calls > 0 else 0,
},
'tool_calls': tool_calls,
'tool_stats': {
'total_calls': total_tool_calls,
'success_calls': success_tool_calls,
'error_calls': error_tool_calls,
'total_duration_ms': total_tool_duration,
'average_duration_ms': int(total_tool_duration / total_tool_calls) if total_tool_calls > 0 else 0,
},
'errors': errors,
'session_duration_seconds': session_duration_seconds,
}
@@ -49,28 +49,6 @@ class MonitoringLLMCall(Base):
message_id = sqlalchemy.Column(sqlalchemy.String(255), nullable=True, index=True) # Associated message ID
class MonitoringToolCall(Base):
"""Tool call records"""
__tablename__ = 'monitoring_tool_calls'
id = sqlalchemy.Column(sqlalchemy.String(255), primary_key=True)
timestamp = sqlalchemy.Column(sqlalchemy.DateTime, nullable=False, index=True)
tool_name = sqlalchemy.Column(sqlalchemy.String(255), nullable=False)
tool_source = sqlalchemy.Column(sqlalchemy.String(50), nullable=False) # native, plugin, mcp, skill
duration = sqlalchemy.Column(sqlalchemy.Integer, nullable=False) # milliseconds
status = sqlalchemy.Column(sqlalchemy.String(50), nullable=False) # success, error
bot_id = sqlalchemy.Column(sqlalchemy.String(255), nullable=False, index=True)
bot_name = sqlalchemy.Column(sqlalchemy.String(255), nullable=False)
pipeline_id = sqlalchemy.Column(sqlalchemy.String(255), nullable=False, index=True)
pipeline_name = sqlalchemy.Column(sqlalchemy.String(255), nullable=False)
session_id = sqlalchemy.Column(sqlalchemy.String(255), nullable=True, index=True)
message_id = sqlalchemy.Column(sqlalchemy.String(255), nullable=True, index=True)
arguments = sqlalchemy.Column(sqlalchemy.Text, nullable=True)
result = sqlalchemy.Column(sqlalchemy.Text, nullable=True)
error_message = sqlalchemy.Column(sqlalchemy.Text, nullable=True)
class MonitoringSession(Base):
"""Session tracking records"""
@@ -1,17 +0,0 @@
from langbot.pkg.entity.persistence import monitoring as persistence_monitoring
from .. import migration
@migration.migration_class(26)
class DBMigrateMonitoringToolCalls(migration.DBMigration):
"""Add monitoring_tool_calls table"""
async def upgrade(self):
"""Upgrade"""
async with self.ap.persistence_mgr.get_db_engine().begin() as conn:
await conn.run_sync(persistence_monitoring.MonitoringToolCall.__table__.create, checkfirst=True)
async def downgrade(self):
"""Downgrade"""
async with self.ap.persistence_mgr.get_db_engine().begin() as conn:
await conn.run_sync(persistence_monitoring.MonitoringToolCall.__table__.drop, checkfirst=True)
+6 -110
View File
@@ -4,7 +4,6 @@ import asyncio
import traceback
import datetime
import json
import time
import aiocqhttp
import pydantic
@@ -17,14 +16,6 @@ from ...utils import image
import langbot_plugin.api.definition.abstract.platform.event_logger as abstract_platform_logger
_GROUP_NAME_CACHE_TTL_SECONDS = 3600
_GROUP_NAME_NEGATIVE_CACHE_TTL_SECONDS = 60
_GROUP_NAME_LOOKUP_TIMEOUT_SECONDS = 2
_GROUP_MEMBER_INFO_CACHE_TTL_SECONDS = 86400
_GROUP_MEMBER_INFO_NEGATIVE_CACHE_TTL_SECONDS = 600
_GROUP_MEMBER_INFO_LOOKUP_TIMEOUT_SECONDS = 2
def _normalize_base64_payload(value: str) -> str:
if value.startswith('base64://'):
return value.removeprefix('base64://')
@@ -33,21 +24,6 @@ def _normalize_base64_payload(value: str) -> str:
return value
def _get_field(data: dict, key: str, default: str = '') -> str:
value = data.get(key)
if value is None:
return default
return str(value)
def _get_group_member_name(sender: dict) -> str:
return _get_field(sender, 'card') or _get_field(sender, 'nickname') or _get_field(sender, 'user_id')
def _get_group_name_placeholder(group_id: typing.Union[int, str]) -> str:
return f'Group {group_id}'
class AiocqhttpMessageConverter(abstract_platform_adapter.AbstractMessageConverter):
@staticmethod
async def yiri2target(
@@ -359,96 +335,16 @@ class AiocqhttpMessageConverter(abstract_platform_adapter.AbstractMessageConvert
class AiocqhttpEventConverter(abstract_platform_adapter.AbstractEventConverter):
def __init__(self):
self._group_name_cache: dict[typing.Union[int, str], tuple[str, float]] = {}
self._group_name_negative_cache: dict[typing.Union[int, str], float] = {}
self._group_member_info_cache: dict[
tuple[typing.Union[int, str], typing.Union[int, str]], tuple[dict, float]
] = {}
self._group_member_info_negative_cache: dict[tuple[typing.Union[int, str], typing.Union[int, str]], float] = {}
@staticmethod
async def yiri2target(event: platform_events.MessageEvent, bot_account_id: int):
return event.source_platform_object
async def _get_group_name(self, group_id: typing.Union[int, str], bot=None) -> str:
now = time.monotonic()
if group_id in self._group_name_cache:
group_name, expires_at = self._group_name_cache[group_id]
if expires_at > now:
return group_name
del self._group_name_cache[group_id]
if group_id in self._group_name_negative_cache:
expires_at = self._group_name_negative_cache[group_id]
if expires_at > now:
return ''
del self._group_name_negative_cache[group_id]
if bot is None:
return ''
try:
group_info = await asyncio.wait_for(
bot.get_group_info(group_id=group_id),
timeout=_GROUP_NAME_LOOKUP_TIMEOUT_SECONDS,
)
except Exception:
self._group_name_negative_cache[group_id] = now + _GROUP_NAME_NEGATIVE_CACHE_TTL_SECONDS
return ''
group_name = _get_field(group_info, 'group_name') if isinstance(group_info, dict) else ''
if group_name:
self._group_name_cache[group_id] = (group_name, now + _GROUP_NAME_CACHE_TTL_SECONDS)
self._group_name_negative_cache.pop(group_id, None)
else:
self._group_name_negative_cache[group_id] = now + _GROUP_NAME_NEGATIVE_CACHE_TTL_SECONDS
return group_name
async def _get_group_member_info(
self,
group_id: typing.Union[int, str],
user_id: typing.Union[int, str],
bot=None,
) -> dict:
now = time.monotonic()
cache_key = (group_id, user_id)
if cache_key in self._group_member_info_cache:
member_info, expires_at = self._group_member_info_cache[cache_key]
if expires_at > now:
return member_info
del self._group_member_info_cache[cache_key]
if cache_key in self._group_member_info_negative_cache:
expires_at = self._group_member_info_negative_cache[cache_key]
if expires_at > now:
return {}
del self._group_member_info_negative_cache[cache_key]
if bot is None:
return {}
try:
member_info = await asyncio.wait_for(
bot.get_group_member_info(group_id=group_id, user_id=user_id),
timeout=_GROUP_MEMBER_INFO_LOOKUP_TIMEOUT_SECONDS,
)
except Exception:
self._group_member_info_negative_cache[cache_key] = now + _GROUP_MEMBER_INFO_NEGATIVE_CACHE_TTL_SECONDS
return {}
if isinstance(member_info, dict) and member_info:
self._group_member_info_cache[cache_key] = (
member_info,
now + _GROUP_MEMBER_INFO_CACHE_TTL_SECONDS,
)
self._group_member_info_negative_cache.pop(cache_key, None)
return member_info
self._group_member_info_negative_cache[cache_key] = now + _GROUP_MEMBER_INFO_NEGATIVE_CACHE_TTL_SECONDS
return {}
async def target2yiri(self, event: aiocqhttp.Event, bot=None):
@staticmethod
async def target2yiri(event: aiocqhttp.Event, bot=None):
yiri_chain = await AiocqhttpMessageConverter.target2yiri(event.message, event.message_id, bot)
if event.message_type == 'group':
permission = 'MEMBER'
group_name = await self._get_group_name(event.group_id, bot) or _get_group_name_placeholder(event.group_id)
special_title = _get_field(event.sender, 'title')
if not special_title:
member_info = await self._get_group_member_info(event.group_id, event.sender['user_id'], bot)
special_title = _get_field(member_info, 'title')
if 'role' in event.sender:
if event.sender['role'] == 'admin':
@@ -458,14 +354,14 @@ class AiocqhttpEventConverter(abstract_platform_adapter.AbstractEventConverter):
converted_event = platform_events.GroupMessage(
sender=platform_entities.GroupMember(
id=event.sender['user_id'], # message_seq 放哪?
member_name=_get_group_member_name(event.sender),
member_name=event.sender['nickname'],
permission=permission,
group=platform_entities.Group(
id=event.group_id,
name=group_name,
name=event.sender['nickname'],
permission=platform_entities.Permission.Member,
),
special_title=special_title,
special_title=event.sender['title'] if 'title' in event.sender else '',
),
message_chain=yiri_chain,
time=event.time,
@@ -489,7 +385,7 @@ class AiocqhttpAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter)
bot: aiocqhttp.CQHttp = pydantic.Field(exclude=True, default_factory=aiocqhttp.CQHttp)
message_converter: AiocqhttpMessageConverter = AiocqhttpMessageConverter()
event_converter: AiocqhttpEventConverter = pydantic.Field(default_factory=AiocqhttpEventConverter)
event_converter: AiocqhttpEventConverter = AiocqhttpEventConverter()
on_websocket_connection_event_cache: typing.List[typing.Callable[[aiocqhttp.Event], None]] = []
+772
View File
@@ -0,0 +1,772 @@
"""itchat-uos adapter for LangBot.
Uses the itchat-uos WeChat Web library to integrate personal WeChat accounts
with LangBot via QR code login.
Reference: https://github.com/littlecodersh/ItChat
UOS fork: https://github.com/why2lyj/ItChat-uos
"""
from __future__ import annotations
import asyncio
import base64
import os
import re
import tempfile
import threading
import time
import traceback
import typing
from itchat.content import TEXT, PICTURE, RECORDING, VIDEO, SHARING
import pydantic
from itchat.core import Core as ItchatCore
try:
import queue
except ImportError:
import Queue as queue
import langbot_plugin.api.definition.abstract.platform.adapter as abstract_platform_adapter
import langbot_plugin.api.definition.abstract.platform.event_logger as abstract_platform_logger
import langbot_plugin.api.entities.builtin.platform.entities as platform_entities
import langbot_plugin.api.entities.builtin.platform.events as platform_events
import langbot_plugin.api.entities.builtin.platform.message as platform_message
from langbot.pkg.platform.logger import EventLogger
class ItchatMessageConverter(abstract_platform_adapter.AbstractMessageConverter):
"""Converts between LangBot MessageChain and itchat message dicts."""
@staticmethod
async def yiri2target(
message_chain: platform_message.MessageChain,
) -> list[dict]:
"""LangBot MessageChain -> list of itchat-sendable items.
Each item is a dict with 'type' and the relevant content field.
The adapter's send_message() will call itchat.send() accordingly.
"""
items: list[dict] = []
for component in message_chain:
if isinstance(component, platform_message.Plain):
if component.text:
items.append({'type': 'text', 'content': component.text})
elif isinstance(component, platform_message.Image):
if component.base64:
items.append({'type': 'image', 'base64': component.base64})
elif component.url:
items.append({'type': 'image', 'url': component.url})
elif isinstance(component, platform_message.Voice):
if component.base64:
items.append({'type': 'voice', 'base64': component.base64})
elif component.url:
items.append({'type': 'voice', 'url': component.url})
elif isinstance(component, platform_message.File):
if component.base64:
items.append({'type': 'file', 'base64': component.base64, 'name': component.name or 'file'})
elif component.url:
items.append({'type': 'file', 'url': component.url, 'name': component.name or 'file'})
elif isinstance(component, platform_message.At):
items.append({'type': 'text', 'content': f'@{component.target} '})
elif isinstance(component, platform_message.AtAll):
items.append({'type': 'text', 'content': '@所有人 '})
elif isinstance(component, platform_message.Forward):
for node in component.node_list:
if node.message_chain:
items.extend(await ItchatMessageConverter.yiri2target(node.message_chain))
elif isinstance(component, platform_message.Unknown):
pass # skip unknown outbound
return items
@staticmethod
def target2yiri(msg: dict) -> platform_message.MessageChain:
"""Convert an itchat msg dict to a LangBot MessageChain."""
components: list[platform_message.MessageComponent] = []
msg_type = msg.get('Type', '')
if msg_type == 'Text':
text = msg.get('Text', '')
if text:
components.append(platform_message.Plain(text=text))
elif msg_type == 'Picture':
try:
temp_dir = tempfile.gettempdir()
file_path = os.path.join(temp_dir, msg.get('FileName', 'image.jpg'))
msg.download(file_path)
if os.path.exists(file_path):
with open(file_path, 'rb') as f:
img_bytes = f.read()
b64 = base64.b64encode(img_bytes).decode('utf-8')
components.append(platform_message.Image(base64=f'data:image/jpeg;base64,{b64}'))
os.remove(file_path)
else:
components.append(platform_message.Unknown(text='[Image download failed]'))
except Exception:
components.append(platform_message.Unknown(text='[Image download failed]'))
elif msg_type == 'Recording':
try:
temp_dir = tempfile.gettempdir()
file_path = os.path.join(temp_dir, msg.get('FileName', 'voice.mp3'))
msg.download(file_path)
if os.path.exists(file_path):
with open(file_path, 'rb') as f:
voice_bytes = f.read()
b64 = base64.b64encode(voice_bytes).decode('utf-8')
components.append(platform_message.Voice(base64=b64))
os.remove(file_path)
else:
components.append(platform_message.Unknown(text='[Voice download failed]'))
except Exception:
components.append(platform_message.Unknown(text='[Voice download failed]'))
elif msg_type == 'Sharing':
text = msg.get('Text', '')
url = msg.get('Url', '')
content = text
if url and url not in text:
content = f'{text}\n{url}' if text else url
if content:
components.append(platform_message.Plain(text=content))
elif msg_type == 'Video':
components.append(platform_message.Unknown(text='[Video]'))
elif msg_type == 'Map':
components.append(platform_message.Unknown(text='[Location]'))
elif msg_type == 'Card':
components.append(platform_message.Unknown(text='[Contact Card]'))
elif msg_type == 'Note':
text = msg.get('Text', '')
if text:
components.append(platform_message.Unknown(text=f'[Note: {text}]'))
else:
text = msg.get('Text', '')
if text:
components.append(platform_message.Plain(text=text))
else:
components.append(platform_message.Unknown(text=f'[Unsupported message type: {msg_type}]'))
return platform_message.MessageChain(components)
class ItchatEventConverter(abstract_platform_adapter.AbstractEventConverter):
"""Converts itchat msg dicts to LangBot events."""
def __init__(self, adapter_ref: typing.Callable[[], typing.Any]):
"""adapter_ref is a callable returning the ItchatAdapter instance."""
self._get_adapter = adapter_ref
@staticmethod
async def yiri2target(event: platform_events.MessageEvent) -> dict:
return event.source_platform_object
def target2yiri(self, msg: dict) -> typing.Optional[platform_events.MessageEvent]:
"""Convert itchat msg to FriendMessage or GroupMessage."""
from_user = msg.get('FromUserName', '')
if not from_user:
return None
adapter = self._get_adapter()
bot_account_id = adapter.bot_account_id
bot_nickname = adapter._bot_nickname
message_chain = ItchatMessageConverter.target2yiri(msg)
if not message_chain:
return None
# Determine if this is a group message
# itchat uses '@@' prefix for chatroom IDs (not '@chatroom' suffix)
is_group = '@@' in from_user
timestamp = msg.get('CreateTime', 0)
if is_group:
# Actual sender within the group
actual_user = msg.get('ActualUserName', '')
actual_nick = msg.get('ActualNickName', '')
if not actual_nick:
actual_nick = actual_user
# Prepend @bot if the bot was mentioned
# itchat uses 'IsAt' (capital I, capital A) in produce_group_chat
if msg.get('IsAt', False):
# Strip @bot_nickname from the text content to avoid LLM confusion
if bot_nickname:
at_pattern = '@' + bot_nickname + ('' if '' in msg.get('Content', '') else ' ')
for component in message_chain:
if isinstance(component, platform_message.Plain):
if component.text.startswith(at_pattern):
component.text = component.text[len(at_pattern) :]
elif at_pattern in component.text:
component.text = component.text.replace(at_pattern, '')
break
message_chain = platform_message.MessageChain(
[platform_message.At(target=bot_account_id)] + list(message_chain)
)
# Try to get group display name
group_obj = msg.get('User', {})
group_name = ''
if hasattr(group_obj, 'NickName'):
group_name = group_obj.NickName
elif isinstance(group_obj, dict):
group_name = group_obj.get('NickName', '')
return platform_events.GroupMessage(
sender=platform_entities.GroupMember(
id=actual_user or actual_nick,
member_name=actual_nick or actual_user,
permission=platform_entities.Permission.Member,
group=platform_entities.Group(
id=from_user,
name=group_name or from_user,
permission=platform_entities.Permission.Member,
),
special_title='',
),
message_chain=message_chain,
time=timestamp,
source_platform_object=msg,
)
else:
# Private / friend message
sender_nick = ''
user_obj = msg.get('User', {})
if hasattr(user_obj, 'NickName'):
sender_nick = user_obj.NickName
elif isinstance(user_obj, dict):
sender_nick = user_obj.get('NickName', '')
return platform_events.FriendMessage(
sender=platform_entities.Friend(
id=from_user,
nickname=sender_nick or from_user,
remark='',
),
message_chain=message_chain,
time=timestamp,
source_platform_object=msg,
)
class ItchatAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
"""LangBot adapter for itchat-uos (WeChat Web)."""
name: str = 'itchat'
config: dict
logger: EventLogger
message_converter: ItchatMessageConverter
event_converter: ItchatEventConverter
listeners: typing.Dict[
typing.Type[platform_events.Event],
typing.Callable[[platform_events.Event, abstract_platform_adapter.AbstractMessagePlatformAdapter], None],
] = {}
_loop: typing.Optional[asyncio.AbstractEventLoop] = pydantic.PrivateAttr(default=None)
_logged_in: typing.Optional[threading.Event] = pydantic.PrivateAttr(default=None)
_itchat_thread: typing.Optional[threading.Thread] = pydantic.PrivateAttr(default=None)
_core: typing.Optional[ItchatCore] = pydantic.PrivateAttr(default=None)
_bot_nickname: str = pydantic.PrivateAttr(default='')
_bot_uuid: typing.Optional[str] = pydantic.PrivateAttr(default=None)
_startup_error: typing.Optional[str] = pydantic.PrivateAttr(default=None)
_connection_status: str = pydantic.PrivateAttr(default='disconnected')
_connection_error: str = pydantic.PrivateAttr(default='')
_last_connected_at: typing.Optional[float] = pydantic.PrivateAttr(default=None)
_last_disconnected_at: typing.Optional[float] = pydantic.PrivateAttr(default=None)
class Config:
arbitrary_types_allowed = True
def __init__(self, config: dict, logger: abstract_platform_logger.AbstractEventLogger):
message_converter = ItchatMessageConverter()
# Event converter needs a reference to self for bot_account_id + nickname
event_converter = ItchatEventConverter(adapter_ref=lambda: self)
super().__init__(
config=config,
logger=logger,
message_converter=message_converter,
event_converter=event_converter,
listeners={},
bot_account_id='',
)
# Initialize private attributes (can't be class-level defaults due to pickle)
self._loop = None
self._logged_in = threading.Event()
self._itchat_thread = None
self._core = ItchatCore()
self._startup_error = None
self._connection_status = 'disconnected'
self._connection_error = ''
self._last_connected_at = None
self._last_disconnected_at = None
@staticmethod
def _get_obj_value(obj: typing.Any, key: str, default: str = '') -> str:
if isinstance(obj, dict):
return obj.get(key, default) or default
return getattr(obj, key, default) or default
@staticmethod
def _safe_status_name(value: str) -> str:
cleaned = re.sub(r'[^A-Za-z0-9_.@-]+', '_', value.strip())
cleaned = cleaned.strip('._')
return cleaned
@staticmethod
def login_status_dir() -> str:
path = os.path.join('data', 'itchat')
os.makedirs(path, exist_ok=True)
return path
@classmethod
def login_status_path_for_account(cls, account_id: str) -> str:
safe_name = cls._safe_status_name(account_id)
if not safe_name:
raise ValueError('account_id is required for itchat login status')
filename = f'{safe_name}.pkl'
return os.path.join(cls.login_status_dir(), filename)
def _login_status_path(self) -> str:
configured_path = self.config.get('login_status_path', '').strip()
if configured_path:
return configured_path
account_id = self.config.get('account_id', '').strip()
if not account_id:
raise ValueError('account_id is required. Please scan the QR code and save this bot first.')
return self.login_status_path_for_account(account_id)
def set_bot_uuid(self, bot_uuid: str):
self._bot_uuid = bot_uuid
def _set_connection_status(self, status: str, error: str = ''):
self._connection_status = status
self._connection_error = error
now = time.time()
if status == 'connected':
self._last_connected_at = now
elif status in {'disconnected', 'error'}:
self._last_disconnected_at = now
def get_runtime_status(self) -> dict:
return {
'connection_status': self._connection_status,
'connection_error': self._connection_error,
'last_connected_at': self._last_connected_at,
'last_disconnected_at': self._last_disconnected_at,
}
def _on_login(self):
"""Called by itchat after successful QR code login."""
try:
# Refresh contacts
self._core.get_friends(update=True)
self._core.get_chatrooms(update=True)
# Get bot's own WeChat info from loginInfo['User']
user_info = self._core.loginInfo.get('User', {})
nick_name = self._get_obj_value(user_info, 'NickName')
user_name = self._get_obj_value(user_info, 'UserName')
# bot_account_id: config override or auto-detected wxid
# Used by AtBotRule for matching At.target
configured_id = self.config.get('account_id', '').strip()
self.bot_account_id = configured_id or user_name or nick_name or 'itchat-bot'
# _bot_nickname: config override or auto-detected nickname
configured_nick = self.config.get('nickname', '').strip()
self._bot_nickname = configured_nick or nick_name
self._set_connection_status('connected')
try:
chatrooms = self._core.search_chatrooms() or []
group_names = []
for c in chatrooms:
name = self._get_obj_value(c, 'NickName', str(c))
if name:
group_names.append(name)
if group_names:
self._log_sync(
f'itchat login as {nick_name} ({user_name}) | Groups ({len(group_names)}): {", ".join(group_names[:10])}{"..." if len(group_names) > 10 else ""}'
)
else:
self._log_sync(f'itchat login as {nick_name} ({user_name}) | No groups found')
except Exception as e:
self._log_sync(f'itchat login as {nick_name} ({user_name}) | Failed to list groups: {e}', 'warning')
except Exception as e:
self.bot_account_id = f'WeChat Bot (Error: {e})'
self._set_connection_status('error', str(e))
finally:
self._logged_in.set()
def _log_sync(self, msg: str, level: str = 'info'):
"""Thread-safe logging from itchat's sync thread."""
try:
if self._loop and not self._loop.is_closed():
log_fn = getattr(self.logger, level)
asyncio.run_coroutine_threadsafe(log_fn(msg), self._loop)
except Exception:
pass
def _drain_msglist(self):
"""Clear all stale messages from the msgList queue.
itchat's load_login_status fetches old messages via get_msg() and
pushes them into msgList. We drain them to avoid replaying history.
"""
try:
q = self._core.msgList
while True:
q.get_nowait()
except queue.Empty:
pass
def _on_qr_callback(self, **kwargs):
"""Called by itchat when QR code is generated or status changes.
Args:
uuid: QR code uuid
status: '200' = logged in, '201' = confirmed on phone, '408' = timeout
qrcode: raw bytes of the QR code PNG image
"""
status = kwargs.get('status', '')
qr_bytes = kwargs.get('qrcode', b'')
if status == '200':
# Login success, no need to show QR
return
# Only show QR on new QR generation (status='0') to avoid spamming
if not qr_bytes or status != '0':
return
try:
b64 = base64.b64encode(qr_bytes).decode('utf-8')
if self._loop and not self._loop.is_closed():
asyncio.run_coroutine_threadsafe(
self.logger.info(
'Please scan the QR code to login WeChat:',
images=[platform_message.Image(base64=f'data:image/png;base64,{b64}')],
),
self._loop,
)
except Exception:
pass
def _on_exit(self):
"""Called by itchat on exit."""
self._set_connection_status('disconnected')
self._log_sync('itchat session exited')
def _register_itchat_handlers(self):
"""Register itchat message decorators by re-registering handlers."""
@self._core.msg_register([TEXT])
def _on_text(msg):
self._dispatch_itchat_message(msg)
@self._core.msg_register([TEXT], isGroupChat=True)
def _on_group_text(msg):
self._dispatch_itchat_message(msg)
@self._core.msg_register([PICTURE])
def _on_picture(msg):
self._dispatch_itchat_message(msg)
@self._core.msg_register([PICTURE], isGroupChat=True)
def _on_group_picture(msg):
self._dispatch_itchat_message(msg)
@self._core.msg_register([RECORDING])
def _on_recording(msg):
self._dispatch_itchat_message(msg)
@self._core.msg_register([RECORDING], isGroupChat=True)
def _on_group_recording(msg):
self._dispatch_itchat_message(msg)
@self._core.msg_register([SHARING])
def _on_sharing(msg):
self._dispatch_itchat_message(msg)
@self._core.msg_register([SHARING], isGroupChat=True)
def _on_group_sharing(msg):
self._dispatch_itchat_message(msg)
@self._core.msg_register([VIDEO])
def _on_video(msg):
self._dispatch_itchat_message(msg)
@self._core.msg_register([VIDEO], isGroupChat=True)
def _on_group_video(msg):
self._dispatch_itchat_message(msg)
def _dispatch_itchat_message(self, msg: dict):
"""Bridge itchat callback (sync, in itchat thread) to async listener."""
try:
event = self.event_converter.target2yiri(msg)
if event is None:
return
event_type = type(event)
if event_type in self.listeners and self._loop:
callback = self.listeners[event_type]
asyncio.run_coroutine_threadsafe(
callback(event, self),
self._loop,
)
except Exception:
self._log_sync(f'Error dispatching itchat message: {traceback.format_exc()}', 'error')
async def send_message(
self,
target_type: str,
target_id: str,
message: platform_message.MessageChain,
):
"""Send a message to a user or group via itchat."""
items = await self.message_converter.yiri2target(message)
loop = asyncio.get_event_loop()
# Merge consecutive text items to avoid splitting messages
merged = []
for item in items:
if item['type'] == 'text' and merged and merged[-1]['type'] == 'text':
merged[-1]['content'] += item['content']
else:
merged.append(item)
for item in merged:
try:
if item['type'] == 'text':
await loop.run_in_executor(None, self._core.send, item['content'], target_id)
elif item['type'] == 'image':
# Save to temp file then send
temp_path = self._save_to_temp(item, 'image')
if temp_path:
await loop.run_in_executor(None, self._core.send, f'@img@{temp_path}', target_id)
self._cleanup_temp(temp_path)
elif item['type'] == 'voice':
temp_path = self._save_to_temp(item, 'voice')
if temp_path:
await loop.run_in_executor(None, self._core.send, f'@fil@{temp_path}', target_id)
self._cleanup_temp(temp_path)
elif item['type'] == 'file':
temp_path = self._save_to_temp(item, 'file')
if temp_path:
await loop.run_in_executor(None, self._core.send, f'@fil@{temp_path}', target_id)
self._cleanup_temp(temp_path)
except Exception:
await self.logger.error(f'Failed to send itchat message: {traceback.format_exc()}')
def _save_to_temp(self, item: dict, prefix: str) -> typing.Optional[str]:
"""Save base64 or URL data to a temp file and return the path."""
try:
if 'base64' in item:
b64_data = item['base64']
# Strip data URI prefix if present
if ',' in b64_data:
b64_data = b64_data.split(',', 1)[1]
file_bytes = base64.b64decode(b64_data)
suffix = '.jpg' if prefix == 'image' else ('.mp3' if prefix == 'voice' else '.bin')
fd, temp_path = tempfile.mkstemp(suffix=suffix, prefix=f'itchat_{prefix}_')
with os.fdopen(fd, 'wb') as f:
f.write(file_bytes)
return temp_path
elif 'url' in item:
import requests
resp = requests.get(item['url'], timeout=30)
if resp.status_code == 200:
suffix = '.jpg' if prefix == 'image' else ('.mp3' if prefix == 'voice' else '.bin')
fd, temp_path = tempfile.mkstemp(suffix=suffix, prefix=f'itchat_{prefix}_')
with os.fdopen(fd, 'wb') as f:
f.write(resp.content)
return temp_path
except Exception:
self._log_sync(f'Failed to save temp file: {traceback.format_exc()}', 'error')
return None
def _cleanup_temp(self, path: str):
"""Remove a temp file."""
try:
if os.path.exists(path):
os.remove(path)
except OSError:
pass
def _prepare_reply_message(
self,
message_source: platform_events.MessageEvent,
message: platform_message.MessageChain,
) -> platform_message.MessageChain:
"""Render group sender mentions with display names while keeping internal IDs stable."""
if not isinstance(message_source, platform_events.GroupMessage):
return message
source_msg = message_source.source_platform_object or {}
actual_user = source_msg.get('ActualUserName', '')
actual_nick = source_msg.get('ActualNickName', '')
if not actual_user or not actual_nick:
return message
components: list[platform_message.MessageComponent] = []
changed = False
for component in message:
if isinstance(component, platform_message.At) and str(component.target) == str(actual_user):
components.append(platform_message.Plain(text=f'@{actual_nick} '))
changed = True
else:
components.append(component)
if not changed:
return message
return platform_message.MessageChain(components)
async def reply_message(
self,
message_source: platform_events.MessageEvent,
message: platform_message.MessageChain,
quote_origin: bool = False,
):
"""Reply to a received message."""
source_msg = message_source.source_platform_object
if not source_msg:
return
# For group messages, reply to the group; for private, reply to the sender
from_user = source_msg.get('FromUserName', '')
if not from_user:
return
await self.send_message('friend', from_user, self._prepare_reply_message(message_source, message))
def register_listener(
self,
event_type: typing.Type[platform_events.Event],
callback: typing.Callable[
[platform_events.Event, abstract_platform_adapter.AbstractMessagePlatformAdapter], None
],
):
self.listeners[event_type] = callback
def unregister_listener(
self,
event_type: typing.Type[platform_events.Event],
callback: typing.Callable[
[platform_events.Event, abstract_platform_adapter.AbstractMessagePlatformAdapter], None
],
):
self.listeners.pop(event_type, None)
async def run_async(self):
"""Start the itchat adapter.
If an account-specific cached session file exists from a previous QR login,
itchat will reuse data/itchat/<account_id>.pkl without requiring a new QR scan.
"""
self._loop = asyncio.get_running_loop()
self._logged_in.clear()
self._startup_error = None
self._set_connection_status('connecting')
await self.logger.info('itchat adapter starting...')
# Register itchat message handlers BEFORE calling itchat.auto_login()
self._register_itchat_handlers()
# Run itchat in a daemon thread (it blocks)
def _run_itchat():
try:
status_path = self._login_status_path()
if not os.path.exists(status_path):
self._startup_error = (
f'No cached WeChat session found at {status_path}. '
'Please scan the QR code in the bot config page first.'
)
self._set_connection_status('error', self._startup_error)
self._log_sync(self._startup_error, 'error')
self._logged_in.set()
return
# Use hotReload to reuse the cached session from QR login
# If no cache exists, fail fast instead of triggering QR login
result = self._core.load_login_status(
status_path, loginCallback=self._on_login, exitCallback=self._on_exit
)
if result.get('BaseResponse', {}).get('Ret') != 0:
self._startup_error = (
f'Cached WeChat session at {status_path} is invalid. '
'Please scan the QR code in the bot config page again.'
)
self._set_connection_status('error', self._startup_error)
self._log_sync(self._startup_error, 'error')
self._logged_in.set()
return
# Session loaded, start message loop
self._log_sync(f'WeChat session loaded from cache: {status_path}')
# Clear stale messages that itchat fetched during hot-reload
self._drain_msglist()
self._core.run(blockThread=True)
self._set_connection_status('disconnected')
self._log_sync('itchat message loop stopped', 'error')
except Exception as e:
error = f'itchat run error: {e}'
self._set_connection_status('error', error)
self._log_sync(error, 'error')
self._logged_in.set()
self._itchat_thread = threading.Thread(target=_run_itchat, daemon=True, name='itchat-thread')
self._itchat_thread.start()
# Wait for login to complete (with timeout)
await asyncio.get_event_loop().run_in_executor(None, lambda: self._logged_in.wait(timeout=300))
if not self._logged_in.is_set():
raise RuntimeError('itchat login timed out (300s)')
if self._startup_error:
raise RuntimeError(self._startup_error)
await self.logger.info(f'itchat adapter running, bot: {self.bot_account_id}')
# Keep the adapter alive
try:
await asyncio.Event().wait()
except asyncio.CancelledError:
pass
async def kill(self) -> bool:
"""Stop the itchat adapter."""
try:
self._core.alive = False
self._core.isLogging = False
except Exception:
pass
self._set_connection_status('disconnected')
await self.logger.info('itchat adapter stopped')
return True
@@ -0,0 +1,75 @@
apiVersion: v1
kind: MessagePlatformAdapter
metadata:
name: itchat
label:
en_US: Itchat WeChat
zh_Hans: 个人微信 (itchat)
zh_Hant: 個人微信 (itchat)
ja_JP: 個人WeChat (itchat)
description:
en_US: Personal WeChat adapter via itchat-uos, supports QR code login and text/image/voice messages
zh_Hans: 基于 itchat-uos 的个人微信适配器,扫码登录,支持文本/图片/语音消息
zh_Hant: 基於 itchat-uos 的個人微信適配器,掃碼登入,支援文字/圖片/語音訊息
icon: wechat.png
spec:
categories:
- china
help_links:
zh: https://github.com/littlecodersh/ItChat
en: https://github.com/littlecodersh/ItChat
config:
- name: qr-login
label:
en_US: Scan QR Login
zh_Hans: 扫码登录
zh_Hant: 掃碼登入
description:
en_US: Scan QR code with WeChat to login. The session will be cached for the adapter to reuse.
zh_Hans: 使用微信扫码登录,登录状态将被缓存供适配器复用
zh_Hant: 使用微信掃碼登入,登入狀態將被快取供適配器復用
type: qr-code-login
login_platform: itchat
required: false
- name: account_id
label:
en_US: Bot Account ID
zh_Hans: 机器人账号标识
zh_Hant: 機器人帳號標識
ja_JP: ボットアカウントID
description:
en_US: Auto-filled after QR login with the WeChat wxid. Used for @-mention matching and to load data/itchat/<account_id>.pkl; do not change it to a nickname.
zh_Hans: 扫码登录后自动填入微信 wxid。用于群聊 @ 匹配,并加载 data/itchat/<account_id>.pkl;不要改成昵称。
zh_Hant: 掃碼登入後自動填入微信 wxid。用於群聊 @ 匹配,並載入 data/itchat/<account_id>.pkl;不要改成暱稱。
ja_JP: QRログイン後にWeChatのwxidが自動入力されます。@メンション判定と data/itchat/<account_id>.pkl の読み込みに使うため、ニックネームへ変更しないでください。
type: string
required: true
default: ""
- name: nickname
label:
en_US: Bot Nickname
zh_Hans: 机器人昵称
zh_Hant: 機器人暱稱
description:
en_US: The display nickname of the bot. Used to strip @nickname from incoming group messages. Auto-filled after QR login.
zh_Hans: 机器人的微信昵称。用于删除群聊消息中的 @昵称 前缀。扫码登录后自动填入
zh_Hant: 機器人的微信暱稱。用於刪除群聊訊息中的 @暱稱 前綴。掃碼登入後自動填入
type: string
required: false
default: ""
- name: hot_reload
label:
en_US: Hot Reload
zh_Hans: 登录缓存
zh_Hant: 登入快取
description:
en_US: Persist login session to avoid repeated QR code scans on restart
zh_Hans: 保存登录状态到本地,重启后无需重新扫码
zh_Hant: 儲存登入狀態到本機,重啟後無需重新掃碼
type: boolean
required: false
default: true
execution:
python:
path: ./itchat.py
attr: ItchatAdapter
+1 -23
View File
@@ -2,7 +2,6 @@ from __future__ import annotations
import typing
import asyncio
import traceback
import uuid
import datetime
import pydantic
@@ -183,28 +182,7 @@ class WecomCSAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
)
async def send_message(self, target_type: str, target_id: str, message: platform_message.MessageChain):
if target_type != 'person':
raise ValueError('WeCom customer service only supports sending messages to person targets')
open_kfid = self.bot_account_id
external_userid = target_id
if '|' in target_id:
open_kfid, external_userid = target_id.split('|', 1)
if external_userid.startswith('u'):
external_userid = external_userid[1:]
if not open_kfid:
raise ValueError('WeCom customer service open_kfid is required before sending messages')
content_list = await WecomMessageConverter.yiri2target(message, self.bot)
for content in content_list:
msgid = f'langbot_{uuid.uuid4().hex}'
if content['type'] == 'text':
await self.bot.send_text_msg(
open_kfid=open_kfid,
external_userid=external_userid,
msgid=msgid,
content=content['content'],
)
pass
def set_bot_uuid(self, bot_uuid: str):
"""设置 bot UUID(用于生成 webhook URL"""
+1 -1
View File
@@ -736,7 +736,7 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
) -> context.EventContext:
event_ctx = context.EventContext.from_event(event)
if not self.is_enable_plugin:
if not self.is_enable_plugin or not hasattr(self, 'handler'):
event_ctx._emitted_plugins = []
event_ctx._response_sources = []
return event_ctx
+8 -101
View File
@@ -25,7 +25,7 @@ from ....core import app
import langbot_plugin.api.entities.builtin.resource.tool as resource_tool
import langbot_plugin.api.entities.builtin.provider.message as provider_message
from ....entity.persistence import mcp as persistence_mcp
from .mcp_stdio import BoxStdioSessionRuntime, MCPServerBoxConfig, MCPSessionErrorPhase, _ColdStartRetry # noqa: F401
from .mcp_stdio import BoxStdioSessionRuntime, MCPServerBoxConfig, MCPSessionErrorPhase # noqa: F401
# Synthesized LLM tools for MCP resources (not from server tools/list).
# Dispatched in MCPLoader.invoke_tool; placeholder func on LLMTool is never used.
@@ -185,16 +185,6 @@ class MCPSessionStatus(enum.Enum):
ERROR = 'error'
class _TransportReconnect(Exception):
"""Internal signal: the Box stdio WS transport dropped but the managed
process is still alive. Triggers a lightweight transport reconnect that
reuses the live process, instead of a full process rebuild.
Reconnect attempts are NOT counted toward the fatal retry budget, so a
long-lived session can survive arbitrarily many transient drops.
"""
class RuntimeMCPSession:
"""运行时 MCP 会话"""
@@ -264,16 +254,6 @@ class RuntimeMCPSession:
self._lifecycle_task = None
self._shutdown_event = asyncio.Event()
self._ready_event = asyncio.Event()
# Set transiently when a WS transport drop should NOT stop the managed
# process (it will be re-attached on the next initialize()).
self._preserve_managed_process = False
# Log buffer for capturing stderr from Box managed process (maxlen=500 keeps
# recent lines without unbounded memory growth)
import collections as _collections
self._log_buffer: _collections.deque = _collections.deque(maxlen=500)
self._last_stderr_text: str = ''
self._box_stdio_runtime = BoxStdioSessionRuntime(self)
self.box_config = self._box_stdio_runtime.config
@@ -419,39 +399,11 @@ class RuntimeMCPSession:
task.cancel()
for task in done:
if task is monitor_task and not self._shutdown_event.is_set():
# The monitor completed. This is EITHER the managed
# process actually exiting OR just the WS transport
# dropping while the process stays alive in the Box
# runtime. Re-check the real process state so a
# transient transport drop reconnects (reusing the live
# process) instead of tearing the process down and
# running a full rebuild+backoff cycle.
process_still_running = False
try:
process_still_running = await self._box_stdio_runtime._managed_process_is_running()
except Exception:
process_still_running = False
if process_still_running:
self.ap.logger.info(
f'MCP server {self.server_name}: transport dropped but '
f'managed process is still running; reconnecting transport'
)
self.error_phase = MCPSessionErrorPhase.RELAY_CONNECT
# Preserve the live process across the finally-block
# cleanup: only the WS transport should be torn down.
self._preserve_managed_process = True
raise _TransportReconnect('Box managed process transport dropped; reconnecting')
self.error_phase = MCPSessionErrorPhase.RUNTIME
raise Exception('Box managed process exited unexpectedly')
else:
await self._shutdown_event.wait()
except _ColdStartRetry:
# Cold-start in progress: set the preserve flag BEFORE the finally
# block runs so it does not stop the live managed process. The outer
# _lifecycle_loop_with_retry will reuse it on the next attempt.
self._preserve_managed_process = True
raise
except Exception as e:
self.status = MCPSessionStatus.ERROR
self.error_message = str(e)
@@ -472,55 +424,14 @@ class RuntimeMCPSession:
except Exception as e:
self.ap.logger.error(f'Error cleaning up MCP session {self.server_name}: {e}\n{traceback.format_exc()}')
finally:
# On a transport-only reconnect the managed process is healthy
# and will be re-attached on the next initialize(); do NOT stop
# it. Any other exit path fully tears the session down.
if getattr(self, '_preserve_managed_process', False):
self._preserve_managed_process = False
else:
await self._cleanup_box_stdio_session()
await self._cleanup_box_stdio_session()
async def _lifecycle_loop_with_retry(self):
"""Wrap _lifecycle_loop with retry and exponential backoff."""
attempt = 0
while attempt <= self._MAX_RETRIES:
for attempt in range(self._MAX_RETRIES + 1):
try:
await self._lifecycle_loop()
return # Normal shutdown, don't retry
except _TransportReconnect as e:
# Transient WS transport drop while the managed process is still
# alive. Reconnect promptly WITHOUT consuming the fatal retry
# budget and WITHOUT stopping the process — initialize() will
# re-attach to the live process. This is what lets a long-lived
# stdio MCP survive repeated brief event-loop stalls / pings.
if self._shutdown_event.is_set():
return
self.ap.logger.info(
f'MCP session {self.server_name}: reconnecting transport ({self._describe_exception(e)})'
)
self.status = MCPSessionStatus.CONNECTING
self.error_message = None
self.error_phase = None
await asyncio.sleep(1)
continue
except _ColdStartRetry as e:
# The managed process is alive but still cold-starting (e.g.
# `npx -y <pkg>` is still installing) and cannot yet answer the
# handshake. Reuse the live process and retry the attach WITHOUT
# consuming the fatal retry budget or stopping the process, so a
# slow cold start is waited out instead of failing. Preserve the
# process across the finally-block cleanup.
if self._shutdown_event.is_set():
return
self._preserve_managed_process = True
self.ap.logger.debug(
f'MCP session {self.server_name}: waiting for cold start ({self._describe_exception(e)})'
)
self.status = MCPSessionStatus.CONNECTING
self.error_message = None
self.error_phase = None
await asyncio.sleep(2)
continue
except Exception as e:
self.retry_count = attempt + 1
if self._shutdown_event.is_set():
@@ -549,7 +460,6 @@ class RuntimeMCPSession:
self.error_message = None
self.error_phase = None
await asyncio.sleep(delay)
attempt += 1
@staticmethod
def _describe_exception(exc: BaseException) -> str:
@@ -1017,14 +927,11 @@ class RuntimeMCPSession:
return self._box_stdio_runtime.uses_box_stdio()
def _build_box_session_id(self) -> str:
# Both live servers and transient config-page tests share ONE Box
# session ('mcp-shared'). A test therefore reuses the already-running
# container (and, for an existing server, its live managed process)
# instead of paying a full per-test session cold-start + dependency
# bootstrap. Isolation between a test and the live servers is provided
# at the *process* level: each server/test has its own process_id and a
# test only ever stops its own process_id (see cleanup_session), so it
# never disturbs another server's process or the shared session itself.
# Transient test sessions get their own isolated Box session so a
# failing/short-lived test can never disturb the shared session that
# hosts live, already-connected MCP servers.
if self.is_transient:
return f'mcp-test-{self.server_uuid}'
return 'mcp-shared'
def _rewrite_path(self, path: str, host_path: str | None) -> str:
@@ -6,7 +6,7 @@ import os
import shutil
import shlex
import threading
from contextlib import suppress, AsyncExitStack
from contextlib import suppress
from typing import TYPE_CHECKING, Any
import pydantic
@@ -74,35 +74,6 @@ class MCPServerBoxConfig(pydantic.BaseModel):
model_config = pydantic.ConfigDict(extra='ignore')
_HANDSHAKE_ATTEMPT_TIMEOUT_SEC = 10.0
class _TransferredStack:
"""Adapts an already-populated AsyncExitStack into an async context manager
so ownership of its resources can be transferred into another exit stack.
Entering is a no-op; exiting closes the wrapped stack (and thus the live WS
transport + ClientSession) when the owning session shuts down."""
def __init__(self, stack: AsyncExitStack):
self._stack = stack
async def __aenter__(self):
return self
async def __aexit__(self, exc_type, exc, tb):
await self._stack.aclose()
return False
class _ColdStartRetry(Exception):
"""Signal: the managed process is alive but not yet answering the MCP
handshake because it is still cold-starting (e.g. `npx -y <pkg>` is still
installing). The outer lifecycle retry treats this like a transient
reconnect: it reuses the live process and does not count toward the fatal
retry budget, so a slow cold start is waited out rather than failing.
"""
class BoxStdioSessionRuntime:
"""Encapsulate Box-backed stdio MCP session orchestration."""
@@ -142,15 +113,7 @@ class BoxStdioSessionRuntime:
read_only_rootfs=self.config.read_only_rootfs if self.config.read_only_rootfs is not None else False,
image=self.config.image,
cpus=self.config.cpus,
# Node.js runtimes (npx/bunx) reserve large virtual address space and
# load WebAssembly modules (llhttp) on startup; the default 512 MB
# cgroup_mem_max is too small and causes OOM kills (return_code=137).
# Auto-bump to 1024 MB when the runner is npx/bunx/pnpm dlx.
memory_mb=(
(self.config.memory_mb or 1024)
if self.server_config.get('command', '') in ('npx', 'bunx', 'pnpm')
else self.config.memory_mb
),
memory_mb=self.config.memory_mb,
pids_limit=self.config.pids_limit,
persistent=True,
)
@@ -210,55 +173,28 @@ class BoxStdioSessionRuntime:
stderr_preview = (result.stderr or '')[:500]
raise Exception(f'Dependency install failed (exit code {result.exit_code}): {stderr_preview}')
# Reuse an already-running managed process instead of rebuilding it.
# The Box runtime keeps the managed process alive across a transient
# WebSocket transport drop, so on a reconnect we only need to re-attach
# the WS below. Rebuilding here would needlessly stop a healthy process
# and re-run the (slow, network-touching) dependency bootstrap.
if not await self._managed_process_is_running():
try:
process_workspace = (
self._build_workspace(host_path=host_path, workdir=process_cwd, mount_path=process_cwd)
if host_path
else workspace
)
payload = process_workspace.build_process_payload(
self.server_config['command'],
self.server_config.get('args', []),
env=self.server_config.get('env', {}),
cwd=process_cwd,
)
if install_cmd:
payload = self._wrap_process_payload_with_python_env(payload, process_cwd)
payload['process_id'] = self.process_id
await workspace.box_service.start_managed_process(workspace.session_id, payload)
except Exception:
self.owner.error_phase = MCPSessionErrorPhase.PROCESS_START
raise
else:
self.ap.logger.info(
f'MCP server {self.server_name}: reusing live managed process '
f'process_id={self.process_id} (transport reconnect)'
)
websocket_url = workspace.get_managed_process_websocket_url(self.process_id)
# Attach the WS transport + MCP session ONCE, on the owner's exit stack,
# in the same task as the serve loop that follows. websocket_client and
# ClientSession use anyio task groups whose cancel scope is bound to the
# frame/stack that entered them, so they must live on the owner exit
# stack (not a deferred/transferred one) or the streams close the moment
# initialize() returns and the next request fails with "Connection
# closed".
#
# A slow (`npx -y <pkg>`) cold start makes this single attempt fail
# while the process is still alive — the package is still installing and
# cannot answer the handshake. We surface that to the outer retry loop
# as a _ColdStartRetry: it must NOT stop the process (it is healthy and
# will be reused) and must NOT consume the fatal retry budget. The next
# attempt re-attaches to the same live process; once it has finished
# cold start the handshake succeeds and stays healthy.
try:
process_workspace = (
self._build_workspace(host_path=host_path, workdir=process_cwd, mount_path=process_cwd)
if host_path
else workspace
)
payload = process_workspace.build_process_payload(
self.server_config['command'],
self.server_config.get('args', []),
env=self.server_config.get('env', {}),
cwd=process_cwd,
)
if install_cmd:
payload = self._wrap_process_payload_with_python_env(payload, process_cwd)
payload['process_id'] = self.process_id
await workspace.box_service.start_managed_process(workspace.session_id, payload)
except Exception:
self.owner.error_phase = MCPSessionErrorPhase.PROCESS_START
raise
try:
websocket_url = workspace.get_managed_process_websocket_url(self.process_id)
transport = await self.owner.exit_stack.enter_async_context(websocket_client(websocket_url))
read_stream, write_stream = transport
self.owner.session = await self.owner.exit_stack.enter_async_context(
@@ -266,19 +202,12 @@ class BoxStdioSessionRuntime:
)
except Exception:
self.owner.error_phase = MCPSessionErrorPhase.RELAY_CONNECT
if not await self._managed_process_has_exited():
# Process is alive but not yet serving (cold start) — reconnect.
raise _ColdStartRetry(f'{self.server_name}: transport not ready during cold start')
raise
try:
await asyncio.wait_for(self.owner.session.initialize(), timeout=_HANDSHAKE_ATTEMPT_TIMEOUT_SEC)
except Exception as exc:
await self.owner.session.initialize()
except Exception:
self.owner.error_phase = MCPSessionErrorPhase.MCP_INIT
if not await self._managed_process_has_exited():
raise _ColdStartRetry(
f'{self.server_name}: handshake not ready during cold start ({type(exc).__name__})'
)
raise
async def monitor_process_health(self) -> None:
@@ -305,74 +234,8 @@ class BoxStdioSessionRuntime:
)
if consecutive_errors >= self.owner._MONITOR_MAX_CONSECUTIVE_ERRORS:
return
# Capture stderr logs from the managed process
if isinstance(info, dict):
stderr_text = info.get('stderr', '') or info.get('stderr_preview', '')
else:
stderr_text = getattr(info, 'stderr', '') or getattr(info, 'stderr_preview', '')
if stderr_text and stderr_text != self.owner._last_stderr_text:
# Find new lines not in the previous snapshot
old_lines = set(self.owner._last_stderr_text.splitlines()) if self.owner._last_stderr_text else set()
new_lines = [l for l in stderr_text.splitlines() if l and l not in old_lines]
self.owner._last_stderr_text = stderr_text
import time as _time
for line in new_lines:
level = (
'error'
if any(k in line.upper() for k in ('ERROR', 'CRITICAL'))
else 'warning'
if 'WARNING' in line.upper()
else 'debug'
if 'DEBUG' in line.upper()
else 'info'
)
self.owner._log_buffer.append({'ts': _time.time(), 'level': level, 'text': line})
await asyncio.sleep(self.owner._MONITOR_POLL_INTERVAL)
async def _managed_process_is_running(self) -> bool:
"""Return True if this server's managed process exists and is running.
Used to decide whether initialize() must (re)start the process or can
simply re-attach the WebSocket transport to a process the Box runtime
kept alive across a transient transport drop.
"""
from langbot_plugin.box.models import BoxManagedProcessStatus
workspace = self._build_workspace()
try:
info = await workspace.get_managed_process(self.process_id)
except Exception:
return False
status = info.get('status', '') if isinstance(info, dict) else getattr(info, 'status', '')
return status in (BoxManagedProcessStatus.RUNNING.value, BoxManagedProcessStatus.RUNNING)
async def _managed_process_has_exited(self) -> bool:
"""Return True only if the process is DEFINITIVELY gone (reports EXITED).
Distinct from ``not _managed_process_is_running()``: a process that has
just been spawned may not yet report RUNNING, and a transient query
error is not proof of exit. During the cold-start handshake retry we
must NOT treat 'not yet running' or 'query failed' as a terminal
failure, or we bail out to the outer rebuild path and churn the
process (relay then rejects the early re-attach with HTTP 400). Only a
successful query that reports EXITED stops the retry loop.
"""
from langbot_plugin.box.models import BoxManagedProcessStatus
workspace = self._build_workspace()
try:
info = await workspace.get_managed_process(self.process_id)
except Exception:
# Unknown — treat as 'still coming up', not exited.
return False
status = info.get('status', '') if isinstance(info, dict) else getattr(info, 'status', '')
return status in (BoxManagedProcessStatus.EXITED.value, BoxManagedProcessStatus.EXITED)
async def _stage_host_path_to_shared_workspace(self, host_path: str) -> str:
source_path = normalize_host_path(host_path)
if not source_path:
@@ -479,20 +342,16 @@ class BoxStdioSessionRuntime:
workspace = self._build_workspace(host_path=None)
# Transient config-page tests now share the same 'mcp-shared' Box
# session as live servers, so we must NOT tear the session down here —
# that would kill every other MCP server in the container. A test is
# isolated at the process level: it ran under its own process_id, so we
# stop only that process, exactly like a live server does below. The
# shared session and all other servers' live processes are untouched.
# (Staged per-test workspace files are still cleaned up.)
# Transient test sessions own their isolated Box session, so tear the
# whole session down rather than leaking it. This cannot affect live
# servers because they live in the separate shared session.
if getattr(self.owner, 'is_transient', False):
try:
await workspace.stop_managed_process(self.process_id)
await workspace.cleanup()
except Exception as exc:
self.ap.logger.warning(
f'MCP server {self.server_name}: failed to stop transient test process '
f'process_id={self.process_id}: {type(exc).__name__}: {exc}'
f'MCP server {self.server_name}: failed to delete transient test session '
f'{self.owner._build_box_session_id()}: {type(exc).__name__}: {exc}'
)
await self._cleanup_staged_workspace()
return
+4 -114
View File
@@ -1,7 +1,6 @@
from __future__ import annotations
import typing
import time
from typing import TYPE_CHECKING
import langbot_plugin.api.entities.builtin.resource.tool as resource_tool
@@ -143,130 +142,21 @@ class ToolManager:
return tools
def _get_query_session_id(self, query: pipeline_query.Query) -> str | None:
launcher_type = getattr(query, 'launcher_type', None)
launcher_id = getattr(query, 'launcher_id', None)
if launcher_type is None or launcher_id is None:
return None
launcher_type_value = launcher_type.value if hasattr(launcher_type, 'value') else launcher_type
return f'{launcher_type_value}_{launcher_id}'
async def _record_tool_call(
self,
*,
name: str,
source: str,
parameters: dict,
query: pipeline_query.Query,
duration_ms: int,
status: str,
result: typing.Any = None,
error_message: str | None = None,
) -> None:
monitoring_service = getattr(self.ap, 'monitoring_service', None)
if not monitoring_service:
return
variables = getattr(query, 'variables', {}) or {}
message_id = variables.get('_monitoring_message_id') if isinstance(variables, dict) else None
bot_name = variables.get('_monitoring_bot_name') if isinstance(variables, dict) else None
pipeline_name = variables.get('_monitoring_pipeline_name') if isinstance(variables, dict) else None
try:
await monitoring_service.record_tool_call(
tool_name=name,
tool_source=source,
duration=duration_ms,
status=status,
bot_id=getattr(query, 'bot_uuid', None),
bot_name=bot_name,
pipeline_name=pipeline_name,
session_id=self._get_query_session_id(query),
message_id=message_id,
arguments=parameters,
result=result,
error_message=error_message,
)
except Exception as e:
self.ap.logger.warning(f'Failed to record tool call: {e}')
async def _invoke_tool_with_monitoring(
self,
*,
source: str,
name: str,
parameters: dict,
query: pipeline_query.Query,
invoke: typing.Callable[[], typing.Awaitable[typing.Any]],
) -> typing.Any:
start_time = time.perf_counter()
try:
result = await invoke()
except Exception as e:
duration_ms = int((time.perf_counter() - start_time) * 1000)
await self._record_tool_call(
name=name,
source=source,
parameters=parameters,
query=query,
duration_ms=duration_ms,
status='error',
error_message=str(e),
)
raise
duration_ms = int((time.perf_counter() - start_time) * 1000)
await self._record_tool_call(
name=name,
source=source,
parameters=parameters,
query=query,
duration_ms=duration_ms,
status='success',
result=result,
)
return result
async def execute_func_call(self, name: str, parameters: dict, query: pipeline_query.Query) -> typing.Any:
from langbot.pkg.telemetry import features as telemetry_features
if await self.native_tool_loader.has_tool(name):
telemetry_features.increment(query, 'tool_calls', 'native')
return await self._invoke_tool_with_monitoring(
source='native',
name=name,
parameters=parameters,
query=query,
invoke=lambda: self.native_tool_loader.invoke_tool(name, parameters, query),
)
return await self.native_tool_loader.invoke_tool(name, parameters, query)
if await self.plugin_tool_loader.has_tool(name):
telemetry_features.increment(query, 'tool_calls', 'plugin')
return await self._invoke_tool_with_monitoring(
source='plugin',
name=name,
parameters=parameters,
query=query,
invoke=lambda: self.plugin_tool_loader.invoke_tool(name, parameters, query),
)
return await self.plugin_tool_loader.invoke_tool(name, parameters, query)
if await self.mcp_tool_loader.has_tool(name):
telemetry_features.increment(query, 'tool_calls', 'mcp')
return await self._invoke_tool_with_monitoring(
source='mcp',
name=name,
parameters=parameters,
query=query,
invoke=lambda: self.mcp_tool_loader.invoke_tool(name, parameters, query),
)
return await self.mcp_tool_loader.invoke_tool(name, parameters, query)
if await self.skill_tool_loader.has_tool(name):
telemetry_features.increment(query, 'tool_calls', 'skill')
return await self._invoke_tool_with_monitoring(
source='skill',
name=name,
parameters=parameters,
query=query,
invoke=lambda: self.skill_tool_loader.invoke_tool(name, parameters, query),
)
return await self.skill_tool_loader.invoke_tool(name, parameters, query)
raise ToolNotFoundError(name)
async def shutdown(self):
-1
View File
@@ -81,7 +81,6 @@ def fake_monitoring_app():
)
app.monitoring_service.get_messages = AsyncMock(return_value=([{'id': 'msg-1', 'content': 'test'}], 100))
app.monitoring_service.get_llm_calls = AsyncMock(return_value=([{'id': 'llm-1'}], 50))
app.monitoring_service.get_tool_calls = AsyncMock(return_value=([{'id': 'tool-1'}], 5))
app.monitoring_service.get_embedding_calls = AsyncMock(return_value=([{'id': 'emb-1'}], 10))
app.monitoring_service.get_sessions = AsyncMock(return_value=([{'session_id': 'sess-1'}], 20))
app.monitoring_service.get_errors = AsyncMock(return_value=([{'id': 'err-1'}], 2))
@@ -280,25 +280,6 @@ class TestMCPServiceCreateMCPServer:
assert server_uuid is not None
assert len(server_uuid) == 36 # UUID format
async def test_create_mcp_server_duplicate_name_raises(self):
"""Rejects duplicate MCP server names."""
# Setup
ap = SimpleNamespace()
ap.persistence_mgr = SimpleNamespace()
ap.instance_config = SimpleNamespace()
ap.instance_config.data = {'system': {'limitation': {'max_extensions': -1}}}
ap.tool_mgr = None
existing_server = _create_mock_mcp_server(name='Existing Server')
ap.persistence_mgr.execute_async = AsyncMock(return_value=_create_mock_result(first_item=existing_server))
ap.persistence_mgr.serialize_model = Mock(return_value={})
service = MCPService(ap)
# Execute & Verify
with pytest.raises(ValueError, match='MCP server already exists: Existing Server'):
await service.create_mcp_server({'name': 'Existing Server'})
async def test_create_mcp_server_loads_server(self):
"""Loads server into tool_mgr when enabled."""
# Setup
@@ -320,7 +301,7 @@ class TestMCPServiceCreateMCPServer:
nonlocal call_count
call_count += 1
if call_count == 1:
return _create_mock_result([]) # Empty result for duplicate-name check
return _create_mock_result([]) # Empty list for limit check
elif call_count == 2:
return Mock() # Insert
return _create_mock_result(first_item=server_entity) # Select created
@@ -1,76 +0,0 @@
from __future__ import annotations
import sys
import types
from importlib import import_module
from types import SimpleNamespace
from unittest.mock import AsyncMock
import pytest
import quart
core_app_module = types.ModuleType('langbot.pkg.core.app')
core_app_module.Application = object
sys.modules.setdefault('langbot.pkg.core.app', core_app_module)
pytestmark = pytest.mark.asyncio
async def _create_test_client(mcp_service: SimpleNamespace):
app = quart.Quart(__name__)
user_service = SimpleNamespace(
verify_jwt_token=AsyncMock(return_value='test@example.com'),
get_user_by_email=AsyncMock(return_value=SimpleNamespace(user='test@example.com')),
)
ap = SimpleNamespace(mcp_service=mcp_service, user_service=user_service)
MCPRouterGroup = import_module('langbot.pkg.api.http.controller.groups.resources.mcp').MCPRouterGroup
group = MCPRouterGroup(ap, app)
await group.initialize()
return app.test_client()
async def test_mcp_server_route_accepts_encoded_slash_name():
mcp_service = SimpleNamespace(
get_mcp_server_by_name=AsyncMock(
return_value={
'uuid': 'test-uuid',
'name': 'pab1it0/prometheus',
'enable': True,
'mode': 'stdio',
'extra_args': {},
}
)
)
client = await _create_test_client(mcp_service)
response = await client.get(
'/api/v1/mcp/servers/pab1it0%2Fprometheus',
headers={'Authorization': 'Bearer test-token'},
)
assert response.status_code == 200
mcp_service.get_mcp_server_by_name.assert_awaited_once_with('pab1it0/prometheus')
payload = await response.get_json()
assert payload['data']['server']['name'] == 'pab1it0/prometheus'
async def test_mcp_resource_route_accepts_encoded_slash_name():
mcp_service = SimpleNamespace(
get_mcp_server_by_name=AsyncMock(),
get_mcp_server_resources=AsyncMock(return_value=[]),
get_mcp_server_resource_templates=AsyncMock(return_value=[]),
get_runtime_info=AsyncMock(return_value={'resource_capabilities': {'subscribe': False}}),
)
client = await _create_test_client(mcp_service)
response = await client.get(
'/api/v1/mcp/servers/pab1it0%2Fprometheus/resources',
headers={'Authorization': 'Bearer test-token'},
)
assert response.status_code == 200
mcp_service.get_mcp_server_by_name.assert_not_awaited()
mcp_service.get_mcp_server_resources.assert_awaited_once_with('pab1it0/prometheus')
payload = await response.get_json()
assert payload['data']['resource_capabilities'] == {'subscribe': False}
@@ -1,12 +1,7 @@
import pytest
import aiocqhttp
import langbot_plugin.api.entities.builtin.platform.message as platform_message
from langbot.pkg.platform.sources.aiocqhttp import (
AiocqhttpAdapter,
AiocqhttpEventConverter,
AiocqhttpMessageConverter,
)
from langbot.pkg.platform.sources.aiocqhttp import AiocqhttpAdapter, AiocqhttpMessageConverter
async def _convert_single(component: platform_message.MessageComponent):
@@ -108,431 +103,3 @@ async def test_forward_image_base64_payload_is_normalized():
'type': 'image',
'data': {'file': 'base64://raw-forward-image'},
}
@pytest.mark.asyncio
async def test_group_message_member_name_prefers_group_card():
event = aiocqhttp.Event(
{
'post_type': 'message',
'message_type': 'group',
'message_id': 1000,
'message': '',
'time': 1776491725,
'group_id': 2000,
'sender': {
'user_id': 3000,
'nickname': 'QQ Nickname',
'card': 'Group Card',
'role': 'member',
'title': 'Special Title',
},
}
)
class Bot:
async def get_group_info(self, group_id):
assert group_id == 2000
return {'group_id': group_id, 'group_name': 'Test Group'}
converted = await AiocqhttpEventConverter().target2yiri(event, Bot())
assert converted.sender.member_name == 'Group Card'
assert converted.sender.group.id == 2000
assert converted.sender.group.name == 'Test Group'
assert converted.sender.special_title == 'Special Title'
@pytest.mark.asyncio
async def test_group_message_member_name_falls_back_to_nickname():
event = aiocqhttp.Event(
{
'post_type': 'message',
'message_type': 'group',
'message_id': 1000,
'message': '',
'time': 1776491725,
'group_id': 2000,
'sender': {
'user_id': 3000,
'nickname': 'QQ Nickname',
'card': '',
'role': 'member',
},
}
)
converted = await AiocqhttpEventConverter().target2yiri(event)
assert converted.sender.member_name == 'QQ Nickname'
@pytest.mark.asyncio
async def test_group_message_special_title_uses_group_member_info_when_sender_title_is_empty():
event = aiocqhttp.Event(
{
'post_type': 'message',
'message_type': 'group',
'message_id': 1000,
'message': '',
'time': 1776491725,
'group_id': 2000,
'sender': {
'user_id': 3000,
'nickname': 'QQ Nickname',
'card': 'Group Card',
'role': 'member',
'title': '',
},
}
)
class Bot:
async def get_group_info(self, group_id):
return {'group_id': group_id, 'group_name': 'Test Group'}
async def get_group_member_info(self, group_id, user_id):
assert group_id == 2000
assert user_id == 3000
return {'group_id': group_id, 'user_id': user_id, 'title': 'Member Title'}
converted = await AiocqhttpEventConverter().target2yiri(event, Bot())
assert converted.sender.special_title == 'Member Title'
@pytest.mark.asyncio
async def test_group_message_special_title_does_not_lookup_when_sender_title_exists():
event = aiocqhttp.Event(
{
'post_type': 'message',
'message_type': 'group',
'message_id': 1000,
'message': '',
'time': 1776491725,
'group_id': 2000,
'sender': {
'user_id': 3000,
'nickname': 'QQ Nickname',
'card': 'Group Card',
'role': 'member',
'title': 'Event Title',
},
}
)
class Bot:
async def get_group_info(self, group_id):
return {'group_id': group_id, 'group_name': 'Test Group'}
async def get_group_member_info(self, group_id, user_id):
raise AssertionError('get_group_member_info should not be called')
converted = await AiocqhttpEventConverter().target2yiri(event, Bot())
assert converted.sender.special_title == 'Event Title'
@pytest.mark.asyncio
async def test_group_message_special_title_member_info_failure_is_cached(monkeypatch):
event = aiocqhttp.Event(
{
'post_type': 'message',
'message_type': 'group',
'message_id': 1000,
'message': '',
'time': 1776491725,
'group_id': 2000,
'sender': {
'user_id': 3000,
'nickname': 'QQ Nickname',
'card': 'Group Card',
'role': 'member',
'title': '',
},
}
)
now = 1000.0
class Bot:
member_info_calls = 0
async def get_group_info(self, group_id):
return {'group_id': group_id, 'group_name': 'Test Group'}
async def get_group_member_info(self, group_id, user_id):
self.member_info_calls += 1
raise RuntimeError('api unavailable')
monkeypatch.setattr('langbot.pkg.platform.sources.aiocqhttp.time.monotonic', lambda: now)
bot = Bot()
converter = AiocqhttpEventConverter()
first = await converter.target2yiri(event, bot)
second = await converter.target2yiri(event, bot)
assert first.sender.special_title == ''
assert second.sender.special_title == ''
assert bot.member_info_calls == 1
@pytest.mark.asyncio
async def test_group_message_special_title_member_info_cache_expires(monkeypatch):
event = aiocqhttp.Event(
{
'post_type': 'message',
'message_type': 'group',
'message_id': 1000,
'message': '',
'time': 1776491725,
'group_id': 2000,
'sender': {
'user_id': 3000,
'nickname': 'QQ Nickname',
'card': 'Group Card',
'role': 'member',
'title': '',
},
}
)
now = 1000.0
class Bot:
member_info_calls = 0
async def get_group_info(self, group_id):
return {'group_id': group_id, 'group_name': 'Test Group'}
async def get_group_member_info(self, group_id, user_id):
self.member_info_calls += 1
return {
'group_id': group_id,
'user_id': user_id,
'title': f'Member Title {self.member_info_calls}',
}
monkeypatch.setattr('langbot.pkg.platform.sources.aiocqhttp.time.monotonic', lambda: now)
bot = Bot()
converter = AiocqhttpEventConverter()
first = await converter.target2yiri(event, bot)
now = 87401.0
second = await converter.target2yiri(event, bot)
assert first.sender.special_title == 'Member Title 1'
assert second.sender.special_title == 'Member Title 2'
assert bot.member_info_calls == 2
@pytest.mark.asyncio
async def test_group_message_special_title_retries_after_negative_cache_expires(monkeypatch):
event = aiocqhttp.Event(
{
'post_type': 'message',
'message_type': 'group',
'message_id': 1000,
'message': '',
'time': 1776491725,
'group_id': 2000,
'sender': {
'user_id': 3000,
'nickname': 'QQ Nickname',
'card': 'Group Card',
'role': 'member',
'title': '',
},
}
)
now = 1000.0
class Bot:
member_info_calls = 0
async def get_group_info(self, group_id):
return {'group_id': group_id, 'group_name': 'Test Group'}
async def get_group_member_info(self, group_id, user_id):
self.member_info_calls += 1
if self.member_info_calls == 1:
raise RuntimeError('api unavailable')
return {'group_id': group_id, 'user_id': user_id, 'title': 'Recovered Title'}
monkeypatch.setattr('langbot.pkg.platform.sources.aiocqhttp.time.monotonic', lambda: now)
bot = Bot()
converter = AiocqhttpEventConverter()
failed = await converter.target2yiri(event, bot)
now = 1601.0
recovered = await converter.target2yiri(event, bot)
assert failed.sender.special_title == ''
assert recovered.sender.special_title == 'Recovered Title'
assert bot.member_info_calls == 2
@pytest.mark.asyncio
async def test_group_message_group_name_is_cached(monkeypatch):
event = aiocqhttp.Event(
{
'post_type': 'message',
'message_type': 'group',
'message_id': 1000,
'message': '',
'time': 1776491725,
'group_id': 2000,
'sender': {
'user_id': 3000,
'nickname': 'QQ Nickname',
'card': 'Group Card',
'role': 'member',
},
}
)
class Bot:
calls = 0
async def get_group_info(self, group_id):
self.calls += 1
assert group_id == 2000
return {'group_id': group_id, 'group_name': 'Cached Group'}
monotonic = 1000.0
monkeypatch.setattr('langbot.pkg.platform.sources.aiocqhttp.time.monotonic', lambda: monotonic)
bot = Bot()
converter = AiocqhttpEventConverter()
first = await converter.target2yiri(event, bot)
second = await converter.target2yiri(event, bot)
assert first.sender.group.name == 'Cached Group'
assert second.sender.group.name == 'Cached Group'
assert bot.calls == 1
@pytest.mark.asyncio
async def test_group_message_group_name_cache_expires(monkeypatch):
event = aiocqhttp.Event(
{
'post_type': 'message',
'message_type': 'group',
'message_id': 1000,
'message': '',
'time': 1776491725,
'group_id': 2000,
'sender': {
'user_id': 3000,
'nickname': 'QQ Nickname',
'card': 'Group Card',
'role': 'member',
},
}
)
now = 1000.0
class Bot:
calls = 0
async def get_group_info(self, group_id):
self.calls += 1
return {'group_id': group_id, 'group_name': f'Group Name {self.calls}'}
monkeypatch.setattr('langbot.pkg.platform.sources.aiocqhttp.time.monotonic', lambda: now)
bot = Bot()
converter = AiocqhttpEventConverter()
first = await converter.target2yiri(event, bot)
now = 4601.0
second = await converter.target2yiri(event, bot)
assert first.sender.group.name == 'Group Name 1'
assert second.sender.group.name == 'Group Name 2'
assert bot.calls == 2
@pytest.mark.asyncio
async def test_group_message_group_name_uses_placeholder_when_lookup_fails(monkeypatch):
event = aiocqhttp.Event(
{
'post_type': 'message',
'message_type': 'group',
'message_id': 1000,
'message': '',
'time': 1776491725,
'group_id': 2000,
'sender': {
'user_id': 3000,
'nickname': 'QQ Nickname',
'card': 'Group Card',
'role': 'member',
},
}
)
now = 1000.0
class Bot:
calls = 0
async def get_group_info(self, group_id):
self.calls += 1
raise RuntimeError('api unavailable')
monkeypatch.setattr('langbot.pkg.platform.sources.aiocqhttp.time.monotonic', lambda: now)
bot = Bot()
converter = AiocqhttpEventConverter()
converted = await converter.target2yiri(event, bot)
cached_failure = await converter.target2yiri(event, bot)
assert converted.sender.group.name == 'Group 2000'
assert cached_failure.sender.group.name == 'Group 2000'
assert bot.calls == 1
@pytest.mark.asyncio
async def test_group_message_group_name_retries_after_negative_cache_expires(monkeypatch):
event = aiocqhttp.Event(
{
'post_type': 'message',
'message_type': 'group',
'message_id': 1000,
'message': '',
'time': 1776491725,
'group_id': 2000,
'sender': {
'user_id': 3000,
'nickname': 'QQ Nickname',
'card': 'Group Card',
'role': 'member',
},
}
)
now = 1000.0
class Bot:
calls = 0
async def get_group_info(self, group_id):
self.calls += 1
if self.calls == 1:
raise RuntimeError('api unavailable')
return {'group_id': group_id, 'group_name': 'Recovered Group'}
monkeypatch.setattr('langbot.pkg.platform.sources.aiocqhttp.time.monotonic', lambda: now)
bot = Bot()
converter = AiocqhttpEventConverter()
failed = await converter.target2yiri(event, bot)
now = 1061.0
recovered = await converter.target2yiri(event, bot)
assert failed.sender.group.name == 'Group 2000'
assert recovered.sender.group.name == 'Recovered Group'
assert bot.calls == 2
@@ -1,91 +0,0 @@
from types import SimpleNamespace
from unittest.mock import AsyncMock
import pytest
import langbot_plugin.api.definition.abstract.platform.event_logger as abstract_platform_logger
import langbot_plugin.api.entities.builtin.platform.message as platform_message
from langbot.pkg.platform.sources.wecomcs import WecomCSAdapter
class DummyLogger(abstract_platform_logger.AbstractEventLogger):
async def info(self, *args, **kwargs):
pass
async def debug(self, *args, **kwargs):
pass
async def warning(self, *args, **kwargs):
pass
async def error(self, *args, **kwargs):
pass
def make_adapter():
return WecomCSAdapter(
config={
'corpid': 'corp-id',
'secret': 'secret',
'token': 'token',
'EncodingAESKey': 'encoding-key',
},
logger=DummyLogger(),
)
@pytest.mark.asyncio
async def test_send_message_sends_text_to_customer_service_user():
adapter = make_adapter()
adapter.bot_account_id = 'kf-test'
adapter.bot = SimpleNamespace(send_text_msg=AsyncMock())
message = platform_message.MessageChain([platform_message.Plain(text='hello')])
await adapter.send_message('person', 'uexternal-user', message)
adapter.bot.send_text_msg.assert_awaited_once()
kwargs = adapter.bot.send_text_msg.await_args.kwargs
assert kwargs['open_kfid'] == 'kf-test'
assert kwargs['external_userid'] == 'external-user'
assert kwargs['content'] == 'hello'
assert kwargs['msgid'].startswith('langbot_')
@pytest.mark.asyncio
async def test_send_message_allows_explicit_open_kfid_in_target_id():
adapter = make_adapter()
adapter.bot = SimpleNamespace(send_text_msg=AsyncMock())
message = platform_message.MessageChain([platform_message.Plain(text='hello')])
await adapter.send_message('person', 'kf-explicit|uexternal-user', message)
kwargs = adapter.bot.send_text_msg.await_args.kwargs
assert kwargs['open_kfid'] == 'kf-explicit'
assert kwargs['external_userid'] == 'external-user'
@pytest.mark.asyncio
async def test_send_message_requires_open_kfid():
adapter = make_adapter()
adapter.bot = SimpleNamespace(send_text_msg=AsyncMock())
message = platform_message.MessageChain([platform_message.Plain(text='hello')])
with pytest.raises(ValueError, match='open_kfid is required'):
await adapter.send_message('person', 'uexternal-user', message)
adapter.bot.send_text_msg.assert_not_called()
@pytest.mark.asyncio
async def test_send_message_rejects_group_targets():
adapter = make_adapter()
adapter.bot_account_id = 'kf-test'
adapter.bot = SimpleNamespace(send_text_msg=AsyncMock())
message = platform_message.MessageChain([platform_message.Plain(text='hello')])
with pytest.raises(ValueError, match='only supports sending messages to person'):
await adapter.send_message('group', 'group-id', message)
adapter.bot.send_text_msg.assert_not_called()
@@ -639,13 +639,10 @@ class TestGetRuntimeInfoDict:
assert info['box_session_id'] == 'mcp-shared'
assert info['box_enabled'] is True
def test_transient_test_shares_session_but_isolated_by_process(self, mcp_module):
"""A transient config-page "test" now shares the same 'mcp-shared' Box
session as live servers (so a test reuses the running container / live
process instead of a cold per-test session bootstrap). Isolation is at
the PROCESS level: the test runs under its own process_id and only ever
stops that process_id, so it cannot disturb another server's live
process or the shared session itself."""
def test_transient_test_session_is_isolated_from_shared(self, mcp_module):
"""A transient test session (config-page "test", no persisted UUID)
must NOT share the live "mcp-shared" Box session. Regression: a failing
test churned the shared session and tore down healthy live servers."""
ap = _make_ap()
ap.box_service.available = True
transient = _make_session(
@@ -673,12 +670,10 @@ class TestGetRuntimeInfoDict:
)
assert transient.is_transient is True
assert live.is_transient is False
# Both share ONE Box session ...
assert transient._build_box_session_id() == 'mcp-shared'
# Isolated session id for the test, shared for the live server.
assert transient._build_box_session_id() == 'mcp-test-gen-uuid-123'
assert live._build_box_session_id() == 'mcp-shared'
assert transient._build_box_session_id() == live._build_box_session_id()
# ... but are isolated by distinct process_ids within that session.
assert transient._box_stdio_runtime.process_id != live._box_stdio_runtime.process_id
assert transient._build_box_session_id() != live._build_box_session_id()
def test_stdio_session_refuses_when_box_unavailable(self, mcp_module):
"""Policy: when Box is configured but unavailable (disabled in config
@@ -829,129 +824,3 @@ async def test_init_box_stdio_server_stages_host_path_in_shared_workspace(mcp_mo
assert process_payload['command'] == 'python'
assert process_payload['args'] == ['/workspace/.mcp/u1/workspace/server.py']
assert process_payload['cwd'] == '/workspace/.mcp/u1/workspace'
@pytest.mark.asyncio
async def test_stdio_handshake_raises_coldstart_retry_while_process_alive(mcp_module, tmp_path, monkeypatch):
"""During a slow (npx) cold start the handshake fails while the managed
process is still alive. initialize() must raise _ColdStartRetry (so the
outer lifecycle loop reuses the live process and retries without stopping it
or consuming the fatal budget), NOT a fatal error."""
from contextlib import asynccontextmanager
mcp_stdio_module = sys.modules['langbot.pkg.provider.tools.loaders.mcp_stdio']
class ColdClientSession:
def __init__(self, *_args):
pass
async def __aenter__(self):
return self
async def __aexit__(self, exc_type, exc, tb):
return False
async def initialize(self):
# Process still cold-starting: handshake fails.
raise Exception('Connection closed')
@asynccontextmanager
async def fake_websocket_client(_url: str):
yield ('read-stream', 'write-stream')
monkeypatch.setattr(mcp_stdio_module, 'ClientSession', ColdClientSession)
monkeypatch.setattr(mcp_stdio_module, 'websocket_client', fake_websocket_client)
monkeypatch.setattr(mcp_stdio_module, '_HANDSHAKE_ATTEMPT_TIMEOUT_SEC', 1.0, raising=False)
ap = _make_ap()
ap.box_service.available = True
ap.box_service.create_session = AsyncMock(return_value={})
ap.box_service.start_managed_process = AsyncMock(return_value={})
ap.box_service.get_managed_process_websocket_url = Mock(return_value='ws://box/p')
session = _make_session(
mcp_module,
{
'name': 'slow',
'uuid': 'slow-uuid',
'mode': 'stdio',
'command': 'npx',
'args': ['-y', 'some-mcp'],
},
ap=ap,
)
# Process is NOT exited (still cold-starting) and not yet running for reuse.
async def _not_exited():
return False
session._box_stdio_runtime._managed_process_has_exited = _not_exited
async def _not_running():
return False
session._box_stdio_runtime._managed_process_is_running = _not_running
with pytest.raises(mcp_stdio_module._ColdStartRetry):
await session._init_box_stdio_server()
# Process was started exactly once (the retry will reuse it, not rebuild).
assert ap.box_service.start_managed_process.await_count == 1
await session.exit_stack.aclose()
@pytest.mark.asyncio
async def test_stdio_handshake_raises_fatal_when_process_exited(mcp_module, tmp_path, monkeypatch):
"""If the handshake fails AND the process has definitively exited, that is a
real failure initialize() must NOT swallow it as a cold-start retry."""
from contextlib import asynccontextmanager
mcp_stdio_module = sys.modules['langbot.pkg.provider.tools.loaders.mcp_stdio']
class DeadClientSession:
def __init__(self, *_args):
pass
async def __aenter__(self):
return self
async def __aexit__(self, exc_type, exc, tb):
return False
async def initialize(self):
raise Exception('Connection closed')
@asynccontextmanager
async def fake_websocket_client(_url: str):
yield ('read-stream', 'write-stream')
monkeypatch.setattr(mcp_stdio_module, 'ClientSession', DeadClientSession)
monkeypatch.setattr(mcp_stdio_module, 'websocket_client', fake_websocket_client)
monkeypatch.setattr(mcp_stdio_module, '_HANDSHAKE_ATTEMPT_TIMEOUT_SEC', 1.0, raising=False)
ap = _make_ap()
ap.box_service.available = True
ap.box_service.create_session = AsyncMock(return_value={})
ap.box_service.start_managed_process = AsyncMock(return_value={})
ap.box_service.get_managed_process_websocket_url = Mock(return_value='ws://box/p')
session = _make_session(
mcp_module,
{'name': 'dead', 'uuid': 'dead-uuid', 'mode': 'stdio', 'command': 'npx', 'args': ['-y', 'x']},
ap=ap,
)
async def _exited():
return True
session._box_stdio_runtime._managed_process_has_exited = _exited
async def _not_running():
return False
session._box_stdio_runtime._managed_process_is_running = _not_running
with pytest.raises(Exception) as ei:
await session._init_box_stdio_server()
assert not isinstance(ei.value, mcp_stdio_module._ColdStartRetry)
await session.exit_stack.aclose()
Generated
+26 -5
View File
@@ -1814,6 +1814,19 @@ wheels = [
{ url = "https://files.pythonhosted.org/packages/cb/b1/3846dd7f199d53cb17f49cba7e651e9ce294d8497c8c150530ed11865bb8/iniconfig-2.3.0-py3-none-any.whl", hash = "sha256:f631c04d2c48c52b84d0d0549c99ff3859c98df65b3101406327ecc7d53fbf12", size = 7484, upload-time = "2025-10-18T21:55:41.639Z" },
]
[[package]]
name = "itchat-uos"
version = "1.5.0.dev0"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "pypng" },
{ name = "pyqrcode" },
{ name = "requests" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/31/2b/0be7e46195dc3c461518b046a6fb9c34c97cd2fd44847b840a0011686575/itchat_uos-1.5.0.dev0-py3-none-any.whl", hash = "sha256:0293b77cab31fa8c9c2144ea8b5636d43836f47855beba0c603ad3fb1ff89625", size = 52507, upload-time = "2022-07-14T02:38:32.507Z" },
]
[[package]]
name = "itsdangerous"
version = "2.2.0"
@@ -2008,7 +2021,7 @@ wheels = [
[[package]]
name = "langbot"
version = "4.10.5"
version = "4.10.4"
source = { editable = "." }
dependencies = [
{ name = "aiocqhttp" },
@@ -2035,6 +2048,7 @@ dependencies = [
{ name = "ebooklib" },
{ name = "gewechat-client" },
{ name = "html2text" },
{ name = "itchat-uos" },
{ name = "langbot-plugin" },
{ name = "langchain" },
{ name = "langchain-core" },
@@ -2123,7 +2137,8 @@ requires-dist = [
{ name = "ebooklib", specifier = ">=0.18" },
{ name = "gewechat-client", specifier = ">=0.1.5" },
{ name = "html2text", specifier = ">=2024.2.26" },
{ name = "langbot-plugin", specifier = "==0.4.9" },
{ name = "itchat-uos", specifier = ">=1.5.0.dev0" },
{ name = "langbot-plugin", specifier = "==0.4.6" },
{ name = "langchain", specifier = ">=1.3.9" },
{ name = "langchain-core", specifier = ">=1.3.3" },
{ name = "langchain-text-splitters", specifier = ">=1.1.2" },
@@ -2187,7 +2202,7 @@ dev = [
[[package]]
name = "langbot-plugin"
version = "0.4.9"
version = "0.4.6"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "aiofiles" },
@@ -2208,9 +2223,9 @@ dependencies = [
{ name = "watchdog" },
{ name = "websockets" },
]
sdist = { url = "https://files.pythonhosted.org/packages/54/52/f85939d3be929696dc5f9f3369f18bd6652981a8985c6e6f0dfe48e04f75/langbot_plugin-0.4.9.tar.gz", hash = "sha256:1cc2882f1be96cd52be6c392c3925e44cf785c97c10f82a321339636da318c6e", size = 334682, upload-time = "2026-07-03T04:57:42.813Z" }
sdist = { url = "https://files.pythonhosted.org/packages/b4/6a/5fdb5365ad04aaa61344e92578d73eb1577af35783b80767c7d6c51cb8b9/langbot_plugin-0.4.6.tar.gz", hash = "sha256:838e3cd45ed795ed4c3299c73f141b217adfa05f09937a01694e7158619e4f6e", size = 334171, upload-time = "2026-06-22T15:06:56.565Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/47/12/1cf535377a81cf01607bbee5e3f4745a8c9f6dbb561887467001a182a533/langbot_plugin-0.4.9-py3-none-any.whl", hash = "sha256:693225a089ba38bc8d0b58d9e0fdf07b027b5be6b5e858cf7b4e56d095b68ba5", size = 221722, upload-time = "2026-07-03T04:57:41.578Z" },
{ url = "https://files.pythonhosted.org/packages/6b/55/7adc2e180a299ed58613e159c64195477e09c05136f949942c6cec5219e8/langbot_plugin-0.4.6-py3-none-any.whl", hash = "sha256:30eb47efc0b703818ac003a5cd67caf720d9749dd503155eb65cce0c28b194a7", size = 217434, upload-time = "2026-06-22T15:06:55.237Z" },
]
[[package]]
@@ -4517,6 +4532,12 @@ wheels = [
{ url = "https://files.pythonhosted.org/packages/bd/24/12818598c362d7f300f18e74db45963dbcb85150324092410c8b49405e42/pyproject_hooks-1.2.0-py3-none-any.whl", hash = "sha256:9e5c6bfa8dcc30091c74b0cf803c81fdd29d94f01992a7707bc97babb1141913", size = 10216, upload-time = "2024-09-29T09:24:11.978Z" },
]
[[package]]
name = "pyqrcode"
version = "1.2.1"
source = { registry = "https://pypi.org/simple" }
sdist = { url = "https://files.pythonhosted.org/packages/37/61/f07226075c347897937d4086ef8e55f0a62ae535e28069884ac68d979316/PyQRCode-1.2.1.tar.gz", hash = "sha256:fdbf7634733e56b72e27f9bce46e4550b75a3a2c420414035cae9d9d26b234d5", size = 36989, upload-time = "2016-06-20T03:28:03.411Z" }
[[package]]
name = "pyreadline3"
version = "3.5.4"
+126 -50
View File
@@ -2,6 +2,7 @@ import { useState, useEffect, useRef, useCallback } from 'react';
import { useNavigate } from 'react-router-dom';
import { Tabs, TabsList, TabsTrigger, TabsContent } from '@/components/ui/tabs';
import { Button } from '@/components/ui/button';
import { Badge } from '@/components/ui/badge';
import { Switch } from '@/components/ui/switch';
import { Label } from '@/components/ui/label';
import {
@@ -26,9 +27,79 @@ import type { BotSessionMonitorHandle } from '@/app/home/bots/components/bot-ses
import { httpClient } from '@/app/infra/http/HttpClient';
import { useSidebarData } from '@/app/home/components/home-sidebar/SidebarDataContext';
import { useTranslation } from 'react-i18next';
import { Settings, FileText, Users, RefreshCw, Trash2 } from 'lucide-react';
import {
Settings,
FileText,
Users,
RefreshCw,
Trash2,
CircleCheck,
CircleAlert,
Loader2,
CircleOff,
} from 'lucide-react';
import { cn } from '@/lib/utils';
import { toast } from 'sonner';
import type { Bot, BotAdapterRuntimeStatus } from '@/app/infra/entities/api';
function getBotRuntimeStatus(bot: Bot | null): BotAdapterRuntimeStatus | null {
return bot?.adapter_runtime_values?.runtime_status ?? null;
}
function RuntimeStatusBadge({
status,
}: {
status: BotAdapterRuntimeStatus | null;
}) {
const { t } = useTranslation();
if (!status) return null;
const value = status?.connection_status ?? 'disconnected';
const config = {
connected: {
label: t('bots.runtimeConnected'),
className: 'border-emerald-500/30 bg-emerald-500/10 text-emerald-700',
icon: CircleCheck,
},
connecting: {
label: t('bots.runtimeConnecting'),
className: 'border-amber-500/30 bg-amber-500/10 text-amber-700',
icon: Loader2,
},
disconnected: {
label: t('bots.runtimeDisconnected'),
className: 'border-muted-foreground/20 bg-muted text-muted-foreground',
icon: CircleOff,
},
error: {
label: t('bots.runtimeError'),
className: 'border-destructive/30 bg-destructive/10 text-destructive',
icon: CircleAlert,
},
}[value];
const Icon = config.icon;
return (
<div className="flex min-w-0 items-center gap-2">
<Badge
variant="outline"
className={cn('h-6 gap-1.5 px-2 text-xs', config.className)}
>
<Icon
className={cn('size-3.5', value === 'connecting' && 'animate-spin')}
/>
{config.label}
</Badge>
{status?.connection_error && (
<span className="max-w-[360px] truncate text-xs text-destructive">
{status.connection_error}
</span>
)}
</div>
);
}
export default function BotDetailContent({ id }: { id: string }) {
const isCreateMode = id === 'new';
@@ -58,16 +129,24 @@ export default function BotDetailContent({ id }: { id: string }) {
// Enable state managed here so the header switch works
const [botEnabled, setBotEnabled] = useState(true);
const [enableLoaded, setEnableLoaded] = useState(false);
const [botDetail, setBotDetail] = useState<Bot | null>(null);
const fetchBotDetail = useCallback(async () => {
if (isCreateMode) return;
const res = await httpClient.getBot(id);
setBotDetail(res.bot);
setBotEnabled(res.bot.enable ?? true);
setEnableLoaded(true);
}, [id, isCreateMode]);
// Fetch bot enable state
useEffect(() => {
if (!isCreateMode) {
httpClient.getBot(id).then((res) => {
setBotEnabled(res.bot.enable ?? true);
setEnableLoaded(true);
});
fetchBotDetail();
const timer = window.setInterval(fetchBotDetail, 5000);
return () => window.clearInterval(timer);
}
}, [id, isCreateMode]);
}, [fetchBotDetail, isCreateMode]);
const handleEnableToggle = useCallback(
async (checked: boolean) => {
@@ -95,9 +174,7 @@ export default function BotDetailContent({ id }: { id: string }) {
function handleFormSubmit() {
// Re-sync enable state after form save (form may update enable too)
httpClient.getBot(id).then((res) => {
setBotEnabled(res.bot.enable ?? true);
});
fetchBotDetail();
refreshBots();
}
@@ -173,6 +250,7 @@ export default function BotDetailContent({ id }: { id: string }) {
</Label>
</div>
)}
<RuntimeStatusBadge status={getBotRuntimeStatus(botDetail)} />
</div>
<Button
type="submit"
@@ -191,47 +269,45 @@ export default function BotDetailContent({ id }: { id: string }) {
onValueChange={setActiveTab}
className="flex flex-1 flex-col min-h-0"
>
<div className="flex shrink-0 items-center gap-1">
<TabsList>
<TabsTrigger value="config" className="gap-1.5">
<Settings className="size-3.5" />
{t('bots.configuration')}
</TabsTrigger>
<TabsTrigger value="logs" className="gap-1.5">
<FileText className="size-3.5" />
{t('bots.logs')}
</TabsTrigger>
<TabsTrigger value="sessions" className="gap-1.5">
<Users className="size-3.5" />
{t('bots.sessionMonitor.title')}
</TabsTrigger>
</TabsList>
{activeTab === 'sessions' && (
<button
type="button"
aria-label={t('bots.sessionMonitor.refresh')}
title={t('bots.sessionMonitor.refresh')}
className="inline-flex h-8 w-8 items-center justify-center rounded-md text-muted-foreground transition-colors hover:bg-accent hover:text-foreground disabled:pointer-events-none disabled:opacity-50"
disabled={isRefreshingSessions}
onClick={() => {
if (isRefreshingSessions) return;
setIsRefreshingSessions(true);
const minDelay = new Promise((r) => setTimeout(r, 500));
Promise.all([
sessionMonitorRef.current?.refreshSessions(),
minDelay,
]).finally(() => setIsRefreshingSessions(false));
}}
>
<RefreshCw
className={cn(
'size-3.5',
isRefreshingSessions && 'animate-spin',
)}
/>
</button>
)}
</div>
<TabsList className="shrink-0">
<TabsTrigger value="config" className="gap-1.5">
<Settings className="size-3.5" />
{t('bots.configuration')}
</TabsTrigger>
<TabsTrigger value="logs" className="gap-1.5">
<FileText className="size-3.5" />
{t('bots.logs')}
</TabsTrigger>
<TabsTrigger value="sessions" className="gap-1.5">
<Users className="size-3.5" />
{t('bots.sessionMonitor.title')}
{activeTab === 'sessions' && (
<button
type="button"
className="inline-flex items-center justify-center ml-0.5"
onPointerDown={(e) => e.stopPropagation()}
onClick={(e) => {
e.stopPropagation();
e.preventDefault();
if (isRefreshingSessions) return;
setIsRefreshingSessions(true);
const minDelay = new Promise((r) => setTimeout(r, 500));
Promise.all([
sessionMonitorRef.current?.refreshSessions(),
minDelay,
]).finally(() => setIsRefreshingSessions(false));
}}
>
<RefreshCw
className={cn(
'size-3 text-muted-foreground hover:text-foreground transition-colors',
isRefreshingSessions && 'animate-spin',
)}
/>
</button>
)}
</TabsTrigger>
</TabsList>
{/* Tab: Configuration */}
<TabsContent
@@ -3,7 +3,6 @@ import React, {
useEffect,
useRef,
useCallback,
useMemo,
forwardRef,
useImperativeHandle,
} from 'react';
@@ -16,14 +15,11 @@ import {
Bot,
Copy,
Check,
ChevronDown,
ChevronRight,
Workflow,
ThumbsUp,
ThumbsDown,
ShieldCheck,
ShieldOff,
Wrench,
} from 'lucide-react';
import { toast } from 'sonner';
import BotAdminsDialog, {
@@ -80,35 +76,6 @@ interface SessionFeedback {
stream_id?: string | null;
}
interface SessionToolCall {
id: string;
timestamp: string;
tool_name: string;
tool_source: string;
duration: number;
status: string;
message_id?: string | null;
arguments?: string | null;
result?: string | null;
error_message?: string | null;
}
type SessionTimelineItem =
| {
id: string;
type: 'message';
timestamp: number;
order: number;
message: SessionMessage;
}
| {
id: string;
type: 'tool';
timestamp: number;
order: number;
toolCall: SessionToolCall;
};
export interface BotSessionMonitorHandle {
refreshSessions: () => Promise<void>;
}
@@ -133,10 +100,6 @@ const BotSessionMonitor = forwardRef<
const [feedbackMap, setFeedbackMap] = useState<
Record<string, SessionFeedback>
>({});
const [toolCalls, setToolCalls] = useState<SessionToolCall[]>([]);
const [expandedToolCallIds, setExpandedToolCallIds] = useState<
Record<string, boolean>
>({});
const messagesContainerRef = useRef<HTMLDivElement>(null);
const { admins, reload: reloadAdmins } = useBotAdmins(botId);
const [adminsDialogOpen, setAdminsDialogOpen] = useState(false);
@@ -226,7 +189,6 @@ const BotSessionMonitor = forwardRef<
const loadMessages = useCallback(
async (sessionId: string) => {
setLoadingMessages(true);
setExpandedToolCallIds({});
try {
const messagesRes = await httpClient.getSessionMessages(sessionId);
const sorted = (messagesRes.messages ?? []).sort(
@@ -235,18 +197,6 @@ const BotSessionMonitor = forwardRef<
);
setMessages(sorted);
try {
const analysisRes = await httpClient.get<{
tool_calls?: SessionToolCall[];
}>(
`/api/v1/monitoring/sessions/${encodeURIComponent(sessionId)}/analysis`,
);
setToolCalls(analysisRes?.tool_calls ?? []);
} catch (analysisError) {
console.error('Failed to load session tool calls:', analysisError);
setToolCalls([]);
}
// Collect user message IDs for feedback matching
const userMsgIds = new Set(
sorted.filter((m) => !m.role || m.role === 'user').map((m) => m.id),
@@ -290,14 +240,11 @@ const BotSessionMonitor = forwardRef<
loadMessages(selectedSessionId);
} else {
setMessages([]);
setToolCalls([]);
setExpandedToolCallIds({});
setFeedbackMap({});
}
}, [selectedSessionId, loadMessages]);
useEffect(() => {
if (messages.length === 0 && toolCalls.length === 0) return;
if (messages.length === 0) return;
// Wait for DOM to render the new messages before scrolling
requestAnimationFrame(() => {
const container = messagesContainerRef.current;
@@ -309,7 +256,7 @@ const BotSessionMonitor = forwardRef<
scrollTarget.scrollTop = scrollTarget.scrollHeight;
}
});
}, [messages, toolCalls]);
}, [messages]);
const parseMessageChain = (content: string): MessageChainComponent[] => {
try {
@@ -484,71 +431,6 @@ const BotSessionMonitor = forwardRef<
return `${diffDays}d`;
};
const formatDuration = (durationMs: number): string => {
if (!durationMs) return '0ms';
if (durationMs < 1000) return `${durationMs}ms`;
return `${(durationMs / 1000).toFixed(2)}s`;
};
const truncateToolDetail = (value?: string | null): string => {
if (!value) return '';
return value.length > 600 ? `${value.slice(0, 600)}...` : value;
};
const toggleToolCallDetails = (toolCallId: string) => {
setExpandedToolCallIds((previous) => ({
...previous,
[toolCallId]: !previous[toolCallId],
}));
};
const feedbackByMessageId = useMemo(() => {
const map: Record<string, SessionFeedback> = {};
for (let index = 0; index < messages.length; index++) {
const msg = messages[index];
if (isUserMessage(msg)) continue;
for (let previousIndex = index - 1; previousIndex >= 0; previousIndex--) {
const previousMessage = messages[previousIndex];
if (isUserMessage(previousMessage)) {
const feedback = feedbackMap[previousMessage.id];
if (feedback) {
map[msg.id] = feedback;
}
break;
}
}
}
return map;
}, [feedbackMap, messages]);
const timelineItems = useMemo<SessionTimelineItem[]>(() => {
const messageItems: SessionTimelineItem[] = messages.map(
(message, index) => ({
id: `message-${message.id}`,
type: 'message',
timestamp: parseTimestamp(message.timestamp).getTime(),
order: index * 2,
message,
}),
);
const toolItems: SessionTimelineItem[] = toolCalls.map(
(toolCall, index) => ({
id: `tool-${toolCall.id}`,
type: 'tool',
timestamp: parseTimestamp(toolCall.timestamp).getTime(),
order: index * 2 + 1,
toolCall,
}),
);
return [...messageItems, ...toolItems].sort(
(a, b) => a.timestamp - b.timestamp || a.order - b.order,
);
}, [messages, toolCalls]);
const selectedSession = sessions.find(
(s) => s.session_id === selectedSessionId,
);
@@ -730,162 +612,29 @@ const BotSessionMonitor = forwardRef<
<div className="text-center text-muted-foreground py-12 text-sm">
{t('bots.sessionMonitor.loading')}
</div>
) : timelineItems.length === 0 ? (
) : messages.length === 0 ? (
<div className="text-center text-muted-foreground py-12 text-sm">
{t('bots.sessionMonitor.noMessages')}
</div>
) : (
timelineItems.map((item) => {
if (item.type === 'tool') {
const call = item.toolCall;
const hasToolDetails = Boolean(
call.arguments || call.result || call.error_message,
);
const expandedToolCall = Boolean(
expandedToolCallIds[call.id],
);
const detailsId = `tool-call-details-${call.id}`;
return (
<div key={item.id} className="flex justify-start">
<div className="max-w-2xl rounded-xl rounded-bl-sm border border-border/60 bg-muted/25 px-2.5 py-1.5 text-xs text-muted-foreground">
<button
type="button"
className={cn(
'flex w-full items-center justify-between gap-3 rounded-md text-left outline-none transition-colors',
hasToolDetails &&
'cursor-pointer hover:bg-muted/40 focus-visible:ring-2 focus-visible:ring-ring',
)}
aria-expanded={
hasToolDetails ? expandedToolCall : undefined
}
aria-controls={
hasToolDetails ? detailsId : undefined
}
aria-disabled={!hasToolDetails}
onClick={() =>
hasToolDetails &&
toggleToolCallDetails(call.id)
}
>
<div className="flex min-w-0 flex-wrap items-center gap-1.5">
{hasToolDetails &&
(expandedToolCall ? (
<ChevronDown className="h-3.5 w-3.5 shrink-0 text-muted-foreground/70" />
) : (
<ChevronRight className="h-3.5 w-3.5 shrink-0 text-muted-foreground/70" />
))}
<Wrench className="h-3.5 w-3.5 shrink-0 text-muted-foreground/70" />
<span className="min-w-0 max-w-[18rem] truncate text-[13px] font-medium text-foreground/75">
{call.tool_name}
</span>
<span className="rounded border border-border/50 bg-background/60 px-1.5 py-0.5 text-[10px] leading-none text-muted-foreground">
{call.tool_source}
</span>
<span
className={cn(
'rounded px-1.5 py-0.5 text-[10px] font-medium leading-none',
call.status === 'success'
? 'bg-green-100/70 text-green-700 dark:bg-green-950/60 dark:text-green-300'
: 'bg-red-100/70 text-red-700 dark:bg-red-950/60 dark:text-red-300',
)}
>
{call.status}
</span>
</div>
<span className="shrink-0 text-[11px] tabular-nums text-muted-foreground/80">
{formatDuration(call.duration)}
</span>
</button>
{hasToolDetails && expandedToolCall && (
<div
id={detailsId}
className="mt-2 space-y-1.5"
>
{(call.arguments || call.result) && (
<div className="space-y-1.5">
{call.arguments && (
<div>
<div className="mb-1 text-[11px] font-medium text-muted-foreground">
{t(
'monitoring.toolCalls.arguments',
{
defaultValue: '参数',
},
)}
</div>
<pre className="whitespace-pre-wrap break-words rounded bg-background/80 p-2 font-mono text-[11px] leading-4 text-muted-foreground">
{truncateToolDetail(call.arguments)}
</pre>
</div>
)}
{call.result && (
<div>
<div className="mb-1 text-[11px] font-medium text-muted-foreground">
{t('monitoring.toolCalls.result', {
defaultValue: '结果',
})}
</div>
<pre className="whitespace-pre-wrap break-words rounded bg-background/80 p-2 font-mono text-[11px] leading-4 text-muted-foreground">
{truncateToolDetail(call.result)}
</pre>
</div>
)}
</div>
)}
{call.error_message && (
<div className="whitespace-pre-wrap break-words rounded bg-red-50 p-2 text-[11px] text-red-600 dark:bg-red-950/40 dark:text-red-400">
{call.error_message}
</div>
)}
</div>
)}
<div className="mt-1.5 flex items-center gap-1.5 text-[11px] text-muted-foreground">
<span>
{t('monitoring.toolCalls.title', {
defaultValue: '工具调用',
})}
</span>
<span className="tabular-nums">
{formatTime(call.timestamp)}
</span>
{hasToolDetails && (
<>
<span>·</span>
<span>
{expandedToolCall
? t(
'monitoring.toolCalls.hideDetails',
{
defaultValue: '隐藏详情',
},
)
: t(
'monitoring.toolCalls.showDetails',
{
defaultValue: '查看详情',
},
)}
</span>
</>
)}
</div>
</div>
</div>
);
}
const msg = item.message;
messages.map((msg, msgIndex) => {
const isUser = isUserMessage(msg);
const isDiscarded =
msg.status === 'discarded' ||
msg.pipeline_id === PIPELINE_DISCARD;
const msgFeedback = feedbackByMessageId[msg.id];
// For bot replies, find feedback linked to the preceding user message
let msgFeedback: SessionFeedback | undefined;
if (!isUser) {
for (let i = msgIndex - 1; i >= 0; i--) {
if (isUserMessage(messages[i])) {
msgFeedback = feedbackMap[messages[i].id];
break;
}
}
}
return (
<div
key={item.id}
key={msg.id}
className={cn(
'flex',
isUser ? 'justify-end' : 'justify-start',
@@ -624,7 +624,10 @@ export default function DynamicFormComponent({
onSuccess={(credentials) => {
for (const [key, value] of Object.entries(credentials)) {
if (value) {
form.setValue(key as keyof FormValues, value as never);
form.setValue(key as keyof FormValues, value as never, {
shouldDirty: true,
shouldValidate: true,
});
}
}
}}
@@ -67,14 +67,16 @@ import SettingsDialog, {
import ToolResourceSelectors from '@/app/home/components/dynamic-form/ToolResourceSelectors';
import { LANGBOT_MODELS_PROVIDER_REQUESTER } from '@/app/home/components/models-dialog/types';
function hasUsableUuid<T extends { uuid?: string | null }>(
item: T,
): item is T & { uuid: string } {
return typeof item.uuid === 'string' && item.uuid.trim().length > 0;
const EMPTY_SELECT_ITEM_VALUE = '__langbot_empty_select_item_value__';
function toSelectValue(value: unknown): string {
return typeof value === 'string' ? value : '';
}
function hasUsableOptionName(option: { name?: string | null }): boolean {
return typeof option.name === 'string' && option.name.trim().length > 0;
function hasNonEmptyUuid<T extends { uuid?: string | null }>(
item: T,
): item is T & { uuid: string } {
return typeof item.uuid === 'string' && item.uuid.length > 0;
}
export default function DynamicFormItemComponent({
@@ -114,7 +116,7 @@ export default function DynamicFormItemComponent({
httpClient
.getProviderLLMModels()
.then((resp) => {
setLlmModels(resp.models.filter(hasUsableUuid));
setLlmModels(resp.models);
})
.catch((err) => {
toast.error(t('models.getModelListError') + err.msg);
@@ -125,7 +127,7 @@ export default function DynamicFormItemComponent({
httpClient
.getProviderEmbeddingModels()
.then((resp) => {
setEmbeddingModels(resp.models.filter(hasUsableUuid));
setEmbeddingModels(resp.models);
})
.catch((err) => {
toast.error(t('embedding.getModelListError') + err.msg);
@@ -136,7 +138,7 @@ export default function DynamicFormItemComponent({
httpClient
.getProviderRerankModels()
.then((resp) => {
setRerankModels(resp.models.filter(hasUsableUuid));
setRerankModels(resp.models);
})
.catch((err) => {
toast.error('Failed to load rerank models: ' + err.msg);
@@ -240,7 +242,7 @@ export default function DynamicFormItemComponent({
httpClient
.getKnowledgeBases()
.then((resp) => {
setKnowledgeBases(resp.bases.filter(hasUsableUuid));
setKnowledgeBases(resp.bases);
})
.catch((err) => {
toast.error(t('knowledge.getKnowledgeBaseListError') + err.msg);
@@ -253,7 +255,7 @@ export default function DynamicFormItemComponent({
httpClient
.getBots()
.then((resp) => {
setBots(resp.bots.filter(hasUsableUuid));
setBots(resp.bots);
})
.catch((err) => {
toast.error(t('bots.getBotListError') + err.msg);
@@ -390,19 +392,33 @@ export default function DynamicFormItemComponent({
</div>
);
case DynamicFormItemType.SELECT:
case DynamicFormItemType.SELECT: {
const hasEmptyOption =
config.options?.some((option) => option.name === '') ?? false;
const selectValue =
hasEmptyOption && field.value === ''
? EMPTY_SELECT_ITEM_VALUE
: toSelectValue(field.value);
return (
<Select value={field.value} onValueChange={field.onChange}>
<Select
value={selectValue}
onValueChange={(value) =>
field.onChange(value === EMPTY_SELECT_ITEM_VALUE ? '' : value)
}
>
<SelectTrigger className="w-full max-w-md bg-[#ffffff] dark:bg-[#2a2a2e]">
<SelectValue placeholder={t('common.select')} />
</SelectTrigger>
<SelectContent>
<SelectGroup>
{config.options?.filter(hasUsableOptionName).map((option) => (
{config.options?.map((option, index) => (
<SelectItem
key={option.name}
value={option.name}
description={option.name}
key={`${option.name}-${index}`}
value={
option.name === '' ? EMPTY_SELECT_ITEM_VALUE : option.name
}
description={option.name || undefined}
>
{extractI18nObject(option.label)}
</SelectItem>
@@ -411,6 +427,7 @@ export default function DynamicFormItemComponent({
</SelectContent>
</Select>
);
}
case DynamicFormItemType.LLM_MODEL_SELECTOR:
// Separate space models from regular models
@@ -465,7 +482,7 @@ export default function DynamicFormItemComponent({
{Object.entries(groupedModels).map(([providerName, models]) => (
<SelectGroup key={providerName}>
<SelectLabel>{providerName}</SelectLabel>
{models.map((model) => (
{models.filter(hasNonEmptyUuid).map((model) => (
<SelectItem key={model.uuid} value={model.uuid}>
<span className="inline-flex items-center gap-1">
{model.name}
@@ -569,7 +586,7 @@ export default function DynamicFormItemComponent({
{providerName}
</span>
</SelectLabel>
{models.map((model) => (
{models.filter(hasNonEmptyUuid).map((model) => (
<SelectItem key={model.uuid} value={model.uuid}>
<span className="inline-flex items-center gap-1">
{model.name}
@@ -666,7 +683,7 @@ export default function DynamicFormItemComponent({
([providerName, models]) => (
<SelectGroup key={providerName}>
<SelectLabel>{providerName}</SelectLabel>
{models.map((model) => (
{models.filter(hasNonEmptyUuid).map((model) => (
<SelectItem key={model.uuid} value={model.uuid}>
{model.name}
</SelectItem>
@@ -758,7 +775,7 @@ export default function DynamicFormItemComponent({
{providerName}
</span>
</SelectLabel>
{models.map((model) => (
{models.filter(hasNonEmptyUuid).map((model) => (
<SelectItem key={model.uuid} value={model.uuid}>
{model.name}
</SelectItem>
@@ -823,7 +840,7 @@ export default function DynamicFormItemComponent({
([providerName, models]) => (
<SelectGroup key={providerName}>
<SelectLabel>{providerName}</SelectLabel>
{models.map((model) => (
{models.filter(hasNonEmptyUuid).map((model) => (
<SelectItem key={model.uuid} value={model.uuid}>
{model.name}
</SelectItem>
@@ -918,7 +935,7 @@ export default function DynamicFormItemComponent({
([providerName, models]) => (
<SelectGroup key={providerName}>
<SelectLabel>{providerName}</SelectLabel>
{models.map((model) => (
{models.filter(hasNonEmptyUuid).map((model) => (
<SelectItem key={model.uuid} value={model.uuid}>
<span className="inline-flex items-center gap-1">
{model.name}
@@ -1023,7 +1040,7 @@ export default function DynamicFormItemComponent({
{providerName}
</span>
</SelectLabel>
{models.map((model) => (
{models.filter(hasNonEmptyUuid).map((model) => (
<SelectItem key={model.uuid} value={model.uuid}>
<span className="inline-flex items-center gap-1">
{model.name}
@@ -1186,8 +1203,7 @@ export default function DynamicFormItemComponent({
case DynamicFormItemType.KNOWLEDGE_BASE_SELECTOR:
// Group KBs by Knowledge Engine name
const validKnowledgeBases = knowledgeBases.filter(hasUsableUuid);
const kbsByEngine = validKnowledgeBases.reduce(
const kbsByEngine = knowledgeBases.reduce(
(acc, kb) => {
const engineName = kb.knowledge_engine?.name
? extractI18nObject(kb.knowledge_engine.name)
@@ -1198,7 +1214,7 @@ export default function DynamicFormItemComponent({
acc[engineName].push(kb);
return acc;
},
{} as Record<string, typeof validKnowledgeBases>,
{} as Record<string, typeof knowledgeBases>,
);
return (
@@ -1206,7 +1222,7 @@ export default function DynamicFormItemComponent({
<SelectTrigger className="min-w-0 bg-[#ffffff] dark:bg-[#2a2a2e]">
{field.value && field.value !== '__none__' ? (
(() => {
const selectedKb = validKnowledgeBases.find(
const selectedKb = knowledgeBases.find(
(kb) => kb.uuid === field.value,
);
return (
@@ -1234,7 +1250,7 @@ export default function DynamicFormItemComponent({
{Object.entries(kbsByEngine).map(([engineName, kbs]) => (
<SelectGroup key={engineName}>
<SelectLabel>{engineName}</SelectLabel>
{kbs.map((base) => (
{kbs.filter(hasNonEmptyUuid).map((base) => (
<SelectItem key={base.uuid} value={base.uuid}>
<div className="flex items-center gap-2">
{base.emoji && (
@@ -1252,8 +1268,7 @@ export default function DynamicFormItemComponent({
case DynamicFormItemType.KNOWLEDGE_BASE_MULTI_SELECTOR:
// Group KBs by Knowledge Engine name for multi-selector
const validMultiKnowledgeBases = knowledgeBases.filter(hasUsableUuid);
const multiKbsByEngine = validMultiKnowledgeBases.reduce(
const multiKbsByEngine = knowledgeBases.reduce(
(acc, kb) => {
const engineName = kb.knowledge_engine?.name
? extractI18nObject(kb.knowledge_engine.name)
@@ -1264,7 +1279,7 @@ export default function DynamicFormItemComponent({
acc[engineName].push(kb);
return acc;
},
{} as Record<string, typeof validMultiKnowledgeBases>,
{} as Record<string, typeof knowledgeBases>,
);
return (
@@ -1273,7 +1288,7 @@ export default function DynamicFormItemComponent({
{field.value && field.value.length > 0 ? (
<div className="min-w-0 space-y-2">
{field.value.map((kbId: string) => {
const currentKb = validMultiKnowledgeBases.find(
const currentKb = knowledgeBases.find(
(base) => base.uuid === kbId,
);
if (!currentKb) return null;
@@ -1359,13 +1374,15 @@ export default function DynamicFormItemComponent({
{engineName}
</div>
{kbs.map((base) => {
const isSelected = tempSelectedKBIds.includes(base.uuid);
const isSelected = tempSelectedKBIds.includes(
base.uuid ?? '',
);
return (
<div
key={base.uuid}
className="flex items-center gap-3 rounded-lg border p-3 hover:bg-accent cursor-pointer"
onClick={() => {
const kbId = base.uuid;
const kbId = base.uuid ?? '';
setTempSelectedKBIds((prev) =>
prev.includes(kbId)
? prev.filter((id) => id !== kbId)
@@ -1427,7 +1444,7 @@ export default function DynamicFormItemComponent({
</SelectTrigger>
<SelectContent>
<SelectGroup>
{bots.filter(hasUsableUuid).map((bot) => (
{bots.filter(hasNonEmptyUuid).map((bot) => (
<SelectItem key={bot.uuid} value={bot.uuid}>
{bot.name}
</SelectItem>
@@ -264,48 +264,6 @@ function saveListExpansionState(state: SidebarListExpansionState) {
// Maximum number of entity sub-items visible before "More" toggle
const MAX_VISIBLE_ITEMS = 5;
const MCP_REFRESH_POLL_INTERVAL_MS = 1000;
const MCP_REFRESH_TIMEOUT_MS = 60000;
function sleep(ms: number) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
async function waitForMCPRefreshTask(taskId: number) {
const deadline = Date.now() + MCP_REFRESH_TIMEOUT_MS;
while (Date.now() < deadline) {
const task = await httpClient.getAsyncTask(taskId);
if (task.runtime.done) return task;
await sleep(MCP_REFRESH_POLL_INTERVAL_MS);
}
throw new Error(`Timed out waiting for MCP refresh task ${taskId}`);
}
async function refreshEnabledMCPConnections() {
const resp = await httpClient.getMCPServers();
const enabledServers = resp.servers.filter((server) => server.enable);
if (enabledServers.length === 0) return;
const taskResults = await Promise.allSettled(
enabledServers.map((server) => httpClient.testMCPServer(server.name, {})),
);
const taskIds: number[] = [];
for (const result of taskResults) {
if (
result.status === 'fulfilled' &&
typeof result.value.task_id === 'number'
) {
taskIds.push(result.value.task_id);
} else if (result.status === 'rejected') {
console.error('Failed to start MCP refresh task:', result.reason);
}
}
await Promise.allSettled(taskIds.map(waitForMCPRefreshTask));
}
// Sort entity items by updatedAt descending (most recent first), items without updatedAt go last
function sortByRecent(items: SidebarEntityItem[]): SidebarEntityItem[] {
@@ -394,19 +352,11 @@ function NavItems({
if (extRefreshing) return;
setExtRefreshing(true);
try {
const results = await Promise.allSettled([
await Promise.all([
sidebarData.refreshPlugins(),
sidebarData.refreshMCPServers(),
sidebarData.refreshSkills(),
refreshEnabledMCPConnections(),
]);
const mcpRefreshResult = results[2];
if (mcpRefreshResult.status === 'rejected') {
console.error(
'Failed to refresh MCP connections:',
mcpRefreshResult.reason,
);
}
await sidebarData.refreshMCPServers();
} finally {
setExtRefreshing(false);
}
@@ -16,7 +16,12 @@ import {
} from 'lucide-react';
import QRCode from 'qrcode';
export type QrLoginPlatform = 'feishu' | 'weixin' | 'dingtalk' | 'wecombot';
export type QrLoginPlatform =
| 'feishu'
| 'weixin'
| 'dingtalk'
| 'wecombot'
| 'itchat';
interface PlatformConfig {
titleKey: string;
@@ -92,6 +97,20 @@ const PLATFORM_CONFIGS: Record<QrLoginPlatform, PlatformConfig> = {
}),
successNoteKey: 'wecombot.robotNameNote',
},
itchat: {
titleKey: 'itchat.scanLogin',
connectingKey: 'itchat.connecting',
scanQRCodeKey: 'itchat.scanQRCode',
waitingKey: 'itchat.waitingForScan',
successKey: 'itchat.loginSuccess',
failedKey: 'itchat.loginFailed',
retryKey: 'itchat.retry',
apiBase: '/api/v1/platform/adapters/itchat/login',
extractSuccess: (data) => ({
account_id: data.wxid || '',
nickname: data.nickname || '',
}),
},
};
interface QrCodeLoginDialogProps {
@@ -116,6 +135,7 @@ export default function QrCodeLoginDialog({
const [state, setState] = useState<DialogState>('connecting');
const [qrDataUrl, setQrDataUrl] = useState('');
const qrDataUrlRef = useRef('');
const [expireIn, setExpireIn] = useState(0);
const [errorMessage, setErrorMessage] = useState('');
const pollTimerRef = useRef<ReturnType<typeof setInterval> | null>(null);
@@ -176,6 +196,7 @@ export default function QrCodeLoginDialog({
cleanedRef.current = false;
setState('connecting');
setQrDataUrl('');
qrDataUrlRef.current = '';
setExpireIn(0);
setErrorMessage('');
@@ -204,12 +225,14 @@ export default function QrCodeLoginDialog({
if (qr_data_url) {
setQrDataUrl(qr_data_url);
qrDataUrlRef.current = qr_data_url;
} else if (qr_url) {
const dataUrl = await QRCode.toDataURL(qr_url, {
width: 224,
margin: 2,
});
setQrDataUrl(dataUrl);
qrDataUrlRef.current = dataUrl;
}
setState('waiting');
@@ -289,6 +312,19 @@ export default function QrCodeLoginDialog({
cleanup();
setExpireIn(0);
setState('expired');
} else if (status === 'waiting') {
// Update QR data URL if regenerated (e.g. itchat QR expiry)
if (rest.qr_data_url && rest.qr_data_url !== qrDataUrlRef.current) {
setQrDataUrl(rest.qr_data_url);
qrDataUrlRef.current = rest.qr_data_url;
}
if (rest.expire_at) {
const remaining = Math.max(
0,
Math.floor(rest.expire_at - Date.now() / 1000),
);
setExpireIn(remaining);
}
}
} catch {
// ignore poll errors
@@ -157,10 +157,6 @@ export default function MCPDetailContent({ id }: { id: string }) {
navigate(`/home/mcp?id=${encodeURIComponent(serverName)}`);
}
const handlePersistedTestComplete = useCallback(async () => {
await refreshMCPServers();
}, [refreshMCPServers]);
function confirmDelete() {
httpClient
.deleteMCPServer(id)
@@ -368,7 +364,6 @@ export default function MCPDetailContent({ id }: { id: string }) {
onRuntimeInfoChange={(runtimeInfo) =>
setDetailRuntimeStatus(runtimeInfo?.status ?? null)
}
onPersistedTestComplete={handlePersistedTestComplete}
/>
</div>
</div>
@@ -41,7 +41,6 @@ import {
} from '@/components/ui/card';
import { httpClient } from '@/app/infra/http/HttpClient';
import { Tabs, TabsContent, TabsList, TabsTrigger } from '@/components/ui/tabs';
import MCPLogs from '@/app/home/mcp/components/mcp-form/MCPLogs';
import MCPReadme from '@/app/home/mcp/components/mcp-form/MCPReadme';
import {
MCPServerRuntimeInfo,
@@ -488,7 +487,6 @@ interface MCPFormProps {
onDirtyChange?: (dirty: boolean) => void;
onTestingChange?: (testing: boolean) => void;
onRuntimeInfoChange?: (runtimeInfo: MCPServerRuntimeInfo | null) => void;
onPersistedTestComplete?: (serverName: string) => void | Promise<void>;
/** Reported when the form cannot be saved because the current mode is
* ``stdio`` and the Box sandbox is disabled/unavailable. Parents that
* render the Save button outside this component should disable it. */
@@ -513,7 +511,6 @@ const MCPForm = forwardRef<MCPFormHandle, MCPFormProps>(function MCPForm(
onDirtyChange,
onTestingChange,
onRuntimeInfoChange,
onPersistedTestComplete,
onSaveBlockedChange,
layout = 'stacked',
sideHeader,
@@ -753,8 +750,6 @@ const MCPForm = forwardRef<MCPFormHandle, MCPFormProps>(function MCPForm(
}
try {
let serverConfig: MCPServer;
const serverName =
isEditMode && initServerName ? initServerName : value.name;
if (value.mode === 'remote') {
const headers: Record<string, string> = {};
@@ -763,7 +758,7 @@ const MCPForm = forwardRef<MCPFormHandle, MCPFormProps>(function MCPForm(
});
serverConfig = {
name: serverName,
name: value.name,
mode: 'remote',
enable: true,
extra_args: {
@@ -779,7 +774,7 @@ const MCPForm = forwardRef<MCPFormHandle, MCPFormProps>(function MCPForm(
});
serverConfig = {
name: serverName,
name: value.name,
mode: 'stdio',
enable: true,
extra_args: {
@@ -823,10 +818,6 @@ const MCPForm = forwardRef<MCPFormHandle, MCPFormProps>(function MCPForm(
// `uvx` with no package (exit 2 / "Connection closed", no detail).
// The form values are kept in sync on every edit and on load, so they
// are always current.
const serverName =
isEditMode && initServerName ? initServerName : form.getValues('name');
const shouldTestPersistedServer =
isEditMode && !!initServerName && !form.formState.isDirty;
const formExtraArgs = form.getValues('extra_args') ?? [];
const formStdioArgs = form.getValues('args') ?? [];
let extraArgsData: MCPServerExtraArgsRemote | MCPServerExtraArgsStdio;
@@ -849,20 +840,12 @@ const MCPForm = forwardRef<MCPFormHandle, MCPFormProps>(function MCPForm(
};
}
const testTarget = shouldTestPersistedServer ? serverName : '_';
const testPayload = shouldTestPersistedServer
? {}
: ({
name: serverName,
mode,
enable: true,
extra_args: extraArgsData,
} as MCPServer);
const { task_id } = await httpClient.testMCPServer(
testTarget,
testPayload,
);
const { task_id } = await httpClient.testMCPServer('_', {
name: form.getValues('name'),
mode,
enable: true,
extra_args: extraArgsData,
} as MCPServer);
if (!task_id) {
throw new Error(t('mcp.noTaskId'));
@@ -888,18 +871,14 @@ const MCPForm = forwardRef<MCPFormHandle, MCPFormProps>(function MCPForm(
resource_count: 0,
resources: [],
});
if (shouldTestPersistedServer) {
await onPersistedTestComplete?.(serverName);
}
} else {
if (shouldTestPersistedServer) {
await loadServerForEdit(serverName);
await onPersistedTestComplete?.(serverName);
if (isEditMode) {
await loadServerForEdit(form.getValues('name'));
} else {
// Transient tests have no persisted server to reload tools from.
// Create mode has no persisted server to reload tools from.
// The backend stashes the discovered runtime info (status +
// tools) in the task metadata before tearing the transient
// session down — surface it so a successful test
// tools) in the test task's metadata before tearing the
// transient session down — surface it so a successful test
// shows the tool list instead of "no tools found".
const runtimeInfoFromTest = taskResp.task_context?.metadata
?.runtime_info as MCPServerRuntimeInfo | undefined;
@@ -1184,14 +1163,11 @@ const MCPForm = forwardRef<MCPFormHandle, MCPFormProps>(function MCPForm(
</Card>
);
const persistedServerName =
isEditMode && initServerName ? initServerName : form.getValues('name');
const runtimePanel = (
<RuntimePanel
mcpTesting={mcpTesting}
runtimeInfo={runtimeInfo}
serverName={persistedServerName}
serverName={form.getValues('name')}
t={t}
/>
);
@@ -1224,9 +1200,6 @@ const MCPForm = forwardRef<MCPFormHandle, MCPFormProps>(function MCPForm(
<TabsTrigger value="resources" className="flex-none px-4">
{resourcesTabLabel}
</TabsTrigger>
<TabsTrigger value="logs" className="flex-none px-4">
{t('mcp.tabLogs')}
</TabsTrigger>
</TabsList>
<TabsContent value="docs" className="mt-4 min-h-0 flex-1 overflow-y-auto">
<MCPReadme readme={readme} />
@@ -1238,7 +1211,7 @@ const MCPForm = forwardRef<MCPFormHandle, MCPFormProps>(function MCPForm(
<RuntimePanel
mcpTesting={mcpTesting}
runtimeInfo={runtimeInfo}
serverName={persistedServerName}
serverName={form.getValues('name')}
content="tools"
t={t}
/>
@@ -1250,14 +1223,11 @@ const MCPForm = forwardRef<MCPFormHandle, MCPFormProps>(function MCPForm(
<RuntimePanel
mcpTesting={mcpTesting}
runtimeInfo={runtimeInfo}
serverName={persistedServerName}
serverName={form.getValues('name')}
content="resources"
t={t}
/>
</TabsContent>
<TabsContent value="logs" className="mt-4 min-h-0 flex-1 overflow-y-auto">
{persistedServerName && <MCPLogs serverName={persistedServerName} />}
</TabsContent>
</Tabs>
) : (
runtimePanel
@@ -1,149 +0,0 @@
import { useCallback, useEffect, useRef, useState } from 'react';
import { httpClient } from '@/app/infra/http/HttpClient';
import { useTranslation } from 'react-i18next';
import { PluginLogEntry } from '@/app/infra/entities/plugin';
import { Button } from '@/components/ui/button';
import { Switch } from '@/components/ui/switch';
import { Label } from '@/components/ui/label';
import {
Select,
SelectContent,
SelectItem,
SelectTrigger,
SelectValue,
} from '@/components/ui/select';
import { RefreshCw } from 'lucide-react';
const LEVEL_OPTIONS = ['ALL', 'DEBUG', 'INFO', 'WARNING', 'ERROR'] as const;
function levelClassName(level: string): string {
switch (level) {
case 'ERROR':
case 'CRITICAL':
return 'text-red-500';
case 'WARNING':
return 'text-amber-500';
case 'DEBUG':
return 'text-gray-400 dark:text-gray-500';
default:
return 'text-gray-700 dark:text-gray-300';
}
}
export default function MCPLogs({ serverName }: { serverName: string }) {
const { t } = useTranslation();
const [logs, setLogs] = useState<PluginLogEntry[]>([]);
const [isLoading, setIsLoading] = useState(false);
const [level, setLevel] = useState<string>('ALL');
const [autoRefresh, setAutoRefresh] = useState(true);
const scrollRef = useRef<HTMLDivElement>(null);
const atBottomRef = useRef(true);
const fetchLogs = useCallback(() => {
setIsLoading(true);
httpClient
.getMcpServerLogs(serverName, 500, level === 'ALL' ? undefined : level)
.then((res) => {
setLogs(res.logs ?? []);
})
.catch(() => {
setLogs([]);
})
.finally(() => {
setIsLoading(false);
});
}, [serverName, level]);
useEffect(() => {
fetchLogs();
}, [fetchLogs]);
// Auto-refresh poll loop.
useEffect(() => {
if (!autoRefresh) return;
const timer = setInterval(fetchLogs, 3000);
return () => clearInterval(timer);
}, [autoRefresh, fetchLogs]);
// Keep view pinned to bottom when the user is already at the bottom.
useEffect(() => {
const el = scrollRef.current;
if (el && atBottomRef.current) {
el.scrollTop = el.scrollHeight;
}
}, [logs]);
function handleScroll() {
const el = scrollRef.current;
if (!el) return;
atBottomRef.current = el.scrollHeight - el.scrollTop - el.clientHeight < 40;
}
return (
<div className="flex h-full flex-col">
<div className="flex shrink-0 flex-wrap items-center gap-2 px-1 pb-3 sm:px-6">
<Select value={level} onValueChange={setLevel}>
<SelectTrigger className="h-8 w-[130px]">
<SelectValue />
</SelectTrigger>
<SelectContent>
{LEVEL_OPTIONS.map((opt) => (
<SelectItem key={opt} value={opt}>
{opt === 'ALL' ? t('mcp.logsLevelAll') : opt}
</SelectItem>
))}
</SelectContent>
</Select>
<Button
type="button"
variant="outline"
size="sm"
className="h-8"
onClick={fetchLogs}
disabled={isLoading}
>
<RefreshCw
className={`mr-1.5 size-3.5 ${isLoading ? 'animate-spin' : ''}`}
/>
{t('mcp.logsRefresh')}
</Button>
<div className="flex items-center gap-2">
<Switch
id="mcp-logs-auto-refresh"
checked={autoRefresh}
onCheckedChange={setAutoRefresh}
/>
<Label
htmlFor="mcp-logs-auto-refresh"
className="cursor-pointer text-sm font-normal text-muted-foreground"
>
{t('mcp.logsAutoRefresh')}
</Label>
</div>
</div>
<div
ref={scrollRef}
onScroll={handleScroll}
className="min-h-0 flex-1 overflow-auto bg-gray-50 px-3 py-3 font-mono text-xs leading-relaxed dark:bg-gray-900/40 sm:px-6"
>
{logs.length === 0 ? (
<div className="py-8 text-center text-sm text-gray-500 dark:text-gray-400">
{t('mcp.logsEmpty')}
</div>
) : (
logs.map((entry, idx) => (
<div
key={`${entry.ts}-${idx}`}
className={`whitespace-pre-wrap break-all ${levelClassName(
entry.level,
)}`}
>
{entry.text}
</div>
))
)}
</div>
</div>
);
}
@@ -1,649 +0,0 @@
import React from 'react';
import { useTranslation } from 'react-i18next';
import {
AlertCircle,
Bot,
ChevronDown,
ChevronRight,
Clock,
Cpu,
Hash,
User,
Wrench,
} from 'lucide-react';
import { cn } from '@/lib/utils';
import { MessageContentRenderer } from './MessageContentRenderer';
import {
ConversationTurn,
hasRenderableMessageContent,
} from '../utils/conversationTurns';
import { MonitoringMessage } from '../types/monitoring';
interface ConversationTurnListProps {
turns: ConversationTurn[];
expandedTurnId: string | null;
onToggleTurn: (turnId: string) => void;
}
function shortId(id?: string) {
if (!id) return '-';
if (id.length <= 12) return id;
return `${id.slice(0, 8)}...${id.slice(-4)}`;
}
function formatDuration(ms: number) {
if (!ms) return '0ms';
if (ms < 1000) return `${ms}ms`;
return `${(ms / 1000).toFixed(2)}s`;
}
function truncateDetail(value?: string) {
if (!value) return '';
return value.length > 1200 ? `${value.slice(0, 1200)}...` : value;
}
function roleLabel(message: MonitoringMessage | undefined) {
const role = message?.role?.toLowerCase();
if (role === 'assistant') return 'assistant';
if (role === 'user') return 'user';
return 'message';
}
function statusClass(level: ConversationTurn['level']) {
if (level === 'error') {
return 'border-red-200 bg-red-50 text-red-700 dark:border-red-900 dark:bg-red-950/40 dark:text-red-300';
}
if (level === 'warning') {
return 'border-yellow-200 bg-yellow-50 text-yellow-700 dark:border-yellow-900 dark:bg-yellow-950/40 dark:text-yellow-300';
}
return 'border-green-200 bg-green-50 text-green-700 dark:border-green-900 dark:bg-green-950/40 dark:text-green-300';
}
function Metric({
icon,
label,
tone = 'default',
}: {
icon: React.ReactNode;
label: string;
tone?: 'default' | 'error';
}) {
return (
<span
className={cn(
'inline-flex h-7 items-center gap-1.5 rounded-md border px-2 text-xs font-medium',
tone === 'error'
? 'border-red-200 bg-red-50 text-red-700 dark:border-red-900 dark:bg-red-950/40 dark:text-red-300'
: 'border-border bg-background text-muted-foreground',
)}
>
{icon}
{label}
</span>
);
}
function MetaItem({ label, value }: { label: string; value?: string }) {
return (
<div className="min-w-0 rounded-md bg-background px-3 py-2">
<div className="text-xs text-muted-foreground">{label}</div>
<div className="truncate text-sm font-medium text-foreground">
{value || '-'}
</div>
</div>
);
}
function MessageLane({
label,
icon,
content,
empty,
maxLines,
}: {
label: string;
icon: React.ReactNode;
content?: string;
empty: string;
maxLines: number;
}) {
return (
<div className="grid grid-cols-[5.25rem_minmax(0,1fr)] items-start gap-3 text-sm sm:grid-cols-[6rem_minmax(0,1fr)]">
<div className="flex h-7 items-center gap-1.5 text-xs font-medium text-muted-foreground">
{icon}
<span>{label}</span>
</div>
<div className="min-w-0 rounded-md bg-muted/45 px-3 py-2 text-foreground">
{content && hasRenderableMessageContent(content) ? (
<MessageContentRenderer content={content} maxLines={maxLines} />
) : (
<span className="italic text-muted-foreground">{empty}</span>
)}
</div>
</div>
);
}
function ExpandedMessage({
message,
label,
}: {
message: MonitoringMessage;
label: string;
}) {
return (
<div className="border-t border-border/70 py-3 first:border-t-0 first:pt-0 last:pb-0">
<div className="mb-2 flex flex-wrap items-center gap-2 text-xs text-muted-foreground">
<span className="rounded-md bg-muted px-2 py-1 font-medium text-foreground">
{label}
</span>
<span>{message.timestamp.toLocaleString()}</span>
<span className="font-mono">ID: {shortId(message.id)}</span>
</div>
<div className="text-sm leading-6 text-foreground">
<MessageContentRenderer content={message.messageContent} maxLines={4} />
</div>
</div>
);
}
export function ConversationTurnList({
turns,
expandedTurnId,
onToggleTurn,
}: ConversationTurnListProps) {
const { t } = useTranslation();
const [expandedToolCallIds, setExpandedToolCallIds] = React.useState<
Record<string, boolean>
>({});
const toggleToolCallDetails = (toolCallKey: string) => {
setExpandedToolCallIds((previous) => ({
...previous,
[toolCallKey]: !previous[toolCallKey],
}));
};
return (
<div className="space-y-4">
<div className="flex items-center justify-between text-sm text-muted-foreground">
<span className="font-medium text-foreground">
{t('monitoring.messageList.turns', {
defaultValue: '{{count}} 轮对话',
count: turns.length,
})}
</span>
</div>
{turns.map((turn) => {
const expanded = expandedTurnId === turn.id;
const firstAssistant = turn.assistantMessages[0];
const assistantOverflow = Math.max(
turn.assistantMessages.length - 1,
0,
);
return (
<div
key={turn.id}
className={cn(
'overflow-hidden rounded-xl border bg-card transition-colors',
turn.level === 'error' && 'border-red-200 dark:border-red-900',
)}
>
<div
role="button"
tabIndex={0}
className="cursor-pointer p-3 outline-none transition-colors hover:bg-accent/60 focus-visible:ring-2 focus-visible:ring-ring sm:p-5"
onClick={() => onToggleTurn(turn.id)}
onKeyDown={(event) => {
if (event.key === 'Enter' || event.key === ' ') {
event.preventDefault();
onToggleTurn(turn.id);
}
}}
>
<div className="flex flex-col gap-4 lg:flex-row lg:items-start lg:justify-between">
<div className="min-w-0 flex-1">
<div className="mb-2 flex min-w-0 items-center gap-2">
{expanded ? (
<ChevronDown className="h-5 w-5 shrink-0 text-muted-foreground" />
) : (
<ChevronRight className="h-5 w-5 shrink-0 text-muted-foreground" />
)}
<span className="truncate font-mono text-xs text-muted-foreground">
Turn: {shortId(turn.id)}
</span>
</div>
<div className="mb-3 flex min-w-0 flex-wrap items-center gap-2">
<span className="truncate text-sm font-medium text-foreground">
{turn.botName}
</span>
<span className="text-muted-foreground"></span>
<span className="truncate text-sm text-muted-foreground">
{turn.pipelineName}
</span>
{turn.runnerName && (
<>
<span className="text-muted-foreground"></span>
<span className="truncate text-sm text-muted-foreground">
{turn.runnerName}
</span>
</>
)}
</div>
<div className="space-y-2">
<MessageLane
label={t('monitoring.messageList.userMessage', {
defaultValue: '用户',
})}
icon={<User className="h-3.5 w-3.5" />}
content={turn.userMessage?.messageContent}
empty={t('monitoring.messageList.noUserMessage', {
defaultValue: '未记录用户输入',
})}
maxLines={2}
/>
<MessageLane
label={
assistantOverflow > 0
? t('monitoring.messageList.assistantMessageCount', {
defaultValue: '助手 +{{count}}',
count: assistantOverflow,
})
: t('monitoring.messageList.assistantMessage', {
defaultValue: '助手',
})
}
icon={<Bot className="h-3.5 w-3.5" />}
content={firstAssistant?.messageContent}
empty={t('monitoring.messageList.noAssistantMessage', {
defaultValue: '未记录助手回复',
})}
maxLines={2}
/>
</div>
</div>
<div className="flex shrink-0 flex-col gap-2 lg:items-end">
<div className="text-xs text-muted-foreground">
{turn.lastActivityAt.toLocaleString()}
</div>
<div
className={cn(
'inline-flex h-7 items-center rounded-md border px-2 text-xs font-medium',
statusClass(turn.level),
)}
>
{turn.level}
</div>
<div className="flex flex-wrap gap-2 lg:justify-end">
<Metric
icon={<Cpu className="h-3.5 w-3.5" />}
label={`${turn.llmCalls.length} LLM`}
/>
{turn.toolCalls.length > 0 && (
<Metric
icon={<Wrench className="h-3.5 w-3.5" />}
label={`${turn.toolCalls.length} tools`}
/>
)}
<Metric
icon={<Hash className="h-3.5 w-3.5" />}
label={`${turn.totalTokens.toLocaleString()} tokens`}
/>
<Metric
icon={<Clock className="h-3.5 w-3.5" />}
label={formatDuration(turn.totalDuration)}
/>
{turn.errors.length > 0 && (
<Metric
icon={<AlertCircle className="h-3.5 w-3.5" />}
label={`${turn.errors.length} errors`}
tone="error"
/>
)}
</div>
</div>
</div>
</div>
{expanded && (
<div className="border-t bg-muted/40 p-3 sm:p-5">
<div className="space-y-5 border-l-2 border-border pl-4 sm:pl-6">
<div className="grid grid-cols-2 gap-2 lg:grid-cols-5">
<MetaItem
label={t('monitoring.messageList.platform', {
defaultValue: '平台',
})}
value={turn.platform}
/>
<MetaItem
label={t('monitoring.messageList.user', {
defaultValue: '用户',
})}
value={turn.userName || turn.userId}
/>
<MetaItem
label={t('monitoring.messageList.runner', {
defaultValue: '执行器',
})}
value={turn.runnerName}
/>
<MetaItem
label={t('monitoring.sessions.sessionId', {
defaultValue: '会话 ID',
})}
value={turn.sessionId}
/>
<MetaItem
label={t('monitoring.messageList.messageCount', {
defaultValue: '消息数',
})}
value={String(turn.messages.length)}
/>
</div>
<section>
<h4 className="mb-3 flex items-center gap-2 text-sm font-semibold text-foreground">
<Bot className="h-4 w-4" />
{t('monitoring.messageList.conversationTrace', {
defaultValue: '消息链路',
})}
</h4>
<div className="rounded-lg bg-background px-3 py-3">
{turn.messages.map((message) => (
<ExpandedMessage
key={message.id}
message={message}
label={t(
`monitoring.messageList.roles.${roleLabel(message)}`,
{
defaultValue:
roleLabel(message) === 'assistant'
? '助手'
: roleLabel(message) === 'user'
? '用户'
: '消息',
},
)}
/>
))}
</div>
</section>
<section>
<h4 className="mb-3 flex items-center gap-2 text-sm font-semibold text-foreground">
<Cpu className="h-4 w-4" />
{t('monitoring.llmCalls.title', {
defaultValue: 'LLM 调用',
})}{' '}
({turn.llmCalls.length})
</h4>
<div className="grid grid-cols-3 gap-2">
<MetaItem
label={t('monitoring.llmCalls.totalTokens', {
defaultValue: '总 Token',
})}
value={turn.totalTokens.toLocaleString()}
/>
<MetaItem
label={t('monitoring.llmCalls.inputTokens', {
defaultValue: '输入 Token',
})}
value={turn.inputTokens.toLocaleString()}
/>
<MetaItem
label={t('monitoring.llmCalls.duration', {
defaultValue: '耗时',
})}
value={formatDuration(turn.totalDuration)}
/>
</div>
<div className="mt-3 rounded-lg bg-background px-3 py-3">
{turn.llmCalls.length > 0 ? (
turn.llmCalls.map((call, index) => (
<div
key={call.id}
className="border-t border-border/70 py-3 first:border-t-0 first:pt-0 last:pb-0"
>
<div className="mb-2 flex flex-wrap items-center justify-between gap-2">
<div className="flex min-w-0 flex-wrap items-center gap-2">
<span className="text-sm font-medium text-foreground">
#{index + 1} {call.modelName}
</span>
<span
className={cn(
'rounded-md px-2 py-1 text-xs font-medium',
call.status === 'success'
? 'bg-green-100 text-green-700 dark:bg-green-950 dark:text-green-300'
: 'bg-red-100 text-red-700 dark:bg-red-950 dark:text-red-300',
)}
>
{call.status}
</span>
</div>
<span className="text-xs text-muted-foreground">
{formatDuration(call.duration)}
</span>
</div>
<div className="flex flex-wrap gap-x-8 gap-y-1 text-xs text-muted-foreground">
<span>In: {call.tokens.input}</span>
<span>Out: {call.tokens.output}</span>
<span>Total: {call.tokens.total}</span>
<span className="font-mono">
ID: {shortId(call.id)}
</span>
</div>
{call.errorMessage && (
<div className="mt-2 whitespace-pre-wrap break-words text-xs text-red-600 dark:text-red-400">
{call.errorMessage}
</div>
)}
</div>
))
) : (
<div className="py-4 text-center text-sm text-muted-foreground">
{t('monitoring.messageList.noLlmCalls', {
defaultValue: '未记录模型调用',
})}
</div>
)}
</div>
</section>
<section>
<h4 className="mb-3 flex items-center gap-2 text-sm font-semibold text-foreground">
<Wrench className="h-4 w-4" />
{t('monitoring.toolCalls.title', {
defaultValue: '工具调用',
})}{' '}
({turn.toolCalls.length})
</h4>
<div className="grid grid-cols-2 gap-2 lg:grid-cols-3">
<MetaItem
label={t('monitoring.toolCalls.totalCalls', {
defaultValue: '调用次数',
})}
value={String(turn.toolCalls.length)}
/>
<MetaItem
label={t('monitoring.toolCalls.duration', {
defaultValue: '工具耗时',
})}
value={formatDuration(turn.totalToolDuration)}
/>
<MetaItem
label={t('monitoring.toolCalls.errorCalls', {
defaultValue: '失败次数',
})}
value={String(
turn.toolCalls.filter(
(call) => call.status === 'error',
).length,
)}
/>
</div>
<div className="mt-3 rounded-lg bg-background px-3 py-3">
{turn.toolCalls.length > 0 ? (
turn.toolCalls.map((call, index) => {
const toolCallKey = `${turn.id}:${call.id}`;
const hasToolDetails = Boolean(
call.arguments || call.result || call.errorMessage,
);
const expandedToolCall = Boolean(
expandedToolCallIds[toolCallKey],
);
const detailsId = `monitoring-tool-call-details-${call.id}`;
return (
<div
key={call.id}
className="border-t border-border/70 py-2 first:border-t-0 first:pt-0 last:pb-0"
>
<button
type="button"
className={cn(
'flex w-full items-start justify-between gap-3 rounded-md px-2 py-2 text-left outline-none transition-colors',
hasToolDetails &&
'cursor-pointer hover:bg-muted/60 focus-visible:ring-2 focus-visible:ring-ring',
)}
aria-expanded={
hasToolDetails ? expandedToolCall : undefined
}
aria-controls={
hasToolDetails ? detailsId : undefined
}
aria-disabled={!hasToolDetails}
onClick={() =>
hasToolDetails &&
toggleToolCallDetails(toolCallKey)
}
>
<div className="flex min-w-0 flex-wrap items-center gap-2">
{hasToolDetails &&
(expandedToolCall ? (
<ChevronDown className="mt-0.5 h-4 w-4 shrink-0 text-muted-foreground" />
) : (
<ChevronRight className="mt-0.5 h-4 w-4 shrink-0 text-muted-foreground" />
))}
<span className="min-w-0 truncate text-sm font-medium text-foreground">
#{index + 1} {call.toolName}
</span>
<span className="rounded-md bg-muted px-2 py-1 text-xs font-medium text-muted-foreground">
{call.toolSource}
</span>
<span
className={cn(
'rounded-md px-2 py-1 text-xs font-medium',
call.status === 'success'
? 'bg-green-100 text-green-700 dark:bg-green-950 dark:text-green-300'
: 'bg-red-100 text-red-700 dark:bg-red-950 dark:text-red-300',
)}
>
{call.status}
</span>
<span className="font-mono text-xs text-muted-foreground">
ID: {shortId(call.id)}
</span>
</div>
<span className="shrink-0 text-xs text-muted-foreground">
{formatDuration(call.duration)}
</span>
</button>
{hasToolDetails && expandedToolCall && (
<div
id={detailsId}
className="mt-1 grid gap-2 px-2 pb-2 text-xs lg:grid-cols-2"
>
{call.arguments && (
<div className="min-w-0 rounded-md bg-muted/50 p-2">
<div className="mb-1 font-medium text-foreground">
{t('monitoring.toolCalls.arguments', {
defaultValue: '参数',
})}
</div>
<pre className="whitespace-pre-wrap break-words font-mono text-muted-foreground">
{truncateDetail(call.arguments)}
</pre>
</div>
)}
{call.result && (
<div className="min-w-0 rounded-md bg-muted/50 p-2">
<div className="mb-1 font-medium text-foreground">
{t('monitoring.toolCalls.result', {
defaultValue: '结果',
})}
</div>
<pre className="whitespace-pre-wrap break-words font-mono text-muted-foreground">
{truncateDetail(call.result)}
</pre>
</div>
)}
{call.errorMessage && (
<div className="min-w-0 whitespace-pre-wrap break-words rounded-md bg-red-50 p-2 text-red-600 dark:bg-red-950/40 dark:text-red-400 lg:col-span-2">
{call.errorMessage}
</div>
)}
</div>
)}
</div>
);
})
) : (
<div className="py-4 text-center text-sm text-muted-foreground">
{t('monitoring.toolCalls.noToolCalls', {
defaultValue: '未记录工具调用',
})}
</div>
)}
</div>
</section>
{turn.errors.length > 0 && (
<section>
<h4 className="mb-3 flex items-center gap-2 text-sm font-semibold text-red-700 dark:text-red-300">
<AlertCircle className="h-4 w-4" />
{t('monitoring.errors.title', {
defaultValue: '错误日志',
})}{' '}
({turn.errors.length})
</h4>
<div className="rounded-lg bg-background px-3 py-3">
{turn.errors.map((error) => (
<div
key={error.id}
className="border-t border-red-200/80 py-3 first:border-t-0 first:pt-0 last:pb-0 dark:border-red-900"
>
<div className="mb-2 flex flex-wrap items-center justify-between gap-2">
<span className="text-sm font-medium text-red-700 dark:text-red-300">
{error.errorType}
</span>
<span className="text-xs text-muted-foreground">
{error.timestamp.toLocaleString()}
</span>
</div>
<div className="whitespace-pre-wrap break-words text-sm text-red-600 dark:text-red-400">
{error.errorMessage}
</div>
</div>
))}
</div>
</section>
)}
</div>
</div>
)}
</div>
);
})}
</div>
);
}
@@ -106,9 +106,6 @@ export function useMonitoringData(filterState: FilterState) {
const llmCalls = Array.isArray(response?.llmCalls)
? response.llmCalls
: [];
const toolCalls = Array.isArray(response?.toolCalls)
? response.toolCalls
: [];
const embeddingCalls = Array.isArray(response?.embeddingCalls)
? response.embeddingCalls
: [];
@@ -119,7 +116,6 @@ export function useMonitoringData(filterState: FilterState) {
const totalCount = response?.totalCount ?? {
messages: messages.length,
llmCalls: llmCalls.length,
toolCalls: toolCalls.length,
embeddingCalls: embeddingCalls.length,
sessions: sessions.length,
errors: errors.length,
@@ -149,10 +145,8 @@ export function useMonitoringData(filterState: FilterState) {
level: string;
platform?: string;
user_id?: string;
user_name?: string;
runner_name?: string;
variables?: string;
role?: string;
}) => ({
id: msg.id,
timestamp: parseUTCTimestamp(msg.timestamp),
@@ -166,10 +160,8 @@ export function useMonitoringData(filterState: FilterState) {
level: msg.level as 'info' | 'warning' | 'error' | 'debug',
platform: msg.platform,
userId: msg.user_id,
userName: msg.user_name,
runnerName: msg.runner_name,
variables: msg.variables,
role: msg.role,
}),
),
llmCalls: llmCalls.map(
@@ -187,7 +179,6 @@ export function useMonitoringData(filterState: FilterState) {
bot_name: string;
pipeline_id: string;
pipeline_name: string;
session_id?: string;
error_message?: string;
message_id?: string;
}) => ({
@@ -206,46 +197,10 @@ export function useMonitoringData(filterState: FilterState) {
botName: call.bot_name,
pipelineId: call.pipeline_id,
pipelineName: call.pipeline_name,
sessionId: call.session_id,
errorMessage: call.error_message,
messageId: call.message_id,
}),
),
toolCalls: toolCalls.map(
(call: {
id: string;
timestamp: string;
tool_name: string;
tool_source: string;
duration: number;
status: string;
bot_id: string;
bot_name: string;
pipeline_id: string;
pipeline_name: string;
session_id?: string;
message_id?: string;
arguments?: string;
result?: string;
error_message?: string;
}) => ({
id: call.id,
timestamp: parseUTCTimestamp(call.timestamp),
toolName: call.tool_name,
toolSource: call.tool_source,
duration: call.duration,
status: call.status as 'success' | 'error',
botId: call.bot_id,
botName: call.bot_name,
pipelineId: call.pipeline_id,
pipelineName: call.pipeline_name,
sessionId: call.session_id,
messageId: call.message_id,
arguments: call.arguments,
result: call.result,
errorMessage: call.error_message,
}),
),
embeddingCalls: embeddingCalls.map(
(call: {
id: string;
@@ -339,7 +294,6 @@ export function useMonitoringData(filterState: FilterState) {
totalCount: {
messages: totalCount.messages,
llmCalls: totalCount.llmCalls,
toolCalls: totalCount.toolCalls ?? toolCalls.length,
embeddingCalls: totalCount.embeddingCalls || 0,
sessions: totalCount.sessions,
errors: totalCount.errors,
@@ -363,7 +317,6 @@ export function useMonitoringData(filterState: FilterState) {
botName: call.botName,
pipelineId: call.pipelineId,
pipelineName: call.pipelineName,
sessionId: call.sessionId,
}),
);
+264 -33
View File
@@ -18,12 +18,60 @@ import { ExportDropdown } from './components/ExportDropdown';
import { useMonitoringFilters } from './hooks/useMonitoringFilters';
import { useMonitoringData } from './hooks/useMonitoringData';
import { useFeedbackData } from './hooks/useFeedbackData';
import { ConversationTurnList } from './components/ConversationTurnList';
import { MessageDetailsCard } from './components/MessageDetailsCard';
import { MessageContentRenderer } from './components/MessageContentRenderer';
import { FeedbackStatsCards } from './components/FeedbackCard';
import { FeedbackList } from './components/FeedbackList';
import { buildConversationTurns } from './utils/conversationTurns';
import { MessageDetails } from './types/monitoring';
import { httpClient } from '@/app/infra/http/HttpClient';
import { LoadingSpinner, LoadingPage } from '@/components/ui/loading-spinner';
interface RawMessageData {
id: string;
timestamp: string;
bot_id: string;
bot_name: string;
pipeline_id: string;
pipeline_name: string;
message_content: string;
session_id: string;
status: string;
level: string;
platform: string;
user_id: string;
runner_name: string;
variables: Record<string, unknown>;
}
interface RawLLMCallData {
id: string;
timestamp: string;
model_name: string;
status: string;
duration: number;
error_message: string | null;
input_tokens: number;
output_tokens: number;
total_tokens: number;
}
interface RawLLMStatsData {
total_calls: number;
total_input_tokens: number;
total_output_tokens: number;
total_tokens: number;
total_duration_ms: number;
average_duration_ms: number;
}
interface RawErrorData {
id: string;
timestamp: string;
error_type: string;
error_message: string;
stack_trace: string | null;
}
function MonitoringPageContent() {
const { t } = useTranslation();
const { filterState, setSelectedBots, setSelectedPipelines, setTimeRange } =
@@ -98,37 +146,115 @@ function MonitoringPageContent() {
setFeedbackRefreshKey((k) => k + 1);
}, [refetch]);
const conversationTurns = useMemo(
() =>
buildConversationTurns(
data?.messages || [],
data?.llmCalls || [],
data?.errors || [],
data?.toolCalls || [],
),
[data?.messages, data?.llmCalls, data?.errors, data?.toolCalls],
const [expandedMessageId, setExpandedMessageId] = useState<string | null>(
null,
);
const [messageDetails, setMessageDetails] = useState<
Record<string, MessageDetails>
>({});
const [loadingDetails, setLoadingDetails] = useState<Record<string, boolean>>(
{},
);
// State for expanded errors
const [expandedErrorId, setExpandedErrorId] = useState<string | null>(null);
const [expandedTurnId, setExpandedTurnId] = useState<string | null>(null);
// State for controlled tabs
const [activeTab, setActiveTab] = useState<string>('messages');
// Function to jump to a message record
const jumpToMessage = (messageId: string) => {
const jumpToMessage = async (messageId: string) => {
setActiveTab('messages');
// Small delay to ensure tab switch completes
setTimeout(() => {
const turn = conversationTurns.find((item) =>
item.messages.some((message) => message.id === messageId),
);
setExpandedTurnId(turn?.id ?? messageId);
toggleMessageExpand(messageId);
}, 100);
};
const toggleTurnExpand = (turnId: string) => {
setExpandedTurnId((current) => (current === turnId ? null : turnId));
const toggleMessageExpand = async (messageId: string) => {
if (expandedMessageId === messageId) {
// Collapse
setExpandedMessageId(null);
} else {
// Expand
setExpandedMessageId(messageId);
// Fetch details if not already loaded
if (!messageDetails[messageId]) {
setLoadingDetails({ ...loadingDetails, [messageId]: true });
try {
// httpClient.get() returns the inner data directly (response.data.data)
const result = await httpClient.get<{
message_id: string;
found: boolean;
message: RawMessageData | null;
llm_calls: RawLLMCallData[];
llm_stats: RawLLMStatsData;
errors: RawErrorData[];
}>(`/api/v1/monitoring/messages/${messageId}/details`);
if (result) {
setMessageDetails((prev) => ({
...prev,
[messageId]: {
messageId: result.message_id,
found: result.found,
message: result.message
? {
id: result.message.id,
timestamp: new Date(result.message.timestamp),
botId: result.message.bot_id,
botName: result.message.bot_name,
pipelineId: result.message.pipeline_id,
pipelineName: result.message.pipeline_name,
messageContent: result.message.message_content,
sessionId: result.message.session_id,
status: result.message.status,
level: result.message.level,
platform: result.message.platform,
userId: result.message.user_id,
runnerName: result.message.runner_name,
variables: result.message.variables,
}
: undefined,
llmCalls: result.llm_calls.map((call: RawLLMCallData) => ({
id: call.id,
timestamp: new Date(call.timestamp),
modelName: call.model_name,
status: call.status,
duration: call.duration,
errorMessage: call.error_message,
tokens: {
input: call.input_tokens || 0,
output: call.output_tokens || 0,
total: call.total_tokens || 0,
},
})),
errors: result.errors.map((error: RawErrorData) => ({
id: error.id,
timestamp: new Date(error.timestamp),
errorType: error.error_type,
errorMessage: error.error_message,
stackTrace: error.stack_trace,
})),
llmStats: {
totalCalls: result.llm_stats.total_calls,
totalInputTokens: result.llm_stats.total_input_tokens,
totalOutputTokens: result.llm_stats.total_output_tokens,
totalTokens: result.llm_stats.total_tokens,
totalDurationMs: result.llm_stats.total_duration_ms,
averageDurationMs: result.llm_stats.average_duration_ms,
},
} as MessageDetails,
}));
}
} catch (error) {
console.error('Failed to fetch message details:', error);
} finally {
setLoadingDetails({ ...loadingDetails, [messageId]: false });
}
}
}
};
const toggleErrorExpand = (errorId: string) => {
@@ -216,22 +342,127 @@ function MonitoringPageContent() {
</div>
)}
{!loading && data && conversationTurns.length > 0 && (
<ConversationTurnList
turns={conversationTurns}
expandedTurnId={expandedTurnId}
onToggleTurn={toggleTurnExpand}
/>
)}
{!loading &&
data &&
data.messages &&
data.messages.length > 0 && (
<div className="space-y-4">
{data.messages
.filter((msg) => {
// Filter out messages with empty content
const content = msg.messageContent?.trim();
return (
content && content !== '[]' && content !== '""'
);
})
.map((msg) => (
<div
key={msg.id}
className="border rounded-xl overflow-hidden transition-all duration-200"
>
{/* Message Header - Always Visible */}
<div
className="p-3 cursor-pointer hover:bg-accent transition-colors sm:p-5"
onClick={() => toggleMessageExpand(msg.id)}
>
<div className="flex items-start justify-between">
<div className="flex items-start flex-1">
{/* Expand Icon */}
<div className="mr-3 mt-0.5">
{expandedMessageId === msg.id ? (
<ChevronDown className="w-5 h-5 text-muted-foreground" />
) : (
<ChevronRight className="w-5 h-5 text-muted-foreground" />
)}
</div>
{!loading && (!data || conversationTurns.length === 0) && (
<div className="flex flex-col items-center justify-center text-muted-foreground py-16 gap-2">
<MessageSquare className="h-[3rem] w-[3rem]" />
<div className="text-sm">
{t('monitoring.messageList.noMessages')}
{/* Message Info */}
<div className="flex-1">
<div className="flex items-center gap-2 mb-1">
<span className="text-xs text-muted-foreground font-mono">
ID: {msg.id}
</span>
</div>
<div className="flex items-center gap-2 mb-2">
<span className="font-medium text-sm text-foreground">
{msg.botName}
</span>
<span className="text-muted-foreground">
</span>
<span className="text-sm text-muted-foreground">
{msg.pipelineName}
</span>
{msg.runnerName && (
<>
<span className="text-muted-foreground">
</span>
<span className="text-sm text-muted-foreground">
{msg.runnerName}
</span>
</>
)}
</div>
<div className="text-base text-foreground">
<MessageContentRenderer
content={msg.messageContent}
maxLines={3}
/>
</div>
</div>
</div>
{/* Status and Timestamp */}
<div className="flex flex-col items-end gap-2 ml-4">
<span className="text-xs text-muted-foreground whitespace-nowrap">
{msg.timestamp.toLocaleString()}
</span>
<span
className={`text-xs px-2 py-1 rounded ${
msg.level === 'error'
? 'bg-red-100 text-red-800 dark:bg-red-900 dark:text-red-200'
: msg.level === 'warning'
? 'bg-yellow-100 text-yellow-800 dark:bg-yellow-900 dark:text-yellow-200'
: 'bg-green-100 text-green-800 dark:bg-green-900 dark:text-green-200'
}`}
>
{msg.level}
</span>
</div>
</div>
</div>
{/* Expanded Details */}
{expandedMessageId === msg.id && (
<div className="border-t p-4 bg-muted">
{loadingDetails[msg.id] && (
<div className="py-4 flex justify-center">
<LoadingSpinner size="sm" text="" />
</div>
)}
{!loadingDetails[msg.id] &&
messageDetails[msg.id] && (
<MessageDetailsCard
details={messageDetails[msg.id]}
/>
)}
</div>
)}
</div>
))}
</div>
</div>
)}
)}
{!loading &&
(!data || !data.messages || data.messages.length === 0) && (
<div className="flex flex-col items-center justify-center text-muted-foreground py-16 gap-2">
<MessageSquare className="h-[3rem] w-[3rem]" />
<div className="text-sm">
{t('monitoring.messageList.noMessages')}
</div>
</div>
)}
</div>
</TabsContent>
@@ -11,10 +11,8 @@ export interface MonitoringMessage {
level: 'info' | 'warning' | 'error' | 'debug';
platform?: string;
userId?: string;
userName?: string;
runnerName?: string;
variables?: string;
role?: 'user' | 'assistant' | string;
}
export interface LLMCall {
@@ -33,29 +31,10 @@ export interface LLMCall {
botName: string;
pipelineId: string;
pipelineName: string;
sessionId?: string;
errorMessage?: string;
messageId?: string;
}
export interface ToolCall {
id: string;
timestamp: Date;
toolName: string;
toolSource: 'native' | 'plugin' | 'mcp' | 'skill' | string;
duration: number;
status: 'success' | 'error';
botId: string;
botName: string;
pipelineId: string;
pipelineName: string;
sessionId?: string;
messageId?: string;
arguments?: string;
result?: string;
errorMessage?: string;
}
export interface EmbeddingCall {
id: string;
timestamp: Date;
@@ -220,7 +199,6 @@ export interface MonitoringData {
overview: OverviewMetrics;
messages: MonitoringMessage[];
llmCalls: LLMCall[];
toolCalls: ToolCall[];
embeddingCalls: EmbeddingCall[];
modelCalls: ModelCall[];
sessions: SessionInfo[];
@@ -230,7 +208,6 @@ export interface MonitoringData {
totalCount: {
messages: number;
llmCalls: number;
toolCalls?: number;
embeddingCalls: number;
sessions: number;
errors: number;
@@ -1,294 +0,0 @@
import {
ErrorLog,
LLMCall,
MonitoringMessage,
ToolCall,
} from '../types/monitoring';
type MessageRole = 'user' | 'assistant' | 'unknown';
export interface ConversationTurn {
id: string;
sessionId: string;
startedAt: Date;
lastActivityAt: Date;
botId: string;
botName: string;
pipelineId: string;
pipelineName: string;
runnerName?: string;
platform?: string;
userId?: string;
userName?: string;
userMessage?: MonitoringMessage;
assistantMessages: MonitoringMessage[];
messages: MonitoringMessage[];
llmCalls: LLMCall[];
toolCalls: ToolCall[];
errors: ErrorLog[];
status: 'success' | 'error' | 'pending';
level: 'info' | 'warning' | 'error' | 'debug';
inputTokens: number;
outputTokens: number;
totalTokens: number;
totalDuration: number;
totalToolDuration: number;
}
function normalizeRole(
message: MonitoringMessage,
llmMessageIds: Set<string>,
): MessageRole {
const role = message.role?.toLowerCase();
if (role === 'user' || role === 'assistant') {
return role;
}
if (llmMessageIds.has(message.id)) {
return 'user';
}
return 'unknown';
}
export function hasRenderableMessageContent(content?: string): boolean {
const trimmed = content?.trim();
if (!trimmed || trimmed === '[]' || trimmed === '""') {
return false;
}
try {
const parsed = JSON.parse(trimmed);
if (typeof parsed === 'string') {
return parsed.trim().length > 0;
}
if (Array.isArray(parsed)) {
return parsed.some(
(component) =>
typeof component !== 'object' ||
component === null ||
component.type !== 'Source',
);
}
} catch {
return true;
}
return true;
}
function createTurn(message: MonitoringMessage): ConversationTurn {
return {
id: message.id,
sessionId: message.sessionId,
startedAt: message.timestamp,
lastActivityAt: message.timestamp,
botId: message.botId,
botName: message.botName,
pipelineId: message.pipelineId,
pipelineName: message.pipelineName,
runnerName: message.runnerName,
platform: message.platform,
userId: message.userId,
userName: message.userName,
assistantMessages: [],
messages: [],
llmCalls: [],
toolCalls: [],
errors: [],
status: message.status,
level: message.level,
inputTokens: 0,
outputTokens: 0,
totalTokens: 0,
totalDuration: 0,
totalToolDuration: 0,
};
}
function updateTurnActivity(turn: ConversationTurn, timestamp: Date) {
if (timestamp.getTime() > turn.lastActivityAt.getTime()) {
turn.lastActivityAt = timestamp;
}
}
function addMessageToTurn(
turn: ConversationTurn,
message: MonitoringMessage,
role: MessageRole,
) {
turn.messages.push(message);
updateTurnActivity(turn, message.timestamp);
if (message.level === 'error') {
turn.level = 'error';
} else if (message.level === 'warning' && turn.level !== 'error') {
turn.level = 'warning';
}
if (message.status === 'error') {
turn.status = 'error';
} else if (message.status === 'pending' && turn.status !== 'error') {
turn.status = 'pending';
}
if (role === 'assistant') {
turn.assistantMessages.push(message);
return;
}
if (!turn.userMessage) {
turn.userMessage = message;
turn.userId = message.userId ?? turn.userId;
turn.userName = message.userName ?? turn.userName;
return;
}
turn.assistantMessages.push(message);
}
function findTurnBySessionTime(
sessionTurns: Map<string, ConversationTurn[]>,
sessionId: string | undefined,
timestamp: Date,
): ConversationTurn | undefined {
if (!sessionId) {
return undefined;
}
const turns = sessionTurns.get(sessionId);
if (!turns?.length) {
return undefined;
}
let nearest = turns[0];
const targetTime = timestamp.getTime();
for (const turn of turns) {
if (turn.startedAt.getTime() <= targetTime) {
nearest = turn;
} else {
break;
}
}
return nearest;
}
export function buildConversationTurns(
messages: MonitoringMessage[],
llmCalls: LLMCall[],
errors: ErrorLog[],
toolCalls: ToolCall[] = [],
): ConversationTurn[] {
const activityMessageIds = new Set([
...llmCalls
.map((call) => call.messageId)
.filter((messageId): messageId is string => Boolean(messageId)),
...toolCalls
.map((call) => call.messageId)
.filter((messageId): messageId is string => Boolean(messageId)),
]);
const visibleMessages = messages
.filter((message) => hasRenderableMessageContent(message.messageContent))
.sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime());
const sessionTurns = new Map<string, ConversationTurn[]>();
const lastTurnBySession = new Map<string, ConversationTurn>();
const messageIdToTurn = new Map<string, ConversationTurn>();
for (const message of visibleMessages) {
const role = normalizeRole(message, activityMessageIds);
const previousTurn = lastTurnBySession.get(message.sessionId);
const shouldStartTurn = role === 'user' || !previousTurn;
const turn = shouldStartTurn ? createTurn(message) : previousTurn;
if (shouldStartTurn) {
const turns = sessionTurns.get(message.sessionId) ?? [];
turns.push(turn);
sessionTurns.set(message.sessionId, turns);
lastTurnBySession.set(message.sessionId, turn);
}
addMessageToTurn(turn, message, role);
messageIdToTurn.set(message.id, turn);
}
const allTurns = Array.from(sessionTurns.values()).flat();
for (const call of llmCalls) {
const turn =
(call.messageId ? messageIdToTurn.get(call.messageId) : undefined) ??
findTurnBySessionTime(sessionTurns, call.sessionId, call.timestamp);
if (!turn) {
continue;
}
turn.llmCalls.push(call);
turn.inputTokens += call.tokens.input;
turn.outputTokens += call.tokens.output;
turn.totalTokens += call.tokens.total;
turn.totalDuration += call.duration;
updateTurnActivity(turn, call.timestamp);
if (call.status === 'error') {
turn.status = 'error';
turn.level = 'error';
}
}
for (const call of toolCalls) {
const turn =
(call.messageId ? messageIdToTurn.get(call.messageId) : undefined) ??
findTurnBySessionTime(sessionTurns, call.sessionId, call.timestamp);
if (!turn) {
continue;
}
turn.toolCalls.push(call);
turn.totalToolDuration += call.duration;
updateTurnActivity(turn, call.timestamp);
if (call.status === 'error') {
turn.status = 'error';
turn.level = 'error';
}
}
for (const error of errors) {
const turn =
(error.messageId ? messageIdToTurn.get(error.messageId) : undefined) ??
findTurnBySessionTime(sessionTurns, error.sessionId, error.timestamp);
if (!turn) {
continue;
}
turn.errors.push(error);
turn.status = 'error';
turn.level = 'error';
updateTurnActivity(turn, error.timestamp);
}
for (const turn of allTurns) {
turn.messages.sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime());
turn.assistantMessages.sort(
(a, b) => a.timestamp.getTime() - b.timestamp.getTime(),
);
turn.llmCalls.sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime());
turn.toolCalls.sort(
(a, b) => a.timestamp.getTime() - b.timestamp.getTime(),
);
turn.errors.sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime());
}
return allTurns.sort(
(a, b) => b.lastActivityAt.getTime() - a.lastActivityAt.getTime(),
);
}
@@ -12,15 +12,63 @@ import {
Monitor,
} from 'lucide-react';
import { useMonitoringData } from '@/app/home/monitoring/hooks/useMonitoringData';
import { ConversationTurnList } from '@/app/home/monitoring/components/ConversationTurnList';
import { buildConversationTurns } from '@/app/home/monitoring/utils/conversationTurns';
import { MessageContentRenderer } from '@/app/home/monitoring/components/MessageContentRenderer';
import { LoadingSpinner } from '@/components/ui/loading-spinner';
import { httpClient } from '@/app/infra/http/HttpClient';
import { MessageDetails } from '@/app/home/monitoring/types/monitoring';
import { parseUTCTimestamp } from '@/app/home/monitoring/utils/dateUtils';
interface PipelineMonitoringTabProps {
pipelineId: string;
onNavigateToMonitoring?: () => void;
}
interface RawMessageData {
id: string;
timestamp: string;
bot_id: string;
bot_name: string;
pipeline_id: string;
pipeline_name: string;
message_content: string;
session_id: string;
status: string;
level: string;
platform: string;
user_id: string;
runner_name: string;
variables: Record<string, unknown>;
}
interface RawLLMCallData {
id: string;
timestamp: string;
model_name: string;
status: string;
duration: number;
error_message: string | null;
input_tokens: number;
output_tokens: number;
total_tokens: number;
}
interface RawLLMStatsData {
total_calls: number;
total_input_tokens: number;
total_output_tokens: number;
total_tokens: number;
total_duration_ms: number;
average_duration_ms: number;
}
interface RawErrorData {
id: string;
timestamp: string;
error_type: string;
error_message: string;
stack_trace: string | null;
}
export default function PipelineMonitoringTab({
pipelineId,
onNavigateToMonitoring,
@@ -40,24 +88,98 @@ export default function PipelineMonitoringTab({
const { data, loading, refetch } = useMonitoringData(filterState);
const conversationTurns = useMemo(
() =>
data
? buildConversationTurns(
data.messages,
data.llmCalls,
data.errors,
data.toolCalls,
)
: [],
[data],
const [expandedMessageId, setExpandedMessageId] = useState<string | null>(
null,
);
const [messageDetails, setMessageDetails] = useState<
Record<string, MessageDetails>
>({});
const [loadingDetails, setLoadingDetails] = useState<Record<string, boolean>>(
{},
);
const [expandedTurnId, setExpandedTurnId] = useState<string | null>(null);
const [expandedErrorId, setExpandedErrorId] = useState<string | null>(null);
const [activeTab, setActiveTab] = useState<string>('messages');
const toggleTurnExpand = (turnId: string) => {
setExpandedTurnId((current) => (current === turnId ? null : turnId));
const toggleMessageExpand = async (messageId: string) => {
if (expandedMessageId === messageId) {
setExpandedMessageId(null);
} else {
setExpandedMessageId(messageId);
if (!messageDetails[messageId]) {
setLoadingDetails((prev) => ({ ...prev, [messageId]: true }));
try {
const result = await httpClient.get<{
message_id: string;
found: boolean;
message: RawMessageData | null;
llm_calls: RawLLMCallData[];
llm_stats: RawLLMStatsData;
errors: RawErrorData[];
}>(`/api/v1/monitoring/messages/${messageId}/details`);
if (result) {
setMessageDetails((prev) => ({
...prev,
[messageId]: {
messageId: result.message_id,
found: result.found,
message: result.message
? {
id: result.message.id,
timestamp: parseUTCTimestamp(result.message.timestamp),
botId: result.message.bot_id,
botName: result.message.bot_name,
pipelineId: result.message.pipeline_id,
pipelineName: result.message.pipeline_name,
messageContent: result.message.message_content,
sessionId: result.message.session_id,
status: result.message.status,
level: result.message.level,
platform: result.message.platform,
userId: result.message.user_id,
runnerName: result.message.runner_name,
variables: result.message.variables,
}
: undefined,
llmCalls: result.llm_calls.map((call: RawLLMCallData) => ({
id: call.id,
timestamp: parseUTCTimestamp(call.timestamp),
modelName: call.model_name,
status: call.status,
duration: call.duration,
errorMessage: call.error_message,
tokens: {
input: call.input_tokens || 0,
output: call.output_tokens || 0,
total: call.total_tokens || 0,
},
})),
errors: result.errors.map((error: RawErrorData) => ({
id: error.id,
timestamp: parseUTCTimestamp(error.timestamp),
errorType: error.error_type,
errorMessage: error.error_message,
stackTrace: error.stack_trace,
})),
llmStats: {
totalCalls: result.llm_stats.total_calls,
totalInputTokens: result.llm_stats.total_input_tokens,
totalOutputTokens: result.llm_stats.total_output_tokens,
totalTokens: result.llm_stats.total_tokens,
totalDurationMs: result.llm_stats.total_duration_ms,
averageDurationMs: result.llm_stats.average_duration_ms,
},
} as MessageDetails,
}));
}
} catch (error) {
console.error('Failed to fetch message details:', error);
} finally {
setLoadingDetails((prev) => ({ ...prev, [messageId]: false }));
}
}
}
};
const toggleErrorExpand = (errorId: string) => {
@@ -68,16 +190,12 @@ export default function PipelineMonitoringTab({
}
};
const jumpToMessage = (messageId: string) => {
const jumpToMessage = async (messageId: string) => {
setActiveTab('messages');
const turn = conversationTurns.find((item) =>
item.messages.some((message) => message.id === messageId),
);
if (turn) {
setExpandedTurnId(turn.id);
}
// Small delay to ensure tab transition completes before expanding
setTimeout(() => {
toggleMessageExpand(messageId);
}, 100);
};
return (
@@ -177,22 +295,142 @@ export default function PipelineMonitoringTab({
</div>
)}
{!loading && data && conversationTurns.length > 0 && (
<ConversationTurnList
turns={conversationTurns}
expandedTurnId={expandedTurnId}
onToggleTurn={toggleTurnExpand}
/>
)}
{!loading && data && data.messages && data.messages.length > 0 && (
<div className="space-y-3">
{data.messages
.filter((msg) => {
const content = msg.messageContent?.trim();
return content && content !== '[]' && content !== '""';
})
.map((msg) => (
<div
key={msg.id}
className="border border-gray-200 dark:border-gray-700 rounded-lg overflow-hidden hover:shadow-md transition-all duration-200"
>
<div
className="p-4 cursor-pointer hover:bg-gray-50 dark:hover:bg-gray-800/50 transition-colors"
onClick={() => toggleMessageExpand(msg.id)}
>
<div className="flex items-start justify-between">
<div className="flex items-start flex-1">
<div className="mr-2 mt-0.5">
{expandedMessageId === msg.id ? (
<ChevronDown className="w-4 h-4 text-gray-500" />
) : (
<ChevronRight className="w-4 h-4 text-gray-500" />
)}
</div>
<div className="flex-1">
<div className="flex items-center gap-2 mb-1">
<span
className={`text-xs px-2 py-0.5 rounded ${
msg.status === 'success'
? 'bg-green-100 text-green-800 dark:bg-green-900 dark:text-green-200'
: msg.status === 'error'
? 'bg-red-100 text-red-800 dark:bg-red-900 dark:text-red-200'
: 'bg-yellow-100 text-yellow-800 dark:bg-yellow-900 dark:text-yellow-200'
}`}
>
{msg.status}
</span>
<span className="text-xs text-gray-500 dark:text-gray-400">
{msg.botName}
</span>
</div>
<div className="text-sm text-gray-700 dark:text-gray-300 line-clamp-2">
<MessageContentRenderer
content={msg.messageContent}
/>
</div>
</div>
</div>
<span className="text-xs text-gray-500 dark:text-gray-400 whitespace-nowrap ml-4">
{msg.timestamp.toLocaleString()}
</span>
</div>
</div>
{!loading && (!data || conversationTurns.length === 0) && (
<div className="text-center text-gray-500 dark:text-gray-400 py-16">
<MessageCircle className="w-16 h-16 mx-auto mb-4 text-gray-300 dark:text-gray-600" />
<p className="text-base font-medium">
{t('monitoring.messageList.noMessages')}
</p>
{expandedMessageId === msg.id && (
<div className="border-t border-gray-200 dark:border-gray-700 p-4 bg-gray-50 dark:bg-gray-900">
{loadingDetails[msg.id] && (
<div className="flex justify-center py-8">
<LoadingSpinner
text={t('monitoring.messageList.loading')}
/>
</div>
)}
{!loadingDetails[msg.id] &&
messageDetails[msg.id] && (
<div className="space-y-4">
{messageDetails[msg.id].errors.length > 0 && (
<div className="bg-red-50 dark:bg-red-900/20 rounded-lg p-3">
<h4 className="text-sm font-semibold text-red-700 dark:text-red-400 mb-2">
{t('monitoring.errors.errorMessage')}
</h4>
{messageDetails[msg.id].errors.map(
(error) => (
<div
key={error.id}
className="text-sm space-y-2"
>
<div className="text-red-600 dark:text-red-400">
{error.errorType}:{' '}
{error.errorMessage}
</div>
{error.stackTrace && (
<pre className="text-xs text-gray-600 dark:text-gray-400 overflow-auto max-h-40 bg-white dark:bg-gray-900 p-2 rounded whitespace-pre-wrap break-words">
{error.stackTrace}
</pre>
)}
</div>
),
)}
</div>
)}
{messageDetails[msg.id].llmCalls.length > 0 && (
<div className="bg-blue-50 dark:bg-blue-900/20 rounded-lg p-3">
<h4 className="text-sm font-semibold text-blue-700 dark:text-blue-400 mb-2">
{t('monitoring.tabs.modelCalls')} (
{messageDetails[msg.id].llmCalls.length})
</h4>
<div className="text-xs text-gray-600 dark:text-gray-400 space-y-1">
<div>
{t('monitoring.llmCalls.totalTokens')}:{' '}
{
messageDetails[msg.id].llmStats
.totalTokens
}
</div>
<div>
{t('monitoring.llmCalls.duration')}:{' '}
{messageDetails[
msg.id
].llmStats.totalDurationMs.toFixed(0)}
ms
</div>
</div>
</div>
)}
</div>
)}
</div>
)}
</div>
))}
</div>
)}
{!loading &&
(!data || !data.messages || data.messages.length === 0) && (
<div className="text-center text-gray-500 dark:text-gray-400 py-16">
<MessageCircle className="w-16 h-16 mx-auto mb-4 text-gray-300 dark:text-gray-600" />
<p className="text-base font-medium">
{t('monitoring.messageList.noMessages')}
</p>
</div>
)}
</TabsContent>
{/* Errors Tab */}
+23 -1
View File
@@ -179,6 +179,28 @@ export interface ApiRespPlatformBot {
bot: Bot;
}
export type BotAdapterConnectionStatus =
| 'connecting'
| 'connected'
| 'disconnected'
| 'error';
export interface BotAdapterRuntimeStatus {
connection_status?: BotAdapterConnectionStatus;
connection_error?: string;
last_connected_at?: number | null;
last_disconnected_at?: number | null;
}
export interface BotAdapterRuntimeValues {
bot_account_id?: string;
webhook_url?: string | null;
webhook_full_url?: string | null;
extra_webhook_full_url?: string | null;
runtime_status?: BotAdapterRuntimeStatus;
[key: string]: unknown;
}
export interface Bot {
uuid?: string;
name: string;
@@ -191,7 +213,7 @@ export interface Bot {
pipeline_routing_rules?: PipelineRoutingRule[];
created_at?: string;
updated_at?: string;
adapter_runtime_values?: object;
adapter_runtime_values?: BotAdapterRuntimeValues;
}
export type RoutingRuleOperator =
-36
View File
@@ -671,21 +671,6 @@ export class BackendClient extends BaseHttpClient {
);
}
public getMcpServerLogs(
serverName: string,
limit: number = 200,
level?: string,
): Promise<{ logs: PluginLogEntry[] }> {
const params = new URLSearchParams();
params.set('limit', String(limit));
if (level) {
params.set('level', level);
}
return this.get(
`/api/v1/mcp/servers/${encodeURIComponent(serverName)}/logs?${params.toString()}`,
);
}
public getPluginAssetURL(
author: string,
name: string,
@@ -1214,10 +1199,8 @@ export class BackendClient extends BaseHttpClient {
level: string;
platform?: string;
user_id?: string;
user_name?: string;
runner_name?: string;
variables?: string;
role?: string;
}>;
llmCalls: Array<{
id: string;
@@ -1233,27 +1216,9 @@ export class BackendClient extends BaseHttpClient {
bot_name: string;
pipeline_id: string;
pipeline_name: string;
session_id?: string;
error_message?: string;
message_id?: string;
}>;
toolCalls: Array<{
id: string;
timestamp: string;
tool_name: string;
tool_source: string;
duration: number;
status: string;
bot_id: string;
bot_name: string;
pipeline_id: string;
pipeline_name: string;
session_id?: string;
message_id?: string;
arguments?: string;
result?: string;
error_message?: string;
}>;
embeddingCalls: Array<{
id: string;
timestamp: string;
@@ -1298,7 +1263,6 @@ export class BackendClient extends BaseHttpClient {
totalCount: {
messages: number;
llmCalls: number;
toolCalls?: number;
embeddingCalls: number;
sessions: number;
errors: number;
+14 -29
View File
@@ -354,6 +354,10 @@ const enUS = {
log: 'Log',
configuration: 'Configuration',
logs: 'Logs',
runtimeConnected: 'Connected',
runtimeConnecting: 'Connecting',
runtimeDisconnected: 'Disconnected',
runtimeError: 'Connection error',
basicInfo: 'Basic Information',
basicInfoDescription: 'Set the bot name and description',
routingConnection: 'Routing & Connection',
@@ -828,12 +832,6 @@ const enUS = {
tabTools: 'Tools',
tabResources: 'Resources',
tabDocs: 'Docs',
tabLogs: 'Logs',
logsLevelAll: 'All levels',
logsRefresh: 'Refresh',
logsAutoRefresh: 'Auto refresh',
logsEmpty:
'No logs yet. Runtime logs from the MCP server will appear here.',
noReadme: 'No documentation available',
parseResultFailed: 'Failed to parse test result',
noResultReturned: 'Test returned no result',
@@ -1335,20 +1333,6 @@ const enUS = {
level: 'Level',
runner: 'Runner',
viewConversation: 'View Conversation',
turns: '{{count}} conversation turns',
userMessage: 'User',
noUserMessage: 'No user input recorded',
assistantMessage: 'Assistant',
assistantMessageCount: 'Assistant +{{count}}',
noAssistantMessage: 'No assistant reply recorded',
messageCount: 'Messages',
conversationTrace: 'Conversation Trace',
noLlmCalls: 'No model calls recorded',
roles: {
user: 'User',
assistant: 'Assistant',
message: 'Message',
},
},
llmCalls: {
title: 'LLM Calls',
@@ -1363,15 +1347,6 @@ const enUS = {
avgDuration: 'Avg Duration',
calls: 'Calls',
},
toolCalls: {
title: 'Tool Calls',
totalCalls: 'Calls',
duration: 'Tool Duration',
errorCalls: 'Failed Calls',
arguments: 'Arguments',
result: 'Result',
noToolCalls: 'No tool calls recorded',
},
tokens: {
totalTokens: 'Total Tokens',
inputTokens: 'Input Tokens',
@@ -1798,6 +1773,16 @@ const enUS = {
loginSuccess: 'Login successful! Token has been filled in',
loginFailed: 'Login failed',
},
itchat: {
scanLogin: 'Scan QR Login',
scanQRCode: 'Scan the QR code below with WeChat to login',
loginSuccess:
'Login successful! Session cached. Save config to start using.',
loginFailed: 'Login failed',
connecting: 'Starting WeChat login...',
waitingForScan: 'Waiting for scan',
retry: 'Retry',
},
dingtalk: {
createApp: 'One-Click Create DingTalk App',
scanQRCode:
-29
View File
@@ -842,12 +842,6 @@ const esES = {
tabTools: 'Herramientas',
tabResources: 'Recursos',
tabDocs: 'Documentación',
tabLogs: 'Registros',
logsLevelAll: 'Todos los niveles',
logsRefresh: 'Actualizar',
logsAutoRefresh: 'Actualización automática',
logsEmpty:
'Aún no hay registros. Los registros de ejecución del servidor MCP aparecerán aquí.',
noReadme: 'No hay documentación disponible',
parseResultFailed: 'Error al analizar el resultado de la prueba',
noResultReturned: 'La prueba no devolvió resultados',
@@ -1369,20 +1363,6 @@ const esES = {
level: 'Nivel',
runner: 'Ejecutor',
viewConversation: 'Ver conversación',
turns: '{{count}} turnos de conversación',
userMessage: 'Usuario',
noUserMessage: 'No se registró entrada del usuario',
assistantMessage: 'Asistente',
assistantMessageCount: 'Asistente +{{count}}',
noAssistantMessage: 'No se registró respuesta del asistente',
messageCount: 'Mensajes',
conversationTrace: 'Flujo de conversación',
noLlmCalls: 'No se registraron llamadas al modelo',
roles: {
user: 'Usuario',
assistant: 'Asistente',
message: 'Mensaje',
},
},
llmCalls: {
title: 'Llamadas LLM',
@@ -1397,15 +1377,6 @@ const esES = {
avgDuration: 'Duración promedio',
calls: 'Llamadas',
},
toolCalls: {
title: 'Llamadas de herramientas',
totalCalls: 'Llamadas',
duration: 'Duración de herramientas',
errorCalls: 'Llamadas fallidas',
arguments: 'Argumentos',
result: 'Resultado',
noToolCalls: 'No se registraron llamadas de herramientas',
},
tokens: {
totalTokens: 'Tokens totales',
inputTokens: 'Tokens de entrada',
+4 -28
View File
@@ -360,6 +360,10 @@ const jaJP = {
log: 'ログ',
configuration: '設定',
logs: 'ログ',
runtimeConnected: '接続済み',
runtimeConnecting: '接続中',
runtimeDisconnected: '切断済み',
runtimeError: '接続エラー',
basicInfo: '基本情報',
basicInfoDescription: 'ボットの名前と説明を設定',
routingConnection: 'ルーティングと接続',
@@ -834,11 +838,6 @@ const jaJP = {
tabTools: 'ツール',
tabResources: 'リソース',
tabDocs: 'ドキュメント',
tabLogs: 'ログ',
logsLevelAll: 'すべてのレベル',
logsRefresh: '更新',
logsAutoRefresh: '自動更新',
logsEmpty: 'ログはありません。MCPサーバーの実行ログがここに表示されます。',
noReadme: 'ドキュメントがありません',
parseResultFailed: 'テスト結果の解析に失敗しました',
noResultReturned: 'テスト結果が返されませんでした',
@@ -1341,20 +1340,6 @@ const jaJP = {
level: 'レベル',
runner: 'ランナー',
viewConversation: '会話詳細を表示',
turns: '{{count}} 会話ターン',
userMessage: 'ユーザー',
noUserMessage: 'ユーザー入力は記録されていません',
assistantMessage: 'アシスタント',
assistantMessageCount: 'アシスタント +{{count}}',
noAssistantMessage: 'アシスタントの返信は記録されていません',
messageCount: 'メッセージ数',
conversationTrace: '会話トレース',
noLlmCalls: 'モデル呼び出しは記録されていません',
roles: {
user: 'ユーザー',
assistant: 'アシスタント',
message: 'メッセージ',
},
},
llmCalls: {
title: 'LLM呼び出し',
@@ -1369,15 +1354,6 @@ const jaJP = {
avgDuration: '平均期間',
calls: '呼び出し',
},
toolCalls: {
title: 'ツール呼び出し',
totalCalls: '呼び出し',
duration: 'ツール時間',
errorCalls: '失敗した呼び出し',
arguments: '引数',
result: '結果',
noToolCalls: 'ツール呼び出しは記録されていません',
},
tokens: {
totalTokens: '総トークン数',
inputTokens: '入力トークン',
-29
View File
@@ -839,12 +839,6 @@ const ruRU = {
tabTools: 'Инструменты',
tabResources: 'Ресурсы',
tabDocs: 'Документация',
tabLogs: 'Журнал',
logsLevelAll: 'Все уровни',
logsRefresh: 'Обновить',
logsAutoRefresh: 'Автообновление',
logsEmpty:
'Журналов пока нет. Здесь будут отображаться журналы выполнения MCP-сервера.',
noReadme: 'Документация отсутствует',
parseResultFailed: 'Не удалось разобрать результат теста',
noResultReturned: 'Тест не вернул результат',
@@ -1345,20 +1339,6 @@ const ruRU = {
level: 'Уровень',
runner: 'Обработчик',
viewConversation: 'Просмотр диалога',
turns: '{{count}} диалоговых ходов',
userMessage: 'Пользователь',
noUserMessage: 'Ввод пользователя не записан',
assistantMessage: 'Ассистент',
assistantMessageCount: 'Ассистент +{{count}}',
noAssistantMessage: 'Ответ ассистента не записан',
messageCount: 'Сообщения',
conversationTrace: 'Ход диалога',
noLlmCalls: 'Вызовы модели не записаны',
roles: {
user: 'Пользователь',
assistant: 'Ассистент',
message: 'Сообщение',
},
},
llmCalls: {
title: 'Вызовы LLM',
@@ -1373,15 +1353,6 @@ const ruRU = {
avgDuration: 'Средняя длительность',
calls: 'Вызовы',
},
toolCalls: {
title: 'Вызовы инструментов',
totalCalls: 'Вызовы',
duration: 'Длительность инструментов',
errorCalls: 'Неудачные вызовы',
arguments: 'Аргументы',
result: 'Результат',
noToolCalls: 'Вызовы инструментов не записаны',
},
tokens: {
totalTokens: 'Всего токенов',
inputTokens: 'Входные токены',
-28
View File
@@ -817,11 +817,6 @@ const thTH = {
tabTools: 'เครื่องมือ',
tabResources: 'ทรัพยากร',
tabDocs: 'เอกสาร',
tabLogs: 'บันทึก',
logsLevelAll: 'ทุกระดับ',
logsRefresh: 'รีเฟรช',
logsAutoRefresh: 'รีเฟรชอัตโนมัติ',
logsEmpty: 'ยังไม่มีบันทึก บันทึกการทำงานของ MCP Server จะแสดงที่นี่',
noReadme: 'ไม่มีเอกสาร',
parseResultFailed: 'ไม่สามารถแยกวิเคราะห์ผลการทดสอบได้',
noResultReturned: 'การทดสอบไม่ส่งผลลัพธ์กลับมา',
@@ -1312,20 +1307,6 @@ const thTH = {
level: 'ระดับ',
runner: 'ตัวประมวลผล',
viewConversation: 'ดูการสนทนา',
turns: '{{count}} รอบการสนทนา',
userMessage: 'ผู้ใช้',
noUserMessage: 'ยังไม่มีการบันทึกข้อความจากผู้ใช้',
assistantMessage: 'ผู้ช่วย',
assistantMessageCount: 'ผู้ช่วย +{{count}}',
noAssistantMessage: 'ยังไม่มีการบันทึกคำตอบจากผู้ช่วย',
messageCount: 'จำนวนข้อความ',
conversationTrace: 'ลำดับการสนทนา',
noLlmCalls: 'ยังไม่มีการบันทึกการเรียกโมเดล',
roles: {
user: 'ผู้ใช้',
assistant: 'ผู้ช่วย',
message: 'ข้อความ',
},
},
llmCalls: {
title: 'การเรียก LLM',
@@ -1340,15 +1321,6 @@ const thTH = {
avgDuration: 'ระยะเวลาเฉลี่ย',
calls: 'การเรียก',
},
toolCalls: {
title: 'การเรียกใช้เครื่องมือ',
totalCalls: 'การเรียก',
duration: 'ระยะเวลาเครื่องมือ',
errorCalls: 'การเรียกที่ล้มเหลว',
arguments: 'อาร์กิวเมนต์',
result: 'ผลลัพธ์',
noToolCalls: 'ยังไม่มีการบันทึกการเรียกใช้เครื่องมือ',
},
tokens: {
totalTokens: 'Token ทั้งหมด',
inputTokens: 'Token อินพุต',
-29
View File
@@ -832,12 +832,6 @@ const viVN = {
tabTools: 'Công cụ',
tabResources: 'Tài nguyên',
tabDocs: 'Tài liệu',
tabLogs: 'Nhật ký',
logsLevelAll: 'Tất cả cấp độ',
logsRefresh: 'Làm mới',
logsAutoRefresh: 'Tự động làm mới',
logsEmpty:
'Chưa có nhật ký. Nhật ký chạy của MCP Server sẽ hiển thị ở đây.',
noReadme: 'Không có tài liệu',
parseResultFailed: 'Phân tích kết quả kiểm tra thất bại',
noResultReturned: 'Kiểm tra không trả về kết quả',
@@ -1338,20 +1332,6 @@ const viVN = {
level: 'Mức',
runner: 'Trình chạy',
viewConversation: 'Xem cuộc trò chuyện',
turns: '{{count}} lượt hội thoại',
userMessage: 'Người dùng',
noUserMessage: 'Chưa ghi nhận đầu vào người dùng',
assistantMessage: 'Trợ lý',
assistantMessageCount: 'Trợ lý +{{count}}',
noAssistantMessage: 'Chưa ghi nhận phản hồi của trợ lý',
messageCount: 'Số tin nhắn',
conversationTrace: 'Luồng hội thoại',
noLlmCalls: 'Chưa ghi nhận lệnh gọi mô hình',
roles: {
user: 'Người dùng',
assistant: 'Trợ lý',
message: 'Tin nhắn',
},
},
llmCalls: {
title: 'Cuộc gọi LLM',
@@ -1366,15 +1346,6 @@ const viVN = {
avgDuration: 'Thời lượng trung bình',
calls: 'Cuộc gọi',
},
toolCalls: {
title: 'Lượt gọi công cụ',
totalCalls: 'Lượt gọi',
duration: 'Thời lượng công cụ',
errorCalls: 'Lượt gọi thất bại',
arguments: 'Tham số',
result: 'Kết quả',
noToolCalls: 'Chưa ghi nhận lượt gọi công cụ',
},
tokens: {
totalTokens: 'Tổng số Token',
inputTokens: 'Token đầu vào',
+13 -28
View File
@@ -339,6 +339,10 @@ const zhHans = {
log: '日志',
configuration: '配置',
logs: '日志',
runtimeConnected: '已连接',
runtimeConnecting: '连接中',
runtimeDisconnected: '已掉线',
runtimeError: '连接错误',
basicInfo: '基础信息',
basicInfoDescription: '设置机器人名称和描述',
routingConnection: '路由与连接',
@@ -794,11 +798,6 @@ const zhHans = {
tabTools: '工具',
tabResources: '资源',
tabDocs: '文档',
tabLogs: '日志',
logsLevelAll: '全部级别',
logsRefresh: '刷新',
logsAutoRefresh: '自动刷新',
logsEmpty: '暂无日志。MCP 服务器的运行日志会显示在这里。',
noReadme: '暂无文档',
parseResultFailed: '解析测试结果失败',
noResultReturned: '测试未返回结果',
@@ -1271,20 +1270,6 @@ const zhHans = {
level: '级别',
runner: '执行器',
viewConversation: '显示对话详情',
turns: '{{count}} 轮对话',
userMessage: '用户',
noUserMessage: '未记录用户输入',
assistantMessage: '助手',
assistantMessageCount: '助手 +{{count}}',
noAssistantMessage: '未记录助手回复',
messageCount: '消息数',
conversationTrace: '消息链路',
noLlmCalls: '未记录模型调用',
roles: {
user: '用户',
assistant: '助手',
message: '消息',
},
},
llmCalls: {
title: 'LLM调用',
@@ -1299,15 +1284,6 @@ const zhHans = {
avgDuration: '平均耗时',
calls: '调用次数',
},
toolCalls: {
title: '工具调用',
totalCalls: '调用次数',
duration: '工具耗时',
errorCalls: '失败次数',
arguments: '参数',
result: '结果',
noToolCalls: '未记录工具调用',
},
tokens: {
totalTokens: '总 Token 数',
inputTokens: '输入 Token',
@@ -1719,6 +1695,15 @@ const zhHans = {
loginSuccess: '登录成功!令牌已自动填入',
loginFailed: '登录失败',
},
itchat: {
scanLogin: '扫码登录微信',
scanQRCode: '请使用微信扫描以下二维码登录',
loginSuccess: '登录成功!会话已缓存,保存配置后即可使用',
loginFailed: '登录失败',
connecting: '正在启动微信登录...',
waitingForScan: '等待扫码中',
retry: '重试',
},
dingtalk: {
createApp: '一键创建钉钉应用',
scanQRCode: '请使用钉钉扫描以下二维码,授权后将自动创建应用并填写凭据',
-28
View File
@@ -793,11 +793,6 @@ const zhHant = {
tabTools: '工具',
tabResources: '資源',
tabDocs: '文件',
tabLogs: '日誌',
logsLevelAll: '全部級別',
logsRefresh: '重新整理',
logsAutoRefresh: '自動重新整理',
logsEmpty: '暫無日誌。MCP 服務器的運行日誌會顯示在這裡。',
noReadme: '暫無文件',
parseResultFailed: '解析測試結果失敗',
noResultReturned: '測試未返回結果',
@@ -1269,20 +1264,6 @@ const zhHant = {
level: '級別',
runner: '執行器',
viewConversation: '顯示對話詳情',
turns: '{{count}} 輪對話',
userMessage: '使用者',
noUserMessage: '未記錄使用者輸入',
assistantMessage: '助手',
assistantMessageCount: '助手 +{{count}}',
noAssistantMessage: '未記錄助手回覆',
messageCount: '訊息數',
conversationTrace: '訊息鏈路',
noLlmCalls: '未記錄模型呼叫',
roles: {
user: '使用者',
assistant: '助手',
message: '訊息',
},
},
llmCalls: {
title: 'LLM呼叫',
@@ -1297,15 +1278,6 @@ const zhHant = {
avgDuration: '平均持續時間',
calls: '呼叫次數',
},
toolCalls: {
title: '工具呼叫',
totalCalls: '呼叫次數',
duration: '工具耗時',
errorCalls: '失敗次數',
arguments: '參數',
result: '結果',
noToolCalls: '未記錄工具呼叫',
},
tokens: {
totalTokens: '總 Token 數',
inputTokens: '輸入 Token',
@@ -1,179 +0,0 @@
import { expect, test } from '@playwright/test';
import { installLangBotApiMocks } from './fixtures/langbot-api';
const botId = 'bot-tool-timeline';
const sessionId = 'person-tool-timeline-user';
const botName = 'Tool Timeline Bot';
const pipelineId = 'pipeline-tool-timeline';
const pipelineName = 'Tool Timeline Pipeline';
function at(minute: number, second = 0) {
return `2026-07-02T10:${String(minute).padStart(2, '0')}:${String(
second,
).padStart(2, '0')}Z`;
}
function sessionMessage(
id: string,
role: 'user' | 'assistant',
minute: number,
content: string,
) {
return {
id,
timestamp: at(minute),
bot_id: botId,
bot_name: botName,
pipeline_id: pipelineId,
pipeline_name: pipelineName,
message_content: content,
session_id: sessionId,
status: 'success',
level: 'info',
platform: role === 'user' ? 'person' : 'bot',
user_id: 'timeline-user',
user_name: 'Timeline User',
runner_name: role === 'assistant' ? 'local-agent' : null,
variables: '{}',
role,
};
}
function toolCall(
id: string,
minute: number,
toolName: string,
duration: number,
status: 'success' | 'error' = 'success',
) {
return {
id,
timestamp: at(minute, 30),
tool_name: toolName,
tool_source: 'native',
duration,
status,
bot_id: botId,
bot_name: botName,
pipeline_id: pipelineId,
pipeline_name: pipelineName,
session_id: sessionId,
message_id: 'user-message',
arguments: JSON.stringify({ target: toolName }),
result: status === 'success' ? JSON.stringify({ ok: true }) : null,
error_message: status === 'error' ? 'Tool execution failed' : null,
};
}
test.describe('bot session monitor tool timeline', () => {
test('renders tool calls as left-side agent events interleaved with messages', async ({
page,
}) => {
await installLangBotApiMocks(page, {
authenticated: true,
monitoringSessions: [
{
session_id: sessionId,
bot_id: botId,
bot_name: botName,
pipeline_id: pipelineId,
pipeline_name: pipelineName,
message_count: 3,
start_time: at(0),
last_activity: at(4),
is_active: true,
platform: 'person',
user_id: 'timeline-user',
user_name: 'Timeline User',
},
],
sessionMessages: {
[sessionId]: [
sessionMessage('user-message', 'user', 0, 'Need a timeline check'),
sessionMessage(
'assistant-step-1',
'assistant',
2,
'Agent step 1: inspected repository files',
),
sessionMessage(
'assistant-step-2',
'assistant',
4,
'Agent step 2: test suite finished',
),
],
},
sessionAnalyses: {
[sessionId]: {
session_id: sessionId,
found: true,
tool_calls: [
toolCall('tool-repo-read', 1, 'repo_file_read', 80),
toolCall('tool-test-run', 3, 'run_test_suite', 140),
],
},
},
});
await page.goto(`/home/bots?id=${botId}`);
await page.getByRole('tab', { name: /Sessions/ }).click();
await page.getByRole('button', { name: /Timeline User/ }).click();
await expect(page.getByText('Need a timeline check')).toBeVisible();
await expect(
page.getByText('repo_file_read', { exact: true }),
).toBeVisible();
await expect(
page.getByText('Agent step 1: inspected repository files'),
).toBeVisible();
await expect(
page.getByText('run_test_suite', { exact: true }),
).toBeVisible();
await expect(
page.getByText('Agent step 2: test suite finished'),
).toBeVisible();
await expect(page.getByText('{"target":"repo_file_read"}')).toHaveCount(0);
await expect(page.getByText('{"ok":true}')).toHaveCount(0);
await expect(
page.locator('div.flex.justify-start').filter({
hasText: 'repo_file_read',
}),
).toHaveCount(1);
await expect(
page.locator('div.flex.justify-start').filter({
hasText: 'run_test_suite',
}),
).toHaveCount(1);
await expect(
page.locator('div.flex.justify-end').filter({
hasText: 'repo_file_read',
}),
).toHaveCount(0);
await expect(
page.locator('div.flex.justify-end').filter({
hasText: 'run_test_suite',
}),
).toHaveCount(0);
const text = await page.locator('body').innerText();
expect(text.indexOf('Need a timeline check')).toBeLessThan(
text.indexOf('repo_file_read'),
);
expect(text.indexOf('repo_file_read')).toBeLessThan(
text.indexOf('Agent step 1: inspected repository files'),
);
expect(
text.indexOf('Agent step 1: inspected repository files'),
).toBeLessThan(text.indexOf('run_test_suite'));
expect(text.indexOf('run_test_suite')).toBeLessThan(
text.indexOf('Agent step 2: test suite finished'),
);
await page.getByText('repo_file_read', { exact: true }).click();
await expect(page.getByText('{"target":"repo_file_read"}')).toBeVisible();
await expect(page.getByText('{"ok":true}').first()).toBeVisible();
});
});
-25
View File
@@ -88,31 +88,6 @@ test.describe('frontend CRUD smoke flows', () => {
).toBeVisible();
});
test('opens pipeline AI capabilities with malformed model options', async ({
page,
}) => {
await installLangBotApiMocks(page, { authenticated: true });
await page.goto('/home/pipelines?id=pipeline-ai');
await expect(page.locator('input[name="basic.name"]')).toBeVisible();
await page.getByRole('button', { name: /^AI$/ }).click();
await expect(page.getByText('Runtime')).toBeVisible();
await expect(
page.locator('[data-slot="card-title"]').filter({
hasText: 'Built-in Agent',
}),
).toBeVisible();
await expect(
page.locator('label').filter({
hasText: 'Model',
}),
).toBeVisible();
await expect(page.getByText('A <Select.Item')).toHaveCount(0);
await expect(page.getByText('500')).toHaveCount(0);
});
test('creates, edits, and deletes a knowledge base', async ({ page }) => {
await installLangBotApiMocks(page, { authenticated: true });
+5 -169
View File
@@ -72,11 +72,7 @@ interface LangBotApiMockState {
counters: Record<string, number>;
knowledgeBases: KnowledgeBaseMock[];
mcpServers: MCPServerMock[];
monitoringData: unknown;
monitoringSessions: unknown[];
pipelines: PipelineMock[];
sessionAnalyses: Record<string, unknown>;
sessionMessages: Record<string, unknown[]>;
skills: SkillMock[];
}
@@ -126,14 +122,12 @@ function emptyMonitoringData() {
},
messages: [],
llmCalls: [],
toolCalls: [],
embeddingCalls: [],
sessions: [],
errors: [],
totalCount: {
messages: 0,
llmCalls: 0,
toolCalls: 0,
embeddingCalls: 0,
sessions: 0,
errors: 0,
@@ -194,102 +188,6 @@ function makePipeline(
};
}
function pipelineMetadata() {
return {
configs: [
{
name: 'ai',
label: {
en_US: 'AI Capabilities',
zh_Hans: 'AI 能力',
},
stages: [
{
name: 'runner',
label: {
en_US: 'Runtime',
zh_Hans: '运行方式',
},
config: [
{
id: 'runner',
name: 'runner',
label: {
en_US: 'Runner',
zh_Hans: '运行器',
},
type: 'select',
required: true,
default: 'local-agent',
options: [
{
name: 'local-agent',
label: {
en_US: 'Built-in Agent',
zh_Hans: '内置 Agent',
},
},
],
},
],
},
{
name: 'local-agent',
label: {
en_US: 'Built-in Agent',
zh_Hans: '内置 Agent',
},
config: [
{
id: 'model',
name: 'model',
label: {
en_US: 'Model',
zh_Hans: '模型',
},
type: 'model-fallback-selector',
required: true,
default: {
primary: 'llm-valid',
fallbacks: [],
},
},
],
},
],
},
],
};
}
function providerModelList() {
return {
models: [
{
uuid: '',
name: 'Broken Empty UUID Model',
provider_uuid: 'provider-empty',
provider: {
uuid: 'provider-empty',
name: 'Broken Provider',
requester: 'mock-provider',
},
},
{
uuid: 'llm-valid',
name: 'Valid Mock Model',
provider_uuid: 'provider-valid',
provider: {
uuid: 'provider-valid',
name: 'Mock Provider',
requester: 'mock-provider',
},
abilities: ['func_call'],
},
],
};
}
function knowledgeEngine() {
return {
plugin_id: 'builtin/minimal-knowledge',
@@ -491,20 +389,8 @@ async function handleBackendApi(route: Route, state: LangBotApiMockState) {
});
}
if (path === '/api/v1/provider/models/llm') {
return fulfillJson(route, providerModelList());
}
if (path === '/api/v1/provider/models/embedding') {
return fulfillJson(route, { models: [] });
}
if (path === '/api/v1/provider/models/rerank') {
return fulfillJson(route, { models: [] });
}
if (path === '/api/v1/pipelines/_/metadata') {
return fulfillJson(route, pipelineMetadata());
return fulfillJson(route, { configs: [] });
}
if (path === '/api/v1/pipelines') {
@@ -803,43 +689,11 @@ async function handleBackendApi(route: Route, state: LangBotApiMockState) {
}
if (path === '/api/v1/monitoring/data') {
return fulfillJson(route, state.monitoringData);
}
if (path === '/api/v1/monitoring/sessions') {
return fulfillJson(route, {
sessions: state.monitoringSessions,
total: state.monitoringSessions.length,
});
}
if (path === '/api/v1/monitoring/messages') {
const sessionId = url.searchParams.get('sessionId') || '';
const messages = state.sessionMessages[sessionId] || [];
return fulfillJson(route, {
messages,
total: messages.length,
});
}
const sessionAnalysisMatch = path.match(
/^\/api\/v1\/monitoring\/sessions\/([^/]+)\/analysis$/,
);
if (sessionAnalysisMatch) {
const sessionId = decodeURIComponent(sessionAnalysisMatch[1]);
return fulfillJson(
route,
state.sessionAnalyses[sessionId] || {
session_id: sessionId,
found: true,
tool_calls: [],
},
);
return fulfillJson(route, emptyMonitoringData());
}
if (path === '/api/v1/monitoring/overview') {
const data = state.monitoringData as { overview?: unknown };
return fulfillJson(route, data.overview || emptyMonitoringData().overview);
return fulfillJson(route, emptyMonitoringData().overview);
}
if (path === '/api/v1/monitoring/token-statistics') {
@@ -944,33 +798,15 @@ async function handleCloudApi(route: Route) {
export async function installLangBotApiMocks(
page: Page,
options: {
authenticated?: boolean;
monitoringData?: unknown;
monitoringSessions?: unknown[];
sessionAnalyses?: Record<string, unknown>;
sessionMessages?: Record<string, unknown[]>;
storage?: JsonRecord;
} = {},
options: { authenticated?: boolean; storage?: JsonRecord } = {},
) {
const {
authenticated = false,
monitoringData,
monitoringSessions,
sessionAnalyses,
sessionMessages,
storage = {},
} = options;
const { authenticated = false, storage = {} } = options;
const state: LangBotApiMockState = {
bots: [],
counters: {},
knowledgeBases: [],
mcpServers: [],
monitoringData: monitoringData || emptyMonitoringData(),
monitoringSessions: monitoringSessions || [],
pipelines: [],
sessionAnalyses: sessionAnalyses || {},
sessionMessages: sessionMessages || {},
skills: [],
};
-453
View File
@@ -1,453 +0,0 @@
import { expect, test } from '@playwright/test';
import { installLangBotApiMocks } from './fixtures/langbot-api';
import { buildConversationTurns } from '../../src/app/home/monitoring/utils/conversationTurns';
import {
ErrorLog,
LLMCall,
MonitoringMessage,
ToolCall,
} from '../../src/app/home/monitoring/types/monitoring';
const bot = {
id: 'bot-monitoring',
name: 'Monitoring Bot',
};
const pipeline = {
id: 'pipeline-monitoring',
name: 'Monitoring Pipeline',
};
function time(minute: number) {
return new Date(`2026-07-02T10:${String(minute).padStart(2, '0')}:00Z`);
}
function message(
id: string,
role: 'user' | 'assistant',
minute: number,
content: string,
sessionId = 'session-agent',
): MonitoringMessage {
return {
id,
timestamp: time(minute),
botId: bot.id,
botName: bot.name,
pipelineId: pipeline.id,
pipelineName: pipeline.name,
messageContent: content,
sessionId,
status: 'success',
level: 'info',
platform: role === 'user' ? 'person' : 'bot',
userId: 'user-1',
userName: 'Playwright User',
runnerName: 'local-agent',
variables: '{}',
role,
};
}
function llmCall(
id: string,
minute: number,
messageId: string | undefined,
input: number,
output: number,
duration: number,
sessionId = 'session-agent',
): LLMCall {
return {
id,
timestamp: time(minute),
modelName: 'gpt-5.5',
tokens: {
input,
output,
total: input + output,
},
duration,
status: 'success',
botId: bot.id,
botName: bot.name,
pipelineId: pipeline.id,
pipelineName: pipeline.name,
sessionId,
messageId,
};
}
function errorLog(id: string, minute: number, messageId: string): ErrorLog {
return {
id,
timestamp: time(minute),
errorType: 'ToolExecutionError',
errorMessage: 'Tool retry failed',
botId: bot.id,
botName: bot.name,
pipelineId: pipeline.id,
pipelineName: pipeline.name,
sessionId: 'session-agent',
messageId,
};
}
function toolCall(
id: string,
minute: number,
messageId: string | undefined,
toolName: string,
duration: number,
sessionId = 'session-agent',
status: 'success' | 'error' = 'success',
): ToolCall {
return {
id,
timestamp: time(minute),
toolName,
toolSource: 'native',
duration,
status,
botId: bot.id,
botName: bot.name,
pipelineId: pipeline.id,
pipelineName: pipeline.name,
sessionId,
messageId,
arguments: JSON.stringify({ query: toolName }),
result: status === 'success' ? JSON.stringify({ ok: true }) : undefined,
errorMessage: status === 'error' ? 'Tool failed' : undefined,
};
}
function rawMessage(message: MonitoringMessage) {
return {
id: message.id,
timestamp: message.timestamp.toISOString(),
bot_id: message.botId,
bot_name: message.botName,
pipeline_id: message.pipelineId,
pipeline_name: message.pipelineName,
message_content: message.messageContent,
session_id: message.sessionId,
status: message.status,
level: message.level,
platform: message.platform,
user_id: message.userId,
user_name: message.userName,
runner_name: message.runnerName,
variables: message.variables,
role: message.role,
};
}
function rawLlmCall(call: LLMCall) {
return {
id: call.id,
timestamp: call.timestamp.toISOString(),
model_name: call.modelName,
input_tokens: call.tokens.input,
output_tokens: call.tokens.output,
total_tokens: call.tokens.total,
duration: call.duration,
cost: call.cost,
status: call.status,
bot_id: call.botId,
bot_name: call.botName,
pipeline_id: call.pipelineId,
pipeline_name: call.pipelineName,
session_id: call.sessionId,
error_message: call.errorMessage,
message_id: call.messageId,
};
}
function rawError(error: ErrorLog) {
return {
id: error.id,
timestamp: error.timestamp.toISOString(),
error_type: error.errorType,
error_message: error.errorMessage,
bot_id: error.botId,
bot_name: error.botName,
pipeline_id: error.pipelineId,
pipeline_name: error.pipelineName,
session_id: error.sessionId,
stack_trace: error.stackTrace,
message_id: error.messageId,
};
}
function rawToolCall(call: ToolCall) {
return {
id: call.id,
timestamp: call.timestamp.toISOString(),
tool_name: call.toolName,
tool_source: call.toolSource,
duration: call.duration,
status: call.status,
bot_id: call.botId,
bot_name: call.botName,
pipeline_id: call.pipelineId,
pipeline_name: call.pipelineName,
session_id: call.sessionId,
message_id: call.messageId,
arguments: call.arguments,
result: call.result,
error_message: call.errorMessage,
};
}
function monitoringScenario() {
const messages = [
message(
'single-user',
'user',
1,
'Standalone question with no reply',
'session-single',
),
message('agent-user-1', 'user', 10, 'Need deployment plan'),
message('agent-assistant-1', 'assistant', 11, 'Agent step 1: inspect repo'),
message('agent-assistant-2', 'assistant', 12, 'Agent step 2: run tests'),
message(
'agent-assistant-3',
'assistant',
13,
'Final answer: deployment plan ready',
),
message('agent-user-2', 'user', 20, 'Continue with rollback plan'),
message('agent-assistant-4', 'assistant', 21, 'Rollback plan ready'),
];
const llmCalls = [
llmCall('agent-call-1', 10, 'agent-user-1', 100, 40, 120),
llmCall('agent-call-2', 11, 'agent-user-1', 200, 60, 220),
llmCall('agent-call-3', 12, 'agent-user-1', 300, 90, 260),
llmCall('agent-call-4', 20, 'agent-user-2', 50, 25, 80),
];
const errors = [errorLog('agent-error-1', 12, 'agent-user-1')];
const toolCalls = [
toolCall('agent-tool-1', 11, 'agent-user-1', 'repo_search', 90),
toolCall('agent-tool-2', 12, 'agent-user-1', 'run_tests', 150),
toolCall('agent-tool-3', 20, 'agent-user-2', 'rollback_lookup', 70),
];
return {
messages,
llmCalls,
toolCalls,
errors,
};
}
function rawMonitoringData() {
const scenario = monitoringScenario();
return {
overview: {
total_messages: scenario.messages.length,
llm_calls: scenario.llmCalls.length,
embedding_calls: 0,
model_calls: scenario.llmCalls.length,
success_rate: 100,
active_sessions: 2,
},
messages: scenario.messages.map(rawMessage),
llmCalls: scenario.llmCalls.map(rawLlmCall),
toolCalls: scenario.toolCalls.map(rawToolCall),
embeddingCalls: [],
sessions: [],
errors: scenario.errors.map(rawError),
totalCount: {
messages: scenario.messages.length,
llmCalls: scenario.llmCalls.length,
toolCalls: scenario.toolCalls.length,
embeddingCalls: 0,
sessions: 0,
errors: scenario.errors.length,
},
};
}
test.describe('monitoring conversation turn grouping', () => {
test('keeps a single user message as one observable turn', () => {
const userOnly = message(
'single-user-only',
'user',
1,
'No answer yet',
'session-user-only',
);
const turns = buildConversationTurns([userOnly], [], []);
expect(turns).toHaveLength(1);
expect(turns[0].id).toBe(userOnly.id);
expect(turns[0].userMessage?.messageContent).toBe('No answer yet');
expect(turns[0].assistantMessages).toHaveLength(0);
expect(turns[0].llmCalls).toHaveLength(0);
expect(turns[0].totalTokens).toBe(0);
});
test('groups multi-step agent execution and multiple replies into one user turn', () => {
const scenario = monitoringScenario();
const turns = buildConversationTurns(
scenario.messages,
scenario.llmCalls,
scenario.errors,
scenario.toolCalls,
);
const agentTurn = turns.find((turn) => turn.id === 'agent-user-1');
expect(agentTurn).toBeTruthy();
expect(agentTurn?.userMessage?.messageContent).toBe('Need deployment plan');
expect(
agentTurn?.assistantMessages.map((item) => item.messageContent),
).toEqual([
'Agent step 1: inspect repo',
'Agent step 2: run tests',
'Final answer: deployment plan ready',
]);
expect(agentTurn?.llmCalls).toHaveLength(3);
expect(agentTurn?.toolCalls).toHaveLength(2);
expect(agentTurn?.errors).toHaveLength(1);
expect(agentTurn?.totalTokens).toBe(790);
expect(agentTurn?.totalDuration).toBe(600);
expect(agentTurn?.totalToolDuration).toBe(240);
});
test('starts a new turn for each later user message in the same session', () => {
const firstUser = message('same-session-user-1', 'user', 1, 'First');
const firstReply = message(
'same-session-reply-1',
'assistant',
2,
'First reply',
);
const secondUser = message('same-session-user-2', 'user', 3, 'Second');
const secondReply = message(
'same-session-reply-2',
'assistant',
4,
'Second reply',
);
const turns = buildConversationTurns(
[firstUser, firstReply, secondUser, secondReply],
[
llmCall('same-session-call-1', 1, firstUser.id, 10, 5, 40),
llmCall('same-session-call-2', 3, secondUser.id, 20, 10, 50),
],
[],
);
expect(turns.map((turn) => turn.id)).toEqual([
'same-session-user-2',
'same-session-user-1',
]);
expect(
turns[0].assistantMessages.map((item) => item.messageContent),
).toEqual(['Second reply']);
expect(
turns[1].assistantMessages.map((item) => item.messageContent),
).toEqual(['First reply']);
});
test('attaches calls without message ids by session time', () => {
const user = message('fallback-user', 'user', 1, 'Use session fallback');
const assistant = message(
'fallback-assistant',
'assistant',
2,
'Fallback reply',
);
const call = llmCall('fallback-call', 2, undefined, 25, 5, 70);
const turns = buildConversationTurns([user, assistant], [call], []);
expect(turns).toHaveLength(1);
expect(turns[0].llmCalls).toHaveLength(1);
expect(turns[0].llmCalls[0].id).toBe(call.id);
expect(turns[0].totalTokens).toBe(30);
});
test('attaches tool calls without message ids by session time', () => {
const user = message('tool-fallback-user', 'user', 1, 'Use tool fallback');
const assistant = message(
'tool-fallback-assistant',
'assistant',
2,
'Tool fallback reply',
);
const call = toolCall(
'tool-fallback-call',
2,
undefined,
'memory_lookup',
45,
);
const turns = buildConversationTurns([user, assistant], [], [], [call]);
expect(turns).toHaveLength(1);
expect(turns[0].toolCalls).toHaveLength(1);
expect(turns[0].toolCalls[0].id).toBe(call.id);
expect(turns[0].totalToolDuration).toBe(45);
});
test('renders user-only, multi-agent, and multi-turn cases in the monitoring page', async ({
page,
}) => {
await installLangBotApiMocks(page, {
authenticated: true,
monitoringData: rawMonitoringData(),
});
await page.goto('/home/monitoring');
await expect(page.getByText('3 conversation turns')).toBeVisible();
await expect(
page.getByText('Standalone question with no reply'),
).toBeVisible();
await expect(page.getByText('No assistant reply recorded')).toBeVisible();
await expect(page.getByText('Need deployment plan')).toBeVisible();
await expect(page.getByText('Agent step 1: inspect repo')).toBeVisible();
await expect(page.getByText('Assistant +2')).toBeVisible();
await expect(page.getByText('3 LLM')).toBeVisible();
await expect(page.getByText('2 tools')).toBeVisible();
await expect(page.getByText('790 tokens')).toBeVisible();
await expect(page.getByText('1 errors')).toBeVisible();
await expect(page.getByText('Continue with rollback plan')).toBeVisible();
await expect(page.getByText('Rollback plan ready')).toBeVisible();
const agentTurn = page
.locator('div[role="button"]')
.filter({ hasText: 'Need deployment plan' });
await expect(agentTurn).toHaveCount(1);
await agentTurn.click();
await expect(page.getByText('Conversation Trace')).toBeVisible();
await expect(page.getByText('Agent step 2: run tests')).toBeVisible();
await expect(
page.getByText('Final answer: deployment plan ready'),
).toBeVisible();
await expect(page.getByText('LLM Calls (3)')).toBeVisible();
await expect(page.getByText('#3 gpt-5.5')).toBeVisible();
await expect(page.getByText('In: 300')).toBeVisible();
await expect(page.getByText('Out: 90')).toBeVisible();
await expect(page.getByText('Total: 390')).toBeVisible();
await expect(page.getByText('Tool Calls (2)')).toBeVisible();
await expect(page.getByText('#1 repo_search')).toBeVisible();
await expect(page.getByText('#2 run_tests')).toBeVisible();
await expect(page.getByText('Arguments')).toHaveCount(0);
await expect(page.getByText('Result')).toHaveCount(0);
await page.getByText('#1 repo_search').click();
await expect(page.getByText('Arguments').first()).toBeVisible();
await expect(page.getByText('Result').first()).toBeVisible();
await expect(page.getByText('Tool retry failed')).toBeVisible();
});
});
@@ -1,195 +0,0 @@
import { expect, test } from '@playwright/test';
import { installLangBotApiMocks } from './fixtures/langbot-api';
const bot = {
id: 'bot-pipeline-monitoring',
name: 'Pipeline Bot',
};
const pipeline = {
id: 'pipeline-monitoring',
name: 'Pipeline Under Test',
};
function at(minute: number) {
return `2026-07-02T10:${String(minute).padStart(2, '0')}:00Z`;
}
function message(
id: string,
role: 'user' | 'assistant',
minute: number,
content: string,
sessionId = 'session-pipeline-agent',
) {
return {
id,
timestamp: at(minute),
bot_id: bot.id,
bot_name: bot.name,
pipeline_id: pipeline.id,
pipeline_name: pipeline.name,
message_content: content,
session_id: sessionId,
status: 'success',
level: 'info',
platform: role === 'user' ? 'person' : 'bot',
user_id: 'pipeline-user',
user_name: 'Pipeline User',
runner_name: 'local-agent',
variables: '{}',
role,
};
}
function llmCall(
id: string,
minute: number,
messageId: string,
input: number,
output: number,
duration: number,
) {
return {
id,
timestamp: at(minute),
model_name: 'gpt-5.5',
input_tokens: input,
output_tokens: output,
total_tokens: input + output,
duration,
cost: 0,
status: 'success',
bot_id: bot.id,
bot_name: bot.name,
pipeline_id: pipeline.id,
pipeline_name: pipeline.name,
session_id: 'session-pipeline-agent',
message_id: messageId,
};
}
function toolCall(id: string, minute: number, messageId: string, name: string) {
return {
id,
timestamp: at(minute),
tool_name: name,
tool_source: 'native',
duration: 120,
status: 'success',
bot_id: bot.id,
bot_name: bot.name,
pipeline_id: pipeline.id,
pipeline_name: pipeline.name,
session_id: 'session-pipeline-agent',
message_id: messageId,
arguments: JSON.stringify({ query: name }),
result: JSON.stringify({ ok: true }),
};
}
function monitoringData() {
const messages = [
message(
'single-user',
'user',
1,
'Pipeline single user message without reply',
'session-pipeline-single',
),
message('agent-user', 'user', 10, 'Pipeline needs a deployment plan'),
message(
'agent-assistant-1',
'assistant',
11,
'Pipeline agent step 1: inspect repository',
),
message(
'agent-assistant-2',
'assistant',
12,
'Pipeline agent step 2: run tests',
),
message(
'agent-assistant-3',
'assistant',
13,
'Pipeline final answer: deployment ready',
),
];
const llmCalls = [
llmCall('pipeline-call-1', 10, 'agent-user', 100, 40, 180),
llmCall('pipeline-call-2', 11, 'agent-user', 140, 50, 220),
];
const toolCalls = [
toolCall('pipeline-tool-1', 11, 'agent-user', 'repo_search'),
toolCall('pipeline-tool-2', 12, 'agent-user', 'run_tests'),
];
return {
overview: {
total_messages: messages.length,
llm_calls: llmCalls.length,
embedding_calls: 0,
model_calls: llmCalls.length,
success_rate: 100,
active_sessions: 2,
},
messages,
llmCalls,
toolCalls,
embeddingCalls: [],
sessions: [],
errors: [],
totalCount: {
messages: messages.length,
llmCalls: llmCalls.length,
toolCalls: toolCalls.length,
embeddingCalls: 0,
sessions: 0,
errors: 0,
},
};
}
test.describe('pipeline monitoring conversation turns', () => {
test('uses conversation turns and folded tool calls in the pipeline dashboard', async ({
page,
}) => {
await installLangBotApiMocks(page, {
authenticated: true,
monitoringData: monitoringData(),
});
await page.goto(`/home/pipelines?id=${pipeline.id}`);
await page.getByRole('tab', { name: 'Dashboard' }).click();
await expect(page.getByText('2 conversation turns')).toBeVisible();
await expect(
page.getByText('Pipeline single user message without reply'),
).toBeVisible();
await expect(
page.getByText('Pipeline needs a deployment plan'),
).toBeVisible();
await expect(
page.getByText('Pipeline agent step 1: inspect repository'),
).toBeVisible();
await expect(page.getByText('Assistant +2')).toBeVisible();
await expect(page.getByText('2 tools')).toBeVisible();
const agentTurn = page
.locator('div[role="button"]')
.filter({ hasText: 'Pipeline needs a deployment plan' });
await expect(agentTurn).toHaveCount(1);
await agentTurn.click();
await expect(page.getByText('Tool Calls (2)')).toBeVisible();
await expect(page.getByText('#1 repo_search')).toBeVisible();
await expect(page.getByText('#2 run_tests')).toBeVisible();
await expect(page.getByText('Arguments')).toHaveCount(0);
await page.getByText('#1 repo_search').click();
await expect(page.getByText('Arguments')).toBeVisible();
await expect(page.getByText('Result')).toBeVisible();
});
});