feat(worker): load_pipeline_service 按 PIPELINE_MODE 门控 poller 注册(web=仅HTTP无poller, worker=仅poller, all=默认)

This commit is contained in:
ymq 2026-08-18 17:15:21 +08:00
parent 8108d7d902
commit 1eba1c3f87

View File

@ -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 + pollerweb = 仅 HTTPworker = 仅 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.pybuild.sh 会调用),此处不再启动时建表。