From 781858b61aa43191fb52b7b7dcbd501759b75ce2 Mon Sep 17 00:00:00 2001 From: ymq Date: Tue, 18 Aug 2026 18:13:50 +0800 Subject: [PATCH] =?UTF-8?q?feat(worker):=20git=20=E9=94=81=E4=BB=8E=20floc?= =?UTF-8?q?k=20=E5=8D=87=E7=BA=A7=20DB=20=E9=94=81=E8=A1=A8(pipeline=5Fgit?= =?UTF-8?q?=5Flocks)=EF=BC=8C=E8=B7=A8=E4=B8=BB=E6=9C=BA=E5=A4=9A=20worker?= =?UTF-8?q?=20=E7=94=9F=E6=95=88=EF=BC=8CTTL=20180s=20=E5=B4=A9=E6=BA=83?= =?UTF-8?q?=E8=87=AA=E9=87=8A=E6=94=BE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pipeline_service/agent_loop.py | 62 ++++++++++++++++++++-------------- 1 file changed, 36 insertions(+), 26 deletions(-) diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index d415d6a..ff03ec8 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -11,7 +11,6 @@ v3.4.0 新增: """ import asyncio -import fcntl import hashlib import json import os @@ -263,37 +262,46 @@ 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') +_GIT_LOCK_TTL = 180 # 秒:覆盖 git 单次操作最坏时长(clone 120s);过期可被原子接管 -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') +def _git_lock_key(repo_dir): + """同 repo 的 git 操作共用一把锁。lock_key = sha256(repo_abs_path)。""" + return hashlib.sha256(os.path.abspath(repo_dir).encode('utf-8')).hexdigest() @asynccontextmanager async def _git_lock(repo_dir, timeout=90): - """git 操作串行锁:只锁 git 那几秒,任务其它部分(LLM/写文件)完全并行。""" - path = _git_lock_file(repo_dir) - fd = open(path, 'w') + """git 操作串行锁(DB 版,跨主机多 worker 生效):只锁 git 那几秒,任务其它部分完全并行。 + + 用 pipeline_git_locks 表 + ON DUPLICATE KEY UPDATE + expires_at TTL: + - 原子抢占:INSERT ... ON DUPLICATE KEY UPDATE,仅当 expires_at 过期才接管他人锁 + - 崩溃自动释放:expires_at 过期后下一个申请者原子接管 + - 释放:DELETE ... WHERE token 匹配(避免误删他人锁) + """ + from appPublic.uniqueID import getID + + key = _git_lock_key(repo_dir) + token = getID() deadline = time.time() + timeout acquired = False + db = _get_db() try: - while True: - try: - fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) - acquired = True - break - except BlockingIOError: + while not acquired: + async with db.sqlorContext("pipeline") as sor: + await sor.sqlExe( + "INSERT INTO pipeline_git_locks (lock_key, token, expires_at) " + "VALUES (${k}$, ${t}$, DATE_ADD(NOW(), INTERVAL " + str(_GIT_LOCK_TTL) + " SECOND)) " + "ON DUPLICATE KEY UPDATE " + "token = IF(expires_at < NOW(), VALUES(token), token), " + "expires_at = IF(expires_at < NOW(), DATE_ADD(NOW(), INTERVAL " + str(_GIT_LOCK_TTL) + " SECOND), expires_at)", + {"k": key, "t": token}) + r = await sor.sqlExe( + "SELECT token FROM pipeline_git_locks WHERE lock_key=${k}$ AND token=${t}$", + {"k": key, "t": token}) + if r: + acquired = True + if not acquired: if time.time() > deadline: raise TimeoutError(f'git lock timeout after {timeout}s: {repo_dir}') await asyncio.sleep(0.3) @@ -301,10 +309,12 @@ async def _git_lock(repo_dir, timeout=90): finally: if acquired: try: - fcntl.flock(fd, fcntl.LOCK_UN) - except OSError: + async with db.sqlorContext("pipeline") as sor: + await sor.sqlExe( + "DELETE FROM pipeline_git_locks WHERE lock_key=${k}$ AND token=${t}$", + {"k": key, "t": token}) + except Exception: pass - fd.close() async def _git_setup(workdir):