diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index 5c3d145..43fff54 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -525,16 +525,26 @@ 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 回收判断。 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", - {"setstate": set_state, "cb": claim_token, "tid": task_id, "state": state}) + "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}) 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: - logger.info(f"claim lost race: task={task_id}") + # 并发抢输(task 已被别人认领)或同项目已有活跃任务被串行阻塞,均属正常,下轮再试 + logger.info(f"claim skipped: task={task_id}") return None return task