fix workflow helper

This commit is contained in:
jxxghp
2025-07-31 07:17:05 +08:00
parent 749aaeb003
commit d8f10e9ac4
+117 -87
View File
@@ -8,6 +8,7 @@ from app.log import logger
from app.utils.http import RequestUtils, AsyncRequestUtils from app.utils.http import RequestUtils, AsyncRequestUtils
from app.utils.singleton import WeakSingleton from app.utils.singleton import WeakSingleton
from app.utils.system import SystemUtils from app.utils.system import SystemUtils
from db.models import Workflow
class WorkflowHelper(metaclass=WeakSingleton): class WorkflowHelper(metaclass=WeakSingleton):
@@ -28,170 +29,198 @@ class WorkflowHelper(metaclass=WeakSingleton):
def __init__(self): def __init__(self):
self.get_user_uuid() self.get_user_uuid()
def workflow_share(self, workflow_id: int, @staticmethod
share_title: str, share_comment: str, share_user: str) -> Tuple[bool, str]: def _check_workflow_share_enabled() -> Tuple[bool, str]:
""" """
分享工作流 检查工作流分享功能是否开启
""" """
if not settings.WORKFLOW_STATISTIC_SHARE: # 使用独立的工作流分享开关 if not settings.WORKFLOW_STATISTIC_SHARE:
return False, "当前没有开启工作流数据共享功能" return False, "当前没有开启工作流数据共享功能"
return True, ""
# 获取工作流信息 @staticmethod
workflow = WorkflowOper().get(workflow_id) def _validate_workflow(workflow: Workflow) -> Tuple[bool, str]:
"""
验证工作流是否可以分享
"""
if not workflow: if not workflow:
return False, "工作流不存在" return False, "工作流不存在"
if not workflow.actions or not workflow.flows: if not workflow.actions or not workflow.flows:
return False, "请分享有动作和流程的工作流" return False, "请分享有动作和流程的工作流"
return True, ""
@staticmethod
def _prepare_workflow_data(workflow: Workflow) -> dict:
"""
准备工作流分享数据
"""
workflow_dict = workflow.to_dict() workflow_dict = workflow.to_dict()
workflow_dict.pop("id", None) workflow_dict.pop("id", None)
workflow_dict.pop("context", None) workflow_dict.pop("context", None)
workflow_dict['actions'] = json.dumps(workflow_dict['actions'] or []) workflow_dict['actions'] = json.dumps(workflow_dict['actions'] or [])
workflow_dict['flows'] = json.dumps(workflow_dict['flows'] or []) workflow_dict['flows'] = json.dumps(workflow_dict['flows'] or [])
return workflow_dict
def _build_share_payload(self, share_title: str, share_comment: str,
share_user: str, workflow_dict: dict) -> dict:
"""
构建分享请求载荷
"""
return {
"share_title": share_title,
"share_comment": share_comment,
"share_user": share_user,
"share_uid": self._share_user_id,
**workflow_dict
}
def _handle_response(self, res, clear_cache: bool = True) -> Tuple[bool, str]:
"""
处理HTTP响应
"""
if res is None:
return False, "连接MoviePilot服务器失败"
# 检查响应状态
success = True if res.status_code == 200 else False
if success:
# 清除缓存
if clear_cache:
cache_backend.clear(region=self._shares_cache_region)
return True, ""
else:
return False, res.json().get("message")
def workflow_share(self, workflow_id: int,
share_title: str, share_comment: str, share_user: str) -> Tuple[bool, str]:
"""
分享工作流
"""
# 检查功能是否开启
enabled, message = self._check_workflow_share_enabled()
if not enabled:
return False, message
# 获取工作流信息
workflow = WorkflowOper().get(workflow_id)
# 验证工作流
valid, message = self._validate_workflow(workflow)
if not valid:
return False, message
# 准备数据
workflow_dict = self._prepare_workflow_data(workflow)
payload = self._build_share_payload(share_title, share_comment, share_user, workflow_dict)
# 发送分享请求 # 发送分享请求
res = RequestUtils(proxies=settings.PROXY or {}, res = RequestUtils(proxies=settings.PROXY or {},
content_type="application/json", content_type="application/json",
timeout=10).post(self._workflow_share, timeout=10).post(self._workflow_share, json=payload)
json={
"share_title": share_title, return self._handle_response(res)
"share_comment": share_comment,
"share_user": share_user,
"share_uid": self._share_user_id,
**workflow_dict
})
if res is None:
return False, "连接MoviePilot服务器失败"
if res.ok:
# 清除 get_shares 的缓存,以便实时看到结果
cache_backend.clear(region=self._shares_cache_region)
return True, ""
else:
return False, res.json().get("message")
async def async_workflow_share(self, workflow_id: int, async def async_workflow_share(self, workflow_id: int,
share_title: str, share_comment: str, share_user: str) -> Tuple[bool, str]: share_title: str, share_comment: str, share_user: str) -> Tuple[bool, str]:
""" """
异步分享工作流 异步分享工作流
""" """
if not settings.WORKFLOW_STATISTIC_SHARE: # 使用独立的工作流分享开关 # 检查功能是否开启
return False, "当前没有开启工作流数据共享功能" enabled, message = self._check_workflow_share_enabled()
if not enabled:
return False, message
# 获取工作流信息 # 获取工作流信息
workflow = await WorkflowOper().async_get(workflow_id) workflow = await WorkflowOper().async_get(workflow_id)
if not workflow:
return False, "工作流不存在"
if not workflow.actions or not workflow.flows: # 验证工作流
return False, "请分享有动作和流程的工作流" valid, message = self._validate_workflow(workflow)
if not valid:
return False, message
workflow_dict = workflow.to_dict() # 准备数据
workflow_dict.pop("id", None) workflow_dict = self._prepare_workflow_data(workflow)
workflow_dict.pop("context", None) payload = self._build_share_payload(share_title, share_comment, share_user, workflow_dict)
workflow_dict['actions'] = json.dumps(workflow_dict['actions'] or [])
workflow_dict['flows'] = json.dumps(workflow_dict['flows'] or [])
# 发送分享请求 # 发送分享请求
res = await AsyncRequestUtils(proxies=settings.PROXY or {}, res = await AsyncRequestUtils(proxies=settings.PROXY or {},
content_type="application/json", content_type="application/json",
timeout=10).post(self._workflow_share, timeout=10).post(self._workflow_share, json=payload)
json={
"share_title": share_title, return self._handle_response(res)
"share_comment": share_comment,
"share_user": share_user,
"share_uid": self._share_user_id,
**workflow_dict
})
if res is None:
return False, "连接MoviePilot服务器失败"
if res.status_code == 200:
# 清除 get_shares 的缓存,以便实时看到结果
cache_backend.clear(region=self._shares_cache_region)
return True, ""
else:
return False, res.json().get("message")
def share_delete(self, share_id: int) -> Tuple[bool, str]: def share_delete(self, share_id: int) -> Tuple[bool, str]:
""" """
删除分享 删除分享
""" """
if not settings.WORKFLOW_STATISTIC_SHARE: # 使用独立的工作流分享开关 # 检查功能是否开启
return False, "当前没有开启工作流数据共享功能" enabled, message = self._check_workflow_share_enabled()
if not enabled:
return False, message
res = RequestUtils(proxies=settings.PROXY or {}, res = RequestUtils(proxies=settings.PROXY or {},
timeout=5).delete_res(f"{self._workflow_share}/{share_id}", timeout=5).delete_res(f"{self._workflow_share}/{share_id}",
params={"share_uid": self._share_user_id}) params={"share_uid": self._share_user_id})
if res is None:
return False, "连接MoviePilot服务器失败" return self._handle_response(res)
if res.ok:
# 清除 get_shares 的缓存,以便实时看到结果
cache_backend.clear(region=self._shares_cache_region)
return True, ""
else:
return False, res.json().get("message")
async def async_share_delete(self, share_id: int) -> Tuple[bool, str]: async def async_share_delete(self, share_id: int) -> Tuple[bool, str]:
""" """
异步删除分享 异步删除分享
""" """
if not settings.WORKFLOW_STATISTIC_SHARE: # 使用独立的工作流分享开关 # 检查功能是否开启
return False, "当前没有开启工作流数据共享功能" enabled, message = self._check_workflow_share_enabled()
if not enabled:
return False, message
res = await AsyncRequestUtils(proxies=settings.PROXY or {}, res = await AsyncRequestUtils(proxies=settings.PROXY or {},
timeout=5).delete_res(f"{self._workflow_share}/{share_id}", timeout=5).delete_res(f"{self._workflow_share}/{share_id}",
params={"share_uid": self._share_user_id}) params={"share_uid": self._share_user_id})
if res is None:
return False, "连接MoviePilot服务器失败" return self._handle_response(res)
if res.status_code == 200:
# 清除 get_shares 的缓存,以便实时看到结果
cache_backend.clear(region=self._shares_cache_region)
return True, ""
else:
return False, res.json().get("message")
def workflow_fork(self, share_id: int) -> Tuple[bool, str]: def workflow_fork(self, share_id: int) -> Tuple[bool, str]:
""" """
复用分享的工作流 复用分享的工作流
""" """
if not settings.WORKFLOW_STATISTIC_SHARE: # 使用独立的工作流分享开关 # 检查功能是否开启
return False, "当前没有开启工作流数据共享功能" enabled, message = self._check_workflow_share_enabled()
if not enabled:
return False, message
res = RequestUtils(proxies=settings.PROXY or {}, timeout=5, headers={ res = RequestUtils(proxies=settings.PROXY or {}, timeout=5, headers={
"Content-Type": "application/json" "Content-Type": "application/json"
}).get_res(self._workflow_fork % share_id) }).get_res(self._workflow_fork % share_id)
if res is None:
return False, "连接MoviePilot服务器失败" return self._handle_response(res, clear_cache=False)
if res.ok:
return True, ""
else:
return False, res.json().get("message")
async def async_workflow_fork(self, share_id: int) -> Tuple[bool, str]: async def async_workflow_fork(self, share_id: int) -> Tuple[bool, str]:
""" """
异步复用分享的工作流 异步复用分享的工作流
""" """
if not settings.WORKFLOW_STATISTIC_SHARE: # 使用独立的工作流分享开关 # 检查功能是否开启
return False, "当前没有开启工作流数据共享功能" enabled, message = self._check_workflow_share_enabled()
if not enabled:
return False, message
res = await AsyncRequestUtils(proxies=settings.PROXY or {}, res = await AsyncRequestUtils(proxies=settings.PROXY or {},
timeout=5, timeout=5,
headers={ headers={
"Content-Type": "application/json" "Content-Type": "application/json"
}).get_res(self._workflow_fork % share_id) }).get_res(self._workflow_fork % share_id)
if res is None:
return False, "连接MoviePilot服务器失败" return self._handle_response(res, clear_cache=False)
if res.status_code == 200:
return True, ""
else:
return False, res.json().get("message")
@cached(region=_shares_cache_region, maxsize=1, skip_empty=True) @cached(region=_shares_cache_region, maxsize=1, skip_empty=True)
def get_shares(self, name: Optional[str] = None, page: Optional[int] = 1, count: Optional[int] = 30) -> List[dict]: def get_shares(self, name: Optional[str] = None, page: Optional[int] = 1, count: Optional[int] = 30) -> List[dict]:
""" """
获取工作流分享数据 获取工作流分享数据
""" """
if not settings.WORKFLOW_STATISTIC_SHARE: # 使用独立的工作流分享开关 enabled, _ = self._check_workflow_share_enabled()
if not enabled:
return [] return []
res = RequestUtils(proxies=settings.PROXY or {}, timeout=15).get_res(self._workflow_shares, params={ res = RequestUtils(proxies=settings.PROXY or {}, timeout=15).get_res(self._workflow_shares, params={
@@ -209,7 +238,8 @@ class WorkflowHelper(metaclass=WeakSingleton):
""" """
异步获取工作流分享数据 异步获取工作流分享数据
""" """
if not settings.WORKFLOW_STATISTIC_SHARE: # 使用独立的工作流分享开关 enabled, _ = self._check_workflow_share_enabled()
if not enabled:
return [] return []
res = await AsyncRequestUtils(proxies=settings.PROXY or {}, timeout=15).get_res(self._workflow_shares, params={ res = await AsyncRequestUtils(proxies=settings.PROXY or {}, timeout=15).get_res(self._workflow_shares, params={