deerflow-code/offline-backend-20260512/backend/tests/test_workflow_runs.py
2026-09-07 18:24:55 +08:00

1350 lines
55 KiB
Python

"""Run lifecycle tests: HTTP APIs, dispatcher + executor, SSE, pause/resume, embed.
Everything runs inside one event loop (httpx ASGI transport rather than
``TestClient``) because the in-memory stores hold ``asyncio`` primitives that
must not be shared across loops.
"""
from __future__ import annotations
import asyncio
import json
from datetime import UTC, datetime, timedelta
from typing import Any
import httpx
import pytest
import pytest_asyncio
from fastapi import FastAPI
from app.gateway.routers import workflow_data_sources, workflow_embed, workflow_planning, workflow_resources
from app.gateway.routers import workflow_runs as runs_router
from app.gateway.workflow_agent_runner import WorkflowAgentRunner
from app.gateway.workflow_dispatcher import WorkflowRunDispatcher
from app.gateway.workflow_executor import WorkflowRunExecutor
from app.gateway.workflow_live_hub import WorkflowLiveHub
from deerflow.config.workflow_config import WorkflowConfig
from deerflow.persistence.workflow_data_sources import MemoryWorkflowDataSourceStore
from deerflow.persistence.workflow_events import MemoryWorkflowEventStore
from deerflow.persistence.workflow_planning import MemoryWorkflowPlanningStore
from deerflow.persistence.workflow_runs import MemoryWorkflowRunStore
from deerflow.persistence.workflows import MemoryWorkflowStore
REPORT_GRAPH: dict[str, Any] = {
"schemaVersion": "1.0",
"id": "wf-demo",
"name": "demo",
"inputSchema": {"required": ["topic"]},
"nodes": [
{"id": "start", "type": "start"},
{
"id": "shape",
"type": "transform",
"config": {"operations": [{"op": "set", "target": "title", "value": "{{ inputs.topic }} 报告"}]},
},
{"id": "end", "type": "output", "config": {"mapping": {"title": "{{ nodes.shape.data.title }}"}}},
],
"edges": [
{"id": "e1", "source": "start", "target": "shape"},
{"id": "e2", "source": "shape", "target": "end"},
],
}
FAIL_GRAPH: dict[str, Any] = {
"schemaVersion": "1.0",
"id": "wf-fail",
"name": "fail",
"inputSchema": {"required": ["topic"]},
"nodes": [
{"id": "start", "type": "start"},
{
"id": "shape",
"type": "transform",
"config": {"operations": [{"op": "set", "target": "title", "value": "{{ inputs.topic }}"}]},
},
{"id": "boom", "type": "transform", "config": {"operations": [{"op": "not_a_real_op"}]}},
{"id": "end", "type": "output", "config": {"mapping": {"title": "{{ nodes.boom.data.title }}"}}},
],
"edges": [
{"id": "e1", "source": "start", "target": "shape"},
{"id": "e2", "source": "shape", "target": "boom"},
{"id": "e3", "source": "boom", "target": "end"},
],
}
PAUSE_GRAPH: dict[str, Any] = {
"schemaVersion": "1.0",
"id": "wf-pause",
"name": "pause",
"nodes": [
{"id": "start", "type": "start"},
{
"id": "ask",
"type": "human_input",
"config": {"prompt": "确认", "formSchema": {"required": ["comment"]}, "actions": ["submit"]},
},
{"id": "end", "type": "output", "config": {"mapping": {"comment": "{{ nodes.ask.data.values.comment }}"}}},
],
"edges": [
{"id": "e1", "source": "start", "target": "ask"},
{"id": "e2", "source": "ask", "target": "end"},
],
}
PLANNING_GRAPH: dict[str, Any] = {
"schemaVersion": "1.0",
"id": "wf-planning",
"name": "planning",
"nodes": [
{"id": "start", "type": "start"},
{
"id": "research",
"type": "agent",
"name": "研究员",
"config": {"agentId": "default", "promptTemplate": "收集事实"},
},
{
"id": "writer",
"type": "agent",
"name": "分析师",
"config": {"agentId": "default", "promptTemplate": "综合结论"},
},
{"id": "end", "type": "output", "config": {"mapping": {"text": "{{ nodes.writer.data.text }}"}}},
],
"edges": [
{"id": "p1", "source": "start", "target": "research"},
{"id": "p2", "source": "research", "target": "writer"},
{"id": "p3", "source": "writer", "target": "end"},
],
}
FEEDBACK_GRAPH: dict[str, Any] = {
"schemaVersion": "1.0",
"id": "wf-feedback",
"name": "feedback",
"inputSchema": {"type": "object", "required": ["topic"], "properties": {"topic": {"type": "string"}}},
"nodes": [
{"id": "start", "type": "start"},
{
"id": "shape",
"type": "transform",
"config": {"operations": [{"op": "set", "target": "brief", "value": "{{ inputs.topic }} 的基础材料"}]},
},
{
"id": "writer",
"type": "agent",
"config": {
"agentId": "default",
"promptTemplate": "请基于 {{ nodes.shape.data.brief }} 写出分析。",
},
},
{"id": "end", "type": "output", "config": {"mapping": {"text": "{{ nodes.writer.data.text }}"}}},
],
"edges": [
{"id": "feedback-1", "source": "start", "target": "shape"},
{"id": "feedback-2", "source": "shape", "target": "writer"},
{"id": "feedback-3", "source": "writer", "target": "end"},
],
}
TYPED_INPUT_GRAPH: dict[str, Any] = {
"schemaVersion": "1.0",
"id": "wf-typed",
"name": "typed",
"inputSchema": {
"type": "object",
"required": ["topic"],
"properties": {"topic": {"type": "string"}, "depth": {"type": "integer"}},
},
"nodes": [
{"id": "start", "type": "start"},
{"id": "end", "type": "output", "config": {"mapping": {"topic": "{{ inputs.topic }}"}}},
],
"edges": [{"id": "e1", "source": "start", "target": "end"}],
}
class Harness:
def __init__(self, app: FastAPI, client: httpx.AsyncClient, config: WorkflowConfig) -> None:
self.app = app
self.client = client
self.config = config
@property
def runs(self) -> MemoryWorkflowRunStore:
return self.app.state.workflow_run_store
async def drive(self, rounds: int = 4) -> None:
"""Dispatch and wait for the launched run tasks, as the loop would."""
executor = self.app.state.workflow_executor
for _ in range(rounds):
await self.app.state.workflow_dispatcher.dispatch_once()
pending = [t for t in executor._tasks.values() if not t.done()]
if pending:
await asyncio.wait(pending, timeout=10)
await asyncio.sleep(0)
async def publish(self, graph: dict[str, Any], *, name: str = "demo") -> tuple[str, str]:
store = self.app.state.workflow_store
definition = await store.create_definition({"name": name, "owner_id": "user-a", "draft_graph": graph})
version = await store.publish_version(
definition["id"],
graph=graph,
graph_hash="hash-1",
published_by="user-a",
change_note="init",
)
return definition["id"], version["id"]
async def start(self, workflow_id: str, version_id: str, **body: Any) -> dict[str, Any]:
payload = {"versionId": version_id, "inputs": {"topic": "储能"}, **body}
response = await self.client.post(f"/api/workflows/{workflow_id}/runs", json=payload)
assert response.status_code == 200, response.text
return response.json()
async def events(self, run_id: str, *, after: int = 0) -> list[dict[str, Any]]:
response = await self.client.get(f"/api/workflows/runs/{run_id}/events?after={after}")
assert response.status_code == 200, response.text
return response.json()["events"]
def _sse_events(body: str) -> list[dict[str, Any]]:
return [json.loads(line[len("data: ") :]) for line in body.splitlines() if line.startswith("data: ")]
@pytest_asyncio.fixture
async def harness(monkeypatch: pytest.MonkeyPatch):
config = WorkflowConfig()
class _Cfg:
workflows = config
@staticmethod
def get_model_config(name: str):
return {"name": name} if name == "planner-model" else None
async def _user(_request=None) -> str:
return "user-a"
async def _optional(_request=None):
return None
for module in (runs_router, workflow_data_sources, workflow_embed, workflow_planning, workflow_resources):
monkeypatch.setattr(module, "get_app_config", lambda: _Cfg())
monkeypatch.setattr(module, "get_current_user", _user)
if hasattr(module, "get_optional_user_from_request"):
monkeypatch.setattr(module, "get_optional_user_from_request", _optional)
app = FastAPI()
app.state.workflow_store = MemoryWorkflowStore()
app.state.workflow_run_store = MemoryWorkflowRunStore()
app.state.workflow_planning_store = MemoryWorkflowPlanningStore()
app.state.workflow_event_store = MemoryWorkflowEventStore()
app.state.workflow_data_source_store = MemoryWorkflowDataSourceStore()
app.state.workflow_live_hub = WorkflowLiveHub()
app.state.checkpointer = None
app.state.store = None
app.state.workflow_executor = WorkflowRunExecutor(
app,
run_store=app.state.workflow_run_store,
event_store=app.state.workflow_event_store,
workflow_store=app.state.workflow_store,
data_source_store=app.state.workflow_data_source_store,
live_publisher=app.state.workflow_live_hub.publish,
config=config,
)
app.state.workflow_dispatcher = WorkflowRunDispatcher(
app.state.workflow_run_store,
app.state.workflow_executor,
event_store=app.state.workflow_event_store,
live_publisher=app.state.workflow_live_hub.publish,
config=config,
)
app.include_router(workflow_data_sources.router)
app.include_router(workflow_embed.router)
app.include_router(workflow_planning.router)
app.include_router(workflow_resources.router)
app.include_router(runs_router.router)
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
yield Harness(app, client, config)
# ── run lifecycle ───────────────────────────────────────────────────────
@pytest.mark.asyncio
async def test_run_start_executes_and_completes(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id, idempotencyKey="k-1")
assert run["status"] == "queued" and run["created"] is True
again = await harness.start(workflow_id, version_id, idempotencyKey="k-1")
assert again["id"] == run["id"] and again["created"] is False
await harness.drive()
detail = (await harness.client.get(f"/api/workflows/runs/{run['id']}")).json()
assert detail["status"] == "completed"
assert detail["output"] == {"title": "储能 报告"}
assert {n["nodeId"]: n["status"] for n in detail["nodeRuns"]} == {
"start": "completed",
"shape": "completed",
"end": "completed",
}
events = await harness.events(run["id"])
seqs = [e["seq"] for e in events]
assert seqs == sorted(seqs) and len(set(seqs)) == len(seqs)
kinds = [e["event"] for e in events]
assert kinds[0] == "run.created" and kinds[1] == "run.started"
assert kinds[-1] == "run.completed"
assert kinds.count("node.completed") == 3
@pytest.mark.asyncio
async def test_missing_required_input_is_rejected_before_queueing(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
response = await harness.client.post(f"/api/workflows/{workflow_id}/runs", json={"versionId": version_id, "inputs": {}})
assert response.status_code == 400
assert response.json()["detail"]["code"] == "WORKFLOW_INPUT_INVALID"
assert await harness.runs.list_runs(limit=10) == []
@pytest.mark.asyncio
async def test_planning_session_selects_an_editable_candidate_then_confirms_one_durable_run(
harness: Harness,
) -> None:
async def planner(**_kwargs: Any) -> dict[str, Any]:
return {
"recommendedStrategy": "parallel_research",
"titles": {"parallel_research": "多角度调研方案"},
}
harness.app.state.workflow_proposal_planner = planner
definition = await harness.app.state.workflow_store.create_definition(
{"name": "planner", "owner_id": "user-a", "draft_graph": PLANNING_GRAPH},
)
created = await harness.client.post(
f"/api/workflows/{definition['id']}/planning-sessions",
json={"query": "分析本周销售下滑并给出改进方案", "modelName": "planner-model"},
)
assert created.status_code == 200, created.text
session = created.json()
assert session["status"] == "proposed"
assert [proposal["strategy"] for proposal in session["proposals"]] == [
"parallel_research",
"collaborative",
"quick_answer",
]
assert session["proposals"][0]["title"] == "多角度调研方案"
assert all(proposal["validation"]["valid"] for proposal in session["proposals"])
for proposal in session["proposals"]:
for node in proposal["graph"]["nodes"]:
if node["type"] == "agent":
assert node["config"]["modelName"] == "planner-model"
parallel_graph = session["proposals"][0]["graph"]
evidence = next(node for node in parallel_graph["nodes"] if node["type"] == "evidence_normalizer")
synthesis = next(node for node in parallel_graph["nodes"] if node["id"] == "writer")
report = next(node for node in parallel_graph["nodes"] if node["type"] == "deep_research_write")
assert evidence["config"]["sources"] == [{"nodeId": "research", "label": "研究员"}]
assert synthesis["config"]["inputBindings"]["evidencePack"] == "{{ nodes.__parallel_evidence__.data.evidencePack }}"
assert report["config"]["evidenceBinding"] == "{{ nodes.__parallel_evidence__.data.evidencePack }}"
assert report["config"]["researchConfig"]["mode"] == "detailed"
assert [role["name"] for role in session["proposals"][0]["roles"]] == ["研究员", "分析师"]
assert session["proposals"][0]["steps"] == ["研究员", "证据归一化", "分析师", "深度研究报告写作"]
selected = session["proposals"][0]
selected_response = await harness.client.post(
f"/api/workflows/planning-sessions/{session['planningSessionId']}/proposals/{selected['proposalId']}/select",
)
assert selected_response.status_code == 200, selected_response.text
selected_session = selected_response.json()
assert selected_session["status"] == "selected"
assert selected_session["selectedProposalId"] == selected["proposalId"]
started = await harness.client.post(
f"/api/workflows/{definition['id']}/runs",
json={
"planningSessionId": session["planningSessionId"],
"proposalId": selected["proposalId"],
"idempotencyKey": "confirm-proposal-1",
},
)
assert started.status_code == 200, started.text
run = started.json()
assert run["created"] is True
assert run["executionMode"] == "planned_workflow"
assert run["planningSessionId"] == session["planningSessionId"]
# A duplicate confirmation never creates a second execution. The planning
# session is the durable source of truth, not a front-end disabled button.
duplicate = await harness.client.post(
f"/api/workflows/{definition['id']}/runs",
json={
"planningSessionId": session["planningSessionId"],
"proposalId": selected["proposalId"],
"idempotencyKey": "confirm-proposal-1",
},
)
assert duplicate.status_code == 200, duplicate.text
assert duplicate.json()["runId"] == run["runId"]
assert duplicate.json()["created"] is False
persisted = await harness.client.get(
f"/api/workflows/planning-sessions/{session['planningSessionId']}",
)
assert persisted.status_code == 200
assert persisted.json()["status"] == "confirmed"
assert persisted.json()["confirmedRunId"] == run["runId"]
@pytest.mark.asyncio
async def test_planning_session_streams_real_controller_milestones_before_result(
harness: Harness,
) -> None:
async def planner(**_kwargs: Any) -> dict[str, Any]:
return {
"recommendedStrategy": "quick_answer",
"strategyOrder": ["quick_answer", "collaborative"],
"agentSelections": {"quick_answer": ["default"], "collaborative": ["default"]},
}
harness.app.state.workflow_proposal_planner = planner
definition = await harness.app.state.workflow_store.create_definition(
{"name": "stream-planner", "owner_id": "user-a", "draft_graph": PLANNING_GRAPH},
)
async with harness.client.stream(
"POST",
f"/api/workflows/{definition['id']}/planning-sessions/stream",
json={"query": "请评估本周销售下滑原因"},
) as response:
assert response.status_code == 200
payload = "".join([chunk async for chunk in response.aiter_text()])
assert "event: progress" in payload
assert '"phase": "connection"' in payload
assert '"phase": "catalog"' in payload
assert '"phase": "controller"' in payload
assert '"phase": "selection"' in payload
assert '"phase": "assembly"' in payload
assert '"phase": "validation"' in payload
assert '"phase": "complete"' in payload
assert "event: result" in payload
assert ": stream-flush" in payload
@pytest.mark.asyncio
async def test_dedicated_workflow_controller_has_no_tools_or_thinking(
harness: Harness,
monkeypatch: pytest.MonkeyPatch,
) -> None:
captured: dict[str, Any] = {}
async def fake_run_agent(_self: Any, **kwargs: Any) -> dict[str, Any]:
captured.update(kwargs)
return {
"text": (
'{"recommendedStrategy":"quick_answer",'
'"agentSelections":{"collaborative":["default"],'
'"quick_answer":["default"]}}'
)
}
monkeypatch.setattr(WorkflowAgentRunner, "run_agent", fake_run_agent)
definition = await harness.app.state.workflow_store.create_definition(
{"name": "controller-boundary", "owner_id": "user-a", "draft_graph": PLANNING_GRAPH},
)
response = await harness.client.post(
f"/api/workflows/{definition['id']}/planning-sessions",
json={"query": "快速确认这份材料是否有缺口", "modelName": "planner-model"},
)
assert response.status_code == 200, response.text
assert captured["agent_id"] == "workflow-planner"
assert captured["disable_tools"] is True
assert captured["thinking_enabled"] is False
assert captured["force_disable_thinking"] is True
assert captured["model_name"] == "planner-model"
@pytest.mark.asyncio
async def test_planner_uses_visible_agent_catalog_when_canvas_is_only_start_and_end(
harness: Harness,
) -> None:
class VisibleAgentStore:
async def list_visible(self, _user_id: str) -> list[dict[str, str]]:
return [
{
"id": "agent-research",
"name": "行业研究员",
"description": "收集行业数据、案例和不确定性。",
},
{
"id": "agent-writer",
"name": "策略分析师",
"description": "比较证据并形成策略建议。",
},
{
"id": "workflow-planner",
"name": "工作流总控",
"description": "不应被选作正式工作流执行角色。",
},
]
async def get_extras_for(self, _agent_ids: list[str]) -> dict[str, dict[str, list[str]]]:
return {
"agent-research": {"skills": ["web-search"]},
"agent-writer": {"skills": []},
}
captured: dict[str, Any] = {}
async def planner(**kwargs: Any) -> dict[str, Any]:
captured.update(kwargs)
return {
"recommendedStrategy": "parallel_research",
"strategyOrder": ["parallel_research", "collaborative", "quick_answer"],
"agentSelections": {
"parallel_research": ["agent-research", "agent-writer"],
"collaborative": ["agent-research", "agent-writer"],
"quick_answer": ["agent-research"],
},
"taskContracts": {
"parallel_research": {
"agent-research": {
"mission": "只收集新能源乘用车市场规模与增长证据。",
"deliverable": "带来源和置信度的市场规模证据笔记。",
"scope": "不要撰写最终报告。",
"handoff": "交给策略分析师的 Evidence Pack。",
},
"agent-writer": {
"mission": "比较 Evidence Pack 中的证据并形成判断。",
"deliverable": "冲突、假设和落地建议。",
"scope": "区分事实与推断。",
"handoff": "交给深度研究报告写作节点。",
},
}
},
"titles": {"parallel_research": "双角色调研报告"},
}
harness.app.state.agent_store = VisibleAgentStore()
harness.app.state.workflow_proposal_planner = planner
resource_response = await harness.client.get("/api/workflows/resources/agents?keyword=总控")
assert resource_response.status_code == 200
assert resource_response.json() == {"agents": []}
definition = await harness.app.state.workflow_store.create_definition(
{
"name": "empty-canvas-planner",
"owner_id": "user-a",
"draft_graph": {
"schemaVersion": "1.0",
"id": "wf-empty-planning",
"nodes": [
{"id": "start", "type": "start"},
{"id": "end", "type": "output", "config": {"mapping": {}}},
],
"edges": [{"id": "start-to-end", "source": "start", "target": "end"}],
},
}
)
response = await harness.client.post(
f"/api/workflows/{definition['id']}/planning-sessions",
json={"query": "研究新能源乘用车市场,并给出落地建议"},
)
assert response.status_code == 200, response.text
assert [item["agentId"] for item in captured["agents"]] == [
"agent-research",
"agent-writer",
]
session = response.json()
assert [proposal["strategy"] for proposal in session["proposals"]] == [
"parallel_research",
"collaborative",
"quick_answer",
]
parallel = session["proposals"][0]
assert parallel["title"] == "双角色调研报告"
assert [role["agentId"] for role in parallel["roles"]] == [
"agent-research",
"agent-writer",
]
graph_nodes = {node["id"]: node for node in parallel["graph"]["nodes"]}
research_prompt = graph_nodes["__planner_agent_1__"]["config"]["promptTemplate"]
writer_prompt = graph_nodes["__planner_agent_3__"]["config"]["promptTemplate"]
assert "只收集新能源乘用车市场规模与增长证据。" in research_prompt
assert "不要撰写最终报告。" in research_prompt
assert "比较 Evidence Pack 中的证据并形成判断。" in writer_prompt
assert "这是工作流总控为你分配的固定岗位" in writer_prompt
assert all(proposal["validation"]["valid"] for proposal in session["proposals"])
@pytest.mark.asyncio
async def test_sql_planning_store_flushes_the_session_before_its_candidates() -> None:
"""A durable planning session can insert its child proposals with FK checks on."""
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
from deerflow.persistence.base import Base
from deerflow.persistence.workflow_planning import SqlWorkflowPlanningStore
engine = create_async_engine("sqlite+aiosqlite:///:memory:")
try:
async with engine.begin() as connection:
await connection.run_sync(Base.metadata.create_all)
store = SqlWorkflowPlanningStore(async_sessionmaker(engine, expire_on_commit=False))
created = await store.create_session(
{
"id": "wps-parent",
"workflow_id": "wf-parent",
"owner_id": "user-a",
"query": "生成一份市场分析",
"proposals": [
{
"id": "wpp-child",
"strategy": "collaborative",
"title": "协同分析",
"graph": {"nodes": [], "edges": []},
}
],
}
)
assert created["id"] == "wps-parent"
assert [proposal["id"] for proposal in created["proposals"]] == ["wpp-child"]
persisted = await store.get_session("wps-parent")
assert persisted is not None
assert persisted["proposals"][0]["session_id"] == "wps-parent"
finally:
await engine.dispose()
@pytest.mark.asyncio
async def test_run_start_rejects_unbound_nodes_before_queueing(harness: Harness) -> None:
graph = {
"schemaVersion": "1.0",
"id": "wf-unbound",
"nodes": [
{"id": "start", "type": "start"},
{"id": "agent", "type": "agent", "config": {"agentId": "default", "promptTemplate": ""}},
{"id": "end", "type": "output", "config": {"mapping": {"text": "{{ nodes.agent.data.text }}"}}},
],
"edges": [
{"id": "e1", "source": "start", "target": "agent"},
{"id": "e2", "source": "agent", "target": "end"},
],
}
definition = await harness.app.state.workflow_store.create_definition(
{"name": "unbound", "owner_id": "user-a", "draft_graph": graph},
)
response = await harness.client.post(
f"/api/workflows/{definition['id']}/runs",
json={"target": "draft", "inputs": {}},
)
assert response.status_code == 400
detail = response.json()["detail"]
assert detail["code"] == "WORKFLOW_RESOURCE_MISSING"
assert detail["nodeId"] == "agent"
assert await harness.runs.list_runs(limit=10) == []
@pytest.mark.asyncio
async def test_directed_agent_task_projects_one_bound_agent_into_a_durable_run(harness: Harness) -> None:
class VisibleAgentStore:
async def list_visible(self, _user_id: str) -> list[dict[str, str]]:
return [{"id": "agent-research"}]
async def get_extras_for(self, _agent_ids: list[str]) -> dict[str, dict[str, list[str]]]:
return {}
harness.app.state.agent_store = VisibleAgentStore()
graph = {
"schemaVersion": "1.0",
"id": "wf-directed",
"nodes": [
{"id": "start", "type": "start"},
{
"id": "100002",
"type": "agent",
"name": "市场研究员",
"config": {
"agentId": "agent-research",
"promptTemplate": "分析用户的市场研究任务",
},
},
{"id": "end", "type": "output", "config": {"mapping": {"text": "{{ nodes.100002.data.text }}"}}},
],
"edges": [
{"id": "e1", "source": "start", "target": "100002"},
{"id": "e2", "source": "100002", "target": "end"},
],
}
definition = await harness.app.state.workflow_store.create_definition(
{"name": "directed", "owner_id": "user-a", "draft_graph": graph},
)
response = await harness.client.post(
f"/api/workflows/{definition['id']}/runs",
json={
"target": "draft",
"input": {"query": "分析新能源行业趋势"},
"targetAgentId": "agent-research",
},
)
assert response.status_code == 200, response.text
run = response.json()
assert run["executionMode"] == "agent_task"
assert run["targetAgentId"] == "agent-research"
assert run["targetAgentName"] == "市场研究员"
stored = await harness.runs.get_run(run["runId"])
assert stored is not None
projected = stored["context"]["draftGraph"]
assert [node["id"] for node in projected["nodes"]] == [
"__agent_task_start__",
"100002",
"__agent_task_output__",
]
agent = projected["nodes"][1]
assert agent["config"]["inputBindings"] == {"workflowInput": "{{ inputs }}"}
assert projected["nodes"][2]["config"]["mapping"] == {
"text": "{{ nodes.100002.data.text }}"
}
@pytest.mark.asyncio
async def test_directed_agent_task_rejects_agents_not_bound_on_the_canvas(harness: Harness) -> None:
definition = await harness.app.state.workflow_store.create_definition(
{"name": "not-bound", "owner_id": "user-a", "draft_graph": REPORT_GRAPH},
)
response = await harness.client.post(
f"/api/workflows/{definition['id']}/runs",
json={"target": "draft", "input": {"query": "你好"}, "targetAgentId": "agent-missing"},
)
assert response.status_code == 400
assert response.json()["detail"]["code"] == "WORKFLOW_AGENT_TARGET_NOT_BOUND"
assert await harness.runs.list_runs(limit=10) == []
@pytest.mark.asyncio
async def test_draft_human_resume_preserves_the_execution_graph_snapshot(harness: Harness) -> None:
definition = await harness.app.state.workflow_store.create_definition(
{"name": "draft-pause", "owner_id": "user-a", "draft_graph": PAUSE_GRAPH},
)
response = await harness.client.post(
f"/api/workflows/{definition['id']}/runs",
json={"target": "draft", "inputs": {}},
)
assert response.status_code == 200, response.text
run_id = response.json()["runId"]
await harness.drive()
paused = await harness.runs.get_run(run_id)
assert paused is not None and paused["status"] == "awaiting_input"
assert [node["id"] for node in paused["context"]["draftGraph"]["nodes"]] == [
"start",
"ask",
"end",
]
resumed = await harness.client.post(
f"/api/workflows/runs/{run_id}/resume",
json={
"resumeToken": paused["resume_token"],
"action": "submit",
"payload": {"comment": "继续执行"},
},
)
assert resumed.status_code == 200, resumed.text
await harness.drive()
finished = await harness.runs.get_run(run_id)
assert finished is not None
assert finished["status"] == "completed"
assert finished["output"] == {"comment": "继续执行"}
assert [node["id"] for node in finished["context"]["draftGraph"]["nodes"]] == [
"start",
"ask",
"end",
]
@pytest.mark.asyncio
async def test_typed_input_schema_violations_return_structured_400(harness: Harness) -> None:
"""Wrong types (not just missing keys) fail at start with field paths."""
workflow_id, version_id = await harness.publish(TYPED_INPUT_GRAPH, name="typed")
response = await harness.client.post(
f"/api/workflows/{workflow_id}/runs",
json={"versionId": version_id, "inputs": {"topic": 123, "depth": "many"}},
)
assert response.status_code == 400
detail = response.json()["detail"]
assert detail["code"] == "WORKFLOW_INPUT_INVALID"
paths = {tuple(e["path"]) for e in detail["details"]["errors"]}
assert ("topic",) in paths and ("depth",) in paths
assert await harness.runs.list_runs(limit=10) == []
ok = await harness.client.post(
f"/api/workflows/{workflow_id}/runs",
json={"versionId": version_id, "inputs": {"topic": "储能", "depth": 2}},
)
assert ok.status_code == 200, ok.text
@pytest.mark.asyncio
async def test_unpublished_workflow_cannot_run(harness: Harness) -> None:
definition = await harness.app.state.workflow_store.create_definition({"name": "draft-only", "owner_id": "user-a", "draft_graph": REPORT_GRAPH})
response = await harness.client.post(f"/api/workflows/{definition['id']}/runs", json={"inputs": {"topic": "x"}})
assert response.status_code == 400
assert response.json()["detail"]["code"] == "WORKFLOW_VERSION_NOT_FOUND"
@pytest.mark.asyncio
async def test_events_replay_from_cursor(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id)
await harness.drive()
everything = await harness.events(run["id"])
tail = await harness.events(run["id"], after=2)
assert [e["seq"] for e in tail] == [e["seq"] for e in everything[2:]]
aliased = await harness.client.get(f"/api/workflows/runs/{run['id']}/events?after_seq=2")
assert aliased.status_code == 200
assert [e["seq"] for e in aliased.json()["events"]] == [e["seq"] for e in everything[2:]]
@pytest.mark.asyncio
async def test_stream_replays_and_closes_on_terminal_event(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id)
await harness.drive()
async with harness.client.stream("GET", f"/api/workflows/runs/{run['id']}/stream") as response:
assert response.status_code == 200
body = "".join([chunk async for chunk in response.aiter_text()])
events = _sse_events(body)
assert [e["seq"] for e in events] == sorted(e["seq"] for e in events)
assert events[-1]["event"] == "run.completed"
assert events[-1]["data"]["output"] == {"title": "储能 报告"}
@pytest.mark.asyncio
async def test_cancel_before_dispatch_reaches_terminal_state(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id)
cancelled = (await harness.client.post(f"/api/workflows/runs/{run['id']}/cancel")).json()
assert cancelled["cancelled"] is True and cancelled["status"] == "cancel_requested"
await harness.drive()
final = (await harness.client.get(f"/api/workflows/runs/{run['id']}")).json()
assert final["status"] == "cancelled"
kinds = [e["event"] for e in await harness.events(run["id"])]
assert "run.cancel_requested" in kinds and "run.cancelled" in kinds
# The graph never ran, so no node ever started.
assert not any(k.startswith("node.") for k in kinds)
repeat = (await harness.client.post(f"/api/workflows/runs/{run['id']}/cancel")).json()
assert repeat["cancelled"] is False
@pytest.mark.asyncio
async def test_human_input_pause_and_resume_round_trip(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(PAUSE_GRAPH, name="pause")
run = await harness.start(workflow_id, version_id, inputs={})
await harness.drive()
paused = (await harness.client.get(f"/api/workflows/runs/{run['id']}")).json()
assert paused["status"] == "awaiting_input"
assert paused["pendingInput"]["nodeId"] == "ask"
assert paused["pendingInput"]["prompt"] == "确认"
# The resume token travels on the event stream only, never on the run view.
assert "resumeToken" not in json.dumps(paused)
awaiting = next(e for e in await harness.events(run["id"]) if e["event"] == "run.awaiting_input")
token = awaiting["data"]["resumeToken"]
bad = await harness.client.post(
f"/api/workflows/runs/{run['id']}/resume",
json={"resumeToken": "nope", "action": "submit", "values": {"comment": "ok"}},
)
assert bad.status_code == 409
assert bad.json()["detail"]["code"] == "WORKFLOW_RESUME_TOKEN_INVALID"
resumed = await harness.client.post(
f"/api/workflows/runs/{run['id']}/resume",
json={"resumeToken": token, "action": "submit", "values": {"comment": "同意"}},
)
assert resumed.status_code == 200 and resumed.json()["status"] == "queued"
await harness.drive()
done = (await harness.client.get(f"/api/workflows/runs/{run['id']}")).json()
assert done["status"] == "completed"
assert done["output"] == {"comment": "同意"}
kinds = [e["event"] for e in await harness.events(run["id"])]
# start was replayed rather than re-executed: three nodes, three completions.
assert kinds.count("node.completed") == 3
assert "run.resumed" in kinds
@pytest.mark.asyncio
async def test_resume_is_rejected_when_not_paused(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id)
response = await harness.client.post(f"/api/workflows/runs/{run['id']}/resume", json={"resumeToken": "x", "action": "submit"})
assert response.status_code == 409
assert response.json()["detail"]["code"] == "WORKFLOW_RUN_NOT_RESUMABLE"
@pytest.mark.asyncio
async def test_concurrent_run_limit(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
harness.config.max_concurrent_runs_per_user = 1
await harness.start(workflow_id, version_id)
second = await harness.client.post(
f"/api/workflows/{workflow_id}/runs",
json={"versionId": version_id, "inputs": {"topic": "b"}},
)
assert second.status_code == 429
assert second.json()["detail"]["code"] == "WORKFLOW_LIMIT_EXCEEDED"
@pytest.mark.asyncio
async def test_other_users_run_is_not_visible(harness: Harness, monkeypatch: pytest.MonkeyPatch) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id)
async def _other(_request=None) -> str:
return "user-b"
monkeypatch.setattr(runs_router, "get_current_user", _other)
response = await harness.client.get(f"/api/workflows/runs/{run['id']}")
assert response.status_code == 403
assert response.json()["detail"]["code"] == "WORKFLOW_FORBIDDEN"
@pytest.mark.asyncio
async def test_expired_lease_is_reclaimed_by_another_worker(harness: Harness) -> None:
"""Restart recovery: a run left ``running`` by a dead worker is picked up."""
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id)
claimed = await harness.runs.claim_run(run["id"], lease_owner="dead-worker", lease_until=datetime.now(UTC) + timedelta(seconds=30))
assert claimed is not None and claimed["status"] == "running"
await harness.runs.renew_lease(run["id"], lease_owner="dead-worker", lease_until=datetime.now(UTC) - timedelta(seconds=1))
await harness.drive()
final = await harness.runs.get_run(run["id"])
assert final["status"] == "completed"
assert final["lease_owner"] is None
assert final["attempt"] == 2
@pytest.mark.asyncio
async def test_disabled_workflows_return_503(harness: Harness) -> None:
harness.config.enabled = False
response = await harness.client.get("/api/workflows/runs/anything")
assert response.status_code == 503
@pytest.mark.asyncio
async def test_retry_failed_run_replays_completed_nodes(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(FAIL_GRAPH, name="fail")
run = await harness.start(workflow_id, version_id)
await harness.drive()
failed = (await harness.client.get(f"/api/workflows/runs/{run['id']}")).json()
assert failed["status"] == "failed"
response = await harness.client.post(f"/api/workflows/runs/{run['id']}/retry")
assert response.status_code == 201, response.text
retried = response.json()
assert retried["created"] is True
assert retried["id"] != run["id"]
assert retried["retryOfRunId"] == run["id"]
assert retried["status"] == "queued"
await harness.drive()
again = (await harness.client.get(f"/api/workflows/runs/{retried['id']}")).json()
assert again["status"] == "failed"
assert again["retryOfRunId"] == run["id"]
resumed = next(e for e in await harness.events(retried["id"]) if e["event"] == "run.resumed")
assert set(resumed["data"]["replayedNodes"]) == {"start", "shape"}
original = (await harness.client.get(f"/api/workflows/runs/{run['id']}")).json()
assert original["status"] == "failed"
@pytest.mark.asyncio
async def test_retry_rejected_unless_failed(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id)
queued = await harness.client.post(f"/api/workflows/runs/{run['id']}/retry")
assert queued.status_code == 409
assert queued.json()["detail"]["code"] == "WORKFLOW_RUN_NOT_RETRYABLE"
await harness.drive()
completed = await harness.client.post(f"/api/workflows/runs/{run['id']}/retry")
assert completed.status_code == 409
assert completed.json()["detail"]["code"] == "WORKFLOW_RUN_NOT_RETRYABLE"
@pytest.mark.asyncio
async def test_completed_agent_feedback_creates_revision_and_reuses_unaffected_nodes(
harness: Harness,
monkeypatch: pytest.MonkeyPatch,
) -> None:
prompts: list[str] = []
async def fake_run_agent(_self: Any, **kwargs: Any) -> dict[str, Any]:
prompt = str(kwargs["prompt"])
prompts.append(prompt)
text = "已根据反馈重写" if "<workflow-user-feedback>" in prompt else "初版分析"
on_delta = kwargs.get("on_delta")
if on_delta is not None:
await on_delta(text)
return {"text": text, "thread_id": "fake-workflow-agent"}
monkeypatch.setattr(WorkflowAgentRunner, "run_agent", fake_run_agent)
workflow_id, version_id = await harness.publish(FEEDBACK_GRAPH, name="feedback")
original = await harness.start(workflow_id, version_id)
await harness.drive()
original_detail = (await harness.client.get(f"/api/workflows/runs/{original['id']}")).json()
assert original_detail["status"] == "completed"
assert original_detail["output"] == {"text": "初版分析"}
response = await harness.client.post(
f"/api/workflows/runs/{original['id']}/feedback",
json={
"nodeId": "writer",
"message": "请从成本、风险和可执行性三个维度重新分析。",
"idempotencyKey": "writer-revision-1",
},
)
assert response.status_code == 201, response.text
revision = response.json()
assert revision["created"] is True
assert revision["revisionOfRunId"] == original["id"]
assert revision["retryOfRunId"] == original["id"]
assert revision["affectedNodeIds"] == ["end", "writer"]
assert revision["reusedNodeIds"] == ["shape", "start"]
assert revision["feedback"]["targetNodeId"] == "writer"
assert revision["feedback"]["reusedNodeIds"] == ["shape", "start"]
before_drive = (await harness.client.get(f"/api/workflows/runs/{revision['id']}")).json()
assert {row["nodeId"] for row in before_drive["nodeRuns"]} == {"start", "shape"}
await harness.drive()
completed = (await harness.client.get(f"/api/workflows/runs/{revision['id']}")).json()
assert completed["status"] == "completed"
assert completed["output"] == {"text": "已根据反馈重写"}
assert {row["nodeId"] for row in completed["nodeRuns"]} == {"start", "shape", "writer", "end"}
assert "<workflow-user-feedback>" in prompts[-1]
assert "成本、风险和可执行性" in prompts[-1]
duplicate = await harness.client.post(
f"/api/workflows/runs/{original['id']}/feedback",
json={
"nodeId": "writer",
"message": "请从成本、风险和可执行性三个维度重新分析。",
"idempotencyKey": "writer-revision-1",
},
)
assert duplicate.status_code == 201
assert duplicate.json()["id"] == revision["id"]
assert duplicate.json()["created"] is False
assert duplicate.json()["reused"] is True
assert (await harness.client.get(f"/api/workflows/runs/{original['id']}")).json()["output"] == {"text": "初版分析"}
@pytest.mark.asyncio
async def test_stream_last_event_id_beats_after_query(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id)
await harness.drive()
everything = await harness.events(run["id"])
header_from = everything[1]["seq"]
async with harness.client.stream(
"GET",
f"/api/workflows/runs/{run['id']}/stream?after=0",
headers={"Last-Event-ID": str(header_from)},
) as response:
assert response.status_code == 200
body = "".join([chunk async for chunk in response.aiter_text()])
events = _sse_events(body)
assert events[0]["seq"] == header_from + 1
assert events[-1]["event"] == "run.completed"
assert [e["seq"] for e in events] == sorted(e["seq"] for e in events)
@pytest.mark.asyncio
async def test_sse_reconnect_storm_is_ordered_and_complete(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id)
await harness.drive()
expected = [e["seq"] for e in await harness.events(run["id"])]
async def reconnect(after: int) -> list[int]:
async with harness.client.stream(
"GET",
f"/api/workflows/runs/{run['id']}/stream?after={after}",
headers={"Last-Event-ID": str(after)},
) as response:
body = "".join([chunk async for chunk in response.aiter_text()])
events = _sse_events(body)
seqs = [e["seq"] for e in events]
assert seqs == sorted(set(seqs))
assert events[-1]["event"] == "run.completed"
return seqs
waves = await asyncio.gather(*[reconnect(after) for after in (0, 0, 1, 2, 3) for _ in range(4)])
full = [w for w in waves if w and w[0] == expected[0]]
assert full
assert all(w == expected for w in full)
@pytest.mark.asyncio
async def test_two_dispatchers_only_one_claims_a_run(harness: Harness) -> None:
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id)
rival = WorkflowRunDispatcher(
harness.app.state.workflow_run_store,
harness.app.state.workflow_executor,
event_store=harness.app.state.workflow_event_store,
live_publisher=harness.app.state.workflow_live_hub.publish,
config=harness.config,
)
started = await asyncio.gather(
harness.app.state.workflow_dispatcher.dispatch_once(),
rival.dispatch_once(),
)
assert sum(started) == 1
pending = [t for t in harness.app.state.workflow_executor._tasks.values() if not t.done()]
if pending:
await asyncio.wait(pending, timeout=10)
final = await harness.runs.get_run(run["id"])
assert final is not None and final["status"] == "completed"
@pytest.mark.asyncio
async def test_retention_cleaner_purges_old_runs_and_events(harness: Harness) -> None:
from app.gateway.workflow_retention import WorkflowRetentionCleaner
workflow_id, version_id = await harness.publish(REPORT_GRAPH)
run = await harness.start(workflow_id, version_id)
await harness.drive()
stale = datetime.now(UTC) - timedelta(days=120)
harness.runs._runs[run["id"]]["finished_at"] = stale.isoformat()
for event in harness.app.state.workflow_event_store._events.get(run["id"], []):
event["created_at"] = stale.isoformat()
cleaner = WorkflowRetentionCleaner(
harness.runs,
harness.app.state.workflow_event_store,
config=harness.config,
)
result = await cleaner.sweep()
assert result["runs"] == 1
assert result["events"] >= 1
assert await harness.runs.get_run(run["id"]) is None
assert await harness.app.state.workflow_event_store.list_after(run["id"]) == []
# ── data sources ────────────────────────────────────────────────────────
@pytest.mark.asyncio
async def test_data_source_secret_is_write_only(harness: Harness) -> None:
created = await harness.client.post(
"/api/workflows/data-sources",
json={
"name": "reporting",
"kind": "sql",
"dsn": "postgresql+asyncpg://reader:s3cret@db.internal:5432/analytics",
"maxRows": 500,
},
)
assert created.status_code == 200, created.text
body = created.json()
assert "s3cret" not in json.dumps(body) and "reader" not in json.dumps(body)
assert body["masked_target"].startswith("postgresql+asyncpg://***@")
assert body["max_rows"] == 500
listed = (await harness.client.get("/api/workflows/data-sources")).json()
assert "s3cret" not in json.dumps(listed)
assert [s["id"] for s in listed["dataSources"]] == [body["id"]]
catalog = (await harness.client.get("/api/workflows/resources/data-sources")).json()
assert catalog["dataSources"][0]["id"] == body["id"]
fetched = (await harness.client.get(f"/api/workflows/data-sources/{body['id']}")).json()
assert fetched["id"] == body["id"] and "s3cret" not in json.dumps(fetched)
introspected = (await harness.client.post(f"/api/workflows/data-sources/{body['id']}/introspect")).json()
assert introspected["tables"] == []
valid = await harness.client.post(
f"/api/workflows/data-sources/{body['id']}/validate-query",
json={"statement": "SELECT 1"},
)
assert valid.status_code == 200 and valid.json()["ok"] is True
denied = await harness.client.post(
f"/api/workflows/data-sources/{body['id']}/validate-query",
json={"statement": "DELETE FROM t"},
)
assert denied.status_code == 400
# The executor is the only reader of the plaintext.
resolved = await harness.app.state.workflow_data_source_store.resolve_dsn(body["id"])
assert resolved["dsn"].endswith("/analytics")
deleted = (await harness.client.delete(f"/api/workflows/data-sources/{body['id']}")).json()
assert deleted["deleted"] is True
@pytest.mark.asyncio
async def test_sql_data_source_requires_a_dsn(harness: Harness) -> None:
response = await harness.client.post("/api/workflows/data-sources", json={"name": "bad", "kind": "sql"})
assert response.status_code == 400
@pytest.mark.asyncio
async def test_http_data_source_keeps_headers_write_only_and_exposes_resource_metadata(
harness: Harness,
) -> None:
created = await harness.client.post(
"/api/workflows/data-sources",
json={
"name": "crm-http",
"description": "查询 CRM 客户",
"kind": "http",
"baseUrl": "https://crm.example.com/v1/customers",
"allowedMethods": ["get", "POST", "GET"],
"headers": {"Authorization": "Bearer header-secret"},
},
)
assert created.status_code == 200, created.text
body = created.json()
assert body["kind"] == "http"
assert body["base_url"] == "https://crm.example.com/v1/customers"
assert body["allowed_methods"] == ["GET", "POST"]
assert body["masked_target"] == "https://crm.example.com/v1/customers"
assert "header-secret" not in json.dumps(body)
listed = (await harness.client.get("/api/workflows/data-sources")).json()
assert "header-secret" not in json.dumps(listed)
catalog = (await harness.client.get("/api/workflows/resources/data-sources")).json()
item = catalog["dataSources"][0]
assert item["dataSourceId"] == body["id"]
assert item["baseUrl"] == "https://crm.example.com/v1/customers"
assert item["allowedMethods"] == ["GET", "POST"]
assert item["readOnly"] is False
assert "header-secret" not in json.dumps(catalog)
# The executor is the sole plaintext reader, including for HTTP headers.
resolved = await harness.app.state.workflow_data_source_store.resolve_dsn(body["id"])
assert resolved["kind"] == "http"
assert resolved["base_url"] == "https://crm.example.com/v1/customers"
assert json.loads(resolved["dsn"])["headers"]["Authorization"] == "Bearer header-secret"
@pytest.mark.asyncio
async def test_shared_data_source_requires_admin(harness: Harness) -> None:
response = await harness.client.post(
"/api/workflows/data-sources",
json={"name": "shared", "kind": "sql", "dsn": "sqlite+aiosqlite:///x.db", "shared": True},
)
assert response.status_code == 403
# ── embed ───────────────────────────────────────────────────────────────
@pytest.mark.asyncio
async def test_embed_ticket_requires_allowlisted_origin(harness: Harness) -> None:
workflow_id, _ = await harness.publish(REPORT_GRAPH)
response = await harness.client.post(
"/api/workflows/embed/tickets",
json={"workflowId": workflow_id, "origin": "https://evil.example.com"},
)
assert response.status_code == 403
assert response.json()["detail"]["code"] == "WORKFLOW_EMBED_ORIGIN_DENIED"
@pytest.mark.asyncio
async def test_embed_ticket_round_trip(harness: Harness, monkeypatch: pytest.MonkeyPatch) -> None:
monkeypatch.setenv("WORKFLOW_SECRET_KEY", "unit-test-embed-signing-key")
harness.config.embed.allowed_origins = ["https://host.example.com"]
workflow_id, _ = await harness.publish(REPORT_GRAPH)
issued = await harness.client.post(
"/api/workflows/embed/tickets",
json={"workflowId": workflow_id, "origin": "https://host.example.com/app"},
)
assert issued.status_code == 200, issued.text
assert issued.json()["frameAncestors"] == ["https://host.example.com"]
ticket = issued.json()["ticket"]
verified = await harness.client.post(
"/api/workflows/embed/verify",
json={"ticket": ticket, "origin": "https://host.example.com"},
)
assert verified.status_code == 200
assert verified.json()["workflowId"] == workflow_id
replay = await harness.client.post("/api/workflows/embed/verify", json={"ticket": ticket})
assert replay.status_code == 401
@pytest.mark.asyncio
async def test_draft_run_executes_current_draft_without_publishing(harness: Harness) -> None:
"""target=draft runs the CURRENT draft graph (version id "draft"), no publish needed."""
store = harness.app.state.workflow_store
definition = await store.create_definition({"name": "draft-run", "owner_id": "user-a", "draft_graph": REPORT_GRAPH})
response = await harness.client.post(
f"/api/workflows/{definition['id']}/runs",
json={"inputs": {"topic": "储能"}, "target": "draft", "idempotencyKey": "draft-1"},
)
assert response.status_code == 200, response.text
run = response.json()
assert run["versionId"] == "draft"
assert run["status"] in ("queued", "running")
await harness.drive()
events = await harness.events(run["runId"])
types = [event["event"] for event in events]
assert "run.started" in types
assert "run.completed" in types
final = await harness.client.get(f"/api/workflows/runs/{run['runId']}")
assert final.status_code == 200
assert final.json()["status"] == "completed"
@pytest.mark.asyncio
async def test_draft_run_without_graph_is_rejected(harness: Harness) -> None:
store = harness.app.state.workflow_store
definition = await store.create_definition(
{"name": "empty-draft", "owner_id": "user-a", "draft_graph": {"nodes": [], "edges": []}}
)
response = await harness.client.post(
f"/api/workflows/{definition['id']}/runs",
json={"inputs": {}, "target": "draft"},
)
assert response.status_code == 400
assert response.json()["detail"]["code"] == "WORKFLOW_DRAFT_EMPTY"