-
Notifications
You must be signed in to change notification settings - Fork 2.8k
Expand file tree
/
Copy pathtool_workflow.py
More file actions
294 lines (257 loc) · 13.9 KB
/
tool_workflow.py
File metadata and controls
294 lines (257 loc) · 13.9 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
# coding=utf-8
"""
@project: MaxKB
@Author:虎虎
@file: tool_workflow.py
@date:2026/3/6 13:59
@desc:
"""
import asyncio
import json
import os
# coding=utf-8
import pickle
import tempfile
import zipfile
from functools import reduce
from typing import Dict, List
import requests
import uuid_utils.compat as uuid
from django.db import transaction
from django.db.models import QuerySet
from django.http import HttpResponse
from django.utils import timezone
from django.utils.translation import gettext_lazy as _
from rest_framework import serializers, status
from application.flow.common import Workflow, WorkflowMode
from application.flow.i_step_node import ToolWorkflowPostHandler
from application.flow.tool_workflow_manage import ToolWorkflowManage
from application.models import ChatRecord
from application.serializers.application import McpServersSerializer, get_mcp_tools
from application.serializers.common import ToolExecute
from common.exception.app_exception import AppApiException
from common.field.common import UploadedFileField
from common.result import result
from common.utils.common import bytes_to_uploaded_file
from common.utils.common import restricted_loads, generate_uuid
from common.utils.logger import maxkb_logger
from common.utils.tool_code import ToolExecutor
from knowledge.models import KnowledgeWorkflow
from system_manage.models import AuthTargetType
from system_manage.serializers.user_resource_permission import UserResourcePermissionSerializer
from tools.models import Tool, ToolScope, ToolWorkflow, ToolWorkflowVersion
from tools.serializers.tool import ToolExportModelSerializer, ToolSerializer
from users.models import User
tool_executor = ToolExecutor()
def hand_node(node, update_tool_map):
if node.get('type') == 'tool-lib-node':
tool_lib_id = (node.get('properties', {}).get('node_data', {}).get('tool_lib_id') or '')
node.get('properties', {}).get('node_data', {})['tool_lib_id'] = update_tool_map.get(tool_lib_id, tool_lib_id)
if node.get('type') == 'search-knowledge-node':
node.get('properties', {}).get('node_data', {})['knowledge_id_list'] = []
if node.get('type') == 'ai-chat-node':
node_data = node.get('properties', {}).get('node_data', {})
mcp_tool_ids = node_data.get('mcp_tool_ids') or []
node_data['mcp_tool_ids'] = [update_tool_map.get(tool_id,
tool_id) for tool_id in mcp_tool_ids]
tool_ids = node_data.get('tool_ids') or []
node_data['tool_ids'] = [update_tool_map.get(tool_id,
tool_id) for tool_id in tool_ids]
if node.get('type') == 'mcp-node':
mcp_tool_id = (node.get('properties', {}).get('node_data', {}).get('mcp_tool_id') or '')
node.get('properties', {}).get('node_data', {})['mcp_tool_id'] = update_tool_map.get(mcp_tool_id,
mcp_tool_id)
class ToolWorkflowModelSerializer(serializers.ModelSerializer):
class Meta:
model = ToolWorkflow
fields = '__all__'
class ToolWorkflowImportRequest(serializers.Serializer):
file = UploadedFileField(required=True, label=_("file"))
class ToolWorkflowActionListQuerySerializer(serializers.Serializer):
user_name = serializers.CharField(required=False, label=_('Name'), allow_blank=True, allow_null=True)
state = serializers.CharField(required=False, label=_("State"), allow_blank=True, allow_null=True)
class ToolWorkflowInstance:
def __init__(self, knowledge_workflow: dict, version: str, tool_list: List[dict]):
self.knowledge_workflow = knowledge_workflow
self.version = version
self.tool_list = tool_list
def get_tool_list(self):
return self.tool_list or []
class ToolWorkflowSerializer(serializers.Serializer):
class Operate(serializers.Serializer):
user_id = serializers.UUIDField(required=True, label=_('user id'))
workspace_id = serializers.CharField(required=True, label=_('workspace id'))
tool_id = serializers.UUIDField(required=True, label=_('tool id'))
def debug(self, instance: Dict, user, with_valid=True):
if with_valid:
self.is_valid(raise_exception=True)
tool_workflow = QuerySet(ToolWorkflow).filter(tool_id=self.data.get("tool_id")).first()
tool_record_id = instance.get('chat_record_id') or str(uuid.uuid7())
took_execute = ToolExecute(self.data.get("tool_id"), tool_record_id,
self.data.get("workspace_id"),
None,
None,
True)
record = took_execute.get_record()
work_flow_manage = ToolWorkflowManage(
Workflow.new_instance(tool_workflow.work_flow, WorkflowMode.TOOL),
{
'chat_record_id': tool_record_id,
'tool_id': self.data.get("tool_id"),
'stream': True,
'workspace_id': self.data.get("workspace_id"),
**instance},
ToolWorkflowPostHandler(took_execute, self.data.get("tool_id")),
is_the_task_interrupted=lambda: False,
child_node=instance.get('child_node'),
start_node_id=instance.get('runtime_node_id'),
start_node_data=instance.get('node_data'),
chat_record=self.to_chat_record(record)
)
r = work_flow_manage.run()
return r
@staticmethod
def to_chat_record(record):
if record is None:
return None
return ChatRecord(
answer_text_list=record.meta.get('answer_text_list'),
details=record.meta.get('details'),
answer_text='',
)
def publish(self, with_valid=True):
if with_valid:
self.is_valid()
user_id = self.data.get('user_id')
workspace_id = self.data.get("workspace_id")
user = QuerySet(User).filter(id=user_id).first()
tool_workflow = QuerySet(ToolWorkflow).filter(tool_id=self.data.get("tool_id"),
workspace_id=workspace_id).first()
work_flow_version = ToolWorkflowVersion(work_flow=tool_workflow.work_flow,
tool_id=self.data.get("tool_id"),
name=timezone.localtime(timezone.now()).strftime(
'%Y-%m-%d %H:%M:%S'),
publish_user_id=user_id,
publish_user_name=user.username,
workspace_id=workspace_id)
work_flow_version.save()
QuerySet(ToolWorkflow).filter(
tool_id=self.data.get("tool_id")
).update(is_publish=True, publish_time=timezone.now())
return True
def edit(self, instance: Dict):
self.is_valid(raise_exception=True)
if instance.get("work_flow"):
QuerySet(ToolWorkflow).update_or_create(tool_id=self.data.get("tool_id"),
create_defaults={'id': uuid.uuid7(),
'tool_id': self.data.get(
"tool_id"),
"workspace_id": self.data.get(
'workspace_id'),
'work_flow': instance.get('work_flow',
{}), },
defaults={
'tool_id': self.data.get("tool_id"),
'workspace_id': self.data.get(
'workspace_id'),
'work_flow': instance.get('work_flow')
})
return self.one()
if instance.get("work_flow_template"):
template_instance = instance.get('work_flow_template')
download_url = template_instance.get('downloadUrl')
# 查找匹配的版本名称
res = requests.get(download_url, timeout=5)
tool = QuerySet(Tool).filter(id=self.data.get("tool_id")).first()
ToolSerializer.Import(data={
'user_id': self.data.get('user_id'),
'workspace_id': self.data.get('workspace_id'),
'folder_id': tool.folder_id,
'file': bytes_to_uploaded_file(res.content, 'file.tool')
}).update_template_workflow(str(self.data.get('tool_id')))
try:
requests.get(template_instance.get('downloadCallbackUrl'), timeout=5)
except Exception as e:
maxkb_logger.error(f"callback appstore tool download error: {e}")
return self.one()
def one(self):
self.is_valid(raise_exception=True)
workflow = QuerySet(ToolWorkflow).filter(tool_id=self.data.get('tool_id')).first()
return {**ToolWorkflowModelSerializer(workflow).data}
class ToolWorkflowMcpSerializer(serializers.Serializer):
tool_id = serializers.UUIDField(required=True, label=_('Tool id'))
user_id = serializers.UUIDField(required=True, label=_("User ID"))
workspace_id = serializers.CharField(required=False, allow_null=True, allow_blank=True, label=_("Workspace ID"))
def is_valid(self, *, raise_exception=False):
super().is_valid(raise_exception=True)
workspace_id = self.data.get('workspace_id')
query_set = QuerySet(Tool).filter(id=self.data.get('tool_id'))
if workspace_id:
query_set = query_set.filter(workspace_id=workspace_id)
if not query_set.exists():
raise AppApiException(500, _('Tool id does not exist'))
def get_mcp_servers(self, instance, with_valid=True):
if with_valid:
self.is_valid(raise_exception=True)
McpServersSerializer(data=instance).is_valid(raise_exception=True)
servers = json.loads(instance.get('mcp_servers'))
for server, config in servers.items():
if config.get('transport') not in ['sse', 'streamable_http']:
raise AppApiException(500, _('Only support transport=sse or transport=streamable_http'))
tools = []
for server in servers:
tools += [
{
'server': server,
'name': tool.name,
'description': tool.description,
'args_schema': tool.args_schema,
}
for tool in asyncio.run(get_mcp_tools({server: servers[server]}))]
return tools
class StoreToolWorkflow(serializers.Serializer):
user_id = serializers.UUIDField(required=True, label=_("User ID"))
name = serializers.CharField(required=False, label=_("tool name"), allow_null=True, allow_blank=True)
def get_appstore_templates(self):
self.is_valid(raise_exception=True)
# 下载zip文件
try:
res = requests.get('https://apps-assets.fit2cloud.com/stable/maxkb.json.zip', timeout=5)
res.raise_for_status()
# 创建临时文件保存zip
with tempfile.NamedTemporaryFile(delete=False, suffix='.zip') as temp_zip:
temp_zip.write(res.content)
temp_zip_path = temp_zip.name
try:
# 解压zip文件
with zipfile.ZipFile(temp_zip_path, 'r') as zip_ref:
# 获取zip中的第一个文件(假设只有一个json文件)
json_filename = zip_ref.namelist()[0]
json_content = zip_ref.read(json_filename)
# 将json转换为字典
tool_store = json.loads(json_content.decode('utf-8'))
tag_dict = {tag['name']: tag['key'] for tag in tool_store['additionalProperties']['tags']}
filter_apps = []
for tool in tool_store['apps']:
if self.data.get('name', '') != '':
if self.data.get('name').lower() not in tool.get('name', '').lower():
continue
if not tool['downloadUrl'].endswith('.tool') or not [tag_dict[tag] for tag in
tool.get('tags')].__contains__(
'workflow_template'):
continue
versions = tool.get('versions', [])
tool['label'] = tag_dict[tool.get('tags')[0]] if tool.get('tags') else ''
tool['version'] = next(
(version.get('name') for version in versions if
version.get('downloadUrl') == tool['downloadUrl']),
)
filter_apps.append(tool)
tool_store['apps'] = filter_apps
return tool_store
finally:
# 清理临时文件
os.unlink(temp_zip_path)
except Exception as e:
maxkb_logger.error(f"fetch appstore tools error: {e}")
return {'apps': [], 'additionalProperties': {'tags': []}}