mirror of
https://github.com/langbot-app/LangBot.git
synced 2026-08-28 21:27:14 +00:00
refactor(bots): remove route execution test
This commit is contained in:
@@ -90,15 +90,9 @@ class BotsRouterGroup(group.RouterGroup):
|
||||
permission=Permission.RESOURCE_VIEW,
|
||||
)
|
||||
async def _(bot_uuid: str, request_context: RequestContext) -> str:
|
||||
return self.success(
|
||||
data=await self.ap.bot_service.list_event_route_statuses(
|
||||
request_context, bot_uuid
|
||||
)
|
||||
)
|
||||
return self.success(data=await self.ap.bot_service.list_event_route_statuses(request_context, bot_uuid))
|
||||
|
||||
async def _dry_run_event_route(
|
||||
bot_uuid: str, request_context: RequestContext
|
||||
) -> str:
|
||||
async def _dry_run_event_route(bot_uuid: str, request_context: RequestContext) -> str:
|
||||
json_data = await quart.request.json
|
||||
if not isinstance(json_data, dict):
|
||||
return self.http_status(400, -1, 'invalid request body')
|
||||
@@ -128,24 +122,6 @@ class BotsRouterGroup(group.RouterGroup):
|
||||
permission=Permission.RESOURCE_VIEW,
|
||||
)(_dry_run_event_route)
|
||||
|
||||
@self.route(
|
||||
'/<bot_uuid>/event-routes/test',
|
||||
methods=['POST'],
|
||||
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
|
||||
permission=Permission.RUNTIME_OPERATE,
|
||||
)
|
||||
async def _(bot_uuid: str, request_context: RequestContext) -> str:
|
||||
json_data = await quart.request.json
|
||||
if not isinstance(json_data, dict):
|
||||
return self.http_status(400, -1, 'invalid request body')
|
||||
result = await self.ap.bot_service.dispatch_test_event_route(
|
||||
request_context,
|
||||
bot_uuid=bot_uuid,
|
||||
event_type=json_data.get('event_type'),
|
||||
payload=json_data.get('event_data', json_data.get('payload')),
|
||||
)
|
||||
return self.success(data=result)
|
||||
|
||||
@self.route(
|
||||
'/<bot_uuid>/send_message',
|
||||
methods=['POST'],
|
||||
|
||||
@@ -21,7 +21,6 @@ class BotService:
|
||||
FAILURE_PROCESSOR_NOT_FOUND = 'processor_not_found'
|
||||
FAILURE_PROCESSOR_INCOMPATIBLE = 'processor_incompatible'
|
||||
FAILURE_INVALID_EVENT = 'invalid_event'
|
||||
FAILURE_BOT_RUNTIME_UNAVAILABLE = 'bot_runtime_unavailable'
|
||||
ROUTE_TRACE_KIND = 'event_route_trace'
|
||||
|
||||
BOT_FIELDS = {
|
||||
@@ -798,68 +797,6 @@ class BotService:
|
||||
'stale_routes': stale_routes,
|
||||
}
|
||||
|
||||
async def dispatch_test_event_route(
|
||||
self,
|
||||
context: TenantContext,
|
||||
bot_uuid: str,
|
||||
event_type: str,
|
||||
payload: dict[str, typing.Any] | None = None,
|
||||
) -> dict[str, typing.Any]:
|
||||
"""Dispatch a synthetic event through the saved Bot runtime route configuration."""
|
||||
event_type = str(event_type or '').strip()
|
||||
if not event_type:
|
||||
return {
|
||||
'dispatched': False,
|
||||
'event_type': '',
|
||||
'failure_code': self.FAILURE_INVALID_EVENT,
|
||||
'reason': 'event_type is required',
|
||||
'suppressed_outputs': [],
|
||||
'route_status': {
|
||||
'routes': [],
|
||||
'unmatched_events': [],
|
||||
'stale_routes': [],
|
||||
},
|
||||
}
|
||||
if payload is not None and not isinstance(payload, dict):
|
||||
return {
|
||||
'dispatched': False,
|
||||
'event_type': event_type,
|
||||
'failure_code': self.FAILURE_INVALID_EVENT,
|
||||
'reason': 'payload must be an object',
|
||||
'suppressed_outputs': [],
|
||||
'route_status': {
|
||||
'routes': [],
|
||||
'unmatched_events': [],
|
||||
'stale_routes': [],
|
||||
},
|
||||
}
|
||||
|
||||
if await self.get_bot(context, bot_uuid, include_secret=False) is None:
|
||||
raise WorkspaceNotFoundError('Bot not found')
|
||||
runtime_bot = await self.ap.platform_mgr.get_bot_by_uuid(context, bot_uuid)
|
||||
if runtime_bot is None:
|
||||
return {
|
||||
'dispatched': False,
|
||||
'event_type': event_type,
|
||||
'failure_code': self.FAILURE_BOT_RUNTIME_UNAVAILABLE,
|
||||
'reason': 'Bot runtime is unavailable',
|
||||
'suppressed_outputs': [],
|
||||
'route_status': await self.list_event_route_statuses(context, bot_uuid),
|
||||
}
|
||||
|
||||
dispatch_result = await runtime_bot.dispatch_test_event(event_type, payload or {})
|
||||
route_status = await self.list_event_route_statuses(context, bot_uuid)
|
||||
return {
|
||||
'dispatched': bool(dispatch_result.get('dispatched')),
|
||||
'event_type': event_type,
|
||||
'status': dispatch_result.get('status'),
|
||||
'binding_id': dispatch_result.get('binding_id'),
|
||||
'failure_code': dispatch_result.get('failure_code'),
|
||||
'reason': dispatch_result.get('reason'),
|
||||
'suppressed_outputs': dispatch_result.get('suppressed_outputs', []),
|
||||
'route_status': route_status,
|
||||
}
|
||||
|
||||
async def send_message(
|
||||
self,
|
||||
context: TenantContext,
|
||||
|
||||
@@ -132,25 +132,6 @@ class LangBotMCPServer:
|
||||
async def list_bot_event_route_statuses(bot_uuid: str) -> str:
|
||||
return _dump(await ap.bot_service.list_event_route_statuses(bot_uuid))
|
||||
|
||||
@mcp.tool(
|
||||
description=(
|
||||
'Dispatch a synthetic event through the saved bot event routes. '
|
||||
'This validates routing without sending real outbound platform messages.'
|
||||
)
|
||||
)
|
||||
async def test_bot_event_route(
|
||||
bot_uuid: str,
|
||||
event_type: str,
|
||||
payload: dict | None = None,
|
||||
) -> str:
|
||||
return _dump(
|
||||
await ap.bot_service.dispatch_test_event_route(
|
||||
bot_uuid=bot_uuid,
|
||||
event_type=event_type,
|
||||
payload=payload,
|
||||
)
|
||||
)
|
||||
|
||||
# ----- Pipelines ----------------------------------------------- #
|
||||
@mcp.tool(description='List all pipelines.')
|
||||
async def list_pipelines() -> str:
|
||||
|
||||
@@ -51,226 +51,6 @@ from langbot_plugin.api.entities.builtin.agent_runner.input import AgentInput
|
||||
from langbot_plugin.api.entities.builtin.agent_runner.delivery import DeliveryContext
|
||||
|
||||
|
||||
class SyntheticRouteTestAdapter:
|
||||
"""Adapter wrapper that suppresses outbound platform delivery for test events."""
|
||||
|
||||
SIDE_EFFECT_API_NAMES = {
|
||||
'send_message',
|
||||
'reply_message',
|
||||
'reply_message_chunk',
|
||||
'create_message_card',
|
||||
'edit_message',
|
||||
'delete_message',
|
||||
'add_reaction',
|
||||
'remove_reaction',
|
||||
'forward_message',
|
||||
'set_group_name',
|
||||
'mute_member',
|
||||
'unmute_member',
|
||||
'kick_member',
|
||||
'leave_group',
|
||||
'approve_friend_request',
|
||||
'approve_group_invite',
|
||||
'upload_file',
|
||||
'call_platform_api',
|
||||
}
|
||||
|
||||
def __init__(self, source: abstract_platform_adapter.AbstractMessagePlatformAdapter):
|
||||
self.source = source
|
||||
self.bot_account_id = getattr(source, 'bot_account_id', '')
|
||||
self.config = getattr(source, 'config', {})
|
||||
self.logger = getattr(source, 'logger', None)
|
||||
self.suppressed_outputs: list[dict[str, typing.Any]] = []
|
||||
|
||||
@staticmethod
|
||||
def _message_to_payload(message: platform_message.MessageChain) -> typing.Any:
|
||||
return message.model_dump() if hasattr(message, 'model_dump') else str(message)
|
||||
|
||||
def _suppress(self, method: str, **payload: typing.Any) -> None:
|
||||
self.suppressed_outputs.append({'method': method, **payload})
|
||||
|
||||
def __getattr__(self, name: str) -> typing.Any:
|
||||
return getattr(self.source, name)
|
||||
|
||||
def get_supported_apis(self) -> list[str]:
|
||||
get_supported_apis = getattr(self.source, 'get_supported_apis', None)
|
||||
if not callable(get_supported_apis):
|
||||
return []
|
||||
return [api_name for api_name in get_supported_apis() if api_name not in self.SIDE_EFFECT_API_NAMES]
|
||||
|
||||
async def send_message(
|
||||
self,
|
||||
target_type: str,
|
||||
target_id: str,
|
||||
message: platform_message.MessageChain,
|
||||
) -> dict[str, typing.Any]:
|
||||
self._suppress(
|
||||
'send_message',
|
||||
target_type=target_type,
|
||||
target_id=target_id,
|
||||
message=self._message_to_payload(message),
|
||||
)
|
||||
return {'suppressed': True}
|
||||
|
||||
async def reply_message(
|
||||
self,
|
||||
message_source: platform_events.MessageEvent,
|
||||
message: platform_message.MessageChain,
|
||||
quote_origin: bool = False,
|
||||
) -> dict[str, typing.Any]:
|
||||
self._suppress(
|
||||
'reply_message',
|
||||
message=self._message_to_payload(message),
|
||||
quote_origin=quote_origin,
|
||||
)
|
||||
return {'suppressed': True}
|
||||
|
||||
async def reply_message_chunk(
|
||||
self,
|
||||
message_source: platform_events.MessageEvent,
|
||||
bot_message: dict,
|
||||
message: platform_message.MessageChain,
|
||||
quote_origin: bool = False,
|
||||
is_final: bool = False,
|
||||
) -> dict[str, typing.Any]:
|
||||
self._suppress(
|
||||
'reply_message_chunk',
|
||||
message=self._message_to_payload(message),
|
||||
quote_origin=quote_origin,
|
||||
is_final=is_final,
|
||||
)
|
||||
return {'suppressed': True}
|
||||
|
||||
async def create_message_card(
|
||||
self,
|
||||
message_id: str | int,
|
||||
event: platform_events.MessageEvent,
|
||||
) -> bool:
|
||||
self._suppress('create_message_card', message_id=str(message_id))
|
||||
return False
|
||||
|
||||
async def is_stream_output_supported(self) -> bool:
|
||||
return False
|
||||
|
||||
async def edit_message(
|
||||
self,
|
||||
chat_type: str,
|
||||
chat_id: typing.Union[int, str],
|
||||
message_id: typing.Union[int, str],
|
||||
new_content: platform_message.MessageChain,
|
||||
) -> None:
|
||||
self._suppress(
|
||||
'edit_message',
|
||||
chat_type=str(chat_type),
|
||||
chat_id=str(chat_id),
|
||||
message_id=str(message_id),
|
||||
new_content=self._message_to_payload(new_content),
|
||||
)
|
||||
|
||||
async def delete_message(
|
||||
self,
|
||||
chat_type: str,
|
||||
chat_id: typing.Union[int, str],
|
||||
message_id: typing.Union[int, str],
|
||||
) -> None:
|
||||
self._suppress(
|
||||
'delete_message',
|
||||
chat_type=str(chat_type),
|
||||
chat_id=str(chat_id),
|
||||
message_id=str(message_id),
|
||||
)
|
||||
|
||||
async def forward_message(
|
||||
self,
|
||||
from_chat_type: str,
|
||||
from_chat_id: typing.Union[int, str],
|
||||
message_id: typing.Union[int, str],
|
||||
to_chat_type: str,
|
||||
to_chat_id: typing.Union[int, str],
|
||||
) -> platform_events.MessageResult:
|
||||
self._suppress(
|
||||
'forward_message',
|
||||
from_chat_type=str(from_chat_type),
|
||||
from_chat_id=str(from_chat_id),
|
||||
message_id=str(message_id),
|
||||
to_chat_type=str(to_chat_type),
|
||||
to_chat_id=str(to_chat_id),
|
||||
)
|
||||
return platform_events.MessageResult(raw={'suppressed': True})
|
||||
|
||||
async def set_group_name(
|
||||
self,
|
||||
group_id: typing.Union[int, str],
|
||||
name: str,
|
||||
) -> None:
|
||||
self._suppress('set_group_name', group_id=str(group_id), name=name)
|
||||
|
||||
async def mute_member(
|
||||
self,
|
||||
group_id: typing.Union[int, str],
|
||||
user_id: typing.Union[int, str],
|
||||
duration: int = 0,
|
||||
) -> None:
|
||||
self._suppress(
|
||||
'mute_member',
|
||||
group_id=str(group_id),
|
||||
user_id=str(user_id),
|
||||
duration=duration,
|
||||
)
|
||||
|
||||
async def unmute_member(
|
||||
self,
|
||||
group_id: typing.Union[int, str],
|
||||
user_id: typing.Union[int, str],
|
||||
) -> None:
|
||||
self._suppress('unmute_member', group_id=str(group_id), user_id=str(user_id))
|
||||
|
||||
async def kick_member(
|
||||
self,
|
||||
group_id: typing.Union[int, str],
|
||||
user_id: typing.Union[int, str],
|
||||
) -> None:
|
||||
self._suppress('kick_member', group_id=str(group_id), user_id=str(user_id))
|
||||
|
||||
async def leave_group(
|
||||
self,
|
||||
group_id: typing.Union[int, str],
|
||||
) -> None:
|
||||
self._suppress('leave_group', group_id=str(group_id))
|
||||
|
||||
async def approve_friend_request(
|
||||
self,
|
||||
request_id: typing.Union[int, str],
|
||||
approve: bool = True,
|
||||
remark: str | None = None,
|
||||
) -> None:
|
||||
self._suppress(
|
||||
'approve_friend_request',
|
||||
request_id=str(request_id),
|
||||
approve=approve,
|
||||
remark=remark,
|
||||
)
|
||||
|
||||
async def approve_group_invite(
|
||||
self,
|
||||
request_id: typing.Union[int, str],
|
||||
approve: bool = True,
|
||||
) -> None:
|
||||
self._suppress(
|
||||
'approve_group_invite',
|
||||
request_id=str(request_id),
|
||||
approve=approve,
|
||||
)
|
||||
|
||||
async def upload_file(self, file_data: bytes, filename: str) -> str:
|
||||
self._suppress('upload_file', filename=filename, size=len(file_data))
|
||||
return f'suppressed:{filename}'
|
||||
|
||||
async def call_platform_api(self, action: str, params: dict | None = None) -> dict:
|
||||
self._suppress('call_platform_api', action=action, params=params or {})
|
||||
return {'suppressed': True}
|
||||
|
||||
|
||||
class RuntimeBot:
|
||||
"""运行时机器人"""
|
||||
|
||||
@@ -568,126 +348,6 @@ class RuntimeBot:
|
||||
"""Return the selected event binding plus per-binding diagnostic steps."""
|
||||
return self._evaluate_eba_event_bindings(self._get_event_bindings(), event, event_type)
|
||||
|
||||
@staticmethod
|
||||
def _build_test_platform_event(
|
||||
event_type: str,
|
||||
payload: dict[str, typing.Any] | None = None,
|
||||
) -> platform_events.EBAEvent:
|
||||
"""Build a synthetic platform event for route validation."""
|
||||
payload = payload or {}
|
||||
now = time.time()
|
||||
common = {
|
||||
'type': event_type,
|
||||
'timestamp': payload.get('timestamp') or now,
|
||||
'adapter_name': payload.get('adapter_name') or 'test-event',
|
||||
'source_platform_object': {'synthetic': True, 'payload': payload},
|
||||
}
|
||||
|
||||
user_id = str(payload.get('user_id') or payload.get('sender_id') or 'test-user')
|
||||
user_name = str(payload.get('user_name') or payload.get('sender_name') or 'Test User')
|
||||
group_id = str(payload.get('group_id') or payload.get('chat_id') or 'test-group')
|
||||
group_name = str(payload.get('group_name') or 'Test Group')
|
||||
|
||||
if event_type == 'message.received':
|
||||
chat_type_value = str(payload.get('chat_type') or 'private')
|
||||
chat_type = (
|
||||
platform_entities.ChatType.GROUP
|
||||
if chat_type_value == platform_entities.ChatType.GROUP.value
|
||||
else platform_entities.ChatType.PRIVATE
|
||||
)
|
||||
chat_id = str(
|
||||
payload.get('chat_id') or (group_id if chat_type == platform_entities.ChatType.GROUP else user_id)
|
||||
)
|
||||
message_text = str(payload.get('message_text') or payload.get('text') or '')
|
||||
message_chain_data = payload.get('message_chain')
|
||||
if message_chain_data is None:
|
||||
message_chain = platform_message.MessageChain([platform_message.Plain(text=message_text)])
|
||||
else:
|
||||
message_chain = platform_message.MessageChain.model_validate(message_chain_data)
|
||||
group = (
|
||||
platform_entities.UserGroup(id=chat_id, name=group_name)
|
||||
if chat_type == platform_entities.ChatType.GROUP
|
||||
else None
|
||||
)
|
||||
return platform_events.MessageReceivedEvent(
|
||||
**common,
|
||||
message_id=str(payload.get('message_id') or f'test-message:{uuid.uuid4()}'),
|
||||
message_chain=message_chain,
|
||||
sender=platform_entities.User(id=user_id, nickname=user_name),
|
||||
chat_type=chat_type,
|
||||
chat_id=chat_id,
|
||||
group=group,
|
||||
)
|
||||
|
||||
if event_type == 'group.member_joined':
|
||||
return platform_events.MemberJoinedEvent(
|
||||
**common,
|
||||
group=platform_entities.UserGroup(id=group_id, name=group_name),
|
||||
member=platform_entities.User(id=user_id, nickname=user_name),
|
||||
inviter=platform_entities.User(
|
||||
id=str(payload.get('inviter_id')),
|
||||
nickname=str(payload.get('inviter_name') or ''),
|
||||
)
|
||||
if payload.get('inviter_id')
|
||||
else None,
|
||||
join_type=payload.get('join_type'),
|
||||
)
|
||||
|
||||
if event_type == 'group.member_left':
|
||||
return platform_events.MemberLeftEvent(
|
||||
**common,
|
||||
group=platform_entities.UserGroup(id=group_id, name=group_name),
|
||||
member=platform_entities.User(id=user_id, nickname=user_name),
|
||||
is_kicked=bool(payload.get('is_kicked', False)),
|
||||
operator=platform_entities.User(
|
||||
id=str(payload.get('operator_id')),
|
||||
nickname=str(payload.get('operator_name') or ''),
|
||||
)
|
||||
if payload.get('operator_id')
|
||||
else None,
|
||||
)
|
||||
|
||||
if event_type == 'platform.specific':
|
||||
return platform_events.PlatformSpecificEvent(
|
||||
**common,
|
||||
action=str(payload.get('action') or 'test'),
|
||||
data=payload.get('data') if isinstance(payload.get('data'), dict) else payload,
|
||||
)
|
||||
|
||||
return platform_events.EBAEvent(**common)
|
||||
|
||||
async def dispatch_test_event(
|
||||
self,
|
||||
event_type: str,
|
||||
payload: dict[str, typing.Any] | None = None,
|
||||
) -> dict[str, typing.Any]:
|
||||
"""Dispatch a synthetic event through the real runtime route path."""
|
||||
event_type = str(event_type or '').strip()
|
||||
if not event_type:
|
||||
raise ValueError('event_type is required')
|
||||
|
||||
event = self._build_test_platform_event(event_type, payload)
|
||||
await self._record_event_route_trace(
|
||||
event_type=event_type,
|
||||
status='test_started',
|
||||
reason='Synthetic test event dispatched from control plane',
|
||||
text=f'Test event {event_type} dispatched from control plane',
|
||||
)
|
||||
test_adapter = SyntheticRouteTestAdapter(self.adapter)
|
||||
outcome = await self._dispatch_eba_event_to_processor(
|
||||
event,
|
||||
typing.cast(abstract_platform_adapter.AbstractMessagePlatformAdapter, test_adapter),
|
||||
)
|
||||
return {
|
||||
'event_type': event_type,
|
||||
'dispatched': outcome['status'] in {'delivered', 'discarded'},
|
||||
'status': outcome['status'],
|
||||
'binding_id': outcome.get('binding_id'),
|
||||
'failure_code': outcome.get('failure_code'),
|
||||
'reason': outcome.get('reason'),
|
||||
'suppressed_outputs': test_adapter.suppressed_outputs,
|
||||
}
|
||||
|
||||
async def _record_event_route_trace(
|
||||
self,
|
||||
*,
|
||||
@@ -1170,7 +830,7 @@ class RuntimeBot:
|
||||
self,
|
||||
envelope: AgentEventEnvelope,
|
||||
outputs: list[provider_message.Message | provider_message.MessageChunk],
|
||||
adapter: abstract_platform_adapter.AbstractMessagePlatformAdapter | SyntheticRouteTestAdapter | None = None,
|
||||
adapter: abstract_platform_adapter.AbstractMessagePlatformAdapter | None = None,
|
||||
) -> None:
|
||||
if not outputs or not envelope.delivery.reply_target:
|
||||
return
|
||||
|
||||
Reference in New Issue
Block a user