From 3687270ca746453795d4794726d7399fabe05b09 Mon Sep 17 00:00:00 2001 From: ymq Date: Sun, 16 Aug 2026 23:39:05 +0800 Subject: [PATCH] =?UTF-8?q?feat(task-capability):=20=E4=BB=BB=E5=8A=A1?= =?UTF-8?q?=E7=8A=B6=E6=80=81=E6=9C=BA=E8=AF=AD=E4=B9=89=E5=8C=96=E8=BF=81?= =?UTF-8?q?=E7=A7=BB=E8=83=BD=E5=8A=9B=20+=20=E9=80=9A=E7=94=A8=E5=AE=A1?= =?UTF-8?q?=E8=AE=A1=E5=8E=9F=E8=AF=AD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - task_capability.py: claim/submit/approve/reject/complete/mark_failed/retry/suspend/revive/set_task_state CAS 原子迁移(防并发/越权覆盖) + tenant 隔离 + record_audit 审计 定位=状态机语义化操作层,CRUD 复用 storage.py,不遮蔽 storage/executor 同名函数 - audit.py: record_audit append-only 通用审计原语(支持 sor 复用避免嵌套连接) - models/audit_log.json: 审计日志表(append-only, tenant 隔离) - init.py: 注册 task_* + record_audit 到 ServerEnv --- models/audit_log.json | 28 +++++ pipeline_service/audit.py | 64 ++++++++++ pipeline_service/init.py | 19 +++ pipeline_service/task_capability.py | 189 ++++++++++++++++++++++++++++ 4 files changed, 300 insertions(+) create mode 100644 models/audit_log.json create mode 100644 pipeline_service/audit.py create mode 100644 pipeline_service/task_capability.py diff --git a/models/audit_log.json b/models/audit_log.json new file mode 100644 index 0000000..04f20d5 --- /dev/null +++ b/models/audit_log.json @@ -0,0 +1,28 @@ +{ + "summary": [ + { + "name": "audit_log", + "title": "审计日志表(append-only)", + "primary": ["id"], + "catelog": "entity" + } + ], + "fields": [ + {"name": "id", "title": "主键ID", "type": "str", "length": 32, "nullable": "no"}, + {"name": "tenant_id", "title": "租户ID", "type": "str", "length": 32, "nullable": "no"}, + {"name": "entity", "title": "实体类型", "type": "str", "length": 64, "nullable": "no"}, + {"name": "entity_id", "title": "实体ID", "type": "str", "length": 64, "nullable": "no"}, + {"name": "action", "title": "动作", "type": "str", "length": 64, "nullable": "no"}, + {"name": "from_state", "title": "原状态", "type": "str", "length": 32}, + {"name": "to_state", "title": "新状态", "type": "str", "length": 32}, + {"name": "who", "title": "操作角色", "type": "str", "length": 64}, + {"name": "agent_id", "title": "操作Agent标识", "type": "str", "length": 64}, + {"name": "detail", "title": "详情", "type": "text"}, + {"name": "created_at", "title": "创建时间", "type": "timestamp", "nullable": "no"} + ], + "indexes": [ + {"name": "idx_al_tenant", "idxtype": "index", "idxfields": ["tenant_id"]}, + {"name": "idx_al_entity", "idxtype": "index", "idxfields": ["entity", "entity_id"]}, + {"name": "idx_al_created", "idxtype": "index", "idxfields": ["created_at"]} + ] +} diff --git a/pipeline_service/audit.py b/pipeline_service/audit.py new file mode 100644 index 0000000..cde161d --- /dev/null +++ b/pipeline_service/audit.py @@ -0,0 +1,64 @@ +"""审计日志 — append-only 通用审计原语。 + +审计粒度(什么该记、记到什么程度)在 skill 里说明;本模块只提供 record_audit +通用原语,所有概念能力(task/question/deliverable)的状态迁移/流转都调它记一条。 + +设计约束: +- append-only:只 INSERT,不 UPDATE/DELETE(审计独立性,防自删)。 +- 租户隔离:每条必带 tenant_id。 +- sor 参数可选:在已有 sqlorContext 里调用时传入避免嵌套开连接。 +""" + +import logging +from sqlor.dbpools import DBPools +from appPublic.uniqueID import getID + +DBNAME = "pipeline" +logger = logging.getLogger("pipeline.audit") + + +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 + + +async def record_audit(tenant_id, entity, entity_id, action, + from_state=None, to_state=None, who=None, + agent_id=None, detail=None, sor=None): + """追加一条审计记录(只追加不改删)。 + + Args: + tenant_id: 租户ID(必带,隔离) + entity: 实体类型(如 pipeline_tasks / pipeline_agent_questions / pipeline_deliverables) + entity_id: 实体ID + action: 动作(如 claim/submit/approve/reject/raise/escalate/resolve/set_state) + from_state / to_state: 状态迁移前后(可空) + who: 操作角色(agent.{role} 或 {orgtype}.{role}) + agent_id: 具体 agent 标识 + detail: 附加详情 + sor: 可选的 sqlor context,传入则复用(避免嵌套开连接) + """ + data = { + 'id': getID(), + 'tenant_id': tenant_id or '', + 'entity': entity or '', + 'entity_id': entity_id or '', + 'action': action or '', + 'from_state': from_state, + 'to_state': to_state, + 'who': who, + 'agent_id': agent_id, + 'detail': detail, + } + if sor is not None: + await sor.C('audit_log', data) + return True + db, dbname = _get_db() + async with db.sqlorContext(dbname) as s: + await s.C('audit_log', data) + return True diff --git a/pipeline_service/init.py b/pipeline_service/init.py index 36669e0..ca870bf 100644 --- a/pipeline_service/init.py +++ b/pipeline_service/init.py @@ -41,6 +41,12 @@ from .communication import ( raise_problem, resolve_problem, escalate_problem, list_problems_for, get_task_qa, ) +from .audit import record_audit +from .task_capability import ( + claim_task, submit_task, approve_task, reject_task, + complete_task, mark_failed, retry_task, suspend_task, + revive_task, set_task_state, +) from .agent_loop import role_agent_run, role_agent_loop, agent_loop, run_agent_loop from .agent_loop import pm_review_run, pm_review_loop, handle_failed_task @@ -550,6 +556,19 @@ def load_pipeline_service(): env.problem_list_for = list_problems_for env.problem_qa = get_task_qa + # 任务能力(状态机语义化迁移;状态机在 task skill)+ 审计原语 + env.record_audit = record_audit + env.claim_task = claim_task + env.submit_task = submit_task + env.approve_task = approve_task + env.reject_task = reject_task + env.complete_task = complete_task + env.mark_failed = mark_failed + env.retry_task = retry_task + env.suspend_task = suspend_task + env.revive_task = revive_task + env.set_task_state = set_task_state + # DevOps: git/shell operations (v3.2.1) env.shell_exec = shell_exec env.skill_import_git = skill_import_git diff --git a/pipeline_service/task_capability.py b/pipeline_service/task_capability.py new file mode 100644 index 0000000..3925e85 --- /dev/null +++ b/pipeline_service/task_capability.py @@ -0,0 +1,189 @@ +"""任务能力 — pipeline_tasks 的状态机语义化迁移(通用、产线无关)。 + +定位:storage.py 是数据 CRUD 层(create_task/list_tasks/get_task/update_task_state 低层改状态), +本模块是状态机语义化操作层 —— 每个迁移 CAS 原子(防并发/防越权覆盖)+ 租户隔离 + 审计。 + +任务的状态机/流转规则在 task skill 里(LLM 读 skill 判断合法性);本模块只固化操作原语。 + +角色规范:role 参数用 agent.{role}(无前缀自动补 agent.);人角色用 {orgtype}.{role}。 +""" + +import logging +from sqlor.dbpools import DBPools +from appPublic.uniqueID import getID +from .audit import record_audit + +DBNAME = "pipeline" +logger = logging.getLogger("pipeline.task_capability") + +# 任务状态(SDLC 默认,状态机语义见 task skill) +S_SUBMITTED = "submitted" # 待认领 +S_RUNNING = "running" # 角色 agent 执行中 +S_REVIEW = "review" # 已提交待审核 +S_APPROVED = "approved" # 审核通过 +S_COMPLETED = "completed" # 全流程完成 +S_WAITING = "waiting" # 挂起等回答 +S_FAILED = "failed" # 失败 + + +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(task_id, tenant_id, from_state, to_state, action, + who=None, agent_id=None, detail=None): + """CAS 状态迁移(清 claimed_by)+ 审计。返回 (ok, message)。""" + if not task_id or not tenant_id: + return False, "缺少 task_id 或 tenant_id" + db, dbname = _get_db() + async with db.sqlorContext(dbname) as sor: + await sor.sqlExe( + "UPDATE pipeline_tasks SET state=${to}$, claimed_by=NULL, updated_at=NOW() " + "WHERE id=${tid}$ AND tenant_id=${tn}$ AND state=${from}$", + {"to": to_state, "tid": task_id, "tn": tenant_id, "from": from_state}) + recs = await sor.R('pipeline_tasks', {'id': task_id}) + await sor.sqlExe("COMMIT", {}) + if not recs: + return False, "任务不存在" + cur = getattr(recs[0], 'state', '') + if cur != to_state: + return False, f"状态迁移失败(CAS): 期望 from={from_state} 实际 state={cur}" + await record_audit(tenant_id, 'pipeline_tasks', task_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 claim_task(tenant_id, role, agent_id, from_state=S_SUBMITTED, + to_state=S_RUNNING): + """认领任务:CAS claimed_by IS NULL 保证并发认领原子;设置 claimed_by=agent_id。 + + 返回 (ok, task_id_or_message)。 + """ + role = _normalize_role(role) + db, dbname = _get_db() + async with db.sqlorContext(dbname) as sor: + claim_token = agent_id or getID() + recs = await sor.sqlExe( + "SELECT id FROM pipeline_tasks " + "WHERE tenant_id=${tn}$ AND state=${st}$ AND role=${role}$ " + "AND (claimed_by IS NULL OR claimed_by='') LIMIT 1", + {"tn": tenant_id, "st": from_state, "role": role}) + await sor.sqlExe("COMMIT", {}) + if not recs: + return False, f"没有可认领的 {role} 任务(state={from_state})" + task_id = getattr(recs[0], 'id', '') + # CAS:claimed_by IS NULL 保证原子(并发 poller 不会双重认领) + await sor.sqlExe( + "UPDATE pipeline_tasks SET state=${to}$, claimed_by=${cb}$, updated_at=NOW() " + "WHERE id=${tid}$ AND tenant_id=${tn}$ AND state=${from}$ AND claimed_by IS NULL", + {"to": to_state, "cb": claim_token, "tid": task_id, "tn": tenant_id, "from": from_state}) + chk = await sor.R('pipeline_tasks', {'id': task_id}) + await sor.sqlExe("COMMIT", {}) + if not chk or getattr(chk[0], 'state', '') != to_state: + return False, "认领竞争失败(已被他人认领)" + await record_audit(tenant_id, 'pipeline_tasks', task_id, 'claim', + from_state=from_state, to_state=to_state, + who=role, agent_id=agent_id, sor=sor) + return True, task_id + + +async def submit_task(task_id, tenant_id, who=None, agent_id=None): + """提交产出:running → review(清 claimed_by)。""" + return await _transition(task_id, tenant_id, S_RUNNING, S_REVIEW, 'submit', + who=who, agent_id=agent_id) + + +async def approve_task(task_id, tenant_id, who=None, agent_id=None, comment=None): + """审核通过:review → approved。""" + return await _transition(task_id, tenant_id, S_REVIEW, S_APPROVED, 'approve', + who=who, agent_id=agent_id, detail=comment) + + +async def reject_task(task_id, tenant_id, who=None, agent_id=None, comment=None): + """审核退回:review → submitted(清 claimed_by,被退角色重新认领)。""" + return await _transition(task_id, tenant_id, S_REVIEW, S_SUBMITTED, 'reject', + who=who, agent_id=agent_id, detail=comment) + + +async def complete_task(task_id, tenant_id, who=None, agent_id=None): + """标记完成:approved → completed(也兼容 review → completed)。""" + ok, msg = await _transition(task_id, tenant_id, S_APPROVED, S_COMPLETED, 'complete', + who=who, agent_id=agent_id) + if ok: + return ok, msg + return await _transition(task_id, tenant_id, S_REVIEW, S_COMPLETED, 'complete', + who=who, agent_id=agent_id) + + +async def mark_failed(task_id, tenant_id, who=None, agent_id=None, error=None): + """标记失败(任意状态 → failed,清 claimed_by)。""" + db, dbname = _get_db() + async with db.sqlorContext(dbname) as sor: + recs = await sor.R('pipeline_tasks', {'id': task_id}) + if not recs: + return False, "任务不存在" + from_state = getattr(recs[0], 'state', '') + await sor.sqlExe( + "UPDATE pipeline_tasks SET state='failed', last_error=${e}$, " + "claimed_by=NULL, updated_at=NOW() WHERE id=${tid}$ AND tenant_id=${tn}$", + {"e": (error or '')[:4000], "tid": task_id, "tn": tenant_id}) + await sor.sqlExe("COMMIT", {}) + await record_audit(tenant_id, 'pipeline_tasks', task_id, 'fail', + from_state=from_state, to_state=S_FAILED, + who=who, agent_id=agent_id, detail=error, sor=sor) + return True, S_FAILED + + +async def retry_task(task_id, tenant_id, who=None, agent_id=None): + """重试:failed → submitted(retry_count+1,清 claimed_by)。""" + db, dbname = _get_db() + async with db.sqlorContext(dbname) as sor: + await sor.sqlExe( + "UPDATE pipeline_tasks SET retry_count=retry_count+1, state='submitted', " + "claimed_by=NULL, last_error=NULL, updated_at=NOW() " + "WHERE id=${tid}$ AND tenant_id=${tn}$ AND state='failed'", + {"tid": task_id, "tn": tenant_id}) + chk = await sor.R('pipeline_tasks', {'id': task_id}) + await sor.sqlExe("COMMIT", {}) + if not chk or getattr(chk[0], 'state', '') != S_SUBMITTED: + return False, "重试失败(任务不在 failed 状态)" + await record_audit(tenant_id, 'pipeline_tasks', task_id, 'retry', + from_state=S_FAILED, to_state=S_SUBMITTED, + who=who, agent_id=agent_id, sor=sor) + return True, S_SUBMITTED + + +async def suspend_task(task_id, tenant_id, who=None, agent_id=None): + """挂起等回答:running → waiting(清 claimed_by)。""" + return await _transition(task_id, tenant_id, S_RUNNING, S_WAITING, 'suspend', + who=who, agent_id=agent_id) + + +async def revive_task(task_id, tenant_id, who=None, agent_id=None): + """恢复认领:waiting → submitted(清 claimed_by,重新认领)。""" + return await _transition(task_id, tenant_id, S_WAITING, S_SUBMITTED, 'revive', + who=who, agent_id=agent_id) + + +async def set_task_state(task_id, tenant_id, from_state, to_state, + who=None, agent_id=None, detail=None): + """通用 CAS 状态迁移兜底(跨产线自定义状态机用)。""" + return await _transition(task_id, tenant_id, from_state, to_state, 'set_state', + who=who, agent_id=agent_id, detail=detail)