1350 lines
55 KiB
Python
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"
|