From e99a952dd0447670ce4e73120308d80b0fbcf938 Mon Sep 17 00:00:00 2001 From: zizi Date: Thu, 10 Sep 2026 19:25:40 +0800 Subject: [PATCH] =?UTF-8?q?W05=20=E4=BB=BB=E5=8A=A1=E7=8A=B6=E6=80=81?= =?UTF-8?q?=E3=80=81=E7=A7=9F=E7=BA=A6=E3=80=81=E4=BA=8B=E4=BB=B6=E4=B8=8E?= =?UTF-8?q?=E6=81=A2=E5=A4=8D=EF=BC=9A=E4=BB=BB=E5=8A=A1=E7=8A=B6=E6=80=81?= =?UTF-8?q?=E6=9C=BA=E3=80=81=E7=9F=AD=E4=BA=8B=E5=8A=A1=E7=A7=9F=E7=BA=A6?= =?UTF-8?q?=E3=80=81=E4=BA=8B=E4=BB=B6=E7=BB=AD=E6=8E=A5=E3=80=81=E6=81=A2?= =?UTF-8?q?=E5=A4=8D=E6=8E=A7=E5=88=B6=E4=B8=8E=E4=BD=9C=E7=94=A8=E5=9F=9F?= =?UTF-8?q?=E4=BA=92=E6=96=A5=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 按 R2 串行阶段整理提交;包内文件为该阶段交付(含后续小增量),状态以工作包清单为准。 --- .agent/skills/操作/查看与恢复任务/SKILL.md | 16 + .agent/skills/操作/查看与恢复任务/目录.md | 3 + src/muse/任务运行/__init__.py | 1 + src/muse/任务运行/事件记录.py | 16 + src/muse/任务运行/任务领取.py | 46 ++ src/muse/任务运行/作用域互斥.py | 31 ++ src/muse/任务运行/存储.py | 512 +++++++++++++++++++++ src/muse/任务运行/恢复控制.py | 44 ++ src/muse/任务运行/接口.py | 316 +++++++++++++ src/muse/任务运行/模型.py | 220 +++++++++ src/muse/任务运行/步骤执行.py | 77 ++++ src/muse/基础设施/__init__.py | 4 + src/muse/基础设施/可观测性.py | 37 ++ src/muse/接入/cli/任务命令.py | 33 ++ src/muse/接入/http/路由/任务运行.py | 210 +++++++++ src/muse/接入/执行器.py | 34 ++ src/muse/编排/__init__.py | 1 + src/muse/编排/接口.py | 24 + src/muse/编排/槽位约束.py | 46 ++ src/muse/编排/流程版本.py | 12 + src/muse/编排/流程登记.py | 77 ++++ tests/单元/test_流程登记与保护.py | 57 +++ tests/集成/test_任务租约与恢复.py | 439 ++++++++++++++++++ 数据库/迁移/V0004__任务运行与证据.sql | 290 ++++++++++++ 数据库/迁移/V0016__流程定义与版本.sql | 20 + 25 files changed, 2566 insertions(+) create mode 100644 .agent/skills/操作/查看与恢复任务/SKILL.md create mode 100644 .agent/skills/操作/查看与恢复任务/目录.md create mode 100644 src/muse/任务运行/__init__.py create mode 100644 src/muse/任务运行/事件记录.py create mode 100644 src/muse/任务运行/任务领取.py create mode 100644 src/muse/任务运行/作用域互斥.py create mode 100644 src/muse/任务运行/存储.py create mode 100644 src/muse/任务运行/恢复控制.py create mode 100644 src/muse/任务运行/接口.py create mode 100644 src/muse/任务运行/模型.py create mode 100644 src/muse/任务运行/步骤执行.py create mode 100644 src/muse/基础设施/__init__.py create mode 100644 src/muse/基础设施/可观测性.py create mode 100644 src/muse/接入/cli/任务命令.py create mode 100644 src/muse/接入/http/路由/任务运行.py create mode 100644 src/muse/接入/执行器.py create mode 100644 src/muse/编排/__init__.py create mode 100644 src/muse/编排/接口.py create mode 100644 src/muse/编排/槽位约束.py create mode 100644 src/muse/编排/流程版本.py create mode 100644 src/muse/编排/流程登记.py create mode 100644 tests/单元/test_流程登记与保护.py create mode 100644 tests/集成/test_任务租约与恢复.py create mode 100644 数据库/迁移/V0004__任务运行与证据.sql create mode 100644 数据库/迁移/V0016__流程定义与版本.sql diff --git a/.agent/skills/操作/查看与恢复任务/SKILL.md b/.agent/skills/操作/查看与恢复任务/SKILL.md new file mode 100644 index 0000000..ff5d3d2 --- /dev/null +++ b/.agent/skills/操作/查看与恢复任务/SKILL.md @@ -0,0 +1,16 @@ +--- +id: inspect-task +name: 查看与恢复任务 +category: operation +contract_version: 1 +commands: [list_tasks, read_task, read_task_events, control_task] +description: 查询持久任务与事件,按作者请求暂停、取消或经登记检查恢复任务。 +--- + +# 查看与恢复任务 + +用 `.venv/bin/python -m muse 任务 <私人作者配置> 列表` 定位任务,再用 `查看 <任务ID>` 或 `事件 <任务ID> --游标 <序号>` 读取状态。任务、步骤和尝试状态分别判断;事件文本不等于业务成功。 + +控制请求写入私人 JSON:`command_id`、`target_ref`、`expected_state`、`action`。action 为“暂停”“取消”或“恢复”,通过 `.venv/bin/python -m muse 任务 <私人作者配置> 控制 <请求文件>` 执行。只携带用户已授权的动作;重试同一请求保留 command_id,参数改变必须用新的命令身份。 + +恢复由服务端按冻结流程版本取得 owner 检查。来源、结构、授权或预算已变化时保留拒绝原因;未知调用先对账,不把它改记为未调用。具体合同见[任务工具与事件](../../../../docs/系统架构/新版设计/接口契约/任务工具与事件.md)。 diff --git a/.agent/skills/操作/查看与恢复任务/目录.md b/.agent/skills/操作/查看与恢复任务/目录.md new file mode 100644 index 0000000..f78a3c3 --- /dev/null +++ b/.agent/skills/操作/查看与恢复任务/目录.md @@ -0,0 +1,3 @@ +| 名称 | 相对地址 | 内容描述 | 使用场景 | 使用要求 | +|------|----------|----------|----------|----------| +| 查看与恢复任务 | [SKILL.md](SKILL.md) | 查询持久任务与事件,按作者请求暂停、取消或经登记检查恢复任务。 | 需要该项已实现操作时 | 遵循配置用途与用户动作授权 | diff --git a/src/muse/任务运行/__init__.py b/src/muse/任务运行/__init__.py new file mode 100644 index 0000000..9298fbe --- /dev/null +++ b/src/muse/任务运行/__init__.py @@ -0,0 +1 @@ +"""持久任务与可靠执行;公开用例见接口模块。""" diff --git a/src/muse/任务运行/事件记录.py b/src/muse/任务运行/事件记录.py new file mode 100644 index 0000000..c296693 --- /dev/null +++ b/src/muse/任务运行/事件记录.py @@ -0,0 +1,16 @@ +"""尝试事件只在有效租约内提交;任务生命周期事件由状态入口产生。""" + +from muse.任务运行.模型 import 事件类型, 任务错误 + +_尝试事件 = frozenset( + {事件类型.文本片段, 事件类型.工具读取, 事件类型.候选保存, 事件类型.检查完成, 事件类型.等待作者} +) + + +def 校验尝试事件(类型: 事件类型, 载荷: dict) -> None: + if 类型 not in _尝试事件: + raise 任务错误("任务和步骤终态事件只能由对应状态入口生成") + if not isinstance(载荷, dict): + raise 任务错误("事件载荷必须是有版本的结构对象") + if 类型 is 事件类型.文本片段 and not isinstance(载荷.get("文本"), str): + raise 任务错误("文本片段事件必须提供文本") diff --git a/src/muse/任务运行/任务领取.py b/src/muse/任务运行/任务领取.py new file mode 100644 index 0000000..59fe141 --- /dev/null +++ b/src/muse/任务运行/任务领取.py @@ -0,0 +1,46 @@ +"""短事务领取与租约校验;步骤和作用域分别验证。""" + +from __future__ import annotations + +from muse.任务运行.作用域互斥 import 取得范围, 核对范围, 释放范围 +from muse.任务运行.存储 import 任务存储 +from muse.任务运行.模型 import 任务错误, 处理器目录, 租约失效, 领取凭证 + + +def 校验领取(存储: 任务存储, 领取: 领取凭证) -> tuple[dict, dict]: + 任务 = 存储.任务行(领取.任务ID, 锁定=True) + if 任务["state"] != "running": + raise 租约失效("任务已停止当前尝试") + 核对范围(存储, 任务, 领取.作用域代次) + 尝试 = 存储.尝试行(领取) + if ( + 尝试 is None + or 尝试["state"] != "running" + or 尝试["step_state"] != "running" + or str(尝试["current_attempt"]) != 领取.尝试ID + or 尝试["processor_id"] != 领取.处理器ID + or 尝试["processor_version"] != 领取.处理器版本 + or 尝试["scope_generation"] != 领取.作用域代次 + ): + raise 租约失效("尝试已过期或被替代") + return 任务, 尝试 + + +def 领取步骤( + 存储: 任务存储, 执行者: str, 能力: list[str], 租期: float, 处理器: 处理器目录 +) -> 领取凭证 | None: + if not 执行者 or 租期 <= 0: + raise 任务错误("执行者和正数租期必须提供") + while 候选 := 存储.领取候选(能力): + 任务 = 候选[0] + 处理器.获取(任务["processor_id"], 任务["processor_version"]) + if 任务["call_state"] in ("sent", "unknown"): + 存储.设置未知(str(任务["task_id"]), 任务["step_id"], str(任务["current_attempt"])) + 释放范围(存储, 任务) + continue + try: + 代次 = 取得范围(存储, 任务, 租期) + except 租约失效: + return None + return 存储.新尝试(任务, 执行者, 租期, 代次) + return None diff --git a/src/muse/任务运行/作用域互斥.py b/src/muse/任务运行/作用域互斥.py new file mode 100644 index 0000000..742425f --- /dev/null +++ b/src/muse/任务运行/作用域互斥.py @@ -0,0 +1,31 @@ +"""作用域持有代次与释放:调用方必须已锁定任务,并处于同一短事务。""" + +from __future__ import annotations + +from typing import TYPE_CHECKING + +from muse.任务运行.模型 import 租约失效 + +if TYPE_CHECKING: + from muse.任务运行.存储 import 任务存储 + + +def 取得范围(存储: 任务存储, 任务: dict, 租期: float) -> int | None: + if 任务["scope_key"] is None: + return None + 代次 = 存储.取得作用域(任务["task_id"], 任务["scope_key"], 租期) + if 代次 is None: + raise 租约失效("任务作用域由另一任务持有") + return 代次 + + +def 核对范围(存储: 任务存储, 任务: dict, 代次: int | None) -> None: + if 任务["scope_key"] is not None and not 存储.有效作用域( + 任务["task_id"], 任务["scope_key"], 代次 + ): + raise 租约失效("作用域已过期或被接管,旧持有者不能写入") + + +def 释放范围(存储: 任务存储, 任务: dict) -> None: + if 任务["scope_key"] is not None: + 存储.释放作用域(任务["task_id"], 任务["scope_key"], 任务["scope_generation"]) diff --git a/src/muse/任务运行/存储.py b/src/muse/任务运行/存储.py new file mode 100644 index 0000000..5a9a7b9 --- /dev/null +++ b/src/muse/任务运行/存储.py @@ -0,0 +1,512 @@ +"""任务运行数据 owner;只使用调用方连接与事务,不调用处理器。""" + +from __future__ import annotations + +import uuid +from dataclasses import asdict +from typing import Any, LiteralString + +import psycopg +from psycopg import sql +from psycopg.rows import dict_row +from psycopg.types.json import Jsonb + +from muse.任务运行.模型 import ( + 事件类型, + 事件续接, + 任务快照, + 任务状态, + 任务请求, + 任务错误, + 内容哈希, + 命令参数冲突, + 执行计划, + 状态冲突, + 运行事件, + 领取凭证, +) +from muse.共享.调用身份 import 用途 + + +class 任务存储: + def __init__(self, 连接: psycopg.Connection, 执行用途: 用途) -> None: + self.连接 = 连接 + self.用途 = 执行用途 + self.命名空间 = "evaluation" if 执行用途 is 用途.评测 else "public" + + def 查询(self, 语句: LiteralString, 参数: tuple = ()) -> Any: + """仅本模块固定 SQL 使用 schema 占位,不接收外部 SQL 或表名。""" + 游标 = self.连接.cursor(row_factory=dict_row) + return 游标.execute(sql.SQL(语句).format(s=sql.Identifier(self.命名空间)), 参数) + + def 登记计划(self, 计划: 执行计划) -> None: + 冻结 = 计划.冻结() + 哈希 = 内容哈希(冻结) + self.查询( + "INSERT INTO {s}.muse_flow_version (run_purpose,flow_id,version,definition," + "definition_hash) " + "VALUES (%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING", + (self.用途.value, 计划.流程ID, 计划.版本, Jsonb(冻结), 哈希), + ) + 已存 = self.查询( + "SELECT definition_hash FROM {s}.muse_flow_version WHERE run_purpose=%s AND " + "flow_id=%s AND version=%s", + (self.用途.value, 计划.流程ID, 计划.版本), + ).fetchone() + if 已存["definition_hash"] != 哈希: + raise 状态冲突("同一流程版本不能覆盖既有定义") + + def 读取计划(self, 流程ID: str, 版本: str) -> 执行计划: + 行 = self.查询( + "SELECT definition FROM {s}.muse_flow_version WHERE run_purpose=%s AND flow_id=%s " + "AND version=%s", + (self.用途.value, 流程ID, 版本), + ).fetchone() + if 行 is None: + raise 任务错误("流程版本尚未发布") + return 执行计划.从快照(行["definition"]) + + def 创建(self, 请求: 任务请求, 计划: 执行计划) -> str: + 冻结 = asdict(请求) + 哈希 = 内容哈希({"请求": 冻结, "流程": 计划.冻结()}) + 身份 = str(uuid.uuid4()) + 新建 = self.查询( + "INSERT INTO {s}.muse_task " + "(task_id,run_purpose,command_id,author_id,request_hash,frozen_input,flow_id," + "flow_version,flow_snapshot,state,scope_key) " + "VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,'queued',%s) ON CONFLICT DO NOTHING RETURNING " + "task_id", + ( + 身份, + self.用途.value, + 请求.命令ID, + 请求.作者, + 哈希, + Jsonb(冻结), + 计划.流程ID, + 计划.版本, + Jsonb(计划.冻结()), + 请求.排他范围.键 if 请求.排他范围 else None, + ), + ).fetchone() + if 新建 is None: + 既有 = self.查询( + "SELECT task_id,request_hash FROM {s}.muse_task WHERE run_purpose=%s AND " + "author_id=%s AND command_id=%s", + (self.用途.value, 请求.作者, 请求.命令ID), + ).fetchone() + if 既有["request_hash"] != 哈希: + raise 状态冲突("同一命令身份对应不同任务输入") + return str(既有["task_id"]) + for 步骤 in 计划.步骤: + self.查询( + "INSERT INTO {s}.muse_step (task_id,step_id,processor_id,processor_version," + "dependencies,state) " + "VALUES (%s,%s,%s,%s,%s,'pending')", + (身份, 步骤.步骤ID, 步骤.处理器ID, 步骤.处理器版本, list(步骤.依赖)), + ) + self.追加事件(身份, 事件类型.任务创建, {"流程ID": 计划.流程ID, "流程版本": 计划.版本}) + return 身份 + + def 任务行(self, 任务ID: str, *, 锁定: bool = False) -> dict: + 语句 = "SELECT * FROM {s}.muse_task WHERE task_id=%s AND run_purpose=%s" + 行 = self.查询(语句 + (" FOR UPDATE" if 锁定 else ""), (任务ID, self.用途.value)).fetchone() + if 行 is None: + raise 任务错误("任务不存在或不属于当前用途") + return 行 + + def 快照(self, 任务ID: str) -> 任务快照: + 行 = self.任务行(任务ID, 锁定=True) + 计划 = 执行计划.从快照(行["flow_snapshot"]) + 步骤 = self.查询( + "SELECT step_id,state,checkpoint,result,current_attempt FROM {s}.muse_step WHERE " + "task_id=%s ORDER BY array_position(%s::text[],step_id)", + (任务ID, [步.步骤ID for 步 in 计划.步骤]), + ).fetchall() + return 任务快照( + 任务ID, + 任务状态(行["state"]), + 行["author_id"], + 行["request_hash"], + 行["frozen_input"], + 计划, + tuple(步骤), + 行["last_sequence"], + ) + + def 列出(self, 作者: str, 作品ID: str | None) -> list[任务快照]: + 行 = self.查询( + "SELECT task_id FROM {s}.muse_task WHERE run_purpose=%s AND author_id=%s " + "AND (%s::text IS NULL OR frozen_input->>'作品ID'=%s) " + "ORDER BY created_at DESC,task_id LIMIT 100", + (self.用途.value, 作者, 作品ID, 作品ID), + ).fetchall() + return [self.快照(str(项["task_id"])) for 项 in 行] + + def 登记控制(self, 任务ID: str, 作者: str, 命令ID: str, 哈希: str) -> dict | None: + 新 = self.查询( + "INSERT INTO {s}.muse_task_control " + "(run_purpose,author_id,command_id,task_id,request_hash) " + "VALUES (%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING RETURNING command_id", + (self.用途.value, 作者, 命令ID, 任务ID, 哈希), + ).fetchone() + if 新: + return None + 已存 = self.查询( + "SELECT request_hash,receipt FROM {s}.muse_task_control WHERE " + "run_purpose=%s AND author_id=%s AND command_id=%s", + (self.用途.value, 作者, 命令ID), + ).fetchone() + if 已存["request_hash"] != 哈希: + raise 命令参数冲突("控制命令 ID 对应的参数不一致") + return 已存["receipt"] + + def 保存控制回执(self, 作者: str, 命令ID: str, 回执: dict) -> None: + self.查询( + "UPDATE {s}.muse_task_control SET receipt=%s WHERE " + "run_purpose=%s AND author_id=%s AND command_id=%s", + (Jsonb(回执), self.用途.value, 作者, 命令ID), + ) + + def 领取候选(self, 能力: list[str]) -> list[dict]: + return self.查询( + "SELECT t.*,s.step_id,s.processor_id,s.processor_version,s.checkpoint," + "s.current_attempt," + "a.call_state,a.lease_until FROM {s}.muse_task t JOIN {s}.muse_step s USING(task_id) " + "LEFT JOIN {s}.muse_attempt a ON a.attempt_id=s.current_attempt " + "LEFT JOIN {s}.muse_task_scope_lease scope ON scope.run_purpose=t.run_purpose AND " + "scope.scope_key=t.scope_key " + "WHERE t.run_purpose=%s AND t.state IN ('queued','running') AND s.processor_id=ANY(%s) " + "AND (t.scope_key IS NULL OR scope.holder_task_id IS NULL OR " + "scope.holder_task_id=t.task_id OR scope.lease_until<=clock_timestamp()) " + "AND (s.state='pending' OR (s.state='running' AND a.lease_until<=clock_timestamp())) " + "AND NOT EXISTS (SELECT 1 FROM {s}.muse_step d WHERE d.task_id=s.task_id " + "AND d.step_id=ANY(s.dependencies) AND d.state<>'completed') " + "ORDER BY t.created_at,s.step_id FOR UPDATE OF t,s SKIP LOCKED LIMIT 1", + (self.用途.value, 能力), + ).fetchall() + + def 取得作用域(self, 任务ID: str, 范围: str, 租期: float) -> int | None: + 行 = self.查询( + "INSERT INTO {s}.muse_task_scope_lease AS lease (run_purpose,scope_key," + "holder_task_id,generation,lease_until) " + "VALUES (%s,%s,%s,1,clock_timestamp()+%s*interval '1 second') " + "ON CONFLICT (run_purpose,scope_key) DO UPDATE SET " + "holder_task_id=EXCLUDED.holder_task_id," + "generation=CASE WHEN lease.holder_task_id=EXCLUDED.holder_task_id AND " + "lease.lease_until>clock_timestamp() " + "THEN lease.generation ELSE lease.generation+1 END, " + "lease_until=GREATEST(lease.lease_until,EXCLUDED.lease_until) " + "WHERE lease.holder_task_id=EXCLUDED.holder_task_id OR " + "lease.lease_until<=clock_timestamp() " + "OR lease.holder_task_id IS NULL RETURNING generation", + (self.用途.value, 范围, 任务ID, 租期), + ).fetchone() + return 行["generation"] if 行 else None + + def 有效作用域(self, 任务ID: str, 范围: str, 代次: int | None) -> bool: + return ( + self.查询( + "SELECT 1 FROM {s}.muse_task_scope_lease WHERE run_purpose=%s AND scope_key=%s " + "AND holder_task_id=%s AND generation=%s AND lease_until>clock_timestamp() FOR " + "UPDATE", + (self.用途.value, 范围, 任务ID, 代次), + ).fetchone() + is not None + ) + + def 释放作用域(self, 任务ID: str, 范围: str, 代次: int | None) -> None: + self.查询( + "UPDATE {s}.muse_task_scope_lease SET holder_task_id=NULL," + "lease_until=clock_timestamp() " + "WHERE run_purpose=%s AND scope_key=%s AND holder_task_id=%s AND generation=%s", + (self.用途.value, 范围, 任务ID, 代次), + ) + + def 新尝试(self, 任务: dict, 执行者: str, 租期: float, 代次: int | None) -> 领取凭证: + if 任务["current_attempt"]: + self.查询( + "UPDATE {s}.muse_attempt SET state='expired' WHERE attempt_id=%s", + (任务["current_attempt"],), + ) + 尝试ID, 令牌 = str(uuid.uuid4()), str(uuid.uuid4()) + 行 = self.查询( + "INSERT INTO {s}.muse_attempt (attempt_id,task_id,step_id,worker_id,lease_token," + "lease_until,scope_generation,state) " + "VALUES (%s,%s,%s,%s,%s,clock_timestamp()+%s*interval '1 second',%s,'running') " + "RETURNING lease_until", + (尝试ID, 任务["task_id"], 任务["step_id"], 执行者, 令牌, 租期, 代次), + ).fetchone() + self.查询( + "UPDATE {s}.muse_step SET state='running',current_attempt=%s WHERE task_id=%s AND " + "step_id=%s", + (尝试ID, 任务["task_id"], 任务["step_id"]), + ) + self.查询( + "UPDATE {s}.muse_task SET state='running',scope_generation=%s," + "updated_at=clock_timestamp() WHERE task_id=%s", + (代次, 任务["task_id"]), + ) + self.追加事件( + str(任务["task_id"]), + 事件类型.步骤开始, + {"执行者": 执行者}, + 步骤ID=任务["step_id"], + 尝试ID=尝试ID, + ) + return 领取凭证( + str(任务["task_id"]), + 任务["step_id"], + 尝试ID, + 令牌, + 执行者, + 行["lease_until"], + 代次, + 任务["processor_id"], + 任务["processor_version"], + 任务["checkpoint"], + ) + + def 尝试行(self, 领取: 领取凭证) -> dict | None: + return self.查询( + "SELECT a.*,s.state AS step_state,s.current_attempt,s.processor_id,s.processor_version " + "FROM {s}.muse_attempt a " + "JOIN {s}.muse_step s USING(task_id,step_id) WHERE a.attempt_id=%s AND a.task_id=%s " + "AND a.step_id=%s AND a.lease_token=%s AND a.worker_id=%s " + "AND a.lease_until>clock_timestamp() FOR UPDATE OF a,s", + (领取.尝试ID, 领取.任务ID, 领取.步骤ID, 领取.租约令牌, 领取.执行者), + ).fetchone() + + def 心跳(self, 领取: 领取凭证, 租期: float, 任务: dict) -> None: + self.查询( + "UPDATE {s}.muse_attempt SET heartbeat_at=clock_timestamp()," + "lease_until=clock_timestamp()+%s*interval '1 second' WHERE attempt_id=%s", + (租期, 领取.尝试ID), + ) + if 任务["scope_key"]: + self.查询( + "UPDATE {s}.muse_task_scope_lease SET lease_until=GREATEST(lease_until," + "clock_timestamp()+%s*interval '1 second') WHERE run_purpose=%s AND " + "scope_key=%s AND holder_task_id=%s AND generation=%s", + (租期, self.用途.value, 任务["scope_key"], 领取.任务ID, 领取.作用域代次), + ) + + def 保存检查点(self, 领取: 领取凭证, 检查点: dict) -> None: + self.查询( + "UPDATE {s}.muse_step SET checkpoint=%s WHERE task_id=%s AND step_id=%s", + (Jsonb(检查点), 领取.任务ID, 领取.步骤ID), + ) + self.追加事件(领取.任务ID, 事件类型.检查点保存, {}, 步骤ID=领取.步骤ID, 尝试ID=领取.尝试ID) + + def 完成步骤(self, 领取: 领取凭证, 输出: dict, 检查点: dict) -> bool: + self.查询( + "UPDATE {s}.muse_step SET state='completed',result=%s,checkpoint=%s WHERE " + "task_id=%s AND step_id=%s", + (Jsonb(输出), Jsonb(检查点), 领取.任务ID, 领取.步骤ID), + ) + self.查询( + "UPDATE {s}.muse_attempt SET state='completed',result=%s WHERE attempt_id=%s", + (Jsonb(输出), 领取.尝试ID), + ) + self.追加事件( + 领取.任务ID, + 事件类型.步骤完成, + {"输出哈希": 内容哈希(输出)}, + 步骤ID=领取.步骤ID, + 尝试ID=领取.尝试ID, + ) + return self.更新任务完成(领取.任务ID) + + def 更新任务完成(self, 任务ID: str) -> bool: + 未完成 = self.查询( + "SELECT 1 FROM {s}.muse_step WHERE task_id=%s AND state<>'completed' LIMIT 1", (任务ID,) + ).fetchone() + if not 未完成: + self.查询( + "UPDATE {s}.muse_task SET state='completed',updated_at=clock_timestamp() WHERE " + "task_id=%s", + (任务ID,), + ) + self.追加事件(任务ID, 事件类型.任务结束, {}) + return 未完成 is None + + def 设置未知(self, 任务ID: str, 步骤ID: str, 尝试ID: str) -> None: + self.查询( + "UPDATE {s}.muse_attempt SET state='unknown',call_state='unknown' WHERE attempt_id=%s", + (尝试ID,), + ) + self.查询( + "UPDATE {s}.muse_step SET state='unknown' WHERE task_id=%s AND step_id=%s", + (任务ID, 步骤ID), + ) + self.查询( + "UPDATE {s}.muse_task SET state='reconciling',updated_at=clock_timestamp() WHERE " + "task_id=%s", + (任务ID,), + ) + self.追加事件(任务ID, 事件类型.调用未知, {}, 步骤ID=步骤ID, 尝试ID=尝试ID) + + def 失败步骤(self, 领取: 领取凭证, 失败码: str) -> None: + self.查询( + "UPDATE {s}.muse_attempt SET state='failed',failure_code=%s WHERE attempt_id=%s", + (失败码, 领取.尝试ID), + ) + self.查询( + "UPDATE {s}.muse_step SET state='failed' WHERE task_id=%s AND step_id=%s", + (领取.任务ID, 领取.步骤ID), + ) + self.查询( + "UPDATE {s}.muse_task SET state='failed',updated_at=clock_timestamp() WHERE task_id=%s", + (领取.任务ID,), + ) + self.追加事件( + 领取.任务ID, + 事件类型.步骤失败, + {"失败码": 失败码}, + 步骤ID=领取.步骤ID, + 尝试ID=领取.尝试ID, + ) + + def 登记调用(self, 领取: 领取凭证, 调用引用: str) -> None: + self.查询( + "UPDATE {s}.muse_attempt SET call_state='sent',call_reference=%s WHERE attempt_id=%s", + (调用引用, 领取.尝试ID), + ) + self.追加事件( + 领取.任务ID, + 事件类型.调用发出, + {"调用引用": 调用引用}, + 步骤ID=领取.步骤ID, + 尝试ID=领取.尝试ID, + ) + + def 控制(self, 任务ID: str, 目标: str) -> None: + self.查询( + "UPDATE {s}.muse_task SET state=%s,updated_at=clock_timestamp() WHERE task_id=%s", + (目标, 任务ID), + ) + if 目标 in ("paused", "cancelled"): + self.查询( + "UPDATE {s}.muse_attempt SET state=%s WHERE task_id=%s AND state='running'", + (目标, 任务ID), + ) + elif 目标 == "queued": + self.查询( + "UPDATE {s}.muse_step SET state='pending' WHERE task_id=%s AND state IN " + "('running','failed','unknown')", + (任务ID,), + ) + + def 未知调用(self, 任务ID: str) -> list[dict]: + return self.查询( + "SELECT * FROM {s}.muse_attempt WHERE task_id=%s AND call_state IN ('sent','unknown')", + (任务ID,), + ).fetchall() + + def 对账(self, 任务ID: str, 尝试ID: str, 输出: dict | None, 回执: str) -> None: + 行 = self.查询( + "SELECT * FROM {s}.muse_attempt WHERE task_id=%s AND attempt_id=%s AND call_state " + "IN ('sent','unknown') FOR UPDATE", + (任务ID, 尝试ID), + ).fetchone() + if 行 is None: + raise 状态冲突("没有对应的待对账调用") + self.查询( + "UPDATE {s}.muse_attempt SET call_state=%s,state=%s,result=%s WHERE attempt_id=%s", + ( + "saved" if 输出 is not None else "not_sent", + "paused" if 输出 is not None else "failed", + Jsonb(输出), + 尝试ID, + ), + ) + self.查询( + "UPDATE {s}.muse_step SET state='pending',checkpoint=checkpoint || %s " + "WHERE task_id=%s AND step_id=%s AND " + "current_attempt=%s", + ( + Jsonb( + { + "已保存调用": { + "调用引用": 行["call_reference"], + "尝试ID": 尝试ID, + "输出": 输出, + "对账回执": 回执, + } + } + ), + 任务ID, + 行["step_id"], + 尝试ID, + ), + ) + self.追加事件( + 任务ID, + 事件类型.调用对账, + {"回执": 回执, "已保存输出": 输出 is not None}, + 步骤ID=行["step_id"], + 尝试ID=尝试ID, + ) + + def 追加事件( + self, + 任务ID: str, + 类型: 事件类型, + 载荷: dict, + *, + 步骤ID: str | None = None, + 尝试ID: str | None = None, + 事件ID: str | None = None, + ) -> 运行事件: + 身份 = 事件ID or str(uuid.uuid4()) + 哈希 = 内容哈希({"类型": 类型.value, "载荷": 载荷, "步骤": 步骤ID, "尝试": 尝试ID}) + 既有 = self.查询("SELECT * FROM {s}.muse_task_event WHERE event_id=%s", (身份,)).fetchone() + if 既有: + if str(既有["task_id"]) != 任务ID or 既有["payload_hash"] != 哈希: + raise 状态冲突("相同事件身份不能对应不同内容") + return self._事件(既有) + 序号 = self.查询( + "UPDATE {s}.muse_task SET last_sequence=last_sequence+1 WHERE task_id=%s RETURNING " + "last_sequence", + (任务ID,), + ).fetchone()["last_sequence"] + 行 = self.查询( + "INSERT INTO {s}.muse_task_event (event_id,task_id,step_id,attempt_id,sequence," + "event_type,payload_version,payload,payload_hash) VALUES (%s,%s,%s,%s,%s,%s,1,%s," + "%s) RETURNING *", + (身份, 任务ID, 步骤ID, 尝试ID, 序号, 类型.value, Jsonb(载荷), 哈希), + ).fetchone() + return self._事件(行) + + def 续接(self, 任务ID: str, 游标: int, 数量: int) -> 事件续接: + 任务 = self.任务行(任务ID, 锁定=True) + 最后 = 任务["last_sequence"] + 范围 = self.查询( + "SELECT min(sequence) AS first FROM {s}.muse_task_event WHERE task_id=%s", (任务ID,) + ).fetchone() + 首 = 范围["first"] or 最后 + 1 + 行 = self.查询( + "SELECT * FROM {s}.muse_task_event WHERE task_id=%s AND sequence>%s ORDER BY " + "sequence LIMIT %s", + (任务ID, 游标, 数量), + ).fetchall() + 缺口 = 游标 < 首 - 1 or 游标 > 最后 + for 下标, 项 in enumerate(行, 游标 + 1): + 缺口 |= 项["sequence"] != 下标 + if not 行 and 游标 < 最后: + 缺口 = True + return 事件续接(tuple(self._事件(项) for 项 in 行) if not 缺口 else (), 缺口, 首, 最后) + + @staticmethod + def _事件(行: dict) -> 运行事件: + return 运行事件( + str(行["event_id"]), + str(行["task_id"]), + 行["step_id"], + str(行["attempt_id"]) if 行["attempt_id"] else None, + 行["sequence"], + 事件类型(行["event_type"]), + 行["occurred_at"], + 行["payload_version"], + 行["payload"], + ) diff --git a/src/muse/任务运行/恢复控制.py b/src/muse/任务运行/恢复控制.py new file mode 100644 index 0000000..ca38d80 --- /dev/null +++ b/src/muse/任务运行/恢复控制.py @@ -0,0 +1,44 @@ +"""作者控制与恢复前检查;未知外部调用必须先对账。""" + +from __future__ import annotations + +from collections.abc import Callable + +from muse.任务运行.作用域互斥 import 释放范围 +from muse.任务运行.存储 import 任务存储 +from muse.任务运行.模型 import 事件类型, 任务快照, 任务状态, 状态冲突 + + +def 控制任务( + 存储: 任务存储, + 任务ID: str, + 作者: str, + 预期: 任务状态, + 动作: str, + 重验: Callable[[任务快照], None] | None = None, +) -> None: + 任务 = 存储.任务行(任务ID, 锁定=True) + if 任务["author_id"] != 作者 or 任务["state"] != 预期.value: + raise 状态冲突("任务作者或预期状态不匹配") + if 预期 in (任务状态.已取消, 任务状态.已完成): + raise 状态冲突("终态任务不能恢复或重新控制") + if 动作 == "取消": + 存储.控制(任务ID, "cancelled") + 类型 = 事件类型.任务取消 + elif 动作 == "暂停" and 预期 in (任务状态.待运行, 任务状态.运行中): + 存储.控制(任务ID, "paused") + 类型 = 事件类型.任务暂停 + elif 动作 == "恢复" and 预期 in (任务状态.已暂停, 任务状态.已失败, 任务状态.待对账): + if 存储.未知调用(任务ID): + raise 状态冲突("调用结果未知,必须先对账再恢复") + if 重验 is None: + raise 状态冲突("恢复需要业务 owner 重验来源、结构、授权和预算") + 重验(存储.快照(任务ID)) + 存储.控制(任务ID, "queued") + 类型 = 事件类型.任务恢复 + else: + raise 状态冲突("动作不适用于当前任务状态") + 释放范围(存储, 任务) + 存储.追加事件(任务ID, 类型, {}) + if 动作 == "恢复": + 存储.更新任务完成(任务ID) diff --git a/src/muse/任务运行/接口.py b/src/muse/任务运行/接口.py new file mode 100644 index 0000000..212dacd --- /dev/null +++ b/src/muse/任务运行/接口.py @@ -0,0 +1,316 @@ +"""任务运行唯一公开入口;所有状态和输出通过持久事务保存。""" + +from __future__ import annotations + +from collections.abc import Iterator +from contextlib import contextmanager + +from muse.任务运行.事件记录 import 校验尝试事件 +from muse.任务运行.任务领取 import 校验领取, 领取步骤 +from muse.任务运行.作用域互斥 import 释放范围 +from muse.任务运行.原文生命周期 import 原文服务, 原文错误, 校验原文字节, 校验期限 +from muse.任务运行.存储 import 任务存储 +from muse.任务运行.工具调用 import 只读工具集, 工具定义, 工具来源, 工具结果, 工具范围 +from muse.任务运行.恢复控制 import 控制任务 +from muse.任务运行.执行合同 import ( + 回合记录, + 工具回传, + 工具请求, + 已准备模型调用, + 模型协议错误, + 模型宿主, + 模型用量, + 模型结果, + 模型请求, +) +from muse.任务运行.模型 import ( + 事件类型, + 事件续接, + 任务快照, + 任务状态, + 任务请求, + 任务错误, + 作用域, + 内容哈希, + 处理器目录, + 执行上下文, + 执行计划, + 步骤处理器, + 步骤结果, + 步骤计划, + 状态冲突, + 租约失效, + 运行事件, + 领取凭证, +) +from muse.任务运行.模型调用 import ( + 校验模型输出, + 校验输出合同, + 模型交付, + 模型执行器, + 请求字节, + 调用计价, +) +from muse.任务运行.步骤执行 import 执行一步 +from muse.任务运行.角色会话 import 组装角色请求, 角色会话 +from muse.任务运行.角色策略 import 角色执行策略, 角色策略目录 +from muse.任务运行.运行证据 import 证据服务 +from muse.任务运行.配置版本 import ( + 凭据引用, + 提供方配置, + 运行配置内容, + 配置快照, + 配置版本管理, + 配置验证证据, +) +from muse.任务运行.预算管理 import ( + 任务预算计划, + 停止任务预算, + 恢复任务预算, + 角色预算, + 预算管理, + 额度策略, +) +from muse.基础设施.数据库.连接 import 数据库工厂 + + +class 任务服务: + def __init__(self, 数据库: 数据库工厂, 处理器: 处理器目录) -> None: + self.数据库 = 数据库 + self.处理器 = 处理器 + self._停止领取 = False + + @contextmanager + def _存储(self) -> Iterator[任务存储]: + with self.数据库.连接() as 连, 连.transaction(): + yield 任务存储(连, self.数据库.用途) + + def 发布计划(self, 计划: 执行计划) -> None: + self.处理器.核对计划(计划) + if not 计划.步骤: + raise 任务错误("流程必须包含步骤") + for 步骤 in 计划.步骤: + 处理器 = self.处理器.获取(步骤.处理器ID, 步骤.处理器版本) + if (处理器.输入合同, 处理器.输出合同) != (步骤.输入合同, 步骤.输出合同): + raise 任务错误("发布计划与登记处理器合同不一致") + with self._存储() as 存储: + 存储.登记计划(计划) + + def 创建任务(self, 请求: 任务请求, 流程ID: str, 版本: str) -> str: + if 请求.执行用途 is not self.数据库.用途: + raise 任务错误("任务用途与装配用途不一致") + with self._存储() as 存储: + 计划 = 存储.读取计划(流程ID, 版本) + self.处理器.核对计划(计划) + for 步骤 in 计划.步骤: + self.处理器.获取(步骤.处理器ID, 步骤.处理器版本) + return 存储.创建(请求, 计划) + + def 读取任务(self, 任务ID: str) -> 任务快照: + with self._存储() as 存储: + return 存储.快照(任务ID) + + def 列出任务(self, 作者: str, *, 作品ID: str | None = None) -> list[任务快照]: + with self._存储() as 存储: + return 存储.列出(作者, 作品ID) + + def 领取步骤(self, 执行者: str, 能力: list[str], *, 租期秒: float = 60) -> 领取凭证 | None: + if self._停止领取: + return None + with self._存储() as 存储: + return 领取步骤(存储, 执行者, 能力, 租期秒, self.处理器) + + def 停止领取(self) -> None: + """仅控制本进程领取循环;已经持久化的尝试仍可完成或由租约恢复。""" + self._停止领取 = True + + def 核对领取(self, 领取: 领取凭证) -> None: + with self._存储() as 存储: + 校验领取(存储, 领取) + + def 重验执行范围(self, 领取: 领取凭证) -> None: + """会话初始、恢复和每轮消费都调用具名owner的当前来源检查。""" + self.核对领取(领取) + self.处理器.重验任务(self.读取任务(领取.任务ID)) + + def 续租(self, 领取: 领取凭证, *, 租期秒: float = 60) -> None: + if 租期秒 <= 0: + raise 任务错误("租期必须大于零") + with self._存储() as 存储: + 任务, _ = 校验领取(存储, 领取) + 存储.心跳(领取, 租期秒, 任务) + + def 保存检查点(self, 领取: 领取凭证, 检查点: dict) -> None: + with self._存储() as 存储: + 校验领取(存储, 领取) + 存储.保存检查点(领取, 检查点) + + def 完成步骤(self, 领取: 领取凭证, 结果: 步骤结果) -> None: + with self._存储() as 存储: + 任务, 尝试 = 校验领取(存储, 领取) + if 尝试["call_state"] in ("sent", "unknown"): + raise 状态冲突("调用结果尚未可靠保存,不能完成步骤") + if 存储.完成步骤(领取, 结果.输出, 结果.检查点): + 停止任务预算(存储, 领取.任务ID) + 释放范围(存储, 任务) + + def 执行一步(self, 领取: 领取凭证, *, 租期秒: float | None = None) -> 步骤结果: + return 执行一步(self, 领取, 租期秒=租期秒) + + def 失败步骤(self, 领取: 领取凭证, 失败码: str) -> None: + with self._存储() as 存储: + 任务, 尝试 = 校验领取(存储, 领取) + if 尝试["call_state"] in ("sent", "unknown"): + 存储.设置未知(领取.任务ID, 领取.步骤ID, 领取.尝试ID) + else: + 存储.失败步骤(领取, 失败码) + 停止任务预算(存储, 领取.任务ID) + 释放范围(存储, 任务) + + def 追加事件(self, 领取: 领取凭证, 类型: 事件类型, 载荷: dict, *, 事件ID: str) -> 运行事件: + 校验尝试事件(类型, 载荷) + with self._存储() as 存储: + 校验领取(存储, 领取) + return 存储.追加事件( + 领取.任务ID, 类型, 载荷, 步骤ID=领取.步骤ID, 尝试ID=领取.尝试ID, 事件ID=事件ID + ) + + def 续接事件(self, 任务ID: str, 游标: int = 0, *, 数量: int = 100) -> 事件续接: + if 游标 < 0 or 数量 < 1: + raise 任务错误("事件游标与数量不合法") + with self._存储() as 存储: + return 存储.续接(任务ID, 游标, 数量) + + def 控制任务( + self, + 任务ID: str, + 作者: str, + 预期状态: 任务状态, + 动作: str, + *, + 命令ID: str, + ) -> dict: + if not 命令ID: + raise 任务错误("任务控制需要明确命令身份") + 哈希 = 内容哈希({"task_id": 任务ID, "expected_state": 预期状态, "action": 动作}) + with self._存储() as 存储: + 任务 = 存储.任务行(任务ID, 锁定=True) + if 任务["author_id"] != 作者: + raise 状态冲突("任务不属于当前作者") + 既有 = 存储.登记控制(任务ID, 作者, 命令ID, 哈希) + if 既有 is not None: + return 既有 + 控制任务(存储, 任务ID, 作者, 预期状态, 动作, self.处理器.重验任务) + if 动作 == "取消": + 停止任务预算(存储, 任务ID) + elif 动作 == "恢复": + 恢复任务预算(存储, 任务ID) + 当前 = 存储.快照(任务ID) + 回执 = { + "command_id": 命令ID, + "task_id": 任务ID, + "state": 当前.状态.value, + "last_sequence": 当前.最后序号, + } + 存储.保存控制回执(作者, 命令ID, 回执) + return 回执 + + def 登记外部调用(self, 领取: 领取凭证, 调用引用: str) -> None: + """发送前持久化调用引用;实际模型治理和预算由 W06 接入。""" + if not 调用引用: + raise 任务错误("调用引用不能为空") + with self._存储() as 存储: + _, 尝试 = 校验领取(存储, 领取) + if 尝试["call_state"] != "not_sent": + raise 状态冲突("当前尝试已登记调用,不能盲目重发") + 存储.登记调用(领取, 调用引用) + + def 对账调用( + self, 任务ID: str, 尝试ID: str, 作者: str, *, 已保存输出: dict | None, 对账回执: str + ) -> None: + """对账只固定调用事实并留恢复检查点;整个步骤由处理器恢复后完成。""" + if not 对账回执: + raise 任务错误("调用对账必须有可靠回执引用") + with self._存储() as 存储: + 任务 = 存储.任务行(任务ID, 锁定=True) + if 任务["author_id"] != 作者 or 任务["state"] not in ("reconciling", "paused"): + raise 状态冲突("任务不处于可对账状态或作者不匹配") + # 受控模型调用需要实际预算结算和原尝试的完整响应证据。 + 模型调用 = 存储.查询( + "SELECT b.* FROM {s}.muse_budget_reservation b JOIN {s}.muse_attempt a " + "ON a.task_id=b.task_id AND a.call_reference=b.call_id " + "WHERE a.task_id=%s AND a.attempt_id=%s", + (任务ID, 尝试ID), + ).fetchone() + if 模型调用: + if 模型调用["state"] != "settled" or not 已保存输出: + raise 状态冲突("模型成本与响应尚未可靠对账") + 证据 = 存储.查询( + "SELECT content FROM {s}.muse_runtime_evidence " + "WHERE task_id=%s AND attempt_id=%s AND evidence_id=%s " + "AND kind='model_response' AND metadata->>'call_id'=%s", + (任务ID, 尝试ID, 已保存输出.get("evidence_id"), 模型调用["call_id"]), + ).fetchone() + if 证据 is None or 证据["content"] is None: + raise 状态冲突("对账输出没有关联完整模型响应证据") + 存储.对账(任务ID, 尝试ID, 已保存输出, 对账回执) + + +__all__ = [ + "提供方配置", + "调用计价", + "配置快照", + "角色会话", + "组装角色请求", + "回合记录", + "工具回传", + "校验原文字节", + "模型执行器", + "模型交付", + "请求字节", + "证据服务", + "原文服务", + "原文错误", + "校验期限", + "只读工具集", + "工具范围", + "工具来源", + "工具定义", + "工具结果", + "角色策略目录", + "角色执行策略", + "任务预算计划", + "角色预算", + "预算管理", + "额度策略", + "凭据引用", + "运行配置内容", + "配置版本管理", + "配置验证证据", + "校验模型输出", + "校验输出合同", + "工具请求", + "模型协议错误", + "模型用量", + "模型结果", + "模型请求", + "模型宿主", + "已准备模型调用", + "任务服务", + "任务请求", + "任务快照", + "任务状态", + "任务错误", + "状态冲突", + "租约失效", + "作用域", + "步骤计划", + "执行计划", + "步骤处理器", + "步骤结果", + "执行上下文", + "领取凭证", + "事件类型", + "运行事件", + "事件续接", +] diff --git a/src/muse/任务运行/模型.py b/src/muse/任务运行/模型.py new file mode 100644 index 0000000..2ffa353 --- /dev/null +++ b/src/muse/任务运行/模型.py @@ -0,0 +1,220 @@ +"""任务、执行计划与租约的稳定合同;不包含存储和外部调用。""" + +from __future__ import annotations + +import hashlib +import json +from collections.abc import Callable +from dataclasses import asdict, dataclass, field +from datetime import datetime +from enum import StrEnum +from typing import Any, Protocol + +from muse.共享.调用身份 import 内容用途, 用途 +from muse.共享.错误 import Muse错误 + + +class 任务错误(Muse错误): + 错误码 = "MUSE_TASK" + + +class 租约失效(任务错误): + 错误码 = "MUSE_LEASE_STALE" + + +class 状态冲突(任务错误): + 错误码 = "MUSE_TASK_STATE_CONFLICT" + + +class 命令参数冲突(状态冲突): + 错误码 = "COMMAND_PAYLOAD_MISMATCH" + + +class 任务状态(StrEnum): + 待运行 = "queued" + 运行中 = "running" + 已暂停 = "paused" + 待对账 = "reconciling" + 已失败 = "failed" + 已取消 = "cancelled" + 已完成 = "completed" + + +class 事件类型(StrEnum): + 任务创建 = "task.created" + 步骤开始 = "step.started" + 文本片段 = "text.delta" + 工具读取 = "tool.read" + 候选保存 = "candidate.saved" + 检查完成 = "check.completed" + 等待作者 = "task.waiting" + 检查点保存 = "step.checkpoint" + 步骤完成 = "step.completed" + 步骤失败 = "step.failed" + 任务暂停 = "task.paused" + 任务恢复 = "task.resumed" + 任务取消 = "task.cancelled" + 任务结束 = "task.completed" + 调用发出 = "call.sent" + 调用未知 = "call.unknown" + 调用对账 = "call.reconciled" + + +def 内容哈希(内容: Any) -> str: + return hashlib.sha256( + json.dumps( + 内容, ensure_ascii=False, sort_keys=True, separators=(",", ":"), allow_nan=False + ).encode() + ).hexdigest() + + +@dataclass(frozen=True, slots=True) +class 作用域: + 对象类型: str + 对象ID: str + 操作族: str + + @property + def 键(self) -> str: + return 内容哈希(asdict(self)) + + +@dataclass(frozen=True, slots=True) +class 任务请求: + 任务类型: str + 命令ID: str + 作者: str + 执行用途: 用途 + 内容用途: 内容用途 + 输入: dict[str, Any] + 角色策略版本: str + 资源发布身份: str + 冻结上下文: dict[str, Any] + 作品ID: str | None = None + 来源ID: str | None = None + 排他范围: 作用域 | None = None + + def __post_init__(self) -> None: + for 名 in ("任务类型", "命令ID", "作者", "角色策略版本", "资源发布身份"): + if not getattr(self, 名): + raise 任务错误(f"任务缺少 {名}") + 必需 = {"source_scope", "schema_versions", "authorization", "budget", "stop_conditions"} + if not 必需 <= self.冻结上下文.keys(): + raise 任务错误("任务缺少来源、结构、授权、预算或停止条件的冻结输入") + if self.排他范围 is not None: + 关联 = {"work": self.作品ID, "source": self.来源ID} + 范围 = self.排他范围 + if not 范围.操作族 or not 范围.对象ID or 关联.get(范围.对象类型) != 范围.对象ID: + raise 任务错误("排他范围必须绑定任务已授权的作品或来源") + + +@dataclass(frozen=True, slots=True) +class 步骤计划: + 步骤ID: str + 处理器ID: str + 处理器版本: str + 依赖: tuple[str, ...] = () + 输入合同: str = "" + 输出合同: str = "" + 角色: str | None = None + 工具: tuple[str, ...] = () + + +@dataclass(frozen=True, slots=True) +class 执行计划: + 流程ID: str + 版本: str + 步骤: tuple[步骤计划, ...] + + def 冻结(self) -> dict[str, Any]: + return json.loads(json.dumps(asdict(self), ensure_ascii=False)) + + @classmethod + def 从快照(cls, 数据: dict[str, Any]) -> 执行计划: + return cls( + 数据["流程ID"], + 数据["版本"], + tuple( + 步骤计划(**{**项, "依赖": tuple(项["依赖"]), "工具": tuple(项["工具"])}) + for 项 in 数据["步骤"] + ), + ) + + +@dataclass(frozen=True, slots=True) +class 任务快照: + 任务ID: str + 状态: 任务状态 + 作者: str + 输入哈希: str + 冻结输入: dict[str, Any] + 流程: 执行计划 + 步骤: tuple[dict[str, Any], ...] + 最后序号: int + + +@dataclass(frozen=True, slots=True) +class 领取凭证: + 任务ID: str + 步骤ID: str + 尝试ID: str + 租约令牌: str + 执行者: str + 到期时间: datetime + 作用域代次: int | None + 处理器ID: str + 处理器版本: str + 检查点: dict[str, Any] + + +@dataclass(frozen=True, slots=True) +class 执行上下文: + 领取: 领取凭证 + 任务: 任务快照 + + +@dataclass(frozen=True, slots=True) +class 步骤结果: + 输出: dict[str, Any] + 检查点: dict[str, Any] = field(default_factory=dict) + + +@dataclass(frozen=True, slots=True) +class 步骤处理器: + 身份: str + 版本: str + 执行: Callable[[执行上下文], 步骤结果] + 输入合同: str + 输出合同: str + 保护职责: str | None = None + 角色: str | None = None + 允许工具: tuple[str, ...] = () + + +class 处理器目录(Protocol): + def 获取(self, 身份: str, 版本: str) -> 步骤处理器: ... + + def 核对计划(self, 计划: 执行计划) -> None: ... + + def 重验任务(self, 任务: 任务快照) -> None: ... + + +@dataclass(frozen=True, slots=True) +class 运行事件: + 事件ID: str + 任务ID: str + 步骤ID: str | None + 尝试ID: str | None + 序号: int + 类型: 事件类型 + 发生时间: datetime + 载荷版本: int + 载荷: dict[str, Any] + + +@dataclass(frozen=True, slots=True) +class 事件续接: + 事件: tuple[运行事件, ...] + 需要快照: bool + 可用首序号: int + 最后序号: int diff --git a/src/muse/任务运行/步骤执行.py b/src/muse/任务运行/步骤执行.py new file mode 100644 index 0000000..bea6909 --- /dev/null +++ b/src/muse/任务运行/步骤执行.py @@ -0,0 +1,77 @@ +"""一步执行拥有续租与完成边界;停止心跳后才提交最终结果。""" + +from __future__ import annotations + +from collections.abc import Iterator +from contextlib import contextmanager +from threading import Event, Thread +from time import monotonic +from typing import TYPE_CHECKING + +from muse.任务运行.模型 import 任务错误, 执行上下文, 步骤结果, 领取凭证 +from muse.基础设施.可观测性 import 记录运行 + +if TYPE_CHECKING: + from muse.任务运行.接口 import 任务服务 + + +@contextmanager +def _保持租约(服务: 任务服务, 领取: 领取凭证, 租期秒: float | None) -> Iterator[None]: + if 租期秒 is None: + yield + return + 完成 = Event() + 失败: list[Exception] = [] + + def 心跳() -> None: + while not 完成.wait(租期秒 / 3): + try: + 服务.续租(领取, 租期秒=租期秒) + except Exception as exc: + 失败.append(exc) + return + + 线程 = Thread(target=心跳, name="muse-lease-heartbeat", daemon=True) + 线程.start() + try: + yield + finally: + 完成.set() + 线程.join(timeout=max(1, 租期秒)) + if 线程.is_alive(): + raise 任务错误("续租尚未结束,不能提交步骤结果") + if 失败: + raise 失败[0] + + +def 执行一步(服务: 任务服务, 领取: 领取凭证, *, 租期秒: float | None = None) -> 步骤结果: + 服务.核对领取(领取) + 处理器 = 服务.处理器.获取(领取.处理器ID, 领取.处理器版本) + 上下文 = 执行上下文(领取, 服务.读取任务(领取.任务ID)) + 开始 = monotonic() + try: + with _保持租约(服务, 领取, 租期秒): + 结果 = 处理器.执行(上下文) + if not isinstance(结果, 步骤结果): + raise TypeError("处理器没有返回步骤结果") + 服务.完成步骤(领取, 结果) + 记录运行( + "step.finished", + 任务ID=领取.任务ID, + 步骤ID=领取.步骤ID, + 尝试ID=领取.尝试ID, + 状态="completed", + 耗时秒=monotonic() - 开始, + ) + return 结果 + except Exception as exc: + 记录运行( + "step.failed", + 任务ID=领取.任务ID, + 步骤ID=领取.步骤ID, + 尝试ID=领取.尝试ID, + 错误码=type(exc).__name__, + 耗时秒=monotonic() - 开始, + ) + 服务.失败步骤(领取, type(exc).__name__) + raise diff --git a/src/muse/基础设施/__init__.py b/src/muse/基础设施/__init__.py new file mode 100644 index 0000000..28e0b09 --- /dev/null +++ b/src/muse/基础设施/__init__.py @@ -0,0 +1,4 @@ +"""基础设施:数据库、凭据、宿主与受控文件等环境实现。 + +边界:只提供环境能力,不放业务规则;业务模块经接口使用。 +""" diff --git a/src/muse/基础设施/可观测性.py b/src/muse/基础设施/可观测性.py new file mode 100644 index 0000000..7fb2e98 --- /dev/null +++ b/src/muse/基础设施/可观测性.py @@ -0,0 +1,37 @@ +"""只记录可关联的技术外壳;正文、模型输入和凭据不进入应用日志。""" + +import json +import logging + +日志 = logging.getLogger("muse") + + +def 记录运行( + 事件: str, + *, + 请求ID: str | None = None, + 任务ID: str | None = None, + 步骤ID: str | None = None, + 尝试ID: str | None = None, + 状态: str | int | None = None, + 错误码: str | None = None, + 耗时秒: float | None = None, +) -> None: + """字段白名单避免把请求字典或异常正文顺手传入日志。""" + 内容 = { + "event": 事件, + "request_id": 请求ID, + "task_id": 任务ID, + "step_id": 步骤ID, + "attempt_id": 尝试ID, + "status": 状态, + "code": 错误码, + "duration_seconds": 耗时秒, + } + 日志.info(json.dumps({k: v for k, v in 内容.items() if v is not None}, ensure_ascii=False)) + + +def 配置日志() -> None: + 日志.setLevel(logging.INFO) + if not 日志.handlers: + 日志.addHandler(logging.StreamHandler()) diff --git a/src/muse/接入/cli/任务命令.py b/src/muse/接入/cli/任务命令.py new file mode 100644 index 0000000..496d61a --- /dev/null +++ b/src/muse/接入/cli/任务命令.py @@ -0,0 +1,33 @@ +"""本地作者查看或控制持久任务;只转交公共接口。""" + +from dataclasses import asdict +from pathlib import Path + +from pydantic import ValidationError + +from muse.共享.错误 import 配置错误 +from muse.启动 import 应用装配 +from muse.接入.主会话 import 任务控制请求, 参数合同错误, 执行任务控制 + + +def 运行任务命令(装配: 应用装配, 动作: str, 目标: str | None, 游标: int = 0) -> dict | list: + if 装配.配置.HTTP is None or 装配.任务运行 is None: + raise 配置错误("任务命令需要已配置的作者身份和任务服务") + 作者 = 装配.配置.HTTP.作者ID + 服务 = 装配.任务运行 + if 动作 == "列表": + return [asdict(项) for 项 in 服务.列出任务(作者)] + if not 目标: + raise 配置错误("任务命令缺少任务身份或请求文件") + if 动作 == "控制": + try: + 请求 = 任务控制请求.model_validate_json(Path(目标).read_text(encoding="utf-8")) + except ValidationError as exc: + raise 参数合同错误(exc.errors(), 前缀=("body",)) from None + return 执行任务控制(服务, 作者, 请求).model_dump(mode="json") + 任务 = 服务.读取任务(目标) + if 任务.作者 != 作者: + raise 配置错误("任务不属于当前作者") + if 动作 == "事件": + return asdict(服务.续接事件(目标, 游标)) + return asdict(任务) diff --git a/src/muse/接入/http/路由/任务运行.py b/src/muse/接入/http/路由/任务运行.py new file mode 100644 index 0000000..3d2fbfb --- /dev/null +++ b/src/muse/接入/http/路由/任务运行.py @@ -0,0 +1,210 @@ +"""任务查询和作者控制的 HTTP 外壳。""" + +from typing import Annotated, Literal + +from fastapi import APIRouter, Body, HTTPException, Request +from pydantic import AwareDatetime, BaseModel, ConfigDict, Field + +from muse.任务运行.接口 import 任务快照, 任务状态 +from muse.共享.调用身份 import 内容用途 +from muse.接入.http.作者会话 import 作者依赖 +from muse.接入.http.依赖 import 取得任务服务, 读取作者任务 +from muse.接入.主会话 import 任务控制回执, 任务控制请求, 执行任务控制 + +路由 = APIRouter(prefix="/api/v1/tasks", tags=["任务"]) + + +class 步骤响应(BaseModel): + step_id: str + state: str + + +class 任务响应(BaseModel): + task_id: str + state: 任务状态 + work_id: str | None + task_type: str + last_sequence: int + steps: list[步骤响应] + + +def 呈现任务(任务: 任务快照) -> 任务响应: + return 任务响应( + task_id=任务.任务ID, + state=任务.状态, + work_id=任务.冻结输入.get("作品ID"), + task_type=任务.冻结输入["任务类型"], + last_sequence=任务.最后序号, + steps=[步骤响应(step_id=步["step_id"], state=步["state"]) for 步 in 任务.步骤], + ) + + +@路由.get("/{task_id}", response_model=任务响应, operation_id="read_task") +def 读取任务( + task_id: str, request: Request, 作者: 作者依赖, work_id: str | None = None +) -> 任务响应: + return 呈现任务(读取作者任务(request, task_id, work_id)) + + +@路由.get("", response_model=list[任务响应], operation_id="list_tasks") +def 列出任务(request: Request, 作者: 作者依赖, work_id: str | None = None) -> list[任务响应]: + return [呈现任务(项) for 项 in 取得任务服务(request).列出任务(作者, 作品ID=work_id)] + + +@路由.post("/{task_id}/controls", response_model=任务控制回执, operation_id="control_task") +def 控制任务(task_id: str, 请求: 任务控制请求, request: Request, 作者: 作者依赖) -> 任务控制回执: + if 请求.target_ref != task_id: + raise HTTPException(422, "请求目标与任务路径不一致") + return 执行任务控制(取得任务服务(request), 作者, 请求) + + +class 原文授权请求(BaseModel): + model_config = ConfigDict(extra="forbid") + command_id: str = Field(min_length=1) + source_version: str = Field(min_length=1) + content_hashes: list[str] = Field(min_length=1) + content_purpose: 内容用途 + retention_mode: Literal["temporary", "archive", "persistent"] + valid_until: AwareDatetime + call_id: str | None = None + call_request_hash: str | None = None + + +class 原文租约请求(BaseModel): + model_config = ConfigDict(extra="forbid") + authorization_id: str + command_id: str = Field(min_length=1) + retain_until: AwareDatetime + min_remaining: float = Field(ge=0) + + +class 原文归档请求(BaseModel): + model_config = ConfigDict(extra="forbid") + authorization_id: str + + +class 会话授权请求(BaseModel): + model_config = ConfigDict(extra="forbid") + command_id: str = Field(min_length=1) + step_id: str = Field(min_length=1) + stage: str = Field(min_length=1) + initial_request_hash: str = Field(pattern="^[0-9a-f]{64}$") + max_model_calls: int = Field(gt=0, strict=True) + max_tool_calls: int = Field(ge=0, strict=True) + valid_until: AwareDatetime + + +@路由.post("/{task_id}/role-session-authorizations", operation_id="approve_role_session") +def 批准角色会话(task_id: str, 请求: 会话授权请求, request: Request, 作者: 作者依赖) -> dict: + 读取作者任务(request, task_id) + 身份 = request.app.state.装配.要求原文().批准会话( + task_id, + 作者, + 请求.command_id, + 步骤ID=请求.step_id, + 阶段=请求.stage, + 初始请求哈希=请求.initial_request_hash, + 最大模型调用=请求.max_model_calls, + 最大工具调用=请求.max_tool_calls, + 有效期=请求.valid_until, + ) + return {"authorization_id": 身份} + + +@路由.delete( + "/{task_id}/raw-authorizations/{authorization_id}", operation_id="revoke_raw_retention" +) +def 撤销原文批准(task_id: str, authorization_id: str, request: Request, 作者: 作者依赖) -> dict: + 读取作者任务(request, task_id) + request.app.state.装配.要求原文().撤销批准(task_id, 作者, authorization_id) + return {"authorization_id": authorization_id, "revoked": True} + + +@路由.post("/{task_id}/raw-authorizations", operation_id="approve_raw_retention") +def 批准原文(task_id: str, 请求: 原文授权请求, request: Request, 作者: 作者依赖) -> dict: + 读取作者任务(request, task_id) + 身份 = request.app.state.装配.要求原文().批准保留( + task_id, + 作者, + 请求.command_id, + 来源版本=请求.source_version, + 哈希=tuple(请求.content_hashes), + 内容用途=请求.content_purpose.value, + 方式=请求.retention_mode, + 有效期=请求.valid_until, + 调用ID=请求.call_id, + 调用请求哈希=请求.call_request_hash, + ) + return {"authorization_id": 身份} + + +@路由.post("/{task_id}/raw-leases", operation_id="create_raw_lease") +def 创建原文租约(task_id: str, 请求: 原文租约请求, request: Request, 作者: 作者依赖) -> dict: + 读取作者任务(request, task_id) + 身份 = request.app.state.装配.要求原文().创建租约( + task_id, + 请求.authorization_id, + 请求.command_id, + 请求.retain_until, + 最少剩余秒=请求.min_remaining, + ) + return {"lease_id": 身份} + + +@路由.put("/{task_id}/raw-leases/{lease_id}/content/{content_hash}", operation_id="write_raw_bytes") +def 写入原文( + task_id: str, + lease_id: str, + content_hash: str, + request: Request, + 作者: 作者依赖, + 数据: Annotated[bytes, Body(media_type="application/octet-stream")], +) -> dict: + 读取作者任务(request, task_id) + request.app.state.装配.要求原文().写入(task_id, lease_id, content_hash, 数据) + return {"content_hash": content_hash, "bytes": len(数据)} + + +@路由.post("/{task_id}/raw-leases/{lease_id}/archive", operation_id="archive_raw_bytes") +def 归档原文( + task_id: str, + lease_id: str, + 请求: 原文归档请求, + request: Request, + 作者: 作者依赖, +) -> dict: + 读取作者任务(request, task_id) + return request.app.state.装配.要求原文().归档(task_id, lease_id, 请求.authorization_id) + + +@路由.delete("/{task_id}/raw-leases/{lease_id}", operation_id="cleanup_raw_lease") +def 清理原文(task_id: str, lease_id: str, request: Request, 作者: 作者依赖) -> dict: + 读取作者任务(request, task_id) + 服务 = request.app.state.装配.要求原文() + 服务.清理(task_id, lease_id, 作者=作者) + return 服务.状态(task_id, lease_id) + + +@路由.post("/{task_id}/raw-recovery", operation_id="recover_raw_leases") +def 恢复原文(task_id: str, request: Request, 作者: 作者依赖) -> list[dict]: + 读取作者任务(request, task_id) + return request.app.state.装配.要求原文().恢复任务原文(task_id, 作者) + + +@路由.get("/{task_id}/evidence/{evidence_id}", operation_id="read_evidence_receipt") +def 读取证据回执(task_id: str, evidence_id: str, request: Request, 作者: 作者依赖) -> dict: + 读取作者任务(request, task_id) + return request.app.state.装配.要求证据().读取回执(task_id, evidence_id) + + +@路由.put("/{task_id}/evidence/{evidence_id}/content", operation_id="backfill_evidence_bytes") +def 补交证据( + task_id: str, + evidence_id: str, + authorization_id: str, + request: Request, + 作者: 作者依赖, + 数据: Annotated[bytes, Body(media_type="application/octet-stream")], +) -> dict: + 读取作者任务(request, task_id) + return request.app.state.装配.要求证据().补交原文(task_id, evidence_id, 数据, authorization_id) diff --git a/src/muse/接入/执行器.py b/src/muse/接入/执行器.py new file mode 100644 index 0000000..4bba89f --- /dev/null +++ b/src/muse/接入/执行器.py @@ -0,0 +1,34 @@ +"""后台领取与心跳;停机停止新领取,当前步骤按登记合同收尾。""" + +from threading import Event + +from muse.任务运行.接口 import 任务服务, 任务错误 + + +class 执行器: + def __init__(self, 服务: 任务服务, 身份: str, 能力: list[str], *, 租期秒: float = 60): + if not 身份 or not 能力 or 租期秒 <= 0: + raise 任务错误("执行器需要身份、已登记能力和正租期") + self.服务 = 服务 + self.身份 = 身份 + self.能力 = 能力 + self.租期秒 = 租期秒 + self.停止信号 = Event() + + def 停止(self) -> None: + self.停止信号.set() + self.服务.停止领取() + + def 运行一次(self) -> bool: + if self.停止信号.is_set(): + return False + 领取 = self.服务.领取步骤(self.身份, self.能力, 租期秒=self.租期秒) + if 领取 is None: + return False + self.服务.执行一步(领取, 租期秒=self.租期秒) + return True + + def 运行(self, *, 轮询秒: float = 1) -> None: + while not self.停止信号.is_set(): + if not self.运行一次(): + self.停止信号.wait(轮询秒) diff --git a/src/muse/编排/__init__.py b/src/muse/编排/__init__.py new file mode 100644 index 0000000..598a400 --- /dev/null +++ b/src/muse/编排/__init__.py @@ -0,0 +1 @@ +"""代码登记的流程与保护约束;公开入口见接口模块。""" diff --git a/src/muse/编排/接口.py b/src/muse/编排/接口.py new file mode 100644 index 0000000..b57d039 --- /dev/null +++ b/src/muse/编排/接口.py @@ -0,0 +1,24 @@ +"""流程定义、登记与发起的公开面;任务状态只由任务运行模块持久化。""" + +from muse.任务运行.接口 import 任务快照, 任务服务, 任务请求 +from muse.编排.槽位约束 import 流程校验错误 +from muse.编排.流程版本 import 流程定义 +from muse.编排.流程登记 import 流程登记 + + +class 流程服务: + def __init__(self, 运行: 任务服务, 登记: 流程登记) -> None: + self.运行 = 运行 + self.登记 = 登记 + + def 发布(self, 定义: 流程定义) -> None: + self.运行.发布计划(self.登记.冻结(定义)) + + def 发起(self, 请求: 任务请求, 流程ID: str, 版本: str) -> str: + return self.运行.创建任务(请求, 流程ID, 版本) + + def 查询(self, 任务ID: str) -> 任务快照: + return self.运行.读取任务(任务ID) + + +__all__ = ["流程服务", "流程定义", "流程登记", "流程校验错误"] diff --git a/src/muse/编排/槽位约束.py b/src/muse/编排/槽位约束.py new file mode 100644 index 0000000..f57e1e2 --- /dev/null +++ b/src/muse/编排/槽位约束.py @@ -0,0 +1,46 @@ +"""依赖是全部满足才可执行;保护节点必须覆盖每个终点的依赖闭包。""" + +from __future__ import annotations + +from muse.任务运行.接口 import 执行计划, 步骤处理器 +from muse.共享.错误 import Muse错误 + + +class 流程校验错误(Muse错误): + 错误码 = "MUSE_FLOW_INVALID" + + +def 核对流程(计划: 执行计划, 处理器: dict[str, 步骤处理器], 必需保护: tuple[str, ...]) -> None: + 节点 = {项.步骤ID: 项 for 项 in 计划.步骤} + if not 节点 or len(节点) != len(计划.步骤): + raise 流程校验错误("流程步骤不能为空或重名") + 祖先: dict[str, set[str]] = {} + 访问中: set[str] = set() + + def 遍历(身份: str) -> set[str]: + if 身份 in 访问中: + raise 流程校验错误("流程依赖存在环") + if 身份 not in 节点: + raise 流程校验错误("依赖节点不存在") + if 身份 in 祖先: + return 祖先[身份] + 访问中.add(身份) + 结果: set[str] = set() + for 依赖 in 节点[身份].依赖: + 结果 |= 遍历(依赖) | {依赖} + if 节点[依赖].输出合同 != 节点[身份].输入合同: + raise 流程校验错误("相邻步骤的输入输出合同不兼容") + 访问中.remove(身份) + 祖先[身份] = 结果 + return 结果 + + for 身份 in 节点: + 遍历(身份) + 处理 = 处理器[身份] + if 节点[身份].角色 != 处理.角色 or not set(节点[身份].工具) <= set(处理.允许工具): + raise 流程校验错误("槽位角色或工具许可不兼容") + 已被依赖 = {依赖 for 项 in 计划.步骤 for 依赖 in 项.依赖} + for 终点 in 节点.keys() - 已被依赖: + 路径保护 = {处理器[身份].保护职责 for 身份 in 祖先[终点] | {终点}} + if not set(必需保护) <= 路径保护: + raise 流程校验错误("存在绕过必需保护的流程终点") diff --git a/src/muse/编排/流程版本.py b/src/muse/编排/流程版本.py new file mode 100644 index 0000000..e6f1802 --- /dev/null +++ b/src/muse/编排/流程版本.py @@ -0,0 +1,12 @@ +"""流程定义是不可变输入,发布后任务保存固定执行计划。""" + +from dataclasses import dataclass + +from muse.任务运行.接口 import 步骤计划 + + +@dataclass(frozen=True, slots=True) +class 流程定义: + 流程ID: str + 版本: str + 步骤: tuple[步骤计划, ...] diff --git a/src/muse/编排/流程登记.py b/src/muse/编排/流程登记.py new file mode 100644 index 0000000..0a85faa --- /dev/null +++ b/src/muse/编排/流程登记.py @@ -0,0 +1,77 @@ +"""登记实际 Python 处理器;配置只能选择稳定身份与版本。""" + +from __future__ import annotations + +from collections.abc import Callable +from dataclasses import replace + +from muse.任务运行.接口 import 任务快照, 执行计划, 步骤处理器 +from muse.编排.槽位约束 import 核对流程, 流程校验错误 +from muse.编排.流程版本 import 流程定义 + + +class 流程登记: + def __init__(self) -> None: + self._处理器: dict[tuple[str, str], 步骤处理器] = {} + self._类型: dict[str, tuple[str, ...]] = {} + self._恢复检查: dict[tuple[str, str], Callable[[任务快照], None]] = {} + + def 登记恢复检查(self, 流程ID: str, 版本: str, 检查: Callable[[任务快照], None]) -> None: + if 流程ID not in self._类型 or not 版本 or not callable(检查): + raise 流程校验错误("恢复检查必须绑定已登记流程的固定版本") + if (流程ID, 版本) in self._恢复检查: + raise 流程校验错误("恢复检查版本已登记,不能覆盖") + self._恢复检查[(流程ID, 版本)] = 检查 + + def 重验任务(self, 任务: 任务快照) -> None: + self.核对计划(任务.流程) + 检查 = self._恢复检查.get((任务.流程.流程ID, 任务.流程.版本)) + if 检查 is None: + raise 流程校验错误("该流程版本的来源、结构、授权及预算重验尚未装配") + 检查(任务) + + def 登记处理器(self, 处理器: 步骤处理器) -> None: + 身份 = (处理器.身份, 处理器.版本) + if not all(身份) or not callable(处理器.执行): + raise 流程校验错误("只允许登记实际实现的处理器") + if 身份 in self._处理器: + raise 流程校验错误("处理器版本已登记,不能覆盖") + self._处理器[身份] = 处理器 + + def 登记类型(self, 流程ID: str, *, 必需保护: tuple[str, ...]) -> None: + if not 流程ID or 流程ID in self._类型: + raise 流程校验错误("流程类型为空或已登记") + self._类型[流程ID] = 必需保护 + + def 获取(self, 身份: str, 版本: str) -> 步骤处理器: + try: + return self._处理器[(身份, 版本)] + except KeyError: + raise 流程校验错误("处理器身份与版本未登记") from None + + def 冻结(self, 定义: 流程定义) -> 执行计划: + if 定义.流程ID not in self._类型 or not 定义.版本: + raise 流程校验错误("流程类型未登记或版本为空") + 处理器 = {项.步骤ID: self.获取(项.处理器ID, 项.处理器版本) for 项 in 定义.步骤} + 计划 = 执行计划( + 定义.流程ID, + 定义.版本, + tuple( + replace( + 项, 输入合同=处理器[项.步骤ID].输入合同, 输出合同=处理器[项.步骤ID].输出合同 + ) + for 项 in 定义.步骤 + ), + ) + 核对流程(计划, 处理器, self._类型[定义.流程ID]) + return 计划 + + def 核对计划(self, 计划: 执行计划) -> None: + if 计划.流程ID not in self._类型: + raise 流程校验错误("流程类型未登记") + 处理器 = {项.步骤ID: self.获取(项.处理器ID, 项.处理器版本) for 项 in 计划.步骤} + for 项 in 计划.步骤: + 处理 = 处理器[项.步骤ID] + if (项.输入合同, 项.输出合同) != (处理.输入合同, 处理.输出合同): + raise 流程校验错误("冻结计划的处理器合同不一致") + 核对流程(计划, 处理器, self._类型[计划.流程ID]) diff --git a/tests/单元/test_流程登记与保护.py b/tests/单元/test_流程登记与保护.py new file mode 100644 index 0000000..01e13ae --- /dev/null +++ b/tests/单元/test_流程登记与保护.py @@ -0,0 +1,57 @@ +"""流程发布前的处理器身份、版本与保护路径合同。""" + +import pytest + +from muse.任务运行.接口 import 步骤处理器, 步骤结果, 步骤计划 +from muse.编排.接口 import 流程定义, 流程校验错误, 流程登记 + + +def 构造登记() -> 流程登记: + 登记 = 流程登记() + 登记.登记处理器( + 步骤处理器("guard", "1", lambda _: 步骤结果({}), "v1", "v1", 保护职责="来源校验") + ) + 登记.登记处理器(步骤处理器("work", "1", lambda _: 步骤结果({}), "v1", "v1")) + 登记.登记类型("example", 必需保护=("来源校验",)) + return 登记 + + +def test_未知处理器与保护绕行拒绝__a50001() -> None: + 登记 = 构造登记() + with pytest.raises(流程校验错误): + 登记.冻结(流程定义("example", "1", (步骤计划("s1", "/tmp/script.py", "1"),))) + with pytest.raises(流程校验错误, match="保护"): + 登记.冻结( + 流程定义( + "example", "1", (步骤计划("guard", "guard", "1"), 步骤计划("work", "work", "1")) + ) + ) + 计划 = 登记.冻结( + 流程定义( + "example", + "1", + (步骤计划("guard", "guard", "1"), 步骤计划("work", "work", "1", ("guard",))), + ) + ) + assert 计划.流程ID == "example" + + +def test_流程依赖环与不兼容槽位拒绝__a50002() -> None: + 登记 = 构造登记() + with pytest.raises(流程校验错误, match="环"): + 登记.冻结( + 流程定义( + "example", + "1", + (步骤计划("g", "guard", "1", ("w",)), 步骤计划("w", "work", "1", ("g",))), + ) + ) + 登记.登记处理器(步骤处理器("incompatible", "1", lambda _: 步骤结果({}), "v2", "v2")) + with pytest.raises(流程校验错误, match="合同"): + 登记.冻结( + 流程定义( + "example", + "1", + (步骤计划("g", "guard", "1"), 步骤计划("w", "incompatible", "1", ("g",))), + ) + ) diff --git a/tests/集成/test_任务租约与恢复.py b/tests/集成/test_任务租约与恢复.py new file mode 100644 index 0000000..9457b56 --- /dev/null +++ b/tests/集成/test_任务租约与恢复.py @@ -0,0 +1,439 @@ +"""W05 公开任务接口的持久状态、独立连接竞争与恢复合同。""" + +from __future__ import annotations + +import dataclasses +import subprocess +import sys +import time +import uuid +from concurrent.futures import ThreadPoolExecutor +from pathlib import Path +from threading import Barrier + +import psycopg +import pytest +from psycopg import sql + +from muse.任务运行.接口 import ( + 事件类型, + 任务服务, + 任务状态, + 任务请求, + 任务错误, + 作用域, + 步骤处理器, + 步骤结果, + 步骤计划, + 状态冲突, + 租约失效, +) +from muse.共享.调用身份 import 内容用途, 用途 +from muse.基础设施.数据库.迁移 import 执行迁移 +from muse.基础设施.数据库.连接 import 数据库工厂 +from muse.编排.接口 import 流程定义, 流程服务, 流程登记 +from muse.配置 import 数据库引用 + +pytestmark = pytest.mark.数据库 +根 = Path(__file__).resolve().parents[2] + + +@pytest.fixture +def 任务库(隔离数据库URL: str, monkeypatch: pytest.MonkeyPatch, tmp_path: Path): + 库名 = f"muse_task_{uuid.uuid4().hex[:12]}" + with psycopg.connect(隔离数据库URL, autocommit=True) as 连: + 连.execute((根 / "数据库/初始化/用途角色.sql").read_text()) + 连.execute( + sql.SQL("CREATE DATABASE {} OWNER muse_maint TEMPLATE template0").format( + sql.Identifier(库名) + ) + ) + 工厂 = {} + try: + for 声明, 角色 in [ + (用途.维护, "muse_maint"), + (用途.生产, "muse_app"), + (用途.评测, "muse_eval"), + ]: + 参数 = psycopg.conninfo.conninfo_to_dict(隔离数据库URL) + 参数.update(dbname=库名, user=角色) + 引用名 = f"MUSE_W05_{声明.name}_URL" + monkeypatch.setenv(引用名, psycopg.conninfo.make_conninfo(**参数)) + 工厂[声明] = 数据库工厂(数据库引用("环境变量", 引用名), 声明) + for 名称 in [ + "V0001__共享标识与版本.sql", + "V0004__任务运行与证据.sql", + "V0016__流程定义与版本.sql", + ]: + (tmp_path / 名称).write_bytes((根 / "数据库/迁移" / 名称).read_bytes()) + with 工厂[用途.维护].连接() as 连: + 执行迁移(连, tmp_path) + yield 工厂 + finally: + with psycopg.connect(隔离数据库URL, autocommit=True) as 连: + 连.execute(sql.SQL("DROP DATABASE {} WITH (FORCE)").format(sql.Identifier(库名))) + + +def 构造服务(工厂: 数据库工厂, 执行=None, 重验=None): + 登记 = 流程登记() + 登记.登记处理器( + 步骤处理器( + "collect", + "1", + 执行 or (lambda 上下文: 步骤结果({"value": 上下文.任务.冻结输入["输入"]["value"]})), + "v1", + "v1", + 保护职责="输入校验", + ) + ) + 登记.登记处理器(步骤处理器("finish", "1", lambda 上下文: 步骤结果({"done": True}), "v1", "v1")) + 登记.登记类型("example", 必需保护=("输入校验",)) + + def 校验合成授权(快照): + if 快照.冻结输入["冻结上下文"]["authorization"] != "grant-1": + raise 状态冲突("合成来源授权已改变") + + 登记.登记恢复检查("example", "1", 重验 or 校验合成授权) + 运行 = 任务服务(工厂, 登记) + 编排 = 流程服务(运行, 登记) + 编排.发布( + 流程定义( + "example", + "1", + (步骤计划("collect", "collect", "1"), 步骤计划("finish", "finish", "1", ("collect",))), + ) + ) + return 运行, 编排 + + +def 请求(命令: str = "cmd", 作品: str = "work-A", 用途标记: 用途 = 用途.生产): + return 任务请求( + "example", + 命令, + "author", + 用途标记, + 内容用途.抽取, + {"value": 7}, + "policy-1", + "release-1", + { + "source_scope": {"work_id": 作品, "as_of": 3}, + "schema_versions": {"entity": "1"}, + "authorization": "grant-1", + "budget": {"remaining": 10}, + "stop_conditions": ["cancelled"], + }, + 作品ID=作品, + 排他范围=作用域("work", 作品, "extract"), + ) + + +def 完成所有(服务: 任务服务) -> None: + while 领取 := 服务.领取步骤("worker", ["collect", "finish"]): + 服务.执行一步(领取) + + +def test_公开编排重启后复用已完成步骤__a51001(任务库, tmp_path: Path) -> None: + 调用次数 = [] + + def 提取(上下文): + 调用次数.append(上下文.领取.步骤ID) + return 步骤结果({"saved": 1}) + + 运行, 编排 = 构造服务(任务库[用途.生产], 提取) + 身份 = 编排.发起(请求(), "example", "1") + assert 编排.发起(请求(), "example", "1") == 身份 + with pytest.raises(状态冲突): + 编排.发起(dataclasses.replace(请求(), 输入={"value": 8}), "example", "1") + 首步 = 运行.领取步骤("old", ["collect", "finish"]) + assert 首步 is not None + 运行.执行一步(首步) + 子进程代码 = """ +import sys +from muse.任务运行.接口 import 任务服务, 步骤处理器, 步骤结果 +from muse.编排.接口 import 流程登记 +from muse.基础设施.数据库.连接 import 数据库工厂 +from muse.配置 import 数据库引用 +from muse.共享.调用身份 import 用途 +登记 = 流程登记() +def 已完成不能再执行(_): + raise AssertionError("已完成步骤被重复执行") +登记.登记处理器(步骤处理器("collect", "1", 已完成不能再执行, "v1", "v1", 保护职责="输入校验")) +登记.登记处理器(步骤处理器("finish", "1", lambda _: 步骤结果({"done": True}), "v1", "v1")) +登记.登记类型("example", 必需保护=("输入校验",)) +运行 = 任务服务(数据库工厂(数据库引用("环境变量", "MUSE_W05_生产_URL"), 用途.生产), 登记) +while 领取 := 运行.领取步骤("new-process", ["collect", "finish"]): + 运行.执行一步(领取) +assert 运行.读取任务(sys.argv[1]).状态.value == "completed" +""" + 子进程 = subprocess.run( + [sys.executable, "-c", 子进程代码, 身份], cwd=tmp_path, capture_output=True, text=True + ) + assert 子进程.returncode == 0, 子进程.stderr + assert 运行.读取任务(身份).状态 is 任务状态.已完成 + assert 调用次数 == ["collect"] + assert all(步骤["result"] is not None for 步骤 in 运行.读取任务(身份).步骤) + + +def test_独立连接竞争同范围而不同范围可并行__a51002(任务库) -> None: + 运行, 编排 = 构造服务(任务库[用途.生产]) + A = 编排.发起(请求("a"), "example", "1") + B = 编排.发起(请求("b"), "example", "1") + C = 编排.发起(请求("c", "work-C"), "example", "1") + 起点 = Barrier(2) + + def 领取(名称): + 起点.wait() + return 运行.领取步骤(名称, ["collect"]) + + with ThreadPoolExecutor(max_workers=2) as 线程: + 结果 = list(线程.map(领取, ["one", "two"])) + 有效 = [项 for 项 in 结果 if 项] + # 同范围争用者可能先让出;下一次领取应直接命中独立范围。 + if len(有效) == 1: + 下一 = 运行.领取步骤("third", ["collect"]) + assert 下一 is not None + 有效.append(下一) + assert len(有效) == 2 + assert C in {项.任务ID for 项 in 有效} + assert len({项.任务ID for 项 in 有效} & {A, B}) == 1 + assert 运行.领取步骤("blocked", ["collect"]) is None + for 项 in 有效: + 运行.完成步骤(项, 步骤结果({"done": 1})) + + +def test_过期接管拒绝旧步骤和作用域代次__a51003(任务库) -> None: + 运行, 编排 = 构造服务(任务库[用途.生产]) + 身份 = 编排.发起(请求(), "example", "1") + 旧 = 运行.领取步骤("old", ["collect"], 租期秒=0.12) + assert 旧 is not None + 运行.保存检查点(旧, {"offset": 4}) + time.sleep(0.16) + 新 = 运行.领取步骤("new", ["collect"]) + assert 新 is not None and 新.任务ID == 身份 + assert 新.尝试ID != 旧.尝试ID and 新.作用域代次 > 旧.作用域代次 + assert 新.检查点 == {"offset": 4} + for 操作 in [lambda: 运行.完成步骤(旧, 步骤结果({"late": 1})), lambda: 运行.续租(旧)]: + with pytest.raises(租约失效): + 操作() + with pytest.raises(租约失效): + 运行.执行一步(dataclasses.replace(新, 处理器ID="finish")) + 运行.完成步骤(新, 步骤结果({"current": 1})) + assert 运行.读取任务(身份).步骤[0]["result"] == {"current": 1} + + +def test_取消后拒绝旧结果并释放作品范围__a51004(任务库) -> None: + 运行, 编排 = 构造服务(任务库[用途.生产]) + 身份 = 编排.发起(请求("old"), "example", "1") + 旧 = 运行.领取步骤("old", ["collect"]) + assert 旧 is not None + 运行.控制任务(身份, "author", 任务状态.运行中, "取消", 命令ID="cancel") + with pytest.raises(租约失效): + 运行.完成步骤(旧, 步骤结果({"late": 1})) + 新任务 = 编排.发起(请求("new"), "example", "1") + 新 = 运行.领取步骤("new", ["collect"]) + assert 新 is not None and 新.任务ID == 新任务 + assert 运行.读取任务(身份).状态 is 任务状态.已取消 + + +def test_暂停恢复重验且检查点持久保留__a51005(任务库) -> None: + 运行, 编排 = 构造服务(任务库[用途.生产]) + 身份 = 编排.发起(请求(), "example", "1") + 旧 = 运行.领取步骤("old", ["collect"]) + assert 旧 is not None + 运行.保存检查点(旧, {"part": 2}) + 运行.控制任务(身份, "author", 任务状态.运行中, "暂停", 命令ID="pause") + + def 拒绝过期来源(快照): + raise 状态冲突("来源已变化") + + 重启, _ = 构造服务(任务库[用途.生产], 重验=拒绝过期来源) + with pytest.raises(状态冲突, match="来源"): + 重启.控制任务(身份, "author", 任务状态.已暂停, "恢复", 命令ID="resume") + assert 重启.读取任务(身份).状态 is 任务状态.已暂停 + + def 校验当前来源(快照): + assert 快照.冻结输入["冻结上下文"]["authorization"] == "grant-1" + + 重启, _ = 构造服务(任务库[用途.生产], 重验=校验当前来源) + 重启.控制任务(身份, "author", 任务状态.已暂停, "恢复", 命令ID="resume") + 新 = 重启.领取步骤("new", ["collect"]) + assert 新 is not None and 新.检查点 == {"part": 2} + + +def test_具名接入恢复由登记owner重验__a5100f(任务库) -> None: + """HTTP/CLI 共用入口不能自报重验;服务按冻结流程找到当前业务检查。""" + from muse.接入.主会话 import 任务控制请求, 执行任务控制 + + 当前来源 = {"版本": 4} + + def 重验合成任务(快照): + 上下文 = 快照.冻结输入["冻结上下文"] + if 上下文["source_scope"]["as_of"] != 当前来源["版本"]: + raise 状态冲突("来源版本已变化") + assert 上下文["schema_versions"] == {"entity": "1"} + assert 上下文["authorization"] == "grant-1" + assert 上下文["budget"]["remaining"] > 0 + + 运行, 编排 = 构造服务(任务库[用途.生产], 重验=重验合成任务) + 身份 = 编排.发起(请求(), "example", "1") + 运行.控制任务(身份, "author", 任务状态.待运行, "暂停", 命令ID="pause") + 控制 = 任务控制请求( + command_id="resume", target_ref=身份, expected_state=任务状态.已暂停, action="恢复" + ) + with pytest.raises(状态冲突, match="来源版本"): + 执行任务控制(运行, "author", 控制) + assert 运行.读取任务(身份).状态 is 任务状态.已暂停 + 当前来源["版本"] = 3 + 回执 = 执行任务控制(运行, "author", 控制) + assert 回执.state is 任务状态.待运行 + assert 执行任务控制(运行, "author", 控制) == 回执 + + +def test_未知调用先对账并复用可靠输出__a51006(任务库) -> None: + 运行, 编排 = 构造服务( + 任务库[用途.生产], + 执行=lambda 上下文: 步骤结果(上下文.领取.检查点["已保存调用"]["输出"]), + ) + 身份 = 编排.发起(请求(), "example", "1") + 旧 = 运行.领取步骤("old", ["collect"], 租期秒=0.12) + assert 旧 is not None + 运行.登记外部调用(旧, "call-1") + time.sleep(0.16) + assert 运行.领取步骤("new", ["collect", "finish"]) is None + assert 运行.读取任务(身份).状态 is 任务状态.待对账 + with pytest.raises(状态冲突, match="对账"): + 运行.控制任务(身份, "author", 任务状态.待对账, "恢复", 命令ID="resume") + 运行.对账调用(身份, 旧.尝试ID, "author", 已保存输出={"recovered": 1}, 对账回执="receipt-1") + assert 运行.读取任务(身份).步骤[0]["state"] == "pending" + 运行.控制任务(身份, "author", 任务状态.待对账, "恢复", 命令ID="resume") + 下一 = 运行.领取步骤("new", ["collect", "finish"]) + assert 下一 is not None and 下一.步骤ID == "collect" + 运行.执行一步(下一) + 收尾 = 运行.领取步骤("new", ["finish"]) + assert 收尾 is not None + 运行.执行一步(收尾) + assert 运行.读取任务(身份).状态 is 任务状态.已完成 + assert 运行.读取任务(身份).步骤[0]["result"] == {"recovered": 1} + + +def test_事件重复幂等且游标缺口要求快照__a51007(任务库) -> None: + 运行, 编排 = 构造服务(任务库[用途.生产]) + 身份 = 编排.发起(请求(), "example", "1") + 领取 = 运行.领取步骤("one", ["collect"]) + assert 领取 is not None + 事件ID = str(uuid.uuid4()) + 事件 = 运行.追加事件(领取, 事件类型.文本片段, {"文本": "片段"}, 事件ID=事件ID) + assert 运行.追加事件(领取, 事件类型.文本片段, {"文本": "片段"}, 事件ID=事件ID) == 事件 + with pytest.raises(状态冲突): + 运行.追加事件(领取, 事件类型.文本片段, {"文本": "不同"}, 事件ID=事件ID) + 续接 = 运行.续接事件(身份) + assert [e.序号 for e in 续接.事件] == list(range(1, 续接.最后序号 + 1)) + assert not 续接.需要快照 + with 任务库[用途.维护].连接() as 连: + 连.execute("DELETE FROM public.muse_task_event WHERE task_id=%s AND sequence=2", (身份,)) + assert 运行.续接事件(身份).需要快照 + assert 运行.续接事件(身份, 事件.序号).事件 == () + + +def test_失败如实保留且停止领取不丢状态__a51008(任务库) -> None: + def 失败处理器(_): + raise RuntimeError("合成失败") + + 运行, 编排 = 构造服务(任务库[用途.生产], 失败处理器) + 身份 = 编排.发起(请求(), "example", "1") + 领取 = 运行.领取步骤("one", ["collect"]) + assert 领取 is not None + with pytest.raises(RuntimeError): + 运行.执行一步(领取) + assert 运行.读取任务(身份).状态 is 任务状态.已失败 + 编排.发起(请求("next", "another"), "example", "1") + 运行.停止领取() + assert 运行.领取步骤("stop", ["collect"]) is None + 重启, _ = 构造服务(任务库[用途.生产]) + assert 重启.领取步骤("restart", ["collect"]) is not None + + +def test_生产评测任务使用独立权限与实例__a51009(任务库) -> None: + 生产, 生产编排 = 构造服务(任务库[用途.生产]) + 评测, 评测编排 = 构造服务(任务库[用途.评测]) + A = 生产编排.发起(请求(), "example", "1") + B = 评测编排.发起(请求(用途标记=用途.评测), "example", "1") + assert A != B + assert 生产.领取步骤("prod", ["collect"]) is not None + assert 评测.领取步骤("eval", ["collect"]) is not None + with pytest.raises(任务错误, match="不存在|用途"): + 评测.读取任务(A) + + +def test_处理器调用期间不持有数据库事务__a5100a(任务库, 隔离数据库URL: str) -> None: + def 处理器(上下文): + with 任务库[用途.维护].连接() as 连: + 目标库 = 连.info.dbname + with psycopg.connect(隔离数据库URL) as 管理员: + 数量 = 管理员.execute( + "SELECT count(*) FROM pg_stat_activity WHERE datname=%s " + "AND pid<>pg_backend_pid() AND state='idle in transaction'", + (目标库,), + ).fetchone()[0] + assert 数量 == 0 + return 步骤结果({"done": True}) + + 运行, 编排 = 构造服务(任务库[用途.生产], 处理器) + 编排.发起(请求(), "example", "1") + 领取 = 运行.领取步骤("one", ["collect"]) + assert 领取 is not None + assert 运行.执行一步(领取).输出 == {"done": True} + + +def test_控制命令重放不重复事件且异参数拒绝__a5100b(任务库) -> None: + 运行, 编排 = 构造服务(任务库[用途.生产]) + 身份 = 编排.发起(请求(), "example", "1") + 首次 = 运行.控制任务(身份, "author", 任务状态.待运行, "取消", 命令ID="cancel-one") + 重放 = 运行.控制任务(身份, "author", 任务状态.待运行, "取消", 命令ID="cancel-one") + assert 重放 == 首次 + assert 运行.读取任务(身份).最后序号 == 首次["last_sequence"] + with pytest.raises(任务错误): + 运行.控制任务(身份, "author", 任务状态.待运行, "暂停", 命令ID="cancel-one") + + +def test_任务列表按作者作品及用途隔离__a5100c(任务库) -> None: + 运行, 编排 = 构造服务(任务库[用途.生产]) + A = 编排.发起(请求("a", "work-A"), "example", "1") + 编排.发起(请求("b", "work-B"), "example", "1") + 编排.发起(dataclasses.replace(请求("c", "work-A"), 作者="other"), "example", "1") + assert [x.任务ID for x in 运行.列出任务("author", 作品ID="work-A")] == [A] + assert 运行.列出任务("missing") == [] + + +def test_任务步骤按冻结流程顺序显示__a5100d(任务库) -> None: + 运行, 编排 = 构造服务(任务库[用途.生产]) + 编排.发布( + 流程定义( + "example", + "2", + ( + 步骤计划("z-first", "collect", "1"), + 步骤计划("a-last", "finish", "1", ("z-first",)), + ), + ) + ) + 身份 = 编排.发起(请求("ordered"), "example", "2") + assert [step["step_id"] for step in 运行.读取任务(身份).步骤] == ["z-first", "a-last"] + + +def test_执行器续租并在停止后不领取新步骤__a5100e(任务库) -> None: + from muse.接入.执行器 import 执行器 + + def 长步骤(_): + time.sleep(0.45) + return 步骤结果({"done": True}) + + 运行, 编排 = 构造服务(任务库[用途.生产], 长步骤) + 身份 = 编排.发起(请求(), "example", "1") + worker = 执行器(运行, "durable-worker", ["collect", "finish"], 租期秒=0.3) + assert worker.运行一次() + worker.停止() + assert not worker.运行一次() + assert 运行.读取任务(身份).步骤[0]["state"] == "completed" + assert 运行.读取任务(身份).步骤[1]["state"] == "pending" diff --git a/数据库/迁移/V0004__任务运行与证据.sql b/数据库/迁移/V0004__任务运行与证据.sql new file mode 100644 index 0000000..03e29f7 --- /dev/null +++ b/数据库/迁移/V0004__任务运行与证据.sql @@ -0,0 +1,290 @@ +-- W05:持久任务、步骤、尝试、独立作用域租约和有序事件。 +-- W06 模型预算与 W08 原文生命周期随所属实现接入;仅用于尚未冻结的开发基线。 +DO $migration$ +DECLARE + namespace text; +BEGIN + FOREACH namespace IN ARRAY ARRAY['public', 'evaluation'] LOOP + EXECUTE format($ddl$ + CREATE TABLE %1$I.muse_task ( + task_id uuid PRIMARY KEY, + run_purpose public.muse_purpose NOT NULL, + command_id text NOT NULL, + author_id text NOT NULL, + request_hash text NOT NULL, + frozen_input jsonb NOT NULL, + flow_id text NOT NULL, + flow_version text NOT NULL, + flow_snapshot jsonb NOT NULL, + state text NOT NULL CHECK (state IN ('queued','running','paused','reconciling','failed','cancelled','completed')), + scope_key text, + scope_generation bigint, + last_sequence bigint NOT NULL DEFAULT 0, + created_at timestamptz NOT NULL DEFAULT clock_timestamp(), + updated_at timestamptz NOT NULL DEFAULT clock_timestamp(), + UNIQUE (run_purpose, author_id, command_id) + ); + CREATE TABLE %1$I.muse_task_control ( + run_purpose public.muse_purpose NOT NULL, + author_id text NOT NULL, + command_id text NOT NULL, + task_id uuid NOT NULL REFERENCES %1$I.muse_task(task_id), + request_hash text NOT NULL, + receipt jsonb, + PRIMARY KEY (run_purpose, author_id, command_id) + ); + CREATE TABLE %1$I.muse_step ( + task_id uuid NOT NULL REFERENCES %1$I.muse_task(task_id), + step_id text NOT NULL, + processor_id text NOT NULL, + processor_version text NOT NULL, + dependencies text[] NOT NULL, + state text NOT NULL CHECK (state IN ('pending','running','completed','failed','unknown')), + current_attempt uuid, + checkpoint jsonb NOT NULL DEFAULT '{}', + result jsonb, + PRIMARY KEY (task_id, step_id) + ); + CREATE TABLE %1$I.muse_attempt ( + attempt_id uuid PRIMARY KEY, + task_id uuid NOT NULL, + step_id text NOT NULL, + worker_id text NOT NULL, + lease_token uuid NOT NULL, + lease_until timestamptz NOT NULL, + scope_generation bigint, + heartbeat_at timestamptz NOT NULL DEFAULT clock_timestamp(), + state text NOT NULL CHECK (state IN ('running','expired','paused','cancelled','completed','failed','unknown')), + call_state text NOT NULL DEFAULT 'not_sent' CHECK (call_state IN ('not_sent','sent','unknown','saved')), + call_reference text, + result jsonb, + failure_code text, + created_at timestamptz NOT NULL DEFAULT clock_timestamp(), + FOREIGN KEY (task_id, step_id) REFERENCES %1$I.muse_step(task_id, step_id) + ); + CREATE TABLE %1$I.muse_task_scope_lease ( + run_purpose public.muse_purpose NOT NULL, + scope_key text NOT NULL, + holder_task_id uuid, + generation bigint NOT NULL CHECK (generation > 0), + lease_until timestamptz NOT NULL, + PRIMARY KEY (run_purpose, scope_key) + ); + CREATE TABLE %1$I.muse_task_event ( + event_id uuid PRIMARY KEY, + task_id uuid NOT NULL REFERENCES %1$I.muse_task(task_id), + step_id text, + attempt_id uuid, + sequence bigint NOT NULL, + event_type text NOT NULL, + occurred_at timestamptz NOT NULL DEFAULT clock_timestamp(), + payload_version integer NOT NULL, + payload jsonb NOT NULL, + payload_hash text NOT NULL, + UNIQUE(task_id, sequence) + ); + CREATE INDEX ON %1$I.muse_task (run_purpose, state, created_at); + CREATE INDEX ON %1$I.muse_step (state, task_id); + CREATE INDEX ON %1$I.muse_attempt (task_id, state, lease_until); + $ddl$, namespace); + END LOOP; +END +$migration$; + +-- W06 预算与配置:额度账户串行化预留,金额事实留在逐调用账本。 +-- W08 原文:授权和租约先持久化,完整归档事务同时写入清单与字节。 +DO $raw_lifecycle$ +DECLARE namespace text; +BEGIN + FOREACH namespace IN ARRAY ARRAY['public', 'evaluation'] LOOP + EXECUTE format($ddl$ + CREATE TABLE %1$I.muse_raw_namespace ( + singleton boolean PRIMARY KEY DEFAULT true CHECK (singleton), + namespace_id uuid NOT NULL DEFAULT gen_random_uuid() + ); + INSERT INTO %1$I.muse_raw_namespace (singleton) VALUES (true); + CREATE TABLE %1$I.muse_raw_orphan_cleanup ( + run_purpose public.muse_purpose NOT NULL, + object_id uuid NOT NULL, + receipt_id uuid NOT NULL DEFAULT gen_random_uuid(), + state text NOT NULL CHECK (state IN ('pending','completed')), + removed boolean, + completed_at timestamptz, + PRIMARY KEY (run_purpose,object_id) + ); + CREATE TABLE %1$I.muse_raw_authorization ( + authorization_id uuid PRIMARY KEY, + task_id uuid NOT NULL REFERENCES %1$I.muse_task(task_id), + command_id text NOT NULL, + approved_by text NOT NULL, + source_version text NOT NULL, + content_hashes text[] NOT NULL CHECK (cardinality(content_hashes)>0), + purpose text NOT NULL, + retention_mode text NOT NULL CHECK (retention_mode IN ('temporary','archive','persistent')), + approved_at timestamptz NOT NULL, + valid_until timestamptz NOT NULL, + request_hash text NOT NULL, + revoked boolean NOT NULL DEFAULT false, + call_id text, + call_request_hash text, + response_hash text, + parent_authorization_id uuid REFERENCES %1$I.muse_raw_authorization(authorization_id), + attempt_id uuid REFERENCES %1$I.muse_attempt(attempt_id), + derivation_kind text CHECK (derivation_kind IN ('model','tool')), + evidence_refs uuid[] NOT NULL DEFAULT '{}', + CHECK ((call_id IS NULL) = (call_request_hash IS NULL)), + CHECK (response_hash IS NULL OR call_id IS NOT NULL), + UNIQUE(task_id,command_id) + ); + CREATE TABLE %1$I.muse_role_session_authorization ( + authorization_id uuid PRIMARY KEY REFERENCES %1$I.muse_raw_authorization(authorization_id), + step_id text NOT NULL, + stage text NOT NULL, + initial_request_hash text NOT NULL, + frozen_input_hash text NOT NULL, + max_model_calls integer NOT NULL CHECK (max_model_calls>0), + max_tool_calls integer NOT NULL CHECK (max_tool_calls>=0), + policy_hash text NOT NULL + ); + CREATE TABLE %1$I.muse_raw_lease ( + lease_id uuid PRIMARY KEY, + task_id uuid NOT NULL REFERENCES %1$I.muse_task(task_id), + authorization_id uuid NOT NULL REFERENCES %1$I.muse_raw_authorization(authorization_id), + command_id text NOT NULL, + created_at timestamptz NOT NULL, + retain_until timestamptz NOT NULL, + state text NOT NULL CHECK (state IN ('open','closed','migrating','migrated')), + archive_id uuid, + archive_authorization_id uuid REFERENCES %1$I.muse_raw_authorization(authorization_id), + cleanup_receipt uuid, + failure_code text, + UNIQUE(task_id,command_id), + CHECK (retain_until>created_at AND retain_until<=created_at+interval '24 hours') + ); + CREATE TABLE %1$I.muse_raw_archive ( + archive_id uuid PRIMARY KEY, + lease_id uuid UNIQUE REFERENCES %1$I.muse_raw_lease(lease_id), + legacy_ref text UNIQUE, + legacy_manifest jsonb, + CHECK ((lease_id IS NULL) <> (legacy_ref IS NULL)), + authorization_id uuid NOT NULL REFERENCES %1$I.muse_raw_authorization(authorization_id), + tree_hash text NOT NULL, + entry_count integer NOT NULL CHECK (entry_count>0), + total_bytes bigint NOT NULL CHECK (total_bytes>=0), + receipt_id uuid NOT NULL, + created_at timestamptz NOT NULL + ); + CREATE TABLE %1$I.muse_raw_archive_item ( + archive_id uuid NOT NULL REFERENCES %1$I.muse_raw_archive(archive_id), + content_hash text NOT NULL, + content bytea NOT NULL, + PRIMARY KEY(archive_id,content_hash) + ); + CREATE TABLE %1$I.muse_runtime_evidence ( + evidence_id uuid PRIMARY KEY, + task_id uuid NOT NULL REFERENCES %1$I.muse_task(task_id), + attempt_id uuid NOT NULL REFERENCES %1$I.muse_attempt(attempt_id), + kind text NOT NULL CHECK (kind IN ('model_input','model_response','tool_result','failure')), + reference_id text NOT NULL, + content_hash text NOT NULL, + content bytea, + authorization_id uuid REFERENCES %1$I.muse_raw_authorization(authorization_id), + outcome text NOT NULL CHECK (outcome IN ('completed','failed','partial')), + metadata jsonb NOT NULL, + revision integer NOT NULL DEFAULT 1, + created_at timestamptz NOT NULL DEFAULT clock_timestamp(), + UNIQUE(attempt_id,kind,reference_id), + CHECK (content IS NULL OR authorization_id IS NOT NULL) + ); + $ddl$, namespace); + EXECUTE format('REVOKE INSERT,UPDATE,DELETE ON %I.muse_raw_namespace FROM muse_app,muse_eval', namespace); + END LOOP; +END +$raw_lifecycle$; + +DO $budget_config$ +DECLARE + namespace text; +BEGIN + FOREACH namespace IN ARRAY ARRAY['public', 'evaluation'] LOOP + EXECUTE format($ddl$ + CREATE TABLE %1$I.muse_budget_account ( + run_purpose public.muse_purpose NOT NULL, + account_id text NOT NULL, + policy_version text NOT NULL, + policy jsonb NOT NULL, + policy_hash text NOT NULL, + PRIMARY KEY (run_purpose, account_id) + ); + CREATE TABLE %1$I.muse_task_budget ( + task_id uuid PRIMARY KEY REFERENCES %1$I.muse_task(task_id), + account_id text NOT NULL, + plan jsonb NOT NULL, + plan_hash text NOT NULL, + stopped boolean NOT NULL DEFAULT false + ); + CREATE TABLE %1$I.muse_budget_reservation ( + run_purpose public.muse_purpose NOT NULL, + call_id text NOT NULL, + account_id text NOT NULL, + task_id uuid NOT NULL REFERENCES %1$I.muse_task(task_id), + attempt_id uuid NOT NULL REFERENCES %1$I.muse_attempt(attempt_id), + role_id text NOT NULL, + request_hash text NOT NULL, + state text NOT NULL CHECK (state IN ('reserved','in_flight','unknown','settled','released')), + reserved_amount numeric(24,6) NOT NULL CHECK (reserved_amount >= 0), + actual_amount numeric(24,6) CHECK (actual_amount >= 0), + window_start timestamptz NOT NULL, + window_end timestamptz NOT NULL, + expires_at timestamptz NOT NULL, + sent_at timestamptz, + receipt_id text, + reason text, + over_budget boolean NOT NULL DEFAULT false, + PRIMARY KEY (run_purpose, call_id), + CHECK ((state = 'settled') = (actual_amount IS NOT NULL)) + ); + CREATE INDEX ON %1$I.muse_budget_reservation (run_purpose,account_id,window_start); + CREATE INDEX ON %1$I.muse_budget_reservation (task_id,role_id); + CREATE TABLE %1$I.muse_runtime_config_version ( + run_purpose public.muse_purpose NOT NULL, + config_id text NOT NULL, + version text NOT NULL, + content jsonb NOT NULL, + content_hash text NOT NULL, + created_at timestamptz NOT NULL DEFAULT clock_timestamp(), + PRIMARY KEY (run_purpose,config_id,version) + ); + CREATE TABLE %1$I.muse_runtime_config_validation ( + receipt_id uuid PRIMARY KEY, + run_purpose public.muse_purpose NOT NULL, + config_id text NOT NULL, + version text NOT NULL, + content_hash text NOT NULL, + validator_id text NOT NULL, + evidence_refs jsonb NOT NULL, + validated_at timestamptz NOT NULL DEFAULT clock_timestamp(), + FOREIGN KEY (run_purpose,config_id,version) + REFERENCES %1$I.muse_runtime_config_version(run_purpose,config_id,version) + ); + CREATE TABLE %1$I.muse_runtime_config_active ( + run_purpose public.muse_purpose NOT NULL, + config_id text NOT NULL, + version text, + validation_receipt uuid REFERENCES %1$I.muse_runtime_config_validation(receipt_id), + approval_ref text, + generation bigint NOT NULL DEFAULT 0, + PRIMARY KEY (run_purpose,config_id) + ); + CREATE TABLE %1$I.muse_task_config_binding ( + task_id uuid PRIMARY KEY REFERENCES %1$I.muse_task(task_id), + config_id text NOT NULL, + version text NOT NULL, + content_hash text NOT NULL, + content jsonb NOT NULL, + validation_receipt uuid NOT NULL REFERENCES %1$I.muse_runtime_config_validation(receipt_id) + ); + $ddl$, namespace); + END LOOP; +END +$budget_config$; diff --git a/数据库/迁移/V0016__流程定义与版本.sql b/数据库/迁移/V0016__流程定义与版本.sql new file mode 100644 index 0000000..8af2e91 --- /dev/null +++ b/数据库/迁移/V0016__流程定义与版本.sql @@ -0,0 +1,20 @@ +-- W05:已校验流程的不可变执行计划;同版本不能覆盖历史任务定义。 +DO $migration$ +DECLARE + namespace text; +BEGIN + FOREACH namespace IN ARRAY ARRAY['public', 'evaluation'] LOOP + EXECUTE format($ddl$ + CREATE TABLE %I.muse_flow_version ( + run_purpose public.muse_purpose NOT NULL, + flow_id text NOT NULL, + version text NOT NULL, + definition jsonb NOT NULL, + definition_hash text NOT NULL, + created_at timestamptz NOT NULL DEFAULT clock_timestamp(), + PRIMARY KEY (run_purpose, flow_id, version) + ) + $ddl$, namespace); + END LOOP; +END +$migration$;