"""真实 ASGI、CLI 与隔离 PostgreSQL 共用任务合同;不使用仓储替身。""" import json import subprocess import sys from pathlib import Path import pytest from fastapi.testclient import TestClient from muse.任务运行.接口 import 任务服务, 任务请求, 步骤处理器, 步骤结果, 步骤计划 from muse.共享.调用身份 import 内容用途, 用途 from muse.接入.http.应用 import 创建应用 from muse.编排.接口 import 流程定义, 流程服务, 流程登记 from muse.配置 import 应用配置, 服务配置 pytestmark = pytest.mark.数据库 def test_具名CLI按字段传参不依赖JSON排列__793c2a(任务接入环境, tmp_path: Path) -> None: 服务, 身份, _, 配置文件 = 任务接入环境 请求文件 = tmp_path / "作者决定.json" 请求文件.write_text( json.dumps( { "action": "取消", "expected_state": "queued", "target_ref": 身份, "command_id": "作者决定-1", }, ensure_ascii=False, ) ) 进程 = subprocess.run( [sys.executable, "-m", "muse", "任务", str(配置文件), "控制", str(请求文件)], cwd="/", capture_output=True, text=True, timeout=20, ) assert 进程.returncode == 0, 进程.stderr assert json.loads(进程.stdout)["command_id"] == "作者决定-1" assert 服务.读取任务(身份).状态.value == "cancelled" @pytest.mark.parametrize( "请求文本,额外参数,退出码", [ ("{}", [], 1), ('["x", 1, null]', [], 1), ("{}", ["--stdin"], 2), ("not json", [], 1), ('{"a": 1}', [], 1), ('{"sql": "CREATE TABLE t (id int)"}', [], 1), ('{"sql": "DELETE FROM t"}', [], 1), ('{"sql": "ALTER TABLE t ADD COLUMN x int", "params": []}', [], 1), ('{"sql": "DELETE FROM t", "params": []}', [], 1), ('{"sql": "UPDATE t SET x=1 WHERE id=1"}', [], 1), ('{"sql": "DELETE FROM t WHERE id=%s", "params": [1]}', [], 1), ], ids=[ "empty", "typed-array", "competing-input", "invalid-json", "unknown-object", "create-sql", "bare-delete", "alter-params", "delete-params", "update-where", "delete-where", ], ) def test_具名CLI拒绝未登记输入且任务不变__a71002( 任务接入环境, tmp_path: Path, 请求文本: str, 额外参数: list[str], 退出码: int, ) -> None: """旧通用 SQL 入口已退出;只接受具名命令合同,不能以 WHERE 或参数列表换取写能力。""" 服务, 身份, _, 配置文件 = 任务接入环境 请求文件 = tmp_path / "请求.json" 请求文件.write_text(请求文本) 进程 = subprocess.run( [sys.executable, "-m", "muse", "任务", str(配置文件), "控制", str(请求文件), *额外参数], cwd="/", capture_output=True, text=True, timeout=20, ) assert 进程.returncode == 退出码, 进程.stdout if 退出码 == 1: assert json.loads(进程.stderr)["code"] == "INVALID_REQUEST" assert 服务.读取任务(身份).状态.value == "queued" @pytest.fixture def 任务接入环境(应用测试库, tmp_path: Path): registry = 流程登记() registry.登记处理器(步骤处理器("synthetic", "1", lambda _: 步骤结果({"ok": True}), "v1", "v1")) registry.登记类型("synthetic", 必需保护=()) service = 任务服务(应用测试库[用途.生产], registry) flow = 流程服务(service, registry) flow.发布(流程定义("synthetic", "1", (步骤计划("合成步骤", "synthetic", "1"),))) req = 任务请求( "合成任务", "create", "author-local", 用途.生产, 内容用途.抽取, {}, "v1", "test", { "source_scope": {}, "schema_versions": {}, "authorization": "synthetic", "budget": {}, "stop_conditions": [], }, 作品ID="work-A", ) task_id = flow.发起(req, "synthetic", "1") token = tmp_path / "口令.txt" token.write_text("synthetic-password") config = 应用配置( 应用测试库[用途.生产].引用, "test", HTTP=服务配置(str(token), 公开地址="http://testserver", 允许来源=("http://testserver",)), ) file = tmp_path / "应用.toml" file.write_text(f"""["数据库"] "取值方式" = "受控存储" "位置" = {json.dumps(str(tmp_path / "production.txt"))} ["资源"] "发布身份" = "test" [HTTP] "口令文件" = {json.dumps(str(token))} """) return service, task_id, config, file def test_HTTP与CLI控制回执和拒绝保持一致__a71001(任务接入环境, tmp_path: Path) -> None: service, task_id, config, file = 任务接入环境 with TestClient(创建应用(config)) as client: assert client.get(f"/api/v1/tasks/{task_id}").status_code == 401 client.headers["Origin"] = "http://testserver" assert ( client.post("/api/v1/session", json={"password": "synthetic-password"}).status_code == 200 ) assert ( client.get("/api/v1/tasks", params={"work_id": "work-A"}).json()[0]["task_id"] == task_id ) assert ( client.get(f"/api/v1/tasks/{task_id}", params={"work_id": "work-B"}).status_code == 403 ) body = { "command_id": "cancel-once", "target_ref": task_id, "expected_state": "queued", "action": "取消", } response = client.post(f"/api/v1/tasks/{task_id}/controls", json=body) assert response.status_code == 200 request = tmp_path / "请求.json" request.write_text(json.dumps(body)) args = [sys.executable, "-m", "muse", "任务", str(file), "控制", str(request)] cli = subprocess.run(args, cwd="/", capture_output=True, text=True, timeout=20) assert cli.returncode == 0, cli.stderr assert json.loads(cli.stdout) == response.json() body["action"] = "暂停" request.write_text(json.dumps(body)) response = client.post(f"/api/v1/tasks/{task_id}/controls", json=body) cli = subprocess.run(args, cwd="/", capture_output=True, text=True, timeout=20) assert response.status_code == 409 and cli.returncode == 1 assert response.json() == json.loads(cli.stderr) events = client.get(f"/api/v1/tasks/{task_id}/events").json() assert [event["event_type"] for event in events["items"]] == [ "task.created", "task.cancelled", ] assert ( client.get(f"/api/v1/tasks/{task_id}/events", params={"cursor": 2}).json()["items"] == [] )