mirror of
https://github.com/langbot-app/LangBot.git
synced 2026-07-20 03:16:14 +00:00
Merge pull request #1188 from RockChinQ/feat/query-variables
feat: add query variables
This commit is contained in:
@@ -72,6 +72,9 @@ class Query(pydantic.BaseModel):
|
|||||||
user_message: typing.Optional[llm_entities.Message] = None
|
user_message: typing.Optional[llm_entities.Message] = None
|
||||||
"""此次请求的用户消息对象,由前置处理器阶段设置"""
|
"""此次请求的用户消息对象,由前置处理器阶段设置"""
|
||||||
|
|
||||||
|
variables: typing.Optional[dict[str, typing.Any]] = None
|
||||||
|
"""变量,由前置处理器阶段设置。在prompt中嵌入或由 Runner 传递到 LLMOps 平台。"""
|
||||||
|
|
||||||
use_model: typing.Optional[entities.LLMModelInfo] = None
|
use_model: typing.Optional[entities.LLMModelInfo] = None
|
||||||
"""使用的模型,由前置处理器阶段设置"""
|
"""使用的模型,由前置处理器阶段设置"""
|
||||||
|
|
||||||
@@ -86,10 +89,31 @@ class Query(pydantic.BaseModel):
|
|||||||
|
|
||||||
# ======= 内部保留 =======
|
# ======= 内部保留 =======
|
||||||
current_stage: "pkg.pipeline.stagemgr.StageInstContainer" = None
|
current_stage: "pkg.pipeline.stagemgr.StageInstContainer" = None
|
||||||
|
"""当前所处阶段"""
|
||||||
|
|
||||||
class Config:
|
class Config:
|
||||||
arbitrary_types_allowed = True
|
arbitrary_types_allowed = True
|
||||||
|
|
||||||
|
# ========== 插件可调用的 API(请求 API) ==========
|
||||||
|
|
||||||
|
def set_variable(self, key: str, value: typing.Any):
|
||||||
|
"""设置变量"""
|
||||||
|
if self.variables is None:
|
||||||
|
self.variables = {}
|
||||||
|
self.variables[key] = value
|
||||||
|
|
||||||
|
def get_variable(self, key: str) -> typing.Any:
|
||||||
|
"""获取变量"""
|
||||||
|
if self.variables is None:
|
||||||
|
return None
|
||||||
|
return self.variables.get(key)
|
||||||
|
|
||||||
|
def get_variables(self) -> dict[str, typing.Any]:
|
||||||
|
"""获取所有变量"""
|
||||||
|
if self.variables is None:
|
||||||
|
return {}
|
||||||
|
return self.variables
|
||||||
|
|
||||||
|
|
||||||
class Conversation(pydantic.BaseModel):
|
class Conversation(pydantic.BaseModel):
|
||||||
"""对话,包含于 Session 中,一个 Session 可以有多个历史 Conversation,但只有一个当前使用的 Conversation"""
|
"""对话,包含于 Session 中,一个 Session 可以有多个历史 Conversation,但只有一个当前使用的 Conversation"""
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import datetime
|
||||||
|
|
||||||
from .. import stage, entities, stagemgr
|
from .. import stage, entities, stagemgr
|
||||||
from ...core import entities as core_entities
|
from ...core import entities as core_entities
|
||||||
@@ -34,7 +35,7 @@ class PreProcessor(stage.PipelineStage):
|
|||||||
|
|
||||||
conversation = await self.ap.sess_mgr.get_conversation(session)
|
conversation = await self.ap.sess_mgr.get_conversation(session)
|
||||||
|
|
||||||
# 从会话取出消息和情景预设到query
|
# 设置query
|
||||||
query.session = session
|
query.session = session
|
||||||
query.prompt = conversation.prompt.copy()
|
query.prompt = conversation.prompt.copy()
|
||||||
query.messages = conversation.messages.copy()
|
query.messages = conversation.messages.copy()
|
||||||
@@ -43,6 +44,11 @@ class PreProcessor(stage.PipelineStage):
|
|||||||
|
|
||||||
query.use_funcs = conversation.use_funcs if query.use_model.tool_call_supported else None
|
query.use_funcs = conversation.use_funcs if query.use_model.tool_call_supported else None
|
||||||
|
|
||||||
|
query.variables = {
|
||||||
|
"session_id": f"{query.session.launcher_type.value}_{query.session.launcher_id}",
|
||||||
|
"conversation_id": conversation.uuid,
|
||||||
|
"msg_create_time": int(query.message_event.time) if query.message_event.time else int(datetime.datetime.now().timestamp()),
|
||||||
|
}
|
||||||
|
|
||||||
# 检查vision是否启用,没启用就删除所有图片
|
# 检查vision是否启用,没启用就删除所有图片
|
||||||
if not self.ap.provider_cfg.data['enable-vision'] or (self.ap.provider_cfg.data['runner'] == 'local-agent' and not query.use_model.vision_supported):
|
if not self.ap.provider_cfg.data['enable-vision'] or (self.ap.provider_cfg.data['runner'] == 'local-agent' and not query.use_model.vision_supported):
|
||||||
|
|||||||
@@ -167,6 +167,10 @@ class DashScopeAPIRunner(runner.RequestRunner):
|
|||||||
image_ids = [] # 用户输入的图片ID列表 (暂不支持)
|
image_ids = [] # 用户输入的图片ID列表 (暂不支持)
|
||||||
|
|
||||||
plain_text, image_ids = await self._preprocess_user_message(query)
|
plain_text, image_ids = await self._preprocess_user_message(query)
|
||||||
|
|
||||||
|
biz_params = {}
|
||||||
|
biz_params.update(self.biz_params)
|
||||||
|
biz_params.update(query.variables)
|
||||||
|
|
||||||
#发送对话请求
|
#发送对话请求
|
||||||
response = dashscope.Application.call(
|
response = dashscope.Application.call(
|
||||||
@@ -176,7 +180,7 @@ class DashScopeAPIRunner(runner.RequestRunner):
|
|||||||
stream=True, # 流式输出
|
stream=True, # 流式输出
|
||||||
incremental_output=True, # 增量输出,使用流式输出需要开启增量输出
|
incremental_output=True, # 增量输出,使用流式输出需要开启增量输出
|
||||||
session_id=query.session.using_conversation.uuid, # 会话ID用于,多轮对话
|
session_id=query.session.using_conversation.uuid, # 会话ID用于,多轮对话
|
||||||
biz_params=self.biz_params # 工作流应用的自定义输入参数传递
|
biz_params=biz_params, # 工作流应用的自定义输入参数传递
|
||||||
# rag_options={ # 主要用于文件交互,暂不支持
|
# rag_options={ # 主要用于文件交互,暂不支持
|
||||||
# "session_file_ids": ["FILE_ID1"], # FILE_ID1 替换为实际的临时文件ID,逗号隔开多个
|
# "session_file_ids": ["FILE_ID1"], # FILE_ID1 替换为实际的临时文件ID,逗号隔开多个
|
||||||
# }
|
# }
|
||||||
|
|||||||
@@ -111,8 +111,12 @@ class DifyServiceAPIRunner(runner.RequestRunner):
|
|||||||
|
|
||||||
basic_mode_pending_chunk = ''
|
basic_mode_pending_chunk = ''
|
||||||
|
|
||||||
|
inputs = {}
|
||||||
|
|
||||||
|
inputs.update(query.variables)
|
||||||
|
|
||||||
async for chunk in self.dify_client.chat_messages(
|
async for chunk in self.dify_client.chat_messages(
|
||||||
inputs={},
|
inputs=inputs,
|
||||||
query=plain_text,
|
query=plain_text,
|
||||||
user=f"{query.session.launcher_type.value}_{query.session.launcher_id}",
|
user=f"{query.session.launcher_type.value}_{query.session.launcher_id}",
|
||||||
conversation_id=cov_id,
|
conversation_id=cov_id,
|
||||||
@@ -162,8 +166,12 @@ class DifyServiceAPIRunner(runner.RequestRunner):
|
|||||||
|
|
||||||
ignored_events = ["agent_message"]
|
ignored_events = ["agent_message"]
|
||||||
|
|
||||||
|
inputs = {}
|
||||||
|
|
||||||
|
inputs.update(query.variables)
|
||||||
|
|
||||||
async for chunk in self.dify_client.chat_messages(
|
async for chunk in self.dify_client.chat_messages(
|
||||||
inputs={},
|
inputs=inputs,
|
||||||
query=plain_text,
|
query=plain_text,
|
||||||
user=f"{query.session.launcher_type.value}_{query.session.launcher_id}",
|
user=f"{query.session.launcher_type.value}_{query.session.launcher_id}",
|
||||||
response_mode="streaming",
|
response_mode="streaming",
|
||||||
@@ -227,14 +235,11 @@ class DifyServiceAPIRunner(runner.RequestRunner):
|
|||||||
|
|
||||||
if not query.session.using_conversation.uuid:
|
if not query.session.using_conversation.uuid:
|
||||||
query.session.using_conversation.uuid = str(uuid.uuid4())
|
query.session.using_conversation.uuid = str(uuid.uuid4())
|
||||||
|
|
||||||
cov_id = query.session.using_conversation.uuid
|
query.variables["conversation_id"] = query.session.using_conversation.uuid
|
||||||
|
|
||||||
plain_text, image_ids = await self._preprocess_user_message(query)
|
plain_text, image_ids = await self._preprocess_user_message(query)
|
||||||
|
|
||||||
# 尝试获取 CreateTime
|
|
||||||
create_time = int(query.message_event.time) if query.message_event.time else int(datetime.datetime.now().timestamp())
|
|
||||||
|
|
||||||
files = [
|
files = [
|
||||||
{
|
{
|
||||||
"type": "image",
|
"type": "image",
|
||||||
@@ -246,13 +251,17 @@ class DifyServiceAPIRunner(runner.RequestRunner):
|
|||||||
|
|
||||||
ignored_events = ["text_chunk", "workflow_started"]
|
ignored_events = ["text_chunk", "workflow_started"]
|
||||||
|
|
||||||
|
inputs = { # these variables are legacy variables, we need to keep them for compatibility
|
||||||
|
"langbot_user_message_text": plain_text,
|
||||||
|
"langbot_session_id": query.variables["session_id"],
|
||||||
|
"langbot_conversation_id": query.variables["conversation_id"],
|
||||||
|
"langbot_msg_create_time": query.variables["msg_create_time"],
|
||||||
|
}
|
||||||
|
|
||||||
|
inputs.update(query.variables)
|
||||||
|
|
||||||
async for chunk in self.dify_client.workflow_run(
|
async for chunk in self.dify_client.workflow_run(
|
||||||
inputs={
|
inputs=inputs,
|
||||||
"langbot_user_message_text": plain_text,
|
|
||||||
"langbot_session_id": f"{query.session.launcher_type.value}_{query.session.launcher_id}",
|
|
||||||
"langbot_conversation_id": cov_id,
|
|
||||||
"langbot_msg_create_time": create_time,
|
|
||||||
},
|
|
||||||
user=f"{query.session.launcher_type.value}_{query.session.launcher_id}",
|
user=f"{query.session.launcher_type.value}_{query.session.launcher_id}",
|
||||||
files=files,
|
files=files,
|
||||||
timeout=self.ap.provider_cfg.data["dify-service-api"]["workflow"]["timeout"],
|
timeout=self.ap.provider_cfg.data["dify-service-api"]["workflow"]["timeout"],
|
||||||
|
|||||||
Reference in New Issue
Block a user