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 会调用),此处不再启动时建表。