diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index 43fff54..d415d6a 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -11,11 +11,15 @@ v3.4.0 新增: """ import asyncio +import fcntl +import hashlib import json import os import re import subprocess import logging +import time +from contextlib import asynccontextmanager logger = logging.getLogger("pipeline.agent_loop") @@ -259,6 +263,50 @@ async def _write_code_file(filepath, content): 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): """确保 git 用户已配置。""" 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'): """git add → commit → push。""" - await _git_setup(workdir) - r1 = await _run_shell('git add -A', workdir, 10) - r2 = await _run_shell(f'git diff --cached --stat', workdir, 10) - if not r2.get('stdout', '').strip(): - return {"rc": 0, "message": "没有变更需要提交"} - r3 = await _run_shell(f'git commit -m "{commit_message}"', workdir, 15) - if r3['rc'] != 0: - return {"rc": r3['rc'], "message": f"commit 失败: {r3['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 with _git_lock(workdir): + await _git_setup(workdir) + r1 = await _run_shell('git add -A', workdir, 10) + r2 = await _run_shell(f'git diff --cached --stat', workdir, 10) + if not r2.get('stdout', '').strip(): + return {"rc": 0, "message": "没有变更需要提交"} + r3 = await _run_shell(f'git commit -m "{commit_message}"', workdir, 15) + if r3['rc'] != 0: + return {"rc": r3['rc'], "message": f"commit 失败: {r3['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'): """克隆仓库到目标目录。已存在则 pull。""" - if os.path.isdir(os.path.join(target_dir, '.git')): - r = await _run_shell(f'git checkout {branch} && git pull origin {branch}', target_dir, 30) - return {"rc": r['rc'], "message": f"已存在,pull: {r['stdout'][:200]}"} - parent = os.path.dirname(target_dir) - os.makedirs(parent, exist_ok=True) - 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]}"} + async with _git_lock(target_dir): + if os.path.isdir(os.path.join(target_dir, '.git')): + r = await _run_shell(f'git checkout {branch} && git pull origin {branch}', target_dir, 30) + return {"rc": r['rc'], "message": f"已存在,pull: {r['stdout'][:200]}"} + parent = os.path.dirname(target_dir) + os.makedirs(parent, exist_ok=True) + 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 ── @@ -525,26 +575,18 @@ async def _claim_task(sor, tenant_id, role, state='submitted', match_role=True, from appPublic.uniqueID import getID claim_token = getID() # 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 回收判断。 + # 同项目并发写冲突由 git 级串行锁解决(见 _git_lock),不在认领层做项目级串行, + # 否则会把整个任务时长(LLM+写文件)都锁死,牺牲项目内并行度。 await sor.sqlExe( "UPDATE pipeline_tasks SET state=${setstate}$, claimed_by=${cb}$, updated_at=NOW() " - "WHERE id=${tid}$ AND state=${state}$ AND claimed_by IS NULL " - "AND NOT EXISTS (" - " 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}) + "WHERE id=${tid}$ AND state=${state}$ AND claimed_by IS NULL", + {"setstate": set_state, "cb": claim_token, "tid": task_id, "state": state}) check = await sor.sqlExe( "SELECT id FROM pipeline_tasks WHERE id=${tid}$ AND state=${setstate}$ AND claimed_by=${cb}$", {"tid": task_id, "setstate": set_state, "cb": claim_token}) if not check: - # 并发抢输(task 已被别人认领)或同项目已有活跃任务被串行阻塞,均属正常,下轮再试 - logger.info(f"claim skipped: task={task_id}") + logger.info(f"claim lost race: task={task_id}") return None return task