pipeline-service/pipeline_service/iteration_capability.py

173 lines
6.8 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""迭代能力 — sd_iterations 的状态机语义化迁移(通用、产线无关)。
定位:迭代 CRUD 走 xls2ui 生成的端点;本模块只做「生命周期状态流转」——
每个迁移 CAS 原子(防并发/越权)+ 租户隔离 + 审计。
迭代的生命周期/流转规则在 iteration skill 里(LLM 读 skill 判断合法性);本模块只固化操作原语。
角色规范:role 参数用 agent.{role}(无前缀自动补 agent.);人角色用 {orgtype}.{role}。
scope 约定:sd_iterations 用 project_id 做范围校验(CAS 的 WHERE 里带 project_id),
与 pipeline_tasks 用 tenant_id 不同——sd_* 业务表用自己的归属列。
"""
import logging
from sqlor.dbpools import DBPools
from appPublic.uniqueID import getID
from .audit import record_audit
DBNAME = "pipeline"
logger = logging.getLogger("pipeline.iteration_capability")
TABLE = "sd_iterations"
# 迭代状态(SDLC 默认,状态机语义见 iteration skill)
S_PLANNING = "planning" # 规划中
S_IN_PROGRESS = "in_progress" # 进行中
S_COMPLETED = "completed" # 已完成
S_CANCELLED = "cancelled" # 已取消
def _get_db():
db = DBPools()
if not db.databases:
from appPublic.jsonConfig import getConfig
config = getConfig()
if config.databases:
db.databases = config.databases
return db, DBNAME
def _normalize_role(role):
"""角色规范:agent 角色补 agent. 前缀;人角色 {orgtype}.{role} 保留原样。"""
role = (role or "").strip()
if not role:
return ""
if "." in role:
return role
return f"agent.{role}"
async def _transition(iteration_id, project_id, from_state, to_state, action,
who=None, agent_id=None, detail=None):
"""CAS 状态迁移 + 审计。返回 (ok, message)。"""
if not iteration_id or not project_id:
return False, "缺少 iteration_id 或 project_id"
db, dbname = _get_db()
async with db.sqlorContext(dbname) as sor:
await sor.sqlExe(
f"UPDATE {TABLE} SET status=${{to}}$, updated_at=NOW() "
"WHERE id=${iid}$ AND project_id=${pid}$ AND status=${from}$",
{"to": to_state, "iid": iteration_id, "pid": project_id, "from": from_state})
recs = await sor.R(TABLE, {'id': iteration_id})
await sor.sqlExe("COMMIT", {})
if not recs:
return False, "迭代不存在"
cur = getattr(recs[0], 'status', '')
if cur != to_state:
return False, f"状态迁移失败(CAS): 期望 from={from_state} 实际 status={cur}"
await record_audit(project_id, TABLE, iteration_id, action,
from_state=from_state, to_state=to_state,
who=who, agent_id=agent_id, detail=detail, sor=sor)
return True, to_state
async def create_iteration(project_id, iteration_name, iteration_type="new_feature",
scope="", priority=5, created_by="", who=None, agent_id=None):
"""创建迭代:新建记录,status=planning。返回 (ok, iteration_id_or_message)。"""
if not project_id:
return False, "缺少 project_id"
if not iteration_name or not iteration_name.strip():
return False, "缺少 iteration_name"
db, dbname = _get_db()
async with db.sqlorContext(dbname) as sor:
iid = getID()
try:
priority = int(priority)
except (TypeError, ValueError):
priority = 5
await sor.C(TABLE, {
'id': iid,
'project_id': project_id,
'iteration_name': iteration_name.strip(),
'iteration_type': iteration_type or 'new_feature',
'scope': scope or '',
'status': S_PLANNING,
'priority': priority,
'created_by': created_by or '',
})
await record_audit(project_id, TABLE, iid, 'create',
to_state=S_PLANNING, who=_normalize_role(who),
agent_id=agent_id, sor=sor)
logger.info("create_iteration: %s project=%s", iid, project_id)
return True, iid
async def start_iteration(iteration_id, project_id, who=None, agent_id=None):
"""开始迭代:planning → in_progress。"""
return await _transition(iteration_id, project_id, S_PLANNING, S_IN_PROGRESS, 'start',
who=who, agent_id=agent_id)
async def complete_iteration(iteration_id, project_id, who=None, agent_id=None):
"""完成迭代:in_progress → completed。"""
return await _transition(iteration_id, project_id, S_IN_PROGRESS, S_COMPLETED, 'complete',
who=who, agent_id=agent_id)
async def cancel_iteration(iteration_id, project_id, who=None, agent_id=None, comment=None):
"""取消迭代:planning/in_progress → cancelled。"""
ok, msg = await _transition(iteration_id, project_id, S_PLANNING, S_CANCELLED, 'cancel',
who=who, agent_id=agent_id, detail=comment)
if ok:
return ok, msg
return await _transition(iteration_id, project_id, S_IN_PROGRESS, S_CANCELLED, 'cancel',
who=who, agent_id=agent_id, detail=comment)
async def set_iteration_state(iteration_id, project_id, from_state, to_state,
who=None, agent_id=None, detail=None):
"""通用 CAS 状态迁移兜底(跨产线自定义状态机用)。"""
return await _transition(iteration_id, project_id, from_state, to_state, 'set_state',
who=who, agent_id=agent_id, detail=detail)
async def list_iterations(project_id, status=None, limit=50) -> list:
"""列出迭代(可按状态过滤)。"""
db, dbname = _get_db()
async with db.sqlorContext(dbname) as sor:
conditions = ["project_id=${pid}$"]
params = {"pid": project_id}
if status:
conditions.append("status=${status}$")
params["status"] = status
where = " AND ".join(conditions)
try:
limit = int(limit)
except (TypeError, ValueError):
limit = 50
sql = (f"SELECT * FROM {TABLE} WHERE {where} "
f"ORDER BY created_at ASC LIMIT {limit}")
recs = await sor.sqlExe(sql, params)
# 释放 SELECT 元数据锁
await sor.sqlExe("COMMIT", {})
result = []
for rec in (recs or []):
result.append(_rec_to_dict(rec))
return result
def _rec_to_dict(rec):
"""把 sqlor 记录对象转成 dict(sqlor 行是 DictObject,必须 dict(rec) 取列)。"""
if isinstance(rec, dict):
return dict(rec)
try:
return dict(rec)
except (TypeError, ValueError):
if hasattr(rec, 'to_dict'):
try:
return rec.to_dict()
except Exception:
pass
return {}