refactor: unify workflow data port access

This commit is contained in:
jxxghp
2026-08-24 04:11:09 +08:00
parent 8555d0da3e
commit 48e796cab7
9 changed files with 52 additions and 17 deletions
+6 -6
View File
@@ -18,7 +18,7 @@ from app.chain import ChainBase
from app.runtime.config import global_vars
from app.runtime.events import Event, eventmanager
from app.application.workflow import get_workflow_manager
from app.application.chain.data import WorkflowPortProxy as WorkflowOper
from app.application.chain.data import get_chain_workflow_port
from app.runtime.log import logger
from app.schemas.workflow import ActionContext
from app.schemas.workflow import ActionFlow
@@ -1185,13 +1185,13 @@ class WorkflowChain(ChainBase):
:param from_begin: 是否从头开始,默认为True
:param progress_callback: 定时服务进度更新回调
"""
workflowoper = WorkflowOper()
workflowoper = get_chain_workflow_port()
def save_step(action: Action, context: ActionContext, execution_state: dict, completed: bool):
"""
保存上下文到数据库
"""
WorkflowOper().step(
get_chain_workflow_port().step(
workflow_id,
action_id=action.id if completed else "",
context=_serialize_workflow_context(context),
@@ -1263,18 +1263,18 @@ class WorkflowChain(ChainBase):
"""
获取工作流列表
"""
return WorkflowOper().list_enabled()
return get_chain_workflow_port().list_enabled()
@staticmethod
def get_timer_workflows() -> List[Workflow]:
"""
获取定时触发的工作流列表
"""
return WorkflowOper().get_timer_triggered_workflows()
return get_chain_workflow_port().get_timer_triggered_workflows()
@staticmethod
def get_event_workflows() -> List[Workflow]:
"""
获取事件触发的工作流列表
"""
return WorkflowOper().get_event_triggered_workflows()
return get_chain_workflow_port().get_event_triggered_workflows()
+4 -4
View File
@@ -6,7 +6,7 @@ from pydantic import BaseModel
from app.runtime.config import global_vars
from app.runtime.events import eventmanager, Event
from app.application.chain.data import WorkflowPortProxy as WorkflowOper
from app.application.chain.data import get_chain_workflow_port
from app.foundation.reflection import ModuleHelper
from app.runtime.log import logger
from app.schemas.workflow import ActionContext
@@ -282,11 +282,11 @@ class WorkFlowManager(metaclass=Singleton):
"""
workflows = []
if workflow_id:
workflow = WorkflowOper().get(workflow_id)
workflow = get_chain_workflow_port().get(workflow_id)
if workflow:
workflows = [workflow]
else:
workflows = WorkflowOper().get_event_triggered_workflows()
workflows = get_chain_workflow_port().get_event_triggered_workflows()
try:
for workflow in workflows:
self.update_workflow_event(workflow)
@@ -359,7 +359,7 @@ class WorkFlowManager(metaclass=Singleton):
"""
try:
# 检查工作流是否存在且启用
workflow = WorkflowOper().get(workflow_id)
workflow = get_chain_workflow_port().get(workflow_id)
if not workflow or workflow.state == 'P':
return
+2 -2
View File
@@ -3,7 +3,7 @@ from app.chain.subscribe import SubscribeChain
from app.application.configuration import get_chain_runtime_config_snapshot
from app.runtime.config import global_vars
from app.domain.context import MediaInfo
from app.application.chain.data import SubscribePortProxy as SubscribeOper
from app.application.chain.data import get_chain_subscribe_port
from app.runtime.log import logger
from app.schemas.workflow import ActionParams
from app.schemas.workflow import ActionContext
@@ -75,7 +75,7 @@ class AddSubscribeAction(BaseAction):
if self._added_subscribes:
logger.info(f"已添加 {len(self._added_subscribes)} 个订阅")
for sid in self._added_subscribes:
context.subscribes.append(SubscribeOper().get(sid))
context.subscribes.append(get_chain_subscribe_port().get(sid))
elif _started:
self._has_error = True
+2 -2
View File
@@ -6,7 +6,7 @@ from pydantic import Field
from app.workflow.actions import BaseAction
from app.runtime.config import global_vars
from app.application.chain.data import TransferHistoryPortProxy as TransferHistoryOper
from app.application.chain.data import get_chain_transfer_history_port
from app.schemas.workflow import ActionParams
from app.schemas.workflow import ActionContext
from app.chain.storage import StorageChain
@@ -67,7 +67,7 @@ class TransferFileAction(BaseAction):
_failed_count = 0
storagechain = StorageChain()
transferchain = TransferChain()
transferhis = TransferHistoryOper()
transferhis = get_chain_transfer_history_port()
if params.source == "downloads":
# 从下载任务中整理文件
for download in context.downloads: