refactor(worker): 撤项目级串行,改 git 级串行锁(flock 同 repo),任务其它部分(LLM/写文件)完全并行

This commit is contained in:
ymq 2026-08-18 17:55:21 +08:00
parent ac85ee67f1
commit ff9db3d310

View File

@ -11,11 +11,15 @@ v3.4.0 新增:
""" """
import asyncio import asyncio
import fcntl
import hashlib
import json import json
import os import os
import re import re
import subprocess import subprocess
import logging import logging
import time
from contextlib import asynccontextmanager
logger = logging.getLogger("pipeline.agent_loop") logger = logging.getLogger("pipeline.agent_loop")
@ -259,6 +263,50 @@ async def _write_code_file(filepath, content):
return False, str(e) return False, str(e)
_GIT_LOCK_DIR = os.environ.get('PIPELINE_GIT_LOCK_DIR', '/tmp/pipeline_git_locks')
def _git_lock_file(repo_dir):
"""同 repo 的 git 操作共用一把锁文件。
本机多 worker 用 flock 串行(锁文件在本地 /tmp,不落 NFS,flock 可靠);
跨主机分布时需换 DB/Redis 锁,锁文件位置可用 PIPELINE_GIT_LOCK_DIR 覆盖。
"""
try:
os.makedirs(_GIT_LOCK_DIR, exist_ok=True)
except OSError:
pass
key = hashlib.sha256(os.path.abspath(repo_dir).encode('utf-8')).hexdigest()[:32]
return os.path.join(_GIT_LOCK_DIR, key + '.lock')
@asynccontextmanager
async def _git_lock(repo_dir, timeout=90):
"""git 操作串行锁:只锁 git 那几秒,任务其它部分(LLM/写文件)完全并行。"""
path = _git_lock_file(repo_dir)
fd = open(path, 'w')
deadline = time.time() + timeout
acquired = False
try:
while True:
try:
fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
acquired = True
break
except BlockingIOError:
if time.time() > deadline:
raise TimeoutError(f'git lock timeout after {timeout}s: {repo_dir}')
await asyncio.sleep(0.3)
yield
finally:
if acquired:
try:
fcntl.flock(fd, fcntl.LOCK_UN)
except OSError:
pass
fd.close()
async def _git_setup(workdir): async def _git_setup(workdir):
"""确保 git 用户已配置。""" """确保 git 用户已配置。"""
r = await _run_shell('git config user.email', workdir, 5) r = await _run_shell('git config user.email', workdir, 5)
@ -269,27 +317,29 @@ async def _git_setup(workdir):
async def _git_commit_push(workdir, commit_message, branch='main'): async def _git_commit_push(workdir, commit_message, branch='main'):
"""git add → commit → push。""" """git add → commit → push。"""
await _git_setup(workdir) async with _git_lock(workdir):
r1 = await _run_shell('git add -A', workdir, 10) await _git_setup(workdir)
r2 = await _run_shell(f'git diff --cached --stat', workdir, 10) r1 = await _run_shell('git add -A', workdir, 10)
if not r2.get('stdout', '').strip(): r2 = await _run_shell(f'git diff --cached --stat', workdir, 10)
return {"rc": 0, "message": "没有变更需要提交"} if not r2.get('stdout', '').strip():
r3 = await _run_shell(f'git commit -m "{commit_message}"', workdir, 15) return {"rc": 0, "message": "没有变更需要提交"}
if r3['rc'] != 0: r3 = await _run_shell(f'git commit -m "{commit_message}"', workdir, 15)
return {"rc": r3['rc'], "message": f"commit 失败: {r3['stderr'][:200]}"} if r3['rc'] != 0:
r4 = await _run_shell(f'git push origin {branch}', workdir, 30) return {"rc": r3['rc'], "message": f"commit 失败: {r3['stderr'][:200]}"}
return {"rc": r4['rc'], "message": f"push {'成功' if r4['rc']==0 else '失败'}: {r4['stderr'][:200]}"} r4 = await _run_shell(f'git push origin {branch}', workdir, 30)
return {"rc": r4['rc'], "message": f"push {'成功' if r4['rc']==0 else '失败'}: {r4['stderr'][:200]}"}
async def _git_clone(repo_url, target_dir, branch='main'): async def _git_clone(repo_url, target_dir, branch='main'):
"""克隆仓库到目标目录。已存在则 pull。""" """克隆仓库到目标目录。已存在则 pull。"""
if os.path.isdir(os.path.join(target_dir, '.git')): async with _git_lock(target_dir):
r = await _run_shell(f'git checkout {branch} && git pull origin {branch}', target_dir, 30) if os.path.isdir(os.path.join(target_dir, '.git')):
return {"rc": r['rc'], "message": f"已存在,pull: {r['stdout'][:200]}"} r = await _run_shell(f'git checkout {branch} && git pull origin {branch}', target_dir, 30)
parent = os.path.dirname(target_dir) return {"rc": r['rc'], "message": f"已存在,pull: {r['stdout'][:200]}"}
os.makedirs(parent, exist_ok=True) parent = os.path.dirname(target_dir)
r = await _run_shell(f'git clone -b {branch} {repo_url} {target_dir}', parent, 120) os.makedirs(parent, exist_ok=True)
return {"rc": r['rc'], "message": f"clone: {r['stdout'][:200] if r['rc']==0 else r['stderr'][:200]}"} r = await _run_shell(f'git clone -b {branch} {repo_url} {target_dir}', parent, 120)
return {"rc": r['rc'], "message": f"clone: {r['stdout'][:200] if r['rc']==0 else r['stderr'][:200]}"}
# ── Prompts ── # ── Prompts ──
@ -525,26 +575,18 @@ async def _claim_task(sor, tenant_id, role, state='submitted', match_role=True,
from appPublic.uniqueID import getID from appPublic.uniqueID import getID
claim_token = getID() claim_token = getID()
# claimed_by IS NULL 保证原子认领(并发 poller / start_agents 不会双重认领); # claimed_by IS NULL 保证原子认领(并发 poller / start_agents 不会双重认领);
# NOT EXISTS 保证同项目(tenant_id)串行:同一项目同一时刻只允许一个任务处于
# 活跃状态(running / review·已认领 / qc_review·已认领),避免多 worker 并发写
# 同一 workspace。跨项目完全并发。derived table 包裹规避 MySQL 1093(同表子查询)。
# updated_at=NOW() 作为心跳,供 stale 回收判断。 # updated_at=NOW() 作为心跳,供 stale 回收判断。
# 同项目并发写冲突由 git 级串行锁解决(见 _git_lock),不在认领层做项目级串行,
# 否则会把整个任务时长(LLM+写文件)都锁死,牺牲项目内并行度。
await sor.sqlExe( await sor.sqlExe(
"UPDATE pipeline_tasks SET state=${setstate}$, claimed_by=${cb}$, updated_at=NOW() " "UPDATE pipeline_tasks SET state=${setstate}$, claimed_by=${cb}$, updated_at=NOW() "
"WHERE id=${tid}$ AND state=${state}$ AND claimed_by IS NULL " "WHERE id=${tid}$ AND state=${state}$ AND claimed_by IS NULL",
"AND NOT EXISTS (" {"setstate": set_state, "cb": claim_token, "tid": task_id, "state": state})
" SELECT 1 FROM (SELECT 1 FROM pipeline_tasks t2 "
" WHERE t2.tenant_id=${tidx}$ AND t2.claimed_by IS NOT NULL "
" AND t2.state IN ('running','review','qc_review')) AS z"
")",
{"setstate": set_state, "cb": claim_token, "tid": task_id,
"state": state, "tidx": tenant_id})
check = await sor.sqlExe( check = await sor.sqlExe(
"SELECT id FROM pipeline_tasks WHERE id=${tid}$ AND state=${setstate}$ AND claimed_by=${cb}$", "SELECT id FROM pipeline_tasks WHERE id=${tid}$ AND state=${setstate}$ AND claimed_by=${cb}$",
{"tid": task_id, "setstate": set_state, "cb": claim_token}) {"tid": task_id, "setstate": set_state, "cb": claim_token})
if not check: if not check:
# 并发抢输(task 已被别人认领)或同项目已有活跃任务被串行阻塞,均属正常,下轮再试 logger.info(f"claim lost race: task={task_id}")
logger.info(f"claim skipped: task={task_id}")
return None return None
return task return task