From 1eba1c3f8783cb12d5b5de275519240eb14736fa Mon Sep 17 00:00:00 2001 From: ymq Date: Tue, 18 Aug 2026 17:15:21 +0800 Subject: [PATCH] =?UTF-8?q?feat(worker):=20load=5Fpipeline=5Fservice=20?= =?UTF-8?q?=E6=8C=89=20PIPELINE=5FMODE=20=E9=97=A8=E6=8E=A7=20poller=20?= =?UTF-8?q?=E6=B3=A8=E5=86=8C(web=3D=E4=BB=85HTTP=E6=97=A0poller,=20worker?= =?UTF-8?q?=3D=E4=BB=85poller,=20all=3D=E9=BB=98=E8=AE=A4)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pipeline_service/init.py | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/pipeline_service/init.py b/pipeline_service/init.py index db25766..fe2c469 100644 --- a/pipeline_service/init.py +++ b/pipeline_service/init.py @@ -475,6 +475,11 @@ async def skill_import_git(repo_url: str, skills_dir: str = 'skills'): def load_pipeline_service(): """注册所有函数到 ServerEnv。任何宿主应用调用此函数即可使用产线引擎。""" + import os as _os + # 运行模式(进程级,由 start.sh 设 PIPELINE_MODE): + # all(默认) = HTTP + poller;web = 仅 HTTP;worker = 仅 poller(分布式多主机 worker 节点) + _run_mode = (_os.environ.get('PIPELINE_MODE', 'all') or 'all').strip().lower() + _run_pollers = _run_mode in ('all', 'worker') # web 模式不注册 poller env = ServerEnv() # Task lifecycle @@ -642,7 +647,7 @@ def load_pipeline_service(): debug("agent poller started") from ahserver.configuredServer import add_startup - add_startup(_role_poller) + add_startup(_role_poller) if _run_pollers else None # PM review poller (v3.3.0): auto-dispatch review-state tasks to PM agent async def _pm_poller(app): @@ -692,7 +697,7 @@ def load_pipeline_service(): asyncio.create_task(_pm_poll_loop()) debug("pm poller started") - add_startup(_pm_poller) + add_startup(_pm_poller) if _run_pollers else None # QC review poller: auto-dispatch qc_review-state tasks to QC agent(合规/质量门禁) async def _qc_poller(app): @@ -742,7 +747,7 @@ def load_pipeline_service(): asyncio.create_task(_qc_poll_loop()) debug("qc poller started") - add_startup(_qc_poller) + add_startup(_qc_poller) if _run_pollers else None # Failed task poller: 失败任务自动重跑(最多3次),超限报故障给用户 async def _failed_poller(app): @@ -785,7 +790,7 @@ def load_pipeline_service(): asyncio.create_task(_fd_poll_loop()) debug("failed poller started") - add_startup(_failed_poller) + add_startup(_failed_poller) if _run_pollers else None # 铁律:运行期不做任何 schema 变更。v2 引擎的建表/加列已全部迁到 # 部署期 scripts/create_tables.py(build.sh 会调用),此处不再启动时建表。