refactor: switch llm_entities to plugin sdk

This commit is contained in:
Junyan Qin
2025-07-13 20:30:17 +08:00
parent 4a319b2b20
commit 6a1de889b4
15 changed files with 76 additions and 378 deletions
+8 -6
View File
@@ -7,8 +7,8 @@ import dashscope
from .. import runner
from ...core import app
from .. import entities as llm_entities
import langbot_plugin.api.entities.builtin.pipeline.query as pipeline_query
import langbot_plugin.api.entities.builtin.provider.message as provider_message
class DashscopeAPIError(Exception):
@@ -90,7 +90,9 @@ class DashScopeAPIRunner(runner.RequestRunner):
return plain_text, image_ids
async def _agent_messages(self, query: pipeline_query.Query) -> typing.AsyncGenerator[llm_entities.Message, None]:
async def _agent_messages(
self, query: pipeline_query.Query
) -> typing.AsyncGenerator[provider_message.Message, None]:
"""Dashscope 智能体对话请求"""
# 局部变量
@@ -143,14 +145,14 @@ class DashScopeAPIRunner(runner.RequestRunner):
# 将参考资料替换到文本中
pending_content = self._replace_references(pending_content, references_dict)
yield llm_entities.Message(
yield provider_message.Message(
role='assistant',
content=pending_content,
)
async def _workflow_messages(
self, query: pipeline_query.Query
) -> typing.AsyncGenerator[llm_entities.Message, None]:
) -> typing.AsyncGenerator[provider_message.Message, None]:
"""Dashscope 工作流对话请求"""
# 局部变量
@@ -208,12 +210,12 @@ class DashScopeAPIRunner(runner.RequestRunner):
# 将参考资料替换到文本中
pending_content = self._replace_references(pending_content, references_dict)
yield llm_entities.Message(
yield provider_message.Message(
role='assistant',
content=pending_content,
)
async def run(self, query: pipeline_query.Query) -> typing.AsyncGenerator[llm_entities.Message, None]:
async def run(self, query: pipeline_query.Query) -> typing.AsyncGenerator[provider_message.Message, None]:
"""运行"""
if self.app_type == 'agent':
async for msg in self._agent_messages(query):
+19 -17
View File
@@ -9,7 +9,7 @@ import base64
from .. import runner
from ...core import app
from .. import entities as llm_entities
import langbot_plugin.api.entities.builtin.provider.message as provider_message
from ...utils import image
import langbot_plugin.api.entities.builtin.pipeline.query as pipeline_query
from libs.dify_service_api.v1 import client, errors
@@ -90,7 +90,9 @@ class DifyServiceAPIRunner(runner.RequestRunner):
return plain_text, image_ids
async def _chat_messages(self, query: pipeline_query.Query) -> typing.AsyncGenerator[llm_entities.Message, None]:
async def _chat_messages(
self, query: pipeline_query.Query
) -> typing.AsyncGenerator[provider_message.Message, None]:
"""调用聊天助手"""
cov_id = query.session.using_conversation.uuid or ''
query.variables['conversation_id'] = cov_id
@@ -132,7 +134,7 @@ class DifyServiceAPIRunner(runner.RequestRunner):
if mode == 'workflow':
if chunk['event'] == 'node_finished':
if chunk['data']['node_type'] == 'answer':
yield llm_entities.Message(
yield provider_message.Message(
role='assistant',
content=self._try_convert_thinking(chunk['data']['outputs']['answer']),
)
@@ -140,7 +142,7 @@ class DifyServiceAPIRunner(runner.RequestRunner):
if chunk['event'] == 'message':
basic_mode_pending_chunk += chunk['answer']
elif chunk['event'] == 'message_end':
yield llm_entities.Message(
yield provider_message.Message(
role='assistant',
content=self._try_convert_thinking(basic_mode_pending_chunk),
)
@@ -153,7 +155,7 @@ class DifyServiceAPIRunner(runner.RequestRunner):
async def _agent_chat_messages(
self, query: pipeline_query.Query
) -> typing.AsyncGenerator[llm_entities.Message, None]:
) -> typing.AsyncGenerator[provider_message.Message, None]:
"""调用聊天助手"""
cov_id = query.session.using_conversation.uuid or ''
query.variables['conversation_id'] = cov_id
@@ -198,7 +200,7 @@ class DifyServiceAPIRunner(runner.RequestRunner):
else:
if pending_agent_message.strip() != '':
pending_agent_message = pending_agent_message.replace('</details>Action:', '</details>')
yield llm_entities.Message(
yield provider_message.Message(
role='assistant',
content=self._try_convert_thinking(pending_agent_message),
)
@@ -209,13 +211,13 @@ class DifyServiceAPIRunner(runner.RequestRunner):
continue
if chunk['tool']:
msg = llm_entities.Message(
msg = provider_message.Message(
role='assistant',
tool_calls=[
llm_entities.ToolCall(
provider_message.ToolCall(
id=chunk['id'],
type='function',
function=llm_entities.FunctionCall(
function=provider_message.FunctionCall(
name=chunk['tool'],
arguments=json.dumps({}),
),
@@ -232,9 +234,9 @@ class DifyServiceAPIRunner(runner.RequestRunner):
image_url = base_url + chunk['url']
yield llm_entities.Message(
yield provider_message.Message(
role='assistant',
content=[llm_entities.ContentElement.from_image_url(image_url)],
content=[provider_message.ContentElement.from_image_url(image_url)],
)
if chunk['event'] == 'error':
raise errors.DifyAPIError('dify 服务错误: ' + chunk['message'])
@@ -246,7 +248,7 @@ class DifyServiceAPIRunner(runner.RequestRunner):
async def _workflow_messages(
self, query: pipeline_query.Query
) -> typing.AsyncGenerator[llm_entities.Message, None]:
) -> typing.AsyncGenerator[provider_message.Message, None]:
"""调用工作流"""
if not query.session.using_conversation.uuid:
@@ -290,14 +292,14 @@ class DifyServiceAPIRunner(runner.RequestRunner):
if chunk['data']['node_type'] == 'start' or chunk['data']['node_type'] == 'end':
continue
msg = llm_entities.Message(
msg = provider_message.Message(
role='assistant',
content=None,
tool_calls=[
llm_entities.ToolCall(
provider_message.ToolCall(
id=chunk['data']['node_id'],
type='function',
function=llm_entities.FunctionCall(
function=provider_message.FunctionCall(
name=chunk['data']['title'],
arguments=json.dumps({}),
),
@@ -311,14 +313,14 @@ class DifyServiceAPIRunner(runner.RequestRunner):
if chunk['data']['error']:
raise errors.DifyAPIError(chunk['data']['error'])
msg = llm_entities.Message(
msg = provider_message.Message(
role='assistant',
content=chunk['data']['outputs']['summary'],
)
yield msg
async def run(self, query: pipeline_query.Query) -> typing.AsyncGenerator[llm_entities.Message, None]:
async def run(self, query: pipeline_query.Query) -> typing.AsyncGenerator[provider_message.Message, None]:
"""运行请求"""
if self.pipeline_config['ai']['dify-service-api']['app-type'] == 'chat':
async for msg in self._chat_messages(query):
+4 -4
View File
@@ -4,15 +4,15 @@ import json
import typing
from .. import runner
from .. import entities as llm_entities
import langbot_plugin.api.entities.builtin.pipeline.query as pipeline_query
import langbot_plugin.api.entities.builtin.provider.message as provider_message
@runner.runner_class('local-agent')
class LocalAgentRunner(runner.RequestRunner):
"""本地Agent请求运行器"""
async def run(self, query: pipeline_query.Query) -> typing.AsyncGenerator[llm_entities.Message, None]:
async def run(self, query: pipeline_query.Query) -> typing.AsyncGenerator[provider_message.Message, None]:
"""运行请求"""
pending_tool_calls = []
@@ -45,7 +45,7 @@ class LocalAgentRunner(runner.RequestRunner):
func_ret = await self.ap.tool_mgr.execute_func_call(func.name, parameters)
msg = llm_entities.Message(
msg = provider_message.Message(
role='tool',
content=json.dumps(func_ret, ensure_ascii=False),
tool_call_id=tool_call.id,
@@ -56,7 +56,7 @@ class LocalAgentRunner(runner.RequestRunner):
req_messages.append(msg)
except Exception as e:
# 工具调用出错,添加一个报错信息到 req_messages
err_msg = llm_entities.Message(role='tool', content=f'err: {e}', tool_call_id=tool_call.id)
err_msg = provider_message.Message(role='tool', content=f'err: {e}', tool_call_id=tool_call.id)
yield err_msg
+4 -4
View File
@@ -7,8 +7,8 @@ import aiohttp
from .. import runner
from ...core import app
from .. import entities as llm_entities
import langbot_plugin.api.entities.builtin.pipeline.query as pipeline_query
import langbot_plugin.api.entities.builtin.provider.message as provider_message
class N8nAPIError(Exception):
@@ -68,7 +68,7 @@ class N8nServiceAPIRunner(runner.RequestRunner):
return plain_text
async def _call_webhook(self, query: pipeline_query.Query) -> typing.AsyncGenerator[llm_entities.Message, None]:
async def _call_webhook(self, query: pipeline_query.Query) -> typing.AsyncGenerator[provider_message.Message, None]:
"""调用n8n webhook"""
# 生成会话ID(如果不存在)
if not query.session.using_conversation.uuid:
@@ -146,7 +146,7 @@ class N8nServiceAPIRunner(runner.RequestRunner):
output_content = json.dumps(response_data, ensure_ascii=False)
# 返回消息
yield llm_entities.Message(
yield provider_message.Message(
role='assistant',
content=output_content,
)
@@ -154,7 +154,7 @@ class N8nServiceAPIRunner(runner.RequestRunner):
self.ap.logger.error(f'n8n webhook call exception: {str(e)}')
raise N8nAPIError(f'n8n webhook call exception: {str(e)}')
async def run(self, query: pipeline_query.Query) -> typing.AsyncGenerator[llm_entities.Message, None]:
async def run(self, query: pipeline_query.Query) -> typing.AsyncGenerator[provider_message.Message, None]:
"""运行请求"""
async for msg in self._call_webhook(query):
yield msg