W05 任务状态、租约、事件与恢复:任务状态机、短事务租约、事件续接、恢复控制与作用域互斥。

按 R2 串行阶段整理提交;包内文件为该阶段交付(含后续小增量),状态以工作包清单为准。
This commit is contained in:
zizi 2026-09-10 19:25:40 +08:00
parent b0623d9048
commit e99a952dd0
25 changed files with 2566 additions and 0 deletions

View File

@ -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)。

View File

@ -0,0 +1,3 @@
| 名称 | 相对地址 | 内容描述 | 使用场景 | 使用要求 |
|------|----------|----------|----------|----------|
| 查看与恢复任务 | [SKILL.md](SKILL.md) | 查询持久任务与事件,按作者请求暂停、取消或经登记检查恢复任务。 | 需要该项已实现操作时 | 遵循配置用途与用户动作授权 |

View File

@ -0,0 +1 @@
"""持久任务与可靠执行;公开用例见接口模块。"""

View File

@ -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 任务错误("文本片段事件必须提供文本")

View File

@ -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

View File

@ -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"])

View File

@ -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"],
)

View File

@ -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)

View File

@ -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__ = [
"提供方配置",
"调用计价",
"配置快照",
"角色会话",
"组装角色请求",
"回合记录",
"工具回传",
"校验原文字节",
"模型执行器",
"模型交付",
"请求字节",
"证据服务",
"原文服务",
"原文错误",
"校验期限",
"只读工具集",
"工具范围",
"工具来源",
"工具定义",
"工具结果",
"角色策略目录",
"角色执行策略",
"任务预算计划",
"角色预算",
"预算管理",
"额度策略",
"凭据引用",
"运行配置内容",
"配置版本管理",
"配置验证证据",
"校验模型输出",
"校验输出合同",
"工具请求",
"模型协议错误",
"模型用量",
"模型结果",
"模型请求",
"模型宿主",
"已准备模型调用",
"任务服务",
"任务请求",
"任务快照",
"任务状态",
"任务错误",
"状态冲突",
"租约失效",
"作用域",
"步骤计划",
"执行计划",
"步骤处理器",
"步骤结果",
"执行上下文",
"领取凭证",
"事件类型",
"运行事件",
"事件续接",
]

View File

@ -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

View File

@ -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

View File

@ -0,0 +1,4 @@
"""基础设施:数据库、凭据、宿主与受控文件等环境实现。
边界:只提供环境能力,不放业务规则;业务模块经接口使用。
"""

View File

@ -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())

View File

@ -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(任务)

View File

@ -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)

View File

@ -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(轮询秒)

View File

@ -0,0 +1 @@
"""代码登记的流程与保护约束;公开入口见接口模块。"""

24
src/muse/编排/接口.py Normal file
View File

@ -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__ = ["流程服务", "流程定义", "流程登记", "流程校验错误"]

View File

@ -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 流程校验错误("存在绕过必需保护的流程终点")

View File

@ -0,0 +1,12 @@
"""流程定义是不可变输入,发布后任务保存固定执行计划。"""
from dataclasses import dataclass
from muse.任务运行.接口 import 步骤计划
@dataclass(frozen=True, slots=True)
class 流程定义:
流程ID: str
版本: str
步骤: tuple[步骤计划, ...]

View File

@ -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])

View File

@ -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",))),
)
)

View File

@ -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"

View File

@ -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$;

View File

@ -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$;