fix(agent): 角色 agent 实例 id——created_by 存真实 id 不再拼字符串

- pipeline_agent_instances 表(项目×角色→稳定实例 ag.xxxx), 幂等注册
- agent/pm/qc 三 poller 派单改用实例 id(原 poller-{role} 拼接, 33字符撞32列宽→1406)
- DDL: 新表 + pipeline_deliverables.created_by 加宽64兜底
This commit is contained in:
ymq 2026-09-01 09:26:46 +08:00
parent a7d522b15f
commit 6f72e777da
2 changed files with 82 additions and 3 deletions

View File

@ -0,0 +1,73 @@
# -*- coding: utf-8 -*-
"""角色 agent 实例注册2026-09-01项目 × 角色 → 稳定实例 id。
背景poller 派单时 agent_id `poller-{role}` 拼出来的字符串
1. 不是身份审计/落库查不到哪个 agent 实例干的
2. `poller-agent.bid_writer_technical` 33 字符超过
pipeline_deliverables.created_by varchar(32) 1406 连环失败
技术章节交付件全部落不了库商务写者恰好 32 字符没事
规则角色规范名约定处理方 = role + agentid
- 每个项目 × 角色第一次派单时注册一个实例id 形如 `ag.{12}`15 字符
任何 32 宽字段都放得下此后复用
- `created_by`/审计 `agent_id` 一律存实例 id
- pm/qc poller 同走此解析角色 agent.pm / agent.qc
"""
import logging
logger = logging.getLogger("pipeline.agent_instance")
async def resolve_agent_instance(project_id, role):
"""解析(必要时注册)项目 × 角色的 agent 实例 id。幂等。
失败兜底返回 `agent.{role}`规范角色名32 字符保证调用方不炸
"""
role = (role or "").strip()
if not role:
return ""
fallback = ("agent." + role) if not role.startswith("agent.") else role
try:
from sqlor.dbpools import DBPools
from appPublic.uniqueID import getID
db = DBPools()
async with db.sqlorContext("pipeline") as sor:
recs = await sor.sqlExe(
"SELECT agent_id FROM pipeline_agent_instances "
"WHERE project_id=${p}$ AND role=${r}$ AND status='active' LIMIT 1",
{"p": project_id or "", "r": role})
await sor.sqlExe("COMMIT", {})
if recs:
aid = getattr(recs[0], "agent_id", "") or ""
if aid:
return aid
aid = "ag." + getID()[:12]
await sor.C("pipeline_agent_instances", {
"id": getID(),
"project_id": project_id or "",
"role": role,
"agent_id": aid,
"status": "active",
})
await sor.sqlExe("COMMIT", {})
logger.info("agent instance registered: project=%s role=%s id=%s",
project_id, role, aid)
return aid
except Exception as e:
# 并发注册撞唯一键 → 重查一次;仍失败用规范角色名兜底
try:
from sqlor.dbpools import DBPools
db = DBPools()
async with db.sqlorContext("pipeline") as sor:
recs = await sor.sqlExe(
"SELECT agent_id FROM pipeline_agent_instances "
"WHERE project_id=${p}$ AND role=${r}$ AND status='active' LIMIT 1",
{"p": project_id or "", "r": role})
await sor.sqlExe("COMMIT", {})
if recs and (getattr(recs[0], "agent_id", "") or ""):
return recs[0].agent_id
except Exception:
pass
logger.warning("resolve_agent_instance fallback role=%s err=%s", role, str(e)[:120])
return fallback[:32]

View File

@ -777,7 +777,9 @@ def load_pipeline_service():
async def _dispatch(tid, pid, r):
try:
await role_agent_loop(pid, r, agent_id=f"poller-{r}")
from .agent_instance import resolve_agent_instance
_aid = await resolve_agent_instance(pid, r)
await role_agent_loop(pid, r, agent_id=_aid)
except Exception as e:
debug(f"agent_poller dispatch error: task={tid} err={e}")
finally:
@ -861,7 +863,9 @@ def load_pipeline_service():
async def _pm_dispatch(tid, pid):
try:
await pm_review_run(pid, agent_id="pm-poller")
from .agent_instance import resolve_agent_instance
_aid = await resolve_agent_instance(pid, "agent.pm")
await pm_review_run(pid, agent_id=_aid)
except Exception as e:
debug(f"pm_poller dispatch error: task={tid} err={e}")
finally:
@ -912,7 +916,9 @@ def load_pipeline_service():
async def _qc_dispatch(tid, pid):
try:
await qc_review_run(pid, agent_id="qc-poller")
from .agent_instance import resolve_agent_instance
_aid = await resolve_agent_instance(pid, "agent.qc")
await qc_review_run(pid, agent_id=_aid)
except Exception as e:
debug(f"qc_poller dispatch error: task={tid} err={e}")
finally: