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